mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-07-31 12:22:14 +03:00
* chore(release): open v3.8.38 development cycle
* fix(executors): strip client_metadata for cerebras and mistral (#4727)
Integrated into release/v3.8.38 (leva 5)
* fix(codebuddy): only send reasoning params when client requests reasoning (#5019)
Integrated into release/v3.8.38 (leva 5)
* fix(sse): keep streaming for forceStream providers when client requests JSON (#5021)
Integrated into release/v3.8.38 (leva 5)
* fix(sse): guard non-JSON SSE lines and duplicate [DONE] (#4937)
Integrated into release/v3.8.38 (leva 5)
* feat(blackbox): refresh provider model catalog (#4935)
Integrated into release/v3.8.38 (leva 5)
* fix(sse): dedupe case-variant Anthropic version/beta headers (#4846)
Integrated into release/v3.8.38 (leva 5)
* feat(sse): Kiro inline <thinking> stream splitter (#4911)
Integrated into release/v3.8.38 (leva 5)
* feat(cursor): parse Composer DeepSeek-style inline tool calls (#4912)
Integrated into release/v3.8.38 (leva 5)
* feat(proxy): auth-less host:port batch import (#4938)
Integrated into release/v3.8.38 (leva 5)
* fix(oauth): support Kiro IDC (organization) token import (#4944)
Integrated into release/v3.8.38 (leva 5)
* fix(translator): preserve cache_control for DashScope OpenAI-compat providers (port from 9router#2069) (#5013)
Integrated into release/v3.8.38 (leva 5)
* fix(tts): resolve Gemini TTS models from catalog (#4934)
Integrated into release/v3.8.38 (leva 5)
* fix(sse): don't cool down the connection on a self-inflicted upstream timeout (504) (#5064)
Integrated into release/v3.8.38 (leva 5)
* fix(sse): robust Anthropic /v1/messages streaming — real ping keepalive + client-disconnect guard (#5063)
Integrated into release/v3.8.38 (leva 5)
* feat(video): add Alibaba DashScope (wan2.7-t2v) provider (#5051)
Integrated into release/v3.8.38 (leva 5)
* fix: preserve model hidden flags (isHidden) across model sync (#5086)
Integrated into release/v3.8.38 (leva 5)
* fix(models): derive model discovery config from registry modelsUrl (#5087)
Integrated into release/v3.8.38 (leva 5)
* fix(compression): replace fileURLToPath(import.meta.url) with runtime anchors for standalone bundle (#5089)
Integrated into release/v3.8.38 (leva 5)
* feat(cc): add summarized thinking display toggle (#5055)
Integrated into release/v3.8.38 (leva 5)
* Harden selected API error responses (#5032)
Integrated into release/v3.8.38 (leva 5)
* chore(quality): rebaseline file-size for leva 5 PR batch drift
6 frozen files grew from merged leva-5 PRs (cursor #4912, kiro #4911,
videoGeneration #5051, default #4727, base #4846, chat #5064); all covered
by per-PR tests. See _rebaseline_2026_06_26_leva5 in the baseline.
* feat(compression): compression playground (Play + Compare tabs) in the studio (#5080)
Integrated into release/v3.8.38
* fix(combo): fail over on empty-content 502 instead of exhausting the provider (#5085) (#5104)
* fix(dashboard): surface detailed credential-validation error in add-connection modal (#5088) (#5106)
* feat(providers): allow local/private provider URLs by default with scoped metadata-safe guard (#5066) (#5107)
* fix(diagnostics): treat non-streaming Claude messages shape as valid output (#5108) (#5116)
* fix(db): translate pt-BR SQLite driver-fallback log lines to English (#5103) (#5115)
* fix(sse): repair release base-reds — malformed-response false positives + header casing + stale tests (#5117)
Repairs the release/v3.8.38 base-reds; unblocks #5078.
* chore(quality): rebaseline file-size for responseSanitizer (#5117) + AddApiKeyModal drift
* fix(translator): forward image tool_result blocks as image_url (#5100)
Base-reds fixed (#5117); image tool_result→image_url. Integrated into release/v3.8.38.
* fix(responses): default text.format for openai-compatible responses providers (#5101)
Base-reds fixed (#5117); default text.format + file-size rebaseline. Integrated into release/v3.8.38.
* feat(dashboard): expose Fusion judgeModel + fusionTuning in the combo editor (#5074)
Base-reds fixed (#5117); Fusion editor + file-size rebaseline. Integrated into release/v3.8.38.
* feat(quota): add opt-in Codex/Claude auto-ping keepalive (#5102)
Base-reds fixed (#5117); auto-ping keepalive + file-size rebaseline. Integrated into release/v3.8.38.
* test(release): relocate 2 orphan test files into the collected flat tests/unit dir (#5120)
Unblocks Lint (test-discovery) on #5078. Integrated into release/v3.8.38.
* fix(translator): preserve reasoning-replay reasoning_content + repair 3 release-green test reds (#5122)
Repairs 3 release-green test reds + test-masking; unblocks #5078.
* test(golden): redact live Node version from provider translate-path snapshot (#5125)
Final golden unblock for #5078.
* test(golden): redact OmniRoute app version from translate-path snapshot (#5126)
Coverage shard golden unblock for #5078.
* Ignore disconnect races during in-band stream error handling (#5007)
Integrated into release/v3.8.38
* Track final connection IDs in failover logs (#5016)
Integrated into release/v3.8.38
* fix(sse): convert Gemini body to OpenAI format in antigravity MITM handler (#4845)
Integrated into release/v3.8.38 (rebased on tip, CHANGELOG re-injected)
* feat(providers): add ZenMux Free session-cookie provider (#5105)
Integrated into release/v3.8.38 (rebased on tip, CHANGELOG re-injected)
* feat(dashboard): click-to-edit model alias in provider page (#5119)
Integrated into release/v3.8.38 (rebased on tip, i18n scope verified, CHANGELOG re-injected)
* feat(mcp): web-session robustness — cookie dedup (PR6) + browser-pool observability (PR7) (#3368) (#5121)
Integrated into release/v3.8.38 (rebased on tip; cookie-dedup branch extracted to findExistingCookieConnection helper → complexity-neutral; CHANGELOG added)
* fix(usage): dedupe request-usage logging and debounce stats (#4940)
Integrated into release/v3.8.38 (rebased on tip; DB-handle hang was stale-base artifact — resetDbInstance already closes the handle, test green 5/5; file-size drift consolidated at release; CHANGELOG re-injected)
* fix(dashboard): key model visibility toggle on canonical providerId (#5091)
Integrated into release/v3.8.38 (retargeted main→release; .tsx visibility-key test green 2/2)
* chore(deps): bump actions/cache from 5.0.5 to 6.0.0 (#5112)
Integrated into release/v3.8.38 (retargeted main→release; workflow-only actions/cache bump — unit failures were stale main base-reds)
* fix(streaming): harden long OpenAI-compatible SSE streams (#5124)
Integrated into release/v3.8.38 (rebased on tip; streamHandler conflict with #5007 disconnect-guard resolved — both coexist, stream-handler 22/22 green)
* feat: Add Grok Build (xAI) provider with OAuth import-token flow (#5020)
Integrated into release/v3.8.38 (rebased on tip; Hard Rule #11 fix — Grok public client_id now via resolvePublicCred(grok_id), 3 literals removed; grok-oauth 7/7 + check:public-creds green)
* feat(providers): add Factory (factory.ai) as a subscription gateway provider (#5065)
Integrated into release/v3.8.38 (rebased on tip; added factory registry test for PR Test Policy + fixed check:env-doc-sync phantom FACTORY_API_KEY; factory loads in PROVIDERS, no Zod issue — that flag was a false positive)
* chore(test): reconcile golden snapshot + apikey count for new providers
#5020 (grok-cli), #5065 (factory), #5105 (zenmux-free) added providers but did
not regenerate tests/snapshots/provider/translate-path.json (now +3 entries) nor
bump the APIKEY_PROVIDERS count (159->160 for the factory gateway). Test-only
reconciliation; no production change.
* fix(resilience): harden quota and model lockout edge cases (#5093)
Integrated into release/v3.8.38 (rebased on tip). TRUST-BUT-VERIFY: dropped the PR's 0dd7df641 'fix unit gates' commit which reverted #5122 reasoning-replay (preserveReasoningContent) + re-introduced #4849 O(n^2) growth, and restored 5 tests it had realigned. Kept only the 3 declared resilience fixes (quota cutoff guard, gemini MIME, model-lockout maxCooldownMs); 23/23 green.
* Hydrate quota cache and scope auto combo candidates (#5015)
Integrated into release/v3.8.38 (rebased on tip). Kept core quota-cache hydration + auto-combo candidate scoping + combos UI; dropped out-of-scope toolCloaking refactor (conflicted with #4813 stripEnumDescriptions — took tip) and the unrelated sse-auth test split. Added quota-cache-hydrate-5015 regression test (Rule #18); combo-account-allowlist 8/8 + hydration 2/2 green.
* chore(quality): reconcile complexity + file-size baselines for v3.8.38 owner-PR batch
complexity 1972->1978 (+6) and file-size providers.ts 1093->1107 / usageHistory.ts
934->983 — drift from the /review-prs merge batch (#4845/#5105/#5020/#4940/#5093/
#5015 + #5121 cookie-dedup helper extraction). check:complexity/check:file-size do
not run on the PR->release fast-path, so the branch accrued unmeasured; all legit
feature/fix growth, not regression. See per-key justifications in each baseline.
* fix(security): exact-host Anthropic baseUrl check (CodeQL js/incomplete-url-substring-sanitization #674) (#5130)
The anthropic-compatible Bearer-fallback gate decided whether a configured baseUrl
targeted the official api.anthropic.com host via a substring `.includes("api.anthropic.com")`.
A look-alike upstream such as `https://api.anthropic.com.evil.test` or
`https://evil.test/?x=api.anthropic.com` matched the substring and was wrongly treated as
official, suppressing the Bearer fallback meant for third-party gateways
(CodeQL #674, js/incomplete-url-substring-sanitization, high).
Replace the substring test with an exported `isOfficialAnthropicBaseUrl()` helper that
parses the URL and compares the hostname for exact equality. Empty baseUrl stays official;
scheme-less hosts are parsed with an assumed https://; an unparseable baseUrl falls back to
third-party (Bearer emitted) as the safer default. Behavior for legitimate official/third-party
baseUrls is unchanged.
Adds tests/unit/anthropic-official-baseurl-host.test.ts covering official, look-alike,
scheme-less, and unparseable inputs plus a static guard that the substring pattern is gone.
* fix(proxy): repair one-click Deno & Cloudflare relay deployments (#5128) (#5132)
* fix(services): embed WS proxy honours LIVE_WS_HOST; reject empty messages early (#5110) (#5133)
* fix(api): resolve /v1/models/{id} case-insensitively (#5082) (#5135)
* fix(providers): add MiniMax M3 & Nemotron 3 Ultra to Cline catalog (#3321) (#5136)
* fix(proxy): make SOCKS5 handshake timeout tunable via SOCKS_HANDSHAKE_TIMEOUT_MS (#5109) (#5137)
* feat(sidebar): add support for colored menu icons (#3812)
Integrated into release/v3.8.38 (recreated on tip — fork had unrelated history; added getSidebarIconAccent regression test, Rule #18). Clean 2-file UI feature.
* fix(providers): complete grok-cli OAuth wiring + zenmux-free web-session metadata
Base-red repair for #5020 (grok-cli) and #5105 (zenmux-free), surfaced by the
full CI on the release PR (#5078) — the PR->release fast-path does not run the
oauth-providers-config / web-session-credentials / provider-consistency gates.
- grok-cli: register in OAUTH_PROVIDERS (providers.ts canonical list, fixes
check:provider-consistency), add OAUTH_PROVIDER_IDS.GROK_CLI + GROK_CLI_CONFIG
in oauth constants (provider config now sourced there, not a local literal),
align oauth-providers-config.test.ts (EXPECTED_PROVIDER_KEYS + config map).
- zenmux-free: declare its web-session credential requirement (full Cookie header)
in WEB_SESSION_CREDENTIAL_REQUIREMENTS.
Local: oauth-providers-config 27/27, web-session-credentials 4/4, grok-cli-oauth
7/7, check:provider-consistency OK, +115 OAUTH_PROVIDERS tests green.
* Fix resilience settings page response mapping (#5139)
Integrated into release/v3.8.38. Thanks @rdself for the fix and the regression test.
* fix(kiro): retire claude-sonnet-4.5 from catalog + pin 400 model-unavailable test (#5140)
Extracted the real change from #5140 (the bot PR regenerated the entire
freeModelCatalog.data.ts + touched package-lock.json; only the targeted
edits are kept here):
- remove claude-sonnet-4.5 from the Kiro registry entry
- remove the matching kiro free-model catalog row
- pin Kiro's verbatim 400 "Invalid model..." to isModelUnavailableError
Closes #4484
* fix(sidebar): drop orphan `settings` accent color (typecheck:core red) (#5142)
SIDEBAR_ICON_ACCENTS is typed Partial<Record<HideableSidebarItemId, string>>,
but `settings` is not a hideable item id (only `settings-general`,
`settings-appearance`, … and `context-settings` exist; there is no item with
`id: "settings"`), so the accent was unreachable. It broke `typecheck:core`
on the release tip ("'settings' does not exist in type …", introduced by
#3812 colored menu icons). Removing the orphan key restores a clean
typecheck:core (rc=0).
* feat: salvage batch 2 — diagnostics null-guard (#5096) + observed quota reset windows (#5025) (#5141)
* fix(diagnostics): null-guard content blocks in detectMalformedNonStream
A null (or non-object) entry in a Claude-native `content` array made the
non-stream classifier throw `TypeError: Cannot read properties of null
(reading 'type')`, crashing the malformed-response detection path. Guard
before type-asserting each block: a null/non-object block is simply skipped.
Two regression tests added (null block among valid blocks → null; only-null
blocks → empty_choices).
Salvaged from closed PR #5096 (base-stale; only the defensive guard — the
Claude-shape recognition it also carried already landed via #5108).
Co-authored-by: herjarsa <herjarsa@users.noreply.github.com>
* feat(quota): persist observed provider quota reset windows
Adds `provider_quota_reset_events` (migration 108) + `db/quotaResetEvents.ts`
to record real upstream weekly-quota window transitions whenever a quota
refresh shows the reset rolling to a new cycle (different day, later resetAt).
`apiKeyUsageLimits` now prefers the observed window start over the inferred
`resetAt − 7d`, falling back to snapshot inference when no event is recorded
yet. `quotaCache.setQuotaCache` records the transition opportunistically.
`recordProviderQuotaResetEventIfChanged` only fires for the primary weekly
window (not daily/sonnet), is idempotent (INSERT OR IGNORE on the unique
window key), and no-ops when the reset didn't actually roll. 4 unit tests
(tests/unit/lib/quota-reset-events.test.ts).
Salvaged from closed PR #5025 (which bundled this with two unrelated
features + a colliding migration 104). Renumbered to 108; module re-exported
from localDb (Rule #2).
Co-authored-by: Witroch4 <175152067+Witroch4@users.noreply.github.com>
---------
Co-authored-by: herjarsa <herjarsa@users.noreply.github.com>
Co-authored-by: Witroch4 <175152067+Witroch4@users.noreply.github.com>
* docs(i18n): sync 3.8.38 CHANGELOG section to 41 mirrors (unblock docs-accuracy) (#5144)
The root CHANGELOG [3.8.38] section grew with this cycle's merged PRs, but the
docs/i18n/<lang>/CHANGELOG.md mirrors were not re-synced — drifting >25% in body
size and failing check:docs-sync (the "Docs accuracy" fast-gate step) for every
open PR against the release.
Ran scripts/release/sync-changelog-i18n.mjs 3.8.38 3.8.37 to copy the root
[3.8.38] section into all 41 mirrors. check:docs-all now passes (exit 0).
Sections are copied verbatim; the per-language translation pass runs at release
time via i18n:run — this only restores the size-sync the gate enforces.
* feat(compression): pure per-step fidelity checker (4 invariants, fail-open)
* feat(compression): fidelityGate config + rejected breakdown fields
* feat(compression): wire per-step fidelity gate into stacked pipeline (opt-in)
* feat(compression): preview route accepts fidelityGate flag (playground)
* feat(compression): playground fidelity-gate toggle + lane rejection display
* docs(compression): note fidelityGate advanced thresholds are intentionally API-omitted
* refactor(compression): extract fidelity-gate step helpers to shrink strategySelector (file-size gate)
bodyToText and gateAdvance moved to fidelityGateStep.ts; StackAccumulator exported.
strategySelector: 889->854 (-35). Residual +6 vs pre-Milestone-B frozen 848 is the
irreducible StackOptions.fidelityGate field + two stacked-loop dispatch reads + import.
Baseline updated to 854 with justification. No cycle introduced (import type only).
940 compression tests pass; typecheck clean.
* test(usage): wire usageHistoryDedup under unit runner brace-list (#5145)
Integrated into release/v3.8.38.
* feat: salvage batch from closed stale PRs (#5038, #5057, #5076) (#5138)
Integrated into release/v3.8.38.
* test(combo): deterministic routing-decision matrix for all 17 strategies (#5146)
Integrated into release/v3.8.38.
* feat(compression): fuzzy near-duplicate dedup (session-dedup 2nd pass + playground toggle) (#5143)
Integrated into release/v3.8.38.
* chore(quality): rebaseline file-size for sidebarVisibility.ts + chat.ts drift (#5147)
Mid-cycle drift on release/v3.8.38 from already-merged PRs that the fast-path
(PR->release skips check:file-size) let accumulate without a bump:
- src/shared/constants/sidebarVisibility.ts 1100->1198 (#3812 colored menu
icons, per-item accent map; #5142 dropped one orphan, net still above frozen)
- src/sse/handlers/chat.ts 1560->1575 (#5064 self-inflicted-timeout cooldown
skip + #5124 long OpenAI-compatible SSE hardening + #5110 embed-WS
LIVE_WS_HOST honour / early empty-message reject)
Each covered by its own PR tests; structural shrink of chat.ts tracked in #3501.
Unblocks the Fast Quality Gates for PRs targeting release/v3.8.38.
* chore(release): finalize v3.8.38 CHANGELOG + cycle reconciliation
- Reconcile [3.8.38]: +18 bullets (compression fidelity-gate/fuzzy-dedup #5143,
quota keepalive #5102, web-session robustness #5121, MiniMax/Nemotron #5136,
model-visibility #5091, failover logs #5016, disconnect races #5007, sidebar
orphan #5142, SRE playbooks salvage #5138, new Security #5130 + Maintenance roll-up)
- Credit salvaged-PR authors (@JxnLexn / @KooshaPari / @herjarsa / @Witroch4)
- Remove phantom bullet for CLOSED-not-merged #5092 (setup aggregator never landed)
- Fix isHidden bullet PR citation #4389 -> #5086 (@herjarsa)
- Back-fill forgotten v3.8.36 bullet: #5026 crypto.randomUUID ID-gen (@hamsa0x7)
- Sync 41 i18n CHANGELOG mirrors; README What's New -> v3.8.38
- Rebaseline cycle drift: eslint 3987->4002, cognitive 833->841, dead-exports
345->346, cyclomatic 1978->1980 (file-size handled by #5147)
* fix(i18n): add missing English UI labels (#5153)
Integrated into release/v3.8.38
* Preserve non-stream reasoning fields for compatible clients (#5155)
Integrated into release/v3.8.38
* feat(compression): ionizer engine — lossy JSON-array sampling reversible via CCR (#5148)
Integrated into release/v3.8.38
* test(combo): gated live smoke for combo strategies (in-process + VPS HTTP) (#5151)
Integrated into release/v3.8.38
* test: refresh release expectations to match current code (#5150)
Integrated into release/v3.8.38 (test-only base-red alignment extracted from #5150)
---------
Co-authored-by: Éder Costa <eder.almeida.costa@gmail.com>
Co-authored-by: José Victor Ferreira <root@josevictor.me>
Co-authored-by: Hernan Javier Ardila Sanchez <hjasgr@gmail.com>
Co-authored-by: fulorgnas <46461624+fulorgnas@users.noreply.github.com>
Co-authored-by: Randi <55005611+rdself@users.noreply.github.com>
Co-authored-by: Jan Leon <Jan.gaschler@gmail.com>
Co-authored-by: R. Beltran <rbeltran8000@gmail.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
Co-authored-by: KooshaPari <42529354+KooshaPari@users.noreply.github.com>
Co-authored-by: Ramel Tecnologia - Rafa Martins <146174365+rafacpti23@users.noreply.github.com>
Co-authored-by: herjarsa <herjarsa@users.noreply.github.com>
Co-authored-by: Witroch4 <175152067+Witroch4@users.noreply.github.com>
665 lines
21 KiB
TypeScript
665 lines
21 KiB
TypeScript
import { trackPendingRequest } from "@/lib/usageDb";
|
|
import { STREAM_IDLE_TIMEOUT_MS } from "../config/constants.ts";
|
|
import { FORMATS } from "../translator/formats.ts";
|
|
import { PENDING_REQUEST_CLEARED_MARKER } from "./stream.ts";
|
|
|
|
// Stream handler with disconnect detection - shared for all providers
|
|
|
|
// Default budget for the pipeWithDisconnect raw-upstream stall watchdog.
|
|
// Inherits STREAM_IDLE_TIMEOUT_MS so a single env knob still governs the
|
|
// max time we tolerate silence from upstream. Reasoning models (Claude
|
|
// thinking, Kiro EventStream binary frames) emit zero post-transform
|
|
// output for long stretches while raw bytes keep arriving — measuring
|
|
// stall on the transform output false-positives on those streams, so
|
|
// the watchdog must track upstream byte activity instead. Ported from
|
|
// decolua/9router#1243.
|
|
const DEFAULT_STREAM_STALL_TIMEOUT_MS = STREAM_IDLE_TIMEOUT_MS;
|
|
|
|
type StreamDisconnectEvent = {
|
|
reason: string;
|
|
duration: number;
|
|
};
|
|
|
|
type StreamErrorEvent = {
|
|
error: unknown;
|
|
message: string;
|
|
statusCode: number;
|
|
duration: number;
|
|
};
|
|
|
|
type StreamControllerOptions = {
|
|
onDisconnect?: (event: StreamDisconnectEvent) => boolean | void;
|
|
onError?: (event: StreamErrorEvent) => boolean | void;
|
|
provider?: string;
|
|
model?: string;
|
|
connectionId?: string | null;
|
|
clientResponseFormat?: string | null;
|
|
clientAbortSignal?: AbortSignal | null;
|
|
};
|
|
|
|
type StreamController = ReturnType<typeof createStreamController>;
|
|
|
|
type StreamErrorStatusKind = "rate_limit" | "authentication" | "permission" | "client" | "server";
|
|
|
|
type StreamErrorStatusMapping = {
|
|
responses: {
|
|
type: string;
|
|
code: string;
|
|
};
|
|
claude: {
|
|
type: string;
|
|
};
|
|
};
|
|
|
|
function isResponsesClientFormat(clientResponseFormat?: string | null): boolean {
|
|
return (
|
|
clientResponseFormat === FORMATS.OPENAI_RESPONSES ||
|
|
clientResponseFormat === FORMATS.OPENAI_RESPONSE
|
|
);
|
|
}
|
|
|
|
function getStreamErrorStatusKind(statusCode: number): StreamErrorStatusKind {
|
|
if (statusCode === 429) return "rate_limit";
|
|
if (statusCode === 401) return "authentication";
|
|
if (statusCode === 403) return "permission";
|
|
if (statusCode >= 400 && statusCode < 500) return "client";
|
|
return "server";
|
|
}
|
|
|
|
function getStreamErrorStatusMapping(statusCode: number): StreamErrorStatusMapping {
|
|
switch (getStreamErrorStatusKind(statusCode)) {
|
|
case "rate_limit":
|
|
return {
|
|
responses: { type: "rate_limit_error", code: "rate_limit_exceeded" },
|
|
claude: { type: "rate_limit_error" },
|
|
};
|
|
case "authentication":
|
|
return {
|
|
responses: { type: "authentication_error", code: "invalid_authentication" },
|
|
claude: { type: "authentication_error" },
|
|
};
|
|
case "permission":
|
|
return {
|
|
responses: { type: "authentication_error", code: "permission_denied" },
|
|
claude: { type: "permission_error" },
|
|
};
|
|
case "client":
|
|
return {
|
|
responses: { type: "invalid_request_error", code: "bad_request" },
|
|
claude: { type: "invalid_request_error" },
|
|
};
|
|
case "server":
|
|
return {
|
|
responses: { type: "server_error", code: "server_error" },
|
|
claude: { type: "api_error" },
|
|
};
|
|
default:
|
|
return {
|
|
responses: { type: "server_error", code: "server_error" },
|
|
claude: { type: "api_error" },
|
|
};
|
|
}
|
|
}
|
|
|
|
function encodeSseEvent(
|
|
data: unknown,
|
|
{
|
|
event,
|
|
includeDone = false,
|
|
}: {
|
|
event?: string;
|
|
includeDone?: boolean;
|
|
} = {}
|
|
) {
|
|
if (event && /[\r\n]/.test(event)) {
|
|
throw new Error("SSE event names must not contain newlines");
|
|
}
|
|
|
|
const encoder = new TextEncoder();
|
|
const prefix = event ? `event: ${event}\n` : "";
|
|
const chunks = [encoder.encode(`${prefix}data: ${JSON.stringify(data)}\n\n`)];
|
|
if (includeDone) {
|
|
chunks.push(encoder.encode("data: [DONE]\n\n"));
|
|
}
|
|
return chunks;
|
|
}
|
|
|
|
// Get HH:MM:SS timestamp
|
|
function getTimeString() {
|
|
return new Date().toLocaleTimeString("en-US", {
|
|
hour12: false,
|
|
hour: "2-digit",
|
|
minute: "2-digit",
|
|
second: "2-digit",
|
|
});
|
|
}
|
|
|
|
function isPendingRequestClearedError(error: unknown): boolean {
|
|
return (
|
|
!!error &&
|
|
typeof error === "object" &&
|
|
(error as Record<string, unknown>)[PENDING_REQUEST_CLEARED_MARKER] === true
|
|
);
|
|
}
|
|
|
|
/**
|
|
* A client disconnect — the caller aborted the request or closed the SSE
|
|
* connection — is NOT a provider failure. It surfaces either as an
|
|
* AbortError/ResponseAborted, or, when OmniRoute then tries to enqueue another
|
|
* chunk into the now-closed response stream, as a "Controller is already closed"
|
|
* TypeError. Treating any of these as an upstream error wrongly cools down the
|
|
* account/connection, so the stream error path uses this to skip the provider
|
|
* failover/cooldown (the chatgpt-web / codex / antigravity executors already
|
|
* guard client aborts the same way).
|
|
*/
|
|
export function isClientDisconnectError(error: unknown): boolean {
|
|
if (!error || typeof error !== "object") return false;
|
|
const name = (error as { name?: unknown }).name;
|
|
if (name === "AbortError" || name === "ResponseAborted") return true;
|
|
const message = (error as { message?: unknown }).message;
|
|
return typeof message === "string" && /Controller is already closed/i.test(message);
|
|
}
|
|
|
|
function getErrorMessage(error: unknown): string {
|
|
if (error instanceof Error && error.message) return error.message;
|
|
if (typeof error === "string" && error.trim().length > 0) return error;
|
|
return "Upstream stream error";
|
|
}
|
|
|
|
function getErrorStatusCode(error: unknown): number {
|
|
if (error && typeof error === "object" && "statusCode" in error) {
|
|
const statusCode = Number((error as { statusCode?: unknown }).statusCode);
|
|
if (Number.isFinite(statusCode) && statusCode >= 400 && statusCode <= 599) {
|
|
return statusCode;
|
|
}
|
|
}
|
|
return 502;
|
|
}
|
|
|
|
function hasClientTerminalSseMarker(text: string, clientResponseFormat?: string | null): boolean {
|
|
if (/(?:^|\r?\n)data:\s*\[DONE\]\s*(?:\r?\n|$)/.test(text)) {
|
|
return true;
|
|
}
|
|
|
|
if (isResponsesClientFormat(clientResponseFormat)) {
|
|
return (
|
|
/(?:^|\r?\n)event:\s*response\.completed\s*(?:\r?\n|$)/.test(text) ||
|
|
/"type"\s*:\s*"response\.completed"/.test(text)
|
|
);
|
|
}
|
|
|
|
if (clientResponseFormat === FORMATS.CLAUDE) {
|
|
return (
|
|
/(?:^|\r?\n)event:\s*message_stop\s*(?:\r?\n|$)/.test(text) ||
|
|
/"type"\s*:\s*"message_stop"/.test(text)
|
|
);
|
|
}
|
|
|
|
return false;
|
|
}
|
|
|
|
/**
|
|
* Create stream controller with abort and disconnect detection
|
|
* @param {object} options
|
|
* @param {function} options.onDisconnect - Callback when client disconnects
|
|
* @param {object} options.log - Logger instance
|
|
* @param {string} options.provider - Provider name
|
|
* @param {string} options.model - Model name
|
|
*/
|
|
/** @param {StreamControllerOptions} options */
|
|
export function createStreamController({
|
|
onDisconnect,
|
|
onError,
|
|
provider,
|
|
model,
|
|
connectionId,
|
|
clientResponseFormat,
|
|
clientAbortSignal,
|
|
}: StreamControllerOptions = {}) {
|
|
const abortController = new AbortController();
|
|
const startTime = Date.now();
|
|
let disconnected = false;
|
|
let pendingRequestCleared = false;
|
|
let cleanupClientAbortSignal: (() => void) | null = null;
|
|
|
|
const logStream = (status) => {
|
|
const duration = Date.now() - startTime;
|
|
const p = provider?.toUpperCase() || "UNKNOWN";
|
|
console.log(
|
|
`[${getTimeString()}] 🌊 [STREAM] ${p} | ${model || "unknown"} | ${duration}ms | ${status}`
|
|
);
|
|
};
|
|
|
|
const clearPendingRequest = (error?: unknown) => {
|
|
if (pendingRequestCleared) return;
|
|
if (
|
|
error &&
|
|
typeof error === "object" &&
|
|
(error as Record<string, unknown>)[PENDING_REQUEST_CLEARED_MARKER] === true
|
|
) {
|
|
pendingRequestCleared = true;
|
|
return;
|
|
}
|
|
|
|
pendingRequestCleared = true;
|
|
if (!model && !provider && !connectionId) return;
|
|
try {
|
|
trackPendingRequest(model || "", provider || "", connectionId ?? null, false);
|
|
} catch {}
|
|
};
|
|
|
|
const cleanupClientAbortListener = () => {
|
|
if (!cleanupClientAbortSignal) return;
|
|
cleanupClientAbortSignal();
|
|
cleanupClientAbortSignal = null;
|
|
};
|
|
|
|
const getClientAbortReason = () => {
|
|
const reason = clientAbortSignal?.reason;
|
|
if (typeof reason === "string" && reason.trim().length > 0) {
|
|
return reason;
|
|
}
|
|
if (reason instanceof Error && reason.message) {
|
|
return reason.message;
|
|
}
|
|
return "request_signal_aborted";
|
|
};
|
|
|
|
const controller = {
|
|
signal: abortController.signal,
|
|
startTime,
|
|
|
|
isConnected: () => !disconnected,
|
|
|
|
// Call when client disconnects
|
|
handleDisconnect: (reason = "client_closed") => {
|
|
if (disconnected) return;
|
|
disconnected = true;
|
|
cleanupClientAbortListener();
|
|
|
|
logStream(`disconnect: ${reason}`);
|
|
|
|
// Decrement pending request counter — the TransformStream flush() won't
|
|
// fire when the client aborts mid-stream, so we must clean up here.
|
|
clearPendingRequest();
|
|
|
|
abortController.abort(reason);
|
|
|
|
onDisconnect?.({ reason, duration: Date.now() - startTime });
|
|
},
|
|
|
|
// Call when stream completes normally
|
|
handleComplete: () => {
|
|
if (disconnected) return;
|
|
disconnected = true;
|
|
cleanupClientAbortListener();
|
|
|
|
logStream("complete");
|
|
},
|
|
|
|
// Call on error
|
|
handleError: (error: unknown) => {
|
|
cleanupClientAbortListener();
|
|
|
|
// A client disconnect is not a provider failure. If the client already went away
|
|
// (disconnected) or the error is a client abort / "Controller is already closed",
|
|
// skip the onError failover/cooldown path — otherwise one cancelled request marks
|
|
// the upstream connection unavailable.
|
|
if (disconnected || isClientDisconnectError(error)) {
|
|
clearPendingRequest(error);
|
|
logStream(disconnected ? "client_disconnect (post-abort)" : "client_disconnect");
|
|
return;
|
|
}
|
|
|
|
const alreadyCleared = isPendingRequestClearedError(error);
|
|
let handled = false;
|
|
if (!alreadyCleared) {
|
|
try {
|
|
handled =
|
|
onError?.({
|
|
error,
|
|
message: getErrorMessage(error),
|
|
statusCode: getErrorStatusCode(error),
|
|
duration: Date.now() - startTime,
|
|
}) === true;
|
|
} catch {}
|
|
}
|
|
|
|
if (!handled) {
|
|
clearPendingRequest(error);
|
|
} else {
|
|
pendingRequestCleared = true;
|
|
}
|
|
|
|
if (error instanceof Error && error.name === "AbortError") {
|
|
logStream("aborted");
|
|
return;
|
|
}
|
|
|
|
if (error instanceof Error) {
|
|
logStream(`error: ${error.message}`);
|
|
return;
|
|
}
|
|
logStream("error: unknown");
|
|
},
|
|
|
|
abort: () => {
|
|
cleanupClientAbortListener();
|
|
abortController.abort();
|
|
},
|
|
clientResponseFormat,
|
|
};
|
|
|
|
if (clientAbortSignal && typeof clientAbortSignal.addEventListener === "function") {
|
|
const handleClientAbort = () => {
|
|
controller.handleDisconnect(getClientAbortReason());
|
|
};
|
|
if (clientAbortSignal.aborted) {
|
|
queueMicrotask(handleClientAbort);
|
|
} else {
|
|
clientAbortSignal.addEventListener("abort", handleClientAbort, { once: true });
|
|
cleanupClientAbortSignal = () => {
|
|
clientAbortSignal.removeEventListener("abort", handleClientAbort);
|
|
};
|
|
}
|
|
}
|
|
|
|
return controller;
|
|
}
|
|
|
|
function buildStreamErrorChunks(
|
|
errorMsg: string,
|
|
statusCode: number,
|
|
clientResponseFormat?: string | null
|
|
) {
|
|
const statusMapping = getStreamErrorStatusMapping(statusCode);
|
|
|
|
if (isResponsesClientFormat(clientResponseFormat)) {
|
|
const errorEvent = {
|
|
type: "response.failed",
|
|
response: {
|
|
id: null,
|
|
status: "failed",
|
|
error: {
|
|
message: errorMsg,
|
|
type: statusMapping.responses.type,
|
|
code: statusMapping.responses.code,
|
|
},
|
|
},
|
|
};
|
|
|
|
return encodeSseEvent(errorEvent, { event: "response.failed" });
|
|
}
|
|
|
|
if (clientResponseFormat === FORMATS.CLAUDE) {
|
|
const errorEvent = {
|
|
type: "error",
|
|
error: {
|
|
type: statusMapping.claude.type,
|
|
message: errorMsg,
|
|
},
|
|
};
|
|
|
|
return encodeSseEvent(errorEvent, { event: "error" });
|
|
}
|
|
|
|
const errorEvent = {
|
|
object: "chat.completion.chunk",
|
|
choices: [
|
|
{
|
|
index: 0,
|
|
delta: {},
|
|
finish_reason: "error",
|
|
},
|
|
],
|
|
error: {
|
|
message: errorMsg,
|
|
type: statusMapping.responses.type,
|
|
code: statusMapping.responses.code,
|
|
},
|
|
};
|
|
|
|
return encodeSseEvent(errorEvent, { includeDone: true });
|
|
}
|
|
|
|
/**
|
|
* Minimal `writable` half used by `pipeWithDisconnect`. The real writable is
|
|
* driven entirely by the upstream-piped readable, so the writer only needs an
|
|
* `abort()` hook for `createDisconnectAwareStream`'s `cancel()` path.
|
|
*
|
|
* `abort()` returns `Promise<void>` to match the native
|
|
* `WritableStreamDefaultWriter.abort()` contract — `cancel()` (and any caller
|
|
* that awaits the writer) gets a real thenable instead of `undefined`, which
|
|
* keeps abort/error handling clean. Ported from decolua/9router@6b624af4.
|
|
*/
|
|
export function createNoopAbortWritable(): {
|
|
getWriter: () => { abort: () => Promise<void> };
|
|
} {
|
|
return { getWriter: () => ({ abort: () => Promise.resolve() }) };
|
|
}
|
|
|
|
/**
|
|
* Create transform stream with disconnect detection
|
|
* Wraps existing transform stream and adds abort capability
|
|
*/
|
|
export function createDisconnectAwareStream(transformStream, streamController) {
|
|
const reader = transformStream.readable.getReader();
|
|
const writer = transformStream.writable.getWriter();
|
|
const terminalDecoder = new TextDecoder();
|
|
let terminalTail = "";
|
|
let clientTerminalSeen = false;
|
|
|
|
const noteClientChunk = (chunk: unknown) => {
|
|
if (clientTerminalSeen) return;
|
|
if (!(chunk instanceof Uint8Array)) return;
|
|
|
|
terminalTail += terminalDecoder.decode(chunk, { stream: true });
|
|
if (terminalTail.length > 4096) {
|
|
terminalTail = terminalTail.slice(-4096);
|
|
}
|
|
clientTerminalSeen = hasClientTerminalSseMarker(
|
|
terminalTail,
|
|
streamController.clientResponseFormat
|
|
);
|
|
};
|
|
|
|
return new ReadableStream(
|
|
{
|
|
async pull(controller) {
|
|
if (!streamController.isConnected()) {
|
|
controller.close();
|
|
return;
|
|
}
|
|
|
|
try {
|
|
const { done, value } = await reader.read();
|
|
if (done) {
|
|
streamController.handleComplete();
|
|
controller.close();
|
|
return;
|
|
}
|
|
controller.enqueue(value);
|
|
noteClientChunk(value);
|
|
} catch (error) {
|
|
if (!streamController.isConnected()) {
|
|
try {
|
|
controller.close();
|
|
} catch {}
|
|
return;
|
|
}
|
|
|
|
if (clientTerminalSeen) {
|
|
streamController.handleComplete();
|
|
try {
|
|
controller.close();
|
|
} catch {}
|
|
return;
|
|
}
|
|
|
|
streamController.handleError(error);
|
|
|
|
// T35: Encapsulate mid-stream errors as SSE events instead of abruptly aborting
|
|
// This prevents TransferEncodingError on the client side
|
|
const errorMsg = getErrorMessage(error);
|
|
const statusCode = getErrorStatusCode(error);
|
|
|
|
try {
|
|
for (const chunk of buildStreamErrorChunks(
|
|
errorMsg,
|
|
statusCode,
|
|
streamController.clientResponseFormat
|
|
)) {
|
|
controller.enqueue(chunk);
|
|
}
|
|
} catch {
|
|
// The downstream may have closed while we were formatting the in-band
|
|
// error event. The original stream error has already been recorded.
|
|
}
|
|
|
|
try {
|
|
controller.close();
|
|
} catch {}
|
|
}
|
|
},
|
|
|
|
async cancel(reason) {
|
|
streamController.handleDisconnect(reason || "cancelled");
|
|
await Promise.allSettled([reader.cancel(reason), writer.abort(reason)]);
|
|
},
|
|
},
|
|
{ highWaterMark: 16384 }
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Pipe provider response through transform with disconnect detection.
|
|
*
|
|
* Stall watchdog tracks raw upstream byte activity, not transform output.
|
|
* Reasoning models (Claude thinking via Kiro, etc.) can produce zero SSE
|
|
* output for long stretches while partial EventStream frames keep arriving;
|
|
* measuring stall on the transform output caused false stalls. Any upstream
|
|
* chunk resets the timer. If no bytes arrive for `stallTimeoutMs`, the
|
|
* stream surfaces a "stream stall timeout" error and aborts.
|
|
*
|
|
* Ported from decolua/9router#1243 by @zakirkun.
|
|
*
|
|
* @param providerResponse - Response from provider
|
|
* @param transformStream - Transform stream for SSE
|
|
* @param streamController - Stream controller from createStreamController
|
|
* @param opts.stallTimeoutMs - Override the stall budget (defaults to
|
|
* STREAM_IDLE_TIMEOUT_MS / DEFAULT_STREAM_STALL_TIMEOUT_MS). `0` disables
|
|
* the watchdog.
|
|
*/
|
|
export function pipeWithDisconnect(
|
|
providerResponse: Response,
|
|
transformStream: TransformStream<Uint8Array, Uint8Array>,
|
|
streamController: StreamController,
|
|
opts: { stallTimeoutMs?: number } = {}
|
|
) {
|
|
const stallTimeoutMs = opts.stallTimeoutMs ?? DEFAULT_STREAM_STALL_TIMEOUT_MS;
|
|
|
|
// Watchdog disabled — preserve legacy behavior verbatim.
|
|
if (!stallTimeoutMs || stallTimeoutMs <= 0) {
|
|
const transformedBody = providerResponse.body.pipeThrough(transformStream);
|
|
return createDisconnectAwareStream(
|
|
{ readable: transformedBody, writable: createNoopAbortWritable() },
|
|
streamController
|
|
);
|
|
}
|
|
|
|
let stallTimer: ReturnType<typeof setTimeout> | null = null;
|
|
// Captured on the upstream tap's `start`, used by the watchdog to error the
|
|
// pipeline so the downstream reader unblocks and emits a clean SSE error
|
|
// event. Without this, aborting the AbortController alone does not unblock
|
|
// a `reader.read()` already suspended on the transform pipe — the request
|
|
// would hang until the upstream finally closed the socket.
|
|
let upstreamTapController: TransformStreamDefaultController<Uint8Array> | null = null;
|
|
// Set when the watchdog fires so the downstream pull() catch (which sees
|
|
// the same error propagated through the pipeline) does not call
|
|
// handleError a second time — pending-cleanup is idempotent but onError
|
|
// callbacks should fire once per error.
|
|
let stallFired = false;
|
|
|
|
const clearStall = () => {
|
|
if (stallTimer) {
|
|
clearTimeout(stallTimer);
|
|
stallTimer = null;
|
|
}
|
|
};
|
|
const armStall = () => {
|
|
clearStall();
|
|
stallTimer = setTimeout(() => {
|
|
stallTimer = null;
|
|
stallFired = true;
|
|
const stallError = new Error("stream stall timeout");
|
|
// Notify the controller (onError callback + pending-request cleanup).
|
|
try {
|
|
streamController.handleError?.(stallError);
|
|
} catch {}
|
|
// Error the pipeline so the downstream reader unblocks. createDisconnect-
|
|
// AwareStream's catch block translates this into buildStreamErrorChunks
|
|
// (sanitized SSE error event with finish_reason:"error", per the format).
|
|
try {
|
|
upstreamTapController?.error(stallError);
|
|
} catch {}
|
|
// Abort the underlying fetch so upstream releases the connection.
|
|
try {
|
|
streamController.abort?.();
|
|
} catch {}
|
|
}, stallTimeoutMs);
|
|
};
|
|
|
|
// Wrap controller so every termination path clears the stall timer.
|
|
// Without this, abort/complete/error/disconnect paths leave the timer armed
|
|
// and a stale abort could fire after the request has already ended.
|
|
const wrappedController: StreamController = {
|
|
...streamController,
|
|
handleComplete: () => {
|
|
clearStall();
|
|
streamController.handleComplete();
|
|
},
|
|
handleError: (e: unknown) => {
|
|
clearStall();
|
|
// Watchdog already fired its own handleError — the inner pull() catch
|
|
// sees the same error propagated through the pipeline; suppress the
|
|
// duplicate to keep onError callbacks single-fire.
|
|
if (stallFired) return;
|
|
streamController.handleError(e);
|
|
},
|
|
handleDisconnect: (reason?: string) => {
|
|
clearStall();
|
|
streamController.handleDisconnect(reason);
|
|
},
|
|
abort: () => {
|
|
clearStall();
|
|
streamController.abort();
|
|
},
|
|
};
|
|
|
|
// Inert tap that resets the stall timer on every raw upstream byte chunk.
|
|
// Sits between the provider body and the SSE transform so reasoning models
|
|
// that buffer many raw bytes into a single emitted event do not look
|
|
// stalled to the watchdog.
|
|
const upstreamTap = new TransformStream<Uint8Array, Uint8Array>({
|
|
start(controller) {
|
|
upstreamTapController = controller;
|
|
armStall();
|
|
},
|
|
transform(chunk, controller) {
|
|
armStall();
|
|
controller.enqueue(chunk);
|
|
},
|
|
flush() {
|
|
clearStall();
|
|
},
|
|
});
|
|
|
|
const transformedBody = providerResponse.body
|
|
.pipeThrough(upstreamTap)
|
|
.pipeThrough(transformStream);
|
|
return createDisconnectAwareStream(
|
|
{ readable: transformedBody, writable: createNoopAbortWritable() },
|
|
wrappedController
|
|
);
|
|
}
|