diff --git a/bun.lock b/bun.lock index 329084661e..849f4f590c 100644 --- a/bun.lock +++ b/bun.lock @@ -22,7 +22,7 @@ ], "overrides": { "@hono/node-server": "2.1.0", - "fast-uri": "^3.1.5", + "fast-uri": "^3.1.6", "hono": "4.13.1", "ip-address": "^10.4.0", }, @@ -187,7 +187,7 @@ "fast-deep-equal": ["fast-deep-equal@3.1.3", "", {}, "sha512-f3qQ9oQy9j2AhBe/H9VC91wLmKBCCU/gDOnKNAYG5hswO7BLKj09Hc5HYNz9cGI++xlpDCIgDaitVs03ATR84Q=="], - "fast-uri": ["fast-uri@3.1.5", "", {}, "sha512-gHwA1O9LDIcKunMKhObS/HimwtehO1nPUECKAu5TpKgaO19fcWEl4bliWe1jWxVFvIXztJjjQ4L8XQ1EU9f7Jw=="], + "fast-uri": ["fast-uri@3.1.7", "", {}, "sha512-dOvZVzjdZdz7phd9v6jCbwxrBW3fK6n8Rc0CtdmM4bumzMnxywBYhuph6J819RRw/ku+rLbelwfMunktuzVVHg=="], "finalhandler": ["finalhandler@2.1.1", "", { "dependencies": { "debug": "^4.4.0", "encodeurl": "^2.0.0", "escape-html": "^1.0.3", "on-finished": "^2.4.1", "parseurl": "^1.3.3", "statuses": "^2.0.1" } }, "sha512-S8KoZgRZN+a5rNwqTxlZZePjT/4cnm0ROV70LedRHZ0p8u9fRID0hJUZQpkKLzro8LfmC8sx23bY6tVNxv8pQA=="], diff --git a/package.json b/package.json index 8940f761fb..03dc32c6be 100644 --- a/package.json +++ b/package.json @@ -74,7 +74,7 @@ }, "overrides": { "@hono/node-server": "2.1.0", - "fast-uri": "^3.1.5", + "fast-uri": "^3.1.6", "hono": "4.13.1", "ip-address": "^10.4.0" }, diff --git a/src/adapters/openai-chat.ts b/src/adapters/openai-chat.ts index 8b7d9c8614..cc2d965d38 100644 --- a/src/adapters/openai-chat.ts +++ b/src/adapters/openai-chat.ts @@ -2030,13 +2030,13 @@ export function createOpenAIChatAdapter(provider: OcxProviderConfig): ProviderAd let parsed: unknown; try { parsed = await response.json(); - } catch (error) { + } catch { tierMetadata?.markResponseUnparseable(); - throw error; + return [{ type: "error", message: "malformed upstream JSON response" }]; } if (parsed === null || typeof parsed !== "object" || Array.isArray(parsed)) { tierMetadata?.markResponseUnparseable(); - throw new Error("upstream response was not a JSON object"); + return [{ type: "error", message: "malformed upstream JSON response" }]; } const json = parsed as Record; if (Object.hasOwn(json, "service_tier")) { diff --git a/src/bridge.ts b/src/bridge.ts index c50cac8c04..95fa259ed4 100644 --- a/src/bridge.ts +++ b/src/bridge.ts @@ -1375,11 +1375,15 @@ export function bridgeToResponsesSSE( if (currentToolCall) failCurrentToolCall(); if (currentWebSearch) closeCurrentWebSearch("failed", []); releasePendingWebSources(); - const failure = responseError( - 500, - "proxy_error", - redactSecretString(err instanceof Error ? err.message : String(err)), - ); + // Unexpected iterator/read exceptions can contain provider URLs, socket details, + // gateway payload fragments, or credential-bearing diagnostics. Preserve only a + // recognized policy identity after secret redaction; arbitrary transport text must + // collapse to a stable public terminal once output has committed. + const thrownMessage = redactSecretString(err instanceof Error ? err.message : String(err)); + const classifiedThrown = adapterFailureFromMessage(thrownMessage); + const failure = isCyberPolicyCode(classifiedThrown.error.code) + ? classifiedThrown.error + : responseError(502, "upstream_error", "Provider stream failed unexpectedly"); emit("response.failed", { response: { ...responseSnapshot("failed", finishedItems), @@ -2064,7 +2068,11 @@ export function formatErrorResponse( options?: { code?: string | null; retryAfter?: string | null }, ): Response { const error = classifyError(status, type, message); - if (isCyberPolicyCode(options?.code)) { + const explicitCode = typeof options?.code === "string" && options.code.trim() + ? options.code.trim() + : null; + if (explicitCode) error.code = explicitCode; + if (isCyberPolicyCode(explicitCode)) { error.code = CYBER_POLICY_ERROR_CODE; error.type = cyberPolicyErrorType(type); } diff --git a/src/combos/cooldown-disk.ts b/src/combos/cooldown-disk.ts new file mode 100644 index 0000000000..6ceb407f70 --- /dev/null +++ b/src/combos/cooldown-disk.ts @@ -0,0 +1,79 @@ +/** Persist only long-lived combo quota cooldowns across proxy restarts. */ +import { existsSync, readFileSync } from "node:fs"; +import { join } from "node:path"; +import { atomicWriteFile, getConfigDir } from "../config"; + +const FILENAME = "combo-quota-cooldowns.json"; +const MAX_FUTURE_MS = 31 * 24 * 60 * 60_000; +const PERSIST_DEBOUNCE_MS = 250; + +type DiskFile = { version: 1; rows: Record }; +let persistTimer: ReturnType | null = null; +let pendingRows: (() => Iterable<[string, number]>) | null = null; +let pendingDirectory: string | null = null; + +export function comboQuotaCooldownStoreDirectory(): string { + return getConfigDir(); +} + +export function readPersistedComboQuotaCooldowns(directory: string, now = Date.now()): Map { + const rows = new Map(); + try { + const path = join(directory, FILENAME); + if (!existsSync(path)) return rows; + const parsed = JSON.parse(readFileSync(path, "utf8")) as DiskFile; + if (!parsed || parsed.version !== 1 || !parsed.rows || typeof parsed.rows !== "object") return rows; + for (const [key, until] of Object.entries(parsed.rows)) { + if (typeof until !== "number" || !Number.isFinite(until)) continue; + if (until <= now || until - now > MAX_FUTURE_MS) continue; + rows.set(key, until); + } + } catch { + // Missing/corrupt best-effort state must never block routing. + } + return rows; +} + +function persistNow(directory: string, rows: Iterable<[string, number]>, now = Date.now()): void { + try { + const out: Record = {}; + for (const [key, until] of rows) { + if (!Number.isFinite(until) || until <= now || until - now > MAX_FUTURE_MS) continue; + out[key] = until; + } + atomicWriteFile(join(directory, FILENAME), `${JSON.stringify({ version: 1, rows: out } satisfies DiskFile)}\n`); + } catch { + // Best-effort persistence only. Runtime failover remains authoritative. + } +} + +export function schedulePersistComboQuotaCooldowns(directory: string, rows: () => Iterable<[string, number]>): void { + pendingDirectory = directory; + pendingRows = rows; + if (persistTimer) clearTimeout(persistTimer); + persistTimer = setTimeout(() => { + persistTimer = null; + const snapshot = pendingRows; + const directory = pendingDirectory; + pendingRows = null; + pendingDirectory = null; + if (snapshot && directory) persistNow(directory, snapshot()); + }, PERSIST_DEBOUNCE_MS); +} + +export function flushComboQuotaCooldownPersistForTests(now = Date.now()): void { + if (persistTimer) clearTimeout(persistTimer); + persistTimer = null; + const snapshot = pendingRows; + const directory = pendingDirectory; + pendingRows = null; + pendingDirectory = null; + if (snapshot && directory) persistNow(directory, snapshot(), now); +} + +export function cancelPendingComboQuotaCooldownPersist(): void { + if (persistTimer) clearTimeout(persistTimer); + persistTimer = null; + pendingRows = null; + pendingDirectory = null; +} diff --git a/src/combos/failover.ts b/src/combos/failover.ts index ae0c044be2..6de381c501 100644 --- a/src/combos/failover.ts +++ b/src/combos/failover.ts @@ -1,6 +1,11 @@ import { classifyError, isCyberPolicyCode } from "../lib/errors"; import type { OcxComboTarget } from "../types"; import { targetKey } from "./types"; +import { + comboQuotaCooldownStoreDirectory, + readPersistedComboQuotaCooldowns, + schedulePersistComboQuotaCooldowns, +} from "./cooldown-disk"; import { captureConfigGeneration, sweepExpiredOnWrite, @@ -13,6 +18,8 @@ interface TargetCooldown { const DEFAULT_COOLDOWN_MS = 60_000; const MAX_COOLDOWN_MS = 10 * 60_000; +/** Long provider quota windows may legitimately span days; keep this bounded. */ +const MAX_QUOTA_COOLDOWN_MS = 31 * 24 * 60 * 60_000; /** Short cooldown for request-rate 429s (for example provider code 1302) that omit Retry-After. */ export const COMBO_REQUEST_RATE_COOLDOWN_MS = 5_000; @@ -36,10 +43,18 @@ const HTTP_MONTH_INDEX: Record = { jul: 6, aug: 7, sep: 8, oct: 9, nov: 10, dec: 11, }; -/** Map<`${comboId}\0${provider/model}`, TargetCooldown> */ +/** Target keys are combo-local; provider keys represent provider-wide evidence shared by every combo. */ +const PROVIDER_COOLDOWN_PREFIX = "\u0001provider\0"; const targetCooldowns = new Map(); +const persistedQuotaCooldownKeys = new Set(); +let quotaCooldownPersistenceDirectory: string | undefined; let lastReconciledGeneration = 0; let liveComboTargets = new Set(); +let liveProviders = new Set(); + +function providerCooldownMapKey(provider: string): string { + return `${PROVIDER_COOLDOWN_PREFIX}${provider}`; +} function cooldownMapKey( comboId: string, @@ -48,6 +63,40 @@ function cooldownMapKey( return `${comboId}\0${targetKey(target)}`; } +function persistedQuotaCooldownRows(): Iterable<[string, number]> { + return [...persistedQuotaCooldownKeys].flatMap(key => { + const state = targetCooldowns.get(key); + return state ? [[key, state.cooldownUntil] as [string, number]] : []; + }); +} + +function scheduleQuotaCooldownPersistence(): void { + if (!quotaCooldownPersistenceDirectory) return; + schedulePersistComboQuotaCooldowns( + quotaCooldownPersistenceDirectory, + () => persistedQuotaCooldownRows(), + ); +} + +/** Hydrate only durable long quota cooldowns; transient transport/rate cooldowns stay process-local. */ +export function hydrateComboQuotaCooldownsFromDisk(now = Date.now()): void { + const directory = comboQuotaCooldownStoreDirectory(); + if (quotaCooldownPersistenceDirectory === directory) return; + for (const key of persistedQuotaCooldownKeys) targetCooldowns.delete(key); + persistedQuotaCooldownKeys.clear(); + quotaCooldownPersistenceDirectory = directory; + for (const [key, cooldownUntil] of readPersistedComboQuotaCooldowns(directory, now)) { + targetCooldowns.set(key, { cooldownUntil }); + persistedQuotaCooldownKeys.add(key); + } +} + +export function resetComboQuotaCooldownPersistenceForTests(): void { + for (const key of persistedQuotaCooldownKeys) targetCooldowns.delete(key); + persistedQuotaCooldownKeys.clear(); + quotaCooldownPersistenceDirectory = undefined; +} + function parseUtcDateParts( year: number, monthName: string, @@ -113,9 +162,10 @@ function parseHttpDate(value: string, now: number): number | undefined { ); } -export function parseRetryAfterMs( +function parseRetryAfterMsWithLimit( value: string | null | undefined, - now = Date.now(), + now: number, + maxMs: number, options?: { preserveImmediate?: boolean }, ): number | undefined { const text = value?.trim(); @@ -126,26 +176,87 @@ export function parseRetryAfterMs( Number.isFinite(seconds) && (seconds > 0 || (options?.preserveImmediate && seconds === 0)) ) { - return Math.min(Math.max(Math.ceil(seconds * 1000), 1), MAX_COOLDOWN_MS); + return Math.min(Math.max(Math.ceil(seconds * 1000), 1), maxMs); } } const timestamp = parseHttpDate(text, now); if (timestamp === undefined) return undefined; const delay = timestamp - now; - if (delay > 0) return Math.min(delay, MAX_COOLDOWN_MS); + if (delay > 0) return Math.min(delay, maxMs); return options?.preserveImmediate ? 1 : undefined; } +export function parseRetryAfterMs( + value: string | null | undefined, + now = Date.now(), + options?: { preserveImmediate?: boolean }, +): number | undefined { + return parseRetryAfterMsWithLimit(value, now, MAX_COOLDOWN_MS, options); +} + +function parseQuotaResetHintMs(message: string): number | undefined { + const match = /\bresets?\s+in\s+((?:\d+(?:\.\d+)?\s*(?:days?|hours?|hrs?|minutes?|mins?|seconds?|secs?)\s*){1,4})/i.exec(message); + if (!match?.[1]) return undefined; + const unitMs: Record = { + day: 86_400_000, days: 86_400_000, + hour: 3_600_000, hours: 3_600_000, hr: 3_600_000, hrs: 3_600_000, + minute: 60_000, minutes: 60_000, min: 60_000, mins: 60_000, + second: 1_000, seconds: 1_000, sec: 1_000, secs: 1_000, + }; + let total = 0; + const parts = match[1].matchAll(/(\d+(?:\.\d+)?)\s*(days?|hours?|hrs?|minutes?|mins?|seconds?|secs?)/gi); + for (const part of parts) { + const value = Number(part[1]); + const unit = unitMs[part[2]!.toLowerCase()]; + if (!Number.isFinite(value) || value < 0 || unit === undefined) return undefined; + total += value * unit; + } + if (!Number.isFinite(total) || total <= 0) return undefined; + return Math.min(Math.ceil(total), MAX_QUOTA_COOLDOWN_MS); +} + +function isDurableQuotaFailure( + status: number | undefined, + message: string, + code?: string | null, +): boolean { + return isProviderScopedQuotaCap(status, message, code) + || QUOTA_LIMIT_CODES.has(normalizedFailureCode(code)); +} + +export function knownComboFailureRecoveryMs(input: { + retryAfter?: string | null; + status?: number; + code?: string | null; + message?: string; + now?: number; +}): number | undefined { + const now = input.now ?? Date.now(); + const durableQuota = isDurableQuotaFailure(input.status, input.message ?? "", input.code); + const maxMs = durableQuota ? MAX_QUOTA_COOLDOWN_MS : MAX_COOLDOWN_MS; + return parseRetryAfterMsWithLimit(input.retryAfter, now, maxMs) + ?? (durableQuota ? parseQuotaResetHintMs(input.message ?? "") : undefined); +} + export function isComboTargetInCooldown( comboId: string, target: Pick, now = Date.now(), ): boolean { + const providerKey = providerCooldownMapKey(target.provider); + const providerEntry = targetCooldowns.get(providerKey); + if (providerEntry) { + if (providerEntry.cooldownUntil > now) return true; + targetCooldowns.delete(providerKey); + if (persistedQuotaCooldownKeys.delete(providerKey)) scheduleQuotaCooldownPersistence(); + } + const key = cooldownMapKey(comboId, target); const entry = targetCooldowns.get(key); if (!entry) return false; if (entry.cooldownUntil <= now) { targetCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) scheduleQuotaCooldownPersistence(); return false; } return true; @@ -171,14 +282,24 @@ export function isTransientRequestRateLimit(input: { return text.includes("rate limit reached for requests"); } -export function remainingComboCooldownMs(comboId: string, now = Date.now()): number | undefined { +export function remainingComboCooldownMs( + comboId: string, + now = Date.now(), + providers?: Iterable, +): number | undefined { const prefix = `${comboId}\0`; + const providerSet = providers === undefined ? undefined : new Set(providers); let soonest: number | undefined; for (const [key, cooldown] of targetCooldowns) { - if (!key.startsWith(prefix)) continue; + const comboLocal = key.startsWith(prefix); + const providerGlobal = providerSet !== undefined + && key.startsWith(PROVIDER_COOLDOWN_PREFIX) + && providerSet.has(key.slice(PROVIDER_COOLDOWN_PREFIX.length)); + if (!comboLocal && !providerGlobal) continue; const remaining = cooldown.cooldownUntil - now; if (remaining <= 0) { targetCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) scheduleQuotaCooldownPersistence(); continue; } if (soonest === undefined || remaining < soonest) soonest = remaining; @@ -186,12 +307,61 @@ export function remainingComboCooldownMs(comboId: string, now = Date.now()): num return soonest; } -export function comboCooldownRetryAfterSeconds(comboId: string, now = Date.now()): string | undefined { - const remainingMs = remainingComboCooldownMs(comboId, now); +export function comboCooldownRetryAfterSeconds( + comboId: string, + now = Date.now(), + providers?: Iterable, +): string | undefined { + const remainingMs = remainingComboCooldownMs(comboId, now, providers); if (remainingMs === undefined) return undefined; return String(Math.max(1, Math.ceil(remainingMs / 1000))); } +export function coolComboProvider( + provider: string, + options?: { + retryAfter?: string | null; + now?: number; + cooldownMs?: number; + writerGeneration?: number; + status?: number; + code?: string | null; + message?: string; + }, +): void { + const now = options?.now ?? Date.now(); + const writerGeneration = options?.writerGeneration ?? captureConfigGeneration(); + if (writerGeneration < lastReconciledGeneration && !liveProviders.has(provider)) return; + const durableQuota = isDurableQuotaFailure( + options?.status, + options?.message ?? "", + options?.code, + ); + const maxCooldownMs = durableQuota ? MAX_QUOTA_COOLDOWN_MS : MAX_COOLDOWN_MS; + const cooldownMs = options?.cooldownMs + ?? knownComboFailureRecoveryMs({ + retryAfter: options?.retryAfter, + status: options?.status, + code: options?.code, + message: options?.message, + now, + }) + ?? (isTransientRequestRateLimit({ + status: options?.status, + code: options?.code, + message: options?.message, + }) ? COMBO_REQUEST_RATE_COOLDOWN_MS : DEFAULT_COOLDOWN_MS); + const key = providerCooldownMapKey(provider); + const cooldownUntil = now + Math.min(Math.max(cooldownMs, 1), maxCooldownMs); + targetCooldowns.set(key, { cooldownUntil }); + const durableLongQuota = durableQuota && cooldownUntil - now > MAX_COOLDOWN_MS; + const persistenceChanged = durableLongQuota + ? (persistedQuotaCooldownKeys.add(key), true) + : persistedQuotaCooldownKeys.delete(key); + if (persistenceChanged) scheduleQuotaCooldownPersistence(); + sweepExpiredOnWrite(now); +} + export function coolComboTarget( comboId: string, target: Pick, @@ -209,24 +379,61 @@ export function coolComboTarget( const writerGeneration = options?.writerGeneration ?? captureConfigGeneration(); const ownerKey = `${comboId}::${targetKey(target)}`; if (writerGeneration < lastReconciledGeneration && !liveComboTargets.has(ownerKey)) return; + const durableQuota = isDurableQuotaFailure( + options?.status, + options?.message ?? "", + options?.code, + ); + const maxCooldownMs = durableQuota ? MAX_QUOTA_COOLDOWN_MS : MAX_COOLDOWN_MS; const cooldownMs = options?.cooldownMs - ?? parseRetryAfterMs(options?.retryAfter, now) + ?? knownComboFailureRecoveryMs({ + retryAfter: options?.retryAfter, + status: options?.status, + code: options?.code, + message: options?.message, + now, + }) ?? (isTransientRequestRateLimit({ status: options?.status, code: options?.code, message: options?.message, }) ? COMBO_REQUEST_RATE_COOLDOWN_MS : DEFAULT_COOLDOWN_MS); - targetCooldowns.set(cooldownMapKey(comboId, target), { - cooldownUntil: now + Math.min(Math.max(cooldownMs, 1), MAX_COOLDOWN_MS), - }); + const key = cooldownMapKey(comboId, target); + const cooldownUntil = now + Math.min(Math.max(cooldownMs, 1), maxCooldownMs); + targetCooldowns.set(key, { cooldownUntil }); + const durableLongQuota = durableQuota && cooldownUntil - now > MAX_COOLDOWN_MS; + const persistenceChanged = durableLongQuota + ? (persistedQuotaCooldownKeys.add(key), true) + : persistedQuotaCooldownKeys.delete(key); + if (persistenceChanged) scheduleQuotaCooldownPersistence(); sweepExpiredOnWrite(now); } export function reconcileComboTargetCooldowns(context: GenerationContext): number { if (context.generation <= lastReconciledGeneration) return 0; + let removed = 0; + let persistenceChanged = false; + for (const key of targetCooldowns.keys()) { + let live = false; + if (key.startsWith(PROVIDER_COOLDOWN_PREFIX)) { + live = context.providerNames.has(key.slice(PROVIDER_COOLDOWN_PREFIX.length)); + } else { + const separator = key.indexOf("\0"); + if (separator >= 0) { + const ownerKey = `${key.slice(0, separator)}::${key.slice(separator + 1)}`; + live = context.comboTargets.has(ownerKey); + } + } + if (live) continue; + targetCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) persistenceChanged = true; + removed += 1; + } + if (persistenceChanged) scheduleQuotaCooldownPersistence(); liveComboTargets = new Set(context.comboTargets); + liveProviders = new Set(context.providerNames); lastReconciledGeneration = context.generation; - return 0; + return removed; } export function sweepExpiredComboTargetCooldowns(now = Date.now()): number { @@ -234,6 +441,7 @@ export function sweepExpiredComboTargetCooldowns(now = Date.now()): number { for (const [key, cooldown] of targetCooldowns) { if (cooldown.cooldownUntil > now) continue; targetCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) scheduleQuotaCooldownPersistence(); removed += 1; } return removed; @@ -242,18 +450,26 @@ export function sweepExpiredComboTargetCooldowns(now = Date.now()): number { export function clearComboTargetCooldowns(comboId?: string): void { if (comboId === undefined) { targetCooldowns.clear(); + const persistenceChanged = persistedQuotaCooldownKeys.size > 0; + persistedQuotaCooldownKeys.clear(); + if (persistenceChanged) scheduleQuotaCooldownPersistence(); liveComboTargets.clear(); + liveProviders.clear(); lastReconciledGeneration = 0; return; } const prefix = `${comboId}\0`; + let persistenceChanged = false; for (const key of targetCooldowns.keys()) { - if (key.startsWith(prefix)) targetCooldowns.delete(key); + if (!key.startsWith(prefix)) continue; + targetCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) persistenceChanged = true; } + if (persistenceChanged) scheduleQuotaCooldownPersistence(); } export type ComboFailureDecision = "hop" | "stop"; -export type ComboFailureCooldownScope = "target" | "provider"; +export type ComboFailureCooldownScope = "none" | "target" | "provider"; function normalizedFailureCode(code?: string | null): string { return code?.trim().toLowerCase().replaceAll("-", "_") ?? ""; @@ -282,7 +498,34 @@ export function comboFailureCooldownScope( message: string, options?: { code?: string | null }, ): ComboFailureCooldownScope { - return isProviderScopedQuotaCap(status, message, options?.code) ? "provider" : "target"; + const code = normalizedFailureCode(options?.code); + const text = message.toLowerCase(); + // Request-shape limits must exclude only this attempt. Cooling the provider globally would + // make an unrelated shorter/simpler request skip a provider that is still healthy. + if ( + code === "free_rate_limited" + || text.includes("err_free_prompt_cap") + || [ + "input_admission_refused", + "context_length_exceeded", + "tool_catalog_too_large", + "cursor_root_envelope_limit", + "target_incompatible", + ].includes(code) + || status === 413 + ) return "none"; + if (isProviderScopedQuotaCap(status, message, code)) return "provider"; + if (status === 401 || status === 402 || status === 403) return "provider"; + if ([ + "invalid_api_key", + "insufficient_quota", + "subscription_required", + "payment_required", + "billing_error", + "insufficient_balance", + "provider_unavailable", + ].includes(code)) return "provider"; + return "target"; } function isModelLifecycleGone( @@ -341,26 +584,36 @@ export function comboFailureDecision( // an upstream already controls other hop signals (429, 5xx), and traversal is finite: policy // tries each candidate once via `tried`, and combo excludes each attempted target. So this is // structured-code-only, not provably local. - if (options?.code === "input_admission_refused" || error.code === "input_admission_refused") { - return "hop"; - } - if (isProviderScopedQuotaCap(status, message, options?.code || error.code)) { - return "hop"; - } - if (["origin_rejected", "context_length_exceeded", "invalid_request_error"].includes(error.code ?? "")) { - return "stop"; - } - if ([401, 403, 404, 408, 429].includes(status) || status >= 500) return "hop"; + const failureCode = normalizedFailureCode(options?.code || error.code); + const targetLocalCodes = new Set([ + "input_admission_refused", + "context_length_exceeded", + "tool_catalog_too_large", + "cursor_root_envelope_limit", + "kiro_profile_required", + "target_incompatible", + "model_not_found", + "model_unavailable", + "unsupported_model", + ]); + if (targetLocalCodes.has(failureCode)) return "hop"; + const lowerMessage = message.toLowerCase(); + if (/\b(?:model not found|model unavailable|unsupported model)\b/.test(lowerMessage)) return "hop"; + if (isProviderScopedQuotaCap(status, message, failureCode)) return "hop"; if ([ "permission_denied", "subscription_required", "invalid_api_key", "insufficient_quota", + "payment_required", + "billing_error", + "insufficient_balance", "rate_limit_exceeded", "server_is_overloaded", "upstream_server_error", - ].includes(error.code ?? "")) { - return "hop"; - } + "provider_unavailable", + ].includes(failureCode)) return "hop"; + if ([401, 402, 403, 404, 408, 410, 413, 425, 429].includes(status) || status >= 500) return "hop"; + if (["origin_rejected", "invalid_request_error"].includes(error.code ?? "")) return "stop"; return "stop"; } diff --git a/src/combos/index.ts b/src/combos/index.ts index 982f87c9e1..cfc8b29bfe 100644 --- a/src/combos/index.ts +++ b/src/combos/index.ts @@ -34,10 +34,13 @@ export { comboCooldownRetryAfterSeconds, COMBO_REQUEST_RATE_COOLDOWN_MS, coolComboTarget, + hydrateComboQuotaCooldownsFromDisk, isComboTargetInCooldown, isTransientRequestRateLimit, + knownComboFailureRecoveryMs, parseRetryAfterMs, remainingComboCooldownMs, + resetComboQuotaCooldownPersistenceForTests, comboFailureDecision, comboFailureCooldownScope, type ComboFailureDecision, diff --git a/src/combos/resolve.ts b/src/combos/resolve.ts index 56d7dd8fd1..c312d28aaa 100644 --- a/src/combos/resolve.ts +++ b/src/combos/resolve.ts @@ -1,7 +1,7 @@ import type { OcxComboTarget, OcxConfig } from "../types"; import { getCachedProviderQuota } from "../providers/quota-routing-cache"; import type { ProviderQuota } from "../providers/quota-types"; -import { coolComboTarget, isComboTargetInCooldown, type ComboFailureCooldownScope } from "./failover"; +import { coolComboProvider, coolComboTarget, isComboTargetInCooldown, type ComboFailureCooldownScope } from "./failover"; import { quotaResetRemainingMs } from "./reset-window"; import { getCombo, resolveComboId, targetKey } from "./types"; import type { NormalizedComboConfig } from "./types"; @@ -155,6 +155,7 @@ export function pickComboTarget( const eligible = (target: Required): boolean => targetProviderIsUsable(config, target) && !cachedProviderQuotaIsExhausted(getCachedProviderQuota(target.provider, now), now) + && !isComboTargetInCooldown(comboId, target, now) && !excluded.has(targetKey(target)) && (options.eligible?.(target) ?? true); @@ -283,12 +284,13 @@ export function advanceComboAfterFailure( } = {}, ): ComboPick | null { noteComboFailure(pick.comboId, pick.target, pick.writerGeneration); - const combo = getCombo(config, pick.comboId); - const cooldownTargets = options.cooldownScope === "provider" && combo - ? combo.targets.filter(target => target.provider === pick.target.provider) - : [pick.target]; - for (const target of cooldownTargets) { - coolComboTarget(pick.comboId, target, { + if (options.cooldownScope === "provider") { + coolComboProvider(pick.target.provider, { + ...options, + writerGeneration: pick.writerGeneration, + }); + } else if (options.cooldownScope !== "none") { + coolComboTarget(pick.comboId, pick.target, { ...options, writerGeneration: pick.writerGeneration, }); diff --git a/src/lib/errors.ts b/src/lib/errors.ts index 624917507c..7577dd6f77 100644 --- a/src/lib/errors.ts +++ b/src/lib/errors.ts @@ -365,6 +365,10 @@ export function inferHttpStatusFromAdapterMessage(message: string): number { lower.includes("temporarily") || lower.includes("server is busy") ) return 503; + // Malformed bytes produced by the upstream are a provider protocol failure, not a bad + // client request. Keep this ahead of the generic "malformed" -> 400 branch so an explicit + // combo can safely fail over when no output has been committed. + if (lower.includes("malformed upstream")) return 502; if ( lower.includes("invalid") || lower.includes("not found") || diff --git a/src/providers/key-cooldown-disk.ts b/src/providers/key-cooldown-disk.ts new file mode 100644 index 0000000000..f10ccec4b9 --- /dev/null +++ b/src/providers/key-cooldown-disk.ts @@ -0,0 +1,77 @@ +/** Persist only long-lived API-key quota cooldowns across proxy restarts. */ +import { existsSync, readFileSync } from "node:fs"; +import { join } from "node:path"; +import { atomicWriteFile, getConfigDir } from "../config"; + +const FILENAME = "provider-key-quota-cooldowns.json"; +const MAX_FUTURE_MS = 31 * 24 * 60 * 60_000; +const PERSIST_DEBOUNCE_MS = 250; +type DiskFile = { version: 1; rows: Record }; + +let persistTimer: ReturnType | null = null; +let pendingRows: (() => Iterable<[string, number]>) | null = null; +let pendingDirectory: string | null = null; + +export function keyQuotaCooldownStoreDirectory(): string { + return getConfigDir(); +} + +export function readPersistedKeyQuotaCooldowns(directory: string, now = Date.now()): Map { + const rows = new Map(); + try { + const path = join(directory, FILENAME); + if (!existsSync(path)) return rows; + const parsed = JSON.parse(readFileSync(path, "utf8")) as DiskFile; + if (!parsed || parsed.version !== 1 || !parsed.rows || typeof parsed.rows !== "object") return rows; + for (const [key, until] of Object.entries(parsed.rows)) { + if (typeof until !== "number" || !Number.isFinite(until)) continue; + if (until <= now || until - now > MAX_FUTURE_MS) continue; + rows.set(key, until); + } + } catch { + // Missing/corrupt best-effort state must never block routing. + } + return rows; +} +function persistNow(directory: string, rows: Iterable<[string, number]>, now = Date.now()): void { + try { + const out: Record = {}; + for (const [key, until] of rows) { + if (!Number.isFinite(until) || until <= now || until - now > MAX_FUTURE_MS) continue; + out[key] = until; + } + atomicWriteFile(join(directory, FILENAME), `${JSON.stringify({ version: 1, rows: out } satisfies DiskFile)}\n`); + } catch { + // Runtime failover remains authoritative if persistence fails. + } +} + +export function schedulePersistKeyQuotaCooldowns(directory: string, rows: () => Iterable<[string, number]>): void { + pendingDirectory = directory; + pendingRows = rows; + if (persistTimer) clearTimeout(persistTimer); + persistTimer = setTimeout(() => { + persistTimer = null; + const snapshot = pendingRows; + const targetDirectory = pendingDirectory; + pendingRows = null; + pendingDirectory = null; + if (snapshot && targetDirectory) persistNow(targetDirectory, snapshot()); + }, PERSIST_DEBOUNCE_MS); +} +export function flushKeyQuotaCooldownPersistForTests(now = Date.now()): void { + if (persistTimer) clearTimeout(persistTimer); + persistTimer = null; + const snapshot = pendingRows; + const targetDirectory = pendingDirectory; + pendingRows = null; + pendingDirectory = null; + if (snapshot && targetDirectory) persistNow(targetDirectory, snapshot(), now); +} + +export function cancelPendingKeyQuotaCooldownPersist(): void { + if (persistTimer) clearTimeout(persistTimer); + persistTimer = null; + pendingRows = null; + pendingDirectory = null; +} diff --git a/src/providers/key-failover.ts b/src/providers/key-failover.ts index 4e9f2e60a4..dda14df082 100644 --- a/src/providers/key-failover.ts +++ b/src/providers/key-failover.ts @@ -10,6 +10,13 @@ */ import { saveConfigPreservingClaudeCode } from "../config"; import type { OcxConfig, OcxProviderConfig, RateLimitRetryPolicy, TransientRetryPolicy } from "../types"; +import { deleteCachedProviderQuota, getCachedProviderQuota } from "./quota-routing-cache"; +import type { ProviderQuota } from "./quota-types"; +import { + keyQuotaCooldownStoreDirectory, + readPersistedKeyQuotaCooldowns, + schedulePersistKeyQuotaCooldowns, +} from "./key-cooldown-disk"; import { resolveProviderTransport, type OcxProviderTransport } from "./xai-transport"; import { sweepExpiredOnWrite } from "../lib/state-store-sweeper"; @@ -21,6 +28,23 @@ interface KeyCooldown { const DEFAULT_COOLDOWN_MS = 60_000; const MAX_COOLDOWN_MS = 10 * 60_000; // cap at 10 min for api-key rotation +const MAX_QUOTA_COOLDOWN_MS = 31 * 24 * 60 * 60_000; + +function exhaustedQuotaRecoveryMs(quota: ProviderQuota | null, now: number): number | undefined { + if (!quota) return undefined; + const resets: number[] = []; + const add = (percent: number | undefined, resetAt: number | undefined) => { + if (typeof percent !== "number" || !Number.isFinite(percent) || percent < 100) return; + if (typeof resetAt !== "number" || !Number.isFinite(resetAt) || resetAt <= now) return; + resets.push(resetAt); + }; + add(quota.fiveHourPercent, quota.fiveHourResetAt); + add(quota.weeklyPercent, quota.weeklyResetAt); + add(quota.monthlyPercent, quota.monthlyResetAt); + for (const window of quota.customWindows ?? []) add(window.percent, window.resetAt); + if (resets.length === 0) return undefined; + return Math.min(Math.max(...resets) - now, MAX_QUOTA_COOLDOWN_MS); +} /** * Default same-target 429 retry policy used when a provider opts in via a bare @@ -45,11 +69,43 @@ const DEFAULT_TRANSIENT_RETRY = { /** Map<`${providerName}\0${keyId}`, KeyCooldown> */ const keyCooldowns = new Map(); +const persistedQuotaCooldownKeys = new Set(); +let quotaCooldownPersistenceDirectory: string | undefined; + +function persistedQuotaCooldownRows(): Iterable<[string, number]> { + return [...persistedQuotaCooldownKeys].flatMap(key => { + const state = keyCooldowns.get(key); + return state ? [[key, state.cooldownUntil] as [string, number]] : []; + }); +} + +function scheduleQuotaCooldownPersistence(): void { + if (!quotaCooldownPersistenceDirectory) return; + schedulePersistKeyQuotaCooldowns(quotaCooldownPersistenceDirectory, () => persistedQuotaCooldownRows()); +} function cooldownKey(providerName: string, keyId: string): string { return `${providerName}\0${keyId}`; } +export function hydrateKeyQuotaCooldownsFromDisk(now = Date.now()): void { + const directory = keyQuotaCooldownStoreDirectory(); + if (quotaCooldownPersistenceDirectory === directory) return; + for (const key of persistedQuotaCooldownKeys) keyCooldowns.delete(key); + persistedQuotaCooldownKeys.clear(); + quotaCooldownPersistenceDirectory = directory; + for (const [key, cooldownUntil] of readPersistedKeyQuotaCooldowns(directory, now)) { + keyCooldowns.set(key, { cooldownUntil }); + persistedQuotaCooldownKeys.add(key); + } +} + +export function resetKeyQuotaCooldownPersistenceForTests(): void { + keyCooldowns.clear(); + persistedQuotaCooldownKeys.clear(); + quotaCooldownPersistenceDirectory = undefined; +} + /** * Parse an upstream `Retry-After` header: numeric seconds (including `0`) or an HTTP-date. * Returns a bounded delay in ms (1..MAX_COOLDOWN_MS), or undefined when the value is @@ -77,10 +133,12 @@ function parseRetryAfterMs(value: string | null | undefined, now = Date.now()): * window expires). Used to skip keys that the upstream just rate-limited during failover. */ function isKeyInCooldown(providerName: string, keyId: string, now = Date.now()): boolean { - const entry = keyCooldowns.get(cooldownKey(providerName, keyId)); + const key = cooldownKey(providerName, keyId); + const entry = keyCooldowns.get(key); if (!entry) return false; if (entry.cooldownUntil <= now) { - keyCooldowns.delete(cooldownKey(providerName, keyId)); + keyCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) scheduleQuotaCooldownPersistence(); return false; } return true; @@ -176,9 +234,10 @@ export function rateLimitRetryDelayMs( * not assign it to an active route wholesale; use `rotateProviderTransportOn429`, which * takes only the swapped key and keeps the routed provider intact. */ -export function rotateKeyOn429( +function rotateKeyAfterFailure( config: OcxConfig, providerName: string, + failureStatus: 401 | 429, retryAfterHeader: string | null | undefined, now = Date.now(), attemptedKey?: string, @@ -186,60 +245,84 @@ export function rotateKeyOn429( const provider = config.providers[providerName]; if (!provider) return null; if (provider.authMode === "oauth" || provider.authMode === "forward") return null; - const pool = provider.apiKeyPool; if (!pool || pool.length < 2) return null; - // Cool the key that ACTUALLY failed. Under concurrent 429s another request may already have - // rotated provider.apiKey — cooling the live key would punish an innocent replacement and can - // exhaust a 2-key pool from a single bad key. CAS semantics: callers pass the key they used. const failedKey = attemptedKey ?? provider.apiKey; + const failedKeyWasActive = failedKey === provider.apiKey; + const quotaRecoveryMs = failureStatus === 429 && failedKeyWasActive + ? exhaustedQuotaRecoveryMs(getCachedProviderQuota(providerName, now), now) + : undefined; const currentEntry = pool.find(e => e.key === failedKey); if (currentEntry) { - const cooldownMs = parseRetryAfterMs(retryAfterHeader, now) ?? DEFAULT_COOLDOWN_MS; - keyCooldowns.set(cooldownKey(providerName, currentEntry.id), { - cooldownUntil: now + cooldownMs, - }); + const headerCooldownMs = failureStatus === 429 ? parseRetryAfterMs(retryAfterHeader, now) : undefined; + const authCooldownMs = failureStatus === 401 ? MAX_COOLDOWN_MS : DEFAULT_COOLDOWN_MS; + const cooldownMs = Math.max(headerCooldownMs ?? 0, quotaRecoveryMs ?? 0, authCooldownMs); + const key = cooldownKey(providerName, currentEntry.id); + const cooldownUntil = now + cooldownMs; + keyCooldowns.set(key, { cooldownUntil }); + const durableLongQuota = quotaRecoveryMs !== undefined && cooldownMs > MAX_COOLDOWN_MS; + const persistenceChanged = durableLongQuota + ? (persistedQuotaCooldownKeys.add(key), true) + : persistedQuotaCooldownKeys.delete(key); + if (persistenceChanged) scheduleQuotaCooldownPersistence(); sweepExpiredOnWrite(now); } - // Lost the race: someone already rotated away from the failed key. If the live key is healthy, - // retry with it as-is instead of rotating a second time. + // CAS: another request may already have rotated away from the exact credential that failed. if (attemptedKey !== undefined && provider.apiKey !== attemptedKey) { const liveEntry = pool.find(e => e.key === provider.apiKey); - if (liveEntry && !isKeyInCooldown(providerName, liveEntry.id, now)) { - return { ...provider }; - } + if (liveEntry && !isKeyInCooldown(providerName, liveEntry.id, now)) return { ...provider }; } - // Pick the next key that is NOT in cooldown const currentIndex = currentEntry ? pool.indexOf(currentEntry) : -1; for (let i = 1; i < pool.length; i++) { const candidate = pool[(currentIndex + i) % pool.length]!; - if (!isKeyInCooldown(providerName, candidate.id, now)) { - // Swap active key - provider.apiKey = candidate.key; - saveConfigPreservingClaudeCode(config); - console.warn( - // Log ids only — labels are user-supplied free text and could carry secret material. - `[key-failover] ${providerName}: 429 on key ${currentEntry?.id ?? "?"}; rotating to key ${candidate.id}`, - ); - return { ...provider }; - } + if (isKeyInCooldown(providerName, candidate.id, now)) continue; + provider.apiKey = candidate.key; + deleteCachedProviderQuota(providerName); + saveConfigPreservingClaudeCode(config); + console.warn( + `[key-failover] ${providerName}: ${failureStatus} on key ${currentEntry?.id ?? "?"}; rotating to key ${candidate.id}`, + ); + return { ...provider }; } - // All keys in cooldown - console.warn(`[key-failover] ${providerName}: all ${pool.length} keys in cooldown; returning 429 to client`); + console.warn( + `[key-failover] ${providerName}: all ${pool.length} keys in cooldown after ${failureStatus}; no replacement available`, + ); return null; } +export function rotateKeyOn429( + config: OcxConfig, + providerName: string, + retryAfterHeader: string | null | undefined, + now = Date.now(), + attemptedKey?: string, +): OcxProviderConfig | null { + return rotateKeyAfterFailure(config, providerName, 429, retryAfterHeader, now, attemptedKey); +} + +export function rotateKeyOn401( + config: OcxConfig, + providerName: string, + now = Date.now(), + attemptedKey?: string, +): OcxProviderConfig | null { + return rotateKeyAfterFailure(config, providerName, 401, null, now, attemptedKey); +} + export function sweepExpiredApiKeyCooldowns(now = Date.now()): number { let removed = 0; + let persistenceChanged = false; for (const [key, cooldown] of keyCooldowns) { if (cooldown.cooldownUntil > now) continue; keyCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) persistenceChanged = true; removed += 1; } + if (persistenceChanged) scheduleQuotaCooldownPersistence(); return removed; } @@ -283,16 +366,38 @@ export function rotateProviderTransportOn429( : null; } +export function rotateProviderTransportOn401( + config: OcxConfig, + providerName: string, + routedProvider: OcxProviderTransport, + options: Omit = {}, +): OcxProviderTransport | null { + const rotated = rotateKeyOn401(config, providerName, options.now, options.attemptedKey); + return rotated + ? resolveProviderTransport( + providerName, + { ...routedProvider, apiKey: rotated.apiKey }, + options.promptCacheKey, + ) + : null; +} + /** Clear cooldown state for a provider (e.g. after manual key management). */ export function clearKeyCooldowns(providerName?: string): void { + let persistenceChanged = false; if (!providerName) { keyCooldowns.clear(); - return; - } - const prefix = `${providerName}\0`; - for (const key of keyCooldowns.keys()) { - if (key.startsWith(prefix)) keyCooldowns.delete(key); + persistenceChanged = persistedQuotaCooldownKeys.size > 0; + persistedQuotaCooldownKeys.clear(); + } else { + const prefix = `${providerName}\0`; + for (const key of keyCooldowns.keys()) { + if (!key.startsWith(prefix)) continue; + keyCooldowns.delete(key); + if (persistedQuotaCooldownKeys.delete(key)) persistenceChanged = true; + } } + if (persistenceChanged) scheduleQuotaCooldownPersistence(); } /** Visible-for-testing: get the cooldown-until timestamp for a key. */ diff --git a/src/providers/quota-routing-cache.ts b/src/providers/quota-routing-cache.ts index 065d7338ca..990e12ffea 100644 --- a/src/providers/quota-routing-cache.ts +++ b/src/providers/quota-routing-cache.ts @@ -6,6 +6,10 @@ export function clearCachedProviderQuotas(): void { quotaCache.clear(); } +export function deleteCachedProviderQuota(provider: string): void { + quotaCache.delete(provider); +} + export function replaceCachedProviderQuotas(reports: ProviderQuotaReport[]): void { quotaCache.clear(); for (const report of reports) { diff --git a/src/providers/quota.ts b/src/providers/quota.ts index ca7e9eb993..2220109d62 100644 --- a/src/providers/quota.ts +++ b/src/providers/quota.ts @@ -590,7 +590,10 @@ async function fetchDeepSeekQuota(provider: string, config: OcxProviderConfig): ? `API balance ($${balance.toFixed(2)} total, $${grantedBalance.toFixed(2)} granted)` : `API balance ($${balance.toFixed(2)})`; return report(provider, "deepseek:balance", { - customWindows: [{ label, percent: 0 }], + // `is_available=false` is DeepSeek's authoritative statement that the account cannot + // serve API calls, even if a small positive nominal balance is still reported. A positive + // usable balance has no meaningful utilization denominator; zero is also exhausted. + customWindows: [{ label, percent: body?.is_available === false || balance === 0 ? 100 : 0 }], updatedAt: Date.now(), }); } @@ -2330,7 +2333,9 @@ async function maybeFetchProviderQuota( if ((provider.authMode ?? "key") === "key" && name === "openrouter") { return fetchOpenRouterQuota(name, provider); } - if ((provider.authMode ?? "key") === "key" && name === "deepseek") { + if ((provider.authMode ?? "key") === "key" + && (registryEntryForProviderDestination(provider)?.id === "deepseek" + || name.toLowerCase() === "deepseek")) { return fetchDeepSeekQuota(name, provider); } if ((provider.authMode ?? "key") === "key" && name === "cline-pass") { diff --git a/src/routing/analytics.ts b/src/routing/analytics.ts index 6dc217e734..b41e49442b 100644 --- a/src/routing/analytics.ts +++ b/src/routing/analytics.ts @@ -116,6 +116,7 @@ interface Bucket extends AnalyticsBreakdownRow { const COOLDOWN_RECOVERY_KINDS = new Set([ "rate-limit-429", + "key-401", "key-429", "oauth-401", "anthropic-oauth-429", diff --git a/src/server/index.ts b/src/server/index.ts index 92495a8b58..8edeae933c 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -21,6 +21,8 @@ import { websocketsEnabled, } from "../config"; import { grokDefaultReasoningEffort } from "../grok/effort"; +import { hydrateComboQuotaCooldownsFromDisk } from "../combos"; +import { hydrateKeyQuotaCooldownsFromDisk } from "../providers/key-failover"; import { flushConfigDirHardening } from "../config/paths"; import { reconcileOAuthProviders } from "../oauth"; import { withCatalogWriteSerialization } from "../codex/catalog-write-serialization"; @@ -645,6 +647,8 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server { const contentType = response.headers.get("content-type")?.toLowerCase() ?? ""; if (!response.ok || !response.body || !contentType.includes("text/event-stream")) { @@ -161,7 +175,14 @@ export async function preflightComboStreamResponse( try { for (;;) { - const next = await reader.read(); + let next: Awaited>; + try { + next = await reader.read(); + } catch (error) { + if (abortSignal?.aborted) throw error; + await reader.cancel("retrying zero-output combo stream read failure").catch(() => undefined); + return { kind: "failed", response: streamReadFailureResponse(response) }; + } if (next.done) { inspector.finish(); } else { diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index ce6abf5c3d..39ab482246 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -2,6 +2,7 @@ import type { Server } from "bun"; import { randomUUID } from "node:crypto"; import { bridgeToResponsesSSE, buildResponseJSON, formatErrorResponse, type ResponsesTerminalStatus } from "../../bridge"; import { formatPassthroughUpstreamError } from "./passthrough-error"; +import { requestPacingOverloadResponse } from "./pacing-overload"; import { createResponsesFieldBackfillBlockRewrite, backfillResponsesFieldsJson, @@ -76,6 +77,7 @@ import { getCombo, isComboTargetInCooldown, comboCooldownRetryAfterSeconds, + knownComboFailureRecoveryMs, NoAvailableComboTargetsError, noteComboSuccess, parseRetryAfterMs, @@ -87,6 +89,7 @@ import { CYBER_POLICY_ERROR_CODE, CYBER_POLICY_FALLBACK_MESSAGE, adapterFailureFromMessage, + httpStatusFromTerminalError, isCyberPolicyCode, isCyberPolicyMessage, } from "../../lib/errors"; @@ -233,6 +236,7 @@ import { hasKeyPoolFailover, rateLimitRetryDelayMs, rateLimitRetryPolicyFor, + rotateProviderTransportOn401, rotateProviderTransportOn429, transientRetryPolicyFor, } from "../../providers/key-failover"; @@ -1439,7 +1443,7 @@ export function decodeRequestErrorResponse(err: unknown, label: string): Respons export function comboUnavailableResponse( message: string, - options?: { retryAfter?: string | null }, + options?: { retryAfter?: string | null; status?: number }, ): Response { const headers = new Headers({ "Content-Type": "application/json" }); const retryAfter = options?.retryAfter?.trim(); @@ -1450,14 +1454,90 @@ export function comboUnavailableResponse( JSON.stringify({ error: { message, type: "server_error", code: "combo_unavailable" }, }), - { status: 503, headers }, + { status: options?.status ?? 503, headers }, ); } -function comboUnavailable(comboId: string, now = Date.now()): Response { - return comboUnavailableResponse(`No available targets for combo: ${comboId}`, { - retryAfter: comboCooldownRetryAfterSeconds(comboId, now), - }); +function comboUnavailable(config: OcxConfig, comboId: string, now = Date.now()): Response { + const providers = getCombo(config, comboId)?.targets.map(target => target.provider) ?? []; + return comboExhaustedResponse( + comboId, + [{ provider: "", category: "no_eligible_target" }], + comboCooldownRetryAfterSeconds(comboId, now, providers), + ); +} + +interface ComboFailureSummary { + provider: string; + category: string; + recoveryAt?: number; +} + +function publicComboFailureCategory(status: number, code: string | undefined, message: string): string { + const normalized = code?.trim().toLowerCase().replaceAll("-", "_") ?? ""; + const text = message.toLowerCase(); + if ( + normalized === "free_rate_limited" + || normalized === "gousagelimiterror" + || normalized === "insufficient_quota" + || text.includes("monthly usage limit") + || text.includes("err_free_prompt_cap") + ) return "quota"; + if (["payment_required", "billing_error", "insufficient_balance", "subscription_required"].includes(normalized) || status === 402) { + return "billing"; + } + if (normalized === "invalid_api_key" || status === 401) return "authentication"; + if (normalized === "permission_denied" || status === 403) return "permission"; + if ([ + "input_admission_refused", + "context_length_exceeded", + "tool_catalog_too_large", + "cursor_root_envelope_limit", + "kiro_profile_required", + "target_incompatible", + "model_not_found", + "model_unavailable", + "unsupported_model", + "model_deprecated", + "model_end_of_life", + "model_eol", + "model_retired", + ].includes(normalized) || status === 404 || status === 410 || status === 413) return "target_unavailable"; + if (status === 408 || status === 425 || status === 429 || status >= 500) return "transient_upstream"; + return "provider_error"; +} + +function comboExhaustedResponse( + comboId: string, + failures: ComboFailureSummary[], + retryAfter?: string, +): Response { + const correlationId = `combo-${crypto.randomUUID()}`; + const counts = new Map(); + for (const failure of failures) counts.set(failure.category, (counts.get(failure.category) ?? 0) + 1); + const causes = [...counts].map(([category, count]) => ({ category, count })); + const recoveries = failures + .filter((failure): failure is ComboFailureSummary & { recoveryAt: number } => + failure.provider.length > 0 && failure.recoveryAt !== undefined) + .map(failure => ({ provider: failure.provider, available_after: new Date(failure.recoveryAt).toISOString() })); + const headers = new Headers({ "Content-Type": "application/json" }); + if (retryAfter) headers.set("Retry-After", retryAfter); + const lastProvider = [...failures].reverse().find(failure => failure.provider.length > 0)?.provider; + console.warn(`[combo] ${comboId}: exhausted correlation=${correlationId} last=${lastProvider ?? "none"} causes=${causes.map(c => `${c.category}:${c.count}`).join(",") || "none"}`); + return new Response(JSON.stringify({ + error: { + message: `All targets are currently unavailable for combo: ${comboId}`, + type: "server_error", + code: "combo_unavailable", + details: { + combo: comboId, + correlation_id: correlationId, + ...(lastProvider ? { last_provider: lastProvider } : {}), + causes, + ...(recoveries.length > 0 ? { recoveries } : {}), + }, + }, + }), { status: 503, headers }); } @@ -1574,6 +1654,52 @@ export function sanitizedRetryAfter(value: string | null, now: number): string | +function isKnownTargetIncompatibilityText(text: string): boolean { + return text.includes("Kiro supports only automatic tool choice or tool_choice:none") + || text.includes("Kiro does not support service tiers") + || text.includes("Kiro does not support Responses structured output") + || /Kiro .+ does not support reasoning effort /.test(text) + || text.includes("ollama-native does not support required or exact named tool_choice") + || /ollama-native does not support reasoning level /.test(text) + || text.includes("ollama-native does not support structured output on Ollama Cloud") + || text.includes("ollama-native does not support forwarded caller credentials") + || text.includes("ollama-native cannot send video content in ") + || text.includes("ollama-native cannot preserve images in ") + || text.includes("azure-openai does not support forward auth mode") + || text.includes("tool_choice requires function ") && text.includes("cannot be represented for this destination"); +} + + +async function comboNonStreamingTerminalFailure(response: Response): Promise { + if (!response.ok || response.headers.get("content-type")?.toLowerCase().includes("text/event-stream")) return null; + let payload: unknown; + try { + payload = await response.clone().json(); + } catch { + return null; + } + if (payload === null || typeof payload !== "object" || Array.isArray(payload)) return null; + const record = payload as Record; + if (record.status !== "failed") return null; + const candidate = record.error ?? record.last_error; + const upstreamError = candidate && typeof candidate === "object" && !Array.isArray(candidate) + ? candidate as { type?: unknown; code?: unknown; message?: unknown } + : undefined; + const error = { + ...(typeof upstreamError?.type === "string" ? { type: upstreamError.type } : {}), + ...(typeof upstreamError?.code === "string" ? { code: upstreamError.code } : {}), + ...(typeof upstreamError?.message === "string" ? { message: upstreamError.message } : {}), + }; + const status = httpStatusFromTerminalError(error); + const message = redactSecretString(error.message ?? "Provider returned a failed response before output").slice(0, 500); + return formatErrorResponse( + status, + error.type ?? "upstream_error", + message, + { ...(error.code ? { code: error.code } : {}) }, + ); +} + export async function consumeComboFailure( response: Response, signal?: AbortSignal, @@ -1600,7 +1726,14 @@ export async function consumeComboFailure( classificationText = fallback; } const cyberFailure = isCyberPolicyCode(upstreamCode) || isCyberPolicyMessage(classificationText); - const normalizedUpstreamCode = cyberFailure ? CYBER_POLICY_ERROR_CODE : upstreamCode; + const targetIncompatible = !cyberFailure + && (upstreamCode === undefined || upstreamCode === "invalid_request_error") + && isKnownTargetIncompatibilityText(classificationText); + const normalizedUpstreamCode = cyberFailure + ? CYBER_POLICY_ERROR_CODE + : targetIncompatible + ? "target_incompatible" + : upstreamCode; const message = cyberFailure ? upstreamMessage ?? (isCyberPolicyCode(upstreamCode) ? CYBER_POLICY_FALLBACK_MESSAGE : classificationText) @@ -2314,7 +2447,7 @@ export async function handleComboResponses( config, { parentThreadId: inboundClientThreadId }, ); - return comboUnavailable(comboId); + return comboUnavailable(config, comboId); } let recovered = false; try { @@ -2351,13 +2484,14 @@ export async function handleComboResponses( } if (!pick) { - return comboUnavailable(comboId); + return comboUnavailable(config, comboId); } // One immutable combo selection trace, before any child dispatch; child // adoption below must never replace it with a concrete child route trace. logCtx.routeDecision = comboRouteDecisionTrace(config, comboId, pick, requestedModel); let lastFailure: Response | null = null; + const failureSummaries: ComboFailureSummary[] = []; while (pick) { if (options.abortSignal?.aborted) return clientCancelledResponse(); const childLog: RequestLogContext = { @@ -2440,7 +2574,12 @@ export async function handleComboResponses( retainCancelledAttempt(); return clientCancelledResponse(); } - throw error; + const pacingOverload = requestPacingOverloadResponse(error); + if (!pacingOverload) throw error; + // A local queue admission failure happened before this combo child reached the + // provider and before any output was committed. Treat it exactly like a retryable + // pre-stream 429 so an explicitly declared combo can try its next target. + response = pacingOverload; } if (options.abortSignal?.aborted) { @@ -2449,12 +2588,17 @@ export async function handleComboResponses( return clientCancelledResponse(); } + if (response.ok && (rawBody as { stream?: unknown } | null)?.stream !== true) { + const terminalFailure = await comboNonStreamingTerminalFailure(response); + if (terminalFailure) response = terminalFailure; + } + if (response.ok && !runTurnAdapterSseResponses.has(response)) { const nativePassthrough = isNativePassthroughSseResponse(response); const eagerRelay = isEagerRelaySseResponse(response); let preflight; try { - preflight = await preflightComboStreamResponse(response, childLog); + preflight = await preflightComboStreamResponse(response, childLog, options.abortSignal); } catch (error) { callbackGate.discard(); if (options.abortSignal?.aborted) { @@ -2536,7 +2680,10 @@ export async function handleComboResponses( (logCtx.attempts ??= []).push(attempt); attemptRetained = true; lastFailure = failure.response; - if (storedPool401ReplayDispatched) { + if (storedPool401ReplayDispatched && attempt.firstOutputMs !== undefined) { + // Once a replay has emitted output, replaying the logical turn on another provider + // would duplicate visible content. A stored replay that failed before first output + // remains safe to classify and fail over inside an explicitly declared combo. adoptFailedChildLog(childLog); return lastFailure; } @@ -2555,9 +2702,22 @@ export async function handleComboResponses( console.warn( `[combo] ${comboId}: ${targetKey(pick.target)} failed with ${failure.response.status} after ${Date.now() - started}ms`, ); + const failureNow = Date.now(); + const knownRecoveryMs = knownComboFailureRecoveryMs({ + retryAfter: failure.retryAfter, + status: failure.response.status, + code: failure.upstreamCode, + message: failure.classificationText, + now: failureNow, + }); + failureSummaries.push({ + provider: pick.target.provider, + category: publicComboFailureCategory(failure.response.status, failure.upstreamCode, failure.classificationText), + ...(knownRecoveryMs !== undefined ? { recoveryAt: failureNow + knownRecoveryMs } : {}), + }); const nextPick = advanceComboAfterFailure(config, pick, { retryAfter: failure.retryAfter, - now: Date.now(), + now: failureNow, cooldownScope: comboFailureCooldownScope(failure.response.status, failure.classificationText, { code: failure.upstreamCode, }), @@ -2569,13 +2729,26 @@ export async function handleComboResponses( if (!nextPick) adoptFailedChildLog(childLog); pick = nextPick; } - if ( - lastFailure?.status === 413 - && (rawBody as { stream?: unknown } | null)?.stream === true - ) { - return streamingContextOverflowResponse(requestedModel, options.translatorBudget); + if (lastFailure?.status === 413) { + if ((rawBody as { stream?: unknown } | null)?.stream === true) { + return streamingContextOverflowResponse(requestedModel, options.translatorBudget); + } + return formatErrorResponse( + 413, + "request_too_large", + `No target in combo "${comboId}" can accept this request`, + { code: "combo_input_too_large" }, + ); } - return lastFailure!; + // Every failure that reached this point was classified as replay-safe and every declared + // target was attempted or excluded. Keep provider payloads in internal attempts/logs, but do + // not leak the last provider's billing/rate-limit/error body into the client conversation. + return comboExhaustedResponse( + comboId, + failureSummaries, + comboCooldownRetryAfterSeconds(comboId, Date.now(), combo.targets.map(target => target.provider)) + ?? lastFailure?.headers.get("retry-after") ?? undefined, + ); } @@ -2935,7 +3108,7 @@ async function handleResponsesInner( logCtx.routeDecision = route.routeDecision; } catch (err) { if (err instanceof NoAvailableComboTargetsError) { - return comboUnavailable(err.comboId); + return comboUnavailable(config, err.comboId); } if (err instanceof NoEligiblePolicyCandidateError) { // Persist the evaluation trace (per-candidate exclusions + the @@ -3045,7 +3218,7 @@ async function handleResponsesInner( logCtx.routeDecision = route.routeDecision; } catch (err) { if (err instanceof NoAvailableComboTargetsError) { - return comboUnavailable(err.comboId); + return comboUnavailable(config, err.comboId); } if (err instanceof NoEligiblePolicyCandidateError) { logCtx.routeDecision = err.trace; @@ -3169,7 +3342,7 @@ async function handleResponsesInner( logCtx.routeDecision = route.routeDecision; } catch (err) { if (err instanceof NoAvailableComboTargetsError) { - return comboUnavailable(err.comboId); + return comboUnavailable(config, err.comboId); } if (err instanceof NoEligiblePolicyCandidateError) { logCtx.routeDecision = err.trace; @@ -3743,9 +3916,31 @@ async function handleResponsesInner( // unstructured 500 — and no request log — depending only on whether a rotation ran first. // Same shape for a tool_choice this proxy cannot honor: the destination rejects a schema the // catalog had to drop, so the selector naming it is a client input error, not a 500. - if (error instanceof NamespaceToolCollisionError || error instanceof XaiToolSchemaCompatibilityError) { + if (error instanceof NamespaceToolCollisionError) { return formatErrorResponse(400, "invalid_request_error", redactSecretString(error.message)); } + const targetIncompatible = error instanceof XaiToolSchemaCompatibilityError + || (error instanceof Error && ( + error.message === "Kiro supports only automatic tool choice or tool_choice:none" + || error.message === "Kiro does not support service tiers" + || error.message === "Kiro does not support Responses structured output" + || /^Kiro .+ does not support reasoning effort /.test(error.message) + || error.message === "ollama-native does not support required or exact named tool_choice" + || /^ollama-native does not support reasoning level /.test(error.message) + || error.message === "ollama-native does not support structured output on Ollama Cloud" + || error.message === "ollama-native does not support forwarded caller credentials" + || /^ollama-native cannot send video content in /.test(error.message) + || /^ollama-native cannot preserve images in /.test(error.message) + || error.message === "azure-openai does not support forward auth mode" + )); + if (targetIncompatible) { + return formatErrorResponse( + 400, + "invalid_request_error", + redactSecretString(error instanceof Error ? error.message : String(error)), + { code: "target_incompatible" }, + ); + } throw error; } if (!isCanonicalOpenAiForwardProvider(route.provider)) { @@ -6008,6 +6203,33 @@ async function handleResponsesInner( continue recovery; } + // Static API-key pools can recover a credential-scoped 401 without abandoning the + // provider. OAuth providers use the refresh/account paths above and never enter here. + while (upstreamResponse.status === 401 && hasKeyPoolFailover(route.provider)) { + const rotated = rotateProviderTransportOn401(config, route.providerName, route.provider, { + now: Date.now(), + attemptedKey: route.provider.apiKey, + promptCacheKey: parsed.options.promptCacheKey, + }); + if (!rotated) break; + try { void upstreamResponse.body?.cancel().catch(() => {}); } catch { /* already consumed/closed */ } + route.provider = rotated; + invalidateSameTargetRequest(); + activeAdapter = resolveAdapter( + resolveWireProtocolOverride(route.providerName, route.modelId, route.provider, inboundWire), + config.cacheRetention, + ); + bindRouteReasoningReplayScope({ + parsed, + providerName: route.providerName, + provider: route.provider, + adapterName: activeAdapter.name, + }); + const result = await rebuildAndRefetch("key-401"); + if ("failed" in result) return result.failed; + upstreamResponse = result; + } + // Same-target 429 wait-and-retry (opt-in `retryOn429`, issue #487). Codex never retries // 429 itself (it retries 5xx only), and single-key pools cannot use the failover below, // so wait (Retry-After or the fixed interval) and replay the IDENTICAL request on the diff --git a/src/server/responses/pacing-overload.ts b/src/server/responses/pacing-overload.ts index 6ed30b85a7..a7858e0cb8 100644 --- a/src/server/responses/pacing-overload.ts +++ b/src/server/responses/pacing-overload.ts @@ -1,13 +1,23 @@ import { formatErrorResponse } from "../../bridge"; -import { RequestPacingQueueOverloadError } from "../../providers/request-pacing"; +import { RequestPacingProviderRemovedError, RequestPacingQueueOverloadError } from "../../providers/request-pacing"; /** Convert local request-pacing admission failures into a retryable HTTP response. */ export function requestPacingOverloadResponse(error: unknown): Response | undefined { - if (!(error instanceof RequestPacingQueueOverloadError)) return undefined; - return formatErrorResponse( - 429, - "rate_limit_error", - error.message, - { retryAfter: String(error.retryAfterSeconds) }, - ); + if (error instanceof RequestPacingQueueOverloadError) { + return formatErrorResponse( + 429, + "rate_limit_error", + error.message, + { retryAfter: String(error.retryAfterSeconds) }, + ); + } + if (error instanceof RequestPacingProviderRemovedError) { + return formatErrorResponse( + 503, + "server_error", + "request pacing provider became unavailable during dispatch", + { code: "provider_unavailable" }, + ); + } + return undefined; } diff --git a/src/server/responses/policy-fallback.ts b/src/server/responses/policy-fallback.ts index a4f06d0fa6..66d31a7ba0 100644 --- a/src/server/responses/policy-fallback.ts +++ b/src/server/responses/policy-fallback.ts @@ -1,4 +1,5 @@ import { comboFailureDecision } from "../../combos/failover"; +import { formatErrorResponse } from "../../bridge"; import { readBoundedResponseBody } from "../../lib/bounded-body"; import { readJsonRequestBody } from "../request-decompress"; import { finishRequestAttempt, type RequestLogContext } from "../request-log"; @@ -144,8 +145,8 @@ export async function handleResponsesWithPolicyFallback( response = await runCore(req, config, logCtx, coreOptions); } catch (error) { const overload = requestPacingOverloadResponse(error); - if (overload) return overload; - throw error; + if (!overload) throw error; + response = overload; } const initialTrace = logCtx.routeDecision; const initialRequestedModel = logCtx.requestedModel; @@ -155,10 +156,15 @@ export async function handleResponsesWithPolicyFallback( candidateKey({ provider: initialTrace.selected.provider, model: initialTrace.selected.model }), ]); - while (!storedPool401ReplayDispatched && await shouldHopPolicyCandidate(response, req.signal)) { + while (!(storedPool401ReplayDispatched && logCtx.activeAttempt?.firstOutputMs !== undefined) + && await shouldHopPolicyCandidate(response, req.signal)) { if (req.signal.aborted) return response; const next = rankPolicyFallbackCandidates(initialTrace, tried)[0]; - if (!next) return response; + if (!next) { + return response.status === 413 + ? formatErrorResponse(413, "request_too_large", "No eligible policy candidate can accept this request", { code: "policy_input_too_large" }) + : formatErrorResponse(response.status, "server_error", "All eligible policy candidates are temporarily unavailable", { code: "policy_unavailable", retryAfter: response.headers.get("retry-after") ?? undefined }); + } tried.add(candidateKey(next)); finishFailedPolicyAttempt(logCtx, response.status); @@ -168,8 +174,8 @@ export async function handleResponsesWithPolicyFallback( response = await runCore(retryRequest, config, logCtx, coreOptions); } catch (error) { const overload = requestPacingOverloadResponse(error); - if (overload) return overload; - throw error; + if (!overload) throw error; + response = overload; } } finally { logCtx.requestedModel = initialRequestedModel; diff --git a/src/usage/log.ts b/src/usage/log.ts index 6a74ae7f96..acbe8a00db 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -46,6 +46,7 @@ export type AttemptRecoveryKind = | "transient-5xx" | "connection-reset" | "oauth-401" + | "key-401" | "key-429" | "rate-limit-429" | "anthropic-oauth-429" @@ -257,6 +258,7 @@ const ATTEMPT_RECOVERY_KINDS = new Set([ "transient-5xx", "connection-reset", "oauth-401", + "key-401", "key-429", "rate-limit-429", "anthropic-oauth-429", diff --git a/tests/combo-stream-preflight.test.ts b/tests/combo-stream-preflight.test.ts index b4422727d2..2db7c22536 100644 --- a/tests/combo-stream-preflight.test.ts +++ b/tests/combo-stream-preflight.test.ts @@ -256,4 +256,19 @@ describe("combo stream preflight", () => { }); expect(JSON.stringify(body)).not.toContain("provider_trace_id"); }); + test("converts a zero-output stream read reset into a retryable HTTP failure", async () => { + const response = new Response(new ReadableStream({ + pull() { + throw Object.assign(new Error("provider-secret socket reset"), { code: "ECONNRESET" }); + }, + }), { headers: { "content-type": "text/event-stream" } }); + + const result = await preflightComboStreamResponse(response, { model: "m1", provider: "a" }); + expect(result.kind).toBe("failed"); + expect(result.response.status).toBe(502); + const body = await result.response.json(); + expect(body.error).toMatchObject({ type: "upstream_error", code: "upstream_reset" }); + expect(JSON.stringify(body)).not.toContain("provider-secret"); + }); + }); diff --git a/tests/combos.test.ts b/tests/combos.test.ts index 511bdae1a2..4745df05b7 100644 --- a/tests/combos.test.ts +++ b/tests/combos.test.ts @@ -22,6 +22,7 @@ import { coolComboTarget, earliestQuotaResetAt, getCombo, + hydrateComboQuotaCooldownsFromDisk, isComboTargetInCooldown, isValidComboId, listComboIds, @@ -35,6 +36,7 @@ import { pickComboTarget, preservesPhysicalComboProvider, resetComboEffortWarningStateForTests, + resetComboQuotaCooldownPersistenceForTests, resolveComboId, targetKey, tryPickComboModel, @@ -43,7 +45,9 @@ import { import { comboFailureCooldownScope, comboFailureDecision, + coolComboProvider, isTransientRequestRateLimit, + reconcileComboTargetCooldowns, } from "../src/combos/failover"; import { comboUnavailableResponse } from "../src/server/responses/core"; import { getConfigPath, readConfigDiagnostics, saveConfig } from "../src/config"; @@ -62,6 +66,7 @@ import { } from "../src/providers/quota-routing-cache"; import { catalogConvergenceFactory } from "./helpers/catalog-convergence"; import { removeTreeWithRetry } from "./helpers/remove-tree"; +import { cancelPendingComboQuotaCooldownPersist, flushComboQuotaCooldownPersistForTests } from "../src/combos/cooldown-disk"; const VALID_COMBO = { targets: [{ provider: "a", model: "m1" }] }; @@ -443,6 +448,109 @@ describe("combo target cooldowns", () => { expect(isComboTargetInCooldown("free", target, 1_000 + 60_000)).toBe(false); }); + test("structured quota-window codes honor explicit long recovery clocks", () => { + const now = Date.parse("2026-09-02T00:00:00.000Z"); + coolComboTarget("free", target, { + now, status: 429, code: "1308", message: "Usage limit reached for 5 hour", retryAfter: "86400", + }); + expect(isComboTargetInCooldown("free", target, now + 23 * 60 * 60_000)).toBe(true); + expect(isComboTargetInCooldown("free", target, now + 24 * 60 * 60_000)).toBe(false); + + clearComboTargetCooldowns(); + coolComboProvider("a", { + now, status: 429, code: "insufficient_quota", message: "quota exhausted", retryAfter: "1209600", + }); + expect(isComboTargetInCooldown("another", target, now + 13 * 24 * 60 * 60_000)).toBe(true); + expect(isComboTargetInCooldown("another", target, now + 14 * 24 * 60 * 60_000)).toBe(false); + }); + + test("keeps an explicit monthly-reset quota failure cooling until its known reset", () => { + const now = Date.parse("2026-09-02T00:00:00.000Z"); + coolComboTarget("free", target, { + now, + status: 429, + code: "GoUsageLimitError", + message: "Monthly usage limit reached. Resets in 14 days.", + }); + expect(isComboTargetInCooldown("free", target, now + 13 * 24 * 60 * 60_000)).toBe(true); + expect(isComboTargetInCooldown("free", target, now + 14 * 24 * 60 * 60_000)).toBe(false); + + clearComboTargetCooldowns("free"); + coolComboTarget("free", target, { now, status: 503, message: "temporary overload" }); + expect(isComboTargetInCooldown("free", target, now + 60_000)).toBe(false); + }); + + test("long provider quota cooldown survives a proxy restart without persisting transient failures", () => { + const previousHome = process.env.OPENCODEX_HOME; + const home = mkdtempSync(join(tmpdir(), "ocx-combo-quota-cooldown-")); + const now = Date.parse("2026-09-02T00:00:00.000Z"); + process.env.OPENCODEX_HOME = home; + try { + resetComboQuotaCooldownPersistenceForTests(); + hydrateComboQuotaCooldownsFromDisk(now); + coolComboTarget("free", target, { + now, + status: 429, + code: "GoUsageLimitError", + message: "Monthly usage limit reached. Resets in 14 days.", + }); + flushComboQuotaCooldownPersistForTests(now); + + resetComboQuotaCooldownPersistenceForTests(); + hydrateComboQuotaCooldownsFromDisk(now + 1_000); + expect(isComboTargetInCooldown("free", target, now + 13 * 24 * 60 * 60_000)).toBe(true); + expect(isComboTargetInCooldown("free", target, now + 14 * 24 * 60 * 60_000)).toBe(false); + } finally { + cancelPendingComboQuotaCooldownPersist(); + resetComboQuotaCooldownPersistenceForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + removeTreeWithRetry(home); + } + }); + + test("provider-wide long quota cooldown survives restart and applies to another combo", () => { + const previousHome = process.env.OPENCODEX_HOME; + const home = mkdtempSync(join(tmpdir(), "ocx-provider-quota-cooldown-")); + const now = Date.parse("2026-09-02T00:00:00.000Z"); + process.env.OPENCODEX_HOME = home; + try { + resetComboQuotaCooldownPersistenceForTests(); + hydrateComboQuotaCooldownsFromDisk(now); + coolComboProvider("a", { + now, status: 429, code: "GoUsageLimitError", + message: "Monthly usage limit reached. Resets in 14 days.", + }); + flushComboQuotaCooldownPersistForTests(now); + resetComboQuotaCooldownPersistenceForTests(); + hydrateComboQuotaCooldownsFromDisk(now + 1_000); + expect(isComboTargetInCooldown("another", target, now + 13 * 24 * 60 * 60_000)).toBe(true); + expect(comboCooldownRetryAfterSeconds("another", now + 1_000, ["a"])).toBe("1209599"); + expect(comboCooldownRetryAfterSeconds("another", now + 1_000, ["b"])).toBeUndefined(); + } finally { + cancelPendingComboQuotaCooldownPersist(); + resetComboQuotaCooldownPersistenceForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + removeTreeWithRetry(home); + } + }); + + test("config reconciliation removes orphaned provider and target cooldowns", () => { + coolComboProvider("a", { now: 1_000, cooldownMs: 60_000 }); + coolComboTarget("free", { provider: "b", model: "m2" }, { now: 1_000, cooldownMs: 60_000 }); + expect(isComboTargetInCooldown("free", target, 1_001)).toBe(true); + expect(isComboTargetInCooldown("free", { provider: "b", model: "m2" }, 1_001)).toBe(true); + const removed = reconcileComboTargetCooldowns({ + generation: 1, providerNames: new Set(["b"]), comboIds: new Set(["free"]), + comboTargets: new Set(["free::b/m3"]), codexAccountIds: new Set(), + oauthAccountKeys: new Set(), configRoots: new Set(), + }); + expect(removed).toBe(2); + expect(isComboTargetInCooldown("free", target, 1_001)).toBe(false); + expect(isComboTargetInCooldown("free", { provider: "b", model: "m2" }, 1_001)).toBe(false); + }); + test("honors explicit Retry-After over the request-rate default", () => { coolComboTarget("free", target, { now: 1_000, @@ -478,11 +586,12 @@ describe("combo failure policy and advancement", () => { for (const status of [401, 403, 404, 408, 429, 500, 503]) { expect(comboFailureDecision(status, "provider failure")).toBe("hop"); } - expect(comboFailureDecision(400, "context_length_exceeded")).toBe("stop"); + expect(comboFailureDecision(400, "context_length_exceeded")).toBe("hop"); expect(comboFailureDecision(403, '{"code":"origin_rejected"}')).toBe("stop"); - expect(comboFailureDecision(413, "request too large")).toBe("stop"); + expect(comboFailureDecision(413, "request too large")).toBe("hop"); expect(comboFailureDecision(409, "conflict")).toBe("stop"); - expect(comboFailureDecision(410, "resource is gone")).toBe("stop"); + expect(comboFailureDecision(425, "too early")).toBe("hop"); + expect(comboFailureDecision(410, "resource is gone")).toBe("hop"); expect(comboFailureDecision(410, "The model has reached its end of life and is no longer available.")).toBe("hop"); expect(comboFailureDecision(410, "The model is scheduled for retirement.")).toBe("hop"); expect(comboFailureDecision(410, "gone", { code: "model_retired" })).toBe("hop"); @@ -497,10 +606,44 @@ describe("combo failure policy and advancement", () => { // verdict by echoing the token, so that shape must NOT hop. expect(comboFailureDecision(413, 'refused', { code: 'input_admission_refused' })).toBe('hop'); expect(comboFailureDecision(400, 'upstream mentions input_admission_refused in prose')).toBe('stop'); - // An UPSTREAM context verdict still stops: retrying that elsewhere is guesswork, and a - // generic 413 with no structured code keeps its existing conservative handling. - expect(comboFailureDecision(400, "context_length_exceeded")).toBe("stop"); - expect(comboFailureDecision(413, "request too large")).toBe("stop"); + // An upstream context verdict is target-local in an explicit combo: another declared + // model can have a larger context window. A generic 413 without context evidence remains + // terminal because replaying an arbitrary oversized request would be guesswork. + expect(comboFailureDecision(400, "context_length_exceeded")).toBe("hop"); + expect(comboFailureDecision(413, "request too large")).toBe("hop"); + }); + + test("HTTP status matrix keeps retryable target failures distinct from terminal request failures", () => { + const matrix = [ + [400, "ordinary invalid request", "stop", "target"], + [401, "authentication failed", "hop", "provider"], + [402, "payment required", "hop", "provider"], + [403, "permission denied", "hop", "provider"], + [404, "model route not found", "hop", "target"], + [408, "upstream request timeout", "hop", "target"], + [409, "request conflict", "stop", "target"], + [410, "resource is gone", "hop", "target"], + [413, "request too large", "hop", "none"], + [422, "ordinary unprocessable request", "stop", "target"], + [425, "too early", "hop", "target"], + [429, "rate limited", "hop", "target"], + [500, "internal server error", "hop", "target"], + [502, "bad gateway", "hop", "target"], + [503, "temporarily unavailable", "hop", "target"], + [504, "gateway timeout", "hop", "target"], + ] as const; + for (const [status, message, decision, scope] of matrix) { + expect(comboFailureDecision(status, message)).toBe(decision); + expect(comboFailureCooldownScope(status, message)).toBe(scope); + } + + // Structured evidence overrides a generic status only when the failure is scoped to a + // different target/provider. Request-shape limits are request-local and must not poison + // later short requests; provider auth/billing failures exclude sibling models too. + expect(comboFailureDecision(400, "maximum context exceeded", { code: "context_length_exceeded" })).toBe("hop"); + expect(comboFailureCooldownScope(400, "maximum context exceeded", { code: "context_length_exceeded" })).toBe("none"); + expect(comboFailureDecision(422, "bad credential", { code: "invalid_api_key" })).toBe("hop"); + expect(comboFailureCooldownScope(422, "bad credential", { code: "invalid_api_key" })).toBe("provider"); }); test("provider-scoped free-tier and monthly quota failures hop without weakening generic 400 handling", () => { @@ -510,7 +653,7 @@ describe("combo failure policy and advancement", () => { message: "This prompt is longer than the free tier allows for a single request.", }}); expect(comboFailureDecision(400, orca, { code: "free_rate_limited" })).toBe("hop"); - expect(comboFailureCooldownScope(400, orca, { code: "free_rate_limited" })).toBe("provider"); + expect(comboFailureCooldownScope(400, orca, { code: "free_rate_limited" })).toBe("none"); expect(comboFailureDecision(400, "ordinary invalid request", { code: "invalid_request_error" })).toBe("stop"); expect(comboFailureCooldownScope(429, "Monthly usage limit reached. Resets in 14 days.", { code: "GoUsageLimitError", @@ -528,6 +671,28 @@ describe("combo failure policy and advancement", () => { })).toBe(true); }); + test("structured provider and target-local failures cannot be swallowed by generic status classification", () => { + const retryable = [ + [402, "insufficient_quota", "payment required"], + [402, "subscription_required", "billing required"], + [422, "invalid_api_key", "bad credential"], + [400, "model_not_found", "model not found"], + [400, "unsupported_model", "unsupported model"], + [400, "context_length_exceeded", "maximum context exceeded"], + [413, "tool_catalog_too_large", "tool catalog too large"], + [400, "cursor_root_envelope_limit", "Cursor request envelope is too large"], + [400, "kiro_profile_required", "Kiro account lacks the required profile"], + [400, "target_incompatible", "this target cannot represent the requested tool choice"], + ] as const; + for (const [status, code, message] of retryable) { + expect(comboFailureDecision(status, message, { code })).toBe("hop"); + } + expect(comboFailureCooldownScope(402, "payment required", { code: "insufficient_quota" })).toBe("provider"); + expect(comboFailureCooldownScope(422, "bad credential", { code: "invalid_api_key" })).toBe("provider"); + expect(comboFailureCooldownScope(400, "model not found", { code: "model_not_found" })).toBe("target"); + expect(comboFailureDecision(422, "ordinary unprocessable request", { code: "invalid_request_error" })).toBe("stop"); + }); + test("failover skips providers with fresh exhausted quota evidence before dispatch", () => { const now = 50_000; const config = baseConfig(); @@ -540,6 +705,30 @@ describe("combo failure policy and advancement", () => { expect(pick?.target.provider).toBe("b"); }); + test("stale exhausted quota fails open instead of blacklisting the provider", () => { + const now = 50_000_000; + const config = baseConfig(); + setCachedProviderQuotaForTests("a", { + monthlyPercent: 100, monthlyResetAt: now + 14 * 24 * 60 * 60_000, + updatedAt: now - 30 * 60_000 - 1, + }); + expect(pickComboTarget(config, "free", { now })?.target.provider).toBe("a"); + }); + + test("balance-only quota skips zero capacity but keeps positive capacity eligible", () => { + const now = 50_000; + const config = baseConfig(); + setCachedProviderQuotaForTests("a", { + customWindows: [{ label: "API balance ($0.00)", percent: 100 }], updatedAt: now, + }); + expect(pickComboTarget(config, "free", { now })?.target.provider).toBe("b"); + clearCachedProviderQuotas(); + setCachedProviderQuotaForTests("a", { + customWindows: [{ label: "API balance ($47.26)", percent: 0 }], updatedAt: now, + }); + expect(pickComboTarget(config, "free", { now })?.target.provider).toBe("a"); + }); + test("elapsed quota reset does not permanently blacklist a provider", () => { const now = 50_000; const config = baseConfig(); @@ -562,6 +751,24 @@ describe("combo failure policy and advancement", () => { expect(pickComboTarget(config, "free", { now })?.target.provider).toBe("b"); }); + test("request-local failures exclude only the attempted target for this turn", () => { + const config = baseConfig({ + combos: { + free: { + targets: [ + { provider: "a", model: "m1" }, + { provider: "b", model: "m2" }, + ], + }, + }, + }); + const first = pickComboTarget(config, "free", { now: 1_000 })!; + const next = advanceComboAfterFailure(config, first, { now: 1_000, cooldownScope: "none" })!; + expect(next.target.provider).toBe("b"); + expect(isComboTargetInCooldown("free", first.target, 1_001)).toBe(false); + expect(pickComboTarget(config, "free", { now: 1_001 })?.target.provider).toBe("a"); + }); + test("provider-scoped cooldown skips sibling models but leaves other providers eligible", () => { const config = baseConfig({ combos: { @@ -582,6 +789,40 @@ describe("combo failure policy and advancement", () => { expect(isComboTargetInCooldown("free", { provider: "b", model: "m2" }, 1_001)).toBe(false); }); + + test("provider-scoped quota cooldown is shared across combos", () => { + const config = baseConfig({ + combos: { + first: { + strategy: "failover", + targets: [ + { provider: "a", model: "m1" }, + { provider: "b", model: "m2" }, + ], + }, + second: { + strategy: "failover", + targets: [ + { provider: "a", model: "m1" }, + { provider: "b", model: "m2" }, + ], + }, + }, + }); + const now = Date.parse("2026-09-02T00:00:00.000Z"); + const first = pickComboTarget(config, "first", { now })!; + const next = advanceComboAfterFailure(config, first, { + now, + cooldownScope: "provider", + status: 429, + code: "GoUsageLimitError", + message: "Monthly usage limit reached. Resets in 14 days.", + })!; + expect(next.target.provider).toBe("b"); + expect(isComboTargetInCooldown("second", { provider: "a", model: "m1" }, now + 1)).toBe(true); + expect(pickComboTarget(config, "second", { now: now + 1 })?.target.provider).toBe("b"); + }); + test("failure clears the active sticky target without adding a success", () => { const config = rrConfig(2, [1, 1]); const combo = getCombo(config, "free")!; diff --git a/tests/error-fidelity.test.ts b/tests/error-fidelity.test.ts index 5bcd58a042..ad7bc39d73 100644 --- a/tests/error-fidelity.test.ts +++ b/tests/error-fidelity.test.ts @@ -28,6 +28,18 @@ async function collectSse(stream: ReadableStream): Promise<{ event?: }); } +describe("explicit internal error codes", () => { + test("formatErrorResponse preserves a trusted explicit non-cyber code", async () => { + const response = formatErrorResponse(400, "invalid_request_error", "target cannot represent tool choice", { + code: "target_incompatible", + }); + expect(response.status).toBe(400); + expect(await response.json()).toMatchObject({ + error: { type: "invalid_request_error", code: "target_incompatible" }, + }); + }); +}); + describe("error fidelity", () => { test("classifyError maps Codex-recognized context/quota/rate failures", () => { expect(classifyError(400, "upstream_error", "Your input exceeds the context window")).toMatchObject({ diff --git a/tests/key-failover.test.ts b/tests/key-failover.test.ts index ea10f8ee67..4cdb04f095 100644 --- a/tests/key-failover.test.ts +++ b/tests/key-failover.test.ts @@ -7,10 +7,21 @@ import { clearKeyCooldowns, getKeyCooldownUntil, hasKeyPoolFailover, + hydrateKeyQuotaCooldownsFromDisk, + resetKeyQuotaCooldownPersistenceForTests, rotateKeyOn429, rotateProviderTransportOn429, } from "../src/providers/key-failover"; import { deriveXaiConvId } from "../src/providers/xai-transport"; +import { + cancelPendingKeyQuotaCooldownPersist, + flushKeyQuotaCooldownPersistForTests, +} from "../src/providers/key-cooldown-disk"; +import { + clearCachedProviderQuotas, + getCachedProviderQuota, + setCachedProviderQuotaForTests, +} from "../src/providers/quota-routing-cache"; import { routeModel } from "../src/router"; import type { OcxConfig, OcxParsedRequest, OcxProviderConfig } from "../src/types"; import { removeTreeWithRetry } from "./helpers/remove-tree"; @@ -42,13 +53,19 @@ function pool3(): OcxProviderConfig["apiKeyPool"] { beforeEach(() => { home = mkdtempSync(join(tmpdir(), "ocx-keyfailover-")); process.env.OPENCODEX_HOME = home; + cancelPendingKeyQuotaCooldownPersist(); + resetKeyQuotaCooldownPersistenceForTests(); clearKeyCooldowns(); + clearCachedProviderQuotas(); }); afterEach(() => { + cancelPendingKeyQuotaCooldownPersist(); + resetKeyQuotaCooldownPersistenceForTests(); + clearKeyCooldowns(); + clearCachedProviderQuotas(); delete process.env.OPENCODEX_HOME; removeTreeWithRetry(home); - clearKeyCooldowns(); }); describe("hasKeyPoolFailover", () => { @@ -71,6 +88,63 @@ describe("rotateKeyOn429", () => { expect(getKeyCooldownUntil("p", "k1", now)).toBe(now + 60_000); }); + test("uses fresh exhausted quota reset for the failed key instead of the 60s default", () => { + const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); + const now = 1_000_000; + setCachedProviderQuotaForTests("p", { + fiveHourPercent: 100, fiveHourResetAt: now + 2 * 60 * 60_000, updatedAt: now, + }); + rotateKeyOn429(config, "p", null, now, "key-alpha-000111222333"); + expect(getKeyCooldownUntil("p", "k1", now)).toBe(now + 2 * 60 * 60_000); + }); + + test("rotating a key invalidates provider-level quota evidence from the failed key", () => { + const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); + const now = 1_000_000; + setCachedProviderQuotaForTests("p", { + fiveHourPercent: 100, fiveHourResetAt: now + 2 * 60 * 60_000, updatedAt: now, + }); + expect(getCachedProviderQuota("p", now)).not.toBeNull(); + rotateKeyOn429(config, "p", null, now, "key-alpha-000111222333"); + expect(getCachedProviderQuota("p", now)).toBeNull(); + }); + + test("stale exhausted quota does not create a long key cooldown", () => { + const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); + const now = 50_000_000; + setCachedProviderQuotaForTests("p", { + fiveHourPercent: 100, fiveHourResetAt: now + 2 * 60 * 60_000, + updatedAt: now - 30 * 60_000 - 1, + }); + rotateKeyOn429(config, "p", null, now, "key-alpha-000111222333"); + expect(getKeyCooldownUntil("p", "k1", now)).toBe(now + 60_000); + }); + + test("long quota cooldown survives a simulated proxy restart", () => { + const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); + const now = 1_000_000; + hydrateKeyQuotaCooldownsFromDisk(now); + setCachedProviderQuotaForTests("p", { + fiveHourPercent: 100, fiveHourResetAt: now + 2 * 60 * 60_000, updatedAt: now, + }); + rotateKeyOn429(config, "p", null, now, "key-alpha-000111222333"); + flushKeyQuotaCooldownPersistForTests(now); + resetKeyQuotaCooldownPersistenceForTests(); + hydrateKeyQuotaCooldownsFromDisk(now + 1_000); + expect(getKeyCooldownUntil("p", "k1", now + 1_000)).toBe(now + 2 * 60 * 60_000); + }); + + test("short generic 429 cooldown is not persisted across restart", () => { + const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); + const now = 1_000_000; + hydrateKeyQuotaCooldownsFromDisk(now); + rotateKeyOn429(config, "p", null, now, "key-alpha-000111222333"); + flushKeyQuotaCooldownPersistForTests(now); + resetKeyQuotaCooldownPersistenceForTests(); + hydrateKeyQuotaCooldownsFromDisk(now + 1_000); + expect(getKeyCooldownUntil("p", "k1", now + 1_000)).toBeNull(); + }); + test("respects Retry-After seconds for the cooldown window", () => { const config = makeConfig({ apiKey: "key-alpha-000111222333", apiKeyPool: pool3() }); const now = 1_000_000; diff --git a/tests/provider-quota.test.ts b/tests/provider-quota.test.ts index 2768ead96f..dcd5c32fa2 100644 --- a/tests/provider-quota.test.ts +++ b/tests/provider-quota.test.ts @@ -834,6 +834,34 @@ describe("fetchProviderQuotaReports", () => { expect(seen[0]?.redirect).toBe("error"); }); + test("DeepSeek quota recognizes a title-cased built-in destination and marks zero balance exhausted", async () => { + globalThis.fetch = (async () => new Response(JSON.stringify({ + is_available: false, + balance_infos: [{ currency: "USD", total_balance: "0", granted_balance: "0", topped_up_balance: "0" }], + }), { status: 200 })) as typeof fetch; + const config = keyQuotaConfig("DeepSeek", "https://api.deepseek.com"); + const result = await fetchProviderQuotaReports(config, true); + expect(result.reports).toHaveLength(1); + expect(result.reports[0]?.provider).toBe("DeepSeek"); + expect(result.reports[0]?.source).toBe("deepseek:balance"); + expect(result.reports[0]?.quota.customWindows).toEqual([{ + label: "API balance ($0.00)", + percent: 100, + }]); + }); + + test("DeepSeek quota treats is_available=false as authoritative exhaustion even with positive balance", async () => { + globalThis.fetch = (async () => new Response(JSON.stringify({ + is_available: false, + balance_infos: [{ currency: "USD", total_balance: "0.01", topped_up_balance: "0.01" }], + }), { status: 200 })) as typeof fetch; + const result = await fetchProviderQuotaReports(keyQuotaConfig("DeepSeek", "https://api.deepseek.com"), true); + expect(result.reports).toHaveLength(1); + expect(result.reports[0]?.quota.customWindows).toEqual([{ + label: "API balance ($0.01)", percent: 100, + }]); + }); + test("DeepSeek quota never sends the key to a non-canonical base URL", async () => { const seen: string[] = []; globalThis.fetch = (async (input: RequestInfo | URL) => { diff --git a/tests/request-pacing.test.ts b/tests/request-pacing.test.ts index bf8cfbb106..7400153f90 100644 --- a/tests/request-pacing.test.ts +++ b/tests/request-pacing.test.ts @@ -2,6 +2,7 @@ import { afterEach, describe, expect, test } from "bun:test"; import { providerRequestPacingStatus, reconcileProviderRequestPacing, + RequestPacingProviderRemovedError, RequestPacingQueueOverloadError, requestPacingIntervalMs, resetProviderRequestPacingForTest, @@ -235,6 +236,14 @@ describe("provider request pacing queue", () => { expect(clock.pendingTimerCount()).toBe(0); }); + test("maps a provider removed during pacing to retryable provider-unavailable", async () => { + const response = requestPacingOverloadResponse(new RequestPacingProviderRemovedError("demo")); + expect(response?.status).toBe(503); + expect(await response?.json()).toMatchObject({ + error: { type: "server_error", code: "provider_unavailable" }, + }); + }); + test("maps pacing admission overload to 429 with Retry-After", async () => { const response = requestPacingOverloadResponse(new RequestPacingQueueOverloadError("demo", "queue_full", 3)); expect(response?.status).toBe(429); diff --git a/tests/responses-context-overflow.test.ts b/tests/responses-context-overflow.test.ts index bf5b1c2ba2..f82af14eaf 100644 --- a/tests/responses-context-overflow.test.ts +++ b/tests/responses-context-overflow.test.ts @@ -199,7 +199,7 @@ describe("Responses provider input overflow", () => { } }); - test("a combo stops on 413 and does not dispatch a second oversized target", async () => { + test("a combo tries the next target after 413 and reports overflow only after exhaustion", async () => { let firstHits = 0; let secondHits = 0; const first = upstream413(() => { firstHits += 1; }); @@ -223,7 +223,7 @@ describe("Responses provider input overflow", () => { const failed = await responseFailed(await request(String(server.url), "combo/fallback", true)); expect((failed.error as { code?: string }).code).toBe("context_length_exceeded"); expect(firstHits).toBe(1); - expect(secondHits).toBe(0); + expect(secondHits).toBe(1); } finally { await server.stop(true); } diff --git a/tests/responses-pool-401-refresh.test.ts b/tests/responses-pool-401-refresh.test.ts index 9b11ffb672..0f5c921352 100644 --- a/tests/responses-pool-401-refresh.test.ts +++ b/tests/responses-pool-401-refresh.test.ts @@ -19,6 +19,8 @@ import { import type { RequestLogContext } from "../src/server/request-log"; import type { OcxConfig } from "../src/types"; import { removeTreeWithRetry } from "./helpers/remove-tree"; +import { clearCachedProviderQuotas } from "../src/providers/quota-routing-cache"; +import { clearComboSelectionState, clearComboTargetCooldowns } from "../src/combos"; /** * #2887: an ordinary stored pool account holding a TIME-VALID access token that upstream @@ -185,6 +187,9 @@ beforeEach(() => { clearAccountNeedsReauth(OTHER_ACCOUNT_ID); clearCodexUpstreamHealth(); clearThreadAccountMap(); + clearCachedProviderQuotas(); + clearComboSelectionState(); + clearComboTargetCooldowns(); writeStoredAccount(); }); @@ -273,12 +278,12 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { expect(harness.refreshes).toEqual(["refresh-grant"]); }); - test("a combo stops after a stored-account replay consumes the recovery budget", async () => { + test("a combo fails over after a stored-account replay exhausts quota before output", async () => { const cfg = recoveryComboConfig(); const harness = installHarness({ responseForSend: (authorization, _sendNumber, url) => { if (url.hostname === "backup.example") { - return Response.json({ id: "must-not-run", object: "response", status: "completed", output: [] }); + return Response.json({ id: "backup-ok", object: "response", status: "completed", output: [] }); } if (authorization === "Bearer rejected-access") { return Response.json({ error: { message: "rejected bearer" } }, { status: 401 }); @@ -296,17 +301,17 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { { model: "", provider: "" } as RequestLogContext, ); - expect(response.status).toBe(429); - expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(response.status).toBe(200); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access", "Bearer backup-test-key"]); expect(harness.refreshes).toEqual(["refresh-grant"]); }); - test("a combo stops when the stored-account replay hits a transport error", async () => { + test("a combo fails over when the stored-account replay hits a pre-output transport error", async () => { const cfg = recoveryComboConfig(); const harness = installHarness({ responseForSend: (authorization, _sendNumber, url) => { if (url.hostname === "backup.example") { - return Response.json({ id: "must-not-run", object: "response", status: "completed", output: [] }); + return Response.json({ id: "backup-ok", object: "response", status: "completed", output: [] }); } if (authorization === "Bearer rejected-access") { return Response.json({ error: { message: "rejected bearer" } }, { status: 401 }); @@ -324,17 +329,17 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { { model: "", provider: "" } as RequestLogContext, ); - expect(response.status).toBe(502); - expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(response.status).toBe(200); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access", "Bearer backup-test-key"]); expect(harness.refreshes).toEqual(["refresh-grant"]); }); - test("a combo stops on a zero-output failure from the stored-account replay stream", async () => { + test("a combo fails over on a zero-output failure from the stored-account replay stream", async () => { const cfg = recoveryComboConfig(); const harness = installHarness({ responseForSend: (authorization, _sendNumber, url) => { if (url.hostname === "backup.example") { - return Response.json({ id: "must-not-run", object: "response", status: "completed", output: [] }); + return Response.json({ id: "backup-ok", object: "response", status: "completed", output: [] }); } if (authorization === "Bearer rejected-access") { return Response.json({ error: { message: "rejected bearer" } }, { status: 401 }); @@ -366,13 +371,8 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { { model: "", provider: "" } as RequestLogContext, ); - expect(response.status).toBe(502); - const failure = await response.clone().json() as { - error?: { code?: string; message?: string }; - }; - expect(failure.error?.code).toBe("upstream_server_error"); - expect(failure.error?.message).toContain("busy"); - expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(response.status).toBe(200); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access", "Bearer backup-test-key"]); expect(harness.refreshes).toEqual(["refresh-grant"]); }); @@ -420,7 +420,7 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { } if (authorization === "Bearer other-access") { alternateAccountSends += 1; - return Response.json({ id: "must-not-run", object: "response", status: "completed", output: [] }); + return Response.json({ id: "backup-ok", object: "response", status: "completed", output: [] }); } return undefined; }, @@ -758,7 +758,7 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); }); - test("a combo stops after a stored replay 4xx that is neither quota nor a gated-model 400", async () => { + test("an explicit combo can fail over after a stored replay fails before first output", async () => { // The contributor's original bound was a single `status >= 400` break in the passthrough loop, // which stopped combo fallback for EVERY stored replay 4xx. Removing it to keep same-account // rescue alive means the outer layers now rely on the dispatch signal instead. This pins that @@ -786,9 +786,14 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { { model: "", provider: "" } as RequestLogContext, ); - expect(response.status).toBe(403); - // The backup target is never sent to, and the account is charged exactly twice. - expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(response.status).toBe(200); + // The refreshed account was tried exactly once; because its replay failed before first + // output, the explicitly declared combo remains free to use its backup provider. + expect(harness.sends).toEqual([ + "Bearer rejected-access", + "Bearer refreshed-access", + "Bearer backup-test-key", + ]); expect(harness.refreshes).toEqual(["refresh-grant"]); }); diff --git a/tests/routing-policy-fallback.test.ts b/tests/routing-policy-fallback.test.ts index acfadb3744..ae79c05275 100644 --- a/tests/routing-policy-fallback.test.ts +++ b/tests/routing-policy-fallback.test.ts @@ -111,9 +111,9 @@ describe("policy candidate fallback", () => { expect(response.status).toBe(400); expect(seenModels).toEqual(["policy/daily"]); }); - test("an upstream context_length_exceeded still stops the chain (#1524)", async () => { - // The mirror-image contract. An upstream verdict is about the REQUEST, so retrying it - // elsewhere is guesswork -- and hopping would burn every candidate on a doomed request. + test("an upstream context_length_exceeded hops to a larger declared candidate (#1524)", async () => { + // In an explicit policy candidate set, context capacity is target-local: another declared + // model may have a larger window, so a pre-output context verdict remains replay-safe. const trace = policyTrace(); const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; const seenModels: string[] = []; @@ -131,7 +131,7 @@ describe("policy candidate fallback", () => { const response = await handleResponsesWithPolicyFallback(request(), {} as OcxConfig, logCtx, {}, { runCore }); expect(response.status).toBe(400); - expect(seenModels).toEqual(["policy/daily"]); + expect(seenModels).toEqual(["policy/daily", "provider-b/model-b", "provider-c/model-c"]); }); test("retries the next policy candidate and keeps distinct physical attempts", async () => { const trace = policyTrace(); @@ -179,7 +179,7 @@ describe("policy candidate fallback", () => { expect(logCtx.activeAttempt).toBe(logCtx.attempts?.[1]); }); - test("a stored Pool 401 replay dispatch stops policy candidate fallback", async () => { + test("a stored Pool 401 replay that fails before output can still hop policy candidates", async () => { const trace = policyTrace(); const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; const seenModels: string[] = []; @@ -192,35 +192,78 @@ describe("policy candidate fallback", () => { const body = await req.json() as { model: string }; seenModels.push(body.model); childLog.routeDecision = trace; - seedAttempt(childLog, "provider-a", "model-a"); - options.onStoredPool401ReplayDispatched?.(); - return Response.json( - { error: { message: "stored replay exhausted", type: "rate_limit_error" } }, - { status: 429 }, - ); + const first = seenModels.length === 1; + seedAttempt(childLog, first ? "provider-a" : "provider-b", first ? "model-a" : "model-b"); + if (first) { + options.onStoredPool401ReplayDispatched?.(); + return Response.json({ error: { message: "stored replay exhausted", type: "rate_limit_error" } }, { status: 429 }); + } + return Response.json({ status: "completed" }); }, }); - expect(response.status).toBe(429); - expect(seenModels).toEqual(["policy/daily"]); + expect(response.status).toBe(200); + expect(seenModels).toEqual(["policy/daily", "provider-b/model-b"]); expect(replaySignals).toBe(1); }); - test("returns local pacing overload without switching policy candidates", async () => { + test("a stored Pool 401 replay stops policy fallback after visible output", async () => { const trace = policyTrace(); const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; let calls = 0; const response = await handleResponsesWithPolicyFallback(request(), {} as OcxConfig, logCtx, {}, { - runCore: async () => { + runCore: async (_req, _config, childLog, options) => { calls += 1; - throw new RequestPacingQueueOverloadError("provider-a", "queue_full", 2); + childLog.routeDecision = trace; + seedAttempt(childLog, "provider-a", "model-a"); + childLog.activeAttempt!.firstOutputMs = 1; + options.onStoredPool401ReplayDispatched?.(); + return Response.json({ error: { message: "late replay failure", type: "rate_limit_error" } }, { status: 429 }); }, }); expect(response.status).toBe(429); - expect(response.headers.get("Retry-After")).toBe("2"); expect(calls).toBe(1); }); + test("local pacing overload hops to the next policy candidate", async () => { + const trace = policyTrace(); + const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; + const seenModels: string[] = []; + const response = await handleResponsesWithPolicyFallback(request(), {} as OcxConfig, logCtx, {}, { + runCore: async (req, _config, childLog) => { + const body = await req.clone().json() as { model: string }; + seenModels.push(body.model); + childLog.routeDecision = trace; + const first = seenModels.length === 1; + seedAttempt(childLog, first ? "provider-a" : "provider-b", first ? "model-a" : "model-b"); + if (first) throw new RequestPacingQueueOverloadError("provider-a", "queue_full", 2); + return Response.json({ status: "completed" }); + }, + }); + expect(response.status).toBe(200); + expect(seenModels).toEqual(["policy/daily", "provider-b/model-b"]); + }); + + test("all retryable policy candidates exhausted returns a sanitized policy-level error", async () => { + const trace = policyTrace(); + const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; + const seenModels: string[] = []; + const response = await handleResponsesWithPolicyFallback(request(), {} as OcxConfig, logCtx, {}, { + runCore: async (req, _config, childLog) => { + const body = await req.clone().json() as { model: string }; + seenModels.push(body.model); + childLog.routeDecision = trace; + seedAttempt(childLog, "provider", body.model); + return Response.json({ error: { message: "provider-secret-detail", type: "rate_limit_error" } }, { status: 429 }); + }, + }); + expect(response.status).toBe(429); + const text = await response.text(); + expect(text).toContain("All eligible policy candidates are temporarily unavailable"); + expect(text).not.toContain("provider-secret-detail"); + expect(seenModels).toEqual(["policy/daily", "provider-b/model-b", "provider-c/model-c"]); + }); + test("does not switch candidates for terminal client/input failures", async () => { const trace = policyTrace(); const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; diff --git a/tests/server-auth.test.ts b/tests/server-auth.test.ts index 37db545e46..e2487d416d 100644 --- a/tests/server-auth.test.ts +++ b/tests/server-auth.test.ts @@ -678,7 +678,8 @@ describe("server local API auth", () => { }, }); - expect(response.status).toBe(429); + expect(response.status).toBe(503); + expect(await response.clone().text()).toContain("combo_unavailable"); expect(acceptedCount).toBe(1); expect(upstreamModels).toEqual(["first-model", "second-model"]); } finally { diff --git a/tests/server-combo-failover-e2e.test.ts b/tests/server-combo-failover-e2e.test.ts index 39053f39fb..b6f7e866d6 100644 --- a/tests/server-combo-failover-e2e.test.ts +++ b/tests/server-combo-failover-e2e.test.ts @@ -37,6 +37,7 @@ import { } from "../src/responses/state"; import { clearCursorThreadContinuityForTests } from "../src/adapters/cursor/thread-continuity"; import { COMPACT_PROMPT, encodeCompactionSummary } from "../src/responses/compaction"; +import { RequestPacingQueueOverloadError } from "../src/providers/request-pacing"; // Full-suite Windows load: startServer + combo rename/delete management flows exceed the // default 5s per-test budget (same flake class as 810fa115 / claude-management-api). @@ -411,6 +412,46 @@ async function within(promise: Promise, ms = 2_000): Promise { } describe("server combo failover 030 activation matrix", () => { + test("target-local context, billing, and unsupported-model failures reach the healthy fallback", async () => { + const hits: string[] = []; + const make = (name: string, response: () => Response) => serve(() => { hits.push(name); return response(); }); + const a = make("context", () => Response.json({ error: { code: "context_length_exceeded", message: "maximum context exceeded" } }, { status: 400 })); + const b = make("billing", () => Response.json({ error: { code: "insufficient_quota", message: "payment required" } }, { status: 402 })); + const c = make("model", () => Response.json({ error: { code: "unsupported_model", message: "unsupported model" } }, { status: 400 })); + const d = make("gone", () => Response.json({ error: { message: "resource gone" } }, { status: 410 })); + const e = make("too-large", () => Response.json({ error: { message: "request too large" } }, { status: 413 })); + const f = make("healthy", () => chatSuccess("deep fallback", "m6")); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "ka"), + b: provider("openai-chat", baseUrl(b), "kb"), + c: provider("openai-chat", baseUrl(c), "kc"), + d: provider("openai-chat", baseUrl(d), "kd"), + e: provider("openai-chat", baseUrl(e), "ke"), + f: provider("openai-chat", baseUrl(f), "kf"), + }); + const response = await post(config); + expect(response.status).toBe(200); + expect(JSON.stringify(await response.json())).toContain("deep fallback"); + expect(hits).toEqual(["context", "billing", "model", "gone", "too-large", "healthy"]); + }); + + test("local pacing overload before dispatch hops to the next combo target", async () => { + let backupHits = 0; + customFetchResponse = async request => { + const model = (JSON.parse(String(request.body)) as { model?: string }).model; + if (model === "m1") throw new RequestPacingQueueOverloadError("a", "queue_full", 2); + backupHits += 1; + return chatSuccess("pacing backup", "m2"); + }; + const config = comboConfig({ + a: provider("test-response", "https://test.invalid/v1", "ka"), + b: provider("test-response", "https://test.invalid/v1", "kb"), + }); + const response = await post(config); + expect(response.status).toBe(200); + expect(JSON.stringify(await response.json())).toContain("pacing backup"); + expect(backupHits).toBe(1); + }); test("dispatches a selected concrete target despite a shadowing combo alias", async () => { const hits: string[] = []; const a = serve(async request => { @@ -480,6 +521,140 @@ describe("server combo failover 030 activation matrix", () => { expect(hits).toEqual(["go:m1", "orca:m2", "backup:m3"]); }); + test("paid DeepSeek last resort wins after free quota failures without leaking intermediate errors", async () => { + const hits: string[] = []; + const monthly = serve(() => { + hits.push("free-monthly"); + return Response.json({ error: { type: "GoUsageLimitError", message: "Monthly usage limit reached. Resets in 14 days." } }, { status: 429 }); + }); + const capped = serve(() => { + hits.push("free-prompt-cap"); + return Response.json({ error: { + type: "invalid_request_error", + code: "free_rate_limited", + message: "This prompt is longer than the free tier allows for a single request.", + metadata: { reason: "err_free_prompt_cap" }, + } }, { status: 400 }); + }); + const paid = serve(() => { + hits.push("DeepSeek-paid"); + return chatSuccess("paid-last-resort-ok", "deepseek-v4-flash-vision-exp"); + }); + const config = comboConfig({ + free1: provider("openai-chat", baseUrl(monthly), "key-free1"), + free2: provider("openai-chat", baseUrl(capped), "key-free2"), + DeepSeek: provider("openai-chat", baseUrl(paid), "key-paid"), + }, [ + { provider: "free1", model: "m1" }, + { provider: "free2", model: "m2" }, + { provider: "DeepSeek", model: "deepseek-v4-flash-vision-exp" }, + ]); + + const response = await post(config); + expect(response.status).toBe(200); + const text = await response.text(); + expect(text).toContain("paid-last-resort-ok"); + expect(text).not.toContain("Monthly usage limit"); + expect(text).not.toContain("free_rate_limited"); + expect(text).not.toContain("err_free_prompt_cap"); + expect(hits).toEqual(["free-monthly", "free-prompt-cap", "DeepSeek-paid"]); + }); + + test("all-target exhaustion with DeepSeek 402 returns only sanitized combo metadata", async () => { + const hits: string[] = []; + const free = serve(() => { + hits.push("free"); + return Response.json({ error: { type: "GoUsageLimitError", message: "FREE-RAW monthly usage limit reached. Resets in 14 days." } }, { status: 429 }); + }); + const paid = serve(() => { + hits.push("DeepSeek"); + return Response.json({ error: { type: "insufficient_balance", code: "insufficient_balance", message: "DEEPSEEK-RAW account balance is zero" } }, { status: 402 }); + }); + const config = comboConfig({ + free: provider("openai-chat", baseUrl(free), "key-free"), + DeepSeek: provider("openai-chat", baseUrl(paid), "key-paid"), + }, [ + { provider: "free", model: "m1" }, + { provider: "DeepSeek", model: "deepseek-v4-flash" }, + ]); + const response = await post(config); + expect(response.status).toBe(503); + const body = await response.json() as any; + expect(body.error.code).toBe("combo_unavailable"); + expect(body.error.details.last_provider).toBe("DeepSeek"); + expect(body.error.details.correlation_id).toMatch(/^combo-/); + expect(body.error.details.causes).toEqual(expect.arrayContaining([ + { category: "quota", count: 1 }, { category: "billing", count: 1 }, + ])); + const text = JSON.stringify(body); + expect(text).not.toContain("FREE-RAW"); + expect(text).not.toContain("DEEPSEEK-RAW"); + expect(text).not.toContain("key-paid"); + expect(hits).toEqual(["free", "DeepSeek"]); + }); + + test("concurrent turns share a newly learned cooldown and call the paid fallback once per logical request", async () => { + let freeHits = 0; + let paidHits = 0; + const firstPaidEntered = deferred(); + const releasePaid = deferred(); + const free = serve(() => { + freeHits += 1; + return Response.json({ error: { message: "rate limited" } }, { + status: 429, + headers: { "retry-after": "120" }, + }); + }); + const paid = serve(async () => { + paidHits += 1; + if (paidHits === 1) firstPaidEntered.resolve(); + await releasePaid.promise; + return chatSuccess(`paid-${paidHits}`, "deepseek-v4-flash-vision-exp"); + }); + const config = comboConfig({ + free: provider("openai-chat", baseUrl(free), "key-free"), + DeepSeek: provider("openai-chat", baseUrl(paid), "key-paid"), + }, [ + { provider: "free", model: "m1" }, + { provider: "DeepSeek", model: "deepseek-v4-flash-vision-exp" }, + ]); + + const first = post(config); + await within(firstPaidEntered.promise); + const second = post(config); + await within((async () => { while (paidHits < 2) await new Promise(resolve => setTimeout(resolve, 1)); })()); + releasePaid.resolve(); + const [r1, r2] = await Promise.all([first, second]); + expect(r1.status).toBe(200); + expect(r2.status).toBe(200); + expect(freeHits).toBe(1); + expect(paidHits).toBe(2); + }); + + test("adapter-local tool-choice incompatibility hops before upstream dispatch", async () => { + let ollamaHits = 0; + const ollama = serve(() => { + ollamaHits += 1; + return Response.json({ message: { role: "assistant", content: "must not dispatch" }, done: true }); + }); + const backup = serve(async request => { + const body = await request.json() as { model?: string }; + return chatSuccess("tool-choice fallback", body.model ?? "m2"); + }); + const config = comboConfig({ + a: provider("ollama-native", baseUrl(ollama), "ka"), + b: provider("openai-chat", baseUrl(backup), "kb"), + }); + + const response = await post(config, { + tools: [{ type: "function", name: "read", description: "read", parameters: { type: "object", properties: {} } }], + tool_choice: "required", + }); + expect(response.status).toBe(200); + expect(JSON.stringify(await response.json())).toContain("tool-choice fallback"); + expect(ollamaHits).toBe(0); + }); + test("ordinary openai-chat 503 hops to backup for non-stream and stream", async () => { const hits: string[] = []; const a = serve(async request => { @@ -507,6 +682,57 @@ describe("server combo failover 030 activation matrix", () => { expect(hits).toEqual(["a:false", "b:false", "a:true", "b:true"]); }); + test("HTTP 200 Orca-style free prompt cap hops request-locally without cooling the provider", async () => { + const hits: string[] = []; + const capped = serve(() => { + hits.push("orca"); + const frame = `data: ${JSON.stringify({ error: { + type: "invalid_request_error", code: "free_rate_limited", + message: "ORCA-RAW free tier prompt cap err_free_prompt_cap", + } })}\n\ndata: [DONE]\n\n`; + return new Response(frame, { headers: { "content-type": "text/event-stream" } }); + }); + const backup = serve(() => { hits.push("backup"); return chatStream("sse-cap-backup"); }); + const config = comboConfig({ + orca: provider("openai-chat", baseUrl(capped), "key-orca"), + backup: provider("openai-chat", baseUrl(backup), "key-b"), + }); + for (let turn = 0; turn < 2; turn += 1) { + const response = await post(config, { stream: true }); + expect(response.status).toBe(200); + const text = JSON.stringify(await collectSse(response)); + expect(text).toContain("sse-cap-backup"); + expect(text).not.toContain("ORCA-RAW"); + } + expect(hits).toEqual(["orca", "backup", "orca", "backup"]); + }); + + test("HTTP 200 stream body reset before output hops through the real adapter", async () => { + const hits: string[] = []; + customFetchResponse = async request => { + const model = (JSON.parse(String(request.body)) as { model?: string }).model; + hits.push(String(model)); + if (model === "m1") { + return new Response(new ReadableStream({ + pull() { + throw Object.assign(new Error("upstream-body-secret reset"), { code: "ECONNRESET" }); + }, + }), { status: 200, headers: { "content-type": "text/event-stream" } }); + } + return chatStream("backup after body reset"); + }; + const config = comboConfig({ + a: provider("test-response", "https://test.invalid/v1", "key-a"), + b: provider("test-response", "https://test.invalid/v1", "key-b"), + }); + const response = await post(config, { stream: true }); + expect(response.status).toBe(200); + const text = JSON.stringify(await collectSse(response)); + expect(text).toContain("backup after body reset"); + expect(text).not.toContain("upstream-body-secret"); + expect(hits).toEqual(["m1", "m2"]); + }); + test("zero-output terminal SSE failure hops before committing the child stream", async () => { const hits: string[] = []; const a = serve(() => { @@ -577,6 +803,113 @@ describe("server combo failover 030 activation matrix", () => { } }); + test("HTTP 200 invalid JSON before output hops to the next combo target", async () => { + const hits: string[] = []; + const a = serve(() => { + hits.push("a"); + return new Response("{invalid-json", { + status: 200, + headers: { "content-type": "application/json" }, + }); + }); + const b = serve(() => { + hits.push("b"); + return chatSuccess("backup after invalid json", "m2"); + }); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "key-a"), + b: provider("openai-chat", baseUrl(b), "key-b"), + }); + + const response = await post(config); + expect(response.status).toBe(200); + expect(JSON.stringify(await response.json())).toContain("backup after invalid json"); + expect(hits).toEqual(["a", "b"]); + }); + + test("HTTP 200 malformed SSE before output hops without leaking the malformed frame", async () => { + const hits: string[] = []; + const a = serve(() => { + hits.push("a"); + return new Response("data: {not-json}\n\n", { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + }); + const b = serve(() => { + hits.push("b"); + return chatStream("backup after malformed sse"); + }); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "key-a"), + b: provider("openai-chat", baseUrl(b), "key-b"), + }); + + const response = await post(config, { stream: true }); + expect(response.status).toBe(200); + const frames = await collectSse(response); + const text = JSON.stringify(frames); + expect(text).toContain("backup after malformed sse"); + expect(text).not.toContain("not-json"); + expect(hits).toEqual(["a", "b"]); + }); + + test("malformed SSE after visible output never replays on a backup target", async () => { + const hits: string[] = []; + const a = serve(() => { + hits.push("a"); + const frames = [ + `data: ${JSON.stringify({ choices: [{ index: 0, delta: { content: "visible-before-malformed" }, finish_reason: null }] })}\n\n`, + "data: {MALFORMED-RAW\n\n", + ].join(""); + return new Response(frames, { headers: { "content-type": "text/event-stream" } }); + }); + const b = serve(() => { hits.push("b"); return chatStream("must-not-replay"); }); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "key-a"), + b: provider("openai-chat", baseUrl(b), "key-b"), + }); + const response = await post(config, { stream: true }); + expect(response.status).toBe(200); + const text = JSON.stringify(await collectSse(response)); + expect(text).toContain("visible-before-malformed"); + expect(text).not.toContain("must-not-replay"); + expect(text).not.toContain("MALFORMED-RAW"); + expect(hits).toEqual(["a"]); + }); + + test("raw stream reset after output never replays on a backup target", async () => { + const hits: string[] = []; + const encoder = new TextEncoder(); + customFetchResponse = async request => { + const model = (JSON.parse(String(request.body)) as { model?: string }).model; + hits.push(String(model)); + if (model !== "m1") return chatStream("must-not-run-backup"); + let sent = false; + return new Response(new ReadableStream({ + pull(controller) { + if (!sent) { + sent = true; + controller.enqueue(encoder.encode(`data: ${JSON.stringify({ choices: [{ index: 0, delta: { content: "visible-once" }, finish_reason: null }] })}\n\n`)); + return; + } + controller.error(Object.assign(new Error("late-upstream-secret reset"), { code: "ECONNRESET" })); + }, + }), { status: 200, headers: { "content-type": "text/event-stream" } }); + }; + const config = comboConfig({ + a: provider("test-response", "https://test.invalid/v1", "key-a"), + b: provider("test-response", "https://test.invalid/v1", "key-b"), + }); + const response = await post(config, { stream: true }); + expect(response.status).toBe(200); + const text = JSON.stringify(await collectSse(response)); + expect(text).toContain("visible-once"); + expect(text).not.toContain("must-not-run-backup"); + expect(text).not.toContain("late-upstream-secret"); + expect(hits).toEqual(["m1"]); + }); + test("terminal SSE failure after output stays on the first target and never replays", async () => { const hits: string[] = []; const a = serve(() => { @@ -1104,7 +1437,8 @@ describe("server combo failover 030 activation matrix", () => { ].join("\n"), { headers: { "content-type": "text/event-stream" } }); const response = await post(config, { stream: true }); - expect(response.status).toBe(502); + expect(response.status).toBe(503); + expect(await response.clone().text()).toContain("combo_unavailable"); expect(getCodexUpstreamHealth(rawAccountId)).toMatchObject({ consecutiveFailures: 1, lastFailureStatus: 502, @@ -1204,8 +1538,8 @@ describe("server combo failover 030 activation matrix", () => { }; const response = await postLogged(config); - expect(response.status).toBe(429); - await response.text(); + expect(response.status).toBe(503); + expect(await response.text()).toContain("combo_unavailable"); expect(calls).toBe(1); expect(getCodexUpstreamHealth(rawAccountId)?.cooldownSource).toBe("retry-after"); }); @@ -1385,6 +1719,70 @@ describe("server combo failover 030 activation matrix", () => { expect(JSON.stringify(await response.json())).toContain("connected backup"); }); + test("exhausts invalid API keys within a provider before hopping to the next combo provider", async () => { + const seenAuth: string[] = []; + const a = serve(req => { + seenAuth.push(req.headers.get("authorization") ?? ""); + return Response.json({ error: { type: "authentication_error", code: "invalid_api_key", message: "bad pooled credential" } }, { status: 401 }); + }); + let backupHits = 0; + const b = serve(() => { + backupHits += 1; + return chatSuccess("healthy provider after key pool", "m2"); + }); + const first = provider("openai-chat", baseUrl(a), "key-alpha-000111222333", { + apiKeyPool: [ + { id: "k1", key: "key-alpha-000111222333", addedAt: 1 }, + { id: "k2", key: "key-beta-444555666777", addedAt: 2 }, + ], + }); + const config = comboConfig({ + pooled: first, + backup: provider("openai-chat", baseUrl(b), "key-backup-000111222333"), + }, [ + { provider: "pooled", model: "m1" }, + { provider: "backup", model: "m2" }, + ]); + const response = await post(config); + expect(response.status).toBe(200); + const text = JSON.stringify(await response.json()); + expect(text).toContain("healthy provider after key pool"); + expect(text).not.toContain("bad pooled credential"); + expect(seenAuth).toEqual(["Bearer key-alpha-000111222333", "Bearer key-beta-444555666777"]); + expect(backupHits).toBe(1); + }); + + test("transport exception matrix hops to backup without leaking intermediate details", async () => { + const cases = [ + Object.assign(new Error("socket reset"), { code: "ECONNRESET" }), + Object.assign(new Error("getaddrinfo ENOTFOUND provider.invalid"), { code: "ENOTFOUND" }), + Object.assign(new Error("getaddrinfo EAI_AGAIN provider.invalid"), { code: "EAI_AGAIN" }), + Object.assign(new Error('ERR_TLS_CERT_ALTNAME_INVALID fetching "https://user:provider-secret@[::1]/v1"'), { code: "ERR_TLS_CERT_ALTNAME_INVALID" }), + Object.assign(new Error("Timeout elapsed"), { name: "TimeoutError" }), + ]; + for (const failure of cases) { + clearComboSelectionState(); + clearComboTargetCooldowns(); + let backupHits = 0; + customFetchResponse = async request => { + const model = (JSON.parse(String(request.body)) as { model?: string }).model; + if (model === "m1") throw failure; + backupHits += 1; + return chatSuccess("transport backup", "m2"); + }; + const config = comboConfig({ + a: provider("test-response", "https://test.invalid/v1", "key-a"), + b: provider("test-response", "https://test.invalid/v1", "key-b"), + }); + const response = await post(config); + expect(response.status).toBe(200); + const text = JSON.stringify(await response.json()); + expect(text).toContain("transport backup"); + expect(text).not.toContain("provider-secret"); + expect(backupHits).toBe(1); + } + }); + test("Azure passthrough 403 hops into ordinary openai-chat", async () => { const hits: string[] = []; const a = serve(() => { @@ -1514,20 +1912,22 @@ describe("server combo failover 030 activation matrix", () => { expect(modelHits.every(hit => hit.hasWebTool)).toBe(true); }); - test("context 400 stops while exhausted retryable targets return the sanitized last status", async () => { + test("context 400 hops while exhausted retryable targets return the sanitized last status", async () => { let stopBackupHits = 0; const context = serve(() => Response.json({ error: { code: "context_length_exceeded", message: "too many tokens" } }, { status: 400 })); const unused = serve(() => { stopBackupHits += 1; - return chatSuccess("must not run"); + return chatSuccess("context backup"); }); const stopConfig = comboConfig({ a: provider("openai-chat", baseUrl(context), "key-a"), b: provider("openai-chat", baseUrl(unused), "key-b"), }); const stopped = await post(stopConfig); - expect(stopped.status).toBe(400); - expect(stopBackupHits).toBe(0); + expect(stopped.status).toBe(200); + expect(stopBackupHits).toBe(1); + expect(JSON.stringify(await stopped.json())).toContain("context backup"); + clearComboTargetCooldowns("free"); const order: string[] = []; const first = serve(() => { @@ -1542,9 +1942,21 @@ describe("server combo failover 030 activation matrix", () => { a: provider("openai-chat", baseUrl(first), "key-a"), b: provider("openai-chat", baseUrl(last), "key-b"), })); - expect(exhausted.status).toBe(404); + expect(exhausted.status).toBe(503); expect(order).toEqual(["a", "b"]); - expect(await exhausted.text()).not.toContain("sk-a-should-redact"); + const exhaustedBody = await exhausted.json() as { + error?: { code?: string; details?: { combo?: string; correlation_id?: string; last_provider?: string; causes?: unknown[] } }; + }; + expect(exhaustedBody.error?.code).toBe("combo_unavailable"); + expect(exhaustedBody.error?.details).toMatchObject({ + combo: "free", + last_provider: "b", + }); + expect(exhaustedBody.error?.details?.correlation_id).toMatch(/^combo-/); + expect(exhaustedBody.error?.details?.causes?.length).toBeGreaterThan(0); + const exhaustedText = JSON.stringify(exhaustedBody); + expect(exhaustedText).not.toContain("missing model"); + expect(exhaustedText).not.toContain("sk-a-should-redact"); }); test("429 Retry-After 120 keeps A cooling at 60 seconds and restores it at 120", async () => { @@ -2483,7 +2895,8 @@ describe("server combo failover 030 activation matrix", () => { ...(includeC ? [{ provider: "c", model: "m3" }] : []), ]; const response = await post(comboConfig(providers, targets)); - expect(response.status).toBe(includeC ? 200 : 401); + expect(response.status).toBe(includeC ? 200 : 503); + if (!includeC) expect(await response.clone().text()).toContain("combo_unavailable"); expect(refreshHits).toBe(0); expect(backupAuth).toEqual(["Bearer key-b"]); expect([...captured].join("\n")).not.toContain("xai-refreshed-secret"); diff --git a/tests/server-key-failover-e2e.test.ts b/tests/server-key-failover-e2e.test.ts index ab6fca7582..3d4cf5ff32 100644 --- a/tests/server-key-failover-e2e.test.ts +++ b/tests/server-key-failover-e2e.test.ts @@ -249,6 +249,50 @@ describe("server 429 key failover (end-to-end)", () => { } }); + test("routed 401 rotates to the next API key before abandoning the provider", async () => { + const seenAuth: string[] = []; + upstream = Bun.serve({ + hostname: "127.0.0.1", port: 0, + fetch(req) { + seenAuth.push(req.headers.get("authorization") ?? ""); + if (seenAuth.length === 1) { + return Response.json({ error: { type: "authentication_error", code: "invalid_api_key", message: "revoked key" } }, { status: 401 }); + } + return Response.json({ + id: "chatcmpl-auth-rotate", object: "chat.completion", + choices: [{ index: 0, message: { role: "assistant", content: "ok after auth rotate" }, finish_reason: "stop" }], + usage: { prompt_tokens: 3, completion_tokens: 2, total_tokens: 5 }, + }); + }, + }); + const config: OcxConfig = { + port: 0, hostname: "127.0.0.1", defaultProvider: "pooled", + providers: { + pooled: { + adapter: "openai-chat", baseUrl: `http://127.0.0.1:${upstream.port}/v1`, + allowPrivateNetwork: true, apiKey: "key-alpha-000111222333", + apiKeyPool: [ + { id: "k1", key: "key-alpha-000111222333", addedAt: 1 }, + { id: "k2", key: "key-beta-444555666777", addedAt: 2 }, + ], + }, + }, + } as OcxConfig; + saveConfig(config); + const server = startServer(0); + try { + const res = await fetch(new URL("/v1/responses", server.url), { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "pooled/some-model", input: "hello", stream: false }), + }); + expect(res.status).toBe(200); + expect(JSON.stringify(await res.json())).toContain("ok after auth rotate"); + expect(seenAuth).toEqual(["Bearer key-alpha-000111222333", "Bearer key-beta-444555666777"]); + } finally { + await server.stop(true); + } + }); + test("reasoning replay misses after a 429 rotates to a different physical key", async () => { const model = "reasoning-model"; const callId = "call_key_rotation";