Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions open-sse/handlers/chatCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5940,6 +5940,11 @@ export async function handleChatCore({
clientResponseFormat,
echoModel,
responseHeaders,
// Same adaptive budget the pre-handoff readiness gate above just used —
// reasoning models that legitimately take a while to say anything keep
// that same patience for their first REAL content, not just their first
// lifecycle frame. See pipeWithDisconnect's own doc comment.
contentStallTimeoutMs: streamReadinessPolicy.timeoutMs,
});

// ── Gamification event (fire-and-forget) ──
Expand Down
9 changes: 8 additions & 1 deletion open-sse/handlers/chatCore/streamingPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,12 @@ export function assembleStreamingPipeline(
clientResponseFormat: Parameters<typeof defaultShape>[0];
echoModel: string | null | undefined;
responseHeaders: Record<string, string>;
/** See pipeWithDisconnect's own doc comment — the post-handoff "stream is
* open but never produced real output" watchdog. Threaded from the same
* adaptive streamReadinessPolicy.timeoutMs already computed for this
* request's pre-handoff readiness gate, so slow-first-content reasoning
* models keep the same generous budget in both phases. */
contentStallTimeoutMs?: number;
},
deps: StreamingPipelineDeps = DEFAULT_DEPS
) {
Expand All @@ -83,7 +89,8 @@ export function assembleStreamingPipeline(
let piiStream = deps.pipeWithDisconnect(
args.providerResponse,
args.transformStream,
args.streamController
args.streamController,
{ contentStallTimeoutMs: args.contentStallTimeoutMs }
);
if (typeof args.createPiiTransform === "function") {
piiStream = piiStream.pipeThrough((args.createPiiTransform as () => TransformStream)());
Expand Down
105 changes: 91 additions & 14 deletions open-sse/utils/streamHandler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -857,12 +857,19 @@ export function pipeWithDisconnect(
providerResponse: Response,
transformStream: TransformStream<Uint8Array, Uint8Array>,
streamController: StreamController,
opts: { stallTimeoutMs?: number; highWaterMark?: number } = {}
opts: { stallTimeoutMs?: number; contentStallTimeoutMs?: number; highWaterMark?: number } = {}
) {
const stallTimeoutMs = opts.stallTimeoutMs ?? DEFAULT_STREAM_STALL_TIMEOUT_MS;

// Watchdog disabled — preserve legacy behavior verbatim.
if (!stallTimeoutMs || stallTimeoutMs <= 0) {
// Disabled unless a caller opts in with an explicit budget (chatCore wires
// the adaptive streamReadinessPolicy.timeoutMs — see its own doc comment).
// No blanket default here: an arbitrary constant picked at this layer,
// without the request's actual model/provider/payload context, would risk
// false-stalling the exact slow-first-content reasoning models the
// readiness policy already knows to grant more patience.
const contentStallTimeoutMs = opts.contentStallTimeoutMs ?? 0;

// Watchdogs disabled — preserve legacy behavior verbatim.
if ((!stallTimeoutMs || stallTimeoutMs <= 0) && contentStallTimeoutMs <= 0) {
const transformedBody = providerResponse.body.pipeThrough(transformStream);
return createDisconnectAwareStream(
{ readable: transformedBody, writable: createNoopAbortWritable() },
Expand Down Expand Up @@ -891,6 +898,7 @@ export function pipeWithDisconnect(
}
};
const armStall = () => {
if (!stallTimeoutMs || stallTimeoutMs <= 0) return;
clearStall();
stallTimer = setTimeout(() => {
stallTimer = null;
Expand Down Expand Up @@ -919,48 +927,117 @@ export function pipeWithDisconnect(
}, stallTimeoutMs);
};

// Wrap controller so every termination path clears the stall timer.
// Without this, abort/complete/error/disconnect paths leave the timer armed
// Second, independent watchdog: fires when the upstream keeps sending raw
// bytes (so armStall() above keeps resetting and never fires) but none of
// them ever carry real model output — only lifecycle/ping frames
// (OpenAI Responses response.in_progress/response.created, bare
// role-only start chunks, etc). ensureStreamReadiness's own gate already
// treats any one of those as "ready" and hands the connection off (see its
// own doc comment and kiro.ts's deliberate early role-only chunk — several
// providers rely on that fast handoff for UX, so tightening readiness
// itself would regress them). Once handed off there was previously nothing
// watching whether the model ever actually said anything: a stalled free
// OpenRouter model (e.g. minimax-m3:free under load) could stream nothing
// but response.in_progress pings indefinitely, relayed byte-for-byte to
// the client, until the CLIENT's own idle timeout eventually gave up --
// sometimes 120s, sometimes 900s depending on the calling task, always
// slower and less informative than OmniRoute failing this attempt itself
// with a clear error the client's own retry/fallback logic can react to
// immediately. Armed ONCE at stream start (not re-armed by lifecycle-only
// bytes, unlike armStall above) and cleared permanently the first time
// real content is observed -- reuses the exact classifier
// (createStreamContentWatcher) createDisconnectAwareStream already trusts
// for its own end-of-stream #8649 empty-content check.
let contentStallTimer: ReturnType<typeof setTimeout> | null = null;
let contentStallFired = false;
const upstreamContentWatcher = createStreamContentWatcher();
const upstreamContentDecoder = new TextDecoder();

const clearContentStall = () => {
if (contentStallTimer) {
clearTimeout(contentStallTimer);
contentStallTimer = null;
}
};
const armContentStall = () => {
if (contentStallTimeoutMs <= 0) return;
contentStallTimer = setTimeout(() => {
contentStallTimer = null;
contentStallFired = true;
const stallError = new Error(
`stream content stall: no model output within ${contentStallTimeoutMs}ms (lifecycle/heartbeat events only)`
);
try {
streamController.handleError?.(stallError);
} catch (e) {
console.debug(`[STREAM-HANDLER] content stall watchdog handleError failed:`, e);
}
try {
upstreamTapController?.error(stallError);
} catch (e) {
console.debug(`[STREAM-HANDLER] content stall watchdog upstream tap error failed:`, e);
}
try {
streamController.abort?.();
} catch (e) {
console.debug(`[STREAM-HANDLER] content stall watchdog abort failed:`, e);
}
}, contentStallTimeoutMs);
};

// Wrap controller so every termination path clears both stall timers.
// Without this, abort/complete/error/disconnect paths leave a timer armed
// and a stale abort could fire after the request has already ended.
const wrappedController: StreamController = {
...streamController,
handleComplete: () => {
clearStall();
clearContentStall();
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;
clearContentStall();
// A 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 || contentStallFired) return;
streamController.handleError(e);
},
handleDisconnect: (reason?: string) => {
clearStall();
clearContentStall();
streamController.handleDisconnect(reason);
},
abort: () => {
clearStall();
clearContentStall();
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.
// Inert tap that resets the byte-stall timer on every raw upstream chunk
// and (independently) clears the content-stall timer the first time a
// chunk carries real output. 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 either watchdog.
const upstreamTap = new TransformStream<Uint8Array, Uint8Array>({
start(controller) {
upstreamTapController = controller;
armStall();
armContentStall();
},
transform(chunk, controller) {
armStall();
if (contentStallTimeoutMs > 0 && !upstreamContentWatcher.sawContent()) {
upstreamContentWatcher.note(upstreamContentDecoder.decode(chunk, { stream: true }));
if (upstreamContentWatcher.sawContent()) clearContentStall();
}
controller.enqueue(chunk);
},
flush() {
clearStall();
clearContentStall();
},
});

Expand Down
130 changes: 130 additions & 0 deletions tests/unit/stream-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -830,3 +830,133 @@ test("pipeWithDisconnect stall watchdog does not fire after normal stream comple
assert.equal(text, "ok");
assert.equal(onErrorCalled, false, "stall watchdog must be cleared on stream completion");
});

// Content stall: a live incident (2026-09-04) where a free OpenRouter model
// (minimax-m3:free, under load) streamed nothing but OpenAI Responses
// response.in_progress heartbeats for ~90s -- real bytes kept arriving
// (so the byte-stall watchdog above never fired) but never any actual model
// output. ensureStreamReadiness's own pre-handoff gate treats a bare
// response.in_progress as "ready" and hands the connection straight to the
// client (deliberately -- kiro.ts's role-only start chunk relies on the same
// behavior), so nothing was left watching the ALREADY-HANDED-OFF stream for
// whether it ever said anything. The client's own idle timeout (120s-900s
// depending on which internal task made the call) was the only thing that
// eventually noticed.
test("pipeWithDisconnect flags a stream that only ever produces lifecycle/heartbeat frames, never real content", async () => {
const source = new ReadableStream({
start(controller) {
controller.enqueue(
encoder.encode(
'data: {"type":"response.in_progress","response":{"id":"resp_1","status":"in_progress","output":[]}}\n\n'
)
);
},
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, {
// Byte-stall stays generous (nothing keeps this test's single heartbeat
// frame alive) — only the content-stall budget under test is tight.
stallTimeoutMs: 5000,
contentStallTimeoutMs: 80,
});

const text = await readStreamText(stream);

assert.ok(
onErrorEvent !== null,
"content stall watchdog must fire when the stream never produces real output"
);
assert.match(onErrorEvent.message, /content stall/i);
assert.match(text, /content stall/i);
assert.match(text, /"finish_reason":"error"/);
});

test("pipeWithDisconnect does NOT flag a stream that starts with lifecycle frames but then produces real content", async () => {
const source = new ReadableStream({
async start(controller) {
controller.enqueue(
encoder.encode('data: {"type":"response.in_progress","response":{"id":"resp_1"}}\n\n')
);
await new Promise((r) => setTimeout(r, 30));
controller.enqueue(
encoder.encode('data: {"type":"response.output_text.delta","delta":"Hello"}\n\n')
);
controller.close();
},
});

let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});

const stream = pipeWithDisconnect(new Response(source), new TransformStream(), streamController, {
stallTimeoutMs: 5000,
// Longer than the 30ms gap to the real content chunk, short enough that
// a false-firing watchdog would still trip well before readStreamText resolves.
contentStallTimeoutMs: 200,
});

const text = await readStreamText(stream);

assert.equal(
onErrorCalled,
false,
"content stall watchdog must not fire once real content arrives"
);
assert.doesNotMatch(text, /content stall/i);
assert.match(text, /Hello/);
});

test("pipeWithDisconnect content stall watchdog is off by default (no contentStallTimeoutMs)", async () => {
// Same lifecycle-only-forever shape as the firing test above, but no
// contentStallTimeoutMs opt-in — legacy behavior (relay forever) preserved.
const source = new ReadableStream({
start(controller) {
controller.enqueue(
encoder.encode('data: {"type":"response.in_progress","response":{"id":"resp_1"}}\n\n')
);
// never closes — the point is that with the watchdog off, nothing here
// should time out on its own; the test just proves no premature error.
},
cancel() {},
});

let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});

const stream = pipeWithDisconnect(new Response(source), new TransformStream(), streamController, {
stallTimeoutMs: 5000,
// contentStallTimeoutMs intentionally omitted.
});

const reader = stream.getReader();
const { value } = await reader.read();
await reader.cancel("test done");

assert.ok(value, "the lifecycle frame must still be relayed");
assert.equal(
onErrorCalled,
false,
"content stall watchdog must stay off without an explicit budget"
);
});
Loading