diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 6f9ef27f5..ef8be8105 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -1,12 +1,15 @@ # Changelog ## [Unreleased] + ### Added - Added `CheckCredentialsOptions.completionProbe` (and `completionTimeoutMs`) so `AuthStorage.checkCredentials` can additionally exercise each credential against the provider's chat-completion endpoint after refresh-on-expiry. Result lands on `CredentialHealthResult.completion` ({ok, reason?, modelId?, latencyMs?}) without disturbing the usage `ok` field. Public types: `CompletionProbe`, `CompletionProbeInput`, `CompletionProbeCredential`, `CredentialCompletionResult`. The probe is invoked even when no `UsageProvider` is registered for the row, and is skipped when OAuth refresh fails (the stale bytes would only mask the upstream failure). ### Changed +- Changed auth-gateway credential resolution to use per-conversation `promptCacheKey`/`sessionId` when calling `AuthStorage.getApiKey`, so repeated turns can keep the same credential until it becomes unavailable +- Changed auth-gateway and pi-native request handling to align `sessionId` with prompt/context identity before credential lookup - Changed Anthropic prompt preparation to downscale image blocks over 2000px when a request includes 20+ images, reducing oversized payloads automatically - Changed OpenAI chat request parsing to accept `name` on `tool` messages and fall back to the matching assistant `tool_calls` name, so parsed tool results now carry a proper tool name when the wire omits it - Changed `checkCredentials` to skip running `completionProbe` when OAuth refresh fails, so stale bearer tokens are never probed and the refresh failure remains the returned `reason` @@ -15,6 +18,9 @@ ### Fixed +- Fixed auth-gateway to classify usage-limit messages such as `usage_limit_reached`, `resource_exhausted`, and Codex-style `Try again in ~X min` text as 429 `rate_limit_error` responses +- Fixed auth-gateway usage-limit handling to honor parsed retry hints and switch to a sibling credential via `markUsageLimitReached` instead of invalidating the rate-limited credential +- Fixed `streamSimple` to retry on usage-limit errors (including message-only error events) before any content is emitted, so `onAuthError` can rotate credentials automatically - Fixed auth-gateway error classification to extract embedded status codes and use word-boundary matching, so `GenerateContentRequest` and similar messages are no longer misreported as rate-limit errors - Fixed `checkCredentials` to handle `completionProbe` exceptions by recording the failure in `CredentialHealthResult.completion.reason` while still returning the usage probe result diff --git a/packages/ai/src/auth-gateway/server.ts b/packages/ai/src/auth-gateway/server.ts index 2c1859ef1..db2ed193d 100644 --- a/packages/ai/src/auth-gateway/server.ts +++ b/packages/ai/src/auth-gateway/server.ts @@ -17,13 +17,14 @@ * POST /v1/messages → Anthropic messages in/out * POST /v1/responses → OpenAI Responses in/out */ -import { logger } from "@oh-my-pi/pi-utils"; +import { extractRetryHint, logger } from "@oh-my-pi/pi-utils"; import type { AuthStorage } from "../auth-storage"; import { Effort } from "../model-thinking"; import * as anthropicMessages from "../providers/anthropic-messages-server"; import * as openaiChat from "../providers/openai-chat-server"; import * as openaiResponses from "../providers/openai-responses-server"; import * as piNative from "../providers/pi-native-server"; +import { isUsageLimitError } from "../rate-limit-utils"; import { streamSimple } from "../stream"; import type { Api, AssistantMessageEventStream, Context, Model, SimpleStreamOptions } from "../types"; import { parseBind } from "../utils/parse-bind"; @@ -231,7 +232,17 @@ export function classifyGatewayError(err: unknown): { status: number; type: stri if ( // Match rate-limit phrasings without colliding with // `GenerateContentRequest`, `accelerate`, `iterate`, `deprecated`, etc. - /\brate[- _]?limit(?:s|ed|ing)?\b|\bquota(?:_exceeded| exceeded)?\b|\btoo[- _]many[- _]requests\b/i.test(message) + /\brate[- _]?limit(?:s|ed|ing)?\b|\bquota(?:_exceeded| exceeded)?\b|\btoo[- _]many[- _]requests\b/i.test( + message, + ) || + // Usage-limit phrasings emit no embedded status. Codex friendly text + // reads "You have hit your ChatGPT usage limit … Try again in ~158 + // min."; pi-ai's central `isUsageLimitError` already encodes every + // known provider variant, so reuse it instead of forking the regex. + // Without this branch the classifier falls through to the default + // 502/upstream_error, which is what callers were seeing when their + // account hit its cap. + isUsageLimitError(message) ) { return { status: 429, type: "rate_limit_error", message }; } @@ -266,9 +277,32 @@ function extractEmbeddedStatus(message: string): number | undefined { return Number.isFinite(code) && code >= 100 && code < 600 ? code : undefined; } +/** + * Hook fired by {@link streamSimple} when the upstream request fails in a + * way that's rotatable — today that's HTTP 401 (credential is bad) and + * usage-limit phrasing matched by {@link isUsageLimitError} (Codex's + * `usage_limit_reached`, Anthropic's `usage_limit_reached`, Google's + * `resource_exhausted`, …). The two cases need different storage actions: + * + * - **usage-limit** → {@link AuthStorage.markUsageLimitReached}. Marks just + * the current session's credential as temporarily blocked (honouring + * `retry-after` / `resets_at` hints when present) and returns `true` only + * when a sibling credential is still available. Burning the credential + * with `invalidateCredentialMatching` here would orphan accounts whose + * reset window is several hours away — exactly the bug this helper exists + * to avoid. + * - **auth-failure** → {@link AuthStorage.invalidateCredentialMatching}. + * Suspect/delete the row so it doesn't get re-picked next request. + * + * In both branches we return the next `getApiKey` result (sticky on the + * same `sessionId`) so streamSimple can transparently retry the pre-emit + * failure with a fresh credential. Returning `undefined` aborts the retry + * and surfaces the original error to the caller. + */ async function refreshGatewayApiKeyAfterAuthError( storage: AuthStorage, model: Model, + sessionId: string, provider: string, oldKey: string, error: unknown, @@ -276,14 +310,33 @@ async function refreshGatewayApiKeyAfterAuthError( format: string, peer: string, ): Promise { - await storage.invalidateCredentialMatching(provider, oldKey, signal); + const message = error instanceof Error ? error.message : String(error); + if (isUsageLimitError(message)) { + const retryAfterMs = extractRetryHint(undefined, message); + const switched = await storage.markUsageLimitReached(provider, sessionId, { + retryAfterMs, + baseUrl: model.baseUrl, + signal, + }); + logger.debug("auth-gateway retrying provider request after usage-limit block", { + format, + provider, + peer, + switched, + retryAfterMs, + error: message, + }); + if (!switched) return undefined; + return storage.getApiKey(provider, sessionId, { modelId: model.id, signal }); + } + await storage.invalidateCredentialMatching(provider, oldKey, { sessionId, signal }); logger.debug("auth-gateway retrying provider request after credential invalidation", { format, provider, peer, - error: error instanceof Error ? error.message : String(error), + error: message, }); - return storage.getApiKey(provider, undefined, { modelId: model.id, signal }); + return storage.getApiKey(provider, sessionId, { modelId: model.id, signal }); } function clientClosedResponse(route: { module: FormatModule }): Response { @@ -336,37 +389,12 @@ async function handleFormatEndpoint( return route.module.formatError(404, "invalid_request_error", `Unknown model: ${modelId}`); } - // pi-ai's stream() does NOT consult AuthStorage — the caller (us) is - // expected to resolve the credential and pass it as `options.apiKey`. - // For OAuth providers this returns the access token (refreshed via the - // broker override on AuthStorage when needed). - let apiKey: string | undefined; - try { - apiKey = await bootOpts.storage.getApiKey(model.provider, undefined, { - modelId: model.id, - signal: controller.signal, - }); - } catch (error) { - if (controller.signal.aborted) return clientClosedResponse(route); - const classified = classifyGatewayError(error); - logger.warn("auth-gateway getApiKey threw", { provider: model.provider, peer, error: classified.message }); - return route.module.formatError(classified.status, classified.type, classified.message); - } - if (controller.signal.aborted) return clientClosedResponse(route); - if (!apiKey) { - return route.module.formatError( - 401, - "authentication_error", - `No credential available for provider ${model.provider}`, - ); - } - - // Parse + validate against the strict format schema, rebuild as omp's - // canonical Context, dispatch through pi-ai's streamSimple, encode the - // canonical event stream back to the inbound format. There is no - // passthrough fast-path — every request flows through pi-ai so that - // credential-specific request shaping (OAuth Claude-Code prefix, beta - // headers, codex websocket transport, …) always applies. + // Parse the wire-format request BEFORE resolving the credential so we + // have a stable per-conversation `sessionId` to thread into AuthStorage. + // Sticky-credential tracking and `markUsageLimitReached` both key off + // this id; without it `getApiKey` would re-roundrobin every request + // and `markUsageLimitReached` would no-op (it can only mark the + // credential it last handed out to that session). let parsed: ParsedFormatRequest; try { parsed = route.module.parseRequest(body, req.headers); @@ -385,12 +413,45 @@ async function handleFormatEndpoint( } if (controller.signal.aborted) return clientClosedResponse(route); + // Sticky credential id: honour the client's `prompt_cache_key` when + // supplied (so external session ids align), otherwise derive from + // modelId + system + tools + first message. Mirrored into + // streamOpts.sessionId / promptCacheKey by `buildStreamOptions`. + const sessionId = parsed.options.promptCacheKey ?? deriveSessionId(parsed.modelId, parsed.context); + parsed.options.promptCacheKey ??= sessionId; + + // pi-ai's stream() does NOT consult AuthStorage — the caller (us) is + // expected to resolve the credential and pass it as `options.apiKey`. + // For OAuth providers this returns the access token (refreshed via the + // broker override on AuthStorage when needed). + let apiKey: string | undefined; + try { + apiKey = await bootOpts.storage.getApiKey(model.provider, sessionId, { + modelId: model.id, + signal: controller.signal, + }); + } catch (error) { + if (controller.signal.aborted) return clientClosedResponse(route); + const classified = classifyGatewayError(error); + logger.warn("auth-gateway getApiKey threw", { provider: model.provider, peer, error: classified.message }); + return route.module.formatError(classified.status, classified.type, classified.message); + } + if (controller.signal.aborted) return clientClosedResponse(route); + if (!apiKey) { + return route.module.formatError( + 401, + "authentication_error", + `No credential available for provider ${model.provider}`, + ); + } + const streamOpts = buildStreamOptions(parsed, model.api, controller.signal); streamOpts.apiKey = apiKey; streamOpts.onAuthError = (provider, oldKey, error) => refreshGatewayApiKeyAfterAuthError( bootOpts.storage, model, + sessionId, provider, oldKey, error, @@ -508,10 +569,17 @@ async function handlePiNative(bootOpts: AuthGatewayBootOptions, req: Request, pe if (!model) { return piNative.formatError(404, "invalid_request_error", `Unknown model: ${parsed.modelId}`); } + // Pi-native already parsed `streamOpts.sessionId` (when set by the + // client); fall back to the derived key so credential-stickiness lines + // up with cache-prefix stickiness — same identity used for both means + // the next turn of this conversation reuses the same credential until + // it hits a usage cap, then markUsageLimitReached can hand off. + const sessionId = parsed.options.sessionId ?? deriveSessionId(parsed.modelId, parsed.context); + parsed.options.sessionId ??= sessionId; let apiKey: string | undefined; try { - apiKey = await bootOpts.storage.getApiKey(model.provider, undefined, { + apiKey = await bootOpts.storage.getApiKey(model.provider, sessionId, { modelId: model.id, signal: controller.signal, }); @@ -539,6 +607,7 @@ async function handlePiNative(bootOpts: AuthGatewayBootOptions, req: Request, pe refreshGatewayApiKeyAfterAuthError( bootOpts.storage, model, + sessionId, provider, oldKey, error, @@ -554,10 +623,7 @@ async function handlePiNative(bootOpts: AuthGatewayBootOptions, req: Request, pe // headers — the client's values win when they collide. const captured = captureRequestHeaders(req.headers); streamOpts.headers = { ...captured, ...(streamOpts.headers ?? {}) }; - // Cache identity: explicit `sessionId` wins, then derive a stable key - // from model + system + tools + first message so Codex prefix caching - // engages on the same logical conversation across turns. - streamOpts.sessionId ??= deriveSessionId(parsed.modelId, parsed.context); + streamOpts.sessionId ??= sessionId; logger.info("auth-gateway request", { format: "pi-native", diff --git a/packages/ai/src/stream.ts b/packages/ai/src/stream.ts index 7bbe24fd0..b61b16011 100644 --- a/packages/ai/src/stream.ts +++ b/packages/ai/src/stream.ts @@ -45,6 +45,7 @@ import { } from "./providers/register-builtins"; import { isSyntheticModel, streamSynthetic } from "./providers/synthetic"; import { streamXAIResponses } from "./providers/xai-responses"; +import { isUsageLimitError } from "./rate-limit-utils"; import type { Api, AssistantMessage, @@ -318,6 +319,18 @@ function extractStatusFromAssistantError(message: AssistantMessage): number | un return extractHttpStatusFromError({ message: message.errorMessage }); } +function isRetryableUpstreamError(error: unknown, status: number | undefined, message: string | undefined): boolean { + // 401 means the credential is bad. Usage-limit phrasing (Codex's + // "You have hit your ChatGPT usage limit", Anthropic's "usage_limit_reached", + // Google's "resource_exhausted") means this account is parked but a + // sibling credential can usually pick the request up. Both are + // rotatable via `onAuthError` — the auth-gateway maps the former to + // `invalidateCredentialMatching` and the latter to `markUsageLimitReached`. + if (status === 401) return true; + void error; + return !!message && isUsageLimitError(message); +} + function createAssistantAuthError(message: AssistantMessage): Error & { status?: number } { const error: Error & { status?: number } = new Error(message.errorMessage ?? "Provider authentication failed"); const status = extractStatusFromAssistantError(message); @@ -359,7 +372,11 @@ export function streamSimple( !emittedReplayUnsafeEvent && captureAuthFailure && event.type === "error" && - extractStatusFromAssistantError(event.error) === 401 + isRetryableUpstreamError( + event.error, + extractStatusFromAssistantError(event.error), + event.error.errorMessage, + ) ) { return { error: createAssistantAuthError(event.error), bufferedEvents, terminalEvent: event }; } @@ -371,7 +388,15 @@ export function streamSimple( flushBuffered(); if (!outer.done) outer.end(await inner.result()); } catch (error) { - if (!emittedReplayUnsafeEvent && captureAuthFailure && extractHttpStatusFromError(error) === 401) { + if ( + !emittedReplayUnsafeEvent && + captureAuthFailure && + isRetryableUpstreamError( + error, + extractHttpStatusFromError(error), + error instanceof Error ? error.message : undefined, + ) + ) { return { error, bufferedEvents }; } flushBuffered(); diff --git a/packages/ai/test/auth-gateway-classify-error.test.ts b/packages/ai/test/auth-gateway-classify-error.test.ts index c3911f8c2..d780b266f 100644 --- a/packages/ai/test/auth-gateway-classify-error.test.ts +++ b/packages/ai/test/auth-gateway-classify-error.test.ts @@ -54,6 +54,27 @@ describe("auth-gateway classifyGatewayError", () => { expect(c.type).toBe("rate_limit_error"); }); + it("classifies Codex 'You have hit your ChatGPT usage limit' as 429", () => { + // Verbatim shape Codex returns from the `usage_limit_reached` branch + // in `parseCodexError`. No embedded `HTTP NNN`/`(NNN)`/`status NNN` + // token, no `rate limit`/`too many requests` wording — only the + // gateway's `isUsageLimitError` branch catches this. Previously it + // fell through to the default 502/upstream_error, which is why the + // `lg` retry loop kept looping instead of switching to another + // credential. + const c = classifyGatewayError( + new Error("You have hit your ChatGPT usage limit (pro plan). Try again in ~158 min."), + ); + expect(c.status).toBe(429); + expect(c.type).toBe("rate_limit_error"); + }); + + it("classifies generic 'usage_limit_reached' code text as 429", () => { + const c = classifyGatewayError(new Error('{"code":"usage_limit_reached","message":"…"}')); + expect(c.status).toBe(429); + expect(c.type).toBe("rate_limit_error"); + }); + it("does not match 'rate' inside camelCase or compound words", () => { // `Generate`, `iterate`, `deprecated`, `accelerate` all contain `rate` as // a substring and used to trip the classifier. diff --git a/packages/ai/test/google-gemini-cli-429.test.ts b/packages/ai/test/google-gemini-cli-429.test.ts index 934e6218d..6008f1fe1 100644 --- a/packages/ai/test/google-gemini-cli-429.test.ts +++ b/packages/ai/test/google-gemini-cli-429.test.ts @@ -80,6 +80,27 @@ describe("extractRetryHint – body text parsing", () => { expect(extractRetryHint(undefined, "try again in 12s")).toBe(12_000); }); + it("parses Codex 'Try again in ~X min.' (usage_limit_reached friendly text)", () => { + // Verbatim shape Codex's parseCodexError builds when usage_limit_reached + // arrives with a `resets_at` minutes-out reset time. Used to fall + // through to undefined → the gateway and TUI both had no retry-after + // signal to honour, so they defaulted to QUOTA_EXHAUSTED's 30-min + // blanket and rotated immediately even when the actual reset window + // was much longer. + expect(extractRetryHint(undefined, "Try again in ~158 min.")).toBe(158 * 60_000); + }); + + it("parses 'try again in X min' / 'X minutes' without the tilde", () => { + expect(extractRetryHint(undefined, "try again in 5 min")).toBe(5 * 60_000); + expect(extractRetryHint(undefined, "try again in 90 minutes")).toBe(90 * 60_000); + }); + + it("parses 'try again in X h' / 'X hour' / 'X hours'", () => { + expect(extractRetryHint(undefined, "try again in 2 h")).toBe(2 * 60 * 60_000); + expect(extractRetryHint(undefined, "try again in 1 hour")).toBe(60 * 60_000); + expect(extractRetryHint(undefined, "try again in 3 hours")).toBe(3 * 60 * 60_000); + }); + it("returns undefined when body contains no recognised delay pattern", () => { expect(extractRetryHint(undefined, "Quota exceeded, please try again later")).toBeUndefined(); }); diff --git a/packages/ai/test/stream-auth-retry.test.ts b/packages/ai/test/stream-auth-retry.test.ts index 6358fb274..d019e57ce 100644 --- a/packages/ai/test/stream-auth-retry.test.ts +++ b/packages/ai/test/stream-auth-retry.test.ts @@ -244,4 +244,131 @@ describe("streamSimple auth retry", () => { expect(keys).toEqual(["old-key", "new-key"]); expect(authCalls).toBe(1); }); + + it("retries on a thrown usage_limit_reached error before any event has been emitted", async () => { + const keys: Array = []; + const errors: unknown[] = []; + registerCustomApi( + API, + (_model: Model, _context: Context, options?: SimpleStreamOptions) => { + keys.push(options?.apiKey); + const stream = new AssistantMessageEventStream(); + queueMicrotask(() => { + if (keys.length === 1) { + stream.fail( + Object.assign( + new Error("You have hit your ChatGPT usage limit (pro plan). Try again in ~158 min."), + { status: 429 }, + ), + ); + return; + } + const message = assistant(["ok"]); + stream.push({ type: "start", partial: message }); + stream.push({ type: "done", reason: "stop", message }); + }); + return stream; + }, + SOURCE_ID, + ); + + const stream = streamSimple(model(), context, { + apiKey: "old-key", + onAuthError: async (_provider, _oldKey, error) => { + errors.push(error); + return "new-key"; + }, + }); + + for await (const _event of stream) { + // drain + } + + expect((await stream.result()).content).toEqual([{ type: "text", text: "ok" }]); + expect(keys).toEqual(["old-key", "new-key"]); + expect(errors).toHaveLength(1); + // The error surfaced to onAuthError carries the original 429 status so + // the gateway's refresh hook can branch on it. + expect((errors[0] as { status?: number }).status).toBe(429); + expect((errors[0] as Error).message).toMatch(/usage limit/i); + }); + + it("retries when a provider emits a usage_limit_reached error event before content", async () => { + const keys: Array = []; + registerCustomApi( + API, + (_model: Model, _context: Context, options?: SimpleStreamOptions) => { + keys.push(options?.apiKey); + const stream = new AssistantMessageEventStream(); + queueMicrotask(() => { + if (keys.length === 1) { + stream.push({ type: "start", partial: assistant() }); + stream.push({ + type: "error", + reason: "error", + // errorStatus deliberately omitted: matches how the codex + // provider serializes the message-only path that used to + // 502 in the gateway. + error: assistantError("You have hit your ChatGPT usage limit (pro plan). Try again in ~158 min."), + }); + return; + } + const message = assistant(["ok"]); + stream.push({ type: "start", partial: message }); + stream.push({ type: "done", reason: "stop", message }); + }); + return stream; + }, + SOURCE_ID, + ); + + const stream = streamSimple(model(), context, { + apiKey: "old-key", + onAuthError: async () => "new-key", + }); + + for await (const _event of stream) { + // drain + } + + expect((await stream.result()).content).toEqual([{ type: "text", text: "ok" }]); + expect(keys).toEqual(["old-key", "new-key"]); + }); + + it("surfaces the original usage_limit error when the retry callback declines", async () => { + const keys: Array = []; + const original = Object.assign(new Error("You have hit your ChatGPT usage limit (pro plan)."), { status: 429 }); + registerCustomApi( + API, + (_model: Model, _context: Context, options?: SimpleStreamOptions) => { + keys.push(options?.apiKey); + const stream = new AssistantMessageEventStream(); + queueMicrotask(() => stream.fail(original)); + return stream; + }, + SOURCE_ID, + ); + + // Callback returns undefined → no sibling credential to rotate to. + // The original failure must reach the caller untouched so the client + // can decide what to do (back off, surface to user, …). + const stream = streamSimple(model(), context, { + apiKey: "old-key", + onAuthError: async () => undefined, + }); + + let caught: unknown; + try { + for await (const _event of stream) { + // drain + } + } catch (error) { + caught = error; + } + + expect(caught).toBe(original); + // Single attempt — the inner stream is only re-invoked when a new key + // is provided. + expect(keys).toEqual(["old-key"]); + }); }); diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 2deb69659..25bcb3589 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -16,6 +16,22 @@ - Modified hashline anchor syntax to require explicit range notation `A-B:` instead of shorthand `A:` for single-line operations - Updated hashline description in settings to clarify pure insert context behavior without arrow notation +### Removed + +- Removed the `edit.hashlineAutoDropPureInsertDuplicates` setting +- Removed the `edit.hashlineAutoDropPureInsertDuplicates` setting from configuration and execution paths + +### Fixed + +- Fixed `eval` tool to resize large displayed images and append dimension notes to text output +- Fixed `write` tool to strip malformed or loose hashline section headers before writing file content +- Fixed `eval` tool image rendering to resize displayed images before returning them and append image-dimension notes to text output +- Fixed `write` tool output sanitation to strip malformed or loose hashline section headers before writing file content +- Fixed `omp auth-broker serve` crashing at startup with `logger.setTransports is not a function` — switched the call site to `import { setTransports } from "@oh-my-pi/pi-utils/logger"`, bypassing the `logger` namespace re-export that some Bun versions failed to expose at runtime +- Fixed `omp auth-gateway` returning `502 upstream_error` and refusing to rotate credentials when a provider responded with a non-401 usage-limit error (Codex `usage_limit_reached`, Anthropic `usage_limit_reached`, Google `resource_exhausted`). `classifyGatewayError` now reuses `pi-ai`'s central `isUsageLimitError` heuristic and reports those failures as `429 rate_limit_error`. `streamSimple`'s pre-emit retry hook fires on usage-limit phrasing in addition to HTTP 401; the gateway's refresh callback branches on the error type and calls `AuthStorage.markUsageLimitReached(provider, sessionId, { retryAfterMs })` — temporarily blocking just the exhausted credential and surfacing the next sibling — instead of `invalidateCredentialMatching`, which would have suspect/deleted the row. The same branching is wired into the coding-agent `streamFn` callback so subscription multi-account rotation works the same on both surfaces. +- Fixed `extractRetryHint` not recognising Codex's `Try again in ~N min.` / `… hour` / `… hours` phrasing, which left the gateway and TUI without a server-suggested retry window when an upstream account hit its usage cap. The shared `try again in` pattern now accepts `min`, `minutes`, `mins`, `h`, `hr`, `hour`, `hours` units in addition to `ms` / `s` / `sec`, and tolerates a leading `~` and embedded whitespace. +- Fixed the auth-gateway threading `sessionId: undefined` into `AuthStorage.getApiKey`, which left `#sessionLastCredential` empty and made `markUsageLimitReached` a no-op for gateway-mediated requests. Both `/v1/chat/completions`-style endpoints and the `/v1/pi/stream` fast path now derive a stable `sessionId` from the client's `prompt_cache_key` (or the existing model+system+tools+first-message hash when absent) and reuse the same identity for credential-stickiness and prefix-cache routing. + ## [15.5.7] - 2026-05-27 ### Added - `providers.openrouterVariant` setting (Settings → Providers → "OpenRouter Routing") to default OpenRouter requests to a routing-variant suffix (`:nitro`, `:floor`, `:online`, `:exacto`). Selectors that already name a variant (e.g. `openrouter/anthropic/claude-haiku:nitro`) keep precedence. diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 2e8bad681..24c1ea4e4 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -10,6 +10,7 @@ import { } from "@oh-my-pi/pi-agent-core"; import { type CredentialDisabledEvent, + isUsageLimitError, type Message, type Model, type SimpleStreamOptions, @@ -23,6 +24,7 @@ import type { Component } from "@oh-my-pi/pi-tui"; import { $env, $flag, + extractRetryHint, getAgentDbPath, getAgentDir, getProjectDir, @@ -1885,13 +1887,36 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} ...streamOptions, openrouterVariant: streamOptions?.openrouterVariant ?? openrouterVariant, onAuthError: async (provider, oldKey, error) => { + const message = error instanceof Error ? error.message : String(error); + // streamSimple invokes this for both 401 auth failures AND + // rotatable usage-limit errors (Codex usage_limit_reached, + // Anthropic usage_limit_reached, etc.). The two need + // different storage actions: a real 401 means the credential + // is bad and should be marked suspect; a usage limit just + // means this account is parked until reset and should be + // temporarily blocked so a sibling can pick the request up. + if (isUsageLimitError(message)) { + const retryAfterMs = extractRetryHint(undefined, message); + const switched = await modelRegistry.authStorage.markUsageLimitReached(provider, agent.sessionId, { + retryAfterMs, + signal: streamOptions?.signal, + }); + logger.debug("Retrying provider request after usage-limit block", { + provider, + switched, + retryAfterMs, + error: message, + }); + if (!switched) return undefined; + return modelRegistry.getApiKeyForProvider(provider, agent.sessionId); + } await modelRegistry.authStorage.invalidateCredentialMatching(provider, oldKey, { signal: streamOptions?.signal, sessionId: agent.sessionId, }); logger.debug("Retrying provider request after credential invalidation", { provider, - error: error instanceof Error ? error.message : String(error), + error: message, }); return modelRegistry.getApiKeyForProvider(provider, agent.sessionId); }, diff --git a/packages/utils/src/fetch-retry.ts b/packages/utils/src/fetch-retry.ts index b521250fe..06961f0bc 100644 --- a/packages/utils/src/fetch-retry.ts +++ b/packages/utils/src/fetch-retry.ts @@ -6,8 +6,10 @@ const QUOTA_RESET_PATTERN = /reset after (?:(\d+)h)?(?:(\d+)m)?(\d+(?:\.\d+)?)s/ const PLEASE_RETRY_PATTERN = /Please retry in ([0-9.]+)(ms|s)/i; // JSON field: "retryDelay": "34.074824224s" const RETRY_DELAY_FIELD_PATTERN = /"retryDelay":\s*"([0-9.]+)(ms|s)"/i; -// "try again in 250ms" / "try again in 12s" / "try again in 12sec" -const TRY_AGAIN_PATTERN = /try again in\s+(\d+(?:\.\d+)?)\s*(ms|s)(?:ec)?/i; +// "try again in 250ms" / "try again in 12s" / "try again in 12sec" / +// "try again in 5 min" / "try again in ~158 min." / "try again in 2h" / +// "try again in 90 minutes" / "try again in 1 hour" +const TRY_AGAIN_PATTERN = /try again in\s+~?\s*([0-9.]+)\s*(ms|sec|s|minutes?|mins?|m|hours?|hrs?|h)\b/i; /** * Server-suggested retry delay extraction. Merges the patterns historically used @@ -22,7 +24,7 @@ const TRY_AGAIN_PATTERN = /try again in\s+(\d+(?:\.\d+)?)\s*(ms|s)(?:ec)?/i; * - `Your quota will reset after 18h31m10s` / `10m15s` / `39s` * - `Please retry in 250ms` / `Please retry in 12s` * - `"retryDelay": "34.074824224s"` (JSON error detail field) - * - `try again in 250ms` / `try again in 12s` / `try again in 12sec` + * - `try again in 250ms` / `try again in 12s` / `try again in 5 min` / `try again in ~158 min` * * Returns `undefined` if no signal is found. */ @@ -68,13 +70,38 @@ export function extractRetryHint(source: Response | Headers | null | undefined, if (match?.[1]) { const value = Number.parseFloat(match[1]); if (Number.isFinite(value) && value > 0) { - return match[2]!.toLowerCase() === "ms" ? value : value * 1000; + const unitMs = unitToMs(match[2]!); + if (unitMs !== undefined) return value * unitMs; } } } return undefined; } +function unitToMs(unit: string): number | undefined { + switch (unit.toLowerCase()) { + case "ms": + return 1; + case "s": + case "sec": + return 1000; + case "m": + case "min": + case "mins": + case "minute": + case "minutes": + return 60_000; + case "h": + case "hr": + case "hrs": + case "hour": + case "hours": + return 60 * 60_000; + default: + return undefined; + } +} + export interface FetchWithRetryOptions extends RequestInit { /** Total fetch attempts (initial + retries). Default `5`. */ maxAttempts?: number;