From 5479ac7d4998415edaca60e842ecea146b76d1d8 Mon Sep 17 00:00:00 2001 From: Markus Hartung Date: Fri, 4 Sep 2026 17:48:59 +0200 Subject: [PATCH] fix(streaming): fail fast when an upstream stream produces only lifecycle/heartbeat events, never real content Live incident (2026-09-04): a free-tier OpenRouter model (minimax-m3:free, under load) streamed nothing but OpenAI Responses response.in_progress heartbeats for ~90 seconds before OpenClaw's own client-side idle timeout (120s for that particular internal task, 900s for the main agent) finally gave up. OmniRoute's own ensureStreamReadiness pre-handoff gate correctly treats a bare response.in_progress as "ready" (deliberately -- kiro.ts's early role-only start chunk relies on the exact same behavior for UX, so tightening readiness itself would regress that), hands the connection to the client, and from that point on nothing was watching whether the model ever actually produced real output. The existing byte-level stall watchdog in pipeWithDisconnect (ported from decolua/9router#1243) doesn't catch this either -- it deliberately tracks raw upstream byte activity (not transform output) to avoid false-positiving on reasoning models, and the heartbeat bytes keep it happy indefinitely. Adds a second, independent watchdog to pipeWithDisconnect: armed once at stream start, cleared permanently the first time real model output is observed (reusing createStreamContentWatcher, the exact classifier createDisconnectAwareStream already trusts for its own end-of-stream #8649 empty-content check), and firing only if the deadline elapses with lifecycle/ping frames only. Off by default (no arbitrary constant picked at this layer); chatCore wires it to the same adaptive streamReadinessPolicy.timeoutMs already computed for the pre-handoff readiness gate, so slow-first-content reasoning models keep the same generous, request-specific budget in both phases. A genuine invisible mid-stream provider swap is architecturally impossible over HTTP SSE passthrough once headers have been sent to the client -- this cannot make combo retry with a different model. What it does do: turn an indefinite silent hang into a fast, explicit, well-formed SSE error the calling agent's own existing retry/fallback logic can react to immediately, instead of waiting out whichever client-side timeout happened to apply to that specific call. --- open-sse/handlers/chatCore.ts | 5 + .../handlers/chatCore/streamingPipeline.ts | 9 +- open-sse/utils/streamHandler.ts | 105 ++++++++++++-- tests/unit/stream-handler.test.ts | 130 ++++++++++++++++++ 4 files changed, 234 insertions(+), 15 deletions(-) diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 622e9340840..c683d0a3d51 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -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) ── diff --git a/open-sse/handlers/chatCore/streamingPipeline.ts b/open-sse/handlers/chatCore/streamingPipeline.ts index bbe0bdcb8ba..b75c96e1165 100644 --- a/open-sse/handlers/chatCore/streamingPipeline.ts +++ b/open-sse/handlers/chatCore/streamingPipeline.ts @@ -69,6 +69,12 @@ export function assembleStreamingPipeline( clientResponseFormat: Parameters[0]; echoModel: string | null | undefined; responseHeaders: Record; + /** 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 ) { @@ -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)()); diff --git a/open-sse/utils/streamHandler.ts b/open-sse/utils/streamHandler.ts index 841346bcfe0..8499a95b800 100644 --- a/open-sse/utils/streamHandler.ts +++ b/open-sse/utils/streamHandler.ts @@ -857,12 +857,19 @@ export function pipeWithDisconnect( providerResponse: Response, transformStream: TransformStream, 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() }, @@ -891,6 +898,7 @@ export function pipeWithDisconnect( } }; const armStall = () => { + if (!stallTimeoutMs || stallTimeoutMs <= 0) return; clearStall(); stallTimer = setTimeout(() => { stallTimer = null; @@ -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 | 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({ 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(); }, }); diff --git a/tests/unit/stream-handler.test.ts b/tests/unit/stream-handler.test.ts index c358ff7c2a2..3f2f988e4ce 100644 --- a/tests/unit/stream-handler.test.ts +++ b/tests/unit/stream-handler.test.ts @@ -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" + ); +});