Files
OmniRoute/open-sse/utils/error.ts
Markus Hartung cc17b304ab fix(sse): Gemini TPM/RPD quota classification + combo cooldown-wait resilience (#8213)
* fix(sse): Gemini TPM classification, combo-cooldown-wait for auto/quota-share, and target-timeout floor

Gemini TPM/RPM 429s were misclassified as QUOTA_EXHAUSTED because
sanitizeErrorMessage() truncates to the first line, hiding Google's
metric name and retry hint on lines 2-3. Added a rawMessage field
(internal-only, never reaches the client) and classifyGeminiQuotaMetricFromText()
to classify from the untruncated text, reordered ahead of the generic
credits/daily-quota checks.

Widened comboCooldownWaitEnabled (wait out a short transient cooldown
instead of crystallizing a 429/503) from quota-share-only to also cover
auto-strategy combos, and raised the wait ceiling to 65s/130s-budget/90s-cap
to match Gemini's ~60s TPM/RPM windows.

The per-target timeout (DEFAULT_COMBO_TARGET_TIMEOUT_MS, 120s) was shorter
than the new 130s cooldown-wait budget, so a target could get cut off
mid-wait with a synthetic 524 instead of completing the retry. Added
resolveComboTargetTimeoutMsForCombo()/isComboCooldownWaitEligible() in
comboConfig.ts to raise the per-target floor to budgetMs+buffer only for
wait-eligible strategies (auto/quota-share), verified live: a 12-request
concurrent burst against a TPM-exhausted combo went from 2/12 succeeding
(10 x 524) to 12/12 succeeding with zero 503/524.

Also: liveGeminiShared.ts's sendAndValidate now fails fast on a 503
instead of retrying past it, and the health dashboard + request logger
surface TPM stats alongside RPM/RPD.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo-exhausted rejection logs now capture request body + attempted models

recordRejectedRequestUsage() (the fast path for combo requests that never
reach handleChatCore, e.g. all targets locked by resilience cooldown)
hardcoded provider: "-" and never passed a request body to saveCallLog(),
so /dashboard/logs entries for these failures were nearly useless for
debugging: no way to see the client's request or which models were tried.

- recordRejectedRequestUsage() now accepts requestBody and persists it
  through the existing saveCallLog() artifact mechanism (same path
  handleChatCore's own logging uses).
- Added summarizeComboAttemptedModels(), which reads the combo's own model
  list (always available, unlike the response's combo-diagnostics headers —
  a model-level resilience-lockout skip never touches the
  exhaustedProviders/exhaustedConnections sets those headers are built
  from) to populate a real "provider" value instead of "-".
- Wired both into the call site in src/sse/handlers/chat.ts.

NOTE: unrelated to the Gemini TPM/combo-cooldown-wait fix on this branch —
landed here per operator request, to be split into its own branch/PR.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* feat(sse): synthetic streaming keep-alive event + 5-minute Gemini cooldown-wait ceiling

Many clients enforce a first-SSE-byte timeout, which made it unsafe to wait
out a longer upstream rate-limit cooldown on a streaming request — the
client would abandon the connection before any bytes arrived. This landed
in two parts:

1. Synthetic startup "thinking" event (OpenAI chat/completions format):
   the already-existing withEarlyStreamKeepalive wrapper (open-sse/utils/
   earlyStreamKeepalive.ts, wired into /v1/chat/completions, /v1/messages,
   /v1/responses since #2544) opens the SSE stream immediately once a
   request runs past its threshold, but only ever sent empty/no-op
   keepalive frames. Added a `startupFrame` option (defaults to
   `keepaliveFrame` — zero behavior change unless a route opts in) so the
   very first frame can carry real content instead. Wired
   OPENAI_STARTUP_THINKING_FRAME (a reasoning_content delta: "OmniRoute:
   got request, sending to provider") into /v1/chat/completions only —
   Claude Messages and Responses API formats both require a preceding
   envelope event (message_start / response.created) that a synthetic
   pre-dispatch frame can't safely fabricate without risking a duplicate
   envelope once the real stream arrives, so those two routes keep their
   existing (safe, proven) keepalive frames unchanged.

2. Raised the "wait out a known cooldown, then retry" ceiling to 5 minutes
   for both retry mechanisms, now that a client-side first-byte timeout is
   no longer a risk on the (opted-in) route:
   - comboCooldownWait (auto/quota-share combos, open-sse/services/combo.ts):
     maxWaitMs hard clamp raised 90s -> 300s (src/lib/resilience/settings/
     normalize.ts); defaults raised to maxWaitMs:90s/maxAttempts:5/
     budgetMs:300s. comboConfig.ts's resolveComboTargetTimeoutMsForCombo
     already derives the per-target timeout floor from budgetMs, so it
     tracks the new ceiling with no further changes.
   - waitForCooldown (direct, non-combo model requests, src/sse/handlers/
     chat.ts): this mechanism had NO cumulative cap before — only a
     per-wait cap (maxRetryWaitMs) and a retry count (maxRetries), so
     maxRetries x maxRetryWaitMs could exceed 5 minutes with no ceiling.
     Added a budgetMs field (mirrors comboCooldownWait) to
     WaitForCooldownSettings/CooldownAwareRetrySettings, threaded a
     requestRetryBudgetLeftMs tracker through chat.ts's requestAttemptLoop
     (mirrors combo.ts's comboCooldownBudgetLeftMs), and made
     getCooldownAwareRetryDecision refuse to wait once the cumulative
     budget is exhausted even if the single wait is under maxRetryWaitMs.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): extend the synthetic keep-alive thinking event to /v1/responses

Live incident (OpenClaw, log id 1784407081908-cbc24f): a /v1/responses
request to gemini/gemma-4-31b-it took 56s to produce a first byte and the
client disconnected (499 request_signal_aborted) — the same client-first-byte-
timeout problem the previous commit fixed for /v1/chat/completions, but
/v1/responses only had the generic bare-comment keepalive (no content), so it
wasn't covered.

Added RESPONSES_STARTUP_THINKING_FRAME: a self-contained synthetic reasoning
item (response.output_item.added -> reasoning_summary_part.added ->
reasoning_summary_text.delta -> reasoning_summary_part.done), opened AND
closed within this one frame rather than left dangling — it never carries a
response_id, so it can't collide with the real upstream response's own
independent response.created lifecycle that follows. Mirrors the abbreviated
delta+part.done close pattern open-sse/utils/stream.ts's own
emitSyntheticResponsesReasoningSummary already uses for real mid-stream
reasoning content.

Wired into src/app/api/v1/responses/route.ts via the startupFrame option
added in the previous commit.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo cooldown-wait vars reset every setTry, crystallizing a bogus 503 instead of waiting

Live incident (log id 1784416706646-51): a request to the "default" combo
(strategy=auto, maxSetRetries=3) hit a real Gemini TPM 429 on both gemma-4
targets, correctly classified as a short 40s rate_limit lockout — then
crystallized a 503 "all upstream accounts are inactive" in 6.9s instead of
ever reaching the cooldown-aware wait.

Root cause: `lastError`/`earliestRetryAfter`/`lastStatus` were declared with
`let` INSIDE the `for (setTry...)` loop body, so they reset to null at the
start of every set-try. When both targets lock out on setTry 0, every
subsequent setTry (1..maxSetRetries) pre-skips both targets via the
isModelLocked check with no real dispatch — so on the FINAL setTry (the only
one whose values the post-loop decision reads, since it's gated behind
`if (setTry < maxSetRetries) continue`), lastStatus was null, hitting the
"!lastStatus" branch (ALL_ACCOUNTS_INACTIVE 503) and completely bypassing the
comboCooldownWaitEnabled / earliestRetryAfter wait logic — even though a
real 429 with a known ~40s retry-after WAS observed on setTry 0.

This bug predates today's Gemini TPM work (any combo with maxSetRetries > 0
whose targets all lock out on the first pass was affected) but was masked in
existing tests: the "auto strategy (2 models...)" regression test uses
maxSetRetries: 0, so it only ever runs ONE setTry iteration and never
exercises the reset-on-retry path. It also explains why the dedicated
12-concurrent-request burst test passed cleanly — with concurrent requests,
timing variance meant some request's FINAL setTry iteration still had a live
target to dispatch to, giving lastStatus/earliestRetryAfter fresh data. A
single isolated request has no such luck.

Fix: hoist lastError/earliestRetryAfter/lastStatus to just inside
dispatchWithCooldownRetry, before the setTry loop, so they persist across
set-tries (still reset fresh on each recursive dispatchWithCooldownRetry()
call after a wait, which is correct). recordedAttempts/fallbackCount/
exhaustedProviders etc. are intentionally left per-iteration (unrelated to
this bug).

New regression test in tests/unit/combo-quota-share-cooldown-wait.test.ts
reproduces the exact live scenario (2 targets, both lock out on setTry 0,
maxSetRetries: 3) — confirmed red (503) against the pre-fix code, green
(200, waits and retries) against the fix.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* test(sse): extend live Gemini workload to Responses API + add large-context TPM test

Two additions to the live Gemini test suite, both live-verified against the
dev instance:

1. sendAndValidate() (tests/integration/liveGeminiShared.ts) now accepts an
   apiFormat: "chat" | "responses" parameter, building the Responses-API
   request shape (input array, max_output_tokens) and parsing its SSE events
   (response.output_text.delta / response.reasoning_summary_text.delta /
   response.completed) via the new readResponsesSSEStream(). Wired into two
   new tests in live-gemini-workload.test.ts ([30]/[31]), mirroring the
   existing Chat Completions streaming coverage. Verified live: 24/25 + 5/5
   payloads succeeded end-to-end through the new code path (the one failure
   was a ~300s test-client fetch timeout unrelated to the Responses API code
   itself — a separate, not-yet-addressed test-harness limitation).

2. genHugeContextMessage() builds a single message large enough (~4
   chars/token estimate) to approach or exceed Gemini's free-tier TPM ceiling
   (16000 input tokens/min for gemma-4) by itself. Every other prompt
   generator in this file tops out around 1-2k tokens — nowhere near that
   ceiling — so none of the existing workload tests ever exercised a REAL TPM
   429, only RPM-style rate limiting. tests/integration/gemini-large-context-tpm.test.ts
   sends two ~12-13k-token requests back-to-back (comfortably exceeding
   16000/min together) to exercise the full path against production Gemini:
   TPM classification, the comboCooldownWait retry, and the synthetic
   keep-alive frame on a genuinely slow request.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): abandoned combo target dispatch now observes its own per-target timeout, fixing a permanent "pending" dashboard leak

Live incident (dashboard log id 1784418258231-14961a, reported as "an ongoing
request even though there's already a 200"): a combo target dispatch
abandoned by comboTargetTimeoutMs (open-sse/services/combo/targetTimeoutRunner.ts)
left a permanent phantom "pending" entry in the dashboard, even after the
overall combo request had already succeeded via a different retry.

Root cause: chatCore.ts's createStreamController — and everything downstream
that depends on it (withRateLimit's Promise.race against Bottleneck,
acquireAccountSemaphore) — only ever watches clientRawRequest.signal, which
is the ORIGINAL client's request signal (set once via buildClientRawRequest
and reused unchanged across every target dispatch in a combo). It has no
connection to targetTimeoutRunner.ts's OWN AbortController
(target.modelAbortSignal), which is what actually fires when
comboTargetTimeoutMs (300s) elapses. src/sse/handlers/chat.ts's
handleSingleModel bridge between combo.ts and handleSingleModelChat received
`target.modelAbortSignal` but silently dropped it — never forwarded it
anywhere. So when a target got abandoned (e.g. stuck inside a wedged
Bottleneck rate-limiter queue, see the WEDGED force-reset log line from the
same incident), its per-target timeout fired and let the COMBO move on and
retry successfully elsewhere — but the abandoned dispatch's own promise
chain never learned it had been superseded, so it hung forever waiting on a
signal that was never going to fire, and trackPendingRequest(false) (the
finalize call) never ran.

Fix: thread target.modelAbortSignal through as a new modelAbortSignal
runtimeOption, and merge it into clientRawRequest.signal (via the existing
mergeAbortSignals helper from open-sse/executors/base.ts) right before
dispatch, so an abandoned target's own promise chain now observes its abort
and can reach its cleanup path — new resolveDispatchClientRawRequest() makes
this mechanically testable in isolation.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo cooldown-wait state recording, rate-limit wedge recovery, OpenAI-format SSE error frames

Five related fixes surfaced by live incidents (dashboard log ids 1784457764961-73,
1784465227489-a2cbc0, 1784504040241-6f8b9a) while validating the Gemini TPM/cooldown-wait
work on this branch against real OpenClaw traffic:

- combo.ts: the model-lockout bail-out branches in dispatchWithCooldownRetry never
  recorded lastStatus, so once every target in a set hit an existing lockout the final
  check crystallized a bogus ALL_ACCOUNTS_INACTIVE 503 instead of reaching the
  cooldown-wait decision, even with a real 429 + short retry-after observed.
- combo.ts/combo/types.ts: the "all credentials cooling down" pre-dispatch rejection
  (buildModelCooldownBody) nests its retry hint as error.retry_after/reset_seconds, not
  the top-level retryAfter every other 429 shape uses — combo's extraction only read the
  latter, so earliestRetryAfter stayed null for this shape even after lastStatus was fixed.
- rateLimitManager.ts: the wedge-recovery watchdog used disconnect(), which releases the
  heartbeat timer but never rejects jobs already QUEUED on that instance — orphaned
  dispatches hung until the outer ~300s per-target timeout, well past real clients'
  patience. Switched to stop({ dropWaitingJobs: true }), safe because the wedge condition
  already requires RUNNING===0 && EXECUTING===0.
- earlyStreamKeepalive.ts: the in-band error frame emitted after committing to a 200 SSE
  stream was hardcoded to Anthropic's `event: error` convention for every route, including
  the OpenAI-format ones (/v1/chat/completions, /v1/responses) where that framing is
  either invisible or malformed to a plain data-line parser. Added per-route
  OPENAI_CHAT_ERROR_FRAME / OPENAI_RESPONSES_ERROR_FRAME and wired them in.
- chatCore.ts: persisted a synthetic clientResponse error body even when the client had
  already disconnected (AbortError) before that body was ever computed — misleading the
  dashboard into showing "what the client received" for a response that was never sent.

Also: RequestLoggerDetail.tsx — Provider/Client Event Stream panes lost their collapse
toggle when StreamSection replaced the collapsible PayloadSection (692d6be80, unifying
active/finished request views) without carrying the toggle over.

Each fix has a TDD regression test with a confirmed red-before-green cycle.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* test(sse): free-tier model + gemma-4 TPM-ceiling benchmark harness

Adds a live benchmark comparing free models OmniRoute exposes across
configured providers plus previously-unexercised no-auth providers
(felo-web, aihorde, opencode, duckduckgo-web — none need a connection
row, they were just never tried). Reuses liveGeminiShared.ts's SSE
parsers and CASE_BUILDERS instead of duplicating them.

Also adds a targeted TPM-stress test firing back-to-back large-context
prompts at the gemma-4-31b model across its 3 free hosts (Gemini,
NVIDIA, AI Horde) to isolate whether the documented 16k-tokens/minute
free-tier ceiling is Gemini-specific enforcement or an inherent
model property.

FORCE_TOOL_CHOICE_REQUIRED is a test-only, default-off env flag added
to liveGeminiShared.ts and live-gemini-agentic-loop.test.ts for an
earlier live A/B comparison of tool_choice: required vs unset — kept
as a reusable knob for future runs.

Co-Authored-By: Markus Hartung <markus.hartung@gmail.com>

* test(sse): benchmark for the 2026-07-22 newly-enabled provider batch

Adds NEWLY_ENABLED_MODELS to freeModelBenchmarkShared.ts (Mistral
Leanstral, OpenRouter's live "free"-tagged roster, OpenCode Zen's
current free models — refetched live from
https://opencode.ai/zen/v1/models since the static catalog had
drifted) and a dedicated workload benchmark test for them.

Co-Authored-By: Markus Hartung <markus.hartung@gmail.com>

* test(sse): sync geminiRateLimitTracker tests with e74a1722b's corrected Gemma 4 limits

e74a1722b updated geminiRateLimits.json's gemma-4-* entries from the stale
15/1500/-1 (rpm/rpd/tpm) to the real published free-tier values
16000/14400/16000, but never updated the tests asserting the old numbers.
Surfaced by running the full test:unit suite as a post-rebase sanity check.

Co-Authored-By: Markus Hartung <markus.hartream@gmail.com>

* chore(quality): file-size baseline for own-growth (#8213)

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

---------

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>
Co-authored-by: Markus Hartung <markus.hartream@gmail.com>
Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>
2026-07-23 05:17:35 -03:00

744 lines
26 KiB
TypeScript

import { CORS_HEADERS } from "./cors.ts";
import { unwrapClinepassEnvelope } from "./clinepassEnvelope.ts";
import { getDefaultErrorMessage, getErrorInfo } from "../config/errorConfig.ts";
import { normalizePayloadForLog } from "@/lib/logPayloads";
import type { ModelCooldownErrorPayload } from "@/types";
/**
* Sanitize an error message to prevent stack trace exposure in API responses.
* Strips stack traces, file paths, and absolute Windows/POSIX paths from
* error messages before they reach the client.
*/
interface ErrorResponseBody {
error: {
message: string;
type?: string;
code?: string;
};
upstream_details?: Record<string, unknown> | null; // sanitized upstream provider body
}
// Length cap protects against pathological inputs even before tokenization.
const MAX_ERROR_LEN = 4096;
const SOURCE_EXT = ["ts", "tsx", "js", "jsx", "mjs", "cjs"] as const;
function looksLikeAbsolutePath(tok: string): boolean {
// POSIX: "/<...>.ts" (optionally followed by :line[:col]).
// Windows: "C:\<...>.ts" or "C:/<...>.ts".
if (tok.length < 4 || tok.length > 2048) return false;
const isPosix = tok.charCodeAt(0) === 0x2f; // '/'
const isWindows = tok.length > 2 && tok.charCodeAt(1) === 0x3a && /[A-Za-z]/.test(tok[0]);
if (!isPosix && !isWindows) return false;
const dot = tok.lastIndexOf(".");
if (dot <= 0 || dot === tok.length - 1) return false;
const ext = tok
.slice(dot + 1)
.split(":", 1)[0]
.toLowerCase();
return (SOURCE_EXT as readonly string[]).includes(ext);
}
/**
* Strip stack-trace tail and absolute source paths from error messages.
*
* Implemented via simple whitespace tokenization (linear time) instead of a
* single complex regex, so CodeQL `js/polynomial-redos` stays clean even when
* the runtime error message is attacker-controlled.
*/
export function sanitizeErrorMessage(message: unknown): string {
let str = typeof message === "string" ? message : String(message ?? "");
if (str.length > MAX_ERROR_LEN) str = str.slice(0, MAX_ERROR_LEN);
const nl = str.indexOf("\n");
const firstLine = nl >= 0 ? str.slice(0, nl) : str;
// Preserve original whitespace by splitting on captured separator.
const parts = firstLine.split(/(\s+)/);
for (let i = 0; i < parts.length; i++) {
if (looksLikeAbsolutePath(parts[i])) parts[i] = "<path>";
}
return parts.join("");
}
const BLOCKED_KEYS = /stack|trace|path|file|cwd|dir|password|secret|token|key/i;
const MAX_DEPTH = 4;
/**
* Recursively sanitize an arbitrary JSON value from an upstream provider body.
* - Strings: run through sanitizeErrorMessage (strips stacks + absolute paths).
* - Keys matching BLOCKED_KEYS are dropped (credential/path guards).
* - Depth capped at MAX_DEPTH to prevent pathological nesting.
* - Arrays capped at 32 elements.
* - Returns null for null/undefined/non-JSON-serializable values.
*/
export function sanitizeUpstreamDetails(value: unknown, depth = 0): unknown {
if (depth > MAX_DEPTH) return "[truncated]";
if (value === null || value === undefined) return null;
if (typeof value === "string") return sanitizeErrorMessage(value);
if (typeof value === "number" || typeof value === "boolean") return value;
if (Array.isArray(value)) {
return value.slice(0, 32).map((v) => sanitizeUpstreamDetails(v, depth + 1));
}
if (typeof value === "object") {
const out: Record<string, unknown> = {};
for (const [k, v] of Object.entries(value as Record<string, unknown>)) {
if (BLOCKED_KEYS.test(k)) continue;
out[k] = sanitizeUpstreamDetails(v, depth + 1);
}
return out;
}
return null;
}
/**
* Build OpenAI-compatible error response body. Message is always sanitized
* so callers do not need to remember to strip stack traces themselves.
* Optional third argument `upstreamDetails` (raw parsed provider body) is
* sanitized by sanitizeUpstreamDetails before inclusion as `upstream_details`.
*/
export function buildErrorBody(
statusCode: number,
message: string,
upstreamDetails?: unknown
): ErrorResponseBody {
const errorInfo = getErrorInfo(statusCode);
const safeMessage = sanitizeErrorMessage(message) || getDefaultErrorMessage(statusCode);
const body: ErrorResponseBody = {
error: {
message: safeMessage,
type: errorInfo.type,
code: errorInfo.code,
},
};
if (upstreamDetails !== undefined && upstreamDetails !== null) {
const sanitized = sanitizeUpstreamDetails(upstreamDetails);
if (sanitized !== null && typeof sanitized === "object" && !Array.isArray(sanitized)) {
body.upstream_details = sanitized as Record<string, unknown>;
}
}
return body;
}
/**
* Sanitized auto-combo diagnostic trace surfaced on a combo terminal failure.
* Contains ONLY provider/model ids, enumerated reason codes, and counts — never
* keys, tokens, cookies, credentials, or upstream bodies. Fields are length- and
* count-capped so the projection is safe to place in HTTP headers too. (QA P0:
* "Add a sanitized combo diagnostic trace … candidate pool count, excluded
* provider/model reasons, selected attempt order, terminal failure summary.")
*/
export interface ComboExclusion {
provider: string;
model?: string;
reason: string;
}
/**
* Next-step suggestion surfaced when a combo cascade fails. Lets the client (e.g.
* the OpenCode plugin) auto-render an actionable hint in the TUI instead of an
* opaque "model stopped producing output" error — fixes the silent-stop pattern
* where the user has no way to recover a session without guessing. Whitelisted to
* a small set so the projection remains bounded.
*/
export type ComboRecoveryAction =
/** Cascade failed because every candidate is exhausted — try a different combo or `auto`. */
| "try-auto"
/** Upstream asks to retry after a cooldown window — wait, then retry the same combo. */
| "wait"
/** Transient failure (network, 5xx) — retry the same combo immediately. */
| "retry"
/** Cascade used every account of every provider — switch to a different combo entirely. */
| "switch-combo";
export interface ComboRecoveryHint {
/** Machine-readable action verb — consumed by clients to render a UI hint. */
action: ComboRecoveryAction;
/** Seconds the client should wait before retrying. Only meaningful when action="wait". */
retry_after_seconds?: number;
/** Human-readable next step — included verbatim in the error body for non-MCP clients. */
next_step: string;
}
export interface ComboExclusion {
provider: string;
model?: string;
reason: string;
}
export interface ComboDiagnostics {
poolSize: number;
attempted: number;
excluded: ComboExclusion[];
attemptOrder: Array<{ provider: string; model: string }>;
terminalReason: string;
/** Optional next-step hint — populated when the dispatcher can recommend a recovery action. */
recovery?: ComboRecoveryHint;
}
function clampDiagStr(v: unknown, max = 128): string {
return typeof v === "string" ? v.slice(0, max).replace(/[\r\n]+/g, " ") : "";
}
/**
* HTTP header values must be Latin1/ByteString (undici throws a TypeError
* otherwise — see #6612). Replace any codepoint outside the Latin1 range
* (0-255) with "?" so header construction never throws. Only used for the
* literal header value; the JSON body keeps the original, unsanitized
* readable text via `sanitizeComboDiagnostics`.
*/
function toHeaderSafeAscii(v: string): string {
let out = "";
for (let i = 0; i < v.length; i++) {
const code = v.charCodeAt(i);
out += code > 255 ? "?" : v[i];
}
return out;
}
/**
* Whitelist sanitizer for the recovery hint. The `action` enum is a closed set;
* `retry_after_seconds` is clamped to a non-negative integer ≤ 3600; `next_step` is
* capped and stripped of CR/LF (would break header parsing). Returns undefined when
* no usable input was supplied so downstream code can branch cleanly on absence.
*/
const RECOVERY_ACTIONS = new Set<ComboRecoveryAction>([
"try-auto",
"wait",
"retry",
"switch-combo",
]);
export function sanitizeRecoveryHint(
r: ComboRecoveryHint | null | undefined
): ComboRecoveryHint | undefined {
if (!r || typeof r !== "object") return undefined;
const action = typeof r.action === "string" ? (r.action as ComboRecoveryAction) : null;
if (!action || !RECOVERY_ACTIONS.has(action)) return undefined;
// Reject empty OR whitespace-only next_step — the value must render usefully as a
// header and as a body field. A whitespace-only string would print as a blank hint.
const next_step = clampDiagStr(r.next_step, 200).trim();
if (!next_step) return undefined;
const hint: ComboRecoveryHint = { action, next_step };
if (typeof r.retry_after_seconds === "number" && Number.isFinite(r.retry_after_seconds)) {
hint.retry_after_seconds = Math.max(0, Math.min(3600, Math.floor(r.retry_after_seconds)));
}
return hint;
}
/**
* Whitelist projection — guarantees only id/reason string primitives + integer
* counts can escape, regardless of what the caller assembled. This is the secret
* containment boundary for the diagnostic trace.
*/
export function sanitizeComboDiagnostics(d: ComboDiagnostics): ComboDiagnostics {
const recovery = sanitizeRecoveryHint(d?.recovery);
const out: ComboDiagnostics = {
poolSize: Number.isFinite(d?.poolSize) ? d.poolSize : 0,
attempted: Number.isFinite(d?.attempted) ? d.attempted : 0,
excluded: (d?.excluded ?? []).slice(0, 64).map((e) => ({
provider: clampDiagStr(e?.provider, 64),
...(e?.model ? { model: clampDiagStr(e.model, 96) } : {}),
reason: clampDiagStr(e?.reason, 64),
})),
attemptOrder: (d?.attemptOrder ?? [])
.slice(0, 64)
.map((a) => ({ provider: clampDiagStr(a?.provider, 64), model: clampDiagStr(a?.model, 96) })),
terminalReason: clampDiagStr(d?.terminalReason, 200),
};
if (recovery) out.recovery = recovery;
return out;
}
/**
* errorResponse variant that attaches a sanitized combo diagnostic trace as BOTH
* `x-omniroute-combo-*` headers and a `diagnostics` field in the OpenAI-shaped
* error body (extra field — backward-compatible with standard error parsers).
* `opts.code`/`opts.type` override the status-derived defaults (e.g. to preserve
* the `ALL_ACCOUNTS_INACTIVE` code on the 503 terminal path). When the diagnostic
* carries a `recovery` hint it is mirrored as `x-omniroute-recovery-action` /
* `x-omniroute-recovery-next-step` / `x-omniroute-retry-after-seconds` headers and as a
* top-level `recovery_hint` field on the body so non-header-aware clients (curl,
* MCP tools, log scrapers) can also pick it up.
*/
export function errorResponseWithComboDiagnostics(
statusCode: number,
message: string,
diagnostics: ComboDiagnostics,
opts: { code?: string; type?: string } = {}
): Response {
const safe = sanitizeComboDiagnostics(diagnostics);
const body = buildErrorBody(statusCode, message) as ErrorResponseBody & {
diagnostics?: ComboDiagnostics;
recovery_hint?: ComboRecoveryHint;
};
if (opts.code) body.error.code = opts.code;
if (opts.type) body.error.type = opts.type;
body.diagnostics = safe;
if (safe.recovery) body.recovery_hint = safe.recovery;
const excludedHeader = toHeaderSafeAscii(
safe.excluded
.map((e) => `${e.provider}${e.model ? `/${e.model}` : ""}:${e.reason}`)
.join(",")
.slice(0, 900)
);
const headers: Record<string, string> = {
"Content-Type": "application/json",
"x-omniroute-combo-pool-size": String(safe.poolSize),
"x-omniroute-combo-attempted": String(safe.attempted),
"x-omniroute-combo-excluded": excludedHeader,
"x-omniroute-combo-terminal-reason": toHeaderSafeAscii(safe.terminalReason.slice(0, 200)),
};
if (safe.recovery) {
headers["x-omniroute-recovery-action"] = safe.recovery.action;
// Header limit of 128 chars — keep next_step compact for fast parsing.
// The body field carries the full 200-char value for richer display.
headers["x-omniroute-recovery-next-step"] = toHeaderSafeAscii(safe.recovery.next_step).slice(
0,
128
);
if (
typeof safe.recovery.retry_after_seconds === "number" &&
safe.recovery.retry_after_seconds > 0
) {
headers["x-omniroute-retry-after-seconds"] = String(safe.recovery.retry_after_seconds);
}
}
return new Response(JSON.stringify(body), {
status: statusCode,
headers,
});
}
/**
* Create error Response object (for non-streaming)
* @param {number} statusCode - HTTP status code
* @param {string} message - Error message
* @returns {Response} HTTP Response object
*/
export function errorResponse(statusCode: number, message: string): Response {
return new Response(JSON.stringify(buildErrorBody(statusCode, sanitizeErrorMessage(message))), {
status: statusCode,
headers: {
"Content-Type": "application/json",
},
});
}
/**
* Write error to SSE stream (for streaming)
* @param {WritableStreamDefaultWriter} writer - Stream writer
* @param {number} statusCode - HTTP status code
* @param {string} message - Error message
*/
export async function writeStreamError(
writer: WritableStreamDefaultWriter<Uint8Array>,
statusCode: number,
message: string
): Promise<void> {
const errorBody = buildErrorBody(statusCode, sanitizeErrorMessage(message));
const encoder = new TextEncoder();
await writer.write(encoder.encode(`data: ${JSON.stringify(errorBody)}\n\n`));
}
function normalizeRetryAfterSeconds(retryAfter?: string | number | Date | null): number {
if (typeof retryAfter === "number" && Number.isFinite(retryAfter)) {
if (retryAfter > 0 && retryAfter < 1_000_000_000) {
return Math.max(Math.ceil(retryAfter), 1);
}
const retryTimeMs = new Date(retryAfter).getTime();
if (Number.isFinite(retryTimeMs)) {
return Math.max(Math.ceil((retryTimeMs - Date.now()) / 1000), 1);
}
}
if (retryAfter instanceof Date || typeof retryAfter === "string") {
const retryTimeMs = new Date(retryAfter).getTime();
if (Number.isFinite(retryTimeMs)) {
return Math.max(Math.ceil((retryTimeMs - Date.now()) / 1000), 1);
}
}
return 1;
}
/**
* Parse Antigravity error message to extract retry time
* Example: "You have exhausted your capacity on this model. Your quota will reset after 2h7m23s."
* @param {string} message - Error message
* @returns {number|null} Retry time in milliseconds, or null if not found
*/
export function parseAntigravityRetryTime(message: unknown): number | null {
if (typeof message !== "string") return null;
// Match patterns like: 2h7m23s, 5m30s, 45s, 1h20m, etc.
const match = message.match(/reset after (\d+h)?(\d+m)?(\d+s)?/i);
if (!match) return null;
let totalMs = 0;
// Extract hours
if (match[1]) {
const hours = parseInt(match[1]);
totalMs += hours * 60 * 60 * 1000;
}
// Extract minutes
if (match[2]) {
const minutes = parseInt(match[2]);
totalMs += minutes * 60 * 1000;
}
// Extract seconds
if (match[3]) {
const seconds = parseInt(match[3]);
totalMs += seconds * 1000;
}
return totalMs > 0 ? totalMs : null;
}
/**
* Parse upstream provider error response
* @param {Response} response - Fetch response from provider
* @param {string} provider - Provider name (for Antigravity-specific parsing)
* @returns {Promise<{statusCode: number, message: string, retryAfterMs: number|null, responseBody: unknown}>}
*/
export async function parseUpstreamError(response: Response, provider: string | null = null) {
let message: unknown = "";
let retryAfterMs: number | null = null;
let responseBody: unknown = null;
let errorCode: unknown = undefined;
let errorType: unknown = undefined;
try {
const text = await response.text();
responseBody = normalizePayloadForLog(text);
// Try parse as JSON
try {
const parsed = JSON.parse(text);
// Handle array responses (e.g., from some Gemini APIs)
const json = (Array.isArray(parsed) && parsed.length > 0 ? parsed[0] : parsed) || {};
// ClinePass wraps upstream errors in a {success:false, error} envelope.
// Extract the upstream error string (an upstream JSON field, not a local
// stack) — still routed through sanitizeErrorMessage/buildErrorBody by
// every consumer below (Rule #12).
const { error: clinepassEnvError } = unwrapClinepassEnvelope(json, provider);
message = clinepassEnvError
? clinepassEnvError.message
: json.error?.message || json.message || json.error || text;
errorCode = json.error?.code || json.code;
errorType = json.error?.type || json.type;
} catch {
message = text;
}
} catch {
message = `Upstream error: ${response.status}`;
responseBody = { _rawText: message };
}
const messageStr = typeof message === "string" ? message : JSON.stringify(message);
const retryAfterHeader = response.headers?.get?.("retry-after");
if (retryAfterHeader && !retryAfterMs) {
const retryAfterSec = Number.parseInt(retryAfterHeader, 10);
if (Number.isFinite(retryAfterSec) && retryAfterSec > 0) {
retryAfterMs = retryAfterSec * 1000;
} else {
const retryAfterDate = new Date(retryAfterHeader).getTime();
if (Number.isFinite(retryAfterDate) && retryAfterDate > Date.now()) {
retryAfterMs = retryAfterDate - Date.now();
}
}
}
// Parse Antigravity-specific retry time from error message
if (provider === "antigravity" && response.status === 429) {
retryAfterMs = parseAntigravityRetryTime(messageStr);
}
// Also parse retry time for other providers (Qwen, etc.) with "quota will reset after XhYmZs" format
if (response.status === 429 && !retryAfterMs) {
retryAfterMs = parseAntigravityRetryTime(messageStr);
}
// Generic providers: "Please retry after 20s"
if (response.status === 429 && !retryAfterMs) {
const retryMatch = messageStr.match(/retry\s+after\s+(\d+)\s*s/i);
if (retryMatch) {
retryAfterMs = Number.parseInt(retryMatch[1], 10) * 1000;
}
}
// Cap maximum retry time at 24 hours to prevent infinite wait
const MAX_RETRY_MS = 24 * 60 * 60 * 1000;
if (retryAfterMs && retryAfterMs > MAX_RETRY_MS) {
retryAfterMs = MAX_RETRY_MS;
}
const responseHeaders: Record<string, string> | null = response.headers
? Object.fromEntries(response.headers.entries())
: null;
return {
statusCode: response.status,
message: messageStr,
errorCode,
errorType,
retryAfterMs,
responseBody,
responseHeaders,
};
}
/**
* Create error result for chatCore handler
* @param {number} statusCode - HTTP status code
* @param {string} message - Error message
* @param {number|null} retryAfterMs - Optional retry-after time in milliseconds
* @returns {{ success: false, status: number, error: string, response: Response, retryAfterMs?: number }}
*/
export function createErrorResult(
statusCode: number,
message: string,
retryAfterMs: number | null = null,
errorCode?: string,
errorType?: string,
upstreamDetails?: unknown
) {
const body = buildErrorBody(statusCode, message, upstreamDetails);
if (errorCode) {
body.error.code = errorCode;
}
if (errorType) {
body.error.type = errorType;
}
const result: {
success: false;
status: number;
error: string;
/**
* #7360: the FULL, un-sanitized upstream message — `error` above is
* truncated to its first line by sanitizeErrorMessage() (correctly, for
* the client-facing response body). Server-side classification
* (checkFallbackError / Gemini TPM-vs-RPD metric detection) needs the
* complete multi-line text — e.g. Google's metric name and retry hint
* live on lines 2-3, after the generic "quota exceeded" preamble on
* line 1. This field NEVER reaches the HTTP response body (`response`
* below is already built from the sanitized `body`); it exists purely
* for internal callers that inspect the returned object.
*/
rawMessage: string;
errorType?: string;
errorCode?: string;
response: Response;
retryAfterMs?: number;
} = {
success: false,
status: statusCode,
error: body.error.message,
rawMessage: message,
errorType,
errorCode,
response: new Response(JSON.stringify(body), {
status: statusCode,
headers: { "Content-Type": "application/json" },
}),
};
// Add retryAfterMs if available (for Antigravity quota errors)
if (retryAfterMs) {
result.retryAfterMs = retryAfterMs;
}
return result;
}
/**
* Create unavailable response when all accounts are rate limited
* @param {number} statusCode - Original error status code
* @param {string} message - Error message (without retry info)
* @param {string} retryAfter - ISO timestamp when earliest account becomes available
* @param {string} retryAfterHuman - Human-readable retry info e.g. "reset after 30s"
* @returns {Response}
*/
export function unavailableResponse(
statusCode: number,
message: string,
retryAfter?: string | number | Date | null,
retryAfterHuman?: string
) {
const retryAfterSec = normalizeRetryAfterSeconds(retryAfter);
const msg = retryAfterHuman ? `${message} (${retryAfterHuman})` : message;
return new Response(JSON.stringify({ error: { message: msg } }), {
status: statusCode,
headers: {
"Content-Type": "application/json",
"Retry-After": String(retryAfterSec),
},
});
}
export function providerCircuitOpenResponse(
provider: string,
retryAfter?: string | number | Date | null
) {
const retryAfterSec = normalizeRetryAfterSeconds(retryAfter);
return new Response(
JSON.stringify({
error: {
message: `Provider ${provider} circuit breaker is open`,
type: "server_error",
code: "provider_circuit_open",
provider,
retry_after: retryAfterSec,
},
}),
{
status: 503,
headers: {
"Content-Type": "application/json",
"Retry-After": String(retryAfterSec),
"X-OmniRoute-Provider-Breaker": "open",
},
}
);
}
export function buildModelCooldownBody({
model,
retryAfterSec,
retryAfterAt,
credentialsCoolingCount,
}: {
model?: string | null;
retryAfterSec: number;
retryAfterAt?: string | null;
credentialsCoolingCount?: number | null;
}): ModelCooldownErrorPayload {
const resolvedModel = typeof model === "string" && model.trim().length > 0 ? model.trim() : null;
const resolvedRetryAfterAt =
typeof retryAfterAt === "string" && retryAfterAt.length > 0 ? retryAfterAt : null;
const resolvedCoolingCount =
typeof credentialsCoolingCount === "number" &&
Number.isFinite(credentialsCoolingCount) &&
credentialsCoolingCount > 0
? Math.floor(credentialsCoolingCount)
: null;
return {
error: {
message: resolvedModel
? `All credentials for model ${resolvedModel} are cooling down`
: "All credentials for the requested model are cooling down",
type: "rate_limit_error",
code: "model_cooldown",
...(resolvedModel ? { model: resolvedModel } : {}),
reset_seconds: Math.max(Math.ceil(retryAfterSec), 1),
...(resolvedRetryAfterAt ? { retry_after: resolvedRetryAfterAt } : {}),
...(resolvedCoolingCount ? { credentials_cooling: resolvedCoolingCount } : {}),
},
};
}
export function modelCooldownResponse({
model,
retryAfter,
retryAfterAt,
credentialsCoolingCount,
}: {
model?: string | null;
retryAfter?: string | number | Date | null;
retryAfterAt?: string | null;
credentialsCoolingCount?: number | null;
}) {
const retryAfterSec = normalizeRetryAfterSeconds(retryAfter);
const resolvedRetryAfterAt =
typeof retryAfterAt === "string" && retryAfterAt.length > 0
? retryAfterAt
: typeof retryAfter === "string" && retryAfter.length > 0
? retryAfter
: null;
return new Response(
JSON.stringify(
buildModelCooldownBody({
model,
retryAfterSec,
retryAfterAt: resolvedRetryAfterAt,
credentialsCoolingCount,
})
),
{
status: 429,
headers: {
"Content-Type": "application/json",
"Retry-After": String(retryAfterSec),
},
}
);
}
/**
* Build an executor-style error result (response + url + headers + transformedBody).
* Shared by web-cookie executors that return the `{ response, url, headers, transformedBody }` shape.
*/
export function makeExecutorErrorResult(
status: number,
message: string,
body: unknown,
url: string
) {
return {
response: new Response(
JSON.stringify({
error: {
message: sanitizeErrorMessage(message),
type: "upstream_error",
code: `HTTP_${status}`,
},
}),
{ status, headers: { "Content-Type": "application/json" } }
),
url,
headers: {} as Record<string, string>,
transformedBody: body,
};
}
/**
* Normalize a cookie string: strip a leading "Cookie:" prefix if present.
*/
export function normalizeCookie(raw: string): string {
return raw?.startsWith("Cookie:") ? raw.slice(7).trim() : raw || "";
}
/**
* Format provider error with context
* @param {Error} error - Original error
* @param {string} provider - Provider name
* @param {string} model - Model name
* @param {number|string} statusCode - HTTP status code or error code
* @returns {string} Formatted error message
*/
export function formatProviderError(
error: { code?: string | number; message?: string; cause?: unknown } | Error,
provider: string,
model: string,
statusCode?: string | number | null
): string {
const providerCode = "code" in error ? error.code : undefined;
const code = statusCode || providerCode || "FETCH_FAILED";
const message = error.message || "Unknown error";
// Expose low-level cause (e.g. UND_ERR_SOCKET, ECONNRESET, ETIMEDOUT) for diagnosing fetch failures
const cause = (error as { cause?: unknown }).cause;
const causeObj =
cause && typeof cause === "object" ? (cause as Record<string, unknown>) : undefined;
const causeCode = typeof causeObj?.code === "string" ? causeObj.code : undefined;
const causeMsg = typeof causeObj?.message === "string" ? causeObj.message : undefined;
const causeStr =
causeCode || causeMsg ? ` (cause: ${[causeCode, causeMsg].filter(Boolean).join(": ")})` : "";
return `[${code}]: ${message}${causeStr}`;
}