From 7af274a92102196e73193c5b6de6f0d2aedfff46 Mon Sep 17 00:00:00 2001 From: Filip <34747899+neriousy@users.noreply.github.com> Date: Mon, 17 Aug 2026 20:36:00 +0200 Subject: [PATCH] fix(core): fall back on oversized websocket requests (#43099) --- .../opencode/src/plugin/openai/ws-pool.ts | 5 +++-- packages/opencode/src/plugin/openai/ws.ts | 10 ++++++---- .../opencode/test/plugin/openai-ws.test.ts | 20 +++++++++++++++++++ 3 files changed, 29 insertions(+), 6 deletions(-) diff --git a/packages/opencode/src/plugin/openai/ws-pool.ts b/packages/opencode/src/plugin/openai/ws-pool.ts index 3cbb29a301..939c2dc232 100644 --- a/packages/opencode/src/plugin/openai/ws-pool.ts +++ b/packages/opencode/src/plugin/openai/ws-pool.ts @@ -110,10 +110,11 @@ export function createWebSocketFetch(options?: CreateWebSocketFetchOptions) { invalidate(entry) } }, - onConnectionInvalid: (error) => { + onConnectionInvalid: (_error, closeCode) => { entry.busy = false entry.lastUsedAt = Date.now() - if (!entry.fallback) recordStreamFailure(entry) + if (closeCode === OpenAIWebSocket.MESSAGE_TOO_BIG_CLOSE_CODE) entry.fallback = true + else if (!entry.fallback) recordStreamFailure(entry) invalidate(entry) resolveFirstEvent(false) }, diff --git a/packages/opencode/src/plugin/openai/ws.ts b/packages/opencode/src/plugin/openai/ws.ts index 578d00b8ce..4335d9215a 100644 --- a/packages/opencode/src/plugin/openai/ws.ts +++ b/packages/opencode/src/plugin/openai/ws.ts @@ -9,6 +9,7 @@ import { ProxyEnv } from "@/util/proxy-env" import { isRecord } from "@/util/record" export const PROTOCOL_HEADER = "responses_websockets=2026-02-06" +export const MESSAGE_TOO_BIG_CLOSE_CODE = 1009 export interface ConnectResponsesWebSocketOptions { url: string @@ -26,7 +27,7 @@ export interface StreamResponsesWebSocketOptions { onComplete?: (event: Record) => void onTerminal?: (event: Record) => void onRetryableTerminal?: (event: Record) => Promise - onConnectionInvalid?: (error: ProviderError.ResponseStreamError) => void + onConnectionInvalid?: (error: ProviderError.ResponseStreamError, closeCode?: number) => void onAbort?: (error: Error) => void } @@ -162,11 +163,11 @@ export function streamResponsesWebSocket(options: StreamResponsesWebSocketOption controller?.close() } - function invalidate(error: ProviderError.ResponseStreamError) { + function invalidate(error: ProviderError.ResponseStreamError, closeCode?: number) { if (completed) return completed = true cleanup() - options.onConnectionInvalid?.(error) + options.onConnectionInvalid?.(error, closeCode) controller?.error(error) } @@ -274,6 +275,7 @@ export function streamResponsesWebSocket(options: StreamResponsesWebSocketOption if (completed) return invalidate( new ProviderError.ResponseStreamError(closeMessage("WebSocket closed before response.completed", code, reason)), + code, ) } @@ -373,7 +375,7 @@ function abortError(signal: AbortSignal | undefined) { function closeMessage(message: string, code: number, reason: Buffer) { const details = [`code ${code}`] - if (code === 1009) details.push("message too big") + if (code === MESSAGE_TOO_BIG_CLOSE_CODE) details.push("message too big") if (reason.length > 0) details.push(reason.toString()) return `${message} (${details.join(": ")})` } diff --git a/packages/opencode/test/plugin/openai-ws.test.ts b/packages/opencode/test/plugin/openai-ws.test.ts index e8025d0a92..e88d0620f5 100644 --- a/packages/opencode/test/plugin/openai-ws.test.ts +++ b/packages/opencode/test/plugin/openai-ws.test.ts @@ -237,6 +237,26 @@ describe("plugin.openai.ws-pool", () => { fetch.close() }) + test("falls back immediately to HTTP when a websocket request is too large", async () => { + let connections = 0 + await using server = await createWebSocketServer((socket) => { + connections += 1 + socket.once("message", () => socket.close(1009, "payload too large")) + }) + const fetch = OpenAIWebSocketPool.createWebSocketFetch({ + url: server.url, + }) + + const first = await fetch(server.url, streamRequest()) + const second = await fetch(server.url, streamRequest()) + + expect(await first.text()).toBe("http") + expect(await second.text()).toBe("http") + expect(connections).toBe(1) + expect(server.httpRequests).toHaveLength(2) + fetch.close() + }) + test("removes HTTP fallback when its session is deleted", async () => { let websocketAttempts = 0 await using server = await createRejectingWebSocketServer(() => websocketAttempts++)