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(