fix(core): fall back on oversized websocket requests (#43099)
This commit is contained in:
@@ -110,10 +110,11 @@ export function createWebSocketFetch(options?: CreateWebSocketFetchOptions) {
|
|||||||
invalidate(entry)
|
invalidate(entry)
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
onConnectionInvalid: (error) => {
|
onConnectionInvalid: (_error, closeCode) => {
|
||||||
entry.busy = false
|
entry.busy = false
|
||||||
entry.lastUsedAt = Date.now()
|
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)
|
invalidate(entry)
|
||||||
resolveFirstEvent(false)
|
resolveFirstEvent(false)
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ import { ProxyEnv } from "@/util/proxy-env"
|
|||||||
import { isRecord } from "@/util/record"
|
import { isRecord } from "@/util/record"
|
||||||
|
|
||||||
export const PROTOCOL_HEADER = "responses_websockets=2026-02-06"
|
export const PROTOCOL_HEADER = "responses_websockets=2026-02-06"
|
||||||
|
export const MESSAGE_TOO_BIG_CLOSE_CODE = 1009
|
||||||
|
|
||||||
export interface ConnectResponsesWebSocketOptions {
|
export interface ConnectResponsesWebSocketOptions {
|
||||||
url: string
|
url: string
|
||||||
@@ -26,7 +27,7 @@ export interface StreamResponsesWebSocketOptions {
|
|||||||
onComplete?: (event: Record<string, unknown>) => void
|
onComplete?: (event: Record<string, unknown>) => void
|
||||||
onTerminal?: (event: Record<string, unknown>) => void
|
onTerminal?: (event: Record<string, unknown>) => void
|
||||||
onRetryableTerminal?: (event: Record<string, unknown>) => Promise<WebSocket | undefined>
|
onRetryableTerminal?: (event: Record<string, unknown>) => Promise<WebSocket | undefined>
|
||||||
onConnectionInvalid?: (error: ProviderError.ResponseStreamError) => void
|
onConnectionInvalid?: (error: ProviderError.ResponseStreamError, closeCode?: number) => void
|
||||||
onAbort?: (error: Error) => void
|
onAbort?: (error: Error) => void
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -162,11 +163,11 @@ export function streamResponsesWebSocket(options: StreamResponsesWebSocketOption
|
|||||||
controller?.close()
|
controller?.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
function invalidate(error: ProviderError.ResponseStreamError) {
|
function invalidate(error: ProviderError.ResponseStreamError, closeCode?: number) {
|
||||||
if (completed) return
|
if (completed) return
|
||||||
completed = true
|
completed = true
|
||||||
cleanup()
|
cleanup()
|
||||||
options.onConnectionInvalid?.(error)
|
options.onConnectionInvalid?.(error, closeCode)
|
||||||
controller?.error(error)
|
controller?.error(error)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -274,6 +275,7 @@ export function streamResponsesWebSocket(options: StreamResponsesWebSocketOption
|
|||||||
if (completed) return
|
if (completed) return
|
||||||
invalidate(
|
invalidate(
|
||||||
new ProviderError.ResponseStreamError(closeMessage("WebSocket closed before response.completed", code, reason)),
|
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) {
|
function closeMessage(message: string, code: number, reason: Buffer) {
|
||||||
const details = [`code ${code}`]
|
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())
|
if (reason.length > 0) details.push(reason.toString())
|
||||||
return `${message} (${details.join(": ")})`
|
return `${message} (${details.join(": ")})`
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -237,6 +237,26 @@ describe("plugin.openai.ws-pool", () => {
|
|||||||
fetch.close()
|
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 () => {
|
test("removes HTTP fallback when its session is deleted", async () => {
|
||||||
let websocketAttempts = 0
|
let websocketAttempts = 0
|
||||||
await using server = await createRejectingWebSocketServer(() => websocketAttempts++)
|
await using server = await createRejectingWebSocketServer(() => websocketAttempts++)
|
||||||
|
|||||||
Reference in New Issue
Block a user