From 42b4ea937f52117cc8bc8b91bdcc01372f41873e Mon Sep 17 00:00:00 2001 From: Sawyer Date: Thu, 10 Sep 2026 18:15:55 -0700 Subject: [PATCH] Retry retryable worker failures once before failing the session The harness gives up after three attempts and the worker dies even though no tool ran yet. One outer retry with backoff recovers the flake without another orchestrator round-trip, guarded against tool replay and bounded so it cannot multiply sends. --- .../run-resolved-provider-failure.test.ts | 409 ++++++++++++++++++ src/subagent/run.ts | 120 ++++- 2 files changed, 517 insertions(+), 12 deletions(-) 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