Skip to content
Closed
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
20 changes: 12 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 { readCanvas } from '../canvas/space-read.js';

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

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

interface AgentNodeServiceDependencies {
getProfileRegistry: () => AgentProfileRegistryPort | null;
readCanvasNodes: (canvasId: string) => StoredNode[] | null;
readCanvasNodes: (
canvasId: string,
) => Promise<StoredNode[] | null> | 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 readCanvas(canvasId);
if (!canvas) return null;
return canvas.state.nodes as StoredNode[];
}
Expand Down Expand Up @@ -153,11 +157,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 readCanvas(canvasId);
if (!canvas) {
throw new AgentNodeCreationError(
'canvas_not_found',
Expand Down Expand Up @@ -201,7 +205,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
37 changes: 22 additions & 15 deletions apps/server/src/modules/agent/agent-thread-resolver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,39 @@ 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(
createResolver([selectable]).resolveFixedAgentNode(
await createResolver([selectable]).resolveFixedAgentNode(
'canvas-a',
'thread-a',
),
).toBeNull();
expect(
createResolver([FIXED_NODE]).resolveFixedAgentNode(
await createResolver([FIXED_NODE]).resolveFixedAgentNode(
'canvas-a',
'thread-other',
),
).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(
createResolver([selectable]).resolveAgentNodeId('canvas-a', 'thread-a'),
await createResolver([selectable]).resolveAgentNodeId(
'canvas-a',
'thread-a',
),
).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 +106,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 +131,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
34 changes: 21 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 { readCanvas, readCanvasNode } from '../canvas/space-read.js';

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

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

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

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

/**
Expand All @@ -74,8 +79,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 +101,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 +153,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
8 changes: 4 additions & 4 deletions apps/server/src/modules/agent/agent-thread.service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ 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',
Expand All @@ -173,11 +173,11 @@ describe('AgentThreadService', () => {
const harness = createHarness({ target: null, persistedBinding: binding });

expect(
harness.service.resolveExternalTarget('canvas-a', 'thread-a'),
await harness.service.resolveExternalTarget('canvas-a', 'thread-a'),
).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 @@ -187,7 +187,7 @@ describe('AgentThreadService', () => {
});

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

Expand Down
16 changes: 8 additions & 8 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> | 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)
? await 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