Files
OmniRoute/tests/unit/stream-handler.test.ts
Diego Rodrigues de Sa e Souza 7b139fdb5e Release v3.8.38 (#5078)
* 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>
2026-06-27 09:07:12 -03:00

704 lines
21 KiB
TypeScript

import test from "node:test";
import assert from "node:assert/strict";
import {
createDisconnectAwareStream,
createNoopAbortWritable,
createStreamController,
pipeWithDisconnect,
} from "../../open-sse/utils/streamHandler.ts";
import { FORMATS } from "../../open-sse/translator/formats.ts";
import {
clearPendingRequests,
getPendingRequests,
trackPendingRequest,
} from "../../src/lib/usage/usageHistory.ts";
const encoder = new TextEncoder();
const decoder = new TextDecoder();
const PENDING_REQUEST_CLEARED_MARKER = "__omniroutePendingRequestCleared";
async function readStreamText(stream) {
const reader = stream.getReader();
const chunks = [];
while (true) {
const { done, value } = await reader.read();
if (done) break;
chunks.push(value);
}
return decoder.decode(
chunks.length === 1 ? chunks[0] : Uint8Array.from(chunks.flatMap((chunk) => Array.from(chunk)))
);
}
test("createDisconnectAwareStream converts upstream errors into SSE error chunks", async () => {
const upstreamError = Object.assign(new Error("provider exploded"), { statusCode: 429 });
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(transformStream, createStreamController());
const text = await readStreamText(stream);
assert.match(text, /"finish_reason":"error"/);
assert.match(text, /"message":"provider exploded"/);
assert.match(text, /"code":"rate_limit_exceeded"/);
assert.match(text, /\[DONE\]/);
});
test("createDisconnectAwareStream treats errors after OpenAI DONE as successful completion", async () => {
let pullCount = 0;
let errorHandled = false;
const transformStream = {
readable: new ReadableStream({
pull(controller) {
pullCount += 1;
if (pullCount === 1) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
return;
}
controller.error(new Error("terminated"));
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({
onError() {
errorHandled = true;
},
})
);
const text = await readStreamText(stream);
assert.equal(text, "data: [DONE]\n\n");
assert.equal(errorHandled, false);
assert.doesNotMatch(text, /finish_reason/);
assert.doesNotMatch(text, /terminated/);
});
test("createDisconnectAwareStream: Gemini 503 high-demand error becomes SSE error chunk with message preserved", async () => {
const geminiMsg =
"[503]: This model is currently experiencing high demand. Spikes in demand are usually temporary. Please try again later.";
const upstreamError = Object.assign(new Error(geminiMsg), { statusCode: 503 });
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(transformStream, createStreamController());
const text = await readStreamText(stream);
assert.match(text, /"finish_reason":"error"/);
assert.match(text, /"message":"\[503\]: This model is currently experiencing high demand/);
assert.match(text, /"type":"server_error"/);
assert.match(text, /"code":"server_error"/);
assert.match(text, /\[DONE\]/);
});
test("createDisconnectAwareStream emits Responses API failure events for Responses clients", async () => {
const upstreamError = Object.assign(new Error("responses stream\ndied"), { statusCode: 503 });
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({ clientResponseFormat: FORMATS.OPENAI_RESPONSES })
);
const text = await readStreamText(stream);
assert.match(text, /event: response\.failed/);
assert.match(text, /"type":"response\.failed"/);
assert.match(text, /"message":"responses stream\\ndied"/);
assert.match(text, /"type":"server_error"/);
assert.match(text, /"code":"server_error"/);
assert.doesNotMatch(text, /chat\.completion\.chunk/);
assert.doesNotMatch(text, /"finish_reason":"error"/);
assert.doesNotMatch(text, /\[DONE\]/);
});
test("createDisconnectAwareStream keeps newlines escaped inside SSE data fields", async () => {
const upstreamError = Object.assign(new Error("line one\nline two\rline three"), {
statusCode: 400,
});
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({ clientResponseFormat: FORMATS.OPENAI_RESPONSES })
);
const text = await readStreamText(stream);
assert.match(text, /^event: response\.failed\ndata: \{"type":"response\.failed"/);
assert.match(text, /"message":"line one\\nline two\\rline three"/);
assert.doesNotMatch(text, /^line two/m);
assert.doesNotMatch(text, /^line three/m);
});
test("createDisconnectAwareStream treats legacy OpenAI response format alias as Responses", async () => {
const upstreamError = Object.assign(new Error("legacy responses alias died"), {
statusCode: 429,
});
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({ clientResponseFormat: FORMATS.OPENAI_RESPONSE })
);
const text = await readStreamText(stream);
assert.match(text, /event: response\.failed/);
assert.match(text, /"type":"rate_limit_error"/);
assert.match(text, /"code":"rate_limit_exceeded"/);
assert.doesNotMatch(text, /chat\.completion\.chunk/);
assert.doesNotMatch(text, /\[DONE\]/);
});
test("createDisconnectAwareStream emits Claude SSE errors for Claude clients", async () => {
const upstreamError = Object.assign(new Error("claude stream died"), { statusCode: 502 });
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({ clientResponseFormat: FORMATS.CLAUDE })
);
const text = await readStreamText(stream);
assert.match(text, /event: error/);
assert.match(text, /"type":"error"/);
assert.match(text, /"type":"api_error"/);
assert.match(text, /"message":"claude stream died"/);
assert.doesNotMatch(text, /"code"/);
assert.doesNotMatch(text, /chat\.completion\.chunk/);
assert.doesNotMatch(text, /"finish_reason":"error"/);
assert.doesNotMatch(text, /\[DONE\]/);
});
test("createDisconnectAwareStream keeps newlines escaped for Claude SSE errors", async () => {
const upstreamError = Object.assign(new Error("claude line one\nclaude line two"), {
statusCode: 502,
});
const transformStream = {
readable: new ReadableStream({
start(controller) {
controller.error(upstreamError);
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const stream = createDisconnectAwareStream(
transformStream,
createStreamController({ clientResponseFormat: FORMATS.CLAUDE })
);
const text = await readStreamText(stream);
assert.match(text, /^event: error\ndata: \{"type":"error"/);
assert.match(text, /"message":"claude line one\\nclaude line two"/);
assert.doesNotMatch(text, /^claude line two/m);
});
test("createDisconnectAwareStream cancel propagates disconnect reason and aborts the writer", async () => {
let aborted = false;
let disconnectEvent = null;
const transformStream = {
readable: new ReadableStream({
pull() {},
cancel() {},
}),
writable: {
getWriter() {
return {
abort() {
aborted = true;
},
};
},
},
};
const controller = createStreamController({
onDisconnect(event) {
disconnectEvent = event;
},
});
const stream = createDisconnectAwareStream(transformStream, controller);
await stream.cancel("client-gone");
await new Promise((resolve) => setTimeout(resolve, 2050));
assert.equal(aborted, true);
assert.equal(controller.isConnected(), false);
assert.equal(disconnectEvent.reason, "client-gone");
assert.ok(disconnectEvent.duration >= 0);
});
test("createNoopAbortWritable: getWriter().abort() returns a resolved Promise (matches WritableStreamDefaultWriter contract)", async () => {
// The mock writable that pipeWithDisconnect hands to createDisconnectAwareStream
// is consumed only via its writer's abort() hook (in the cancel() path). The
// native WritableStreamDefaultWriter.abort() returns Promise<void>; the mock
// must match that contract so cancel/error handling can await it instead of
// receiving `undefined`. Ported from decolua/9router@6b624af4.
const writable = createNoopAbortWritable();
const writer = writable.getWriter();
const aborted = writer.abort();
assert.ok(aborted instanceof Promise, "abort() must return a Promise, not undefined");
// Awaiting must resolve cleanly to undefined (Promise<void>), never reject.
assert.equal(await aborted, undefined);
});
test("createNoopAbortWritable: cancelling a stream wired through it awaits the abort promise without throwing", async () => {
// End-to-end seam: the noop writable is what pipeWithDisconnect injects. Wire
// it into createDisconnectAwareStream exactly as production does and drive the
// cancel() path. With abort() returning undefined (the pre-fix shape) this
// still completes, but a thenable abort keeps the cancel/error path clean.
const transformStream = {
readable: new ReadableStream({
pull() {},
cancel() {},
}),
writable: createNoopAbortWritable(),
};
const stream = createDisconnectAwareStream(transformStream, createStreamController());
await assert.doesNotReject(stream.cancel("client-gone"));
});
test("createDisconnectAwareStream uses the default cancel reason when none is provided", async () => {
let disconnectEvent = null;
const transformStream = {
readable: new ReadableStream({
cancel() {},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const controller = createStreamController({
onDisconnect(event) {
disconnectEvent = event;
},
});
const stream = createDisconnectAwareStream(transformStream, controller);
await stream.cancel();
assert.equal(disconnectEvent.reason, "cancelled");
});
test("createDisconnectAwareStream closes immediately when the controller is already disconnected", async () => {
const controller = createStreamController();
controller.handleDisconnect("preclosed");
const stream = createDisconnectAwareStream(
{
readable: new ReadableStream({
pull(inner) {
inner.enqueue(encoder.encode("ignored"));
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
},
controller
);
const reader = stream.getReader();
const first = await reader.read();
assert.equal(first.done, true);
});
test("createStreamController aborts after delayed disconnect and tolerates abort/unknown errors", async () => {
const controller = createStreamController();
const errorOnlyController = createStreamController();
controller.handleDisconnect();
controller.handleDisconnect("ignored-repeat");
errorOnlyController.handleError(new DOMException("aborted", "AbortError"));
errorOnlyController.handleError({ statusCode: 418 });
await new Promise((resolve) => setTimeout(resolve, 2050));
assert.equal(controller.signal.aborted, true);
assert.equal(controller.isConnected(), false);
assert.equal(errorOnlyController.signal.aborted, false);
});
test("pipeWithDisconnect pipes transformed bytes and marks the controller complete", async () => {
const source = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode("hello"));
controller.close();
},
});
const providerResponse = new Response(source);
const controller = createStreamController();
const stream = pipeWithDisconnect(providerResponse, new TransformStream(), controller);
const text = await readStreamText(stream);
assert.equal(text, "hello");
assert.equal(controller.isConnected(), false);
});
test("pipeWithDisconnect clears pending requests when the upstream stream errors", async () => {
clearPendingRequests();
const provider = "openai";
const model = "gpt-stream-error";
const connectionId = "conn-stream-error";
const modelKey = `${model} (${provider})`;
trackPendingRequest(model, provider, connectionId, true);
const source = new ReadableStream({
start(controller) {
controller.error(Object.assign(new Error("socket closed"), { statusCode: 502 }));
},
});
const stream = pipeWithDisconnect(
new Response(source),
new TransformStream(),
createStreamController({ provider, model, connectionId })
);
const text = await readStreamText(stream);
const pending = getPendingRequests();
assert.match(text, /"message":"socket closed"/);
assert.equal(pending.byModel[modelKey], 0);
assert.equal(pending.details[connectionId], undefined);
});
test("pipeWithDisconnect lets controller onError own pending cleanup", async () => {
clearPendingRequests();
const provider = "openai";
const model = "gpt-stream-error-owned";
const connectionId = "conn-stream-error-owned";
const modelKey = `${model} (${provider})`;
let errorEvent = null;
trackPendingRequest(model, provider, connectionId, true);
const source = new ReadableStream({
start(controller) {
controller.error(Object.assign(new Error("terminated"), { statusCode: 502 }));
},
});
const stream = pipeWithDisconnect(
new Response(source),
new TransformStream(),
createStreamController({
provider,
model,
connectionId,
onError(event) {
errorEvent = event;
return true;
},
})
);
const text = await readStreamText(stream);
const pending = getPendingRequests();
assert.match(text, /"message":"terminated"/);
assert.equal(errorEvent?.statusCode, 502);
assert.equal(pending.byModel[modelKey], 1);
assert.equal(pending.byAccount[connectionId][modelKey], 1);
});
test("pipeWithDisconnect does not double-clear transform errors already accounted for", async () => {
clearPendingRequests();
const provider = "openai";
const model = "gpt-marked-error";
const connectionId = "conn-marked-error";
const modelKey = `${model} (${provider})`;
trackPendingRequest(model, provider, connectionId, true);
trackPendingRequest(model, provider, connectionId, true);
trackPendingRequest(model, provider, connectionId, false);
const markedError = Object.assign(new Error("already cleared"), {
[PENDING_REQUEST_CLEARED_MARKER]: true,
});
const source = new ReadableStream({
start(controller) {
controller.error(markedError);
},
});
const stream = pipeWithDisconnect(
new Response(source),
new TransformStream(),
createStreamController({ provider, model, connectionId })
);
await readStreamText(stream);
const pending = getPendingRequests();
assert.equal(pending.byModel[modelKey], 1);
assert.equal(pending.byAccount[connectionId][modelKey], 1);
});
test("createDisconnectAwareStream ignores reader errors after client disconnect", async () => {
let readableController!: ReadableStreamDefaultController;
let onErrorCalled = false;
const transformStream = {
readable: new ReadableStream({
start(controller) {
readableController = controller;
},
}),
writable: {
getWriter() {
return {
abort() {},
};
},
},
};
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});
const stream = createDisconnectAwareStream(transformStream, streamController);
const reader = stream.getReader();
const readPromise = reader.read();
streamController.handleDisconnect("ResponseAborted");
readableController.error(new Error("Invalid state: Controller is already closed"));
const result = await readPromise;
assert.equal(result.done, true);
assert.equal(onErrorCalled, false, "disconnect races must not be recorded as upstream errors");
});
// Stall detection: tied to RAW upstream byte activity, not transform output.
// Ports decolua/9router#1243 — reasoning models (Claude thinking, Kiro
// EventStream binary frames) can stream raw bytes for long stretches while
// the SSE transform produces zero output as it accumulates a frame. The
// stall watchdog must NOT fire on those slow-but-progressing streams.
test("pipeWithDisconnect does NOT flag a slow but progressing upstream as stalled (no false positive)", async () => {
// Upstream emits 3 small chunks 30ms apart (90ms total). The transform
// never forwards any output (simulates a translator buffering a frame
// boundary that has not yet completed). The stall budget is 200ms — well
// above the 30ms gap between upstream bytes, so a byte-activity watchdog
// should never fire. A transform-output-activity watchdog would
// false-stall here.
const source = new ReadableStream({
async start(controller) {
controller.enqueue(encoder.encode("a"));
await new Promise((r) => setTimeout(r, 30));
controller.enqueue(encoder.encode("b"));
await new Promise((r) => setTimeout(r, 30));
controller.enqueue(encoder.encode("c"));
await new Promise((r) => setTimeout(r, 30));
controller.close();
},
});
// Black-hole transform — consumes every byte, emits nothing until flush.
const swallowingTransform = new TransformStream({
transform() {
/* drop chunk — output stream is silent */
},
flush(controller) {
controller.enqueue(encoder.encode("done"));
},
});
let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});
const stream = pipeWithDisconnect(new Response(source), swallowingTransform, streamController, {
stallTimeoutMs: 200,
});
const text = await readStreamText(stream);
// No stall error — final flush output reaches the client cleanly.
assert.equal(text, "done");
assert.equal(
onErrorCalled,
false,
"stall watchdog must NOT fire on a slow but progressing upstream"
);
assert.doesNotMatch(text, /stall/i);
assert.doesNotMatch(text, /"finish_reason":"error"/);
});
test("pipeWithDisconnect flags a truly stalled upstream (no bytes for the full stall budget)", async () => {
// Upstream emits one byte and then goes silent forever. Stall budget is
// 80ms — the watchdog must fire and surface a stream-stall error.
const source = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode("x"));
// never enqueue again, never close — simulate a truly hung upstream
},
cancel() {
// upstream cancel hook so the stall abort path can release the source
},
});
let onErrorEvent = null;
const streamController = createStreamController({
onError(event) {
onErrorEvent = event;
return true;
},
});
const stream = pipeWithDisconnect(new Response(source), new TransformStream(), streamController, {
stallTimeoutMs: 80,
});
const text = await readStreamText(stream);
assert.ok(onErrorEvent !== null, "stall watchdog must fire when upstream stops sending bytes");
assert.match(onErrorEvent.message, /stall/i);
assert.match(text, /stall/i);
assert.match(text, /"finish_reason":"error"/);
});
test("pipeWithDisconnect stall watchdog does not fire after normal stream completion", async () => {
// Upstream completes quickly. The stall timer must be cleared on
// completion so a stale abort cannot fire after the request has ended.
const source = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode("ok"));
controller.close();
},
});
let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});
const stream = pipeWithDisconnect(new Response(source), new TransformStream(), streamController, {
stallTimeoutMs: 50,
});
const text = await readStreamText(stream);
// Wait past the stall budget — no late stall error must surface.
await new Promise((r) => setTimeout(r, 120));
assert.equal(text, "ok");
assert.equal(onErrorCalled, false, "stall watchdog must be cleared on stream completion");
});