From 8385ea2812d08a618024dec8283bb92f635ec72b Mon Sep 17 00:00:00 2001 From: Daniel Cadenas Date: Fri, 31 Jul 2026 00:13:58 -0300 Subject: [PATCH 1/2] fix(provider): recover stalled model streams --- packages/core/src/v1/config/provider.ts | 4 +- packages/opencode/src/provider/error.ts | 12 +- packages/opencode/src/provider/provider.ts | 73 ++----- .../opencode/src/provider/stream-liveness.ts | 110 +++++++++++ packages/opencode/src/session/message-v2.ts | 12 ++ .../test/provider/header-timeout.test.ts | 62 +++++- .../test/provider/stream-liveness.test.ts | 178 ++++++++++++++++++ .../test/session/processor-effect.test.ts | 61 +++++- packages/opencode/test/session/retry.test.ts | 14 ++ 9 files changed, 463 insertions(+), 63 deletions(-) create mode 100644 packages/opencode/src/provider/stream-liveness.ts create mode 100644 packages/opencode/test/provider/stream-liveness.test.ts diff --git a/packages/core/src/v1/config/provider.ts b/packages/core/src/v1/config/provider.ts index f860a2b4ac7e..c6540ae6cd3d 100644 --- a/packages/core/src/v1/config/provider.ts +++ b/packages/core/src/v1/config/provider.ts @@ -114,9 +114,9 @@ export const Info = Schema.Struct({ description: "Timeout in milliseconds to wait for response headers. Provider integrations may set defaults. Set to false to disable timeout.", }), - chunkTimeout: Schema.optional(PositiveInt).annotate({ + chunkTimeout: Schema.optional(Schema.Union([Schema.Finite, Schema.Literal(false)])).annotate({ description: - "Timeout in milliseconds between streamed SSE chunks for this provider. If no chunk arrives within this window, the request is aborted.", + "Timeout in milliseconds between streamed SSE chunks for this provider. Set to false, zero, or a negative number to disable the timeout.", }), }), [Schema.Record(Schema.String, Schema.Any)], diff --git a/packages/opencode/src/provider/error.ts b/packages/opencode/src/provider/error.ts index 21149a2cf389..5f4811b48067 100644 --- a/packages/opencode/src/provider/error.ts +++ b/packages/opencode/src/provider/error.ts @@ -13,13 +13,23 @@ export class HeaderTimeoutError extends Error { } export class ResponseStreamError extends Error { - public override readonly name = "ProviderResponseStreamError" + public override readonly name: string = "ProviderResponseStreamError" constructor(message: string, options?: ErrorOptions) { super(message, options) } } +export class ResponseStreamTimeoutError extends ResponseStreamError { + public override readonly name = "ProviderResponseStreamTimeoutError" + + constructor(public readonly ms: number) { + super( + `Provider response stream timed out after ${ms}ms of inactivity. Check provider connectivity or increase provider.options.chunkTimeout before retrying.`, + ) + } +} + function isOpenAiErrorRetryable(e: APICallError) { const status = e.statusCode if (!status) return e.isRetryable diff --git a/packages/opencode/src/provider/provider.ts b/packages/opencode/src/provider/provider.ts index 85b7fd3978e8..43538da2fd3a 100644 --- a/packages/opencode/src/provider/provider.ts +++ b/packages/opencode/src/provider/provider.ts @@ -31,57 +31,10 @@ import { ModelV2 } from "@opencode-ai/core/model" import { ModelStatus } from "./model-status" import { RuntimeFlags } from "@/effect/runtime-flags" import { ProviderError } from "./error" +import { StreamLiveness } from "./stream-liveness" const OPENAI_HEADER_TIMEOUT_DEFAULT = 300_000 -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 body = new ReadableStream({ - async pull(ctrl) { - const part = await new Promise>>((resolve, reject) => { - const id = setTimeout(() => { - const err = new ProviderError.ResponseStreamError("SSE read timed out") - ctl.abort(err) - void reader.cancel(err) - reject(err) - }, ms) - - reader.read().then( - (part) => { - clearTimeout(id) - resolve(part) - }, - (err) => { - clearTimeout(id) - reject(err) - }, - ) - }) - - if (part.done) { - ctrl.close() - return - } - - ctrl.enqueue(part.value) - }, - async cancel(reason) { - ctl.abort(reason) - await reader.cancel(reason) - }, - }) - - return new Response(body, { - headers: new Headers(res.headers), - status: res.status, - statusText: res.statusText, - }) -} - function timeoutController(ms: number) { const ctl = new AbortController() const id = setTimeout(() => ctl.abort(new ProviderError.HeaderTimeoutError(ms)), ms) @@ -1170,6 +1123,7 @@ interface State { sdk: Map modelLoaders: Record varsLoaders: Record + streamLiveness: StreamLiveness.Detector } export class Service extends Context.Service()("@opencode/Provider") {} @@ -1357,6 +1311,7 @@ const layer = Layer.effect( [providerID: string]: CustomVarsLoader } = {} const sdk = new Map() + const streamLiveness = StreamLiveness.create() const discoveryLoaders: { [providerID: string]: CustomDiscoverModels } = {} @@ -1664,6 +1619,7 @@ const layer = Layer.effect( sdk, modelLoaders, varsLoaders, + streamLiveness, } }), ) @@ -1735,7 +1691,13 @@ const layer = Layer.effect( if (existing) return existing const customFetch = options["fetch"] - const chunkTimeout = options["chunkTimeout"] + const configuredChunkTimeout = options["chunkTimeout"] + const fixedChunkTimeout = + typeof configuredChunkTimeout === "number" && Number.isFinite(configuredChunkTimeout) + ? configuredChunkTimeout + : undefined + const streamTimeoutDisabled = + configuredChunkTimeout === false || (fixedChunkTimeout !== undefined && fixedChunkTimeout <= 0) const headerTimeout = options["headerTimeout"] delete options["chunkTimeout"] delete options["headerTimeout"] @@ -1743,13 +1705,13 @@ const layer = Layer.effect( options["fetch"] = async (input: any, init?: BunFetchRequestInit) => { const fetchFn = customFetch ?? fetch const opts = init ?? {} - const chunkAbortCtl = typeof chunkTimeout === "number" && chunkTimeout > 0 ? new AbortController() : undefined + const streamAbortController = streamTimeoutDisabled ? undefined : new AbortController() const headerTimeoutMs = headerTimeout === false ? undefined : headerTimeout const headerTimeoutCtl = typeof headerTimeoutMs === "number" ? timeoutController(headerTimeoutMs) : undefined const signals: AbortSignal[] = [] if (opts.signal) signals.push(opts.signal) - if (chunkAbortCtl) signals.push(chunkAbortCtl.signal) + if (streamAbortController) signals.push(streamAbortController.signal) if (headerTimeoutCtl) signals.push(headerTimeoutCtl.signal) if (options["timeout"] !== undefined && options["timeout"] !== null && options["timeout"] !== false) signals.push(AbortSignal.timeout(options["timeout"])) @@ -1763,8 +1725,13 @@ const layer = Layer.effect( timeout: false, }).finally(() => headerTimeoutCtl?.clear()) - if (!chunkAbortCtl) return res - return wrapSSE(res, chunkTimeout, chunkAbortCtl) + if (!streamAbortController) return res + return s.streamLiveness.wrap({ + response: res, + bucket: `${model.providerID}:${model.api.npm}`, + controller: streamAbortController, + timeout: fixedChunkTimeout, + }) } const bundledLoader = BUNDLED_PROVIDERS[model.api.npm] diff --git a/packages/opencode/src/provider/stream-liveness.ts b/packages/opencode/src/provider/stream-liveness.ts new file mode 100644 index 000000000000..bae6b0948439 --- /dev/null +++ b/packages/opencode/src/provider/stream-liveness.ts @@ -0,0 +1,110 @@ +import { ProviderError } from "./error" + +export type Policy = { + initial: number + minimum: number + maximum: number + multiplier: number + historySize: number +} + +export const defaultPolicy = { + initial: 900_000, + minimum: 900_000, + maximum: 1_800_000, + multiplier: 2, + historySize: 32, +} satisfies Policy + +export type Detector = ReturnType + +export function create(policy: Policy = defaultPolicy, now = () => performance.now()) { + const histories = new Map() + + function deadline(bucket: string) { + const values = histories.get(bucket) + if (!values?.length) return policy.initial + return Math.min(policy.maximum, Math.max(policy.minimum, Math.max(...values) * policy.multiplier)) + } + + function observe(bucket: string, elapsed: number) { + if (!Number.isFinite(elapsed) || elapsed < 0) return + const values = histories.get(bucket) ?? [] + values.push(elapsed) + if (values.length > policy.historySize) values.shift() + histories.set(bucket, values) + } + + function wrap(input: { + response: Response + bucket: string + controller: AbortController + timeout?: number | false + }) { + const fixed = + typeof input.timeout === "number" && Number.isFinite(input.timeout) ? input.timeout : undefined + if (input.timeout === false || (fixed !== undefined && fixed <= 0)) return input.response + if (!input.response.body) return input.response + if (!input.response.headers.get("content-type")?.includes("text/event-stream")) return input.response + + const ms = fixed ?? deadline(input.bucket) + const reader = input.response.body.getReader() + let maximum = 0 + const body = new ReadableStream({ + async pull(controller) { + const started = now() + const part = await new Promise>>((resolve, reject) => { + const id = setTimeout(() => { + const error = new ProviderError.ResponseStreamTimeoutError(ms) + input.controller.abort(error) + void reader.cancel(error).catch(() => {}) + reject(error) + }, ms) + + const read = () => + reader.read().then( + (part) => { + if (!part.done && part.value.byteLength === 0) { + void read() + return + } + clearTimeout(id) + resolve(part) + }, + (error) => { + clearTimeout(id) + reject(error) + }, + ) + void read() + }) + + maximum = Math.max(maximum, now() - started) + if (part.done) { + observe(input.bucket, maximum) + controller.close() + return + } + controller.enqueue(part.value) + }, + async cancel(reason) { + input.controller.abort(reason) + await reader.cancel(reason) + }, + }) + + return new Response(body, { + headers: new Headers(input.response.headers), + status: input.response.status, + statusText: input.response.statusText, + }) + } + + return { + deadline, + observe, + wrap, + } +} + +export * as StreamLiveness from "./stream-liveness" diff --git a/packages/opencode/src/session/message-v2.ts b/packages/opencode/src/session/message-v2.ts index 9b3f2c46f405..79faa358e630 100644 --- a/packages/opencode/src/session/message-v2.ts +++ b/packages/opencode/src/session/message-v2.ts @@ -665,6 +665,18 @@ export function fromError( }, { cause: e }, ).toObject() + case e instanceof ProviderError.ResponseStreamTimeoutError: + return new APIError( + { + message: e.message, + isRetryable: true, + metadata: { + code: e.name, + timeoutMs: String(e.ms), + }, + }, + { cause: e }, + ).toObject() case e instanceof ProviderError.ResponseStreamError: return new APIError( { diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index fc5ab04e108b..859d219bd42d 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -1,5 +1,5 @@ import { afterEach, expect } from "bun:test" -import { createServer, type Server } from "node:http" +import { createServer, type Server, type ServerResponse } from "node:http" import { streamText } from "ai" import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" @@ -46,11 +46,11 @@ it.live("headerTimeout does not abort delayed SSE body after headers arrive", () }), ) -it.live("chunkTimeout raises a response stream error when SSE body stalls", () => +it.live("chunkTimeout raises a typed timeout when SSE body never yields or closes", () => Effect.gen(function* () { const server = yield* Effect.acquireRelease( - Effect.promise(() => delayedBodyServer(250)), - (server) => Effect.sync(() => server.server.close()), + Effect.promise(() => hangingBodyServer()), + (server) => Effect.sync(() => server.close()), ) yield* provideTmpdirInstance( @@ -72,14 +72,43 @@ it.live("chunkTimeout raises a response stream error when SSE body stalls", () = } catch (error) { return error } + return undefined }) - expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) + expect(error).toBeInstanceOf(ProviderError.ResponseStreamTimeoutError) + if (!(error instanceof ProviderError.ResponseStreamTimeoutError)) throw new Error("expected stream timeout") + expect(error.ms).toBe(50) }), { config: providerConfig(server.url, { chunkTimeout: 50 }) }, ) }), ) +for (const chunkTimeout of [false, 0] as const) { + it.live(`chunkTimeout ${chunkTimeout} disables the SSE body watchdog`, () => + Effect.gen(function* () { + const server = yield* Effect.acquireRelease( + Effect.promise(() => delayedBodyServer(100)), + (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), + messages: [{ role: "user", content: "hello" }], + }) + + expect(yield* Effect.promise(() => result.text)).toBe("late") + }), + { config: providerConfig(server.url, { chunkTimeout }) }, + ) + }), + ) +} + it.live("headerTimeout aborts when response headers do not arrive", () => Effect.gen(function* () { const server = yield* Effect.acquireRelease( @@ -154,7 +183,7 @@ it.live("OpenAI Codex headerTimeout default can be disabled by config", () => }), ) -it.live("OpenAI API auth gets default headerTimeout", () => +it.live("OpenAI keeps its header default without a provider-only stream default", () => Effect.gen(function* () { yield* withAuthContent( Effect.gen(function* () { @@ -163,6 +192,7 @@ it.live("OpenAI API auth gets default headerTimeout", () => const provider = yield* Provider.Service const openai = yield* provider.getProvider(ProviderV2.ID.openai) expect(openai.options.headerTimeout).toBe(300_000) + expect(openai.options.chunkTimeout).toBeUndefined() }), ) }), @@ -211,6 +241,26 @@ async function delayedBodyServer(delay: number): Promise<{ server: Server; url: return { server, url: `http://127.0.0.1:${address.port}` } } +async function hangingBodyServer(): Promise<{ close(): void; url: string }> { + const responses = new Set() + const server = createServer((_, res) => { + responses.add(res) + res.on("close", () => responses.delete(res)) + res.writeHead(200, { "content-type": "text/event-stream" }) + res.flushHeaders() + }) + 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 { + close() { + responses.forEach((response) => response.destroy()) + server.close() + }, + url: `http://127.0.0.1:${address.port}`, + } +} + function withAuthContent(self: Effect.Effect, value: Record = defaultAuthContent()) { return Effect.acquireUseRelease( Effect.sync(() => { diff --git a/packages/opencode/test/provider/stream-liveness.test.ts b/packages/opencode/test/provider/stream-liveness.test.ts new file mode 100644 index 000000000000..c8b50cb10665 --- /dev/null +++ b/packages/opencode/test/provider/stream-liveness.test.ts @@ -0,0 +1,178 @@ +import { describe, expect, test } from "bun:test" +import { StreamLiveness } from "@/provider/stream-liveness" + +const policy = { + initial: 30, + minimum: 10, + maximum: 100, + multiplier: 2, + historySize: 2, +} satisfies StreamLiveness.Policy + +describe("stream liveness deadline", () => { + test("uses the cold-start deadline without observations", () => { + expect(StreamLiveness.create(policy).deadline("anthropic:@ai-sdk/anthropic")).toBe(30) + }) + + test("doubles the rolling maximum and applies both bounds", () => { + const detector = StreamLiveness.create(policy) + detector.observe("provider:transport", 2) + expect(detector.deadline("provider:transport")).toBe(10) + detector.observe("provider:transport", 60) + expect(detector.deadline("provider:transport")).toBe(100) + }) + + test("evicts observations outside the bounded history", () => { + const detector = StreamLiveness.create(policy) + detector.observe("provider:transport", 40) + detector.observe("provider:transport", 20) + detector.observe("provider:transport", 5) + expect(detector.deadline("provider:transport")).toBe(40) + }) + + test("isolates provider and transport buckets", () => { + const detector = StreamLiveness.create(policy) + detector.observe("openai:@ai-sdk/openai", 40) + expect(detector.deadline("anthropic:@ai-sdk/anthropic")).toBe(30) + }) +}) + +describe("stream liveness response body", () => { + test("fails a post-header SSE body that never yields or closes", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 20 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response(new ReadableStream({ pull: () => new Promise(() => {}) }), { + headers: { "content-type": "text/event-stream" }, + }), + }) + + await expect(response.text()).rejects.toHaveProperty("name", "ProviderResponseStreamTimeoutError") + }) + + test("raw heartbeat bytes reset the pending-read deadline", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 30 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response( + new ReadableStream({ + start(controller) { + setTimeout(() => controller.enqueue(new TextEncoder().encode(": heartbeat\n\n")), 10) + setTimeout(() => controller.close(), 20) + }, + }), + { headers: { "content-type": "text/event-stream" } }, + ), + }) + + expect(await response.text()).toBe(": heartbeat\n\n") + }) + + test("empty chunks do not count as body progress", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 20 }) + let id: ReturnType | undefined + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response( + new ReadableStream({ + start(controller) { + id = setInterval(() => controller.enqueue(new Uint8Array()), 5) + }, + cancel() { + clearInterval(id) + }, + }), + { headers: { "content-type": "text/event-stream" } }, + ), + }) + + await expect(response.text()).rejects.toHaveProperty("name", "ProviderResponseStreamTimeoutError") + }) + + test("commits one maximum-gap observation only on normal EOF", async () => { + const readings = [0, 7, 7, 20] + const detector = StreamLiveness.create(policy, () => readings.shift() ?? 20) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new Uint8Array([1])) + controller.close() + }, + }), + { headers: { "content-type": "text/event-stream" } }, + ), + }) + + await response.arrayBuffer() + expect(detector.deadline("test:transport")).toBe(26) + }) + + test("does not train history after timeout", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 20 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response(new ReadableStream({ pull: () => new Promise(() => {}) }), { + headers: { "content-type": "text/event-stream" }, + }), + }) + + await expect(response.text()).rejects.toBeDefined() + expect(detector.deadline("test:transport")).toBe(20) + }) + + test("does not train history after consumer cancellation", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 20 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response(new ReadableStream(), { + headers: { "content-type": "text/event-stream" }, + }), + }) + + await response.body?.cancel() + expect(detector.deadline("test:transport")).toBe(20) + }) + + test("leaves non-SSE and explicitly disabled responses unchanged", () => { + const detector = StreamLiveness.create(policy) + const json = new Response("{}") + expect( + detector.wrap({ response: json, bucket: "test:transport", controller: new AbortController() }), + ).toBe(json) + + const sse = new Response("", { headers: { "content-type": "text/event-stream" } }) + expect( + detector.wrap({ + response: sse, + bucket: "test:transport", + controller: new AbortController(), + timeout: false, + }), + ).toBe(sse) + }) + + test("uses a positive fixed override instead of the adaptive deadline", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 100 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response(new ReadableStream({ pull: () => new Promise(() => {}) }), { + headers: { "content-type": "text/event-stream" }, + }), + timeout: 8, + }) + + await expect(response.text()).rejects.toMatchObject({ + name: "ProviderResponseStreamTimeoutError", + ms: 8, + }) + }) +}) diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index 052477d0a2e7..18921a697f1f 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -70,7 +70,7 @@ const cfg = { }, } -function providerCfg(url: string) { +function providerCfg(url: string, options: Record = {}) { return { ...cfg, provider: { @@ -80,6 +80,7 @@ function providerCfg(url: string) { options: { ...cfg.provider.test.options, baseURL: url, + ...options, }, }, }, @@ -651,6 +652,64 @@ it.live("session.processor effect tests retry OpenAI-compatible midstream server ), ) +it.live("session.processor effect tests retry response stream inactivity timeouts", () => + provideTmpdirServer( + ({ dir, llm }) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + const events = yield* EventV2Bridge.Service + + yield* llm.hang + yield* llm.text("after timeout") + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "retry stalled stream") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const states: number[] = [] + const off = yield* events.listen((evt) => { + if (evt.type !== SessionStatus.Event.Status.type) return Effect.void + const data = evt.data as typeof SessionStatus.Event.Status.data.Type + if (data.sessionID === chat.id && data.status.type === "retry") states.push(data.status.attempt) + return Effect.void + }) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies SessionV1.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "retry stalled stream" }], + tools: {}, + }) + + yield* off + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("continue") + expect(yield* llm.calls).toBe(2) + expect(states).toStrictEqual([1]) + expect(parts.some((part) => part.type === "text" && part.text === "after timeout")).toBe(true) + expect(handle.message.error).toBeUndefined() + }), + { config: (url) => providerCfg(url, { chunkTimeout: 50 }) }, + ), +) + it.live("session.processor effect tests publish retry status updates", () => provideTmpdirServer( ({ dir, llm }) => diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index e21b12c0895e..4a91c48e26a8 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -246,6 +246,20 @@ describe("session.retry.retryable", () => { }) }) + test("retries response stream inactivity timeouts", () => { + const request = MessageV2.fromError(new ProviderError.ResponseStreamTimeoutError(900_000), { providerID }) + expect(SessionV1.APIError.isInstance(request)).toBe(true) + if (!SessionV1.APIError.isInstance(request)) throw new Error("expected APIError") + expect(request.data.metadata).toEqual({ + code: "ProviderResponseStreamTimeoutError", + timeoutMs: "900000", + }) + expect(SessionRetry.retryable(request, retryProvider)).toEqual({ + message: + "Provider response stream timed out after 900000ms of inactivity. Check provider connectivity or increase provider.options.chunkTimeout before retrying.", + }) + }) + test("retries websocket stream transport errors", () => { const request = MessageV2.fromError( new ProviderError.ResponseStreamError("WebSocket closed before response.completed (code 1006: Connection ended)"), From 9bc30bbdb9bf208fc793f6141d0b5f1444403fac Mon Sep 17 00:00:00 2001 From: Daniel Cadenas Date: Fri, 31 Jul 2026 12:05:57 -0300 Subject: [PATCH 2/2] fix(provider): detect unlabeled model streams --- packages/opencode/src/provider/provider.ts | 11 ++++- .../opencode/src/provider/stream-liveness.ts | 3 +- .../test/provider/header-timeout.test.ts | 40 ++++++++++++++++++- .../test/provider/stream-liveness.test.ts | 20 ++++++++++ 4 files changed, 70 insertions(+), 4 deletions(-) diff --git a/packages/opencode/src/provider/provider.ts b/packages/opencode/src/provider/provider.ts index 43538da2fd3a..af946bf82838 100644 --- a/packages/opencode/src/provider/provider.ts +++ b/packages/opencode/src/provider/provider.ts @@ -18,7 +18,7 @@ import { iife } from "@/util/iife" import { Global } from "@opencode-ai/core/global" import path from "path" import { pathToFileURL } from "url" -import { Effect, Layer, Context, Schema, Types } from "effect" +import { Effect, Layer, Context, Option, Schema, Types } from "effect" import { EffectBridge } from "@/effect/bridge" import { InstanceState } from "@/effect/instance-state" import { EffectPromise } from "@/effect/promise" @@ -44,6 +44,14 @@ function timeoutController(ms: number) { } } +const decodeJson = Schema.decodeUnknownOption(Schema.UnknownFromJsonString) + +function streamsResponse(body: BodyInit | null | undefined) { + if (typeof body !== "string") return false + const value = Option.getOrUndefined(decodeJson(body)) + return isRecord(value) && value.stream === true +} + function googleVertexAnthropicBaseURL(project: string | undefined, location: string | undefined) { if (!project) return if (location !== "eu" && location !== "us") return @@ -1731,6 +1739,7 @@ const layer = Layer.effect( bucket: `${model.providerID}:${model.api.npm}`, controller: streamAbortController, timeout: fixedChunkTimeout, + stream: streamsResponse(opts.body), }) } diff --git a/packages/opencode/src/provider/stream-liveness.ts b/packages/opencode/src/provider/stream-liveness.ts index bae6b0948439..0cc67a184869 100644 --- a/packages/opencode/src/provider/stream-liveness.ts +++ b/packages/opencode/src/provider/stream-liveness.ts @@ -40,12 +40,13 @@ export function create(policy: Policy = defaultPolicy, now = () => performance.n bucket: string controller: AbortController timeout?: number | false + stream?: boolean }) { const fixed = typeof input.timeout === "number" && Number.isFinite(input.timeout) ? input.timeout : undefined if (input.timeout === false || (fixed !== undefined && fixed <= 0)) return input.response if (!input.response.body) return input.response - if (!input.response.headers.get("content-type")?.includes("text/event-stream")) return input.response + if (!input.stream && !input.response.headers.get("content-type")?.includes("text/event-stream")) return input.response const ms = fixed ?? deadline(input.bucket) const reader = input.response.body.getReader() diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index 859d219bd42d..cb9c73d52187 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -83,6 +83,42 @@ it.live("chunkTimeout raises a typed timeout when SSE body never yields or close }), ) +it.live("chunkTimeout uses request stream intent when the response omits its content type", () => + Effect.gen(function* () { + const server = yield* Effect.acquireRelease( + Effect.promise(() => hangingBodyServer(false)), + (server) => Effect.sync(() => 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), + abortSignal: AbortSignal.timeout(200), + onError() {}, + messages: [{ role: "user", content: "hello" }], + }) + + const error = yield* Effect.promise(async () => { + try { + for await (const part of result.fullStream) { + if (part.type === "error") return part.error + } + } catch (error) { + return error + } + return undefined + }) + expect(error).toBeInstanceOf(ProviderError.ResponseStreamTimeoutError) + }), + { config: providerConfig(server.url, { chunkTimeout: 50 }) }, + ) + }), +) + for (const chunkTimeout of [false, 0] as const) { it.live(`chunkTimeout ${chunkTimeout} disables the SSE body watchdog`, () => Effect.gen(function* () { @@ -241,12 +277,12 @@ async function delayedBodyServer(delay: number): Promise<{ server: Server; url: return { server, url: `http://127.0.0.1:${address.port}` } } -async function hangingBodyServer(): Promise<{ close(): void; url: string }> { +async function hangingBodyServer(contentType = true): Promise<{ close(): void; url: string }> { const responses = new Set() const server = createServer((_, res) => { responses.add(res) res.on("close", () => responses.delete(res)) - res.writeHead(200, { "content-type": "text/event-stream" }) + res.writeHead(200, contentType ? { "content-type": "text/event-stream" } : {}) res.flushHeaders() }) await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)) diff --git a/packages/opencode/test/provider/stream-liveness.test.ts b/packages/opencode/test/provider/stream-liveness.test.ts index c8b50cb10665..f19e4f195f6b 100644 --- a/packages/opencode/test/provider/stream-liveness.test.ts +++ b/packages/opencode/test/provider/stream-liveness.test.ts @@ -51,6 +51,26 @@ describe("stream liveness response body", () => { await expect(response.text()).rejects.toHaveProperty("name", "ProviderResponseStreamTimeoutError") }) + test("uses request stream intent when the response omits its content type", async () => { + const detector = StreamLiveness.create({ ...policy, initial: 10 }) + const response = detector.wrap({ + bucket: "test:transport", + controller: new AbortController(), + response: new Response( + new ReadableStream({ + async pull(controller) { + await Bun.sleep(30) + controller.enqueue(new TextEncoder().encode("late")) + controller.close() + }, + }), + ), + stream: true, + }) + + await expect(response.text()).rejects.toHaveProperty("name", "ProviderResponseStreamTimeoutError") + }) + test("raw heartbeat bytes reset the pending-read deadline", async () => { const detector = StreamLiveness.create({ ...policy, initial: 30 }) const response = detector.wrap({