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
6 changes: 6 additions & 0 deletions .changeset/a2a-request-timeout.md
Original file line number Diff line number Diff line change
@@ -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.
49 changes: 44 additions & 5 deletions packages/a2a/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -40,12 +49,13 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient {
headers: Record<string, string>,
signal: AbortSignal,
maxRetries: number,
requestTimeoutMs: number | undefined,
allowedEndpointHosts?: readonly string[],
): Promise<SdkClient> {
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);
Expand All @@ -61,17 +71,19 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient {
abortSignal?: AbortSignal,
allowedEndpointHosts?: readonly string[],
maxRetries = 0,
requestTimeoutMs?: number,
): Promise<A2ATaskResult> {
const deadline = Date.now() + timeoutMs;
const timeout = new AbortController();
const timer = setTimeout(() => timeout.abort(), timeoutMs);
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.
Expand All @@ -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();
}
}

Expand All @@ -99,6 +114,7 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient {
request.abortSignal,
request.allowedEndpointHosts,
request.maxRetries,
request.requestTimeoutMs,
),

/** See {@link A2AClient.resumeTask} */
Expand All @@ -111,6 +127,7 @@ export function createA2AClient(options: A2AClientOptions = {}): A2AClient {
request.abortSignal,
request.allowedEndpointHosts,
request.maxRetries,
request.requestTimeoutMs,
),
};
}
Expand Down Expand Up @@ -146,6 +163,26 @@ function raceDeliveryBound<T>(start: () => Promise<T>, 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<T>(
start: () => Promise<T>,
signal: AbortSignal,
requestTimeoutMs: number | undefined,
): Promise<T> {
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`
Expand Down Expand Up @@ -176,6 +213,7 @@ async function settle(
first: unknown,
deadline: number,
signal: AbortSignal,
requestTimeoutMs: number | undefined,
): Promise<unknown> {
let task = first as { id?: string; status?: { state?: unknown } };
let waitMs = 100;
Expand All @@ -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
Expand Down
70 changes: 70 additions & 0 deletions packages/a2a/test/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
16 changes: 16 additions & 0 deletions packages/orchestrator/src/a2a/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down Expand Up @@ -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;
}
7 changes: 6 additions & 1 deletion packages/orchestrator/src/a2a/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 2 additions & 0 deletions packages/orchestrator/src/execution/nodes/a2a.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 } : {}),
Expand All @@ -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 } : {}),
Expand Down
17 changes: 17 additions & 0 deletions packages/orchestrator/test/a2a-node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 } = {};

Expand Down
Loading