mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-07-31 04:12:10 +03:00
* 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>
379 lines
16 KiB
TypeScript
379 lines
16 KiB
TypeScript
/**
|
|
* Early SSE keepalive wrapper for streaming route handlers.
|
|
*
|
|
* Strict HTTP clients (notably Codex CLI's `reqwest`, which has a ~5s idle-read
|
|
* timeout) drop the connection if no bytes arrive shortly after the request.
|
|
* OmniRoute, however, holds the streaming response until `ensureStreamReadiness`
|
|
* observes the upstream's first useful byte — which can exceed 5s for reasoning
|
|
* models that "think" before emitting any token (#2544). `curl` has no such
|
|
* idle timeout, so it was never affected, which is why the bug looked
|
|
* client-specific.
|
|
*
|
|
* This wrapper keeps the connection warm without disturbing the handler's
|
|
* internal logic (combo failover, stream readiness, account cooldown all still
|
|
* run inside the handler before it resolves):
|
|
*
|
|
* - Fast path: if the handler resolves within `thresholdMs`, its `Response`
|
|
* is returned verbatim — identical status, headers, and body. There is zero
|
|
* behavior change for normal latency, so metadata headers and non-200 error
|
|
* statuses are fully preserved for the common case.
|
|
*
|
|
* - Slow path: if the handler is still pending after `thresholdMs`, a 200
|
|
* `text/event-stream` response is opened immediately and SSE comment
|
|
* heartbeats are emitted every `intervalMs` until the handler resolves; its
|
|
* body is then forwarded. If the handler ultimately fails, a structured
|
|
* `event: error` frame is emitted in-band (the response is already committed
|
|
* to 200, so the HTTP status can no longer change).
|
|
*/
|
|
|
|
const ENCODER = new TextEncoder();
|
|
const KEEPALIVE_FRAME = ENCODER.encode(": omniroute-keepalive\n\n");
|
|
// OpenAI-compatible keepalive: a syntactically valid empty streaming chunk.
|
|
// Some OpenAI-compatible clients parse every non-empty SSE line as JSON and
|
|
// reject legal SSE comments before their first provider chunk arrives.
|
|
export const OPENAI_KEEPALIVE_FRAME = ENCODER.encode(
|
|
'data: {"id":"omniroute-keepalive","object":"chat.completion.chunk","created":0,"model":"omniroute","choices":[{"index":0,"delta":{},"finish_reason":null}]}\n\n'
|
|
);
|
|
// #7360 follow-up: the FIRST frame of the slow path carries visible content
|
|
// instead of an empty delta, framed as a reasoning/thinking chunk (the same
|
|
// shape OmniRoute already emits for real upstream reasoning — see
|
|
// open-sse/translator/response/claude-to-openai.ts's createChunk) so
|
|
// reasoning-aware clients render it instead of silently ignoring it. This is
|
|
// what lets OmniRoute safely wait out a longer Gemini rate-limit cooldown
|
|
// (see comboCooldownWait/waitForCooldown budgetMs) without the client's own
|
|
// first-event/idle-read timeout firing — the client sees a real byte
|
|
// immediately, it just says we're still working on it.
|
|
const STARTUP_THINKING_TEXT = "OmniRoute: got request, sending to provider";
|
|
export const OPENAI_STARTUP_THINKING_FRAME = ENCODER.encode(
|
|
`data: ${JSON.stringify({
|
|
id: "omniroute-keepalive",
|
|
object: "chat.completion.chunk",
|
|
created: 0,
|
|
model: "omniroute",
|
|
choices: [
|
|
{ index: 0, delta: { reasoning_content: STARTUP_THINKING_TEXT }, finish_reason: null },
|
|
],
|
|
})}\n\n`
|
|
);
|
|
// Anthropic Messages-format keepalive: a REAL `ping` SSE event, not a comment.
|
|
// Anthropic clients (Claude Code, the Anthropic SDK) reset their stream/first-token
|
|
// watchdog on real SSE events but ignore SSE comments (`: ...`), so on a slow first
|
|
// token the comment frame lets the client abort and retry the stream. Anthropic's own
|
|
// API emits `event: ping` for exactly this reason; the /v1/messages route mirrors it.
|
|
export const ANTHROPIC_PING_FRAME = ENCODER.encode('event: ping\ndata: {"type":"ping"}\n\n');
|
|
// Responses API keepalive: a self-contained, self-closed synthetic reasoning
|
|
// item (added -> summary_part.added -> text.delta -> summary_part.done),
|
|
// matching the abbreviated close pattern open-sse/utils/stream.ts's own
|
|
// emitSyntheticResponsesReasoningSummary already uses for real mid-stream
|
|
// reasoning. Closed within this one frame (not left dangling open) since the
|
|
// real upstream response — once it arrives — starts its own independent
|
|
// response.created lifecycle from scratch; this placeholder item never
|
|
// carries a response_id and isn't meant to be continued.
|
|
const RESPONSES_STARTUP_ITEM_ID = "rs_omniroute_keepalive";
|
|
export const RESPONSES_STARTUP_THINKING_FRAME = ENCODER.encode(
|
|
[
|
|
{
|
|
event: "response.output_item.added",
|
|
data: {
|
|
type: "response.output_item.added",
|
|
output_index: 0,
|
|
item: { id: RESPONSES_STARTUP_ITEM_ID, type: "reasoning", summary: [] },
|
|
},
|
|
},
|
|
{
|
|
event: "response.reasoning_summary_part.added",
|
|
data: {
|
|
type: "response.reasoning_summary_part.added",
|
|
item_id: RESPONSES_STARTUP_ITEM_ID,
|
|
output_index: 0,
|
|
summary_index: 0,
|
|
part: { type: "summary_text", text: "" },
|
|
},
|
|
},
|
|
{
|
|
event: "response.reasoning_summary_text.delta",
|
|
data: {
|
|
type: "response.reasoning_summary_text.delta",
|
|
item_id: RESPONSES_STARTUP_ITEM_ID,
|
|
output_index: 0,
|
|
summary_index: 0,
|
|
delta: STARTUP_THINKING_TEXT,
|
|
},
|
|
},
|
|
{
|
|
event: "response.reasoning_summary_part.done",
|
|
data: {
|
|
type: "response.reasoning_summary_part.done",
|
|
item_id: RESPONSES_STARTUP_ITEM_ID,
|
|
output_index: 0,
|
|
summary_index: 0,
|
|
part: { type: "summary_text", text: STARTUP_THINKING_TEXT },
|
|
},
|
|
},
|
|
]
|
|
.map((e) => `event: ${e.event}\ndata: ${JSON.stringify(e.data)}\n\n`)
|
|
.join("")
|
|
);
|
|
// Anthropic Messages API default — Anthropic's own spec really does use a named
|
|
// `event: error` SSE frame, so this is correct there. It is WRONG for the OpenAI-
|
|
// format routes below: Chat Completions and Responses streaming never use the SSE
|
|
// `event:` field at all, only bare `data: {...}` lines — a naive line-based parser
|
|
// (the kind most OpenAI-compatible clients use, not a full EventSource) can silently
|
|
// drop an unrecognized `event:` line and/or desync on the `data:` line that follows,
|
|
// so this error would never surface to the client at all (log ids
|
|
// 1784465227489-a2cbc0 / 1784457764961-73 territory: a client that gives up with no
|
|
// visible reason). See OPENAI_CHAT_ERROR_FRAME / OPENAI_RESPONSES_ERROR_FRAME below
|
|
// for the per-format-correct alternatives.
|
|
const ERROR_FRAME = ENCODER.encode(
|
|
`event: error\ndata: ${JSON.stringify({
|
|
error: { message: "Upstream stream failed before completion.", type: "stream_error" },
|
|
})}\n\n`
|
|
);
|
|
// Chat Completions convention: a plain `data:` line, no `event:` field. This
|
|
// matches what the openai-node SDK's stream iterator actually checks for — it
|
|
// inspects each parsed chunk for a top-level `error` key regardless of any SSE
|
|
// event name (there isn't one to check, since real OpenAI chat completions
|
|
// streams never send `event:` lines).
|
|
export const OPENAI_CHAT_ERROR_FRAME = ENCODER.encode(
|
|
`data: ${JSON.stringify({
|
|
error: { message: "Upstream stream failed before completion.", type: "stream_error" },
|
|
})}\n\n`
|
|
);
|
|
// Responses API convention: also a plain `data:` line, but the discriminator is
|
|
// the `type` field INSIDE the JSON payload (matching every other Responses API
|
|
// event — response.output_text.delta, response.completed, etc.), not an SSE
|
|
// `event:` field.
|
|
export const OPENAI_RESPONSES_ERROR_FRAME = ENCODER.encode(
|
|
`data: ${JSON.stringify({
|
|
type: "error",
|
|
code: null,
|
|
message: "Upstream stream failed before completion.",
|
|
param: null,
|
|
})}\n\n`
|
|
);
|
|
|
|
export type EarlyStreamKeepaliveOptions = {
|
|
/** Wait this long for the handler before committing to a keepalive stream. */
|
|
thresholdMs?: number;
|
|
/** Keepalive cadence once committed (must stay under the client idle timeout). */
|
|
intervalMs?: number;
|
|
/** Client request signal — propagated so a client disconnect cancels the upstream read. */
|
|
signal?: AbortSignal | null;
|
|
/**
|
|
* Frame emitted on each keepalive tick. Defaults to an SSE comment
|
|
* (`: omniroute-keepalive`). Anthropic-format routes (/v1/messages) must pass
|
|
* `ANTHROPIC_PING_FRAME` instead, because Anthropic clients ignore SSE comments
|
|
* for their stream watchdog and only a real `event: ping` keeps them from aborting.
|
|
*/
|
|
keepaliveFrame?: Uint8Array;
|
|
/**
|
|
* Frame emitted ONCE, immediately, as the very first byte of the slow path —
|
|
* before the recurring `keepaliveFrame` ticks start. Defaults to
|
|
* `keepaliveFrame` when omitted (today's behavior, unchanged). Pass a
|
|
* content-bearing frame (e.g. `OPENAI_STARTUP_THINKING_FRAME`) so the client
|
|
* sees visible progress instead of an empty/no-op keepalive on the first byte.
|
|
*/
|
|
startupFrame?: Uint8Array;
|
|
/** Extra headers to include in the keepalive response (e.g. X-Correlation-Id). */
|
|
extraHeaders?: Record<string, string>;
|
|
/**
|
|
* Frame emitted if the handler ultimately fails (or the upstream stream dies
|
|
* mid-flight with zero bytes forwarded) after the slow path has already
|
|
* committed to HTTP 200. Defaults to the Anthropic-style `event: error` frame
|
|
* (correct for /v1/messages). OpenAI-format routes (/v1/chat/completions,
|
|
* /v1/responses) MUST pass OPENAI_CHAT_ERROR_FRAME / OPENAI_RESPONSES_ERROR_FRAME
|
|
* instead — see the doc comment on the default ERROR_FRAME above for why.
|
|
*/
|
|
errorFrame?: Uint8Array;
|
|
};
|
|
|
|
type SettledHandler = { ok: true; response: Response } | { ok: false; error: unknown };
|
|
|
|
export async function withEarlyStreamKeepalive(
|
|
handlerPromise: Promise<Response>,
|
|
options: EarlyStreamKeepaliveOptions = {}
|
|
): Promise<Response> {
|
|
const thresholdMs = Math.max(0, options.thresholdMs ?? 2_000);
|
|
const intervalMs = Math.max(250, options.intervalMs ?? 2_500);
|
|
const signal = options.signal ?? null;
|
|
const keepaliveFrame = options.keepaliveFrame ?? KEEPALIVE_FRAME;
|
|
const startupFrame = options.startupFrame ?? keepaliveFrame;
|
|
const extraHeaders = options.extraHeaders ?? {};
|
|
const errorFrame = options.errorFrame ?? ERROR_FRAME;
|
|
// Single source of truth for whether THIS route's error framing uses a named SSE
|
|
// `event: error` line (Anthropic) or a plain `data:` line (OpenAI Chat Completions /
|
|
// Responses) — derived from errorFrame itself so the dynamic real-upstream-body case
|
|
// below stays consistent with the static default-message case without a second option.
|
|
const errorFrameUsesNamedEvent = new TextDecoder().decode(errorFrame).startsWith("event:");
|
|
|
|
// Settle into a tagged result so neither race branch leaves an unhandled
|
|
// rejection when the threshold timer wins.
|
|
const settled: Promise<SettledHandler> = handlerPromise.then(
|
|
(response) => ({ ok: true as const, response }),
|
|
(error) => ({ ok: false as const, error })
|
|
);
|
|
|
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
|
const raced = await Promise.race([
|
|
settled.then((result) => ({ kind: "settled" as const, result })),
|
|
new Promise<{ kind: "timeout" }>((resolve) => {
|
|
timer = setTimeout(() => resolve({ kind: "timeout" }), thresholdMs);
|
|
}),
|
|
]);
|
|
if (timer) clearTimeout(timer);
|
|
|
|
if (raced.kind === "settled") {
|
|
// Fast path — return verbatim, or rethrow so the route's normal error handling runs.
|
|
if (raced.result.ok) return raced.result.response;
|
|
throw raced.result.error;
|
|
}
|
|
|
|
// Slow path — open the SSE stream now and keep it warm until the handler resolves.
|
|
// Cleanup state is hoisted so both start() and cancel() (client disconnect) can stop
|
|
// the keepalive loop and cancel the upstream read.
|
|
let stopKeepalive = () => {};
|
|
let upstreamReader: ReadableStreamDefaultReader<Uint8Array> | null = null;
|
|
let aborted = false;
|
|
|
|
const stream = new ReadableStream<Uint8Array>({
|
|
async start(controller) {
|
|
let stopped = false;
|
|
const interval = setInterval(() => {
|
|
if (stopped) return;
|
|
try {
|
|
controller.enqueue(keepaliveFrame);
|
|
} catch {
|
|
stopped = true;
|
|
clearInterval(interval);
|
|
}
|
|
}, intervalMs);
|
|
if (interval && typeof interval === "object" && "unref" in interval) {
|
|
interval.unref?.();
|
|
}
|
|
// First frame immediately on commit so the client sees a byte right away.
|
|
// Use `startupFrame` (e.g. OPENAI_STARTUP_THINKING_FRAME / ANTHROPIC_PING_FRAME)
|
|
// — an SSE comment here would be ignored by Anthropic clients' watchdog on a
|
|
// sub-interval gap, defeating the keepalive for exactly the case it targets.
|
|
try {
|
|
controller.enqueue(startupFrame);
|
|
} catch {
|
|
/* consumer already gone */
|
|
}
|
|
|
|
stopKeepalive = () => {
|
|
stopped = true;
|
|
clearInterval(interval);
|
|
};
|
|
|
|
const onAbort = () => {
|
|
if (aborted) return;
|
|
aborted = true;
|
|
stopKeepalive();
|
|
upstreamReader?.cancel().catch(() => {});
|
|
try {
|
|
controller.close();
|
|
} catch {
|
|
/* already closed */
|
|
}
|
|
};
|
|
signal?.addEventListener("abort", onAbort, { once: true });
|
|
// addEventListener does not replay an abort that happened before registration.
|
|
// Checking after registration closes that gap without missing a concurrent abort.
|
|
if (signal?.aborted) onAbort();
|
|
|
|
try {
|
|
const result = await settled;
|
|
stopKeepalive();
|
|
if (aborted) {
|
|
// The synthetic keepalive response can be cancelled before the handler resolves.
|
|
// Cancel the eventual real response so its upstream work and lifecycle hooks finish.
|
|
if (result.ok && result.response.body) {
|
|
await result.response.body.cancel().catch(() => undefined);
|
|
}
|
|
return;
|
|
}
|
|
|
|
if (!result.ok) {
|
|
// Handler rejected — emit a generic error frame (never the raw error/stack).
|
|
controller.enqueue(errorFrame);
|
|
} else {
|
|
const response = result.response;
|
|
const contentType = (response.headers.get("content-type") || "").toLowerCase();
|
|
const isSse = contentType.includes("text/event-stream");
|
|
|
|
if (response.body && isSse) {
|
|
// Real SSE stream — forward it verbatim.
|
|
upstreamReader = response.body.getReader();
|
|
let bytesForwarded = 0;
|
|
try {
|
|
while (true) {
|
|
const { done, value } = await upstreamReader.read();
|
|
if (done) break;
|
|
if (value) {
|
|
controller.enqueue(value);
|
|
bytesForwarded += value.byteLength;
|
|
}
|
|
}
|
|
} catch (readErr) {
|
|
// Upstream stream failed mid-flight. Only emit an error frame if
|
|
// NO content was forwarded yet — otherwise the client already
|
|
// received partial content and a late error frame would corrupt
|
|
// the SSE stream. Silently close instead; the client will see
|
|
// the stream end naturally.
|
|
if (bytesForwarded === 0) {
|
|
controller.enqueue(errorFrame);
|
|
}
|
|
}
|
|
} else {
|
|
// Non-SSE response (e.g. a JSON error) reached us after we already
|
|
// committed to a 200 event-stream, so the HTTP status can no longer
|
|
// change. Frame the (already-sanitized) body as an in-band error event
|
|
// instead of forwarding raw JSON, which would be malformed SSE.
|
|
const text = response.body ? await response.text().catch(() => "") : "";
|
|
const dataLine =
|
|
text.trim() ||
|
|
JSON.stringify({ error: { message: "stream_error", type: "stream_error" } });
|
|
const framed = errorFrameUsesNamedEvent
|
|
? `event: error\ndata: ${dataLine}\n\n`
|
|
: `data: ${dataLine}\n\n`;
|
|
controller.enqueue(ENCODER.encode(framed));
|
|
}
|
|
}
|
|
} catch {
|
|
// Defensive: never surface a raw error/stack to the client.
|
|
if (!aborted) {
|
|
try {
|
|
controller.enqueue(errorFrame);
|
|
} catch {
|
|
/* consumer gone */
|
|
}
|
|
}
|
|
} finally {
|
|
stopKeepalive();
|
|
signal?.removeEventListener("abort", onAbort);
|
|
try {
|
|
controller.close();
|
|
} catch {
|
|
/* already closed */
|
|
}
|
|
}
|
|
},
|
|
cancel() {
|
|
// Consumer (Next.js → client) went away — stop keepalives and release the upstream.
|
|
aborted = true;
|
|
stopKeepalive();
|
|
upstreamReader?.cancel().catch(() => {});
|
|
},
|
|
});
|
|
|
|
return new Response(stream, {
|
|
status: 200,
|
|
headers: {
|
|
"Content-Type": "text/event-stream; charset=utf-8",
|
|
"Cache-Control": "no-cache, no-transform",
|
|
Connection: "keep-alive",
|
|
...extraHeaders,
|
|
},
|
|
});
|
|
}
|