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
67 changes: 49 additions & 18 deletions packages/core/src/aisdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,42 +23,73 @@ export interface LanguageEvent {
language?: LanguageModelV3
}

// Reports whether a chunk contains the start of a non-comment SSE line.
// Comment heartbeats (`: keepalive`) keep a stalled generation "warm" at the
// byte level without ever dispatching an event, so they must not extend the
// chunk timeout. Tracks the head byte of the current line so a field line
// split across chunk boundaries still counts as progress.
function sseFieldScanner() {
let atLineHead = true
return (chunk: Uint8Array) => {
let progress = false
let i = 0
while (i < chunk.length) {
if (atLineHead) {
const head = chunk[i]!
atLineHead = false
if (head !== 0x3a && head !== 0x0a && head !== 0x0d) progress = true
}
const nl = chunk.indexOf(0x0a, i)
if (nl === -1) break
i = nl + 1
atLineHead = true
}
return progress
}
}

function wrapSSE(res: Response, ms: number, ctl: AbortController) {
if (typeof ms !== "number" || ms <= 0) return res
if (!res.body) return res
if (!res.headers.get("content-type")?.includes("text/event-stream")) return res

const reader = res.body.getReader()
const progressed = sseFieldScanner()
let reject: ((error: unknown) => void) | undefined
let stall: ReturnType<typeof setTimeout> | undefined

const timeout = () => {
const error = new Error("SSE read timed out")
ctl.abort(error)
reader.cancel(error).catch(() => {})
reject?.(error)
}
const arm = () => {
clearTimeout(stall)
stall = setTimeout(timeout, ms)
}

const body = new ReadableStream<Uint8Array>({
async pull(ctrl) {
const part = await new Promise<Awaited<ReturnType<typeof reader.read>>>((resolve, reject) => {
const id = setTimeout(() => {
const err = new Error("SSE read timed out")
ctl.abort(err)
reader.cancel(err).catch(() => {})
reject(err)
}, ms)

reader.read().then(
(part) => {
clearTimeout(id)
resolve(part)
},
(err) => {
clearTimeout(id)
reject(err)
},
)
if (stall === undefined) arm()
const part = await new Promise<Awaited<ReturnType<typeof reader.read>>>((resolve, rejectPromise) => {
reject = rejectPromise
reader.read().then(resolve, rejectPromise)
}).finally(() => {
reject = undefined
})

if (part.done) {
clearTimeout(stall)
ctrl.close()
return
}

if (progressed(part.value)) arm()
ctrl.enqueue(part.value)
},
async cancel(reason) {
clearTimeout(stall)
ctl.abort(reason)
await reader.cancel(reason)
},
Expand Down
67 changes: 49 additions & 18 deletions packages/opencode/src/provider/provider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,42 +34,73 @@ import { ProviderError } from "./error"

const OPENAI_HEADER_TIMEOUT_DEFAULT = 300_000

// Reports whether a chunk contains the start of a non-comment SSE line.
// Comment heartbeats (`: keepalive`) keep a stalled generation "warm" at the
// byte level without ever dispatching an event, so they must not extend the
// chunk timeout. Tracks the head byte of the current line so a field line
// split across chunk boundaries still counts as progress.
function sseFieldScanner() {
let atLineHead = true
return (chunk: Uint8Array) => {
let progress = false
let i = 0
while (i < chunk.length) {
if (atLineHead) {
const head = chunk[i]!
atLineHead = false
if (head !== 0x3a && head !== 0x0a && head !== 0x0d) progress = true
}
const nl = chunk.indexOf(0x0a, i)
if (nl === -1) break
i = nl + 1
atLineHead = true
}
return progress
}
}

function wrapSSE(res: Response, ms: number, ctl: AbortController) {
if (typeof ms !== "number" || ms <= 0) return res
if (!res.body) return res
if (!res.headers.get("content-type")?.includes("text/event-stream")) return res

const reader = res.body.getReader()
const progressed = sseFieldScanner()
let reject: ((error: unknown) => void) | undefined
let stall: ReturnType<typeof setTimeout> | undefined

const timeout = () => {
const error = new ProviderError.ResponseStreamError("SSE read timed out")
ctl.abort(error)
reader.cancel(error).catch(() => {})
reject?.(error)
}
const arm = () => {
clearTimeout(stall)
stall = setTimeout(timeout, ms)
}

const body = new ReadableStream<Uint8Array>({
async pull(ctrl) {
const part = await new Promise<Awaited<ReturnType<typeof reader.read>>>((resolve, reject) => {
const id = setTimeout(() => {
const err = new ProviderError.ResponseStreamError("SSE read timed out")
ctl.abort(err)
reader.cancel(err).catch(() => {})
reject(err)
}, ms)

reader.read().then(
(part) => {
clearTimeout(id)
resolve(part)
},
(err) => {
clearTimeout(id)
reject(err)
},
)
if (stall === undefined) arm()
const part = await new Promise<Awaited<ReturnType<typeof reader.read>>>((resolve, rejectPromise) => {
reject = rejectPromise
reader.read().then(resolve, rejectPromise)
}).finally(() => {
reject = undefined
})

if (part.done) {
clearTimeout(stall)
ctrl.close()
return
}

if (progressed(part.value)) arm()
ctrl.enqueue(part.value)
},
async cancel(reason) {
clearTimeout(stall)
ctl.abort(reason)
await reader.cancel(reason)
},
Expand Down
50 changes: 50 additions & 0 deletions packages/opencode/test/provider/header-timeout.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,39 @@ it.live("configured chunkTimeout raises a retryable response stream error when S
}),
)

// Regression for #43519: SSE comment heartbeats (`: keepalive`) reset the
// per-chunk timer without carrying any event, so a provider or LB holding a
// stalled generation open with comments kept it alive forever. The timeout
// must only extend when a real SSE field line arrives.
it.live("chunkTimeout fires while a stalled stream is kept warm by comment heartbeats", () =>
Effect.gen(function* () {
const server = yield* Effect.acquireRelease(
Effect.promise(() => keepaliveBodyServer(10)),
(server) => Effect.sync(() => server.server.close()),
)

yield* provideTmpdirInstance(
() =>
Effect.gen(function* () {
const provider = yield* Provider.Service
const model = yield* provider.getModel(ProviderV2.ID.make("test"), ModelV2.ID.make("test-model"))
const result = streamText({
model: yield* provider.getLanguage(model),
onError() {},
messages: [{ role: "user", content: "hello" }],
})

const error = yield* Effect.promise(() => firstStreamError(result.fullStream))
expect(error).toBeInstanceOf(ProviderError.ResponseStreamError)
expect(
SessionRetry.retryable(MessageV2.fromError(error, { providerID: model.providerID }), model.providerID),
).toEqual({ message: "SSE read timed out" })
}),
{ config: providerConfig(server.url, { chunkTimeout: 100 }) },
)
}),
)

it.live("chunkTimeout can be disabled with false", () =>
Effect.gen(function* () {
const server = yield* Effect.acquireRelease(
Expand Down Expand Up @@ -390,6 +423,23 @@ async function delayedBodyServer(delay: number): Promise<{ server: Server; url:
return { server, url: `http://127.0.0.1:${address.port}` }
}

// Sends one partial SSE event, then holds the response open with `: keepalive`
// comments every `interval` ms — bytes keep flowing but no SSE field line ever
// completes the generation.
async function keepaliveBodyServer(interval: number): Promise<{ server: Server; url: string }> {
const server = createServer((_, res) => {
res.writeHead(200, { "content-type": "text/event-stream" })
res.flushHeaders()
res.write('data: {"choices":[{"delta":{"role":"assistant","content":""}}]}\n\n')
const timer = setInterval(() => res.write(": keepalive\n\n"), interval)
res.on("close", () => clearInterval(timer))
})
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve))
const address = server.address()
if (!address || typeof address === "string") throw new Error("server did not bind to a TCP port")
return { server, url: `http://127.0.0.1:${address.port}` }
}

function withAuthContent<A, E, R>(self: Effect.Effect<A, E, R>, value: Record<string, unknown> = defaultAuthContent()) {
return Effect.acquireUseRelease(
Effect.sync(() => {
Expand Down
Loading