Files
kilocode/packages/opencode/test/plugin/openai-ws.test.ts
T
+25 ef6b152ff8 OpenCode v1.16.2 (#12088)
* feat(worktree): add managed workspace cloning (#30117)

* test(tui): skip crashing keymap textarea renderer

* fix(core): allow skipping migration execution

* fix(opencode): remove automatic full session diffs (#30127)

* chore: generate

* refactor(worktree): move project out of repository

* zen: deepseek flash

* fix(tui): remount session view on session switch (#30129)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* go: minimax m3

* refactor(opencode): inline local provider helpers (#30169)

* refactor(opencode): simplify provider setup flow (#30173)

* fix(app): show project sessions before path sync resolves (#30167)

Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>

* fix(core): preserve session metadata migration identity (#30176)

* refactor(session): align namespace imports and inline trivial helpers (#30180)

* opencode(run): add queued prompt management (#30103)

Direct run mode previously made submitted follow-up prompts irrevocable while a response was still running. Let users edit or remove queued prompts before dispatch without interrupting the active turn.

* chore: generate

* fix(acp): honor session/cancel by aborting the running turn (#30145)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(tui): prevent prompt corruption when pasting near wide characters (#29710)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Simon Klee <hello@simonklee.dk>

* fix(opencode): avoid nullable webfetch format schema (#30215)

* chore: generate

* fix(core): contain lsp warmup defects (#30226)

* add run --replay mode (#30239)

* chore: generate

* chore: update nix node_modules hashes

* fix(stats): restore leaderboard spacing

* fix(stats): center top models dot grid

* fix(stats): stabilize top models hover

* fix(stats): align big-pickle provider resolution (#30274)

* feat(app): v2 desktop UI improvements (#29689)

Co-authored-by: Brendan Allan <git@brendonovich.dev>
Co-authored-by: Brendan Allan <14191578+Brendonovich@users.noreply.github.com>

* chore: generate

* fix(tui): clarify inline subagent rows (#30051)

* fix(tui): handle events across workspaces (#30281)

* feat(core): update Copilot for token-based billing (#30181)

* fix(tui): keep background marker with subagent label (#30271)

* fix(tui): keep retry attempt before message (#30275)

* chore: generate

* fix(opencode): enforce storage path invariants (#29666)

* chore: generate

* feat(core): add location-based permission service (#30287)

* chore: generate

* fix(tui): preserve live parts during session hydration (#30300)

* fix(app): restore deferred MCP status updates (#30220)

* fix: export v2 stylesheets and declare core node types (#30312)

* chore: update nix node_modules hashes

* fix(app): avoid suspending on pending child path (#30314)

* fix(opencode): remove sunsetted gpt-5.2 and gpt-5.3-codex from allowed models for codex subscriptions (#30316)

* chore: generate

* feat(core): expose session location

* chore: generate

* fix(opencode): preserve websocket api errors (#30321)

* refactor(core): simplify session pagination

* feat(core): add location filesystem contract

* feat(core): add dummy location filesystem layer

* chore: generate

* feat(opencode): add filesystem read and list routes

* chore: generate

* infra: stats

* sync

* feat(app): inset new layout session panels (#30342)

* fix(app): tab title truncation and close button positioning (#30349)

* tui: show model context in run footer (#30380)

* tui: revert OpenTUI upgrade to 0.2.16 (#30383)

* chore: update nix node_modules hashes

* feat(core): add managed repository cache (#30408)

* chore: generate

* chore: generate

* sync

* feat(stats): add cache ratio section

* feat(core): add flagged project references (#30414)

* chore: generate

* feat(core): support named migrations (#30418)

* fix(stats): clean retired provider rows during sync (#30420)

* fix(stats): mention opencode go in top models copy

* feat(core): expose project reference filesystem access (#30423)

* chore: generate

* sync

* fix(tui): scope diff viewer to session directory (#30426)

* test: widen provider header timeout margin (#30427)

* fix(plugin): restore private git install fallback (#30430)

* fix(stats): remove leaderboard nav link

* chore(opencode): remove scout agent (#30435)

* chore: generate

* feat(stats): improve cache ratio chart

* chore: generate

* fix(effect-drizzle-sqlite): preserve transaction begin errors (#30448)

* chore: bump effect beta to 74 (#30449)

* Revert "tui: revert OpenTUI upgrade to 0.2.16 (#30383)" (#30452)

* chore: update nix node_modules hashes

* refactor(opencode): improve startup time by 38% (#30453)

Co-authored-by: starptech <starptech@starptechs-MBP.fritz.box>

* chore: generate

* fix(opencode): patch empty Gemini replay messages (#30463)

* chore: generate

* refactor(core): consolidate filesystem services (#30447)

* chore: generate

* run: enable interactive replay by default (#30465)

* chore: update nix node_modules hashes

* refactor(opencode): remove JSON storage migration (#30461)

* chore: generate

* chore: update nix node_modules hashes

* fix(tui): stop idle background task spinner (#30484)

* refactor(core): move v1 schemas into core (#30473)

* chore: generate

* fix: task id passed to background job for continuation (#30485)

* chore: generate

* feat(core): project copying and tracking directories (#30139)

* chore: generate

* fix(opencode): preserve signed thinking during anthropic reorder (#30182)

* Revert "fix(opencode): preserve signed thinking during anthropic reorder" (#30502)

* fix: rm tool reorder logic from old bug (#30483)

* chore: generate

* feat(app): polish home projects list UI (#30436)

* feat(app): polish select-v2 component (#30446)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* fix(github): enforce existing git author identity (#30507)

* feat(app): new update button  (#30460)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* fix(opencode): fallback to sh for curl upgrade (#30499)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(ui): render whole-file patches as complete diffs (#30516)

* chore: generate

* feat(app): add servers tab to settings dialog (#29675)

* refactor(core): consolidate pty service (#30537)

* chore: generate

* tui: truncate sidebar file paths (#30531)

* chore: update nix node_modules hashes

* feat(stats): add geo breakdown (#30456)

* chore: generate

* chore: update nix node_modules hashes

* fix(acp): classify apply_patch as edit (#30564)

* fix(acp): classify task as think (#30565)

* fix(acp): include external directory permission context (#30567)

* fix(acp): clean read tool display content (#30569)

* fix(tui): route question responses by session directory (#30578)

* fix(stats): serve stats og image from banner

* docs(go): add Qwen3.7 Plus model (#30594)

* fix(openai): preserve websocket idle state (#30586)

* refactor(core): remove ai sdk option fields (#30581)

* chore: generate

* test(core): cover v1 provider option lowering (#30599)

* chore: generate

* refactor(core): nest model api id (#30603)

* fix(core): expose azure openai xhigh efforts (#30620)

* feat(core): add skill registry and file agent loading (#30617)

* chore: generate

* chore: update nix node_modules hashes

* fix(stats): count all go usage

* chore: remove zed extension and automation (#30628)

* fix(opencode): preserve variant for delegated tasks (#30630)

* zen: update nvidia tos

* fix(opencode): route SAP AI Core reasoning variants through modelParams (#30482)

* chore: generate

* fix(app): hide unavailable titlebar update (#30642)

* feat(app): v2 thinking level selector (#30646)

* fix(app,ui): session review reactivity and VCS query cache (#30660)

* feat(core): add embedded v2 session runtime and tool foundation (#30632)

* chore: generate

* chore: update nix node_modules hashes

* docs: correct compaction prune default (#30670)

* fix(opencode): avoid shell cancel race (#30641)

* feat: bump bedrock and add proper mantle support for openai models through aws bedrock (#30464)

* test: wait for shell truncation readiness (#30679)

* chore: update nix node_modules hashes

* refactor(opencode): clean up task tool prompts (#30687)

* feat(core): add command registry (#30624)

* chore: generate

* fix(acp): replay loaded session transcript (#30645)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(core): reset pre-launch session projections (#30728)

* feat(tui): improve experimental session switcher (#30738)

* fix(opencode): respect disabled auto compaction on overflow (#30749)

* zen: nemotron 3 ultra

* fix(enterprise): install hono standard validator peer (#30740)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* fix build

* chore: update nix node_modules hashes

* make scripts executable

* fix(tui): show toast when variant_list keybind used with no variants (#30724)

* fix(opencode): `ACP.loadSession` should replay all messages (#30761)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(opencode): attribute task child agent on creation (#30786)

* fix(tui): add Vue syntax highlighting (#30802)

* fix: bump @openrouter/ai-sdk-provider to 2.9.0 (#30800)

* feat(core): moving sessions (#30640)

* chore: generate

* tweak: background agent prompting to avoid polling issues (#30790)

* upgrade opentui to 0.3.2 (#30748)

* chore: update nix node_modules hashes

* feat(desktop): surface local server startup failures (#30822)

* ci: publish

* refactor(core): make v2 session inputs event sourced (#30785)

* chore: generate

* fix(llm): normalize OpenAI function tool schemas

* chore: generate

* feat(stats): refresh stats routes and homepage (#30419)

* fix(stats): sort metric charts by top usage

* feat(core): add public native API (#30828)

* chore: generate

* feat(app): color themes (#30824)

Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>

* chore: generate

* sync release versions for v1.16.0

* feat(core): attach global native tools (#30832)

* chore: generate

* feat(core): add Snowflake Cortex provider (#29901)

Co-authored-by: Cortex Code <noreply@snowflake.com>

* chore: generate

* feat(core): persist v2 session context epochs (#30789)

* chore: generate

* feat(tui): allow backgrounding synchronous subagents (#30488)

* fix(app): improve tab handling (#30669)

* chore: generate

* fix(tui): prioritize models slash autocomplete (#30848)

* fix(tui): route permission replies to session directory (#30851)

Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>

* fix(cli): harden daemon lifecycle (#30844)

* chore: generate

* feat(app): improve desktop multi-server support (#30678)

Co-authored-by: Brendan Allan <git@brendonovich.dev>

* chore: generate

* fix(app): handle tab overflow and scrolling in titlebar (#30886)

* fix(app): tab overflow (#30894)

* tui: guard path formatting inputs (#30469)

Fixes #27726, #25216, #24856, #24294, #17071, #29164, #24837, #16865, #14279, #29895

* opencode/run: refresh themes after terminal reloads (#30917)

* chore: generate

* fix(tui): fall back to local cwd when editor spawns in attach mode (#30583)

* docs: update Go Qwen tiered pricing (#30936)

* chore: generate

* feat(tui): add diff hunk navigation (#30935)

* chore: rm fuzzy search on references (#30931)

* fix: use mapError instead of orDie for context snapshot decoding (#30905)

Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>

* fix(core): recover corrupted models cache (#30947)

* chore: bun install (#30968)

* fix(opencode): resolve Bedrock hang by using node build conditions (#30873)

* fix(workflows): retry nix-hashes compute-hash on transient failure (#30743)

* fix(stats): scroll model charts to latest on mobile

* fix(opencode): prevent destructive edit matches (#30932)

* chore: generate

* fix(core): respect v2 default agents (#30969)

* chore: generate

* test(opencode): remove disposal event wait race (#30971)

* test(opencode): remove shell timeout output race (#30974)

* fix(opencode): gate reasoning summaries by provider (#30973)

* feat(core): admit v2 skill guidance (#30843)

* fix(workflows): serialize desktop release uploads (#30978)

* fix(stats): add mobile chart end spacing

* release: v1.16.2

* refactor: kilo compat for v1.16.2

* fix(opencode): address v1.16.2 merge regressions

* chore: update kilo-vscode visual regression baselines

* fix(opencode): restore Kilo behavior after v1.16.2 merge

* fix(opencode): retry Windows migration cleanup

* test(opencode): restore clone and macOS watcher coverage

* fix(opencode): address second-pass review for #12099

Preserve imported usage and retry partial JSON migrations. Refresh active dependency patches, remove the obsolete GCP patch, and regenerate Kilo HttpApi branding.

---------

Co-authored-by: Dax <mail@thdxr.com>
Co-authored-by: Dax Raad <d@ironbay.co>
Co-authored-by: opencode-agent[bot] <opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Frank <frank@anoma.ly>
Co-authored-by: opencode-agent[bot] <219766164+opencode-agent[bot]@users.noreply.github.com>
Co-authored-by: Aiden Cline <63023139+rekram1-node@users.noreply.github.com>
Co-authored-by: Michael Hart <mhart@cloudflare.com>
Co-authored-by: LukeParkerDev <10430890+Hona@users.noreply.github.com>
Co-authored-by: Simon Klee <hello@simonklee.dk>
Co-authored-by: smagnuso <smagnuso@gmail.com>
Co-authored-by: Shoubhit Dash <shoubhit2005@gmail.com>
Co-authored-by: Orca丶 <93272799+dauphinYan@users.noreply.github.com>
Co-authored-by: Adam <2363879+adamdotdevin@users.noreply.github.com>
Co-authored-by: Aarav Sareen <96787824+arvsrn@users.noreply.github.com>
Co-authored-by: Brendan Allan <git@brendonovich.dev>
Co-authored-by: Brendan Allan <14191578+Brendonovich@users.noreply.github.com>
Co-authored-by: Kit Langton <kit.langton@gmail.com>
Co-authored-by: James Long <longster@gmail.com>
Co-authored-by: Dustin Deus <deusdustin@gmail.com>
Co-authored-by: starptech <starptech@starptechs-MBP.fritz.box>
Co-authored-by: Ulises Jeremias <ulisescf.24@gmail.com>
Co-authored-by: Jack <jack@anoma.ly>
Co-authored-by: Jérôme Benoit <jerome.benoit@sap.com>
Co-authored-by: Ariane Emory <97994360+ariane-emory@users.noreply.github.com>
Co-authored-by: LIU Xinyu <contact@lxy.cc>
Co-authored-by: Colin McDonnell <colinmcd94@gmail.com>
Co-authored-by: Sebastian <hasta84@gmail.com>
Co-authored-by: opencode <opencode@sst.dev>
Co-authored-by: Kamesh Sampath <kamesh.sampath@hotmail.com>
Co-authored-by: Cortex Code <noreply@snowflake.com>
Co-authored-by: pcadena-lila <pcadena@lila.ai>
Co-authored-by: weiconghe <46336277+weiconghe@users.noreply.github.com>
Co-authored-by: alberto <914199+alblez@users.noreply.github.com>
Co-authored-by: kilo-maintainer[bot] <kilo-maintainer[bot]@users.noreply.github.com>
2026-07-13 18:00:35 +02:00

878 lines
32 KiB
TypeScript

import { describe, expect, test } from "bun:test"
import { EventEmitter } from "node:events"
import { createServer, type IncomingMessage, type Server as HttpServer } from "node:http"
import net, { type AddressInfo, type Socket } from "node:net"
import WebSocket, { WebSocketServer } from "ws"
import { APICallError } from "ai"
import { ProviderError } from "../../src/provider/error"
import { OpenAIWebSocket } from "../../src/plugin/openai/ws"
import { OpenAIWebSocketPool, TITLE_HEADER } from "../../src/plugin/openai/ws-pool"
describe("plugin.openai.ws", () => {
test("derives websocket URLs and sends auth plus protocol headers", async () => {
let headers: IncomingMessage["headers"] | undefined
await using server = await createWebSocketServer((_socket, request) => {
headers = request.headers
})
const socket = await OpenAIWebSocket.connectResponsesWebSocket({
url: server.wsUrl,
headers: { authorization: "Bearer test", "content-length": "123" },
})
expect(OpenAIWebSocket.toWebSocketUrl("http://example.com/v1/responses")).toBe("ws://example.com/v1/responses")
expect(OpenAIWebSocket.toWebSocketUrl("https://example.com/v1/responses")).toBe("wss://example.com/v1/responses")
expect(headers?.authorization).toBe("Bearer test")
expect(headers?.["openai-beta"]).toBe(OpenAIWebSocket.PROTOCOL_HEADER)
expect(headers?.["content-length"]).toBeUndefined()
socket.terminate()
})
test("enforces websocket connect timeout", async () => {
await using server = await createHangingTcpServer()
await expect(
OpenAIWebSocket.connectResponsesWebSocket({
url: server.wsUrl,
headers: {},
timeout: 20,
}),
).rejects.toThrow("WebSocket connect timed out")
})
test("surfaces websocket upgrade rejection messages", async () => {
await using server = await createRejectingWebSocketServer(() => {})
await expect(
OpenAIWebSocket.connectResponsesWebSocket({
url: server.wsUrl,
headers: {},
}),
).rejects.toThrow("Expected 101 status code")
})
test("enforces websocket send idle timeout", async () => {
const socket = new (class extends EventEmitter {
send(_data: string, _callback: (error?: Error) => void) {}
})() as unknown as WebSocket
const invalid: string[] = []
const response = OpenAIWebSocket.streamResponsesWebSocket({
socket,
body: { stream: true, input: "hi" },
idleTimeout: 20,
onConnectionInvalid: (error) => invalid.push(error.message),
})
expect((await readTextError(response.text())).message).toContain("idle timeout sending websocket request")
expect(invalid).toEqual(["idle timeout sending websocket request"])
})
test("streams websocket events as SSE and handles response.done", async () => {
let requestBody: unknown
await using server = await createWebSocketServer((socket) => {
socket.once("message", (data) => {
requestBody = JSON.parse(data.toString())
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "hello" }))
socket.send(JSON.stringify({ type: "response.done", response: { id: "resp_123" } }))
socket.close(1000, "done")
})
})
const socket = await OpenAIWebSocket.connectResponsesWebSocket({
url: server.wsUrl,
headers: { authorization: "Bearer test", "content-length": "123" },
})
const completed: Record<string, unknown>[] = []
const response = OpenAIWebSocket.streamResponsesWebSocket({
socket,
body: { stream: true, background: true, input: "hi" },
onComplete: (event) => completed.push(event),
})
expect(await response.text()).toBe(
'data: {"type":"response.output_text.delta","delta":"hello"}\n\ndata: {"type":"response.done","response":{"id":"resp_123"}}\n\ndata: [DONE]\n\n',
)
expect(requestBody).toEqual({ type: "response.create", input: "hi" })
expect(completed).toHaveLength(1)
expect(completed[0]?.type).toBe("response.done")
})
test("errors the SSE stream when the server closes before a terminal event", async () => {
const invalid: Error[] = []
await using server = await createWebSocketServer((socket) => {
socket.once("message", () => {
socket.close(1009, "payload too large")
})
})
const socket = await OpenAIWebSocket.connectResponsesWebSocket({ url: server.wsUrl, headers: {} })
const response = OpenAIWebSocket.streamResponsesWebSocket({
socket,
body: { stream: true, input: "hi" },
onConnectionInvalid: (error) => invalid.push(error),
})
expect((await readTextError(response.text())).message).toContain(
"WebSocket closed before response.completed (code 1009: message too big: payload too large)",
)
expect(invalid[0]).toBeInstanceOf(ProviderError.ResponseStreamError)
expect(invalid.map((error) => error.message)).toEqual([
"WebSocket closed before response.completed (code 1009: message too big: payload too large)",
])
})
test("rejects unexpected binary websocket frames", async () => {
const invalid: string[] = []
await using server = await createWebSocketServer((socket) => {
socket.once("message", () => {
socket.send(Buffer.from("not json text"))
})
})
const socket = await OpenAIWebSocket.connectResponsesWebSocket({ url: server.wsUrl, headers: {} })
const response = OpenAIWebSocket.streamResponsesWebSocket({
socket,
body: { stream: true, input: "hi" },
onConnectionInvalid: (error) => invalid.push(error.message),
})
expect((await readTextError(response.text())).message).toContain("Unexpected binary WebSocket frame")
expect(invalid).toEqual(["Unexpected binary WebSocket frame"])
})
})
describe("plugin.openai.ws-pool", () => {
test("reuses one healthy websocket for sequential requests", async () => {
let connections = 0
let messages = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.on("message", () => {
messages += 1
socket.send(JSON.stringify({ type: "response.completed", response: { id: `resp_${messages}` } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toContain("data: [DONE]")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(1)
expect(messages).toBe(2)
fetch.close()
})
test("rotates a socket that exceeds max connection age", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.on("message", () => {
socket.send(JSON.stringify({ type: "response.completed", response: { id: `resp_${connections}` } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
maxConnectionAge: 0,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toContain("data: [DONE]")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(2)
fetch.close()
})
test("falls back to HTTP after websocket setup retries are exhausted", async () => {
const attempts: string[] = []
await using server = await createRejectingWebSocketServer(() => attempts.push("websocket"))
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
connectTimeout: 100,
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest({ [TITLE_HEADER]: "false" }))
expect(await readTextError(first.text())).toBeInstanceOf(ProviderError.ResponseStreamError)
const second = await fetch(server.url, streamRequest({ [TITLE_HEADER]: "false" }))
const third = await fetch(server.url, streamRequest({ [TITLE_HEADER]: "false" }))
expect(await second.text()).toBe("http")
expect(await third.text()).toBe("http")
expect(attempts).toEqual(["websocket", "websocket"])
expect(server.httpRequests).toHaveLength(2)
expect(server.httpRequests[0]?.headers[TITLE_HEADER]).toBeUndefined()
expect(server.httpRequests[1]?.headers[TITLE_HEADER]).toBeUndefined()
fetch.close()
})
test("keeps HTTP fallback active after its idle timeout", async () => {
let websocketAttempts = 0
await using server = await createRejectingWebSocketServer(() => websocketAttempts++)
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
connectTimeout: 100,
idleTimeout: 20,
streamRetries: 0,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toBe("http")
await new Promise((resolve) => setTimeout(resolve, 50))
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(websocketAttempts).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++)
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
connectTimeout: 100,
streamRetries: 0,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toBe("http")
fetch.remove("session-1")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(websocketAttempts).toBe(2)
expect(server.httpRequests).toHaveLength(2)
fetch.close()
})
test("terminates active websocket connections when their session is deleted", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
if (connections === 1) {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
return
}
socket.send(JSON.stringify({ type: "response.completed", response: { id: "resp_after_remove" } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest())
const firstText = first.text()
fetch.remove("session-1")
expect((await readTextError(firstText)).message).toContain("WebSocket closed before response.completed")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(2)
fetch.close()
})
test("prunes idle websocket connections after completed responses", async () => {
let connections = 0
let closed = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("close", () => closed++)
socket.once("message", () => {
socket.send(JSON.stringify({ type: "response.completed", response: { id: `resp_${connections}` } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
idleTimeout: 20,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toContain("data: [DONE]")
await waitFor(() => closed === 1, "idle websocket was not pruned")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(2)
fetch.close()
})
test("invalidates but does not reuse a socket after terminal failure frames", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
socket.send(JSON.stringify({ type: connections === 1 ? "response.failed" : "response.completed" }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest())
expect(await first.text()).toContain('data: {"type":"response.failed"}')
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain('data: {"type":"response.completed"}')
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(0)
fetch.close()
})
test("returns initial websocket error frames as HTTP-style API errors", async () => {
const error = {
type: "invalid_request_error",
message: "The model is not supported when using Codex with a ChatGPT account.",
}
const event = {
type: "error",
status: 400,
error,
headers: {
"x-codex-primary-window-minutes": 15,
ignored: { nested: true },
},
}
await using server = await createWebSocketServer((socket) => {
socket.once("message", () => {
socket.send(JSON.stringify(event))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const response = await fetch(server.url, streamRequest())
expect(response.status).toBe(400)
expect(response.headers.get("content-type")).toContain("application/json")
expect(response.headers.get("x-codex-primary-window-minutes")).toBe("15")
expect(response.headers.get("ignored")).toBeNull()
expect(await response.json()).toEqual(event)
fetch.close()
})
test("fails mid-stream wrapped websocket errors as HTTP-style API errors", async () => {
const event = {
type: "error",
status_code: 429,
error: {
type: "usage_limit_reached",
message: "The usage limit has been reached",
},
headers: {
"x-codex-primary-used-percent": "100.0",
},
}
await using server = await createWebSocketServer((socket) => {
socket.once("message", () => {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
socket.send(JSON.stringify(event))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const response = await fetch(server.url, streamRequest())
const error = await readTextError(response.text())
expect(APICallError.isInstance(error)).toBe(true)
if (!APICallError.isInstance(error)) throw new Error("Expected APICallError")
expect(error.statusCode).toBe(429)
expect(error.responseHeaders).toEqual({ "x-codex-primary-used-percent": "100.0" })
expect(error.responseBody).toBe(JSON.stringify(event))
fetch.close()
})
test("retries websocket connection limit errors on the next stream attempt", async () => {
let connections = 0
let messages = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
messages += 1
if (connections === 1) {
socket.send(
JSON.stringify({
type: "error",
status: 400,
error: {
type: "invalid_request_error",
code: "websocket_connection_limit_reached",
message: "Responses websocket connection limit reached",
},
}),
)
return
}
socket.send(JSON.stringify({ type: "response.completed", response: { id: "resp_retry" } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("Responses websocket connection limit reached")
const second = await fetch(server.url, streamRequest())
const text = await second.text()
expect(text).not.toContain("websocket_connection_limit_reached")
expect(text).toContain('data: {"type":"response.completed","response":{"id":"resp_retry"}}')
expect(text).toContain("data: [DONE]")
expect(connections).toBe(2)
expect(messages).toBe(2)
expect(server.httpRequests).toHaveLength(0)
fetch.close()
})
test("falls back to HTTP after websocket connection limit retries are exhausted", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
socket.send(
JSON.stringify({
type: "error",
status: 400,
error: {
type: "invalid_request_error",
code: "websocket_connection_limit_reached",
message: "Responses websocket connection limit reached",
},
}),
)
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
streamRetries: 2,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("Responses websocket connection limit reached")
const second = await fetch(server.url, streamRequest())
expect((await readTextError(second.text())).message).toContain("Responses websocket connection limit reached")
const third = await fetch(server.url, streamRequest())
const fourth = await fetch(server.url, streamRequest())
expect(await third.text()).toBe("http")
expect(await fourth.text()).toBe("http")
expect(connections).toBe(3)
expect(server.httpRequests).toHaveLength(2)
fetch.close()
})
test("shares the websocket retry budget across stream and connection limit failures", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
if (connections === 1) {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
socket.terminate()
return
}
socket.send(
JSON.stringify({
type: "error",
error: {
code: "websocket_connection_limit_reached",
message: "Responses websocket connection limit reached",
},
}),
)
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("WebSocket closed before response.completed")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(1)
fetch.close()
})
test("retries websocket idle failures before first event then falls back to HTTP", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
idleTimeout: 100, // kilocode_change - leave enough time for WebSocket callbacks on loaded CI runners
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("idle timeout waiting for websocket")
const second = await fetch(server.url, streamRequest())
const third = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(await third.text()).toBe("http")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(2)
fetch.close()
})
test("keeps websocket retry state until the failed stream becomes idle", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
idleTimeout: 500,
streamRetries: 1,
})
await new Promise((resolve) => setTimeout(resolve, 250))
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("idle timeout waiting for websocket")
await new Promise((resolve) => setTimeout(resolve, 300))
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(1)
fetch.close()
})
test("retries failed websocket streams before using HTTP fallback", async () => {
await using server = await createWebSocketServer((socket) => {
socket.once("message", () => {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
idleTimeout: 100, // kilocode_change - leave enough time for WebSocket callbacks on loaded CI runners
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("idle timeout waiting for websocket")
const second = await fetch(server.url, streamRequest())
expect((await readTextError(second.text())).message).toContain("idle timeout waiting for websocket")
const third = await fetch(server.url, streamRequest())
expect(await third.text()).toBe("http")
expect(server.httpRequests).toHaveLength(1)
fetch.close()
})
test("resets websocket stream failures after a completed response", async () => {
let connections = 0
let requests = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.on("message", () => {
requests += 1
if (requests === 1 || requests === 3) {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
socket.terminate()
return
}
socket.send(JSON.stringify({ type: "response.completed", response: { id: `resp_${requests}` } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("WebSocket closed before response.completed")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
const third = await fetch(server.url, streamRequest())
expect((await readTextError(third.text())).message).toContain("WebSocket closed before response.completed")
const fourth = await fetch(server.url, streamRequest())
expect(await fourth.text()).toContain("data: [DONE]")
expect(connections).toBe(3)
expect(requests).toBe(4)
expect(server.httpRequests).toHaveLength(0)
fetch.close()
})
test("falls back to HTTP for missing session and title requests", async () => {
await using server = await createWebSocketServer(() => {})
const fetch = OpenAIWebSocketPool.createWebSocketFetch()
const missingSession = await fetch(server.url, {
method: "POST",
headers: { [TITLE_HEADER]: "false" },
body: JSON.stringify({ stream: true }),
})
const title = await fetch(server.url, streamRequest({ [TITLE_HEADER]: "true" }))
expect(await missingSession.text()).toBe("http")
expect(await title.text()).toBe("http")
expect(server.httpRequests).toHaveLength(2)
expect(server.httpRequests[0]?.headers[TITLE_HEADER]).toBeUndefined()
expect(server.httpRequests[1]?.headers[TITLE_HEADER]).toBeUndefined()
fetch.close()
})
test("falls back to HTTP while a websocket lane is busy", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
})
})
const abort = new AbortController()
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest({}, abort.signal))
const firstText = first.text()
await waitFor(() => connections === 1, "websocket did not connect")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(server.httpRequests).toHaveLength(1)
expect(connections).toBe(1)
abort.abort(new Error("stop"))
expect((await readTextError(firstText)).message).toContain("stop")
fetch.close()
})
test("reserves a websocket lane while its socket is connecting", async () => {
await using server = await createHangingTcpServer()
await using fallback = await createHttpServer()
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
connectTimeout: 20,
streamRetries: 0,
})
const first = fetch(fallback.url, streamRequest())
await waitFor(() => server.connections() === 1, "first websocket did not begin connecting")
const second = fetch(fallback.url, streamRequest())
expect(await (await second).text()).toBe("http")
expect(await (await first).text()).toBe("http")
expect(server.connections()).toBe(1)
expect(fallback.httpRequests).toHaveLength(2)
fetch.close()
})
test("retries unexpected closes before first event then falls back to HTTP", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
socket.close(1001, "server shutdown")
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
streamRetries: 1,
})
const first = await fetch(server.url, streamRequest())
expect((await readTextError(first.text())).message).toContain("WebSocket closed before response.completed")
const second = await fetch(server.url, streamRequest())
const third = await fetch(server.url, streamRequest())
expect(await second.text()).toBe("http")
expect(await third.text()).toBe("http")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(2)
fetch.close()
})
test("does not keep HTTP fallback active after aborting a websocket response", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
if (connections === 1) {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
return
}
socket.send(JSON.stringify({ type: "response.completed", response: { id: "resp_456" } }))
})
})
const abort = new AbortController()
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest({}, abort.signal))
const firstText = first.text()
await waitFor(() => connections === 1, "first websocket did not connect")
abort.abort(new Error("stop"))
expect((await readTextError(firstText)).message).toContain("stop")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(0)
fetch.close()
})
test("releases the websocket lane when the response body is cancelled", async () => {
let connections = 0
await using server = await createWebSocketServer((socket) => {
connections += 1
socket.once("message", () => {
if (connections === 1) {
socket.send(JSON.stringify({ type: "response.output_text.delta", delta: "started" }))
return
}
socket.send(JSON.stringify({ type: "response.completed", response: { id: "resp_after_cancel" } }))
})
})
const fetch = OpenAIWebSocketPool.createWebSocketFetch({
url: server.url,
})
const first = await fetch(server.url, streamRequest())
await waitFor(() => connections === 1, "first websocket did not connect")
await first.body!.cancel("stop")
const second = await fetch(server.url, streamRequest())
expect(await second.text()).toContain("data: [DONE]")
expect(connections).toBe(2)
expect(server.httpRequests).toHaveLength(0)
fetch.close()
})
})
function streamRequest(headers?: Record<string, string>, signal?: AbortSignal): RequestInit {
return {
method: "POST",
headers: {
"session-id": "session-1",
authorization: "Bearer test",
...headers,
},
body: JSON.stringify({ stream: true, input: "hi" }),
signal,
}
}
async function readTextError(promise: Promise<string>) {
// Bun 1.3.14 hangs on expect(response.text()).rejects for streams errored from ws callbacks.
return promise.then(
() => {
throw new Error("Expected response text to reject")
},
(error) => {
expect(error).toBeInstanceOf(Error)
return error as Error
},
)
}
async function createWebSocketServer(onConnection: (socket: WebSocket, request: IncomingMessage) => void) {
const http = await createHttpServer()
const server = new WebSocketServer({ server: http.server })
server.on("connection", onConnection)
return websocketServerHandle(server, http)
}
async function createHangingTcpServer() {
const sockets = new Set<Socket>()
let connections = 0
const server = net.createServer((socket) => {
connections += 1
sockets.add(socket)
socket.on("close", () => sockets.delete(socket))
})
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve))
const address = server.address() as AddressInfo
return {
url: `http://127.0.0.1:${address.port}/v1/responses`,
wsUrl: `ws://127.0.0.1:${address.port}/v1/responses`,
connections: () => connections,
async [Symbol.asyncDispose]() {
for (const socket of sockets) socket.destroy()
server.close()
},
}
}
async function createRejectingWebSocketServer(onAttempt: () => void) {
const http = await createHttpServer()
const server = new WebSocketServer({
server: http.server,
verifyClient(_info, callback) {
onAttempt()
callback(false, 401, "denied")
},
})
return websocketServerHandle(server, http)
}
async function createHttpServer() {
const httpRequests: IncomingMessage[] = []
const server = createServer((request, response) => {
httpRequests.push(request)
response.writeHead(200, { "content-type": "text/plain" })
response.end("http")
})
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve))
const address = server.address() as AddressInfo
return {
server,
httpRequests,
url: `http://127.0.0.1:${address.port}/v1/responses`,
async [Symbol.asyncDispose]() {
await closeHttpServer(server)
},
}
}
function websocketServerHandle(server: WebSocketServer, http: Awaited<ReturnType<typeof createHttpServer>>) {
return {
url: http.url,
wsUrl: http.url.replace(/^http/, "ws"),
httpRequests: http.httpRequests,
async [Symbol.asyncDispose]() {
for (const socket of server.clients) socket.terminate()
server.close()
http.server.close()
},
}
}
function closeHttpServer(server: HttpServer) {
return new Promise<void>((resolve, reject) => server.close((error) => (error ? reject(error) : resolve())))
}
async function waitFor(predicate: () => boolean, message: string) {
const started = Date.now()
while (!predicate()) {
if (Date.now() - started > 1_000) throw new Error(message)
await new Promise((resolve) => setTimeout(resolve, 1))
}
}