Files
OmniRoute/open-sse/utils/streamHandler.ts
Diego Rodrigues de Sa e Souza 555b21d296 Release v3.8.37 (#5053)
* chore(release): open v3.8.37 development cycle

* chore(ci): harden release flow — ratchet decoupling, fast-path drift gates, build-scope guard, heap default (#5054)

Implements improvements 1-4 from the v3.8.36 release benchmark (_tasks/release-bench/v3.8.36/PLANO-MELHORIA.md):

1. Quality Ratchet decoupled from flaky coverage (ci.yml): the shard→coverage→ratchet
   chain meant a single flaky Coverage Shard SKIPPED the whole Quality Ratchet on the
   release PR (v3.8.36 #4854), so cycle drift only surfaced post-merge in #5029. The job
   now runs on !cancelled(); coverage download is continue-on-error and the ratchet runs
   --allow-missing, so the DETERMINISTIC gates (eslint/complexity/cognitive/duplication/
   codeql) stay blocking even when coverage is unavailable.

2. Fast-path drift gates (quality.yml PR→release): added check:complexity,
   check:cognitive-complexity, and a new lightweight check:pack-policy (pack-artifact
   unexpected-files check WITHOUT a build, via --policy-only) so drift + stray-tarball-file
   regressions are caught/rebaselined PER-PR instead of cascading onto the release PR.

3. Build heap default 4096→8192 MB (build-next-isolated.mjs): the clean graph peaks
   ~3.9 GB and brushed the old 4 GB ceiling; 8 GB gives headroom. Comment notes heap is
   NOT the fix for a poisoned scope (run check:build-scope instead).

4. check:build-scope gate (new): fails if .ts/.tsx/.js/.jsx files in the tsconfig scope
   exceed a threshold — catches worktrees/cruft leaking into the build scope (the v3.8.36
   OOM root cause: 355,215 vs 4,547 files) BEFORE it detonates next build. Wired into the
   fast-path.

* fix(auth): only trust forwarding headers from loopback TCP peers (#4689)

Integrated into release/v3.8.37 — loopback-gated forwarding headers (IP spoofing fix). Cherry-picked onto current release tip; ipUtils.test.ts 9/9 green.

* fix(codex): treat OAuth 401 as unrecoverable refresh failure (#4686)

Integrated into release/v3.8.37 — codex OAuth 401 treated as unrecoverable refresh. Cherry-picked onto release tip; token-refresh-service.test.ts 38/38 green.

* fix(translator): preserve reasoning_effort for non-Copilot Responses clients (#4688)

Integrated into release/v3.8.37 — preserve reasoning_effort for non-Copilot Responses clients. Cherry-picked onto release tip; tests 47/47 green.

* fix(translator): coerce tool descriptions to strings in OpenAI normalization (#4675)

Integrated into release/v3.8.37 — coerce tool descriptions to strings in OpenAI normalization. Cherry-picked onto release tip; tests 3/3 green.

* feat(sse): x-omniroute-strip-reasoning header to drop reasoning_content (#4678)

Integrated into release/v3.8.37 — x-omniroute-strip-reasoning header. Cherry-picked onto release tip (resolved chatCore.ts/headers.ts adjacency conflict, kept resolveCompressionHeader + isStripReasoningRequested); tests 8/8 green.

* fix(combo): flatten Anthropic tool messages + tool history to prevent upstream 503 (#4648)

Integrated into release/v3.8.37 — flattenToolHistory helper (combo anti-503). Cherry-picked onto release tip; tests 9/9 green.

* feat(headroom): proxy lifecycle management + dashboard UI (Docker sidecar supported) (#4649)

Integrated into release/v3.8.37 — headroom proxy lifecycle (status/start/stop, local-only + spawn-capable per Rules #15/#17). Cherry-picked onto release tip; lifecycle 7/7 + route-guard 43/43 + check:cycles green.

* feat(cli): multi-model support for Factory Droid CLI (#4682)

Integrated into release/v3.8.37 — Factory Droid multi-model support. Cherry-picked onto release tip (kept readJsoncConfig + droidCustomModels imports); droid-custom-models 11/11 green.

* fix(providers): require Default Model in compatible-provider API-key setup (#4641)

Integrated into release/v3.8.37 — require Default Model in compatible-provider API-key setup. Cherry-picked fix + test-move onto release tip (kept release providerSpecificData + QuotaScrapingFields; fixed moved-test import path; baseline rebaseline unneeded, 865<866); UI test 2/2 green.

* fix(dashboard): stop double-masking already-masked API key in list (E2E 3/9 regression) (#4671)

Integrated into release/v3.8.37 — render server-masked key verbatim (drop redundant maskKey call). Note: release's maskKey already guards '****' (since v3.8.34), so this is a safe simplification; added a contract test pinning the **** passthrough invariant (2/2 green, would fail against the pre-guard maskKey = the historical double-mask bug).

* chore(quality): rebaseline file-size for rc17 PR batch drift

Own growth from the merged rc17 PRs (#4678/#4686/#4688) at existing chokepoints —
cohesive, not extractable:
- open-sse/handlers/responseSanitizer.ts 1103->1122 (SanitizeOpenAIResponseOptions + stripReasoning, #4678)
- open-sse/services/tokenRefresh.ts 2070->2090 (codex 401 unrecoverable-refresh guard, #4686)
- tests/unit/token-refresh-service.test.ts 1322->1353 (401 regression case, #4686)
- tests/unit/translator-openai-responses-req.test.ts 1047->1050 (reasoning_effort assertion, #4688)

* docs(env): document HEADROOM_URL in .env.example + ENVIRONMENT.md

The headroom proxy lifecycle (#4649) reads HEADROOM_URL (src/lib/headroom/detect.ts,
default http://localhost:8787) but it was missing from the env contract, tripping
check:env-doc-sync. Adds the var to both .env.example (commented, has a default) and
the Proxy Health table in ENVIRONMENT.md.

* fix(sse): stream writer mock abort() returns a Promise (#4788)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(cli): fall back to default data dir when DATA_DIR is not writable (#4767)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(oauth): verify Cursor installation on Linux before auto-import (#4770)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): track Ollama streaming usage from raw NDJSON chunks (#4754)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): strip enumDescriptions from antigravity tool schema (#4740)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): include low-level cause details in formatProviderError (#4741)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(translator): strip x-anthropic-billing-header in claude-to-openai (#4728)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): gate Kiro image attachments behind a Claude-capability check (#4763)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): read Antigravity usage from the response.usageMetadata envelope (#4785)

Integrated into release/v3.8.37 — Antigravity response.usageMetadata envelope. Cherry-picked onto release tip (resolved test-tail adjacency with #4754 Ollama block); usage-extractor 23/23 green.

* fix(api): fall back to existing access token for any OAuth provider on refresh failure (#4786)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(cli): verify launchd registration + skip self-SIGTERM in macOS autostart (#4765)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(executors): anthropic-compatible-* gateways get Bearer alongside x-api-key (#4729)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): json_schema fallback for OpenAI-compatible providers (#4766)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): use workos auth token shape for cline (#4787)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* feat(sse): parse Gemini CLI 429 retryDelay from structured RetryInfo (#4738)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): finalize tool_calls finish_reason on early stream end in OpenAI Responses translator (#4764)

Integrated into release/v3.8.37 — computeFinishReason finalizes tool_calls on early stream end (Responses translator). Cherry-picked onto release tip; responses-translation-fixes 29/29 green.

* test(sse): golden-lock provider.ts translate-path across all providers (#4734)

Integrated into release/v3.8.37 — golden-lock for provider.ts translate-path. Cherry-picked onto release tip; snapshot regenerated against the current provider set (UPDATE_GOLDEN=1, 167 entries); golden test 3/3 deterministic.

* chore(quality): rebaseline file-size for rc17 leva2 PR batch drift

Own growth from the merged leva2 PRs (cohesive, not extractable):
- src/lib/usage/providerLimits.ts 950->955 (#4786)
- open-sse/executors/default.ts NEW frozen @828 (#4729 + #4766 + #4787 header branches)
- open-sse/translator/request/openai-to-kiro.ts 807->814 (#4763)
- open-sse/translator/response/openai-responses.ts 923->937 (#4764)
- tests/unit/executor-default-base.test.ts 1339->1440 (#4766)
- tests/unit/translator-openai-to-kiro.test.ts 918->980 (#4763)

* fix(dashboard): align Engine Combos editor engines with API schema (#4955) (#5062)

The named-combos pipeline dropdown offered four engines (headroom,
session-dedup, ccr, llmlingua) that stackedPipelineStepSchema rejects, so
selecting one made PUT /api/context/combos/[id] return HTTP 400 while
saveCombo swallowed the non-OK response (if (!res.ok) return). Editing the
default 'Standard Savings' combo and changing an engine reproduced the 400.

- Add canonical STACKED_PIPELINE_ENGINE_INTENSITIES next to the schema as the
  single source of truth; the client dropdown imports it so it can never drift
  from the discriminated union the API validates against.
- Surface save errors and empty-name/empty-pipeline validation in the editor
  instead of failing silently.
- Add a parity unit test asserting the UI engine map equals the schema union
  and that every (engine, intensity) the UI emits is accepted.

* fix(sse): filter nameless hosted tools when converting Responses API to Chat format (#4789)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(dashboard): keep desktop sidebar visible via explicit CSS class (#4812)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip enumDescriptions from Antigravity tool schemas (#4813)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(dashboard): resolve passthrough model aliases by providerId in ModelSelectModal (#4815)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(oauth): allow per-connection refresh lead-time override via providerSpecificData.refreshLeadMs (#4818)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip X-Stainless-* headers and normalize SDK User-Agent for OpenAI-compatible endpoints (#4820)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip Gemini built-in tools when functionDeclarations present in Antigravity envelope (#4821)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(api): surface a Docker-localhost hint on provider-node validation connection errors (#4822)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): resolve bare model names to connection defaultModel before upstream calls (#4825)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(build): trace-include sql.js sql-wasm.wasm in standalone bundle (#4839)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip Composer <|final|> sentinel markers leaking after Composer reasoning (#4842)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(config): sync full SiliconFlow model list into registry (#4844)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): close reasoning before message content in Responses stream (#4848)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): reject unsupported Kiro [1m] context suffix (#4816)

Integrated into release/v3.8.37 — cherry-picked onto release tip; test-tail conflict with #4763 resolved (kept both image + [1m] test blocks); CHANGELOG re-merged; 29/29 green.

* fix(db): validate HuggingFace tokens via whoami-v2 auth probe (#4819)

Integrated into release/v3.8.37 — defining commit re-homed onto the god-file-split validation module (validateHuggingFaceProvider in validation/openaiFormat.ts + map wiring); 115/115 green.

* fix(sse): make anthropic-version default-guard case-insensitive (#4823)

Integrated into release/v3.8.37 — conflict with #4729 Bearer-fallback resolved (kept both Bearer fallback + case-insensitive anthropic-version guard); 48/48 green.

* fix(sse): sanitize Kiro tool schemas to avoid 400 "Improperly formed request" (#4847)

Integrated into release/v3.8.37 — conflict in kiro-to-openai.ts resolved (kept release fallbackToolCallId + adopted #1375 toolNameMap remap); 7/7 green.

* feat(sse): add GPT-4 to the GitHub Copilot provider (#4798)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(sse): add GPT-4o mini to GitHub Copilot provider (#4797)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(api): add MiniMax-M3 pricing row (#4814)

Integrated into release/v3.8.37 — pricing row re-homed onto god-file-split pricing/regional.ts (pricing.ts is now a barrel); 4/4 green.

* fix(cli): save runtime deps with --save-exact so a sibling install can't prune them (#4841)

Integrated into release/v3.8.37 — trayRuntime conflict resolved (kept release SYSTRAY_SPEC + added --save-exact); 2/2 green.

* fix(sse): preserve required fields in antigravity tool schemas (#4843)

Integrated into release/v3.8.37 — conflict resolved (kept #4740/#4813 enumDescriptions strip + typed normalizeSchemaTypes, added required-preservation helpers; test-tail merged keeping both enumDescriptions + required tests); 7/7 green.

* chore(quality): rebaseline file-size for rc17b leva3 PR batch drift

* fix(sse): strip reasoning blobs from agentic context to prevent O(n^2) token growth (#4849)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): unwrap Qoder HTTP 200 SSE error envelope so fallback can trigger (#4850)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip temperature for Claude models with extended thinking (#4853)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): emit valid concatenable kiro tool_calls.arguments deltas (#4855)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(sse): add toggleable tool-source diagnostics (#4856)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): redact api key from the AUTH debug log in the chat handler (#4858)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): forward AI SDK image parts in Responses translator (#4859)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): resolve custom combos by id and case-insensitive name (#4446) (#4869)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): exclude WS bridge controller-closed error from provider breaker (#4602) (#4870)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(providers): add xAI Grok inbound translators and thinking patcher (#4910)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(embeddings): add dimensions override field to embedding combos (#4913)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(oauth): Codex bulk-import endpoint — POST /api/oauth/codex/import (#4914)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(antigravity): retry transient upstream failures (#4941)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): surface malformed HTTP-200 upstream responses (#4942)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): normalize Codex custom tools (apply_patch) to { input: string } schema (#4862)

Integrated into release/v3.8.37 — conflict in request/openai-responses.ts resolved (kept #4789 nameless-tool skip + added #1007 custom-tool {input:string} normalization); 48/48 green incl. #4789/#4859 regression.

* fix(sse): dense, deterministic output ordering in Responses API response.completed (#4906)

Integrated into release/v3.8.37 — manual integration with #4862 in response/openai-responses.ts (custom-tool funcItem + dense recordCompletedItem). Fixed a latent #4848 interaction: the close-reasoning-before-message guard force-closed <think>-tag reasoning prematurely, which dense output (#4906) then snapshotted as a partial buffer ("plan" vs "planning") — scoped the guard to native reasoning_content (!inThinking) in BOTH transformer + translator paths. Full Responses suite 203/203 green incl. #4848/#4862 regression.

* feat(sse): auto-promote successful combo model to position #1 (#4852)

Integrated into release/v3.8.37 — dropped the stale file-size-baseline.json hunk (re-derived against the rc17b rebaseline); code+test applied clean; 13/13 green.

* feat(providers): add Pioneer AI (Fastino Labs) provider (#4909)

Integrated into release/v3.8.37 — providers.ts apikey block re-homed onto god-file-split src/shared/constants/providers/apikey/frontier-labs.ts (inline APIKEY_PROVIDERS no longer exists); registry/pioneer + providers/index.ts applied clean; 6/6 green.

* add DGrid AI gateway provider (#4931)

Integrated into release/v3.8.37 — rebased the contributor's commit onto the release tip; providers.ts god-file-split conflict resolved by relocating the dgrid APIKEY_PROVIDERS entry into apikey/gateways.ts; CHANGELOG added. 7/7 green. Thanks @dgridOP!

* chore(quality): rebaseline file-size for rc17b leva4 PR batch drift

* docs(routing): sync combo strategy docs for Fusion (17 strategies) (#5067)

Fusion (16th strategy, panel fan-out + judge synthesis) and headroom
shipped but the strategy-count docs were stale (14/15) and omitted both.
Update every combo-strategy reference to the canonical 17, add fusion +
headroom to all strategy tables, and add a dedicated Fusion section to
AUTO-COMBO.md documenting judgeModel / fusionTuning config + an example.

- CLAUDE.md, README.md, FEATURES.md, RESILIENCE_GUIDE.md,
  ARCHITECTURE.md, OPEN_SSE_ARCHITECTURE.md, OMNIROUTE_VS_ALTERNATIVES.md,
  docs/README.md, request-pipeline.mmd: 14/15 -> 17, list fusion + headroom
- docs/routing/AUTO-COMBO.md: strategy table + new Fusion strategy section
- docs/openapi.yaml: add reset-window, headroom, fusion to the strategy enum

* fix(oauth): classify /api/oauth/cursor/auto-import as local-only (route-guard) (#5070)

The Cursor auto-import route runs execFile("which", ["cursor"]) to verify a
local Cursor install before importing credentials — a child-process spawn. The
check:route-guard-membership gate (Hard Rules #15/#17) flagged it as an
unclassified spawn-capable route: reachable past the loopback gate, an
RCE-via-tunnel surface (a leaked JWT over a tunnel could trigger the spawn).

Classify the specific path in LOCAL_ONLY_API_PREFIXES so loopback enforcement
runs unconditionally before any auth check. Scoped to the exact path — the rest
of /api/oauth/ (browser redirect/callback flows) stays remote-reachable.

TDD: added a failing-then-passing assertion in route-guard-local-prefix.test.ts
(classification + an over-broadening guard proving sibling OAuth paths stay
remote). check:route-guard-membership now reports 0 new gaps.

* chore(release): v3.8.37 — 2026-06-26

---------

Co-authored-by: dgridOP <dgrid_op@outlook.com>
2026-06-26 02:51:06 -03:00

549 lines
16 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
const DISCONNECT_ABORT_DELAY_MS = 2_000;
// 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) => void;
onError?: (event: StreamErrorEvent) => boolean | void;
provider?: string;
model?: string;
connectionId?: string | null;
clientResponseFormat?: string | 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
);
}
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;
}
/**
* 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,
}: StreamControllerOptions = {}) {
const abortController = new AbortController();
const startTime = Date.now();
let disconnected = false;
let abortTimeout: ReturnType<typeof setTimeout> | null = null;
let pendingRequestCleared = false;
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 {}
};
return {
signal: abortController.signal,
startTime,
isConnected: () => !disconnected,
// Call when client disconnects
handleDisconnect: (reason = "client_closed") => {
if (disconnected) return;
disconnected = true;
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();
// Delay abort to allow cleanup
abortTimeout = setTimeout(() => {
abortController.abort();
}, DISCONNECT_ABORT_DELAY_MS);
onDisconnect?.({ reason, duration: Date.now() - startTime });
},
// Call when stream completes normally
handleComplete: () => {
if (disconnected) return;
disconnected = true;
logStream("complete");
if (abortTimeout) {
clearTimeout(abortTimeout);
abortTimeout = null;
}
},
// Call on error
handleError: (error: unknown) => {
if (abortTimeout) {
clearTimeout(abortTimeout);
abortTimeout = null;
}
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: () => abortController.abort(),
clientResponseFormat,
};
}
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();
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);
} catch (error) {
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);
for (const chunk of buildStreamErrorChunks(
errorMsg,
statusCode,
streamController.clientResponseFormat
)) {
controller.enqueue(chunk);
}
controller.close();
}
},
cancel(reason) {
streamController.handleDisconnect(reason || "cancelled");
reader.cancel();
setTimeout(() => {
writer.abort();
}, DISCONNECT_ABORT_DELAY_MS).unref?.();
},
},
{ 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
);
}