diff --git a/packages/agent/.changes/res-1271-retry-owner.md b/packages/agent/.changes/res-1271-retry-owner.md new file mode 100644 index 0000000000..96e772f5e8 --- /dev/null +++ b/packages/agent/.changes/res-1271-retry-owner.md @@ -0,0 +1 @@ +- Removed the unused `maxRetryDelayMs` agent option; retry delays are owned by the session retry loop. diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index 0b7b5ae774..9c7960bc4a 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -112,7 +112,6 @@ export interface AgentOptions { sessionId?: string; thinkingBudgets?: ThinkingBudgets; transport?: Transport; - maxRetryDelayMs?: number; toolExecution?: ToolExecutionMode; } @@ -216,7 +215,6 @@ export class Agent { public sessionId?: string; public thinkingBudgets?: ThinkingBudgets; public transport: Transport; - public maxRetryDelayMs?: number; public toolExecution: ToolExecutionMode; constructor(options: AgentOptions = {}) { @@ -237,7 +235,6 @@ export class Agent { this.sessionId = options.sessionId; this.thinkingBudgets = options.thinkingBudgets; this.transport = options.transport ?? "auto"; - this.maxRetryDelayMs = options.maxRetryDelayMs; this.toolExecution = options.toolExecution ?? "parallel"; } @@ -470,7 +467,6 @@ export class Agent { onResponse: this.onResponse, transport: this.transport, thinkingBudgets: this.thinkingBudgets, - maxRetryDelayMs: this.maxRetryDelayMs, toolExecution: this.toolExecution, beforeToolCall: this.beforeToolCall, afterToolCall: this.afterToolCall, diff --git a/packages/agent/src/proxy.ts b/packages/agent/src/proxy.ts index 69920ce033..83f3b65385 100644 --- a/packages/agent/src/proxy.ts +++ b/packages/agent/src/proxy.ts @@ -57,7 +57,6 @@ type ProxySerializableStreamOptions = Pick< | "metadata" | "transport" | "thinkingBudgets" - | "maxRetryDelayMs" >; export interface ProxyStreamOptions extends ProxySerializableStreamOptions { @@ -96,7 +95,6 @@ function buildProxyRequestOptions(options: ProxyStreamOptions): ProxySerializabl metadata: options.metadata, transport: options.transport, thinkingBudgets: options.thinkingBudgets, - maxRetryDelayMs: options.maxRetryDelayMs, }; } diff --git a/packages/ai/.changes/res-1271-retry-owner.md b/packages/ai/.changes/res-1271-retry-owner.md new file mode 100644 index 0000000000..5eebccf802 --- /dev/null +++ b/packages/ai/.changes/res-1271-retry-owner.md @@ -0,0 +1,4 @@ +- Removed SDK-internal provider retries (OpenAI, Azure, Anthropic, Bedrock, Codex): providers make a single attempt and report a structured failure so the agent's own retry loop owns every retry. +- Added the server-requested Retry-After delay to structured stream-failure diagnostics. +- Added structured stream-failure diagnostics to the OpenAI-completions and Codex providers, including friendly messages for Codex nested usage-limit error payloads. +- Removed the `maxRetries` and `maxRetryDelayMs` stream options. diff --git a/packages/ai/src/providers/amazon-bedrock.ts b/packages/ai/src/providers/amazon-bedrock.ts index 0a540c337e..455836a8c8 100644 --- a/packages/ai/src/providers/amazon-bedrock.ts +++ b/packages/ai/src/providers/amazon-bedrock.ts @@ -115,6 +115,7 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream", BedrockOpt const config: BedrockRuntimeClientConfig = { profile: options.profile, + maxAttempts: 1, }; const configuredRegion = getConfiguredBedrockRegion(options); const hasConfiguredProfile = hasConfiguredBedrockProfile(); diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index a3fcb41521..658dceb6d9 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -521,7 +521,6 @@ export const streamAnthropic: StreamFunction<"anthropic-messages", AnthropicOpti const requestOptions = { ...(options?.signal ? { signal: options.signal } : {}), ...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}), - ...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}), }; const response = await client.messages.create({ ...params, stream: true }, requestOptions).asResponse(); await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model); @@ -861,6 +860,7 @@ function createClient( if (model.provider === "cloudflare-ai-gateway") { const client = new Anthropic({ + maxRetries: 0, apiKey: null, authToken: null, baseURL: resolveCloudflareBaseUrl(model), @@ -884,6 +884,7 @@ function createClient( if (model.provider === "github-copilot") { const client = new Anthropic({ + maxRetries: 0, apiKey: null, authToken: apiKey, baseURL: model.baseUrl, @@ -905,6 +906,7 @@ function createClient( if (isOAuthToken(apiKey)) { const client = new Anthropic({ + maxRetries: 0, apiKey: null, authToken: apiKey, baseURL: model.baseUrl, @@ -926,6 +928,7 @@ function createClient( } const client = new Anthropic({ + maxRetries: 0, apiKey, baseURL: model.baseUrl, dangerouslyAllowBrowser: true, diff --git a/packages/ai/src/providers/azure-openai-responses.ts b/packages/ai/src/providers/azure-openai-responses.ts index 6e869afc44..b52632086d 100644 --- a/packages/ai/src/providers/azure-openai-responses.ts +++ b/packages/ai/src/providers/azure-openai-responses.ts @@ -93,7 +93,6 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" const requestOptions = { ...(options?.signal ? { signal: options.signal } : {}), ...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}), - ...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}), }; const { data: openaiStream, response } = await client.responses.create(params, requestOptions).withResponse(); await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model); @@ -231,6 +230,7 @@ function createClient(model: Model<"azure-openai-responses">, apiKey: string, op dangerouslyAllowBrowser: true, defaultHeaders: headers, baseURL: baseUrl, + maxRetries: 0, }); } diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index cf9f113a86..b1eca527d1 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -40,13 +40,12 @@ import { } from "../utils/diagnostics.js"; import { AssistantMessageEventStream } from "../utils/event-stream.js"; import { headersToRecord } from "../utils/headers.js"; +import { parseRetryAfterMs, recordStreamFailure } from "../utils/stream-failure.js"; import { convertResponsesMessages, convertResponsesTools, processResponsesStream } from "./openai-responses-shared.js"; import { buildBaseOptions } from "./simple-options.js"; const DEFAULT_CODEX_BASE_URL = "https://chatgpt.com/backend-api"; const JWT_CLAIM_PATH = "https://api.openai.com/auth" as const; -const MAX_RETRIES = 3; -const BASE_DELAY_MS = 1000; const CODEX_TOOL_CALL_PROVIDERS = new Set(["openai", "openai-codex", "opencode"]); const WEBSOCKET_MESSAGE_TOO_BIG_CLOSE_CODE = 1009; @@ -87,27 +86,6 @@ interface RequestBody { [key: string]: unknown; } -function isRetryableError(status: number, errorText: string): boolean { - if (status === 429 || status === 500 || status === 502 || status === 503 || status === 504) { - return true; - } - return /rate.?limit|overloaded|service.?unavailable|upstream.?connect|connection.?refused/i.test(errorText); -} - -function sleep(ms: number, signal?: AbortSignal): Promise { - return new Promise((resolve, reject) => { - if (signal?.aborted) { - reject(new Error("Request was aborted")); - return; - } - const timeout = setTimeout(resolve, ms); - signal?.addEventListener("abort", () => { - clearTimeout(timeout); - reject(new Error("Request was aborted")); - }); - }); -} - export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses", OpenAICodexResponsesOptions> = ( model: Model<"openai-codex-responses">, context: Context, @@ -211,61 +189,31 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" } } - let response: Response | undefined; - let lastError: Error | undefined; + if (options?.signal?.aborted) { + throw new Error("Request was aborted"); + } - for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) { - if (options?.signal?.aborted) { + let response: Response; + try { + response = await fetch(resolveCodexUrl(model.baseUrl), { + method: "POST", + headers: sseHeaders, + body: bodyJson, + signal: options?.signal, + }); + } catch (error) { + if ( + options?.signal?.aborted || + (error instanceof Error && (error.name === "AbortError" || error.message === "Request was aborted")) + ) { throw new Error("Request was aborted"); } - - try { - response = await fetch(resolveCodexUrl(model.baseUrl), { - method: "POST", - headers: sseHeaders, - body: bodyJson, - signal: options?.signal, - }); - await options?.onResponse?.( - { status: response.status, headers: headersToRecord(response.headers) }, - model, - ); - - if (response.ok) { - break; - } - - const errorText = await response.text(); - if (attempt < MAX_RETRIES && isRetryableError(response.status, errorText)) { - const delayMs = BASE_DELAY_MS * 2 ** attempt; - await sleep(delayMs, options?.signal); - continue; - } - - const fakeResponse = new Response(errorText, { - status: response.status, - statusText: response.statusText, - }); - const info = await parseErrorResponse(fakeResponse); - throw new Error(info.friendlyMessage || info.message); - } catch (error) { - if (error instanceof Error) { - if (error.name === "AbortError" || error.message === "Request was aborted") { - throw new Error("Request was aborted"); - } - } - lastError = error instanceof Error ? error : new Error(String(error)); - if (attempt < MAX_RETRIES && !lastError.message.includes("usage limit")) { - const delayMs = BASE_DELAY_MS * 2 ** attempt; - await sleep(delayMs, options?.signal); - continue; - } - throw lastError; - } + throw error; } + await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model); - if (!response?.ok) { - throw lastError ?? new Error("Failed after retries"); + if (!response.ok) { + throw await parseErrorResponse(response); } if (!response.body) { @@ -288,6 +236,7 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; output.errorMessage = error instanceof Error ? error.message : String(error); + recordStreamFailure(model, output, error); stream.push({ type: "error", reason: output.stopReason, error: output }); stream.end(); } @@ -437,12 +386,25 @@ async function processStream( class CodexApiError extends Error { readonly code?: string; + readonly status?: number; + readonly retryAfterMs?: number; readonly payload?: Record; - constructor(message: string, options?: { code?: string; payload?: Record; cause?: unknown }) { + constructor( + message: string, + options?: { + code?: string; + status?: number; + retryAfterMs?: number; + payload?: Record; + cause?: unknown; + }, + ) { super(message); this.name = "CodexApiError"; this.code = options?.code; + this.status = options?.status; + this.retryAfterMs = options?.retryAfterMs; this.payload = options?.payload; this.cause = options?.cause; } @@ -469,12 +431,24 @@ async function* mapCodexEvents(events: AsyncIterable>): if (!type) continue; if (type === "error") { - const code = (event as { code?: string }).code || ""; - const message = (event as { message?: string }).message || ""; - throw new CodexApiError(`Codex error: ${message || code || JSON.stringify(event)}`, { - code: code || undefined, - payload: event, - }); + // Errors arrive flat ({ code, message }) or nested under event.error. + const flatCode = typeof event.code === "string" ? event.code : ""; + const flatMessage = typeof event.message === "string" ? event.message : ""; + const nested = event.error && typeof event.error === "object" ? (event.error as CodexErrorPayload) : undefined; + const statusCode = (event as { status_code?: unknown }).status_code; + const status = typeof statusCode === "number" ? statusCode : undefined; + const code = flatCode || nested?.code || nested?.type || undefined; + const usageLimit = nested ? codexUsageLimitMessage(nested, status) : undefined; + const message = flatMessage || nested?.message || ""; + throw new CodexApiError( + usageLimit?.friendlyMessage ?? `Codex error: ${message || code || JSON.stringify(event)}`, + { + code, + status, + retryAfterMs: usageLimit?.retryAfterMs, + payload: event, + }, + ); } if (type === "response.failed") { @@ -1183,33 +1157,58 @@ async function processWebSocketStream( } } -async function parseErrorResponse(response: Response): Promise<{ message: string; friendlyMessage?: string }> { +interface CodexErrorPayload { + code?: string; + type?: string; + message?: string; + plan_type?: string; + resets_at?: number; +} + +function codexUsageLimitMessage( + err: CodexErrorPayload, + status?: number, +): { friendlyMessage: string; retryAfterMs?: number } | undefined { + const code = err.code || err.type || ""; + if (!/usage_limit_reached|usage_not_included|rate_limit_exceeded/i.test(code) && status !== 429) { + return undefined; + } + const plan = err.plan_type ? ` (${err.plan_type.toLowerCase()} plan)` : ""; + const retryAfterMs = err.resets_at !== undefined ? Math.max(0, err.resets_at * 1000 - Date.now()) : undefined; + const when = retryAfterMs !== undefined ? ` Try again in ~${Math.round(retryAfterMs / 60000)} min.` : ""; + return { + friendlyMessage: `You have hit your ChatGPT usage limit${plan}.${when}`.trim(), + retryAfterMs, + }; +} + +async function parseErrorResponse(response: Response): Promise { const raw = await response.text(); let message = raw || response.statusText || "Request failed"; - let friendlyMessage: string | undefined; + let code: string | undefined; + let retryAfterMs = parseRetryAfterMs(response.headers); try { - const parsed = JSON.parse(raw) as { - error?: { code?: string; type?: string; message?: string; plan_type?: string; resets_at?: number }; - }; + const parsed = JSON.parse(raw) as { error?: CodexErrorPayload }; const err = parsed?.error; if (err) { - const code = err.code || err.type || ""; - if (/usage_limit_reached|usage_not_included|rate_limit_exceeded/i.test(code) || response.status === 429) { - const plan = err.plan_type ? ` (${err.plan_type.toLowerCase()} plan)` : ""; - const mins = err.resets_at - ? Math.max(0, Math.round((err.resets_at * 1000 - Date.now()) / 60000)) - : undefined; - const when = mins !== undefined ? ` Try again in ~${mins} min.` : ""; - friendlyMessage = `You have hit your ChatGPT usage limit${plan}.${when}`.trim(); + code = err.code || err.type || undefined; + const usageLimit = codexUsageLimitMessage(err, response.status); + if (usageLimit) { + message = usageLimit.friendlyMessage; + // Neither server delay (Retry-After header, resets_at body) may undercut the other. + if (usageLimit.retryAfterMs !== undefined) { + retryAfterMs = Math.max(retryAfterMs ?? 0, usageLimit.retryAfterMs); + } + } else { + message = err.message || message; } - message = err.message || friendlyMessage || message; } } catch { // Unparseable error body: fall back to the raw message. } - return { message, friendlyMessage }; + return new CodexApiError(message, { code, status: response.status, retryAfterMs }); } function extractAccountId(token: string): string { diff --git a/packages/ai/src/providers/openai-completions.ts b/packages/ai/src/providers/openai-completions.ts index e39d959178..64bb3be89e 100644 --- a/packages/ai/src/providers/openai-completions.ts +++ b/packages/ai/src/providers/openai-completions.ts @@ -35,6 +35,7 @@ import { AssistantMessageEventStream } from "../utils/event-stream.js"; import { headersToRecord } from "../utils/headers.js"; import { parseStreamingJson } from "../utils/json-parse.js"; import { sanitizeSurrogates } from "../utils/sanitize-unicode.js"; +import { recordStreamFailure } from "../utils/stream-failure.js"; import { isCloudflareProvider, resolveCloudflareBaseUrl } from "./cloudflare.js"; import { buildCopilotDynamicHeaders, hasCopilotVisionInput } from "./github-copilot-headers.js"; import { buildBaseOptions } from "./simple-options.js"; @@ -181,7 +182,6 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions", OpenA const requestOptions = { ...(options?.signal ? { signal: options.signal } : {}), ...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}), - ...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}), }; const { data: openaiStream, response } = await client.chat.completions .create(params, requestOptions) @@ -476,6 +476,7 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions", OpenA // Some providers via OpenRouter give additional information in this field. const rawMetadata = (error as any)?.error?.metadata?.raw; if (rawMetadata) output.errorMessage += `\n${rawMetadata}`; + recordStreamFailure(model, output, error); stream.push({ type: "error", reason: output.stopReason, error: output }); stream.end(); } @@ -565,6 +566,7 @@ function createClient( baseURL: isCloudflareProvider(model.provider) ? resolveCloudflareBaseUrl(model) : model.baseUrl, dangerouslyAllowBrowser: true, defaultHeaders, + maxRetries: 0, }); } diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index 0370d63fcc..53e447b20a 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -101,7 +101,6 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses", OpenAIRes const requestOptions = { ...(options?.signal ? { signal: options.signal } : {}), ...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}), - ...(options?.maxRetries !== undefined ? { maxRetries: options.maxRetries } : {}), }; const { data: openaiStream, response } = await client.responses.create(params, requestOptions).withResponse(); await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model); @@ -212,6 +211,7 @@ function createClient( baseURL: isCloudflareProvider(model.provider) ? resolveCloudflareBaseUrl(model) : model.baseUrl, dangerouslyAllowBrowser: true, defaultHeaders, + maxRetries: 0, }); } diff --git a/packages/ai/src/providers/simple-options.ts b/packages/ai/src/providers/simple-options.ts index 76f6ea8b44..51f5189074 100644 --- a/packages/ai/src/providers/simple-options.ts +++ b/packages/ai/src/providers/simple-options.ts @@ -14,8 +14,6 @@ export function buildBaseOptions(model: Model, options?: SimpleStreamOption onPayload: options?.onPayload, onResponse: options?.onResponse, timeoutMs: options?.timeoutMs, - maxRetries: options?.maxRetries, - maxRetryDelayMs: options?.maxRetryDelayMs, metadata: options?.metadata, }; } diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index d7a679fac0..f5744ba97e 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -118,19 +118,6 @@ export interface StreamOptions { * For example, OpenAI and Anthropic SDK clients default to 10 minutes. */ timeoutMs?: number; - /** - * Maximum retry attempts for providers/SDKs that support client-side retries. - * For example, OpenAI and Anthropic SDK clients default to 2. - */ - maxRetries?: number; - /** - * Maximum delay in milliseconds to wait for a retry when the server requests a long wait. - * If the server's requested delay exceeds this value, the request fails immediately - * with an error containing the requested delay, allowing higher-level retry logic - * to handle it with user visibility. - * Default: 60000 (60 seconds). Set to 0 to disable the cap. - */ - maxRetryDelayMs?: number; /** * Optional metadata to include in API requests. * Providers extract the fields they understand and ignore the rest. diff --git a/packages/ai/src/utils/stream-failure.ts b/packages/ai/src/utils/stream-failure.ts index af8b2a7ea5..ec112607fc 100644 --- a/packages/ai/src/utils/stream-failure.ts +++ b/packages/ai/src/utils/stream-failure.ts @@ -25,6 +25,8 @@ export interface StreamFailureInfo { providerErrorType?: string; status?: number; requestId?: string; + /** Server-requested wait before retrying (Retry-After header or reset info), in milliseconds. */ + retryAfterMs?: number; /** Truncated raw provider payload for post-mortems. */ raw?: string; } @@ -70,7 +72,10 @@ export function classifyStreamFailure(providerErrorType?: string, status?: numbe return "safety"; } if (type.includes("overloaded") || status === 529) return "overloaded"; - if (type.includes("rate_limit") || type.includes("throttl") || status === 429) return "rate_limit"; + // usage_not_included is Codex's plan-entitlement rejection, not bad credentials. + if (/rate_limit|usage_limit|usage_not_included|throttl/.test(type) || status === 429) { + return "rate_limit"; + } if (/authentication|permission|unauthorized/.test(type) || status === 401 || status === 403) return "auth"; if (type.includes("invalid_request") || type.includes("not_found_error") || status === 400 || status === 404) { return "invalid_request"; @@ -126,6 +131,7 @@ function extractStreamFailureParts(error: unknown): { info: StreamFailureInfo; d request_id?: unknown; headers?: unknown; error?: unknown; + retryAfterMs?: unknown; $metadata?: { requestId?: unknown }; }; @@ -150,15 +156,11 @@ function extractStreamFailureParts(error: unknown): { info: StreamFailureInfo; d : undefined; const headers = err.headers; - const headerRequestId = - headers && typeof (headers as Headers).get === "function" - ? ((headers as Headers).get("request-id") ?? (headers as Headers).get("x-request-id")) - : headers && typeof headers === "object" - ? ((headers as Record)["request-id"] ?? - (headers as Record)["x-request-id"]) - : undefined; + const headerRequestId = headerValue(headers, "request-id") ?? headerValue(headers, "x-request-id"); const rawRequestId = err.requestID ?? err.request_id ?? err.$metadata?.requestId ?? headerRequestId; const requestId = typeof rawRequestId === "string" ? rawRequestId : undefined; + const retryAfterMs = + typeof err.retryAfterMs === "number" && err.retryAfterMs >= 0 ? err.retryAfterMs : parseRetryAfterMs(headers); return { info: { @@ -166,11 +168,38 @@ function extractStreamFailureParts(error: unknown): { info: StreamFailureInfo; d providerErrorType, status, requestId, + ...(retryAfterMs !== undefined ? { retryAfterMs } : {}), }, detail: typeof bodyMessage === "string" ? bodyMessage : undefined, }; } +function headerValue(headers: unknown, name: string): string | undefined { + if (!headers || typeof headers !== "object") return undefined; + if (typeof (headers as Headers).get === "function") { + return (headers as Headers).get(name) ?? undefined; + } + // Record-shaped headers must match case-insensitively, like real Headers. + for (const [key, value] of Object.entries(headers as Record)) { + if (key.toLowerCase() === name && typeof value === "string") { + return value; + } + } + return undefined; +} + +/** Parse Retry-After / Retry-After-Ms headers into a millisecond wait. */ +export function parseRetryAfterMs(headers: unknown): number | undefined { + const ms = Number(headerValue(headers, "retry-after-ms")); + if (Number.isFinite(ms) && ms >= 0) return ms; + const raw = headerValue(headers, "retry-after"); + if (!raw) return undefined; + const seconds = Number(raw); + if (Number.isFinite(seconds) && seconds >= 0) return seconds * 1000; + const date = Date.parse(raw); + return Number.isNaN(date) ? undefined : Math.max(0, date - Date.now()); +} + /** * Best-effort extraction of structured failure info from any thrown value: * StreamFailureError, provider SDK errors (Anthropic/OpenAI APIError, AWS SDK diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index e05b9d9831..ea0fc90941 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -8,7 +8,7 @@ import { streamOpenAICodexResponses, streamSimpleOpenAICodexResponses, } from "../src/providers/openai-codex-responses.js"; -import type { Context, Model } from "../src/types.js"; +import type { AssistantMessage, Context, Model } from "../src/types.js"; const originalFetch = global.fetch; const originalWebSocket = globalThis.WebSocket; @@ -1005,4 +1005,108 @@ describe("openai-codex streaming", () => { lastPreviousResponseId: "resp_1", }); }); + + function codexTestModel(): Model<"openai-codex-responses"> { + return { + id: "gpt-5.1-codex", + name: "GPT-5.1 Codex", + api: "openai-codex-responses", + provider: "openai-codex", + baseUrl: "https://chatgpt.com/backend-api", + reasoning: true, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 400000, + maxTokens: 128000, + }; + } + + /** Stub prompt-cache URLs plus a custom /codex/responses handler; returns the request counter. */ + function stubCodexFetch(respond: () => Response): { responsesRequests: number } { + process.env.PI_CODING_AGENT_DIR = mkdtempSync(join(tmpdir(), "pi-codex-stream-")); + const counter = { responsesRequests: 0 }; + global.fetch = vi.fn(async (input: string | URL) => { + const url = typeof input === "string" ? input : input.toString(); + if (url === "https://api.github.com/repos/openai/codex/releases/latest") { + return new Response(JSON.stringify({ tag_name: "rust-v0.0.0" }), { status: 200 }); + } + if (url.startsWith("https://raw.githubusercontent.com/openai/codex/")) { + return new Response("PROMPT", { status: 200, headers: { etag: '"etag"' } }); + } + if (url === "https://chatgpt.com/backend-api/codex/responses") { + counter.responsesRequests++; + return respond(); + } + return new Response("not found", { status: 404 }); + }) as typeof fetch; + return counter; + } + + async function runCodexErrorTurn(): Promise { + const context: Context = { + systemPrompt: "You are a helpful assistant.", + messages: [{ role: "user", content: "Say hello", timestamp: Date.now() }], + }; + return streamOpenAICodexResponses(codexTestModel(), context, { apiKey: mockToken(), transport: "sse" }).result(); + } + + function failureDetails(result: AssistantMessage): { kind?: string; retryAfterMs?: number } | undefined { + const last = result.diagnostics?.at(-1); + return last?.type === "provider_stream_failure" + ? (last as { details?: { kind?: string; retryAfterMs?: number } }).details + : undefined; + } + + it("throws a structured failure after a single attempt on HTTP 500", async () => { + const counter = stubCodexFetch( + () => new Response(JSON.stringify({ error: { type: "server_error", message: "boom" } }), { status: 500 }), + ); + + const result = await runCodexErrorTurn(); + + expect(counter.responsesRequests).toBe(1); + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toBe("boom"); + expect(failureDetails(result)).toMatchObject({ kind: "server_error", status: 500 }); + }); + + it("maps nested streaming usage-limit error payloads to a friendly rate-limit failure", async () => { + const resetsAt = Math.round(Date.now() / 1000) + 2 * 3600; + const sse = `data: ${JSON.stringify({ + type: "error", + status_code: 429, + error: { type: "usage_limit_reached", message: "Usage limit reached", plan_type: "Plus", resets_at: resetsAt }, + })}\n\n`; + stubCodexFetch(() => new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } })); + + const result = await runCodexErrorTurn(); + + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toContain( + "You have hit your ChatGPT usage limit (plus plan). Try again in ~120 min.", + ); + const details = failureDetails(result); + expect(details?.kind).toBe("rate_limit"); + expect(details?.retryAfterMs).toBeGreaterThan(0); + // resets_at has second granularity, so allow the rounding slack. + expect(details?.retryAfterMs).toBeLessThanOrEqual(2 * 3600 * 1000 + 1000); + }); + + it("waits for the longer of Retry-After header and usage-limit reset", async () => { + const resetsAt = Math.round(Date.now() / 1000) + 10; + stubCodexFetch( + () => + new Response( + JSON.stringify({ + error: { type: "usage_limit_reached", message: "Usage limit reached", resets_at: resetsAt }, + }), + { status: 429, headers: { "retry-after": "60" } }, + ), + ); + + const result = await runCodexErrorTurn(); + + expect(result.stopReason).toBe("error"); + expect(failureDetails(result)?.retryAfterMs).toBe(60000); + }); }); diff --git a/packages/ai/test/provider-retry-ownership.test.ts b/packages/ai/test/provider-retry-ownership.test.ts new file mode 100644 index 0000000000..1e7d156cdc --- /dev/null +++ b/packages/ai/test/provider-retry-ownership.test.ts @@ -0,0 +1,57 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { getModel } from "../src/models.js"; +import { streamOpenAICompletions } from "../src/providers/openai-completions.js"; +import type { Context, Model } from "../src/types.js"; + +const originalFetch = global.fetch; + +afterEach(() => { + global.fetch = originalFetch; + vi.restoreAllMocks(); +}); + +function completionsModel(): Model<"openai-completions"> { + const { compat: _compat, ...baseModel } = getModel("openai", "gpt-4o-mini")!; + return { ...baseModel, api: "openai-completions" } as Model<"openai-completions">; +} + +const context: Context = { + messages: [{ role: "user", content: "hi", timestamp: Date.now() }], +}; + +async function streamToError(model: Model<"openai-completions">) { + const stream = streamOpenAICompletions(model, context, { apiKey: "test-key" }); + for await (const event of stream) { + if (event.type === "error") return event.error; + } + throw new Error("expected an error event"); +} + +describe("provider retry ownership", () => { + it.each([ + [ + "makes exactly one request on a 500 and records a structured stream failure", + { type: "server_error", message: "boom" }, + { status: 500 }, + { kind: "server_error", status: 500 }, + ], + [ + "surfaces the server-requested Retry-After delay on rate limits", + { type: "rate_limit_error", message: "slow down" }, + { status: 429, headers: { "retry-after": "30" } }, + { kind: "rate_limit", status: 429, retryAfterMs: 30000 }, + ], + ] as const)("%s", async (_name, errorBody, init, expectedDetails) => { + const fetchMock = vi.fn(async () => new Response(JSON.stringify({ error: errorBody }), init)); + global.fetch = fetchMock as typeof fetch; + + const output = await streamToError(completionsModel()); + + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(output.stopReason).toBe("error"); + expect(output.diagnostics?.[0]).toMatchObject({ + type: "provider_stream_failure", + details: expectedDetails, + }); + }); +}); diff --git a/packages/ai/test/stream-failure.test.ts b/packages/ai/test/stream-failure.test.ts index 5009de0160..4579737a8a 100644 --- a/packages/ai/test/stream-failure.test.ts +++ b/packages/ai/test/stream-failure.test.ts @@ -38,6 +38,8 @@ describe("classifyStreamFailure", () => { ["overloaded_error", undefined, "overloaded"], [undefined, 529, "overloaded"], ["rate_limit_error", undefined, "rate_limit"], + ["usage_limit_reached", undefined, "rate_limit"], + ["usage_not_included", 403, "rate_limit"], [undefined, 429, "rate_limit"], ["refusal", undefined, "refusal"], ["sensitive", undefined, "safety"], @@ -102,6 +104,25 @@ describe("extractStreamFailureInfo", () => { expect(extractStreamFailureInfo(awsError)).toMatchObject({ requestId: "aws_req" }); }); + test.each([ + ["Headers seconds", new Headers({ "retry-after": "120" }), 120000], + ["retry-after-ms precedence", { "retry-after-ms": "1500", "retry-after": "2" }, 1500], + ["record with mixed case", { "Retry-After": "120" }, 120000], + ] as const)("extracts the server-requested retry delay: %s", (_name, headers, expected) => { + const error = Object.assign(new Error("429"), { status: 429, headers }); + expect(extractStreamFailureInfo(error)).toMatchObject({ kind: "rate_limit", retryAfterMs: expected }); + }); + + test("parses an HTTP-date Retry-After relative to now", () => { + const withDate = Object.assign(new Error("429"), { + status: 429, + headers: new Headers({ "retry-after": new Date(Date.now() + 60000).toUTCString() }), + }); + const dateMs = extractStreamFailureInfo(withDate).retryAfterMs; + expect(dateMs).toBeGreaterThan(0); + expect(dateMs).toBeLessThanOrEqual(60000); + }); + test("falls back to classifying the message text", () => { expect(extractStreamFailureInfo(new Error("provider overloaded, retry later")).kind).toBe("overloaded"); expect(extractStreamFailureInfo("not an error").kind).toBe("unknown"); diff --git a/packages/coding-agent/.changes/res-1271-retry-owner.md b/packages/coding-agent/.changes/res-1271-retry-owner.md new file mode 100644 index 0000000000..6fa4e93702 --- /dev/null +++ b/packages/coding-agent/.changes/res-1271-retry-owner.md @@ -0,0 +1,4 @@ +- Changed auto-retry to honor provider Retry-After and usage-limit reset delays, capped by `retry.provider.maxRetryDelayMs`; longer requested waits fail immediately with an informative error instead of sleeping invisibly inside provider SDKs. +- Removed the `retry.provider.maxRetries` setting; provider SDKs no longer retry internally, so `retry.maxRetries` is the single retry knob. +- Changed structured `invalid_request`/`refusal` provider failures to fail immediately instead of being retried once. +- Added the shared retry policy to side questions, compaction and branch summarization, and refinement calls, which run outside the session auto-retry loop and would otherwise make exactly one attempt. diff --git a/packages/coding-agent/docs/settings.md b/packages/coding-agent/docs/settings.md index 8e1c147260..1fa0cfbef1 100644 --- a/packages/coding-agent/docs/settings.md +++ b/packages/coding-agent/docs/settings.md @@ -141,10 +141,9 @@ prime-agent --offline | `retry.maxRetries` | number | `3` | Maximum agent-level retry attempts | | `retry.baseDelayMs` | number | `2000` | Base delay for agent-level exponential backoff (2s, 4s, 8s) | | `retry.provider.timeoutMs` | number | SDK default | Provider/SDK request timeout in milliseconds | -| `retry.provider.maxRetries` | number | SDK default | Provider/SDK retry attempts | -| `retry.provider.maxRetryDelayMs` | number | `60000` | Max server-requested delay before failing (60s) | +| `retry.provider.maxRetryDelayMs` | number | `60000` | Max server-requested retry delay before failing (60s) | -When a provider requests a retry delay longer than `retry.provider.maxRetryDelayMs` (e.g., Google's "quota will reset after 5h"), the request fails immediately with an informative error instead of waiting silently. Set to `0` to disable the cap. +When a provider requests a retry delay longer than `retry.provider.maxRetryDelayMs` (e.g. a usage-limit reset hours away), auto-retry stops immediately with an informative error instead of waiting. Set to `0` to disable the cap. ```json { @@ -154,7 +153,6 @@ When a provider requests a retry delay longer than `retry.provider.maxRetryDelay "baseDelayMs": 2000, "provider": { "timeoutMs": 3600000, - "maxRetries": 0, "maxRetryDelayMs": 60000 } } diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index be375cc048..c58012873d 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -187,6 +187,16 @@ import { import type { ModelRegistry } from "./model-registry.js"; import { throwIfPromptAdmissionCancelled } from "./prompt-admission.js"; import { expandPromptTemplate, type PromptTemplate } from "./prompt-templates.js"; +import { + isAgentLifecycleFailure, + isFauxProviderQueueExhausted, + isPermanentProviderFailureKind, + providerRetryDelay, + providerRetryPolicy, + providerStreamFailureDetails, + providerStreamFailureKind, + providerStreamFailureRetryAfterMs, +} from "./provider-retry.js"; import { type AutoRefineReason, type AutoRefineReview, @@ -7625,6 +7635,7 @@ export class AgentSession { signal, this.thinkingLevel, summaryCall, + providerRetryPolicy(this.settingsManager), )); } @@ -8144,6 +8155,7 @@ export class AgentSession { headers, signal, this.thinkingLevel, + providerRetryPolicy(this.settingsManager), ); } @@ -8379,7 +8391,7 @@ export class AgentSession { history, model, apiKey, - options, + { ...options, retry: providerRetryPolicy(this.settingsManager) }, headers, signal, this.thinkingLevel, @@ -9661,7 +9673,6 @@ export class AgentSession { sessionId: childSessionManager.getSessionId(), thinkingBudgets: this.settingsManager.getThinkingBudgets(), transport: this.settingsManager.getTransport(), - maxRetryDelayMs: this.settingsManager.getProviderRetrySettings().maxRetryDelayMs, toolExecution: this.agent.toolExecution, }); @@ -10979,34 +10990,23 @@ export class AgentSession { } private _isFauxProviderQueueExhausted(message: AssistantMessage): boolean { - return message.provider === "faux" && message.errorMessage === "No more faux responses queued"; + return isFauxProviderQueueExhausted(message); } private _isAgentLifecycleFailure(message: AssistantMessage): boolean { - return message.diagnostics?.some((diagnostic) => diagnostic.type === "agent_lifecycle_failure") ?? false; + return isAgentLifecycleFailure(message); } private _getProviderStreamFailureDetails(message: AssistantMessage): Record | undefined { - const failure = message.diagnostics?.find((diagnostic) => diagnostic.type === "provider_stream_failure"); - const details = failure?.details; - if (!details || typeof details !== "object") { - return undefined; - } - return details; + return providerStreamFailureDetails(message); } private _getProviderStreamFailureKind(message: AssistantMessage): string | undefined { - const kind = this._getProviderStreamFailureDetails(message)?.kind; - return typeof kind === "string" ? kind : undefined; - } - - private _isStructuredPermanentProviderFailure(message: AssistantMessage): boolean { - const kind = this._getProviderStreamFailureKind(message); - return kind === "auth" || kind === "invalid_request" || kind === "refusal"; + return providerStreamFailureKind(message); } private _isStructuredPermanentProviderRetryExhausted(message: AssistantMessage): boolean { - return this._retryAttempt > 0 && this._isStructuredPermanentProviderFailure(message); + return isPermanentProviderFailureKind(this._getProviderStreamFailureKind(message), this._retryAttempt); } private _getProviderStreamFailureAuthStatus(message: AssistantMessage): number | undefined { @@ -11161,7 +11161,27 @@ export class AgentSession { return false; } - const delayMs = settings.baseDelayMs * 2 ** (this._retryAttempt - 1); + // Server-requested waits are honored, capped by retry.provider.maxRetryDelayMs (0 disables). + const maxRetryDelayMs = this.settingsManager.getProviderRetrySettings().maxRetryDelayMs; + const delay = providerRetryDelay(this._retryAttempt, providerStreamFailureRetryAfterMs(message), { + baseDelayMs: settings.baseDelayMs, + maxRetryDelayMs, + }); + if (delay.kind === "exceeds-cap") { + this._markProviderAuthStaleForRetryFailure(message, options); + this._emit({ + type: "auto_retry_end", + success: false, + attempt: this._retryAttempt - 1, + finalError: `Provider requested a ${Math.ceil(delay.retryAfterMs / 1000)}s wait before retrying (above retry.provider.maxRetryDelayMs=${maxRetryDelayMs}ms): ${message.errorMessage || "unknown error"}`, + }); + this._retryAttempt = 0; + this._retryAuthFailureSources = []; + this._resolveRetry(); + return false; + } + + const delayMs = delay.delayMs; // Park now: the retry re-issues the failed call and must reuse its Idempotency-Key. // Payload hooks mutate the wire body after the hash point, so reuse is forfeited. if (!this._extensionRunner.hasHandlers("before_provider_request")) { @@ -11729,6 +11749,7 @@ export class AgentSession { customInstructions, replaceInstructions, reserveTokens: branchSummarySettings.reserveTokens, + retry: providerRetryPolicy(this.settingsManager), }); if (result.aborted) { return { cancelled: true, aborted: true }; diff --git a/packages/coding-agent/src/core/compaction/branch-summarization.ts b/packages/coding-agent/src/core/compaction/branch-summarization.ts index 3917a45d11..d6a5df0315 100644 --- a/packages/coding-agent/src/core/compaction/branch-summarization.ts +++ b/packages/coding-agent/src/core/compaction/branch-summarization.ts @@ -14,6 +14,7 @@ import { createCompactionSummaryMessage, createCustomMessage, } from "../messages.js"; +import { completeWithProviderRetry, type ProviderRetryPolicy } from "../provider-retry.js"; import type { ReadonlySessionManager, SessionEntry } from "../session-manager.js"; import { estimateTokens } from "./compaction.js"; import { @@ -71,6 +72,7 @@ export interface GenerateBranchSummaryOptions { customInstructions?: string; /** If true, customInstructions replaces the default prompt instead of being appended */ replaceInstructions?: boolean; + retry?: ProviderRetryPolicy; /** Tokens reserved for prompt + LLM response (default 16384) */ reserveTokens?: number; } @@ -250,7 +252,16 @@ export async function generateBranchSummary( entries: SessionEntry[], options: GenerateBranchSummaryOptions, ): Promise { - const { model, apiKey, headers, signal, customInstructions, replaceInstructions, reserveTokens = 16384 } = options; + const { + model, + apiKey, + headers, + signal, + customInstructions, + replaceInstructions, + retry, + reserveTokens = 16384, + } = options; const contextWindow = model.contextWindow || 128000; const tokenBudget = contextWindow - reserveTokens; @@ -280,10 +291,14 @@ export async function generateBranchSummary( timestamp: Date.now(), }, ]; - const response = await completeSimple( - model, - { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, - { apiKey, headers, signal, maxTokens: 2048 }, + const response = await completeWithProviderRetry( + () => + completeSimple( + model, + { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, + { apiKey, headers, signal, maxTokens: 2048 }, + ), + { policy: retry, signal }, ); if (response.stopReason === "aborted") { return { aborted: true }; diff --git a/packages/coding-agent/src/core/compaction/compaction.ts b/packages/coding-agent/src/core/compaction/compaction.ts index a47152b350..7aeb060d1f 100644 --- a/packages/coding-agent/src/core/compaction/compaction.ts +++ b/packages/coding-agent/src/core/compaction/compaction.ts @@ -14,6 +14,7 @@ import { createCompactionSummaryMessage, createCustomMessage, } from "../messages.js"; +import { completeWithProviderRetry, type ProviderRetryPolicy } from "../provider-retry.js"; import { buildSessionContext, type CompactionEntry, type SessionEntry } from "../session-manager.js"; import { addAssistantUsage, emptyUsage } from "../usage.js"; import { @@ -523,6 +524,7 @@ export async function generateSummary( customInstructions?: string, previousSummary?: string, thinkingLevel?: ThinkingLevel, + retry?: ProviderRetryPolicy, ): Promise { const maxTokens = Math.floor(0.8 * reserveTokens); @@ -549,10 +551,14 @@ export async function generateSummary( ? { maxTokens, signal, apiKey, headers, reasoning: thinkingLevel } : { maxTokens, signal, apiKey, headers }; - const response = await completeSimple( - model, - { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, - completionOptions, + const response = await completeWithProviderRetry( + () => + completeSimple( + model, + { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, + completionOptions, + ), + { policy: retry, signal }, ); if (response.stopReason === "error") { @@ -692,6 +698,7 @@ export async function compact( signal?: AbortSignal, thinkingLevel?: ThinkingLevel, summaryCall: SummaryCallRunner = (call) => call(headers), + retry?: ProviderRetryPolicy, ): Promise { const { firstKeptEntryId, @@ -721,6 +728,7 @@ export async function compact( customInstructions, previousSummary, thinkingLevel, + retry, ), ) : Promise.resolve({ summary: "No prior history." }), @@ -733,6 +741,7 @@ export async function compact( callHeaders, signal, thinkingLevel, + retry, ), ), ]); @@ -750,6 +759,7 @@ export async function compact( customInstructions, previousSummary, thinkingLevel, + retry, ), ); slices.push(result); @@ -788,6 +798,7 @@ async function generateTurnPrefixSummary( headers?: Record, signal?: AbortSignal, thinkingLevel?: ThinkingLevel, + retry?: ProviderRetryPolicy, ): Promise { const maxTokens = Math.floor(0.5 * reserveTokens); // Smaller budget for turn prefix const llmMessages = convertToLlm(messages); @@ -801,12 +812,16 @@ async function generateTurnPrefixSummary( }, ]; - const response = await completeSimple( - model, - { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, - model.reasoning && thinkingLevel && thinkingLevel !== "off" - ? { maxTokens, signal, apiKey, headers, reasoning: thinkingLevel } - : { maxTokens, signal, apiKey, headers }, + const response = await completeWithProviderRetry( + () => + completeSimple( + model, + { systemPrompt: SUMMARIZATION_SYSTEM_PROMPT, messages: summarizationMessages }, + model.reasoning && thinkingLevel && thinkingLevel !== "off" + ? { maxTokens, signal, apiKey, headers, reasoning: thinkingLevel } + : { maxTokens, signal, apiKey, headers }, + ), + { policy: retry, signal }, ); if (response.stopReason === "error") { diff --git a/packages/coding-agent/src/core/provider-retry.ts b/packages/coding-agent/src/core/provider-retry.ts new file mode 100644 index 0000000000..b69c18afb5 --- /dev/null +++ b/packages/coding-agent/src/core/provider-retry.ts @@ -0,0 +1,127 @@ +import type { AssistantMessage } from "@earendil-works/pi-ai"; +import { sleep } from "../utils/sleep.js"; +import type { SettingsManager } from "./settings-manager.js"; + +/** + * The single retry policy (permanent kinds, Retry-After-aware capped delays), + * shared by the AgentSession auto-retry loop and the one-shot completion + * consumers (side questions, compaction, refinement, session summaries). + */ +export interface ProviderRetryPolicy { + enabled: boolean; + maxRetries: number; + baseDelayMs: number; + /** Max server-requested retry delay before giving up; 0 disables the cap. */ + maxRetryDelayMs: number; +} + +export function providerRetryPolicy(settingsManager: SettingsManager): ProviderRetryPolicy { + return { + ...settingsManager.getRetrySettings(), + maxRetryDelayMs: settingsManager.getProviderRetrySettings().maxRetryDelayMs, + }; +} + +/** Local listener/lifecycle crashes are not provider failures; never retry them. */ +export function isAgentLifecycleFailure(message: AssistantMessage): boolean { + return message.diagnostics?.some((diagnostic) => diagnostic.type === "agent_lifecycle_failure") ?? false; +} + +/** The faux test provider's queue running dry is deterministic; retrying it only stalls tests. */ +export function isFauxProviderQueueExhausted(message: AssistantMessage): boolean { + return message.provider === "faux" && message.errorMessage === "No more faux responses queued"; +} + +export function providerStreamFailureDetails(message: AssistantMessage): Record | undefined { + const failure = message.diagnostics?.find((diagnostic) => diagnostic.type === "provider_stream_failure"); + const details = failure?.details; + if (!details || typeof details !== "object") { + return undefined; + } + return details; +} + +export function providerStreamFailureKind(message: AssistantMessage): string | undefined { + const kind = providerStreamFailureDetails(message)?.kind; + return typeof kind === "string" ? kind : undefined; +} + +export function providerStreamFailureRetryAfterMs(message: AssistantMessage): number | undefined { + const value = providerStreamFailureDetails(message)?.retryAfterMs; + return typeof value === "number" && value >= 0 ? value : undefined; +} + +/** Deterministic rejections never retry; auth gets one retry before it can be marked stale. */ +export function isPermanentProviderFailureKind(kind: string | undefined, retriesPerformed: number): boolean { + if (kind === "invalid_request" || kind === "refusal") { + return true; + } + return retriesPerformed > 0 && kind === "auth"; +} + +export type ProviderRetryDelay = { kind: "wait"; delayMs: number } | { kind: "exceeds-cap"; retryAfterMs: number }; + +/** Node caps timers at 2^31-1 ms; longer delays overflow setTimeout and fire after ~1ms. */ +const MAX_TIMER_DELAY_MS = 2_147_483_647; + +/** Delay before retry `attempt` (1-based), honoring a server-requested wait. */ +export function providerRetryDelay( + attempt: number, + retryAfterMs: number | undefined, + policy: Pick, +): ProviderRetryDelay { + if (retryAfterMs !== undefined && policy.maxRetryDelayMs > 0 && retryAfterMs > policy.maxRetryDelayMs) { + return { kind: "exceeds-cap", retryAfterMs }; + } + return { + kind: "wait", + delayMs: Math.min(Math.max(policy.baseDelayMs * 2 ** (attempt - 1), retryAfterMs ?? 0), MAX_TIMER_DELAY_MS), + }; +} + +/** + * One-shot completion with the shared retry policy, for consumers outside the + * AgentSession auto-retry loop (provider SDKs never retry internally). + */ +export async function completeWithProviderRetry( + attemptCompletion: () => Promise, + options?: { policy?: ProviderRetryPolicy; signal?: AbortSignal }, +): Promise { + const policy = options?.policy ?? DEFAULT_PROVIDER_RETRY_POLICY; + const maxRetries = policy.enabled ? policy.maxRetries : 0; + let retriesPerformed = 0; + for (;;) { + const message = await attemptCompletion(); + if (message.stopReason !== "error") { + return message; + } + if (options?.signal?.aborted) { + // A cancel that raced the failure is an abort, not a provider failure. + return { ...message, stopReason: "aborted" }; + } + if (retriesPerformed >= maxRetries || isAgentLifecycleFailure(message) || isFauxProviderQueueExhausted(message)) { + return message; + } + const kind = providerStreamFailureKind(message); + if (isPermanentProviderFailureKind(kind, retriesPerformed)) { + return message; + } + const delay = providerRetryDelay(retriesPerformed + 1, providerStreamFailureRetryAfterMs(message), policy); + if (delay.kind === "exceeds-cap") { + return message; + } + try { + await sleep(delay.delayMs, options?.signal); + } catch { + return { ...message, stopReason: "aborted" }; + } + retriesPerformed++; + } +} + +export const DEFAULT_PROVIDER_RETRY_POLICY: ProviderRetryPolicy = { + enabled: true, + maxRetries: 3, + baseDelayMs: 2000, + maxRetryDelayMs: 60000, +}; diff --git a/packages/coding-agent/src/core/refinement/refinement.ts b/packages/coding-agent/src/core/refinement/refinement.ts index 76cc701556..c99197e1fc 100644 --- a/packages/coding-agent/src/core/refinement/refinement.ts +++ b/packages/coding-agent/src/core/refinement/refinement.ts @@ -16,6 +16,7 @@ import { completeSimple } from "@earendil-works/pi-ai"; import { getAgentDir } from "../../config.js"; import { serializeConversation } from "../compaction/utils.js"; import { convertToLlm } from "../messages.js"; +import { completeWithProviderRetry, type ProviderRetryPolicy } from "../provider-retry.js"; import type { CustomEntry } from "../session-manager.js"; export const REFINEMENT_CUSTOM_TYPE = "prime-agent.refinement"; @@ -105,6 +106,7 @@ export interface RefineOptions { instructions?: string; rollbackId?: string; global?: boolean; + retry?: ProviderRetryPolicy; } export type AutoRefineReason = "turn_interval" | "compact"; @@ -924,13 +926,17 @@ export async function planRefinement( // Keep the refinement request non-reasoning regardless of the interactive session // thinking level so the model uses its output budget for the JSON object. void thinkingLevel; - const response = await completeSimple( - model, - { - systemPrompt: REFINEMENT_SYSTEM_PROMPT, - messages: [{ role: "user", content: [{ type: "text", text: userPrompt }], timestamp: Date.now() }], - }, - { maxTokens: refinementMaxOutputTokens(model), signal, apiKey, headers }, + const response = await completeWithProviderRetry( + () => + completeSimple( + model, + { + systemPrompt: REFINEMENT_SYSTEM_PROMPT, + messages: [{ role: "user", content: [{ type: "text", text: userPrompt }], timestamp: Date.now() }], + }, + { maxTokens: refinementMaxOutputTokens(model), signal, apiKey, headers }, + ), + { policy: options.retry, signal }, ); if (response.stopReason === "error") { @@ -970,6 +976,7 @@ export async function reviewAutoRefine( headers?: Record, signal?: AbortSignal, thinkingLevel?: ThinkingLevel, + retry?: ProviderRetryPolicy, ): Promise { const conversationText = serializeConversation(convertToLlm(messages)).slice(-40_000); const userPrompt = [ @@ -990,13 +997,17 @@ ${conversationText} // Auto-refine review requires parseable JSON. Keep it non-reasoning so // reasoning-capable models use final text budget for the JSON object. void thinkingLevel; - const response = await completeSimple( - model, - { - systemPrompt: AUTO_REFINE_REVIEW_SYSTEM_PROMPT, - messages: [{ role: "user", content: [{ type: "text", text: userPrompt }], timestamp: Date.now() }], - }, - { maxTokens: autoRefineReviewMaxOutputTokens(model), signal, apiKey, headers }, + const response = await completeWithProviderRetry( + () => + completeSimple( + model, + { + systemPrompt: AUTO_REFINE_REVIEW_SYSTEM_PROMPT, + messages: [{ role: "user", content: [{ type: "text", text: userPrompt }], timestamp: Date.now() }], + }, + { maxTokens: autoRefineReviewMaxOutputTokens(model), signal, apiKey, headers }, + ), + { policy: retry, signal }, ); if (response.stopReason === "error") { throw new Error(`Auto-refine review failed: ${response.errorMessage || "Unknown error"}`); diff --git a/packages/coding-agent/src/core/sdk.ts b/packages/coding-agent/src/core/sdk.ts index 0202492932..723209fef0 100644 --- a/packages/coding-agent/src/core/sdk.ts +++ b/packages/coding-agent/src/core/sdk.ts @@ -293,8 +293,6 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} ...options, apiKey: auth.apiKey, timeoutMs: options?.timeoutMs ?? providerRetrySettings.timeoutMs, - maxRetries: options?.maxRetries ?? providerRetrySettings.maxRetries, - maxRetryDelayMs: options?.maxRetryDelayMs ?? providerRetrySettings.maxRetryDelayMs, headers: auth.headers || options?.headers ? { ...auth.headers, ...options?.headers } : undefined, }); }, @@ -326,7 +324,6 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} followUpMode: settingsManager.getFollowUpMode(), transport: settingsManager.getTransport(), thinkingBudgets: settingsManager.getThinkingBudgets(), - maxRetryDelayMs: settingsManager.getProviderRetrySettings().maxRetryDelayMs, }); if (hasExistingSession) { diff --git a/packages/coding-agent/src/core/settings-manager.ts b/packages/coding-agent/src/core/settings-manager.ts index e2b7ee8e77..350739d89c 100644 --- a/packages/coding-agent/src/core/settings-manager.ts +++ b/packages/coding-agent/src/core/settings-manager.ts @@ -29,8 +29,7 @@ export interface AutoRefineSettings { export interface ProviderRetrySettings { timeoutMs?: number; // SDK/provider request timeout in milliseconds - maxRetries?: number; // SDK/provider retry attempts - maxRetryDelayMs?: number; // default: 60000 (max server-requested delay before failing) + maxRetryDelayMs?: number; // default: 60000 (max server-requested retry delay before failing; 0 disables the cap) } export interface RetrySettings { @@ -951,10 +950,9 @@ export class SettingsManager { }; } - getProviderRetrySettings(): { timeoutMs?: number; maxRetries?: number; maxRetryDelayMs: number } { + getProviderRetrySettings(): { timeoutMs?: number; maxRetryDelayMs: number } { return { timeoutMs: this.settings.retry?.provider?.timeoutMs, - maxRetries: this.settings.retry?.provider?.maxRetries, maxRetryDelayMs: this.settings.retry?.provider?.maxRetryDelayMs ?? 60000, }; } diff --git a/packages/coding-agent/src/core/side-question.ts b/packages/coding-agent/src/core/side-question.ts index 9a3d0e3f0b..80e561fbf0 100644 --- a/packages/coding-agent/src/core/side-question.ts +++ b/packages/coding-agent/src/core/side-question.ts @@ -1,5 +1,10 @@ import { Agent, type AgentMessage } from "@earendil-works/pi-agent-core"; import type { AssistantMessage, UserMessage } from "@earendil-works/pi-ai"; +import { + completeWithProviderRetry, + DEFAULT_PROVIDER_RETRY_POLICY, + type ProviderRetryPolicy, +} from "./provider-retry.js"; import { unwrapSemanticEdgeStreamFn } from "./semantic-edges.js"; export type SideQuestionStatus = "running" | "complete" | "cancelled" | "error"; @@ -46,6 +51,7 @@ export function startSideQuestion( question: string, onEvent: (event: SideQuestionEvent) => void | Promise, previousTurns: SideQuestionTurn[] = [], + retry: ProviderRetryPolicy = DEFAULT_PROVIDER_RETRY_POLICY, ): SideQuestionRun { const model = parent.state.model; if (!model) { @@ -99,13 +105,13 @@ export function startSideQuestion( sessionId: parent.sessionId, thinkingBudgets: parent.thinkingBudgets, transport: "sse", - maxRetryDelayMs: parent.maxRetryDelayMs, toolExecution: parent.toolExecution, }); let answer = ""; let abortRequested = false; let started = false; + const retryAbortController = new AbortController(); const emit = (status: SideQuestionStatus, errorMessage?: string) => onEvent({ id, question, answer, status, ...(errorMessage ? { errorMessage } : {}) }); @@ -130,7 +136,26 @@ export function startSideQuestion( return; } started = true; - await sideAgent.prompt(prompt); + // Standalone side agents bypass the session auto-retry loop; retry here instead. + let promptedOnce = false; + await completeWithProviderRetry( + async () => { + if (promptedOnce) { + // Session-loop recovery: drop the failed assistant turn and re-run. + sideAgent.state.messages = sideAgent.state.messages.slice(0, -1); + await sideAgent.continue(); + } else { + promptedOnce = true; + await sideAgent.prompt(prompt); + } + const last = sideAgent.state.messages.at(-1); + if (last?.role !== "assistant") { + throw new Error(sideAgent.state.errorMessage || "Side question produced no assistant message"); + } + return last as AssistantMessage; + }, + { policy: retry, signal: retryAbortController.signal }, + ); if (abortRequested) { await emit("cancelled"); return; @@ -153,6 +178,7 @@ export function startSideQuestion( done, abort() { abortRequested = true; + retryAbortController.abort(); if (started) { sideAgent.abort(); } diff --git a/packages/coding-agent/src/modes/agent-connection/in-process-agent-connection.ts b/packages/coding-agent/src/modes/agent-connection/in-process-agent-connection.ts index b3d4986e28..ba98dec0c4 100644 --- a/packages/coding-agent/src/modes/agent-connection/in-process-agent-connection.ts +++ b/packages/coding-agent/src/modes/agent-connection/in-process-agent-connection.ts @@ -15,6 +15,7 @@ import type { } from "../../core/cron-jobs.js"; import type { ExtensionUIContext } from "../../core/extensions/types.js"; import type { AcpMcpServerConfig } from "../../core/mcp/acp-mcp-types.js"; +import { providerRetryPolicy } from "../../core/provider-retry.js"; import type { RefinementResult } from "../../core/refinement/index.js"; import { type DeleteSessionFileResult, deleteSessionFile } from "../../core/session-file-actions.js"; import { SessionManager } from "../../core/session-manager.js"; @@ -392,6 +393,7 @@ export class InProcessAgentConnection implements AgentConnection { question, (event) => this.emit({ type: "side_question_event", event }), previousTurns, + providerRetryPolicy(this.session.settingsManager), ); this.sideQuestionRuns.set(id, run); const removeRun = () => { diff --git a/packages/coding-agent/src/modes/daemon/daemon-mode.ts b/packages/coding-agent/src/modes/daemon/daemon-mode.ts index a052ee863f..26ce08d2a5 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-mode.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-mode.ts @@ -89,6 +89,7 @@ import { } from "../../core/cron-jobs.js"; import { ORPHAN_PROCESS_JOURNAL_ENV } from "../../core/orphan-process-journal.js"; import { PromptAdmissionCancelledError, waitForPromptAdmission } from "../../core/prompt-admission.js"; +import { providerRetryPolicy } from "../../core/provider-retry.js"; import type { CreateRlmSubagentRuntimeOptions, SubagentRuntimeHost } from "../../core/rlm-runtime.js"; import { canPassivateSession, @@ -4395,6 +4396,7 @@ export class AgentDaemon { } }, command.previousTurns, + providerRetryPolicy(state.runtime.session.settingsManager), ); this.sideQuestionRuns.set(command.sideQuestionId, { run, diff --git a/packages/coding-agent/src/modes/daemon/daemon-session-summarizer.ts b/packages/coding-agent/src/modes/daemon/daemon-session-summarizer.ts index f2dba51f67..99328286fb 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-session-summarizer.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-session-summarizer.ts @@ -2,6 +2,7 @@ import type { AgentMessage } from "@earendil-works/pi-agent-core"; import type { Api, Model } from "@earendil-works/pi-ai"; import { completeSimple } from "@earendil-works/pi-ai"; import type { ModelRegistry } from "../../core/model-registry.js"; +import { completeWithProviderRetry } from "../../core/provider-retry.js"; import type { AgentStatus, AgentTaskState } from "../../core/session-manager.js"; import type { ActiveSessionState } from "./active-session-state.js"; @@ -162,19 +163,24 @@ export async function generateAgentStatus(params: GenerateAgentStatusParams): Pr return undefined; } try { - const response = await completeSimple( - model, - { - systemPrompt: AGENT_STATUS_SYSTEM_PROMPT, - messages: [ + // One failed attempt would settle an idle session to a stale needs_input verdict. + const response = await completeWithProviderRetry( + () => + completeSimple( + model, { - role: "user" as const, - content: [{ type: "text" as const, text: buildStatusContext(messages, isWorking) }], - timestamp: Date.now(), + systemPrompt: AGENT_STATUS_SYSTEM_PROMPT, + messages: [ + { + role: "user" as const, + content: [{ type: "text" as const, text: buildStatusContext(messages, isWorking) }], + timestamp: Date.now(), + }, + ], }, - ], - }, - { maxTokens: SUMMARY_MAX_TOKENS, apiKey: auth.apiKey, headers: auth.headers, signal }, + { maxTokens: SUMMARY_MAX_TOKENS, apiKey: auth.apiKey, headers: auth.headers, signal }, + ), + { signal }, ); if (response.stopReason === "error") { return undefined; diff --git a/packages/coding-agent/test/acp-cold-cli.test.ts b/packages/coding-agent/test/acp-cold-cli.test.ts index 38a3760e1a..5e6b5568eb 100644 --- a/packages/coding-agent/test/acp-cold-cli.test.ts +++ b/packages/coding-agent/test/acp-cold-cli.test.ts @@ -27,7 +27,9 @@ afterEach(async () => { await new Promise((done) => server.close(() => done())); } for (const dir of tempDirs.splice(0)) { - rmSync(dir, { recursive: true, force: true }); + // The CLI's daemon/kernel children may still be flushing caches under this + // dir when the child exits; retry the ENOTEMPTY/EBUSY window instead of failing cleanup. + rmSync(dir, { recursive: true, force: true, maxRetries: 10, retryDelay: 100 }); } }); @@ -56,6 +58,9 @@ async function driveAcpTurn(baseUrl: string): Promise { const projectDir = join(tempRoot, "project"); mkdirSync(agentDir, { recursive: true }); mkdirSync(projectDir, { recursive: true }); + // One deterministic failed attempt: the session-layer auto-retry would + // otherwise re-issue the rejected request before reporting the failure. + writeFileSync(join(agentDir, "settings.json"), JSON.stringify({ retry: { enabled: false } }), "utf-8"); writeFileSync( join(agentDir, "models.json"), JSON.stringify({ @@ -83,6 +88,7 @@ async function driveAcpTurn(baseUrl: string): Promise { "--model", "cold/test-model", "--no-session", + "--no-tools", "--offline", "--daemon-socket", join(tempRoot, "d.sock"), diff --git a/packages/coding-agent/test/provider-retry.test.ts b/packages/coding-agent/test/provider-retry.test.ts new file mode 100644 index 0000000000..8b28ae2991 --- /dev/null +++ b/packages/coding-agent/test/provider-retry.test.ts @@ -0,0 +1,60 @@ +import type { AssistantMessage } from "@earendil-works/pi-ai"; +import { describe, expect, it } from "vitest"; +import { completeWithProviderRetry, providerRetryDelay } from "../src/core/provider-retry.js"; + +function providerError(): AssistantMessage { + return { + role: "assistant", + content: [], + api: "openai-completions", + provider: "openai", + model: "test-model", + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "error", + errorMessage: "500 Internal Server Error", + timestamp: Date.now(), + }; +} + +describe("completeWithProviderRetry", () => { + it("returns an aborted result instead of the provider error when cancelled during backoff", async () => { + const controller = new AbortController(); + setTimeout(() => controller.abort(), 10); + + const result = await completeWithProviderRetry(async () => providerError(), { + policy: { enabled: true, maxRetries: 3, baseDelayMs: 60_000, maxRetryDelayMs: 0 }, + signal: controller.signal, + }); + + expect(result.stopReason).toBe("aborted"); + }); + + it("makes a single attempt when the policy disables retries", async () => { + let attempts = 0; + const result = await completeWithProviderRetry( + async () => { + attempts++; + return providerError(); + }, + { policy: { enabled: false, maxRetries: 3, baseDelayMs: 1, maxRetryDelayMs: 60_000 } }, + ); + + expect(attempts).toBe(1); + expect(result.stopReason).toBe("error"); + }); + + it("clamps uncapped server delays to Node's max timer instead of overflowing setTimeout", () => { + const ninetyDaysMs = 90 * 24 * 3600 * 1000; + expect(providerRetryDelay(1, ninetyDaysMs, { baseDelayMs: 2000, maxRetryDelayMs: 0 })).toEqual({ + kind: "wait", + delayMs: 2_147_483_647, + }); + }); +}); diff --git a/packages/coding-agent/test/suite/agent-session-retry-events.test.ts b/packages/coding-agent/test/suite/agent-session-retry-events.test.ts index 6e08632645..eb1ea99fbc 100644 --- a/packages/coding-agent/test/suite/agent-session-retry-events.test.ts +++ b/packages/coding-agent/test/suite/agent-session-retry-events.test.ts @@ -37,6 +37,19 @@ function structuredProviderFailure(kind: "auth" | "invalid_request" | "refusal") }; } +function rateLimitedFailure(retryAfterMs: number): AssistantMessage { + return { + ...fauxAssistantMessage("", { stopReason: "error", errorMessage: "429 rate limited" }), + diagnostics: [ + { + type: "provider_stream_failure", + timestamp: Date.now(), + details: { kind: "rate_limit", status: 429, retryAfterMs }, + }, + ], + }; +} + type SessionRetryCompactionInternals = { _retryAttempt: number; _retryPromise: Promise | undefined; @@ -264,25 +277,69 @@ describe("AgentSession retry and event characterization", () => { expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([true]); }); - for (const kind of ["auth", "invalid_request", "refusal"] as const) { - it(`retries structured permanent provider ${kind} failures once`, async () => { + it("retries structured provider auth failures once", async () => { + const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } } }); + harnesses.push(harness); + harness.setResponses([ + structuredProviderFailure("auth"), + structuredProviderFailure("auth"), + fauxAssistantMessage("unused"), + ]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(2); + expect(harness.eventsOfType("auto_retry_start").map((event) => event.attempt)).toEqual([1]); + expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([false]); + expect(harness.session.isRetrying).toBe(false); + }); + + for (const kind of ["invalid_request", "refusal"] as const) { + it(`does not retry structured permanent provider ${kind} failures`, async () => { const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } } }); harnesses.push(harness); - harness.setResponses([ - structuredProviderFailure(kind), - structuredProviderFailure(kind), - fauxAssistantMessage("unused"), - ]); + harness.setResponses([structuredProviderFailure(kind), fauxAssistantMessage("unused")]); await harness.session.prompt("test"); - expect(harness.faux.state.callCount).toBe(2); - expect(harness.eventsOfType("auto_retry_start").map((event) => event.attempt)).toEqual([1]); - expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([false]); + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); expect(harness.session.isRetrying).toBe(false); }); } + it("waits at least the provider-requested Retry-After delay before retrying", async () => { + const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } } }); + harnesses.push(harness); + harness.setResponses([rateLimitedFailure(50), fauxAssistantMessage("recovered")]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(2); + expect(harness.eventsOfType("auto_retry_start").map((event) => event.delayMs)).toEqual([50]); + expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([true]); + }); + + it("fails without retrying when the provider-requested delay exceeds maxRetryDelayMs", async () => { + const harness = await createHarness({ + settings: { + retry: { enabled: true, maxRetries: 3, baseDelayMs: 1, provider: { maxRetryDelayMs: 100 } }, + }, + }); + harnesses.push(harness); + harness.setResponses([rateLimitedFailure(3_600_000), fauxAssistantMessage("unused")]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); + const retryEnd = harness.eventsOfType("auto_retry_end"); + expect(retryEnd).toHaveLength(1); + expect(retryEnd[0]?.success).toBe(false); + expect(retryEnd[0]?.finalError).toContain("maxRetryDelayMs"); + expect(harness.session.isRetrying).toBe(false); + }); + it("keeps retry state active when overflow compaction will retry", async () => { const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } } }); harnesses.push(harness); diff --git a/packages/coding-agent/test/suite/regressions/4509-side-questions.test.ts b/packages/coding-agent/test/suite/regressions/4509-side-questions.test.ts index e10da95722..99520639e8 100644 --- a/packages/coding-agent/test/suite/regressions/4509-side-questions.test.ts +++ b/packages/coding-agent/test/suite/regressions/4509-side-questions.test.ts @@ -113,6 +113,38 @@ describe("ENG-4509 side questions", () => { } }); + it("retries a transient provider error once and completes", async () => { + const harness = await createHarness(); + try { + harness.setResponses([fauxAssistantMessage("main answer")]); + await harness.session.prompt("Main context message."); + + harness.setResponses([ + fauxAssistantMessage("", { stopReason: "error", errorMessage: "500 Internal Server Error" }), + fauxAssistantMessage("recovered side answer"), + ]); + const callsBefore = harness.faux.state.callCount; + + const events: SideQuestionEvent[] = []; + const run = startSideQuestion( + harness.session.agent, + "retry-1", + "Does this survive a transient failure?", + (event) => { + events.push(event); + }, + [], + { enabled: true, maxRetries: 3, baseDelayMs: 1, maxRetryDelayMs: 60_000 }, + ); + await run.done; + + expect(harness.faux.state.callCount - callsBefore).toBe(2); + expect(events.at(-1)).toMatchObject({ status: "complete", answer: "recovered side answer" }); + } finally { + harness.cleanup(); + } + }); + it("can finish while the main agent is still working", async () => { const harness = await createHarness(); const mainStarted = deferred();