From aaa96f17f3a9ee16102ed632b79c066708e57535 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sun, 30 Aug 2026 22:09:33 -0600 Subject: [PATCH 1/2] fix(agent): kind-aware retry with Retry-After and a real delay cap Rate limits, overloads, and auth failures all retried on the same fixed 2s/4s/8s ladder while provider SDK retries stacked underneath, so one failed turn could produce sessionRetries x sdkRetries upstream requests per agent. maxRetryDelayMs was threaded through StreamOptions, Agent, the proxy, and settings without any provider reading it. The session loop is now the only retry layer: provider maxRetries defaults to 0 while retry.enabled is true, and retry.provider.maxRetryDelayMs is the enforced ceiling on provider-requested waits. rate_limit and overloaded back off from a longer base with full jitter, honor an upstream Retry-After, and fail fast when the requested wait exceeds the cap; auth aborts the ladder immediately; every other kind keeps today's behavior. Retry-After now travels in the provider_stream_failure diagnostic, including Anthropic overloads delivered as mid-stream SSE errors. The dead maxRetryDelayMs plumbing is removed. fixes #24 --- .pylon/features.yaml | 15 ++ .pylon/upstream-review.md | 11 + .../agent/.changes/24-kind-aware-retry.md | 1 + packages/agent/src/agent.ts | 4 - packages/agent/src/proxy.ts | 2 - packages/ai/.changes/24-kind-aware-retry.md | 2 + packages/ai/src/providers/anthropic.ts | 9 +- packages/ai/src/providers/simple-options.ts | 1 - packages/ai/src/types.ts | 8 - packages/ai/src/utils/stream-failure.ts | 44 +++- .../ai/test/anthropic-sse-parsing.test.ts | 20 ++ packages/ai/test/stream-failure.test.ts | 69 ++++++ .../.changes/24-kind-aware-retry.md | 3 + packages/coding-agent/docs/settings.md | 8 +- .../coding-agent/src/core/agent-session.ts | 28 ++- .../coding-agent/src/core/retry-backoff.ts | 93 ++++++++ packages/coding-agent/src/core/sdk.ts | 12 +- .../coding-agent/src/core/settings-manager.ts | 4 +- .../coding-agent/src/core/side-question.ts | 1 - .../coding-agent/test/retry-backoff.test.ts | 99 ++++++++ .../suite/agent-session-retry-events.test.ts | 15 +- .../regressions/24-kind-aware-retry.test.ts | 221 ++++++++++++++++++ .../4491-provider-stale-after-401.test.ts | 12 +- 23 files changed, 640 insertions(+), 42 deletions(-) create mode 100644 packages/agent/.changes/24-kind-aware-retry.md create mode 100644 packages/ai/.changes/24-kind-aware-retry.md create mode 100644 packages/coding-agent/.changes/24-kind-aware-retry.md create mode 100644 packages/coding-agent/src/core/retry-backoff.ts create mode 100644 packages/coding-agent/test/retry-backoff.test.ts create mode 100644 packages/coding-agent/test/suite/regressions/24-kind-aware-retry.test.ts diff --git a/.pylon/features.yaml b/.pylon/features.yaml index 33797f3c2d..284da3ca0d 100644 --- a/.pylon/features.yaml +++ b/.pylon/features.yaml @@ -281,3 +281,18 @@ decisions: revisit_when: - Prime routes subagent and derived-request provider hooks through the owning session's extension runner. - An upstream extension contract exposes per-agent provider identity that supersedes the scoped session-id view. + kind-aware-provider-retry: + area: runtime-reliability + state: shipped + owner: pylon-prime-integration + decision: redesign + pylon_refs: + - https://github.com/pylon-code/prime-agent/issues/22 + - https://github.com/pylon-code/prime-agent/issues/24 + upstream_refs: + - https://github.com/PrimeIntellect-ai/prime-agent/tree/a903d4b6768f484bd6d459b7b0aa7dee38e461e2 + fork_change: kind-aware-session-retry-v1 + upstream_support: Prime through a903d4b6768f classifies provider stream failures but retries every one of them on the same fixed 2s/4s/8s ladder, stacks provider SDK retries underneath the session loop, never reads Retry-After, and threads maxRetryDelayMs through StreamOptions, Agent, and the proxy without any provider reading it. + revisit_when: + - Prime upstream makes the session retry loop failure-kind aware and honors Retry-After with a delay ceiling. + - Prime upstream makes exactly one retry layer own policy so provider SDK retries cannot multiply session retries. diff --git a/.pylon/upstream-review.md b/.pylon/upstream-review.md index 2decd48991..f533e0a56f 100644 --- a/.pylon/upstream-review.md +++ b/.pylon/upstream-review.md @@ -109,6 +109,17 @@ This ledger records Prime upstream evidence and the decision taken for each over - Cross-repository merge order: Prime issue #17 and a reproducible artifact, Pylon issue #190 consuming the exact post-attach proof, then Comet issue #7. Exact committed-head API/security/test review and trusted hosted CI remain mandatory. Revisit when upstream offers an equivalent generation-scoped proof and both consumers can remove this token without weakening fail-closed negotiation. +## 2026-08-30 — kind-aware provider retry candidate + +- Upstream baseline: `PrimeIntellect-ai/prime-agent@a903d4b6768f484bd6d459b7b0aa7dee38e461e2`; latest audited release remains `v0.8.1`. This candidate does not advance `reviewed_upstream_commit`. +- Searched current Prime issues, pull requests, and source for retry backoff, `Retry-After`, and `maxRetryDelayMs`. Upstream has no kind-aware retry work and still carries the same dead `maxRetryDelayMs` threading through `packages/ai/src/types.ts`, `packages/ai/src/providers/simple-options.ts`, `packages/agent/src/agent.ts`, `packages/agent/src/proxy.ts`, `packages/coding-agent/src/core/sdk.ts`, and `packages/coding-agent/src/core/side-question.ts`. The behavior is therefore not superseded by an upstream API. +- `kind-aware-provider-retry`: **redesign**. Prime classifies provider stream failures but retries every one on the same fixed 2s/4s/8s ladder, stacks SDK retries under the session loop, and never reads `Retry-After`. Pylon replaces both layers with one policy: the session loop owns retry, provider `maxRetries` defaults to `0` while `retry.enabled` is true, and `retry.provider.maxRetryDelayMs` becomes the enforced ceiling on provider-requested waits instead of a silent no-op setting. +- Policy: `auth` aborts the ladder immediately; `rate_limit` and `overloaded` back off from 15s and 10s with full jitter floored at `retry.baseDelayMs` and capped by `maxRetryDelayMs`, honoring `Retry-After` plus up to 1s of spread when the provider sends one, and failing fast when the requested wait exceeds the cap; every other kind keeps the historical ladder unchanged. +- Wire compatibility: the `auto_retry_start` and `auto_retry_end` session events keep their exact shapes, so there is no daemon protocol, capability, or schema-revision change. `retry.provider.maxRetryDelayMs` keeps its settings key and existing `retry.maxDelayMs` migration. +- Behavior change accepted deliberately: structured auth failures previously retried once before marking auth stale. They now abort on the first failure while still marking the auth source stale and appending login guidance. The affected assertions in `agent-session-retry-events.test.ts` and regression `4491-provider-stale-after-401.test.ts` were updated; the unstructured 401 path continues to exercise mid-backoff credential rotation and cancellation. +- Validation: `npm run check` clean. Focused runs pass `packages/ai` stream-failure 35/35 and Anthropic SSE parsing 7/7; `packages/coding-agent` retry-backoff 10/10, issue #24 regression 7/7, and the retry/auth suites 38/38; `packages/agent` 60/60; settings and telemetry 55/55. +- Revisit when Prime upstream makes its session retry loop failure-kind aware, honors `Retry-After` with a delay ceiling, and gives exactly one layer ownership of retry policy. + ## 2026-08-31 — negotiated proof shipped and mixed-version snapshot catch-up follow-up - Upstream evidence remains fully audited through `PrimeIntellect-ai/prime-agent@a903d4b6768f484bd6d459b7b0aa7dee38e461e2`, the product base used by PR #19. The only later upstream-main commit currently visible is `c382f09856d4a8c8d2b765179657047d58691f25` (PR #1893, terminal Mermaid rendering); its changed paths do not overlap daemon snapshot, worker, supervisor, framing, or recovery code and it does not supersede this boundary. diff --git a/packages/agent/.changes/24-kind-aware-retry.md b/packages/agent/.changes/24-kind-aware-retry.md new file mode 100644 index 0000000000..80fffa962d --- /dev/null +++ b/packages/agent/.changes/24-kind-aware-retry.md @@ -0,0 +1 @@ +- Removed the unused `maxRetryDelayMs` agent and proxy option; it was threaded to providers that never read it ([#24](https://github.com/pylon-code/prime-agent/issues/24)). 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/24-kind-aware-retry.md b/packages/ai/.changes/24-kind-aware-retry.md new file mode 100644 index 0000000000..9c01aeb320 --- /dev/null +++ b/packages/ai/.changes/24-kind-aware-retry.md @@ -0,0 +1,2 @@ +- Added upstream `Retry-After` capture to provider stream-failure diagnostics, including Anthropic overloads delivered as mid-stream SSE errors, so higher-level retry can honor it ([#24](https://github.com/pylon-code/prime-agent/issues/24)). +- Removed the unused `maxRetryDelayMs` stream option; no provider ever read it ([#24](https://github.com/pylon-code/prime-agent/issues/24)). diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index 55596aa0c3..50fac7f728 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -36,6 +36,7 @@ import { classifyStreamFailure, formatStreamFailureMessage, recordStreamFailure, + retryAfterMsFromHeaders, StreamFailureError, streamFailureFromStopReason, streamFailureMessage, @@ -384,7 +385,7 @@ async function* iterateSseMessages( } /** Turn an in-stream `error` SSE event (how Anthropic delivers overloads etc.) into a classified failure. */ -function anthropicSseError(data: string, requestId?: string): StreamFailureError { +function anthropicSseError(data: string, requestId?: string, retryAfterMs?: number): StreamFailureError { let errorType: string | undefined; let detail: string | undefined; try { @@ -400,6 +401,7 @@ function anthropicSseError(data: string, requestId?: string): StreamFailureError kind: classifyStreamFailure(errorType), providerErrorType: errorType, requestId, + retryAfterMs, raw: truncateRawPayload(data), }; return new StreamFailureError(streamFailureMessage(info, detail), info); @@ -416,10 +418,13 @@ async function* iterateAnthropicEvents( let sawMessageStart = false; let sawMessageEnd = false; + // Overload/rate-limit SSE errors arrive on an otherwise-200 response, so the + // only Retry-After the caller ever sees is the one on the response headers. + const retryAfterMs = retryAfterMsFromHeaders(response.headers); for await (const sse of iterateSseMessages(response.body, signal)) { if (sse.event === "error") { - throw anthropicSseError(sse.data, requestId); + throw anthropicSseError(sse.data, requestId, retryAfterMs); } if (!ANTHROPIC_MESSAGE_EVENTS.has(sse.event ?? "")) { diff --git a/packages/ai/src/providers/simple-options.ts b/packages/ai/src/providers/simple-options.ts index 76f6ea8b44..e3448afcf4 100644 --- a/packages/ai/src/providers/simple-options.ts +++ b/packages/ai/src/providers/simple-options.ts @@ -15,7 +15,6 @@ export function buildBaseOptions(model: Model, options?: SimpleStreamOption 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..5bb3034865 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -123,14 +123,6 @@ export interface StreamOptions { * 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..76df320ded 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; + /** Upstream `Retry-After` translated to milliseconds, when the provider sent one. */ + retryAfterMs?: number; /** Truncated raw provider payload for post-mortems. */ raw?: string; } @@ -114,6 +116,39 @@ export function truncateRawPayload(raw: string): string { return raw.length > MAX_RAW_LENGTH ? `${raw.slice(0, MAX_RAW_LENGTH)}…` : raw; } +/** + * Translate an HTTP `Retry-After` value (delta-seconds or HTTP-date) into a + * non-negative millisecond wait. Returns undefined when the header is absent or + * unparseable so callers fall back to their own backoff. + */ +export function parseRetryAfterMs(value: string | null | undefined, now = Date.now()): number | undefined { + const raw = value?.trim(); + if (!raw) return undefined; + const seconds = Number(raw); + if (Number.isFinite(seconds)) return seconds > 0 ? Math.round(seconds * 1000) : 0; + const retryAt = Date.parse(raw); + return Number.isFinite(retryAt) ? Math.max(0, retryAt - now) : undefined; +} + +type HeaderSource = Headers | Record; + +function readHeader(headers: unknown, name: string): string | undefined { + if (!headers || typeof headers !== "object") return undefined; + const source = headers as HeaderSource; + if (typeof (source as Headers).get === "function") { + return (source as Headers).get(name) ?? undefined; + } + const record = source as Record; + const key = Object.keys(record).find((candidate) => candidate.toLowerCase() === name); + const value = key === undefined ? undefined : record[key]; + return typeof value === "string" ? value : undefined; +} + +/** Read `Retry-After` off any header carrier a provider SDK error might expose. */ +export function retryAfterMsFromHeaders(headers: unknown, now = Date.now()): number | undefined { + return parseRetryAfterMs(readHeader(headers, "retry-after"), now); +} + function extractStreamFailureParts(error: unknown): { info: StreamFailureInfo; detail?: string } { if (error instanceof StreamFailureError) return { info: error.info }; if (!(error instanceof Error)) return { info: { kind: "unknown" } }; @@ -150,13 +185,7 @@ 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 = readHeader(headers, "request-id") ?? readHeader(headers, "x-request-id"); const rawRequestId = err.requestID ?? err.request_id ?? err.$metadata?.requestId ?? headerRequestId; const requestId = typeof rawRequestId === "string" ? rawRequestId : undefined; @@ -166,6 +195,7 @@ function extractStreamFailureParts(error: unknown): { info: StreamFailureInfo; d providerErrorType, status, requestId, + retryAfterMs: retryAfterMsFromHeaders(headers), }, detail: typeof bodyMessage === "string" ? bodyMessage : undefined, }; diff --git a/packages/ai/test/anthropic-sse-parsing.test.ts b/packages/ai/test/anthropic-sse-parsing.test.ts index 147eb874d4..08f5f93c74 100644 --- a/packages/ai/test/anthropic-sse-parsing.test.ts +++ b/packages/ai/test/anthropic-sse-parsing.test.ts @@ -153,6 +153,26 @@ describe("Anthropic raw SSE parsing", () => { expect(result.usage.cost.cacheWrite).toBeCloseTo(testCase.expectedCacheWriteCost); }); + it("carries Retry-After from an overload delivered as a mid-stream SSE error", async () => { + const model = getModel("anthropic", "claude-haiku-4-5"); + const response = new Response( + `event: error\ndata: ${JSON.stringify({ type: "error", error: { type: "overloaded_error", message: "Overloaded" } })}\n`, + { status: 200, headers: { "content-type": "text/event-stream", "retry-after": "30" } }, + ); + + const result = await streamAnthropic( + model, + { messages: [{ role: "user", content: "Say hello.", timestamp: Date.now() }] }, + { client: createFakeAnthropicClient(response) }, + ).result(); + + expect(result.stopReason).toBe("error"); + expect(result.diagnostics?.[0]).toMatchObject({ + type: "provider_stream_failure", + details: { kind: "overloaded", retryAfterMs: 30_000 }, + }); + }); + it("preserves configured cache write pricing for non-Anthropic models", async () => { const model = getModel("minimax", "MiniMax-M2.7-highspeed"); const response = createSseResponse( diff --git a/packages/ai/test/stream-failure.test.ts b/packages/ai/test/stream-failure.test.ts index 5009de0160..7cf06e1d6d 100644 --- a/packages/ai/test/stream-failure.test.ts +++ b/packages/ai/test/stream-failure.test.ts @@ -5,7 +5,9 @@ import { classifyStreamFailure, extractStreamFailureInfo, formatStreamFailureMessage, + parseRetryAfterMs, recordStreamFailure, + retryAfterMsFromHeaders, StreamFailureError, streamFailureFromStopReason, } from "../src/utils/stream-failure.js"; @@ -102,12 +104,64 @@ describe("extractStreamFailureInfo", () => { expect(extractStreamFailureInfo(awsError)).toMatchObject({ requestId: "aws_req" }); }); + test("carries Retry-After from rate-limit response headers", () => { + const sdkError = Object.assign(new Error("429 rate limited"), { + status: 429, + headers: new Headers({ "retry-after": "30", "request-id": "req_429" }), + }); + expect(extractStreamFailureInfo(sdkError)).toMatchObject({ + kind: "rate_limit", + status: 429, + requestId: "req_429", + retryAfterMs: 30_000, + }); + }); + + test("leaves retryAfterMs unset when the provider sent no Retry-After", () => { + const sdkError = Object.assign(new Error("529 overloaded"), { status: 529 }); + expect(extractStreamFailureInfo(sdkError).retryAfterMs).toBeUndefined(); + }); + 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"); }); }); +describe("parseRetryAfterMs", () => { + test("reads delta-seconds", () => { + expect(parseRetryAfterMs("30")).toBe(30_000); + expect(parseRetryAfterMs(" 1.5 ")).toBe(1500); + }); + + test("clamps non-positive delta-seconds to zero", () => { + expect(parseRetryAfterMs("0")).toBe(0); + expect(parseRetryAfterMs("-5")).toBe(0); + }); + + test("reads an HTTP-date relative to now", () => { + const now = Date.parse("2026-01-01T00:00:00Z"); + expect(parseRetryAfterMs("Thu, 01 Jan 2026 00:00:45 GMT", now)).toBe(45_000); + expect(parseRetryAfterMs("Wed, 31 Dec 2025 23:59:00 GMT", now)).toBe(0); + }); + + test("returns undefined for missing or unparseable values", () => { + expect(parseRetryAfterMs(undefined)).toBeUndefined(); + expect(parseRetryAfterMs(null)).toBeUndefined(); + expect(parseRetryAfterMs(" ")).toBeUndefined(); + expect(parseRetryAfterMs("soon")).toBeUndefined(); + }); +}); + +describe("retryAfterMsFromHeaders", () => { + test("reads Headers objects and plain records", () => { + expect(retryAfterMsFromHeaders(new Headers({ "retry-after": "12" }))).toBe(12_000); + expect(retryAfterMsFromHeaders({ "Retry-After": "12" })).toBe(12_000); + expect(retryAfterMsFromHeaders({})).toBeUndefined(); + expect(retryAfterMsFromHeaders(undefined)).toBeUndefined(); + }); +}); + describe("formatStreamFailureMessage", () => { test("condenses a classified SDK error to a one-liner instead of the raw payload", () => { const sdkError = Object.assign( @@ -160,6 +214,21 @@ describe("recordStreamFailure", () => { }); }); + test("persists Retry-After in the diagnostic so session retry can honor it", () => { + setLogSink(() => {}); + const output = makeOutput({ errorMessage: "Provider rate limit exceeded" }); + recordStreamFailure( + model, + output, + new StreamFailureError("x", { kind: "rate_limit", status: 429, retryAfterMs: 30_000 }), + ); + + expect(output.diagnostics?.[0]).toMatchObject({ + type: "provider_stream_failure", + details: { kind: "rate_limit", retryAfterMs: 30_000 }, + }); + }); + test("does nothing for user aborts", () => { const logged: unknown[] = []; setLogSink((entry) => logged.push(entry)); diff --git a/packages/coding-agent/.changes/24-kind-aware-retry.md b/packages/coding-agent/.changes/24-kind-aware-retry.md new file mode 100644 index 0000000000..21b5f0b9de --- /dev/null +++ b/packages/coding-agent/.changes/24-kind-aware-retry.md @@ -0,0 +1,3 @@ +- Fixed rate-limit and overload retries hammering an already-degraded provider: they now back off from a longer base with jitter, honor an upstream `Retry-After`, and fail fast when the requested wait exceeds `retry.provider.maxRetryDelayMs` ([#24](https://github.com/pylon-code/prime-agent/issues/24)). +- Changed provider authentication failures to stop retrying immediately instead of burning the retry ladder on a credential that cannot recover by waiting ([#24](https://github.com/pylon-code/prime-agent/issues/24)). +- Changed provider/SDK retries to default to 0 while agent-level retry is enabled, so the two retry layers no longer multiply into repeated upstream requests per failed turn ([#24](https://github.com/pylon-code/prime-agent/issues/24)). diff --git a/packages/coding-agent/docs/settings.md b/packages/coding-agent/docs/settings.md index 8e1c147260..0bf6461a21 100644 --- a/packages/coding-agent/docs/settings.md +++ b/packages/coding-agent/docs/settings.md @@ -141,10 +141,14 @@ 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.maxRetries` | number | `0` while `retry.enabled` | Provider/SDK retry attempts | | `retry.provider.maxRetryDelayMs` | number | `60000` | Max server-requested 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. +Exactly one layer retries. While `retry.enabled` is true, provider/SDK retries default to `0` so a failed turn cannot fan out into `agent retries x SDK retries` upstream requests. Set `retry.provider.maxRetries` explicitly to override. + +Rate-limit and overload failures do not use the plain `retry.baseDelayMs` ladder. They back off from a longer base (15s for rate limits, 10s for overloads), apply jitter so concurrent agents do not resume in lockstep, and honor an upstream `Retry-After` when the provider sends one. Failures the provider did not classify keep the plain ladder. + +When a provider asks for a retry delay longer than `retry.provider.maxRetryDelayMs` (e.g., Google's "quota will reset after 5h"), the turn fails immediately with an informative error instead of waiting silently. Set to `0` to disable the cap. Authentication failures never retry: waiting cannot fix an expired or rejected credential. ```json { diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index e6013046d7..ba50b67fcb 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -226,6 +226,7 @@ import { } from "./refinement/index.js"; import { resolveConfigValue } from "./resolve-config-value.js"; import type { ResourceExtensionPaths, ResourceLoader } from "./resource-loader.js"; +import { computeRetryBackoff } from "./retry-backoff.js"; import { type CreateRlmSubagentRuntimeOptions, createDefaultRlmSubagentSessionName, @@ -10054,7 +10055,6 @@ export class AgentSession { sessionId: childSessionManager.getSessionId(), thinkingBudgets: this.settingsManager.getThinkingBudgets(), transport: this.settingsManager.getTransport(), - maxRetryDelayMs: this.settingsManager.getProviderRetrySettings().maxRetryDelayMs, toolExecution: this.agent.toolExecution, }); @@ -11360,6 +11360,12 @@ export class AgentSession { return typeof kind === "string" ? kind : undefined; } + private _getProviderStreamFailureRetryAfterMs(message: AssistantMessage): number | undefined { + const value = this._getProviderStreamFailureDetails(message)?.retryAfterMs; + const parsed = typeof value === "string" ? Number(value) : typeof value === "number" ? value : Number.NaN; + return Number.isFinite(parsed) && parsed >= 0 ? parsed : undefined; + } + private _isStructuredPermanentProviderFailure(message: AssistantMessage): boolean { const kind = this._getProviderStreamFailureKind(message); return kind === "auth" || kind === "invalid_request" || kind === "refusal"; @@ -11507,8 +11513,24 @@ export class AgentSession { this._retryAttempt++; - if (this._retryAttempt > settings.maxRetries) { + // The ladder is exhausted, or the failure kind/Retry-After says waiting + // cannot help. Either way this is the last word on the turn. + const decision = + this._retryAttempt > settings.maxRetries + ? undefined + : computeRetryBackoff({ + kind: this._getProviderStreamFailureKind(message), + attempt: this._retryAttempt, + baseDelayMs: settings.baseDelayMs, + maxRetryDelayMs: this.settingsManager.getProviderRetrySettings().maxRetryDelayMs, + retryAfterMs: this._getProviderStreamFailureRetryAfterMs(message), + }); + + if (decision === undefined || decision.type === "abort") { this._markProviderAuthStaleForRetryFailure(message, options); + if (decision?.type === "abort" && message.errorMessage) { + message.errorMessage = `${message.errorMessage} (${decision.reason})`; + } this._emit({ type: "auto_retry_end", success: false, @@ -11521,7 +11543,7 @@ export class AgentSession { return false; } - const delayMs = settings.baseDelayMs * 2 ** (this._retryAttempt - 1); + const delayMs = decision.delayMs; this._emit({ type: "auto_retry_start", diff --git a/packages/coding-agent/src/core/retry-backoff.ts b/packages/coding-agent/src/core/retry-backoff.ts new file mode 100644 index 0000000000..e21cc4c460 --- /dev/null +++ b/packages/coding-agent/src/core/retry-backoff.ts @@ -0,0 +1,93 @@ +/** + * Session-level retry policy. + * + * The session loop is the only layer that retries a failed turn (provider SDK + * retries are disabled while it is enabled, see `sdk.ts`), so this module owns + * the whole decision: how long to wait, when to honor an upstream `Retry-After`, + * and when to stop laddering entirely. + * + * Overload and rate-limit failures back off far harder than the generic ladder + * and spread concurrent agents apart with jitter, because N subagents sharing + * one upstream account otherwise retry in lockstep and re-create the burst that + * caused the failure. Kinds this module does not classify keep the historical + * fixed exponential ladder. + */ + +/** Base backoff for 429s. Rate-limit windows are long; a 2s ladder just re-burns quota. */ +const RATE_LIMIT_BASE_DELAY_MS = 15_000; +/** Base backoff for 529/overloaded. Shorter than a rate limit, still far above the generic ladder. */ +const OVERLOADED_BASE_DELAY_MS = 10_000; +/** Spread added on top of an honored `Retry-After` so concurrent agents do not resume together. */ +const RETRY_AFTER_JITTER_MS = 1_000; + +export interface RetryBackoffInput { + /** `provider_stream_failure` diagnostic kind, when the provider classified the failure. */ + kind?: string; + /** 1-based retry attempt number. */ + attempt: number; + /** Configured `retry.baseDelayMs`; also the floor for jittered waits. */ + baseDelayMs: number; + /** Configured `retry.provider.maxRetryDelayMs`; 0 disables the cap. */ + maxRetryDelayMs: number; + /** Upstream `Retry-After` in milliseconds, when the provider sent one. */ + retryAfterMs?: number; + /** Injectable for deterministic tests. */ + random?: () => number; +} + +export type RetryBackoffDecision = + | { type: "wait"; delayMs: number; honoredRetryAfter: boolean } + | { type: "abort"; reason: string }; + +function kindBaseDelayMs(kind: string | undefined): number | undefined { + if (kind === "rate_limit") return RATE_LIMIT_BASE_DELAY_MS; + if (kind === "overloaded") return OVERLOADED_BASE_DELAY_MS; + return undefined; +} + +function formatSeconds(ms: number): string { + return `${Math.round(ms / 1000)}s`; +} + +/** + * Auth failures never recover by waiting: the credential is expired, revoked, or + * wrong, and every extra attempt is another rejected upstream request. + */ +export function isNonRetryableFailureKind(kind: string | undefined): boolean { + return kind === "auth"; +} + +export function computeRetryBackoff(input: RetryBackoffInput): RetryBackoffDecision { + if (isNonRetryableFailureKind(input.kind)) { + return { type: "abort", reason: "Provider authentication failed; not retrying" }; + } + + const base = kindBaseDelayMs(input.kind); + if (base === undefined) { + return { + type: "wait", + delayMs: input.baseDelayMs * 2 ** (input.attempt - 1), + honoredRetryAfter: false, + }; + } + + const random = input.random ?? Math.random; + const cap = input.maxRetryDelayMs > 0 ? input.maxRetryDelayMs : Number.POSITIVE_INFINITY; + + if (input.retryAfterMs !== undefined) { + if (input.retryAfterMs > cap) { + return { + type: "abort", + reason: `Provider asked to wait ${formatSeconds(input.retryAfterMs)}, longer than the ${formatSeconds(cap)} retry delay cap (retry.provider.maxRetryDelayMs)`, + }; + } + const spread = Math.round(random() * RETRY_AFTER_JITTER_MS); + return { type: "wait", delayMs: Math.min(cap, input.retryAfterMs + spread), honoredRetryAfter: true }; + } + + // Full jitter over the exponential window, floored at the configured base + // delay so a lucky draw never turns a rate limit into an immediate retry. + const window = Math.min(cap, base * 2 ** (input.attempt - 1)); + const floor = Math.min(input.baseDelayMs, window); + return { type: "wait", delayMs: Math.max(floor, Math.round(random() * window)), honoredRetryAfter: false }; +} diff --git a/packages/coding-agent/src/core/sdk.ts b/packages/coding-agent/src/core/sdk.ts index acd3fff984..a688a05ab3 100644 --- a/packages/coding-agent/src/core/sdk.ts +++ b/packages/coding-agent/src/core/sdk.ts @@ -291,12 +291,19 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} throw new Error(auth.error); } const providerRetrySettings = settingsManager.getProviderRetrySettings(); + // Exactly one layer owns retry policy. While the session-level loop is + // enabled, provider SDK retries stay off so a failed turn cannot fan out + // into sessionRetries x sdkRetries upstream attempts; an explicit + // `retry.provider.maxRetries` still wins. + const providerMaxRetries = + options?.maxRetries ?? + providerRetrySettings.maxRetries ?? + (settingsManager.getRetryEnabled() ? 0 : undefined); return streamSimple(model, context, { ...options, apiKey: auth.apiKey, timeoutMs: options?.timeoutMs ?? providerRetrySettings.timeoutMs, - maxRetries: options?.maxRetries ?? providerRetrySettings.maxRetries, - maxRetryDelayMs: options?.maxRetryDelayMs ?? providerRetrySettings.maxRetryDelayMs, + maxRetries: providerMaxRetries, headers: auth.headers || options?.headers ? { ...auth.headers, ...options?.headers } : undefined, }); }, @@ -308,7 +315,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 0b8d309b0d..32de81135c 100644 --- a/packages/coding-agent/src/core/settings-manager.ts +++ b/packages/coding-agent/src/core/settings-manager.ts @@ -29,8 +29,8 @@ 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) + maxRetries?: number; // SDK/provider retry attempts; defaults to 0 while agent-level retry is enabled + maxRetryDelayMs?: number; // default: 60000 (max provider-requested delay the agent will wait before failing) } export interface RetrySettings { diff --git a/packages/coding-agent/src/core/side-question.ts b/packages/coding-agent/src/core/side-question.ts index 8f1c4cc812..59ce5fa821 100644 --- a/packages/coding-agent/src/core/side-question.ts +++ b/packages/coding-agent/src/core/side-question.ts @@ -112,7 +112,6 @@ export function startSideQuestion( sessionId: parent.sessionId === undefined ? undefined : `${parent.sessionId}/side:${id}`, thinkingBudgets: parent.thinkingBudgets, transport: "sse", - maxRetryDelayMs: parent.maxRetryDelayMs, toolExecution: parent.toolExecution, }); diff --git a/packages/coding-agent/test/retry-backoff.test.ts b/packages/coding-agent/test/retry-backoff.test.ts new file mode 100644 index 0000000000..82103b9a33 --- /dev/null +++ b/packages/coding-agent/test/retry-backoff.test.ts @@ -0,0 +1,99 @@ +import { describe, expect, it } from "vitest"; +import { computeRetryBackoff, isNonRetryableFailureKind } from "../src/core/retry-backoff.js"; + +const BASE = { baseDelayMs: 2000, maxRetryDelayMs: 60_000 }; + +describe("computeRetryBackoff", () => { + it("keeps the plain exponential ladder for kinds it does not classify", () => { + for (const kind of [undefined, "server_error", "unknown", "malformed_response"]) { + expect([1, 2, 3].map((attempt) => computeRetryBackoff({ ...BASE, kind, attempt }))).toEqual([ + { type: "wait", delayMs: 2000, honoredRetryAfter: false }, + { type: "wait", delayMs: 4000, honoredRetryAfter: false }, + { type: "wait", delayMs: 8000, honoredRetryAfter: false }, + ]); + } + }); + + it("aborts immediately on auth failures", () => { + expect(computeRetryBackoff({ ...BASE, kind: "auth", attempt: 1 })).toEqual({ + type: "abort", + reason: "Provider authentication failed; not retrying", + }); + expect(isNonRetryableFailureKind("auth")).toBe(true); + expect(isNonRetryableFailureKind("server_error")).toBe(false); + }); + + it("backs off rate limits from a longer base with full jitter", () => { + expect(computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, random: () => 1 })).toEqual({ + type: "wait", + delayMs: 15_000, + honoredRetryAfter: false, + }); + expect(computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, random: () => 0.5 })).toEqual({ + type: "wait", + delayMs: 7500, + honoredRetryAfter: false, + }); + }); + + it("floors a jittered wait at the configured base delay", () => { + expect(computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, random: () => 0 })).toEqual({ + type: "wait", + delayMs: 2000, + honoredRetryAfter: false, + }); + }); + + it("backs off overloads from their own base and grows exponentially", () => { + expect(computeRetryBackoff({ ...BASE, kind: "overloaded", attempt: 1, random: () => 1 }).type).toBe("wait"); + expect(computeRetryBackoff({ ...BASE, kind: "overloaded", attempt: 1, random: () => 1 })).toMatchObject({ + delayMs: 10_000, + }); + expect(computeRetryBackoff({ ...BASE, kind: "overloaded", attempt: 2, random: () => 1 })).toMatchObject({ + delayMs: 20_000, + }); + }); + + it("caps the jittered window at maxRetryDelayMs", () => { + expect( + computeRetryBackoff({ ...BASE, kind: "overloaded", attempt: 4, maxRetryDelayMs: 30_000, random: () => 1 }), + ).toMatchObject({ delayMs: 30_000 }); + }); + + it("honors Retry-After and spreads concurrent agents apart", () => { + expect( + computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, retryAfterMs: 30_000, random: () => 0 }), + ).toEqual({ type: "wait", delayMs: 30_000, honoredRetryAfter: true }); + expect( + computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, retryAfterMs: 30_000, random: () => 1 }), + ).toEqual({ type: "wait", delayMs: 31_000, honoredRetryAfter: true }); + }); + + it("fails fast when Retry-After exceeds the cap instead of waiting silently", () => { + expect(computeRetryBackoff({ ...BASE, kind: "rate_limit", attempt: 1, retryAfterMs: 3_600_000 })).toEqual({ + type: "abort", + reason: "Provider asked to wait 3600s, longer than the 60s retry delay cap (retry.provider.maxRetryDelayMs)", + }); + }); + + it("treats maxRetryDelayMs 0 as no cap", () => { + expect( + computeRetryBackoff({ + ...BASE, + kind: "rate_limit", + attempt: 1, + maxRetryDelayMs: 0, + retryAfterMs: 3_600_000, + random: () => 0, + }), + ).toEqual({ type: "wait", delayMs: 3_600_000, honoredRetryAfter: true }); + }); + + it("ignores Retry-After for kinds that keep the plain ladder", () => { + expect(computeRetryBackoff({ ...BASE, kind: "server_error", attempt: 1, retryAfterMs: 3_600_000 })).toEqual({ + type: "wait", + delayMs: 2000, + honoredRetryAfter: false, + }); + }); +}); 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..2047f81022 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 @@ -264,7 +264,20 @@ 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("does not retry structured provider auth failures at all", async () => { + const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } } }); + harnesses.push(harness); + harness.setResponses([structuredProviderFailure("auth"), fauxAssistantMessage("unused")]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); + 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(`retries structured permanent provider ${kind} failures once`, 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/24-kind-aware-retry.test.ts b/packages/coding-agent/test/suite/regressions/24-kind-aware-retry.test.ts new file mode 100644 index 0000000000..ca1d7b88d4 --- /dev/null +++ b/packages/coding-agent/test/suite/regressions/24-kind-aware-retry.test.ts @@ -0,0 +1,221 @@ +import { existsSync, mkdirSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { type AssistantMessage, fauxAssistantMessage, registerFauxProvider } from "@earendil-works/pi-ai"; +import { afterEach, describe, expect, it } from "vitest"; +import { AuthStorage } from "../../../src/core/auth-storage.js"; +import { ModelRegistry } from "../../../src/core/model-registry.js"; +import { createAgentSession } from "../../../src/core/sdk.js"; +import { SessionManager } from "../../../src/core/session-manager.js"; +import { type Settings, SettingsManager } from "../../../src/core/settings-manager.js"; +import { createTestResourceLoader } from "../../utilities.js"; +import { createHarness, type Harness } from "../harness.js"; + +function providerFailure(details: Record, errorMessage: string): AssistantMessage { + return { + ...fauxAssistantMessage("", { stopReason: "error", errorMessage }), + diagnostics: [{ type: "provider_stream_failure", timestamp: Date.now(), details }], + }; +} + +function rateLimited(retryAfterMs?: number): AssistantMessage { + return providerFailure( + { kind: "rate_limit", status: 429, retryAfterMs }, + "Provider rate limit exceeded (rate_limit_error, 429)", + ); +} + +/** Capture the scheduled backoff without serving it: the sleep is abortable. */ +async function captureRetryDelayMs(harness: Harness): Promise { + let delayMs: number | undefined; + const sawRetryStart = new Promise((resolve) => { + const unsubscribe = harness.session.subscribe((event) => { + if (event.type === "auto_retry_start") { + delayMs = event.delayMs; + unsubscribe(); + resolve(); + } + }); + }); + + const promptPromise = harness.session.prompt("test"); + await sawRetryStart; + harness.session.abortRetry(); + await promptPromise; + return delayMs; +} + +describe("issue #24 kind-aware retry backoff", () => { + const harnesses: Harness[] = []; + const cleanups: Array<() => void> = []; + + afterEach(() => { + while (harnesses.length > 0) { + harnesses.pop()?.cleanup(); + } + while (cleanups.length > 0) { + cleanups.pop()?.(); + } + }); + + it("honors an upstream Retry-After instead of the 2s ladder", async () => { + const harness = await createHarness({ + settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 2000 } }, + }); + harnesses.push(harness); + harness.setResponses([rateLimited(30_000), rateLimited(30_000), rateLimited(30_000), rateLimited(30_000)]); + + const delayMs = await captureRetryDelayMs(harness); + + expect(harness.faux.state.callCount).toBe(1); + expect(delayMs).toBeGreaterThanOrEqual(30_000); + expect(delayMs).toBeLessThanOrEqual(31_000); + }); + + it("fails fast when Retry-After exceeds retry.provider.maxRetryDelayMs", async () => { + const harness = await createHarness({ + settings: { + retry: { enabled: true, maxRetries: 3, baseDelayMs: 1, provider: { maxRetryDelayMs: 60_000 } }, + }, + }); + harnesses.push(harness); + harness.setResponses([rateLimited(3_600_000), rateLimited(3_600_000)]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); + expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([false]); + expect(harness.eventsOfType("auto_retry_end")[0]?.finalError).toContain( + "Provider asked to wait 3600s, longer than the 60s retry delay cap", + ); + }); + + it("aborts the ladder on auth failures instead of burning every attempt", async () => { + const harness = await createHarness({ + settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } }, + }); + harnesses.push(harness); + harness.setResponses([ + providerFailure({ kind: "auth", status: 401 }, "401 Unauthorized: token expired"), + fauxAssistantMessage("must not be reached"), + ]); + + await harness.session.prompt("test"); + + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); + expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([false]); + expect(harness.eventsOfType("auto_retry_end")[0]?.finalError).toContain( + "Provider authentication failed; not retrying", + ); + expect(harness.session.isRetrying).toBe(false); + }); + + it("jitters rate-limit backoff well above the plain ladder", async () => { + const controlHarness = await createHarness({ + settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 100 } }, + }); + harnesses.push(controlHarness); + controlHarness.setResponses([ + providerFailure({ kind: "server_error", status: 500 }, "500 Internal Server Error"), + fauxAssistantMessage("unused"), + ]); + expect(await captureRetryDelayMs(controlHarness)).toBe(100); + + const delays: number[] = []; + for (let run = 0; run < 4; run++) { + const harness = await createHarness({ + settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 100 } }, + }); + harnesses.push(harness); + harness.setResponses([rateLimited(), rateLimited()]); + const delayMs = await captureRetryDelayMs(harness); + expect(delayMs).toBeDefined(); + delays.push(delayMs as number); + } + + for (const delayMs of delays) { + expect(delayMs).toBeGreaterThanOrEqual(100); + expect(delayMs).toBeLessThanOrEqual(15_000); + } + expect(new Set(delays).size).toBeGreaterThan(1); + }); +}); + +describe("issue #24 single retry layer", () => { + const cleanups: Array<() => void> = []; + + afterEach(() => { + while (cleanups.length > 0) { + cleanups.pop()?.(); + } + }); + + async function captureProviderMaxRetries(settings: Partial): Promise { + const tempDir = join(tmpdir(), `pi-issue24-${Date.now()}-${Math.random().toString(36).slice(2)}`); + const agentDir = join(tempDir, "agent"); + mkdirSync(agentDir, { recursive: true }); + const faux = registerFauxProvider(); + const model = faux.getModel(); + const authStorage = AuthStorage.inMemory(); + authStorage.setRuntimeApiKey(model.provider, "faux-key"); + const modelRegistry = ModelRegistry.inMemory(authStorage); + modelRegistry.registerProvider(model.provider, { + baseUrl: model.baseUrl, + apiKey: "faux-key", + api: faux.api, + models: faux.models.map((registered) => ({ + id: registered.id, + name: registered.name, + api: registered.api, + reasoning: registered.reasoning, + input: registered.input, + cost: registered.cost, + contextWindow: registered.contextWindow, + maxTokens: registered.maxTokens, + baseUrl: registered.baseUrl, + })), + }); + + const { session } = await createAgentSession({ + cwd: tempDir, + agentDir, + model, + authStorage, + modelRegistry, + settingsManager: SettingsManager.inMemory(settings), + sessionManager: SessionManager.inMemory(), + resourceLoader: createTestResourceLoader(), + }); + cleanups.push(() => { + session.dispose(); + faux.unregister(); + if (existsSync(tempDir)) rmSync(tempDir, { recursive: true, force: true }); + }); + + let observed: number | undefined; + faux.setResponses([ + (_context, streamOptions) => { + observed = streamOptions?.maxRetries; + return fauxAssistantMessage("ok"); + }, + ]); + await session.prompt("hello"); + return observed; + } + + it("disables provider SDK retries while the session retry loop owns policy", async () => { + await expect(captureProviderMaxRetries({ retry: { enabled: true } })).resolves.toBe(0); + }); + + it("keeps an explicit provider retry override", async () => { + await expect(captureProviderMaxRetries({ retry: { enabled: true, provider: { maxRetries: 2 } } })).resolves.toBe( + 2, + ); + }); + + it("leaves provider retries at the SDK default when session retry is disabled", async () => { + await expect(captureProviderMaxRetries({ retry: { enabled: false } })).resolves.toBeUndefined(); + }); +}); diff --git a/packages/coding-agent/test/suite/regressions/4491-provider-stale-after-401.test.ts b/packages/coding-agent/test/suite/regressions/4491-provider-stale-after-401.test.ts index 48fc2a7ace..7db70a43fd 100644 --- a/packages/coding-agent/test/suite/regressions/4491-provider-stale-after-401.test.ts +++ b/packages/coding-agent/test/suite/regressions/4491-provider-stale-after-401.test.ts @@ -51,7 +51,7 @@ describe("issue #4491 provider stale after repeated 401", () => { } }); - it("retries structured provider auth failures once, then marks current auth stale", async () => { + it("stops on the first structured provider auth failure and marks current auth stale", async () => { const harness = await createHarness({ settings: { retry: { enabled: true, maxRetries: 2, baseDelayMs: 1 } }, }); @@ -60,8 +60,8 @@ describe("issue #4491 provider stale after repeated 401", () => { await harness.session.prompt("hello"); - expect(harness.faux.state.callCount).toBe(2); - expect(harness.eventsOfType("auto_retry_start").map((event) => event.attempt)).toEqual([1]); + expect(harness.faux.state.callCount).toBe(1); + expect(harness.eventsOfType("auto_retry_start")).toEqual([]); expect(harness.eventsOfType("auto_retry_end").map((event) => event.success)).toEqual([false]); expect(harness.eventsOfType("auth_stale")).toHaveLength(1); @@ -148,7 +148,7 @@ describe("issue #4491 provider stale after repeated 401", () => { settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 100 } }, }); harnesses.push(harness); - harness.setResponses([provider401Message(), provider401Message()]); + harness.setResponses([bareProvider401Message(), bareProvider401Message()]); const sawRetryStart = new Promise((resolve) => { const unsubscribe = harness.session.subscribe((event) => { @@ -176,7 +176,7 @@ describe("issue #4491 provider stale after repeated 401", () => { settings: { retry: { enabled: true, maxRetries: 1, baseDelayMs: 5 } }, }); harnesses.push(harness); - harness.setResponses([provider401Message(), provider401Message()]); + harness.setResponses([bareProvider401Message(), bareProvider401Message()]); let changedCredentials = false; harness.session.subscribe((event) => { @@ -204,7 +204,7 @@ describe("issue #4491 provider stale after repeated 401", () => { settings: { retry: { enabled: true, maxRetries: 2, baseDelayMs: 1 } }, }); harnesses.push(harness); - harness.setResponses([provider401Message(), provider500Message(), provider500Message()]); + harness.setResponses([bareProvider401Message(), provider500Message(), provider500Message()]); await harness.session.prompt("hello"); From c01f6b9f6a5d4c38b8a7f9bea3f738d258c32cff Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sun, 30 Aug 2026 22:45:45 -0600 Subject: [PATCH 2/2] fix: prefer in-frame retry_after for mid-stream SSE errors A stream's headers are sent before the failure is known, so proxies that compute the wait mid-stream (Meridian, rynfar/meridian#907) carry it in the error frame as seconds. Refs #24. --- packages/ai/src/providers/anthropic.ts | 12 ++++++++- .../ai/test/anthropic-sse-parsing.test.ts | 25 +++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index 50fac7f728..fd01db03ed 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -389,11 +389,21 @@ function anthropicSseError(data: string, requestId?: string, retryAfterMs?: numb let errorType: string | undefined; let detail: string | undefined; try { - const parsed = parseJsonWithRepair<{ error?: { type?: string; message?: string }; request_id?: string }>(data); + const parsed = parseJsonWithRepair<{ + error?: { type?: string; message?: string; retry_after?: number }; + request_id?: string; + }>(data); errorType = parsed.error?.type; detail = parsed.error?.message; // Proxies may strip the request-id header; the error body carries it too. requestId ??= typeof parsed.request_id === "string" ? parsed.request_id : undefined; + // A stream's headers are sent before the error is known, so proxies that + // compute a wait mid-stream (e.g. Meridian) carry it in the frame as + // seconds. The in-frame value describes this exact failure; prefer it. + const frameRetryAfter = parsed.error?.retry_after; + if (typeof frameRetryAfter === "number" && Number.isFinite(frameRetryAfter) && frameRetryAfter >= 0) { + retryAfterMs = Math.round(frameRetryAfter * 1000); + } } catch { detail = data; } diff --git a/packages/ai/test/anthropic-sse-parsing.test.ts b/packages/ai/test/anthropic-sse-parsing.test.ts index 08f5f93c74..4985fcd94d 100644 --- a/packages/ai/test/anthropic-sse-parsing.test.ts +++ b/packages/ai/test/anthropic-sse-parsing.test.ts @@ -173,6 +173,31 @@ describe("Anthropic raw SSE parsing", () => { }); }); + it("prefers an in-frame retry_after over the response header for a mid-stream SSE error", async () => { + const model = getModel("anthropic", "claude-haiku-4-5"); + // A stream's headers go out before the failure is known, so proxies + // (e.g. Meridian) put the wait in the error frame itself, in seconds. + const response = new Response( + `event: error\ndata: ${JSON.stringify({ + type: "error", + error: { type: "rate_limit_error", message: "Rate limited", retry_after: 45 }, + })}\n`, + { status: 200, headers: { "content-type": "text/event-stream", "retry-after": "30" } }, + ); + + const result = await streamAnthropic( + model, + { messages: [{ role: "user", content: "Say hello.", timestamp: Date.now() }] }, + { client: createFakeAnthropicClient(response) }, + ).result(); + + expect(result.stopReason).toBe("error"); + expect(result.diagnostics?.[0]).toMatchObject({ + type: "provider_stream_failure", + details: { kind: "rate_limit", retryAfterMs: 45_000 }, + }); + }); + it("preserves configured cache write pricing for non-Anthropic models", async () => { const model = getModel("minimax", "MiniMax-M2.7-highspeed"); const response = createSseResponse(