diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 326853dd8..bc89d2e2c 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -74,6 +74,9 @@ - Fixed OpenAI Codex usage telemetry blocking explicitly allowed ChatGPT Team credentials when a weekly `used_percent` rounded to 100, which could route multi-account sessions to an actually exhausted sibling instead ([#7617](https://github.com/can1357/oh-my-pi/issues/7617)). - Fixed OpenAI Codex GPT-5.x requests sending optional `reasoning.summary`, `reasoning.context`, and `text.verbosity` controls by default, reducing Codex `server_error` disconnects from unsupported request shapes. ([#4949](https://github.com/can1357/oh-my-pi/issues/4949)) - Classified concurrent-request caps separately from quota exhaustion so they use a short retry backoff without burning a credential, and rotate credentials for account-scoped 403 caps such as Devin's overall message limit. +### Added + +- Added `retryTransientCompletion` for oneshot (non-agent-loop) LLM calls. `completeSimple` reports a transient provider failure — Anthropic `overloaded_error` / `rate_limit_error`, HTTP 429/500/502/503/529 — by **resolving** with `stopReason: "error"` rather than throwing, so callers that wrapped it in a try/catch-based retry never actually retried. The helper handles both shapes (resolved error-stop and thrown error), classifies with the existing `AIError` predicates so the retryable set stays defined in one place, honors `retry-after` / `x-ratelimit-reset*` from response headers (supplied via `getResponseHeaders`, since an `AssistantMessage` carries none) or from the error text, and returns/rethrows the failure unchanged once attempts are exhausted so existing caller fallbacks keep working. Aborts surface immediately. ## [17.2.7] - 2026-08-03 diff --git a/packages/ai/src/index.ts b/packages/ai/src/index.ts index ec472a32c..d2f9a9a00 100644 --- a/packages/ai/src/index.ts +++ b/packages/ai/src/index.ts @@ -6,6 +6,7 @@ export * from "./auth-gateway/types"; export * from "./auth-retry"; export * from "./auth-storage"; export * from "./error/rate-limit"; +export * from "./oneshot-retry"; export * from "./provider-details"; export * from "./providers/anthropic"; export * from "./providers/anthropic-client"; diff --git a/packages/ai/src/oneshot-retry.ts b/packages/ai/src/oneshot-retry.ts new file mode 100644 index 000000000..8a576af68 --- /dev/null +++ b/packages/ai/src/oneshot-retry.ts @@ -0,0 +1,213 @@ +import { extractRetryHint } from "@oh-my-pi/pi-utils"; +import * as AIError from "./error"; +import type { AssistantMessage } from "./types"; +import { getHeadersFromError, getRetryAfterMsFromHeaders, type HeadersLike } from "./utils/retry-after"; + +/** + * Transient-failure retry for **oneshot** (non-agent-loop) completions. + * + * Why this exists: `streamSimple`/`completeSimple` retry *auth* failures + * (credential rotation) but deliberately surface *transient* provider failures + * — Anthropic `overloaded_error`, `rate_limit_error`, HTTP 429/500/502/503/529 + * — as a **resolved** `AssistantMessage` with `stopReason: "error"`. For the + * main agent turn that is correct: `TurnRecovery` owns recovery there, and it + * must refuse to replay once tool calls or visible text have streamed. + * + * Oneshots have no such hazard. A summary, title, handoff, or image + * description produces no side effects, so re-issuing the whole request is + * safe and is almost always what the caller wants. Before this helper every + * oneshot call site had to re-implement that decision, and most did not — + * failing on the first blip, or swallowing it into `null` so a transient + * overload was indistinguishable from a legitimate empty result. + * + * Classification reuses the existing provider predicates (`AIError`), so the + * set of retryable Anthropic failures stays defined in exactly one place. + * Usage limits are included: unlike the provider loop — which excludes them so + * credential rotation can own them — a oneshot has no rotation layer above it, + * and the retry hint the provider supplies (`retry-after`, "try again in ~5m") + * is honored, so waiting is the correct response. + */ +export interface OneshotRetryOptions { + /** Total attempts, including the first. Default 3. Values < 1 are treated as 1. */ + maxAttempts?: number; + /** First backoff step in ms; doubles per attempt. Default 500. */ + baseDelayMs?: number; + /** + * Upper bound for a single wait. Default 30_000. A provider retry hint + * longer than this aborts the retry instead of parking the caller — the + * error surfaces so higher-level recovery (or the user) can decide. + */ + maxDelayMs?: number; + /** + * Stops further attempts. Two distinct paths, both preserving the caller's + * intent: an abort already visible when an attempt settles surfaces that + * attempt's own result (`completeSimple` reports `stopReason: "aborted"`), + * while an abort that lands during the backoff wait rejects with the abort + * reason — a user cancel stays a cancel and is never relabelled as the + * provider failure we happened to be waiting on. + * + * This helper does NOT pass the signal into `run` — cancelling the in-flight + * request is the closure's job, because a per-attempt deadline must be + * rebuilt on every attempt. Construct it inside `run` + * (`signal: AbortSignal.timeout(MS)`, or `AbortSignal.any([outer, perAttempt])`); + * a deadline captured outside would fire once and then abort every retry, + * silently turning this helper into a single attempt. + */ + signal?: AbortSignal; + /** + * Headers of the attempt that just failed, used to honor `retry-after`. + * + * Load-bearing: a transient Anthropic failure arrives as a **resolved** + * `AssistantMessage`, and `AssistantMessage` carries no headers — so without + * this the real `retry-after` / `x-ratelimit-reset` values on a 429/529 are + * invisible and only the (usually hint-free) error text is available. + * Callers that already capture headers via `SimpleStreamOptions.onResponse` + * should return the latest capture here; it is read once per failed attempt. + * Thrown errors need no wiring — headers are recovered from the error itself. + */ + getResponseHeaders?: () => HeadersLike; + /** Observability hook. Fires immediately before sleeping. */ + onRetry?: (info: OneshotRetryInfo) => void; +} + +export interface OneshotRetryInfo { + /** 1-based index of the attempt that just failed. */ + attempt: number; + maxAttempts: number; + delayMs: number; + /** True when `delayMs` came from a provider retry hint rather than backoff. */ + fromRetryHint: boolean; + errorMessage: string; + /** `AIError` classification bits of the failure. */ + errorId: number; +} + +const DEFAULT_MAX_ATTEMPTS = 3; +const DEFAULT_BASE_DELAY_MS = 500; +const DEFAULT_MAX_DELAY_MS = 30_000; +/** Cap on pure backoff growth. A provider hint may still exceed this, up to `maxDelayMs`. */ +const BACKOFF_CEILING_MS = 8_000; + +function backoffDelayMs(attempt: number, baseDelayMs: number): number { + const growth = Math.min(baseDelayMs * 2 ** (attempt - 1), BACKOFF_CEILING_MS); + // 75-100% jitter, matching the provider loop and TurnRecovery, so a fleet of + // concurrent oneshots does not re-converge on the same instant. + return Math.round(growth * (0.75 + Math.random() * 0.25)); +} + +/** Retryable when the provider says transient, or when it says "wait, then retry". */ +function isRetryableOneshotFailure(errorId: number): boolean { + return ( + AIError.is(errorId, AIError.Flag.Transient) || + AIError.is(errorId, AIError.Flag.UsageLimit) || + AIError.retriable(errorId) + ); +} + +function sleep(delayMs: number, signal?: AbortSignal): Promise { + if (delayMs <= 0) return Promise.resolve(); + const { promise, resolve, reject } = Promise.withResolvers(); + const timer = setTimeout(() => { + signal?.removeEventListener("abort", onAbort); + resolve(); + }, delayMs); + const onAbort = () => { + clearTimeout(timer); + reject(signal?.reason ?? new AIError.AbortError("oneshot retry aborted")); + }; + if (signal) { + if (signal.aborted) { + clearTimeout(timer); + return Promise.reject(signal.reason ?? new AIError.AbortError("oneshot retry aborted")); + } + signal.addEventListener("abort", onAbort, { once: true }); + } + return promise; +} + +/** + * Run a oneshot completion, retrying transient provider failures. + * + * Handles both failure shapes: a resolved `AssistantMessage` carrying + * `stopReason: "error"` (what `completeSimple` produces) and a thrown error + * (what the raw HTTP helpers produce). A non-retryable failure is returned or + * rethrown unchanged, so existing caller error handling keeps working — this + * only removes the *first-blip* failure mode. + */ +export async function retryTransientCompletion( + run: (attempt: number) => Promise, + options?: OneshotRetryOptions, +): Promise { + const maxAttempts = Math.max(1, options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS); + const baseDelayMs = options?.baseDelayMs ?? DEFAULT_BASE_DELAY_MS; + const maxDelayMs = options?.maxDelayMs ?? DEFAULT_MAX_DELAY_MS; + const signal = options?.signal; + + for (let attempt = 1; ; attempt++) { + let message: AssistantMessage | undefined; + let thrown: unknown; + try { + message = await run(attempt); + if (message.stopReason !== "error") return message; + } catch (error) { + thrown = error; + } + // A caller abort is never a transient failure — surface it immediately so + // cancellation stays responsive. + if (signal?.aborted) { + if (thrown !== undefined) throw thrown; + return message as AssistantMessage; + } + + const errorId = + thrown !== undefined ? AIError.classify(thrown) : AIError.classifyMessage(message as AssistantMessage); + if (AIError.is(errorId, AIError.Flag.Abort) || AIError.is(errorId, AIError.Flag.UserInterrupt)) { + if (thrown !== undefined) throw thrown; + return message as AssistantMessage; + } + const errorMessage = + thrown !== undefined + ? thrown instanceof Error + ? thrown.message + : String(thrown) + : ((message as AssistantMessage).errorMessage ?? "unknown error"); + const lastAttempt = attempt >= maxAttempts; + if (lastAttempt || !isRetryableOneshotFailure(errorId)) { + if (thrown !== undefined) throw thrown; + return message as AssistantMessage; + } + + // Headers first: a real Anthropic 429/529 carries `retry-after` / + // `x-ratelimit-reset*` in the response, and the resolved AssistantMessage + // has none — the caller supplies them via `getResponseHeaders`. Thrown + // errors (e.g. AnthropicApiError) carry their own headers. + const headers: HeadersLike = thrown !== undefined ? getHeadersFromError(thrown) : options?.getResponseHeaders?.(); + const headerHintMs = getRetryAfterMsFromHeaders(headers); + const textHintMs = extractRetryHint(undefined, errorMessage); + const hintMs = + headerHintMs === undefined && textHintMs === undefined + ? undefined + : Math.max(headerHintMs ?? 0, textHintMs ?? 0); + // An over-cap hint means "come back much later"; parking a oneshot that + // long is worse than surfacing the failure to the caller. + if (hintMs !== undefined && hintMs > maxDelayMs) { + if (thrown !== undefined) throw thrown; + return message as AssistantMessage; + } + const backoff = backoffDelayMs(attempt, baseDelayMs); + const delayMs = Math.min(Math.max(hintMs ?? 0, backoff), maxDelayMs); + + options?.onRetry?.({ + attempt, + maxAttempts, + delayMs, + fromRetryHint: hintMs !== undefined && hintMs >= backoff, + errorMessage, + errorId, + }); + // Aborting mid-backoff rejects with the caller's abort reason: a user + // cancel must stay a cancel, not get relabelled as the provider failure we + // happened to be waiting on. + await sleep(delayMs, signal); + } +} diff --git a/packages/ai/test/oneshot-retry.test.ts b/packages/ai/test/oneshot-retry.test.ts new file mode 100644 index 000000000..32b8880c1 --- /dev/null +++ b/packages/ai/test/oneshot-retry.test.ts @@ -0,0 +1,277 @@ +import { describe, expect, it } from "bun:test"; +import { retryTransientCompletion } from "@oh-my-pi/pi-ai/oneshot-retry"; +import type { AssistantMessage, Usage } from "@oh-my-pi/pi-ai/types"; + +/** + * Defends the contract every oneshot LLM call site now depends on: + * `completeSimple` reports a transient provider failure by RESOLVING with + * `stopReason: "error"` rather than throwing, so a retry layer that only + * catches exceptions silently never fires. Before this helper, an Anthropic + * `overloaded_error` / 429 / 529 on a summary, title, handoff or image + * description failed on the first blip — or was swallowed into `null`, making a + * transient overload indistinguishable from a legitimate empty result. + * + * These tests pin the four properties the call sites rely on: transient + * error-stops are re-issued, non-transient ones are not, the final failure is + * handed back unchanged (so existing `null`/throw fallbacks still work), and a + * caller abort wins immediately. + */ + +const emptyUsage = (): Usage => + ({ + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }) as unknown as Usage; + +function message(overrides: Partial = {}): AssistantMessage { + return { + role: "assistant", + content: [], + api: "anthropic-messages", + provider: "anthropic", + model: "claude-sonnet-4-6", + usage: emptyUsage(), + stopReason: "stop", + timestamp: 0, + ...overrides, + } as AssistantMessage; +} + +const overloaded = (): AssistantMessage => + message({ + stopReason: "error", + errorStatus: 529, + errorMessage: "Anthropic stream error (overloaded_error): Overloaded", + }); + +const rateLimited = (): AssistantMessage => + message({ stopReason: "error", errorStatus: 429, errorMessage: "rate_limit_error: too many requests" }); + +// Keep the suite fast: the helper's real backoff floor is 500ms. +const fast = { baseDelayMs: 1, maxAttempts: 3 } as const; + +describe("retryTransientCompletion", () => { + it("re-issues an Anthropic 529 error-stop and returns the eventual success", async () => { + const results = [overloaded(), overloaded(), message({ stopReason: "stop" })]; + let calls = 0; + const final = await retryTransientCompletion(() => { + calls += 1; + return Promise.resolve(results.shift()!); + }, fast); + + expect(calls).toBe(3); + expect(final.stopReason).toBe("stop"); + }); + + it("re-issues a 429 rate-limit error-stop", async () => { + let calls = 0; + const final = await retryTransientCompletion(() => { + calls += 1; + return Promise.resolve(calls === 1 ? rateLimited() : message()); + }, fast); + + expect(calls).toBe(2); + expect(final.stopReason).toBe("stop"); + }); + + it("returns the failing message unchanged once attempts are exhausted, so caller fallbacks still apply", async () => { + let calls = 0; + const final = await retryTransientCompletion(() => { + calls += 1; + return Promise.resolve(overloaded()); + }, fast); + + expect(calls).toBe(3); + expect(final.stopReason).toBe("error"); + expect(final.errorMessage).toContain("overloaded_error"); + }); + + it("does not retry a non-transient provider error", async () => { + let calls = 0; + const final = await retryTransientCompletion( + () => { + calls += 1; + return Promise.resolve( + message({ + stopReason: "error", + errorStatus: 400, + errorMessage: "invalid_request_error: messages: at least one message is required", + }), + ); + }, + { ...fast, maxAttempts: 5 }, + ); + + expect(calls).toBe(1); + expect(final.stopReason).toBe("error"); + }); + + it("retries a thrown transient error and rethrows the last one when exhausted", async () => { + let calls = 0; + const attempt = retryTransientCompletion(() => { + calls += 1; + const error = new Error("529 overloaded_error: Overloaded") as Error & { status?: number }; + error.status = 529; + throw error; + }, fast); + + await expect(attempt).rejects.toThrow(/overloaded_error/); + expect(calls).toBe(3); + }); + + it("does not retry a thrown non-transient error", async () => { + let calls = 0; + const attempt = retryTransientCompletion(() => { + calls += 1; + throw new Error("invalid_request_error: bad tool schema"); + }, fast); + + await expect(attempt).rejects.toThrow(/invalid_request_error/); + expect(calls).toBe(1); + }); + + it("stops immediately when the caller aborts", async () => { + const controller = new AbortController(); + let calls = 0; + const final = await retryTransientCompletion( + () => { + calls += 1; + controller.abort(); + return Promise.resolve(overloaded()); + }, + { ...fast, signal: controller.signal }, + ); + + expect(calls).toBe(1); + expect(final.stopReason).toBe("error"); + }); + + it("rejects with the abort reason when the caller cancels during backoff", async () => { + // The cancel lands while we are waiting, not while an attempt is in flight: + // it must stay a cancellation rather than being reported as the provider + // failure we happened to be sleeping on. + const controller = new AbortController(); + const reason = new Error("user pressed escape"); + let calls = 0; + const attempt = retryTransientCompletion( + () => { + calls += 1; + setTimeout(() => controller.abort(reason), 5); + return Promise.resolve(overloaded()); + }, + { maxAttempts: 3, baseDelayMs: 200, signal: controller.signal }, + ); + + await expect(attempt).rejects.toThrow("user pressed escape"); + expect(calls).toBe(1); + }); + + it("reports each retry through onRetry so callers can log the wait", async () => { + const seen: number[] = []; + let calls = 0; + await retryTransientCompletion( + () => { + calls += 1; + return Promise.resolve(calls === 1 ? overloaded() : message()); + }, + { ...fast, onRetry: info => seen.push(info.attempt) }, + ); + + expect(seen).toEqual([1]); + }); + + it("surfaces the failure instead of parking when the provider asks for longer than maxDelayMs", async () => { + let calls = 0; + const final = await retryTransientCompletion( + () => { + calls += 1; + return Promise.resolve( + message({ + stopReason: "error", + errorStatus: 429, + errorMessage: "rate_limit_error: please retry in 600s", + }), + ); + }, + { ...fast, maxDelayMs: 1_000 }, + ); + + expect(calls).toBe(1); + expect(final.stopReason).toBe("error"); + }); + + it("honors a retry-after-ms response header over the backoff floor", async () => { + // The header is the only place a real Anthropic 429 carries its wait: the + // resolved AssistantMessage has no headers, so a helper that reads only the + // error text would silently fall back to plain backoff. + let calls = 0; + let observedDelay = -1; + await retryTransientCompletion( + () => { + calls += 1; + return Promise.resolve(calls === 1 ? rateLimited() : message()); + }, + { + maxAttempts: 2, + baseDelayMs: 1, + getResponseHeaders: () => ({ "retry-after-ms": "120" }), + onRetry: info => { + observedDelay = info.delayMs; + }, + }, + ); + + expect(calls).toBe(2); + expect(observedDelay).toBe(120); + }); + + it("surfaces the failure when a retry-after header exceeds maxDelayMs", async () => { + let calls = 0; + const final = await retryTransientCompletion( + () => { + calls += 1; + return Promise.resolve(rateLimited()); + }, + { + maxAttempts: 3, + baseDelayMs: 1, + maxDelayMs: 1_000, + getResponseHeaders: () => ({ "retry-after": "300" }), + }, + ); + + expect(calls).toBe(1); + expect(final.stopReason).toBe("error"); + }); + + it("recovers retry-after from a thrown provider error's own headers", async () => { + let calls = 0; + let observedDelay = -1; + const attempt = retryTransientCompletion( + () => { + calls += 1; + const error = new Error("529 overloaded_error: Overloaded") as Error & { + status?: number; + headers?: Record; + }; + error.status = 529; + error.headers = { "retry-after-ms": "90" }; + throw error; + }, + { + maxAttempts: 2, + baseDelayMs: 1, + onRetry: info => { + observedDelay = info.delayMs; + }, + }, + ); + + await expect(attempt).rejects.toThrow(/overloaded_error/); + expect(calls).toBe(2); + expect(observedDelay).toBe(90); + }); +});