Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/server/src/modules/agent/agent-node.service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ function createHarness(options?: {
: null,
listSelectableProfileIds: () => options?.selectableIds ?? ['profile-a'],
}),
readCanvasNodes: () =>
readCanvasNodes: async () =>
options?.nodes === undefined
? [
{ id: NOTE_ID, type: 'note' },
Expand Down
18 changes: 10 additions & 8 deletions apps/server/src/modules/agent/agent-node.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ import {
import { getLogger } from '../../utils/logger.js';
import { executeCanvasCommandsOnHost } from '../canvas/canvas-command-router.js';
import { buildSpatialBundle } from '../canvas/canvas-spatial.js';
import { getCanvasStore } from '../storage/index.js';
import { space } from '../storage/index.js';

import type { ExecuteOnServerOutput } from '../canvas/canvas-executor.js';

Expand Down Expand Up @@ -87,16 +87,18 @@ interface StoredNode {

interface AgentNodeServiceDependencies {
getProfileRegistry: () => AgentProfileRegistryPort | null;
readCanvasNodes: (canvasId: string) => StoredNode[] | null;
readCanvasNodes: (canvasId: string) => Promise<StoredNode[] | null>;
execute: (input: {
canvasId: string;
commands: readonly CanvasCommand[];
originator: { source: 'system' };
}) => Promise<ExecuteOnServerOutput>;
}

function defaultReadCanvasNodes(canvasId: string): StoredNode[] | null {
const canvas = getCanvasStore(canvasId).read();
async function defaultReadCanvasNodes(
canvasId: string,
): Promise<StoredNode[] | null> {
const canvas = await space(canvasId).read();
if (!canvas) return null;
return canvas.state.nodes as StoredNode[];
}
Expand Down Expand Up @@ -153,11 +155,11 @@ function resolveAnchor(
return nodeId;
}

export function resolveAgentNodePosition(
export async function resolveAgentNodePosition(
canvasId: string,
parentNodeId?: CanvasNodeId,
): Point {
const canvas = getCanvasStore(canvasId).read();
): Promise<Point> {
const canvas = await space(canvasId).read();
if (!canvas) {
throw new AgentNodeCreationError(
'canvas_not_found',
Expand Down Expand Up @@ -201,7 +203,7 @@ export class AgentNodeService {
throw error;
}

const nodes = this.dependencies.readCanvasNodes(input.canvasId);
const nodes = await this.dependencies.readCanvasNodes(input.canvasId);
if (!nodes) {
throw new AgentNodeCreationError(
'canvas_not_found',
Expand Down
44 changes: 24 additions & 20 deletions apps/server/src/modules/agent/agent-thread-resolver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ function createResolver(
content = '',
) {
return new AgentThreadResolver({
readCanvasNodes: () => nodes,
readNodeContent: () => content,
readCanvasNodes: async () => nodes,
readNodeContent: async () => content,
});
}

Expand All @@ -41,8 +41,8 @@ const FIXED_NODE = {
};

describe('AgentThreadResolver', () => {
it('resolves a fixed external Agent Node from Canvas storage', () => {
const target = createResolver(
it('resolves a fixed external Agent Node from Canvas storage', async () => {
const target = await createResolver(
[FIXED_NODE],
'Existing prompt',
).resolveFixedAgentNode('canvas-a', 'thread-a');
Expand All @@ -65,36 +65,36 @@ describe('AgentThreadResolver', () => {
});
});

it('falls back for selectable or unrelated threads', () => {
it('falls back for selectable or unrelated threads', async () => {
const selectable = {
...FIXED_NODE,
data: { ...FIXED_NODE.data, agentBindingPolicy: 'selectable' },
};
expect(
await expect(
createResolver([selectable]).resolveFixedAgentNode(
'canvas-a',
'thread-a',
),
).toBeNull();
expect(
).resolves.toBeNull();
await expect(
createResolver([FIXED_NODE]).resolveFixedAgentNode(
'canvas-a',
'thread-other',
),
).toBeNull();
).resolves.toBeNull();
});

it('resolves any Question Node as a possible parent', () => {
it('resolves any Question Node as a possible parent', async () => {
const selectable = {
...FIXED_NODE,
data: { ...FIXED_NODE.data, agentBindingPolicy: 'selectable' },
};
expect(
await expect(
createResolver([selectable]).resolveAgentNodeId('canvas-a', 'thread-a'),
).toBe('node-agent');
).resolves.toBe('node-agent');
});

it('resolves a Huabu Agent binding for invocation', () => {
it('resolves a Huabu Agent binding for invocation', async () => {
const internal = {
...FIXED_NODE,
data: {
Expand All @@ -103,19 +103,23 @@ describe('AgentThreadResolver', () => {
},
};
expect(
createResolver([internal]).resolveFixedAgentNode('canvas-a', 'thread-a')
?.agentBinding,
(
await createResolver([internal]).resolveFixedAgentNode(
'canvas-a',
'thread-a',
)
)?.agentBinding,
).toEqual({ kind: 'internal' });
});

it('rejects duplicate threads and corrupt fixed-node metadata', () => {
it('rejects duplicate threads and corrupt fixed-node metadata', async () => {
const duplicateResolver = createResolver([
FIXED_NODE,
{ ...FIXED_NODE, id: 'node-agent-2' },
]);
expect(() =>
await expect(
duplicateResolver.resolveFixedAgentNode('canvas-a', 'thread-a'),
).toThrowError(
).rejects.toThrowError(
expect.objectContaining<Partial<AgentThreadResolutionError>>({
code: 'duplicate_thread',
}),
Expand All @@ -124,9 +128,9 @@ describe('AgentThreadResolver', () => {
const corruptResolver = createResolver([
{ ...FIXED_NODE, data: { ...FIXED_NODE.data, agentBinding: null } },
]);
expect(() =>
await expect(
corruptResolver.resolveFixedAgentNode('canvas-a', 'thread-a'),
).toThrowError(
).rejects.toThrowError(
expect.objectContaining<Partial<AgentThreadResolutionError>>({
code: 'invalid_binding',
}),
Expand Down
29 changes: 16 additions & 13 deletions apps/server/src/modules/agent/agent-thread-resolver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import {
InvalidAgentLaunchOverridesError,
parseAgentLaunchOverrides,
} from './agent-launch-overrides.js';
import { getCanvasStore } from '../storage/index.js';
import { space } from '../storage/index.js';

interface StoredNode {
id: string;
Expand All @@ -23,8 +23,8 @@ interface StoredNode {
}

interface ResolverDependencies {
readCanvasNodes: (canvasId: string) => StoredNode[] | null;
readNodeContent: (canvasId: string, nodeId: string) => string | null;
readCanvasNodes: (canvasId: string) => Promise<StoredNode[] | null>;
readNodeContent: (canvasId: string, nodeId: string) => Promise<string | null>;
}

export interface FixedAgentNodeTarget {
Expand Down Expand Up @@ -55,12 +55,12 @@ export class AgentThreadResolutionError extends Error {
}

const DEFAULT_DEPENDENCIES: ResolverDependencies = {
readCanvasNodes: (canvasId) => {
const canvas = getCanvasStore(canvasId).read();
readCanvasNodes: async (canvasId) => {
const canvas = await space(canvasId).read();
return canvas ? (canvas.state.nodes as StoredNode[]) : null;
},
readNodeContent: (canvasId, nodeId) =>
getCanvasStore(canvasId).readNode(nodeId)?.content ?? null,
readNodeContent: async (canvasId, nodeId) =>
(await space(canvasId).nodes.read(nodeId))?.record.content ?? null,
};

/**
Expand All @@ -74,8 +74,11 @@ export class AgentThreadResolver {
private readonly dependencies: ResolverDependencies = DEFAULT_DEPENDENCIES,
) {}

resolveAgentNodeId(canvasId: string, threadId: string): CanvasNodeId | null {
const nodes = this.dependencies.readCanvasNodes(canvasId);
async resolveAgentNodeId(
canvasId: string,
threadId: string,
): Promise<CanvasNodeId | null> {
const nodes = await this.dependencies.readCanvasNodes(canvasId);
if (!nodes) {
throw new AgentThreadResolutionError(
'canvas_not_found',
Expand All @@ -93,11 +96,11 @@ export class AgentThreadResolver {
return node?.type === 'question' ? (node.id as CanvasNodeId) : null;
}

resolveFixedAgentNode(
async resolveFixedAgentNode(
canvasId: string,
threadId: string,
): FixedAgentNodeTarget | null {
const nodes = this.dependencies.readCanvasNodes(canvasId);
): Promise<FixedAgentNodeTarget | null> {
const nodes = await this.dependencies.readCanvasNodes(canvasId);
if (!nodes) {
throw new AgentThreadResolutionError(
'canvas_not_found',
Expand Down Expand Up @@ -145,7 +148,7 @@ export class AgentThreadResolver {
throw error;
}

const content = this.dependencies.readNodeContent(canvasId, node.id);
const content = await this.dependencies.readNodeContent(canvasId, node.id);
if (content === null) {
throw new AgentThreadResolutionError(
'missing_node_content',
Expand Down
16 changes: 8 additions & 8 deletions apps/server/src/modules/agent/agent-thread.service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ function createHarness(options?: {
return emptyInternalStream();
});
const service = new AgentThreadService({
resolveFixedAgentNode: () =>
resolveFixedAgentNode: async () =>
options && 'target' in options ? (options.target ?? null) : TARGET,
resolvePersistedExternalBinding: () =>
options && 'persistedBinding' in options
Expand Down Expand Up @@ -164,20 +164,20 @@ describe('AgentThreadService', () => {
expect(externalBindingFromWorkloadSpec({ binding: {} })).toBeNull();
});

it('resolves a persisted external Thread without a fixed Agent Node', () => {
it('resolves a persisted external Thread without a fixed Agent Node', async () => {
const binding = {
kind: 'external' as const,
profileId: 'profile-selectable',
alias: 'Selectable Agent',
};
const harness = createHarness({ target: null, persistedBinding: binding });

expect(
await expect(
harness.service.resolveExternalTarget('canvas-a', 'thread-a'),
).toEqual({ binding, fixedTarget: null });
).resolves.toEqual({ binding, fixedTarget: null });
});

it('prefers the fixed Agent Node binding when one exists', () => {
it('prefers the fixed Agent Node binding when one exists', async () => {
const harness = createHarness({
persistedBinding: {
kind: 'external',
Expand All @@ -186,9 +186,9 @@ describe('AgentThreadService', () => {
},
});

expect(
await expect(
harness.service.resolveExternalTarget('canvas-a', 'thread-a'),
).toEqual({ binding: TARGET.agentBinding, fixedTarget: TARGET });
).resolves.toEqual({ binding: TARGET.agentBinding, fixedTarget: TARGET });
});

it('uses persisted fixed binding and overrides under one leased lifecycle', async () => {
Expand Down Expand Up @@ -356,7 +356,7 @@ describe('AgentThreadService', () => {
return events([{ type: 'done', data: { message: 'Done' } }]);
});
const service = new AgentThreadService({
resolveFixedAgentNode: () => TARGET,
resolveFixedAgentNode: async () => TARGET,
resolvePersistedExternalBinding: () => null,
waitForTurnRelease: vi.fn().mockResolvedValue(undefined),
acquireTurn: vi.fn(() => vi.fn()),
Expand Down
14 changes: 7 additions & 7 deletions apps/server/src/modules/agent/agent-thread.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ interface AgentThreadServiceDependencies {
resolveFixedAgentNode: (
canvasId: string,
threadId: string,
) => FixedAgentNodeTarget | null;
) => Promise<FixedAgentNodeTarget | null>;
resolvePersistedExternalBinding: (
canvasId: string,
threadId: string,
Expand Down Expand Up @@ -160,20 +160,20 @@ export class AgentThreadService {
private readonly dependencies: AgentThreadServiceDependencies = DEFAULT_DEPENDENCIES,
) {}

resolveFixedTarget(
async resolveFixedTarget(
canvasId: string | undefined,
threadId: string,
): FixedAgentNodeTarget | null {
): Promise<FixedAgentNodeTarget | null> {
return canvasId
? this.dependencies.resolveFixedAgentNode(canvasId, threadId)
: null;
}

resolveExternalTarget(
async resolveExternalTarget(
canvasId: string,
threadId: string,
): ExternalAgentThreadTarget | null {
const fixedTarget = this.resolveFixedTarget(canvasId, threadId);
): Promise<ExternalAgentThreadTarget | null> {
const fixedTarget = await this.resolveFixedTarget(canvasId, threadId);
if (fixedTarget) {
return fixedTarget.agentBinding.kind === 'external'
? { binding: fixedTarget.agentBinding, fixedTarget }
Expand All @@ -197,7 +197,7 @@ export class AgentThreadService {
): Promise<AgentThreadInvocation> {
const fixedTarget =
options.fixedTarget === undefined
? this.resolveFixedTarget(options.canvasId, options.threadId)
? await this.resolveFixedTarget(options.canvasId, options.threadId)
: options.fixedTarget;
const binding: AgentBinding = fixedTarget?.agentBinding ??
options.requestBinding ?? { kind: 'internal' };
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/modules/agent/agent.route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -546,7 +546,7 @@ const agentRoutes: FastifyPluginAsync = async (
} = parsed.data;

const resolvedThreadId = getOrCreateThreadId(threadId);
const fixedTarget = agentThreadService.resolveFixedTarget(
const fixedTarget = await agentThreadService.resolveFixedTarget(
canvasId,
resolvedThreadId,
);
Expand Down
Loading
Loading