fix: resolved auth-gateway handling of 429 usage-limit responses
- Classified usage-limit gateway responses as `429 rate_limit_error` in auth handling paths. - Aligned auth-gateway and pi-native key retrieval with derived `sessionId` for `getApiKey` lookups. - Handled usage-limit auth failures by rotating credentials with retry hints and returning undefined when none available. - Replaced stream auth checks with retryable-upstream logic for 401 and usage-limit errors before content. - Expanded `extractRetryHint` parsing for `~`, `sec`, `ms`, and minute/hour units. - Added coverage for classifyGatewayError, retry-hint parsing variants, and stream-auth retry edge cases.
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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<Api>,
|
||||
sessionId: string,
|
||||
provider: string,
|
||||
oldKey: string,
|
||||
error: unknown,
|
||||
@@ -276,14 +310,33 @@ async function refreshGatewayApiKeyAfterAuthError(
|
||||
format: string,
|
||||
peer: string,
|
||||
): Promise<string | undefined> {
|
||||
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",
|
||||
|
||||
@@ -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<TApi extends Api>(
|
||||
!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<TApi extends Api>(
|
||||
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();
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
|
||||
@@ -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<string | undefined> = [];
|
||||
const errors: unknown[] = [];
|
||||
registerCustomApi(
|
||||
API,
|
||||
(_model: Model<Api>, _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<string | undefined> = [];
|
||||
registerCustomApi(
|
||||
API,
|
||||
(_model: Model<Api>, _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<string | undefined> = [];
|
||||
const original = Object.assign(new Error("You have hit your ChatGPT usage limit (pro plan)."), { status: 429 });
|
||||
registerCustomApi(
|
||||
API,
|
||||
(_model: Model<Api>, _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"]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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);
|
||||
},
|
||||
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user