diff --git a/src/subagent/run-resolved-provider-failure.test.ts b/src/subagent/run-resolved-provider-failure.test.ts index 2b5b1b67b..c2e79abb8 100644 --- a/src/subagent/run-resolved-provider-failure.test.ts +++ b/src/subagent/run-resolved-provider-failure.test.ts @@ -11,6 +11,7 @@ import { type ResolvedProviderFailureError, } from "../inference-error-message.js"; import type { InferenceErrorLike } from "../inference-gateway-error.js"; +import { MAX_BLIND_WAIT_MS } from "../agent/retry-policy.js"; import { createPermissionGate } from "../permission/gate.js"; import { createFleetMailbox, @@ -153,6 +154,109 @@ function runParams(cwd: string): RunSubAgentParams { }; } +interface RetryAttemptScript { + error?: InferenceErrorLike; + replyText?: string; + toolCalls?: string[]; +} + +function deferred(): { promise: Promise; resolve: () => void } { + let resolve!: () => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +async function withScriptedProviderRun( + scripts: RetryAttemptScript[], + callback: (run: Run, cwd: string, sendCount: () => number) => Promise, +): Promise { + const cwd = await mkdtemp(join(tmpdir(), "outer-retry-")); + let sendCount = 0; + const started = scripts.map(() => deferred()); + const eventsDone = scripts.map(() => deferred()); + try { + return await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + ({ + send: async () => { + const index = Math.min(sendCount, scripts.length - 1); + sendCount += 1; + started[index]?.resolve(); + await eventsDone[index]?.promise; + const script = scripts[index] ?? {}; + return { + type: "reply", + reply: script.replyText ?? "recovered", + turn: { role: "assistant", content: [] }, + }; + }, + stream: () => + (async function* (): AsyncGenerator { + let seq = 1; + for (const [index, script] of scripts.entries()) { + await started[index]?.promise; + for (const name of script.toolCalls ?? []) { + yield { + type: "tool.start", + seq: seq++, + data: { call: { name, arguments: {} } }, + } as unknown as ReactorEmittedEvent; + } + if (script.error !== undefined) { + yield { + type: "inference.error", + seq: seq++, + data: { + error: script.error, + partial: { text: "" }, + }, + } as unknown as ReactorEmittedEvent; + } else { + yield { + type: "inference.start", + seq: seq++, + data: { + sourceId: "test", + model: "test-model", + input: [], + }, + } as unknown as ReactorEmittedEvent; + } + yield { + type: "connector.reply", + seq: seq++, + data: { content: script.replyText ?? "" }, + } as unknown as ReactorEmittedEvent; + eventsDone[index]?.resolve(); + } + })(), + deliver: () => undefined, + close: async () => undefined, + setSource: () => undefined, + setSources: () => undefined, + history: async () => [], + checkpoints: async () => [], + readAt: async () => [], + blobReader: {}, + }) as unknown as Awaited< + ReturnType + >, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + return callback(runSubAgent, cwd, () => sendCount); + }, + ); + } finally { + await rm(cwd, { recursive: true, force: true }); + } +} + describe("resolved sub-agent provider failures", () => { test("runSubAgent rejects a raw director reply after inference.error", async () => { const { caught, observed } = await withResolvedProviderRun( @@ -336,4 +440,309 @@ describe("resolved sub-agent provider failures", () => { ); }, ); + + test("outer retry recovers when a retryable failure clears on the 2nd send", async () => { + const progress: { description: string; toolName: string }[] = []; + const { result, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 502, + }, + }, + { replyText: "## Summary\nrecovered" }, + ], + async (run, cwd, sendCount) => ({ + result: await run({ + ...runParams(cwd), + onProgress: (info) => { + progress.push(info); + }, + }), + sends: sendCount(), + }), + ); + + expect(sends).toBe(2); + expect(result.report).toContain("recovered"); + const retryNotice = progress.find((info) => info.toolName === "retry"); + expect(retryNotice?.description).toContain("2/2"); + }); + + test("outer retry does not retry a fatal failure", async () => { + const { caught, sends } = await withScriptedProviderRun( + [{ error: { category: "fatal", message: RAW_DIAGNOSTIC } }], + async (run, cwd, sendCount) => { + try { + await run(runParams(cwd)); + } catch (error) { + return { caught: error, sends: sendCount() }; + } + throw new Error("expected runSubAgent to reject"); + }, + ); + + expect(sends).toBe(1); + expect(isResolvedProviderFailureError(caught)).toBe(true); + expect((caught as ResolvedProviderFailureError).category).toBe("fatal"); + }); + + test("outer retry rethrows the last error unchanged once exhausted", async () => { + const first = { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 503, + } satisfies InferenceErrorLike; + const second = { + category: "retryable", + message: "still failing", + statusCode: 503, + } satisfies InferenceErrorLike; + const { caught, sends } = await withScriptedProviderRun( + [{ error: first }, { error: second }], + async (run, cwd, sendCount) => { + try { + await run(runParams(cwd)); + } catch (error) { + return { caught: error, sends: sendCount() }; + } + throw new Error("expected runSubAgent to reject"); + }, + ); + + expect(sends).toBe(2); + expect(isResolvedProviderFailureError(caught)).toBe(true); + const resolved = caught as ResolvedProviderFailureError; + expect(resolved.category).toBe("retryable"); + expect(resolved.statusCode).toBe(503); + expect(resolved.message).toBe( + "test-provider Provider failed (retryable). Try again.", + ); + }); + + test("outer retry does not retry after a tool already executed", async () => { + const { caught, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 502, + }, + toolCalls: ["read_file"], + }, + ], + async (run, cwd, sendCount) => { + try { + await run(runParams(cwd)); + } catch (error) { + return { caught: error, sends: sendCount() }; + } + throw new Error("expected runSubAgent to reject"); + }, + ); + + expect(sends).toBe(1); + expect(isResolvedProviderFailureError(caught)).toBe(true); + expect((caught as ResolvedProviderFailureError).category).toBe("retryable"); + }); + + test("interrupt during backoff sleep salvages as interrupted, not a provider failure", async () => { + let interrupt: (() => void) | undefined; + let sawRetryNotice = false; + const { result, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 502, + }, + }, + ], + async (run, cwd, sendCount) => ({ + result: await run({ + ...runParams(cwd), + onAgentReady: (handles) => { + interrupt = handles.interrupt; + }, + onProgress: (info) => { + if (info.toolName !== "retry") return; + sawRetryNotice = true; + interrupt?.(); + }, + }), + sends: sendCount(), + }), + ); + + expect(sawRetryNotice).toBe(true); + expect(sends).toBe(1); + expect(result.stopReason).toBe("interrupted"); + expect(result.interrupted).toBe(true); + }); + + test("outer retry honors a short 429 retryAfterMs on the 2nd send", async () => { + let retryDelayMs: number | undefined; + const { result, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 429, + retryAfterMs: 10, + }, + }, + { replyText: "## Summary\nrecovered" }, + ], + async (run, cwd, sendCount) => ({ + result: await run({ + ...runParams(cwd), + onProgress: (info) => { + if (info.toolName !== "retry") return; + const match = info.description.match(/in (\d+)ms/); + if (match !== null) retryDelayMs = Number(match[1]); + }, + }), + sends: sendCount(), + }), + ); + + expect(retryDelayMs).toBe(10); + expect(sends).toBe(2); + expect(result.report).toContain("recovered"); + }); + + test("outer retry caps a long 429 retryAfterMs at MAX_BLIND_WAIT_MS", async () => { + let interrupt: (() => void) | undefined; + let retryDelayMs: number | undefined; + const { result, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 429, + retryAfterMs: 120_000, + }, + }, + ], + async (run, cwd, sendCount) => ({ + result: await run({ + ...runParams(cwd), + onAgentReady: (handles) => { + interrupt = handles.interrupt; + }, + onProgress: (info) => { + if (info.toolName !== "retry") return; + const match = info.description.match(/in (\d+)ms/); + if (match !== null) retryDelayMs = Number(match[1]); + interrupt?.(); + }, + }), + sends: sendCount(), + }), + ); + + expect(retryDelayMs).toBeDefined(); + expect(retryDelayMs as number).toBeLessThanOrEqual(MAX_BLIND_WAIT_MS); + expect(sends).toBe(1); + expect(result.stopReason).toBe("interrupted"); + }); + + test("outer retry does not retry a raw send rejection without inference.error", async () => { + const cwd = await mkdtemp(join(tmpdir(), "outer-retry-raw-")); + let sends = 0; + const rawError = new Error("raw send boom"); + try { + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + ({ + send: async () => { + sends += 1; + throw rawError; + }, + stream: () => + (async function* (): AsyncGenerator { + yield { + type: "inference.start", + seq: 1, + data: { sourceId: "test", model: "test-model", input: [] }, + } as unknown as ReactorEmittedEvent; + yield { + type: "connector.reply", + seq: 2, + data: { content: "" }, + } as unknown as ReactorEmittedEvent; + })(), + deliver: () => undefined, + close: async () => undefined, + setSource: () => undefined, + setSources: () => undefined, + history: async () => [], + checkpoints: async () => [], + readAt: async () => [], + blobReader: {}, + }) as unknown as Awaited< + ReturnType + >, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + try { + await runSubAgent(runParams(cwd)); + } catch (error) { + expect(sends).toBe(1); + expect(error).toBe(rawError); + expect(isResolvedProviderFailureError(error)).toBe(false); + return; + } + throw new Error("expected runSubAgent to reject"); + }, + ); + } finally { + await rm(cwd, { recursive: true, force: true }); + } + }); + + test("outer retry exhaustion surfaces the second attempt's fields", async () => { + const { caught, sends } = await withScriptedProviderRun( + [ + { + error: { + category: "retryable", + message: RAW_DIAGNOSTIC, + statusCode: 502, + }, + }, + { + error: { + category: "retryable", + message: "still failing", + statusCode: 504, + }, + }, + ], + async (run, cwd, sendCount) => { + try { + await run(runParams(cwd)); + } catch (error) { + return { caught: error, sends: sendCount() }; + } + throw new Error("expected runSubAgent to reject"); + }, + ); + + expect(sends).toBe(2); + expect(isResolvedProviderFailureError(caught)).toBe(true); + const resolved = caught as ResolvedProviderFailureError; + expect(resolved.category).toBe("retryable"); + expect(resolved.statusCode).toBe(504); + }); }); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index 88916c5b8..18a8cb82b 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -93,8 +93,12 @@ import { consumeStream } from "../session/stream-consumer.js"; import { createCycleTextRecorder } from "../session/stream-journal.js"; import { onTurnBoundary } from "../agent/reactor-events.js"; import { refreshInferenceSourceBundle } from "./refresh-inference-source.js"; -import { createResolvedProviderFailureError } from "../inference-error-message.js"; +import { + createResolvedProviderFailureError, + isResolvedProviderFailureError, +} from "../inference-error-message.js"; import type { InferenceErrorLike } from "../inference-gateway-error.js"; +import { MAX_BLIND_WAIT_MS } from "../agent/retry-policy.js"; import { createRunEventSettlement } from "./run-event-settlement.js"; import type { CapabilityFilter } from "../agent/profiles.js"; @@ -199,6 +203,41 @@ export function assertReplySend( throw error; } +// Bounded outer retry for the main send (CL-7677). The harness policy in +// vendor/intx-inference/src/retry-policy.ts already retries retryable faults +// up to 3 times per send; this loop covers the case where the harness gives +// up and the failure terminalizes the worker anyway. A terminal worker death +// is not a non-terminal director turn (CL-6910), so the outer budget is one +// retry, not a second full schedule — harness-weighted worst case is 2 outer +// x 3 inner sends. The single delay matches the first step of the harness +// backoff (500ms before attempt 2 in vendor/intx-inference/src/retry-policy.ts) +// with the same jitter applied below: the outer budget is one retry, so only +// the first backoff step is ever reached — a table would be dead weight. +// createCorbitsRetryPolicy decides per-attempt retry inside a live send and +// must not drive this loop. +// Retrying after a tool already executed is unsafe — the second send would +// replay side effects — so any recorded tool start vetoes the retry. +const MAX_OUTER_ATTEMPTS = 2; +const OUTER_RETRY_DELAY_MS = 500; + +function sleepUnlessAborted( + ms: number, + signals: readonly AbortSignal[], +): Promise { + if (ms <= 0 || signals.some((signal) => signal.aborted)) + return Promise.resolve(); + return new Promise((resolve) => { + const timer = setTimeout(done, ms); + function done(): void { + clearTimeout(timer); + for (const signal of signals) signal.removeEventListener("abort", done); + resolve(); + } + for (const signal of signals) + signal.addEventListener("abort", done, { once: true }); + }); +} + export type { NestedDispatchDeps, RunSubAgentParams, @@ -1136,6 +1175,9 @@ async function runSubAgentInner( ...(error.statusCode !== undefined ? { statusCode: error.statusCode } : {}), + ...(error.retryAfterMs !== undefined + ? { retryAfterMs: error.retryAfterMs } + : {}), }; } if (onTurnBoundary(event)) { @@ -1365,18 +1407,72 @@ async function runSubAgentInner( // signal so either one stops this send() call, while only // runController's abort is wired to closeOnAbort/teardown. const sendOpts = { signal: sendAbortSignal() }; - const fresh = await refreshInferenceSourceBundle( - bundle.sources, - bundle.defaultSource, - params.catalog, - ); - agent.setSources(fresh.sources, fresh.defaultSource); - const result = await sendWithProviderFailure(fullPrompt, sendOpts); - if (terminalProviderError !== undefined) { - throw createResolvedProviderFailureError( - params.provider.providerName, - terminalProviderError, + const outerStartedAt = Date.now(); + let attempt = 0; + let result: Awaited>; + for (;;) { + attempt += 1; + const fresh = await refreshInferenceSourceBundle( + bundle.sources, + bundle.defaultSource, + params.catalog, ); + agent.setSources(fresh.sources, fresh.defaultSource); + try { + result = await sendWithProviderFailure(fullPrompt, sendOpts); + if (terminalProviderError !== undefined) { + throw createResolvedProviderFailureError( + params.provider.providerName, + terminalProviderError, + ); + } + } catch (sendError) { + const resolved = isResolvedProviderFailureError(sendError) + ? sendError + : undefined; + const toolsUsed = + toolNamesUsed.length > 0 || telemetryRollup.tool_call_count > 0; + if ( + resolved === undefined || + resolved.category !== "retryable" || + toolsUsed || + attempt >= MAX_OUTER_ATTEMPTS + ) { + throw sendError; + } + const elapsedMs = Date.now() - outerStartedAt; + const outerRemainingMs = MAX_BLIND_WAIT_MS - elapsedMs; + if (outerRemainingMs <= 0) throw sendError; + if (runController.signal.aborted || runController.deadlineHit()) + ensureNotAborted(); + if (thisTurnInterrupt.signal.aborted) + throw abortError(thisTurnInterrupt.signal); + const retryAfterMs = + terminalProviderError?.statusCode === 429 + ? terminalProviderError.retryAfterMs + : undefined; + const backoffMs = + retryAfterMs !== undefined + ? Math.min(retryAfterMs, MAX_BLIND_WAIT_MS) + : Math.round(OUTER_RETRY_DELAY_MS * (0.8 + Math.random() * 0.4)); + const delayMs = Math.min(backoffMs, outerRemainingMs); + params.onProgress?.({ + description: + `${params.description} (provider retry ` + + `${attempt + 1}/${MAX_OUTER_ATTEMPTS} in ${delayMs}ms)`, + toolName: "retry", + }); + await sleepUnlessAborted(delayMs, [ + runController.signal, + thisTurnInterrupt.signal, + ]); + if (runController.signal.aborted || runController.deadlineHit()) + ensureNotAborted(); + if (thisTurnInterrupt.signal.aborted) + throw abortError(thisTurnInterrupt.signal); + continue; + } + break; } if (result.type !== "reply") assertReplySend(result); // A successful non-empty reply must not be clobbered by a late cancel that