From b9234e31edc25d9f1334b27aaf9b4dd9c918484f Mon Sep 17 00:00:00 2001 From: xin zhuang <65798732+1052326311@users.noreply.github.com> Date: Thu, 20 Aug 2026 16:42:48 +0800 Subject: [PATCH 1/3] fix(provider): ignore SSE comment heartbeats for chunk timeout --- packages/opencode/src/provider/provider.ts | 16 ++++++- .../test/provider/header-timeout.test.ts | 43 +++++++++++++++++++ 2 files changed, 58 insertions(+), 1 deletion(-) diff --git a/packages/opencode/src/provider/provider.ts b/packages/opencode/src/provider/provider.ts index 0200b212b21e..62a95060258b 100644 --- a/packages/opencode/src/provider/provider.ts +++ b/packages/opencode/src/provider/provider.ts @@ -40,15 +40,28 @@ function wrapSSE(res: Response, ms: number, ctl: AbortController) { if (!res.headers.get("content-type")?.includes("text/event-stream")) return res const reader = res.body.getReader() + const decoder = new TextDecoder() + let buffer = "" + let deadline: number | undefined + + function observedDataEvent(value: Uint8Array) { + buffer += decoder.decode(value, { stream: true }) + const events = buffer.split(/\r?\n\r?\n/) + buffer = events.pop() ?? "" + return events.some((event) => /^data\s*:/m.test(event)) + } + const body = new ReadableStream({ async pull(ctrl) { + if (deadline === undefined) deadline = Date.now() + ms const part = await new Promise>>((resolve, reject) => { + const remaining = Math.max(0, deadline - Date.now()) const id = setTimeout(() => { const err = new ProviderError.ResponseStreamError("SSE read timed out") ctl.abort(err) void reader.cancel(err) reject(err) - }, ms) + }, remaining) reader.read().then( (part) => { @@ -67,6 +80,7 @@ function wrapSSE(res: Response, ms: number, ctl: AbortController) { return } + if (observedDataEvent(part.value)) deadline = Date.now() + ms ctrl.enqueue(part.value) }, async cancel(reason) { diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index fc5ab04e108b..6dda537c7750 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -80,6 +80,36 @@ it.live("chunkTimeout raises a response stream error when SSE body stalls", () = }), ) +it.live("chunkTimeout ignores SSE comment heartbeats", () => + Effect.gen(function* () { + const server = yield* Effect.acquireRelease( + Effect.promise(() => keepaliveBodyServer(20)), + (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(async () => { + for await (const part of result.fullStream) { + if (part.type === "error") return part.error + } + }) + expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) + }), + { config: providerConfig(server.url, { chunkTimeout: 50 }) }, + ) + }), +) + it.live("headerTimeout aborts when response headers do not arrive", () => Effect.gen(function* () { const server = yield* Effect.acquireRelease( @@ -211,6 +241,19 @@ async function delayedBodyServer(delay: number): Promise<{ server: Server; url: return { server, url: `http://127.0.0.1:${address.port}` } } +async function keepaliveBodyServer(interval: number): Promise<{ server: Server; url: string }> { + const server = createServer((_, res) => { + res.writeHead(200, { "content-type": "text/event-stream" }) + res.write('data: {"choices":[{"delta":{"content":"partial"}}]}\n\n') + const keepalive = setInterval(() => res.write(": keepalive\n\n"), interval) + res.on("close", () => clearInterval(keepalive)) + }) + await new Promise((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(self: Effect.Effect, value: Record = defaultAuthContent()) { return Effect.acquireUseRelease( Effect.sync(() => { From 0ffeeb9b4bfdbfb51efc0a92a2179d42f5856143 Mon Sep 17 00:00:00 2001 From: 1052326311 <65798732+1052326311@users.noreply.github.com> Date: Fri, 21 Aug 2026 12:13:56 +0800 Subject: [PATCH 2/3] test(provider): validate SSE timeout error --- packages/opencode/src/provider/provider.ts | 5 +++-- packages/opencode/test/provider/header-timeout.test.ts | 8 ++++++-- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/packages/opencode/src/provider/provider.ts b/packages/opencode/src/provider/provider.ts index 62a95060258b..10c322f9a07e 100644 --- a/packages/opencode/src/provider/provider.ts +++ b/packages/opencode/src/provider/provider.ts @@ -53,9 +53,10 @@ function wrapSSE(res: Response, ms: number, ctl: AbortController) { const body = new ReadableStream({ async pull(ctrl) { - if (deadline === undefined) deadline = Date.now() + ms + const current = deadline ?? Date.now() + ms + deadline = current const part = await new Promise>>((resolve, reject) => { - const remaining = Math.max(0, deadline - Date.now()) + const remaining = Math.max(0, current - Date.now()) const id = setTimeout(() => { const err = new ProviderError.ResponseStreamError("SSE read timed out") ctl.abort(err) diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index 6dda537c7750..e0d712231a43 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -99,8 +99,12 @@ it.live("chunkTimeout ignores SSE comment heartbeats", () => }) const error = yield* Effect.promise(async () => { - for await (const part of result.fullStream) { - if (part.type === "error") return part.error + try { + for await (const part of result.fullStream) { + if (part.type === "error") return part.error + } + } catch (error) { + return error } }) expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) From 8de252f7aed6b8fc4d93a650be397fd385c7c433 Mon Sep 17 00:00:00 2001 From: 1052326311 <65798732+1052326311@users.noreply.github.com> Date: Fri, 21 Aug 2026 12:18:25 +0800 Subject: [PATCH 3/3] test(provider): return after closed timeout stream --- packages/opencode/test/provider/header-timeout.test.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index e0d712231a43..100450fdf043 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -106,6 +106,7 @@ it.live("chunkTimeout ignores SSE comment heartbeats", () => } catch (error) { return error } + return undefined }) expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) }),