Files
OmniRoute/open-sse/utils/streamHandler.ts
Diego Rodrigues de Sa e Souza 7d57d9f4a1 fix(providers): retire common ChatGPT Web provider (#11754)
Rebased onto the current release/v3.8.51 tip as part of a combined provider-retirement/provenance merge batch (Designer Web, Felo Web, Runtime, GPL-derived removal, Qwen Web already landed). Large conflict set (this is the biggest PR in the batch — the common ChatGPT Web provider touches chat, images, count-tokens, session leases, and combos). Conflicts resolved:

- `open-sse/config/providers/registry/chatgpt-web/*`, `open-sse/executors/chatgpt-web*`, `open-sse/handlers/imageGeneration/providers/chatgptWeb.ts`, and their tests: kept deleted, matching the PR's stated scope.
- `open-sse/config/providers/registry/minimax/web/index.ts`, `open-sse/handlers/imageGeneration/providers/geminiWeb.ts`, `open-sse/executors/gemini-web.ts`'s stale image-mode branch: base-drift collisions against already-merged sibling retirements (#11691, #11708) — kept deleted / dropped the dead code, since this PR's own branch forked before those merged.
- `src/shared/constants/reservedProviderPrefixes.ts`, `open-sse/executors/index.ts`, `executorProxy.ts`, `virtualFactory.ts`, `autoStrategy.ts`, `src/lib/db/providers.ts`, `src/sse/handlers/chat.ts`: combined the Designer + Runtime (Felo/Qwen) + common-ChatGPT-Web retirement guard calls at each shared chokepoint — compute-once-then-OR pattern, consistent with prior combinations in this batch.
- `src/sse/services/model.ts` / `src/sse/handlers/chatHelpers.ts`: adopted this PR's new `getModelInfoOrRetirementResponse()` central wrapper (a real improvement over ad-hoc try/catch), and extended it to also catch the Designer + Runtime retirement errors it didn't originally cover, so the consolidation doesn't regress the other two mechanisms.
- `src/app/api/v1/images/edits/route.ts`: this PR moved the retirement check earlier (before `enforceApiKeyPolicy`) but left the old later call+catch block in place from base drift — removed the now-redundant duplicate `resolveImageRouteModel()` call and merged the Designer catch into the earlier one.
- `open-sse/config/imageRegistry.ts`, `tests/snapshots/executors/executor-map.json` (`keyCount` recomputed to 133), `tests/snapshots/provider/translate-path.json`: same "both sides inserted a different retired provider at the same slot" pattern — resolved by dropping both.
- `tests/unit/chatcore-executor-proxy.test.ts`, `provider-node-reserved-prefix.test.ts`, `combo-auto-candidate-expansion.test.ts`, `messages-count-tokens-route.test.ts`, `virtual-auto-combo.test.ts`: split into independent per-mechanism test blocks (established pattern); `virtual-auto-combo.test.ts`'s old "includes cookie web-session providers" positive-inclusion test (which used chatgpt-web as its example) was retired along with the provider and replaced by this PR's negative-exclusion test for the same slot.
- `docs/architecture/ARCHITECTURE.md`, `CODEBASE_DOCUMENTATION.md` (+ 4 i18n mirrors), `README.md`, `FREE-TIERS-GUIDE.md`, `docs/diagrams/free-tier-budget.svg`, `docs/screenshots/free-tier-budget-card.svg`, `docs/reference/PROVIDER_REFERENCE.md`: recomputed every stale count from the real merged state — 104 executors (`countFiles` gate logic), 351 providers (regenerated via `gen:provider-reference`), 152/351 `hasFree` entries, 445/438/7 free-tier catalog rows, 13 ToS-avoid providers, budget-card regenerated via its real generator script. One doc conflict (`oauth/` module list) needed picking HEAD's side specifically — theirs still listed the already-removed `raycast` module instead of the real `openference`.
- `config/quality/test-masking-allowlist.json`: additive merge of the PR's 17 `_deletedWithReplacement` entries alongside the batch's existing ones (one real duplicate-key mistake in my first pass, caught and fixed via a `object_pairs_hook` duplicate-key check before finalizing).

Also fixed two real, unrelated-to-my-merge issues surfaced by the focused suite:
- `tests/unit/resolve-web-provider-host.test.ts`: the PR's own test had a typo — it asserted `perplexity-web`'s resolved host as `"perplexity.ai"`, but the provider's registered `website` is `"https://www.perplexity.ai"` and the resolver returns the URL's `host` verbatim (no www-stripping), so the correct value is `"www.perplexity.ai"` (consistent with the same test's own `url` assertion).
- `tests/unit/hard-session-lease-bypass-inventory.test.ts`: this golden call-site inventory was already stale on the pristine post-#11713 tip (confirmed via a throwaway probe worktree) — `src/lib/db/providers.ts`'s 3 connection-fallback sites and a third `src/app/api/providers/route.ts` site were never added to the golden list by the earlier-merged #11698/#11720 PRs. Updated it to the real current inventory (dated inline comments explain each delta and which PR introduced it), plus this PR's own legitimate deltas (image-edits duplicate-call removal, `ChatGptWebExecutor.execute()` site removed).

Focused suite green (433/433 across executor-proxy, reserved-prefix, hard-session-lease-bypass-inventory, resolve-web-provider-host, retirement/runtime-block/source-retirement/management-retirement/image-handler-retirement, migration-168, combo-auto-candidate-expansion, virtual-auto-combo, executor-map-golden and siblings), plus `typecheck:core`, `check-file-size`, and `check-changelog-integrity` clean. Thanks for the thorough provenance-hold retirement work — appreciated.
2026-08-28 06:52:46 -03:00

962 lines
33 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";
import { createCompletedResponsesToolHandoffWatcher } from "./responsesToolHandoff.ts";
import { createStreamContentWatcher, type StreamContentWatcher } from "./streamReadiness.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;
allowCompletedToolHandoffGrace?: boolean;
clientDisconnectGracePeriodMs?: number;
};
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 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 {
const errorName =
error && typeof error === "object" && typeof (error as { name?: unknown }).name === "string"
? (error as { name: string }).name
: "";
if (errorName === "TimeoutError" || errorName === "BodyTimeoutError") {
return 504;
}
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 isDeadlineAbortReason(reason: unknown): reason is Error {
return (
reason instanceof Error &&
(reason.name === "TimeoutError" || reason.name === "BodyTimeoutError")
);
}
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)
);
}
// OpenAI chat completions: some providers omit `data: [DONE]` (already
// matched above) and terminate with a finish_reason chunk instead. A
// non-null finish_reason value is that terminal signal — a bare
// `finish_reason: null` delta chunk must NOT count (#10443).
if (clientResponseFormat === FORMATS.OPENAI) {
return /"finish_reason"\s*:\s*"[^"]+"/.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,
allowCompletedToolHandoffGrace = false,
clientDisconnectGracePeriodMs = 0,
}: StreamControllerOptions = {}) {
const abortController = new AbortController();
const startTime = Date.now();
let disconnected = false;
let clientTerminalSeen = false;
let completedToolHandoffSeen = false;
let completedToolHandoffDrain: (() => void) | null = null;
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 (e) {
console.error(
`[${getTimeString()}] [streamHandler] trackPendingRequest decrement failed — counter may drift`,
e
);
}
};
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;
if (clientTerminalSeen) {
controller.handleComplete();
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();
const deferUpstreamAbort =
allowCompletedToolHandoffGrace &&
clientDisconnectGracePeriodMs > 0 &&
completedToolHandoffSeen &&
completedToolHandoffDrain !== null;
if (deferUpstreamAbort) {
completedToolHandoffDrain?.();
} else {
abortController.abort(reason);
}
onDisconnect?.({ reason, duration: Date.now() - startTime });
},
// Call when stream completes normally
handleComplete: () => {
if (disconnected) return;
disconnected = true;
cleanupClientAbortListener();
logStream("complete");
},
markClientTerminalSeen: () => {
clientTerminalSeen = true;
},
markCompletedToolHandoffSeen: () => {
completedToolHandoffSeen = true;
},
registerCompletedToolHandoffDrain: (drain: () => void) => {
completedToolHandoffDrain = drain;
},
shouldDeferCompletedToolHandoff: () =>
allowCompletedToolHandoffGrace &&
clientDisconnectGracePeriodMs > 0 &&
completedToolHandoffSeen &&
completedToolHandoffDrain !== null,
// 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 (e) {
console.debug(`[STREAM-HANDLER] onError callback error:`, e);
}
}
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,
clientDisconnectGracePeriodMs,
};
if (clientAbortSignal && typeof clientAbortSignal.addEventListener === "function") {
const handleClientAbort = () => {
const reason = clientAbortSignal.reason;
if (isDeadlineAbortReason(reason)) {
// An AbortSignal can represent an OmniRoute-owned deadline as well as
// a caller disconnect. Preserve deadline failures as 504; classifying
// them as client disconnects writes a misleading 499 to the call log.
abortController.abort(reason);
controller.handleError(reason);
return;
}
controller.handleDisconnect(getClientAbortReason());
};
if (clientAbortSignal.aborted) {
queueMicrotask(handleClientAbort);
} else {
clientAbortSignal.addEventListener("abort", handleClientAbort, { once: true });
cleanupClientAbortSignal = () => {
clientAbortSignal.removeEventListener("abort", handleClientAbort);
};
}
}
return controller;
}
export 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,
},
};
// #7699 — emit message_stop after event:error so Anthropic SDK / Claude Code
// see a proper terminal frame instead of a silent mid-response close.
// Without message_stop, clients report "Connection closed mid-response."
return [
...encodeSseEvent(errorEvent, { event: "error" }),
...encodeSseEvent({ type: "message_stop" }, { event: "message_stop" }),
];
}
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 });
}
/**
* Synthesized terminal frames for a graceful truncation (#7699): the upstream
* ended without a terminal marker AFTER content was already forwarded to the
* client. Instead of an `event: error` frame (which would discard the partial
* content and report a mid-response failure), emit a clean Claude completion —
* `message_delta` carrying `stop_reason: "max_tokens"` followed by
* `message_stop` — so Anthropic SDK / Claude Code treat the response as a
* budget-limited finish and keep everything already received.
*/
export function buildGracefulTruncationChunks(clientResponseFormat?: string | null): Uint8Array[] {
if (clientResponseFormat !== FORMATS.CLAUDE) return [];
return [
...encodeSseEvent(
{
type: "message_delta",
delta: { stop_reason: "max_tokens", stop_sequence: null },
usage: { input_tokens: 0, output_tokens: 0 },
},
{ event: "message_delta" }
),
...encodeSseEvent({ type: "message_stop" }, { event: "message_stop" }),
];
}
/**
* 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
*/
/**
* Why a finished upstream stream should still be reported as a failure, or null
* when the close was clean. Two distinct silent-close shapes:
*
* - **#7699, no terminal marker.** Scoped to Claude (`/v1/messages`), which is
* the issue's real scope: Anthropic's SSE spec permits a mid-stream
* `event: error`, and Claude clients treat a stream ending without
* `message_stop` as an error. When content already reached the client this is
* NOT a provider failure — the partial response is valid and must be kept — so
* it resolves to a graceful truncation (`stop_reason: max_tokens`). For every
* other format (plain OpenAI chat completions included) a
* done-without-recognized-marker close is NOT necessarily a drop — many
* formats have no `[DONE]` equivalent — so synthesising an error there would
* be a false positive.
*
* - **#8649, no content at all.** The stream terminated properly and carried no
* model output. Unlike the marker case this is not format-dependent: a
* completed stream with zero content is a failure everywhere, and it is the
* streaming twin of the non-streaming `isEmptyContentResponse` check. Only
* applies to bodies that actually looked like SSE, and terminal states where
* emptiness is legitimate (length / tool_calls / content_filter / max_tokens /
* tool_use) are excluded by the watcher. If the stream already carried a
* substantive SSE `error` / `response.failed` / Claude `event:error`, stand
* down — same spirit as Claude #3685 `lifecycle.hasError` and readiness #8972
* (do not invent empty content on top of an actionable error).
*/
type SilentCloseOutcome = { kind: "truncated" } | { kind: "error"; reason: string };
function resolveSilentCloseOutcome(input: {
bytesWereForwarded: boolean;
clientTerminalSeen: boolean;
clientResponseFormat?: string | null;
contentWatcher: StreamContentWatcher;
}): SilentCloseOutcome | null {
if (!input.bytesWereForwarded) return null;
if (!input.clientTerminalSeen) {
if (input.clientResponseFormat === FORMATS.CLAUDE && input.contentWatcher.sawContent()) {
// #7699 — upstream dropped after content reached the client on a Claude
// stream. Keep the partial response: emit a clean max_tokens completion
// instead of an error frame so Anthropic SDK / Claude Code don't report
// a mid-response break.
return { kind: "truncated" };
}
// #10443: every known path that produces OpenAI chat chunks emits a
// terminal — the response translators (gemini/claude/kiro/cursor-to-openai)
// all emit a finish_reason chunk, the non-standard executors (kiro, cursor,
// nlpcloud, poe-web, copilot-m365-web, chipotle, gitlab)
// enqueue `data: [DONE]` themselves, and standard OpenAI-compatible
// upstreams end with finish_reason + [DONE] per spec. So a close that
// forwarded content but no terminal marker is an upstream drop, not a
// legitimate end. Guard on sawContent() so the #8649 empty-content
// verdict below keeps its more precise shape for content-free closes.
if (input.clientResponseFormat === FORMATS.OPENAI && input.contentWatcher.sawContent()) {
return { kind: "error", reason: "Upstream stream ended without a terminal marker" };
}
// Responses-format clients (Codex CLI and other /v1/responses consumers):
// a healthy OpenAI Responses stream ALWAYS terminates with an explicit
// `response.completed` event — it is the format's only terminal marker and
// carries the final status/usage. Content forwarded without it is an
// upstream drop, the same class as #10443 for chat completions; surface a
// synthetic response.failed instead of a silent close so clients report
// the break instead of waiting on a completion event that never comes.
if (isResponsesClientFormat(input.clientResponseFormat) && input.contentWatcher.sawContent()) {
return { kind: "error", reason: "Upstream stream ended without a terminal marker" };
}
}
const watcher = input.contentWatcher;
if (watcher.sawError()) return null;
if (watcher.sawSseFrame() && !watcher.sawContent() && !watcher.sawLegitEmptyTerminal()) {
return { kind: "error", reason: "Provider returned empty content" };
}
return null;
}
export function createDisconnectAwareStream(transformStream, streamController) {
const reader = transformStream.readable.getReader();
const writer = transformStream.writable.getWriter();
const terminalDecoder = new TextDecoder();
const contentDecoder = new TextDecoder();
const contentWatcher = createStreamContentWatcher();
const completedToolHandoffWatcher = createCompletedResponsesToolHandoffWatcher();
const toolHandoffDecoder = new TextDecoder();
let terminalTail = "";
let clientTerminalSeen = false;
let bytesWereForwarded = false;
let completedToolHandoffDrainStarted = false;
const drainCompletedToolHandoff = () => {
if (completedToolHandoffDrainStarted) return;
completedToolHandoffDrainStarted = true;
const gracePeriodMs = Math.max(0, Number(streamController.clientDisconnectGracePeriodMs) || 0);
const timeoutReason = "completed_tool_handoff_grace_expired";
const timeout = setTimeout(() => {
streamController.abort();
void Promise.allSettled([reader.cancel(timeoutReason), writer.abort(timeoutReason)]);
}, gracePeriodMs);
void (async () => {
try {
while (true) {
const { done } = await reader.read();
if (done) break;
}
streamController.handleComplete();
} catch (error) {
streamController.handleError(error);
} finally {
clearTimeout(timeout);
}
})();
};
streamController.registerCompletedToolHandoffDrain?.(drainCompletedToolHandoff);
const noteClientChunk = (chunk: unknown) => {
if (!(chunk instanceof Uint8Array)) return;
bytesWereForwarded = true;
// Runs past clientTerminalSeen: the frame that carries the terminal marker
// can carry the only content too, and #8649 needs the whole stream scanned.
contentWatcher.note(contentDecoder.decode(chunk, { stream: true }));
if (
isResponsesClientFormat(streamController.clientResponseFormat) &&
completedToolHandoffWatcher.note(toolHandoffDecoder.decode(chunk, { stream: true }))
) {
streamController.markCompletedToolHandoffSeen?.();
}
if (clientTerminalSeen) return;
terminalTail += terminalDecoder.decode(chunk, { stream: true });
// Scan before bounding retained state: a compaction terminal frame can
// exceed the tail budget because encrypted_content is carried inline.
clientTerminalSeen = hasClientTerminalSseMarker(
terminalTail,
streamController.clientResponseFormat
);
if (terminalTail.length > 4096) {
terminalTail = terminalTail.slice(-4096);
}
if (clientTerminalSeen) {
streamController.markClientTerminalSeen?.();
}
};
return new ReadableStream(
{
async pull(controller) {
if (!streamController.isConnected()) {
controller.close();
return;
}
try {
const { done, value } = await reader.read();
if (done) {
contentWatcher.finish();
const silentClose = resolveSilentCloseOutcome({
bytesWereForwarded,
clientTerminalSeen,
clientResponseFormat: streamController.clientResponseFormat,
contentWatcher,
});
if (silentClose?.kind === "truncated") {
// #7699 — the upstream dropped without a terminal marker after
// content reached the client. Keep the partial response: emit a
// clean `max_tokens` completion instead of an error frame so
// Anthropic SDK / Claude Code don't report a mid-response break.
streamController.handleComplete();
try {
for (const chunk of buildGracefulTruncationChunks(
streamController.clientResponseFormat
)) {
controller.enqueue(chunk);
}
} catch {
// downstream may have closed; stream already marked complete
}
} else if (silentClose) {
streamController.handleError(
Object.assign(new Error(silentClose.reason), { statusCode: 502 })
);
try {
for (const chunk of buildStreamErrorChunks(
silentClose.reason,
502,
streamController.clientResponseFormat
)) {
controller.enqueue(chunk);
}
} catch {
// downstream may have closed; original error already recorded
}
} else {
streamController.handleComplete();
}
try {
controller.close();
} catch {
// Expected: downstream may have already closed
}
return;
}
controller.enqueue(value);
noteClientChunk(value);
} catch (error) {
if (!streamController.isConnected()) {
try {
controller.close();
} catch {
// Expected: downstream may have already closed
}
return;
}
if (clientTerminalSeen) {
streamController.handleComplete();
try {
controller.close();
} catch {
// Expected: downstream may have already closed
}
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 {
// Closing an already-closed/aborted controller after client disconnect is expected.
}
}
},
async cancel(reason) {
const deferCompletedToolHandoff =
streamController.shouldDeferCompletedToolHandoff?.() === true;
if (clientTerminalSeen) {
streamController.handleComplete();
} else {
streamController.handleDisconnect(reason || "cancelled");
}
if (deferCompletedToolHandoff) return;
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 (e) {
console.debug(`[STREAM-HANDLER] stall watchdog handleError failed:`, e);
}
// 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 (e) {
console.debug(`[STREAM-HANDLER] stall watchdog upstream tap error failed:`, e);
}
// Abort the underlying fetch so upstream releases the connection.
try {
streamController.abort?.();
} catch (e) {
console.debug(`[STREAM-HANDLER] stall watchdog abort failed:`, e);
}
}, 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
);
}