From 359621b8430714882e3cfd8f12a99a58f1a70068 Mon Sep 17 00:00:00 2001 From: felixzsh Date: Fri, 2 Oct 2026 17:04:06 -0500 Subject: [PATCH] fix(ai): bound native stream stalls on framed events The HTTP transport applied the chunk timeout to the raw response byte stream before framing. Framing drops SSE comment keepalives and other empty events, so any such byte reset the stall timer and a provider could hold a stalled generation warm forever (#43519). This matched the native path but not the AI SDK path, whose wrapSSE resets only on parsed events. Bound the timeout on framed progress instead. A stream that only sends keepalives now still trips the chunk timeout, while real events keep resetting it. --- packages/ai/src/route/transport/http.ts | 21 ++++++----- packages/ai/test/http-timeout.test.ts | 50 +++++++++++++++++++++++++ 2 files changed, 61 insertions(+), 10 deletions(-) diff --git a/packages/ai/src/route/transport/http.ts b/packages/ai/src/route/transport/http.ts index 5e330f61d5c4..f6f709ea7117 100644 --- a/packages/ai/src/route/transport/http.ts +++ b/packages/ai/src/route/transport/http.ts @@ -120,16 +120,17 @@ export const httpJson = (input: HttpJsonInput): HttpJs const http = RequestExecutor.responseHttp(response) const remaining = Duration.subtract(total, Duration.millis((yield* Clock.currentTimeMillis) - started)) return { - frames: prepared.framing.frame( - RequestExecutor.responseStream(response).pipe( - Stream.timeoutOrElse({ - duration: timeoutDuration(request.http?.chunkTimeout), - orElse: () => Stream.fail(timeout("read", "Timed out waiting for response data", http)), - }), - Stream.interruptWhen( - Effect.sleep(remaining).pipe( - Effect.andThen(Effect.fail(timeout("read", "Timed out waiting for the response to complete", http))), - ), + // Bound the stall on framed progress, not raw bytes. Framing drops SSE + // comment keepalives (`: keepalive`) and other empty events, so counting + // bytes would let a provider hold a stalled generation warm forever. + frames: prepared.framing.frame(RequestExecutor.responseStream(response)).pipe( + Stream.timeoutOrElse({ + duration: timeoutDuration(request.http?.chunkTimeout), + orElse: () => Stream.fail(timeout("read", "Timed out waiting for response data", http)), + }), + Stream.interruptWhen( + Effect.sleep(remaining).pipe( + Effect.andThen(Effect.fail(timeout("read", "Timed out waiting for the response to complete", http))), ), ), ), diff --git a/packages/ai/test/http-timeout.test.ts b/packages/ai/test/http-timeout.test.ts index 6b45a8070079..88cf21f4a0f8 100644 --- a/packages/ai/test/http-timeout.test.ts +++ b/packages/ai/test/http-timeout.test.ts @@ -49,6 +49,31 @@ const stalledServer = Effect.gen(function* () { return { layer, stalled, resume: () => resume() } }) +// Sends one content chunk, then only what the test asks for. Lets a test prove +// that bytes framing drops, such as SSE comment keepalives, are not progress. +const heartbeatServer = Effect.gen(function* () { + const started = yield* Deferred.make() + const encoder = new TextEncoder() + let emit = (_text: string) => {} + const layer = dynamicResponse((input) => + Effect.sync(() => + input.respond( + new ReadableStream({ + start(controller) { + emit = (text) => controller.enqueue(encoder.encode(text)) + emit(sseRaw(`data: ${JSON.stringify(deltaChunk({ content: "Hi" }))}`)) + }, + pull() { + Deferred.doneUnsafe(started, Effect.void) + }, + }), + SSE, + ), + ), + ) + return { layer, started, emit: (text: string) => emit(text) } +}) + describe("HTTP transport timeouts", () => { it.effect("fails when response headers take longer than five minutes", () => Effect.gen(function* () { @@ -87,6 +112,31 @@ describe("HTTP transport timeouts", () => { }), ) + it.effect("ignores SSE comment keepalives when bounding a stalled stream", () => + Effect.gen(function* () { + const server = yield* heartbeatServer + const fiber = yield* LLMClient.generate( + LLM.request({ model, prompt: "Hello", http: { chunkTimeout: 20_000 } }), + ).pipe(Effect.provide(server.layer), Effect.flip, Effect.forkChild({ startImmediately: true })) + yield* Deferred.await(server.started) + yield* Effect.yieldNow + // A comment every 10s keeps the socket warm without carrying progress, so + // the 20s chunk bound must still fire. + yield* TestClock.adjust("10 seconds") + server.emit(sseRaw(": keepalive")) + yield* Effect.yieldNow + yield* TestClock.adjust("10 seconds") + const error = yield* Fiber.join(fiber) + + expect(error.reason).toMatchObject({ + _tag: "Transport", + transport: "http", + operation: "read", + code: "Timeout", + }) + }), + ) + it.effect("applies a configured header timeout", () => Effect.gen(function* () { const fiber = yield* LLMClient.generate(