From 792cc4b8614b7f454a401a0136c8038f6d93e1bd Mon Sep 17 00:00:00 2001 From: luvs01 Date: Sat, 8 Aug 2026 23:58:29 +0900 Subject: [PATCH] fix upstream response body stall timeout --- src/lib/response-body-inactivity.ts | 71 ++++++++++++++++++++++++++ src/server/index.ts | 1 + src/server/responses/core.ts | 6 +++ tests/response-body-inactivity.test.ts | 42 +++++++++++++++ 4 files changed, 120 insertions(+) create mode 100644 src/lib/response-body-inactivity.ts create mode 100644 tests/response-body-inactivity.test.ts diff --git a/src/lib/response-body-inactivity.ts b/src/lib/response-body-inactivity.ts new file mode 100644 index 000000000..502466c17 --- /dev/null +++ b/src/lib/response-body-inactivity.ts @@ -0,0 +1,71 @@ +import { idleDeadline } from "./abort"; + +/** + * Bound response-body silence after fetch has returned its headers. The deadline is armed only + * while an upstream read is pending, so a downstream client applying backpressure does not count + * as an upstream stall. Non-empty chunks reset the window; normal completion leaves the upstream + * controller untouched. + */ +export function guardResponseBodyInactivity( + response: Response, + upstream: AbortController, + inactivityMs: number, +): Response { + if (!response.body || inactivityMs <= 0) return response; + + const reader = response.body.getReader(); + let settled = false; + let output: ReadableStreamDefaultController | undefined; + const timeoutError = new DOMException( + `Upstream response body stalled for ${inactivityMs}ms`, + "TimeoutError", + ); + const idle = idleDeadline(inactivityMs, () => { + if (settled) return; + settled = true; + upstream.abort(timeoutError); + reader.cancel(timeoutError).catch(() => {}); + try { output?.error(timeoutError); } catch { /* already torn down */ } + }); + + const body = new ReadableStream({ + async pull(controller) { + output = controller; + try { + idle.reset(); + for (;;) { + const { done, value } = await reader.read(); + if (settled) return; + if (done) { + settled = true; + idle.cancel(); + controller.close(); + return; + } + if (value.byteLength === 0) continue; + idle.pause(); + controller.enqueue(value); + return; + } + } catch (error) { + if (settled) return; + settled = true; + idle.cancel(); + try { controller.error(error); } catch { /* already torn down */ } + } + }, + cancel(reason) { + if (settled) return; + settled = true; + idle.cancel(); + upstream.abort(reason); + reader.cancel(reason).catch(() => {}); + }, + }); + + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }); +} diff --git a/src/server/index.ts b/src/server/index.ts index 7d08fc422..f6307d1f2 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -1,3 +1,4 @@ +export { guardResponseBodyInactivity } from "../lib/response-body-inactivity"; import { markActivity } from "../lib/sidecar-tracker"; import { buildWarmupCompletionFrames, diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 2d00ab0c4..7f4cd3e22 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -117,6 +117,7 @@ import type { WsData } from "../ws-bridge"; import { codexAccountSelectionForTurn, registerTurn, trackStreamLifetime, unregisterTurn } from "../lifecycle"; import { redactSecretString } from "../../lib/redact"; import { readBoundedResponseBody } from "../../lib/bounded-body"; +import { guardResponseBodyInactivity } from "../../lib/response-body-inactivity"; import type { AdmissionLease } from "../../lib/admission"; import { supportedLadderFor } from "../effort-policy"; import { isThreadSpawnRequest } from "../effort-policy"; @@ -1783,6 +1784,9 @@ async function handleResponsesInner( const upstream = new AbortController(); linkAbortSignal(upstream, options.abortSignal); const connectMs = config.connectTimeoutMs ?? 200_000; + const bodyInactivityMs = typeof config.stallTimeoutSec === "number" && Number.isFinite(config.stallTimeoutSec) && config.stallTimeoutSec > 0 + ? Math.floor(config.stallTimeoutSec * 1000) + : 300_000; let upstreamResponse: Response; const transportFailureResponse = (err: unknown): Response => { upstream.abort(); @@ -1941,6 +1945,7 @@ async function handleResponsesInner( } } } + upstreamResponse = guardResponseBodyInactivity(upstreamResponse, upstream, bodyInactivityMs); const headers = sanitizePassthroughHeaders(upstreamResponse.headers); const resolvedModel = headers.get("openai-model")?.trim(); if (resolvedModel) logCtx.resolvedModel = resolvedModel; @@ -2957,6 +2962,7 @@ async function handleResponsesInner( } break; } + upstreamResponse = guardResponseBodyInactivity(upstreamResponse, upstream, stallTimeoutMs); if (!upstreamResponse.ok) { if (options.comboAttempt) { const failure = await consumeComboFailure(upstreamResponse, options.abortSignal) diff --git a/tests/response-body-inactivity.test.ts b/tests/response-body-inactivity.test.ts new file mode 100644 index 000000000..d2bb20052 --- /dev/null +++ b/tests/response-body-inactivity.test.ts @@ -0,0 +1,42 @@ +import { describe, expect, test } from "bun:test"; +import { guardResponseBodyInactivity } from "../src/lib/response-body-inactivity"; + +describe("upstream response body inactivity", () => { + test("guards both passthrough and adapter response bodies after headers", async () => { + const core = await Bun.file(new URL("../src/server/responses/core.ts", import.meta.url)).text(); + expect(core).toContain("guardResponseBodyInactivity(upstreamResponse, upstream, bodyInactivityMs)"); + expect(core).toContain("guardResponseBodyInactivity(upstreamResponse, upstream, stallTimeoutMs)"); + }); + + test("aborts a headers-only upstream", async () => { + const upstream = new AbortController(); + const response = guardResponseBodyInactivity(new Response( + new ReadableStream({ pull() { return new Promise(() => {}); } }), + ), upstream, 20); + + await expect(response.text()).rejects.toMatchObject({ name: "TimeoutError" }); + expect(upstream.signal.aborted).toBe(true); + expect(upstream.signal.reason).toMatchObject({ name: "TimeoutError" }); + }); + + test("resets on bytes and ignores downstream backpressure", async () => { + const encoder = new TextEncoder(); + const upstream = new AbortController(); + let controller!: ReadableStreamDefaultController; + const guarded = guardResponseBodyInactivity(new Response(new ReadableStream({ + start(value) { controller = value; }, + })), upstream, 30); + const reader = guarded.body!.getReader(); + + controller.enqueue(encoder.encode("one")); + controller.enqueue(encoder.encode("two")); + expect(new TextDecoder().decode((await reader.read()).value)).toBe("one"); + await Bun.sleep(45); + expect(upstream.signal.aborted).toBe(false); + + expect(new TextDecoder().decode((await reader.read()).value)).toBe("two"); + controller.close(); + expect((await reader.read()).done).toBe(true); + expect(upstream.signal.aborted).toBe(false); + }); +});