diff --git a/.changeset/a2a-request-timeout.md b/.changeset/a2a-request-timeout.md new file mode 100644 index 00000000..ee58e6f7 --- /dev/null +++ b/.changeset/a2a-request-timeout.md @@ -0,0 +1,6 @@ +--- +"@cycgraph/orchestrator": patch +"@cycgraph/a2a": patch +--- + +The A2A registry entry's `timeout_ms` is now enforced: the `a2a` node passes it as `requestTimeoutMs`, and `@cycgraph/a2a` applies it to each connection attempt and status poll on its own. A remote that stalls one of those calls now fails fast instead of holding the node for the full `task_timeout_ms`; the blocking `message/send` stays bounded by `task_timeout_ms` only, so long-running remote tasks are unaffected. diff --git a/packages/a2a/src/client.ts b/packages/a2a/src/client.ts index e16503e8..3d60f7f6 100644 --- a/packages/a2a/src/client.ts +++ b/packages/a2a/src/client.ts @@ -28,6 +28,15 @@ export interface A2AClientOptions { * polls: a remote that accepts the connection and stalls inside * `message/send` would otherwise hang the node forever. * + * The request's `requestTimeoutMs`, when set, additionally bounds EACH + * request on its own — every connection attempt (Agent Card resolution + * and client construction) and every status poll — so a remote that + * stalls one call fails fast instead of holding the node for the whole + * `timeoutMs`. A request that outruns it rejects with a transport error; + * a timed-out connection attempt counts as a failed attempt and is + * retried like any other. The `message/send` is exempt: it may block + * until the remote task finishes, so only `timeoutMs` bounds it. + * * A failed client construction is retried up to the request's * `maxRetries` times with exponential backoff (1s, 2s, 4s…, capped at * 10s), inside the same budget. `message/send` is never retried. @@ -40,12 +49,13 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { headers: Record, signal: AbortSignal, maxRetries: number, + requestTimeoutMs: number | undefined, allowedEndpointHosts?: readonly string[], ): Promise { for (let attempt = 0; ; attempt++) { try { - return await raceDeliveryBound( - () => create(agentCardUrl, headers, signal, allowedEndpointHosts), signal); + return await raceRequestBound( + () => create(agentCardUrl, headers, signal, allowedEndpointHosts), signal, requestTimeoutMs); } catch (error) { if (signal.aborted || attempt >= maxRetries) throw error; await backoff(Math.min(1_000 * 2 ** attempt, MAX_CONNECT_BACKOFF_MS), signal); @@ -61,6 +71,7 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { abortSignal?: AbortSignal, allowedEndpointHosts?: readonly string[], maxRetries = 0, + requestTimeoutMs?: number, ): Promise { const deadline = Date.now() + timeoutMs; const timeout = new AbortController(); @@ -68,10 +79,11 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { const signal = abortSignal ? AbortSignal.any([timeout.signal, abortSignal]) : timeout.signal; try { - const client = await connect(agentCardUrl, headers, signal, maxRetries, allowedEndpointHosts); + const client = await connect( + agentCardUrl, headers, signal, maxRetries, requestTimeoutMs, allowedEndpointHosts); // Cast: the generated request type demands fields the server defaults. const task = await raceDeliveryBound(() => client.sendMessage({ message } as never), signal); - return toResult(await settle(client, task, deadline, signal)); + return toResult(await settle(client, task, deadline, signal, requestTimeoutMs)); } catch (error) { // A rejection here means no task was ever observed (settle absorbs // the bound once one exists), so the bound itself is the outcome. @@ -85,6 +97,9 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { throw error; } finally { clearTimeout(timer); + // Releases any request a per-request timeout abandoned: the SDK + // client is scoped to this delivery, so nothing legitimate is left. + timeout.abort(); } } @@ -99,6 +114,7 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { request.abortSignal, request.allowedEndpointHosts, request.maxRetries, + request.requestTimeoutMs, ), /** See {@link A2AClient.resumeTask} */ @@ -111,6 +127,7 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient { request.abortSignal, request.allowedEndpointHosts, request.maxRetries, + request.requestTimeoutMs, ), }; } @@ -146,6 +163,26 @@ function raceDeliveryBound(start: () => Promise, signal: AbortSignal): Pro signal.reason instanceof Error ? signal.reason : new Error('aborted')); } +/** + * Race a call against the delivery signal and, when `requestTimeoutMs` is + * set, against a timeout of its own. The delivery bound wins when both + * fire, so a caller abort or an exhausted budget is still reported as + * such; the per-request timeout alone rejects with a transport error. + */ +function raceRequestBound( + start: () => Promise, + signal: AbortSignal, + requestTimeoutMs: number | undefined, +): Promise { + if (requestTimeoutMs === undefined) return raceDeliveryBound(start, signal); + const request = new AbortController(); + const timer = setTimeout(() => request.abort(), requestTimeoutMs); + return raceAbort(start, AbortSignal.any([signal, request.signal]), () => { + if (signal.aborted) return signal.reason instanceof Error ? signal.reason : new Error('aborted'); + return new Error(`A2A request did not respond within the ${requestTimeoutMs}ms request timeout`); + }).finally(() => clearTimeout(timer)); +} + /** * Build the outbound message. Parts use the SDK's in-memory `$case` form; * the SDK itself serializes to the wire's named-field form. A `taskId` @@ -176,6 +213,7 @@ async function settle( first: unknown, deadline: number, signal: AbortSignal, + requestTimeoutMs: number | undefined, ): Promise { let task = first as { id?: string; status?: { state?: unknown } }; let waitMs = 100; @@ -202,7 +240,8 @@ async function settle( return task; } try { - task = await raceDeliveryBound(() => client.getTask({ name: `tasks/${task.id}` } as never), signal); + task = await raceRequestBound( + () => client.getTask({ name: `tasks/${task.id}` } as never), signal, requestTimeoutMs); } catch (error) { // The bound fired while a poll was in flight: the last observed task // is still the honest answer. A non-abort rejection is a real diff --git a/packages/a2a/test/client.test.ts b/packages/a2a/test/client.test.ts index 85b166f6..50ac4cd3 100644 --- a/packages/a2a/test/client.test.ts +++ b/packages/a2a/test/client.test.ts @@ -658,6 +658,76 @@ describe('createA2AClient', () => { expect(attempts).toBe(1); }); + it('lets a blocking message/send outlive the request timeout within the task budget', async () => { + vi.useFakeTimers(); + try { + const client = createA2AClient({ + createClient: async () => ({ + sendMessage: () => new Promise((resolve) => + setTimeout(() => resolve({ id: 't', status: { state: 'completed' }, artifacts: [] }), 60_000)), + }) as never, + }); + + const pending = client.runTask({ + agentCardUrl: 'https://x/card.json', headers: {}, input: {}, timeoutMs: 600_000, requestTimeoutMs: 5_000, + }); + await vi.advanceTimersByTimeAsync(60_000); + const result = await pending; + + expect(result.state).toBe('completed'); + } finally { + vi.useRealTimers(); + } + }); + + it('fails a stalled status poll at the request timeout rather than the task budget', async () => { + vi.useFakeTimers(); + try { + const client = createA2AClient({ + createClient: async () => ({ + sendMessage: async () => ({ id: 't', status: { state: 'working' } }), + getTask: () => new Promise(() => {}), + }) as never, + }); + + const pending = client.runTask({ + agentCardUrl: 'https://x/card.json', headers: {}, input: {}, timeoutMs: 600_000, requestTimeoutMs: 5_000, + }); + const outcome = expect(pending).rejects.toThrow('did not respond within the 5000ms request timeout'); + await vi.advanceTimersByTimeAsync(5_100); + + await outcome; + } finally { + vi.useRealTimers(); + } + }); + + it('retries a connection attempt that outruns the request timeout', async () => { + vi.useFakeTimers(); + try { + let attempts = 0; + const client = createA2AClient({ + createClient: async () => { + attempts += 1; + if (attempts < 2) return new Promise(() => {}); + return { sendMessage: async () => ({ id: 't', status: { state: 'completed' }, artifacts: [] }) } as never; + }, + }); + + const pending = client.runTask({ + agentCardUrl: 'https://x/card.json', headers: {}, input: {}, timeoutMs: 600_000, + maxRetries: 1, requestTimeoutMs: 2_000, + }); + await vi.advanceTimersByTimeAsync(3_000); + const result = await pending; + + expect(attempts).toBe(2); + expect(result.state).toBe('completed'); + } finally { + vi.useRealTimers(); + } + }); + it('stops retrying a connection when the budget runs out during backoff', async () => { vi.useFakeTimers(); try { diff --git a/packages/orchestrator/src/a2a/client.ts b/packages/orchestrator/src/a2a/client.ts index 8e8e4b2c..4fa7a14d 100644 --- a/packages/orchestrator/src/a2a/client.ts +++ b/packages/orchestrator/src/a2a/client.ts @@ -89,6 +89,17 @@ export interface A2ATaskRequest { * resend could start the remote task twice. Omitted means no retries. */ maxRetries?: number; + /** + * Bound on EACH request of the delivery — every connection attempt + * (Agent Card resolution and client construction) and every status + * poll — from the registry entry's `timeout_ms`. A request that outruns + * it fails as a transport error, so a remote that stalls one call fails + * fast instead of holding the node for all of `timeoutMs`. The message + * send is exempt: it may block until the remote task finishes, so only + * `timeoutMs` bounds it. Omitted means only `timeoutMs` bounds the + * requests. + */ + requestTimeoutMs?: number; } /** @@ -141,4 +152,9 @@ export interface A2AResumeRequest { * contract as {@link A2ATaskRequest.maxRetries}. */ maxRetries?: number; + /** + * Per-request bound from the registry entry's `timeout_ms`. Same + * contract as {@link A2ATaskRequest.requestTimeoutMs}. + */ + requestTimeoutMs?: number; } diff --git a/packages/orchestrator/src/a2a/schema.ts b/packages/orchestrator/src/a2a/schema.ts index 8b20d639..06cad1ab 100644 --- a/packages/orchestrator/src/a2a/schema.ts +++ b/packages/orchestrator/src/a2a/schema.ts @@ -117,7 +117,12 @@ export const A2AServerEntrySchema = z.object({ auth: A2AAuthSchema.default({ type: 'none' }), /** Agent IDs allowed to use this server. Omit for unrestricted access. */ allowed_agents: z.array(z.string()).optional(), - /** Connection timeout in milliseconds. */ + /** + * Per-request timeout in milliseconds: bounds each connection attempt + * (Agent Card resolution and client construction) and each status poll + * on its own, within `task_timeout_ms`. The message send is exempt, since + * it may block until the remote task finishes. + */ timeout_ms: z.number().int().positive().max(3_600_000).default(30_000), /** * How long to wait for a task to reach a terminal or interrupted state. diff --git a/packages/orchestrator/src/execution/nodes/a2a.ts b/packages/orchestrator/src/execution/nodes/a2a.ts index ff52e79e..1a65f299 100644 --- a/packages/orchestrator/src/execution/nodes/a2a.ts +++ b/packages/orchestrator/src/execution/nodes/a2a.ts @@ -150,6 +150,7 @@ export async function executeA2ANode( taskId: resumingTaskId, response: ctx.state.memory.human_response, timeoutMs, + requestTimeoutMs: server.timeout_ms, maxRetries: server.max_retries, ...endpointHosts, ...(ctx.abortSignal ? { abortSignal: ctx.abortSignal } : {}), @@ -160,6 +161,7 @@ export async function executeA2ANode( input, ...(config.skill_id ? { skillId: config.skill_id } : {}), timeoutMs, + requestTimeoutMs: server.timeout_ms, maxRetries: server.max_retries, ...endpointHosts, ...(ctx.abortSignal ? { abortSignal: ctx.abortSignal } : {}), diff --git a/packages/orchestrator/test/a2a-node.test.ts b/packages/orchestrator/test/a2a-node.test.ts index 0dcd807e..e3e455d7 100644 --- a/packages/orchestrator/test/a2a-node.test.ts +++ b/packages/orchestrator/test/a2a-node.test.ts @@ -194,6 +194,23 @@ describe('executeA2ANode', () => { expect(capture.request.maxRetries).toBe(2); }); + it('forwards the entry request timeout to the client', async () => { + const capture: { request?: any } = {}; + const registry = await registryWith({ timeoutMs: 5_000 }); + + await executeA2ANode(node(), stateView(), 1, await ctxWith(fakeClient({}, capture), registry)); + + expect(capture.request.requestTimeoutMs).toBe(5_000); + }); + + it('forwards the default request timeout when the entry sets none', async () => { + const capture: { request?: any } = {}; + + await executeA2ANode(node(), stateView(), 1, await ctxWith(fakeClient({}, capture))); + + expect(capture.request.requestTimeoutMs).toBe(30_000); + }); + it('sends no allowed endpoint hosts when the entry lists none', async () => { const capture: { request?: any } = {};