Skip to content
Open
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
21 changes: 11 additions & 10 deletions packages/ai/src/route/transport/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -120,16 +120,17 @@ export const httpJson = <Body, Frame>(input: HttpJsonInput<Body, Frame>): 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))),
),
),
),
Expand Down
50 changes: 50 additions & 0 deletions packages/ai/test/http-timeout.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>()
const encoder = new TextEncoder()
let emit = (_text: string) => {}
const layer = dynamicResponse((input) =>
Effect.sync(() =>
input.respond(
new ReadableStream<Uint8Array>({
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* () {
Expand Down Expand Up @@ -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(
Expand Down
Loading