diff --git a/devlog/_plan/260905_http_upstream_ws_parity/012_protocol_outcome.md b/devlog/_plan/260905_http_upstream_ws_parity/012_protocol_outcome.md new file mode 100644 index 00000000000..5f7e216e762 --- /dev/null +++ b/devlog/_plan/260905_http_upstream_ws_parity/012_protocol_outcome.md @@ -0,0 +1,7 @@ +# Protocol foundation outcome + +PR3643 merged at `87083e03422b6096d150232cbdf6066038f53383`, with source head `fded48f491809d781068bc410f70a6986f175355`. The actual merge tree matched the checked source tree and fetched `dev` ancestry was verified. The owner explicitly requested immediate admin merge without waiting for CI; pending/cancelled CI was not reported as green. + +Source-bound remote checks passed:145 tests,1 pre-existing skip,0 failures,715 assertions; typecheck and real synthetic HTTP-to-WS QA passed, with process/listener/home cleanup. The earlier hosted CLI-test timeout remains unassigned to a cause; an auxiliary non-reproduction did not close it. No installation, service, home configuration or credential mutation was performed. + +Next is020: bounded connection lifecycle/reuse. It must preserve the implemented metadata/stream contract and the subsequently landed `beforeDispatch` admission checks. No provider billing conclusion follows from either phase. diff --git a/devlog/_plan/260905_http_upstream_ws_parity/020_lifecycle.md b/devlog/_plan/260905_http_upstream_ws_parity/020_lifecycle.md index 0e16ac72967..fd016d00252 100644 --- a/devlog/_plan/260905_http_upstream_ws_parity/020_lifecycle.md +++ b/devlog/_plan/260905_http_upstream_ws_parity/020_lifecycle.md @@ -2,11 +2,22 @@ Depends on: protocol cycle and its verified request/metadata owner. P must re-read this document and current source after that PR lands. +## Landed-source refresh and loop specification + +Protocol PR3643 is landed; this phase starts from published `bf58ef1824e7b827b2a6bc1a5effb5d36ce80180`. Class C4, spec-satisfaction loop. Goal: eligible full HTTP requests reuse a canonical upstream socket without mixing exchanges. No-code/configuration cannot provide reuse because the existing transport unconditionally closes every terminal; unrelated provider pools speak different protocols. Reuse the existing exchange implementation, not a second relay. Keep frontend transport, native identity, full histories, admission/pacing, and installation unchanged. + +Resource scope: main implementation/audits under the user's no-other-task-communication instruction; no model override, new dependency, provider call or local full suite. Use disposable remote focused tests, typecheck and real loopback HTTP/WS QA plus full current-head CI. Initial wall-clock audit horizon is six hours from the explicit follow-up request. Evidence lives in this unit and the bound goalplan. A main audit is labelled as such; automatic PR review is a separate source. Prior immediate-admin permission closed PR3643 without waiting; it is not a green CI result for this new change. + +Source refresh adds a load-bearing requirement: `beforeDispatch` now guards credentials both before dialing and immediately before each frame. Reuse must call the fresh request's guard, including on a warm socket, and a refusal cannot enter HTTP fallback. Existing exports and Response markers remain compatible. The verified protocol command selects WS, account-attribution, metadata-integrity, reframing, cancellation and core/Lab tests; new lifecycle coverage gets an explicitly registered test file. Positive/negative activation evidence below, not a count alone, closes this phase. + ## File-change map | Operation | Path | Exact change | | --- | --- | --- | | NEW | `src/server/responses/codex-ws-session.ts` | Own one WS connection, exclusive in-flight exchange, per-exchange listeners, bounded queue, and terminal/cancel cleanup. Extract the existing one-shot state machine rather than duplicating it. | +| NEW | `src/server/responses/codex-ws-exchange.ts` | Extract the existing single-request relay/metadata/fallback state machine; both retained and one-shot sessions call this exact owner. | +| NEW | `src/server/responses/codex-ws-wire.ts` | Own unchanged frame limits, event normalization and Response markers; facade re-exports preserve existing callers without a circular import. | +| NEW | `src/server/responses/codex-ws-correlation.ts` | Per-exchange response/item correlation for retained sessions; a first incompatible exchange remains one-shot, reused incompatible traffic fails closed. | | NEW | `src/server/responses/codex-ws-pool.ts` | Own bounded idle sessions, canonical eligibility/keying, idle/max-age expiry, admission fallback and shutdown registration. No configuration/auth-store imports. | | MODIFY | `src/server/responses/codex-ws-request.ts` | Project genuine turn-state/turn-metadata headers into absent per-frame metadata slots before final serialization and byte-cap checks; identity consumes that exact prepared frame. | | MODIFY | `src/server/responses/ws-upstream.ts` | Keep the existing public entrypoint as compatibility facade; acquire an eligible idle canonical session or use the existing one-shot behavior, then send the prepared full frame. | @@ -36,6 +47,30 @@ Keep the actual declaration shaped to existing Bun/Web types; no dependency or f ## Identity and eligibility contract +### Source comparison amendment (2026-09-05) + +Reference checkout: `openai/codex` at `d2d5b70241fb448044c1c088a977cc720d70443a`. +`core/src/client.rs:1358` checks unchanged request properties and an exact input +prefix before deriving an incremental payload; `:1893` pairs that payload with +the actual previous response id. `codex-api/src/endpoint/responses_websocket.rs:299` +holds an exclusive stream lock and `:826` finishes the serial read on completion. +This patch implements socket reuse, not that input-delta algorithm. + +The [official WebSocket guide](https://developers.openai.com/api/docs/guides/websocket-mode) +describes connection-local continuation caches, optional named lanes, and full-input +recovery after cache loss. It documents the public API, not a guarantee for the +ChatGPT backend's private beta protocol. Our conservative subset is serial, +default-lane, complete-input creates. Non-null `previous_response_id`, explicit +`stream_id`, `generate`, or active background mode do not enter this pool. Their +existing one-shot behavior is unchanged; this is not new continuation support. +No steering or cross-lane fork is emitted. A named-lane frame cannot qualify a +connection for reuse. Earlier five-minute/32-exchange retirement is a local resource +policy, not an OpenAI limit, and never discards a continuation id invented by us. + +Remaining validation before readiness: real backend compatibility is not proven +by public docs or mocks; rotating immutable handshake headers may prevent reuse. +No billing/quota causation claim follows from this work. + Reuse is canonical-URL-only and requires usable selected outbound auth/account identity plus explicit thread and turn identities. Missing either identity stays one-shot, so unrelated native turns cannot share a session. A mere model slug or account log label is not a reuse key. Compute an in-memory nonlogged digest of the selected credential/account, conversation/turn scope, actual model/tier, and immutable handshake policy. No raw credential, account id or prompt is emitted in logs, receipt data, exported diagnostics, or persisted cache. Different credentials, account, model/tier, originator/beta/attestation policy, or incompatible handshake headers must never reuse a socket. @@ -58,10 +93,12 @@ The initial implementation sends each complete HTTP request as a complete `respo - Hard cap: 32 retained canonical sessions; at most one active exchange per retained session. On a busy key, use a separately owned one-shot connection, not an unbounded waiter queue or concurrent send on that socket. Global turn admission remains authoritative. - Idle TTL: 30 seconds. Maximum connection age: 5 minutes. Named constants live in the pool owner; fake-clock tests cross exact boundaries. +- Maximum successful exchanges per retained socket: 32, bounding remembered response ids. Expired or superseded active exchanges may finish but are retired at release; age expiry does not kill an in-flight generation merely to free capacity. Correlation ids are bounded to 4096 bytes and item tracking to 10000 items; no prompt/output history is retained. - No timer before first activation. Expiry uses bounded owned timers with `unref` where available; every timer/listener is cleared on disposal. Register one shutdown hook on activation and detach when the pool is fully disposed. - Evict oldest idle entries before retaining a new one. Never evict/steal a live exchange merely to make room; use the existing one-shot bounded path. - Successful terminal closes the exchange stream and releases a reusable socket only after its bounded terminal frame is enqueued. Failed/incomplete/error outcomes are conservatively disposed, not reused. - Request abort removes that exchange's listener, errors its body exactly once, closes its socket, and releases its ownership. A completed request's later abort must not close a session leased to a successor request. +- Per-exchange listeners and quota callbacks detach before release. Session-level listeners handle idle unsolicited data and physical closure only. Explicit pool shutdown settles active requests without HTTP fallback; unexpected pre-send upgrade failure retains the original fallback. Optional socket ref/unref hints do not replace deterministic cleanup. - Closing/error sockets are removed immediately. Reconnect/retry is allowed only before a frame was accepted for send; once inference may have started, do not fall back to HTTP and double-generate. Keep existing send-throw/upgrade failure semantics only when the no-send condition is proven. - Per-frame and per-exchange queue limits remain the existing limits. Connection reuse does not retain completed queues or prior output. - Shutdown closes idle and active pool-owned sockets, settles all requests, and unregisters timers. It cannot import Lab, block synchronous startup, or make unrelated providers start a timer. @@ -97,4 +134,28 @@ Security analysis and negative-case reasoning are maintained in ignored scratch. ## Delivery -Run the focused transport and integration suite, typecheck, privacy/secret checks, docs build, independent adversarial review, and the coordinator-approved full check. The PR remains pending until exact-head required checks and required review are satisfied. Land after protocol, prove ancestry, then close the unit and final goal. No production service restart or link occurs. +Run the focused transport and integration suite, typecheck, privacy/secret checks, +and exact-head CI. Main audits obey the user's no-other-task-communication boundary; +do not represent them as independent security review. Publish with `--no-verify` +as a draft PR targeting dev while verification or review remains outstanding. +The subsequent owner instruction authorizes real selected-account verification +and merging after the remaining checks. Preserve the earlier PR-only publication +record as history, not a current merge prohibition. Start live checks with at most +24 creates, 512 requested output tokens each, concurrency at most two, and a +60-minute diagnostic horizon. Use a separate local process and read-only selected +credentials; no refresh, persisted credentials, secret logs or live proxy changes. +Keep required CI/review evidence truthful and prove fetched merge ancestry. +No production service restart or link occurs. + +The independent A-B-A review trace was disproven in an ignored exact-head probe: +return-null cannot fall through to entry replacement. Add permanent facade +coverage to guard that intended retirement behavior; do not change correct pool +ownership merely to satisfy the proposed explanation. + +Live backend validation also exposed pretty-printed error frames: the existing +relay prefixed only the first physical JSON line with SSE `data:`, so ordinary +SSE readers could not parse the error. Normalize physical CR/LF JSON formatting +inside the existing wire owner before SSE framing, preserving the parsed fields +and all size limits. Cover both errors and completed responses; compact native +frames remain byte-preserved. This is an observed wire fix, not an inferred +quota-accounting change. diff --git a/docs-site/src/content/docs/reference/configuration/server.md b/docs-site/src/content/docs/reference/configuration/server.md index c6994f74a2e..832817a2058 100644 --- a/docs-site/src/content/docs/reference/configuration/server.md +++ b/docs-site/src/content/docs/reference/configuration/server.md @@ -19,7 +19,7 @@ runs helper features around provider requests. | `oauthOpenBrowser?` | `boolean` | `true` | Whether a login may open a browser on the machine running the proxy. Absent and `true` both open, so an existing install is unchanged; only an explicit `false` declines. Decline when you need the authorization link in a different browser profile, or when the dashboard is not on the proxy's machine — the login still starts and the URL is still returned and displayed. `POST /api/oauth/login` and `POST /api/codex-auth/login` accept a per-request `openBrowser` boolean that overrides this, and the dashboard exposes the same choice beside the login button. Device-code flows never open a browser either way. | | `connectTimeoutMs?` | `number` | `200000` | Per-attempt DNS/TCP/TLS/final-header deadline; it ends before body generation. | | `shutdownTimeoutMs?` | `number` | `5000` | Graceful drain deadline before active turns are aborted. | -| `websockets?` | `boolean` | `false` | Advertise and admit the client-facing Responses WebSocket path. False keeps clients on HTTP/SSE; it does not disable an eligible canonical ChatGPT upstream WS optimization. | +| `websockets?` | `boolean` | `false` | Advertise and admit the client-facing Responses WebSocket path. False keeps clients on HTTP/SSE; it does not disable an eligible canonical ChatGPT upstream WS optimization. Complete-input requests may reuse an upstream connection within the same selected credential, account, thread and turn; changed handshake policy or missing identity keeps requests on separate connections. This does not trim HTTP input or create previous-response IDs. | | `corsAllowOrigins?` | `string[]` | `[]` | Additional exact origins allowed by CORS. Loopback origins are always allowed. Authority-based browser extension origins such as `chrome-extension://` are supported; `*` is not a wildcard. Firefox and Safari regenerate the extension UUID (per install / per browser launch), so update the entry when the origin changes. | | `apiKeys?` | `OcxApiKey[]` | `[]` | Generated `ocx_…` credentials accepted by management and data-plane auth on non-loopback binds. Dashboard-managed. | | `storageCleanupPolicy?` | `StorageCleanupPolicy` | disabled | Opt-in archived-session cleanup policy. Never enabled implicitly. | diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 7d312309ce7..29dd2c5f1ca 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -1248,6 +1248,7 @@ "winsw-stop-hardening.test.ts": "windows", "winsw.test.ts": "service", "ws-endpoint.test.ts": "responses", + "ws-upstream-reuse.test.ts": "responses", "ws-upstream.test.ts": "responses", "xai-client.test.ts": "images", "xai-oauth-retry.test.ts": "providers/xai", diff --git a/src/server/responses/codex-ws-correlation.ts b/src/server/responses/codex-ws-correlation.ts new file mode 100644 index 00000000000..d4f467cb5c9 --- /dev/null +++ b/src/server/responses/codex-ws-correlation.ts @@ -0,0 +1,65 @@ +export const CODEX_WS_ID_MAX_BYTES = 4096; +export const CODEX_WS_MAX_TRACKED_ITEMS = 10_000; +const MAX_TRACKED_ID_BYTES = 1024 * 1024; + +function record(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} +function id(value: unknown): value is string { + return typeof value === "string" && value.length > 0 && !/[\u0000-\u001f\u007f]/.test(value) + && Buffer.byteLength(value) <= CODEX_WS_ID_MAX_BYTES; +} + +/** A cold incompatible response may finish one-shot; a reused socket cannot mix owners. */ +export class CodexWsCorrelation { + private responseId: string | null = null; + private reusable = true; + private readonly items = new Set(); + private itemBytes = 0; + + constructor(private readonly strict: boolean, private readonly previouslyCompleted: (id: string) => boolean) {} + + private mismatch(): void { + this.reusable = false; + if (this.strict) throw new Error("codex websocket response identity mismatch"); + } + + accept(event: Record): void { + if (event.stream_id !== undefined) { this.mismatch(); return; } + if (event.type === "error") return; + const response = record(event.response) ? event.response : undefined; + if (event.type === "response.created") { + const next = response?.id; + if (!id(next) || this.responseId !== null || this.previouslyCompleted(next)) { + this.mismatch(); + return; + } + this.responseId = next; + return; + } + if (!this.responseId) { this.mismatch(); return; } + if ((response && response.id !== this.responseId) + || (event.response_id !== undefined && event.response_id !== this.responseId)) this.mismatch(); + const item = record(event.item) ? event.item : undefined; + if (event.type === "response.output_item.added") { + if (!id(item?.id) || this.items.has(item.id)) { this.mismatch(); return; } + this.itemBytes += Buffer.byteLength(item.id); + if (this.items.size >= CODEX_WS_MAX_TRACKED_ITEMS || this.itemBytes > MAX_TRACKED_ID_BYTES) { + throw new Error("codex websocket correlation exceeds its bounded item budget"); + } + this.items.add(item.id); + return; + } + const itemId = event.item_id ?? item?.id; + if (itemId !== undefined && (!id(itemId) || !this.items.has(itemId))) this.mismatch(); + if (typeof event.type === "string" && event.type.endsWith(".delta") && itemId === undefined) this.mismatch(); + } + + completed(event: Record): string | null { + const response = record(event.response) ? event.response : undefined; + return this.reusable && event.type === "response.completed" && response?.status === "completed" + && response.id === this.responseId ? this.responseId : null; + } + + finish(): void { this.items.clear(); this.itemBytes = 0; } +} diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts new file mode 100644 index 00000000000..2f41be02fab --- /dev/null +++ b/src/server/responses/codex-ws-exchange.ts @@ -0,0 +1,261 @@ +import { MAX_CLIENT_SSE_FRAME_BYTES } from "../sse-frame-buffer"; +import { CodexWsMetadata, type CodexWsQuotaObserver } from "./codex-ws-metadata"; +import { CODEX_RESPONSES_HTTP_URL, type PreparedCodexWsRequest } from "./codex-ws-request"; +import { CodexWsCorrelation } from "./codex-ws-correlation"; +import type { CodexWsSession } from "./codex-ws-session"; +import { UPGRADE_DEADLINE_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, + MAX_CODEX_WS_QUEUE_BYTES, markCodexWsResponse, normalizeResponsesWsRelayEvent, closedBeforeTerminalMessage } from "./codex-ws-wire"; + +interface ExchangeOptions { + session: CodexWsSession; + url: string; + init: RequestInit; + prepared: PreparedCodexWsRequest; + sseFallback: typeof globalThis.fetch; + onQuota?: CodexWsQuotaObserver; + beforeDispatch?: (headers: Headers) => void; +} + +/** The sole SSE exchange state machine for both one-shot and retained sockets. */ +export function codexWsExchange(options: ExchangeOptions): Promise { + const { session, url, init, prepared, sseFallback, onQuota, beforeDispatch } = options; + const { frameText, headers } = prepared; + const signal = init.signal ?? undefined; + return new Promise((resolve, reject) => { + const ws = session.socket; + + let opened = session.opened; + let settledPreOpen = false; + let sent = false; + let received = false; + let responseCommitted = false; + let terminal = false; + let controller: ReadableStreamDefaultController | null = null; + const encoder = new TextEncoder(); + const metadata = url === CODEX_RESPONSES_HTTP_URL ? new CodexWsMetadata(onQuota) : null; + const correlation = session.retainable ? new CodexWsCorrelation(session.reused, id => session.hasCompleted(id)) : null; + let detachOwner = () => {}; + let preludeTimer: ReturnType | undefined; + const stream = new ReadableStream({ + start(c) { controller = c; }, + cancel() { + if (terminal) return; + terminal = true; + cleanup(); + session.dispose(); + }, + }, new ByteLengthQueuingStrategy({ highWaterMark: MAX_CODEX_WS_QUEUE_BYTES })); + + const cleanup = () => { + clearTimeout(upgradeTimer); + clearTimeout(preludeTimer); + signal?.removeEventListener("abort", onAbort); + metadata?.finish(); + correlation?.finish(); + detachOwner(); + ws.removeEventListener("open", onOpen); + ws.removeEventListener("message", onMessage); + ws.removeEventListener("close", onClose); + ws.removeEventListener("error", onError); + }; + + const commitResponse = () => { + if (responseCommitted) return; + responseCommitted = true; + clearTimeout(preludeTimer); + const responseHeaders = metadata?.snapshot() ?? new Headers(); + responseHeaders.set("content-type", "text/event-stream; charset=utf-8"); + const response = new Response(stream, { status: 200, headers: responseHeaders }); + metadata?.commit(); + markCodexWsResponse(response, Boolean(metadata && onQuota)); + resolve(response); + }; + + const failStream = (error: unknown) => { + if (terminal) return; + terminal = true; + // A frame may already be executing upstream. Settle as a body failure, + // never a fetch rejection/5xx that the pre-stream wrapper could resend. + if (sent) commitResponse(); + cleanup(); + try { controller?.error(typeof error === "string" ? new Error(error) : error); } catch { /* stream already done */ } + session.dispose(); + }; + + const upgradeTimer = setTimeout(() => { + if (opened || settledPreOpen) return; + settledPreOpen = true; + cleanup(); + session.dispose(); + resolve(sseFallback(url, init)); + }, UPGRADE_DEADLINE_MS); + + const cancelExchange = (reason: unknown) => { + if (terminal || settledPreOpen) return; + if (!sent) { + settledPreOpen = true; + terminal = true; + cleanup(); + session.dispose(); + reject(reason); + return; + } + failStream(reason); + }; + const onAbort = () => cancelExchange(signal?.reason ?? new DOMException("The operation was aborted.", "AbortError")); + signal?.addEventListener("abort", onAbort, { once: true }); + + const onOpen = () => { + if (settledPreOpen) return; + clearTimeout(upgradeTimer); + opened = true; + try { + beforeDispatch?.(new Headers(headers)); + } catch (error) { + // Settle and detach before close: a synchronous close event must not resend over SSE. + settledPreOpen = true; + terminal = true; + cleanup(); + ws.removeEventListener("open", onOpen); + ws.removeEventListener("message", onMessage); + ws.removeEventListener("close", onClose); + ws.removeEventListener("error", onError); + session.dispose(); + reject(error); + return; + } + if (terminal || settledPreOpen || signal?.aborted) return; + sent = true; + try { + ws.send(frameText); + } catch { + if (received || responseCommitted) { + if (terminal) session.dispose(); + failStream("codex websocket send failed after response activity"); + return; + } + // send() throwing means the frame never left, so no upstream turn + // started and the SSE resend cannot double-generate. Falling back + // (instead of erroring a synthetic 200 body) keeps the pre-stream + // HTTP error/refresh/failover machinery in charge. + settledPreOpen = true; + sent = false; + cleanup(); + session.dispose(); + resolve(sseFallback(url, init)); + return; + } + if (!metadata) commitResponse(); + else if (!responseCommitted && !terminal) { + preludeTimer = setTimeout(() => failStream("codex websocket response prelude timed out"), CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS); + } + }; + + const onMessage = (event: MessageEvent) => { + if (!controller || terminal) return; + received = true; + const text = typeof event.data === "string" ? event.data : ""; + if (!text) return; + // UTF-8 byte length is always at least the JS string length. Reject this + // cheap lower bound before parsing so an obviously oversized frame does + // not create another large object graph. + if (text.length > MAX_CODEX_WS_FRAME_BYTES) { + failStream("codex websocket frame exceeds the response size limit"); + return; + } + const rawEncodedText = encoder.encode(text); + if (rawEncodedText.byteLength > MAX_CODEX_WS_FRAME_BYTES) { + failStream("codex websocket frame exceeds the response size limit"); + return; + } + const normalized = normalizeResponsesWsRelayEvent(text); + if (!normalized) return; + const { type } = normalized; + let relayText = normalized.text; + let controlFrame = false; + if (metadata) { + try { + const sanitized = metadata.consume(normalized.payload, rawEncodedText.byteLength); + if (sanitized !== null) { + relayText = sanitized; + controlFrame = true; + } + } catch (error) { + failStream(error); + return; + } + } + const encodedText = relayText === text ? rawEncodedText : encoder.encode(relayText); + if (encodedText.byteLength > MAX_CODEX_WS_FRAME_BYTES) { + failStream("codex websocket frame exceeds the response size limit"); + return; + } + if (!controlFrame && !type.startsWith("response.") && type !== "error") return; + if (!controlFrame) { + try { correlation?.accept(normalized.payload); } catch (error) { failStream(error); return; } + commitResponse(); + } + const prefix = encoder.encode(`event: ${type}\ndata: `); + const suffix = encoder.encode("\n\n"); + const frameBytes = prefix.byteLength + encodedText.byteLength + suffix.byteLength; + if (frameBytes > MAX_CLIENT_SSE_FRAME_BYTES) { + failStream("codex websocket frame exceeds the response size limit"); + return; + } + const availableBytes = controller.desiredSize ?? 0; + if (frameBytes > availableBytes) { + failStream("codex websocket response exceeded the buffered queue limit"); + return; + } + const sseFrame = new Uint8Array(frameBytes); + sseFrame.set(prefix); + sseFrame.set(encodedText, prefix.byteLength); + sseFrame.set(suffix, prefix.byteLength + encodedText.byteLength); + try { + controller.enqueue(sseFrame); + } catch { + failStream("codex websocket response stream closed while enqueueing"); + return; + } + if (type === "response.completed" || type === "response.failed" || type === "response.incomplete" || type === "error") { + const completedId = correlation?.completed(normalized.payload) ?? null; + terminal = true; + cleanup(); + try { controller.close(); } catch { /* already closed */ } + session.release(completedId); + } + }; + + const onClose = (event: unknown) => { + cleanup(); + if (!opened) { + if (settledPreOpen) return; + settledPreOpen = true; + // Upgrade rejected (401/403/429/5xx). Retry over plain SSE so the real + // HTTP status reaches the existing refresh/rotation handlers. No turn + // started upstream, so the resend cannot double-generate. + resolve(sseFallback(url, init)); + return; + } + if (sent && !terminal) failStream(closedBeforeTerminalMessage(event)); + }; + + const onError = () => { + if (terminal || settledPreOpen) return; + if (!opened && !sent) { + settledPreOpen = true; + terminal = true; + cleanup(); + session.dispose(); + resolve(sseFallback(url, init)); + } else failStream("codex websocket transport error"); + }; + detachOwner = session.bindOwner(reason => cancelExchange(reason)); + ws.addEventListener("open", onOpen); + ws.addEventListener("message", onMessage); + ws.addEventListener("close", onClose); + ws.addEventListener("error", onError); + if (signal?.aborted) onAbort(); + else if (session.opened) onOpen(); + }); +} diff --git a/src/server/responses/codex-ws-pool.ts b/src/server/responses/codex-ws-pool.ts new file mode 100644 index 00000000000..378cf2d4a39 --- /dev/null +++ b/src/server/responses/codex-ws-pool.ts @@ -0,0 +1,162 @@ +import { createHmac, randomBytes } from "node:crypto"; +import { registerOptionalShutdownHook } from "../../lib/optional-shutdown-hooks"; +import { CODEX_RESPONSES_HTTP_URL } from "./codex-ws-request"; +import { CODEX_WS_ID_MAX_BYTES } from "./codex-ws-correlation"; +import { CodexWsSession } from "./codex-ws-session"; + +export const CODEX_WS_POOL_MAX_SESSIONS = 32; +export const CODEX_WS_POOL_IDLE_MS = 30_000; +export const CODEX_WS_POOL_MAX_AGE_MS = 5 * 60_000; +const MUTABLE_HEADERS = new Set(["x-codex-turn-state", "x-codex-turn-metadata"]); +let processKey: Buffer | undefined; +let poolSequence = 0; + +export interface CodexWsReuseIdentity { key: string; scope: string } +function value(value: unknown): value is string { + return typeof value === "string" && value.trim().length > 0 + && !/[\u0000-\u001f\u007f]/.test(value) && Buffer.byteLength(value) <= CODEX_WS_ID_MAX_BYTES; +} +function record(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} +function digest(input: unknown): string { + processKey ??= randomBytes(32); + return createHmac("sha256", processKey).update(JSON.stringify(input)).digest("hex"); +} + +/** Identity comes from the selected outgoing request, never a model label or caller hint. */ +export function codexWsReuseIdentity(url: string, headers: Record, frameText: string): CodexWsReuseIdentity | null { + if (url !== CODEX_RESPONSES_HTTP_URL) return null; + let body: unknown; + try { body = JSON.parse(frameText); } catch { return null; } + if (!record(body) || !record(body.client_metadata)) return null; + // Reuse complete HTTP creates only. Cache-dependent continuation, warmup and + // named WS lanes need a different lifecycle/recovery contract. + if (body.previous_response_id != null || Object.hasOwn(body, "stream_id") + || Object.hasOwn(body, "generate") || body.background === true) return null; + const metadata = body.client_metadata; + const bodyThread = metadata.thread_id; + const headerThread = headers["thread-id"]; + if (bodyThread !== undefined && !value(bodyThread)) return null; + if (headerThread !== undefined && !value(headerThread)) return null; + if (bodyThread !== undefined && headerThread !== undefined && bodyThread !== headerThread) return null; + const thread = bodyThread ?? headerThread; + const turn = metadata.turn_id; + const account = headers["chatgpt-account-id"]; + const authorization = headers.authorization; + if (![thread, turn, account, authorization, body.model].every(value)) return null; + if (body.service_tier !== undefined && !value(body.service_tier)) return null; + const immutable = Object.entries(headers).filter(([name]) => !MUTABLE_HEADERS.has(name)).sort(([a], [b]) => a.localeCompare(b)); + if (immutable.length > 128 || immutable.some(([, field]) => !value(field)) + || immutable.reduce((bytes, [name, field]) => bytes + Buffer.byteLength(name) + Buffer.byteLength(field), 0) > 32 * 1024) return null; + const scope = digest([url, account, thread, turn]); + const lite = metadata.ws_request_header_x_openai_internal_codex_responses_lite; + if (lite !== undefined && lite !== "true" && lite !== "false") return null; + return { scope, key: digest([scope, authorization, body.model, body.service_tier ?? null, lite ?? null, immutable]) }; +} + +interface Entry { identity: CodexWsReuseIdentity; session: CodexWsSession; createdAt: number; idleAt: number; retired: boolean } +interface PoolOptions { now?: () => number; maxSessions?: number; idleMs?: number; maxAgeMs?: number } + +/** Bounded retained sockets only. Busy/capacity misses keep the existing one-shot path. */ +export class CodexWsPool { + private readonly entries = new Map(); + private timer?: ReturnType; + private detachShutdown?: () => void; + private readonly hookKey = `codex-upstream-ws-pool-${++poolSequence}`; + private readonly now: () => number; + private readonly maxSessions: number; + private readonly idleMs: number; + private readonly maxAgeMs: number; + constructor(options: PoolOptions = {}) { + this.now = options.now ?? Date.now; + this.maxSessions = options.maxSessions ?? CODEX_WS_POOL_MAX_SESSIONS; + this.idleMs = options.idleMs ?? CODEX_WS_POOL_IDLE_MS; + this.maxAgeMs = options.maxAgeMs ?? CODEX_WS_POOL_MAX_AGE_MS; + } + + acquire(identity: CodexWsReuseIdentity, url: string, headers: Record): CodexWsSession | null { + this.sweep(); + for (const entry of this.entries.values()) { + if (entry.identity.scope !== identity.scope || entry.identity.key === identity.key) continue; + entry.retired = true; + if (!entry.session.busy) this.remove(entry); + } + const existing = this.entries.get(identity.key); + if (existing) { + if (existing.retired || existing.session.busy) return null; + if (existing.session.reserve()) { this.arm(); return existing.session; } + this.remove(existing); + } + if (this.entries.size >= this.maxSessions) { + const oldest = [...this.entries.values()].filter(entry => !entry.session.busy).sort((a, b) => a.idleAt - b.idleAt)[0]; + if (!oldest) return null; + this.remove(oldest); + } + const createdAt = this.now(); + const session = new CodexWsSession(url, headers, true, () => this.changed(entry)); + const entry: Entry = { identity, session, createdAt, idleAt: createdAt, retired: false }; + session.reserve(); + this.entries.set(identity.key, entry); + this.detachShutdown ??= registerOptionalShutdownHook(this.hookKey, () => this.dispose()); + return session; + } + + private changed(entry: Entry): void { + if (this.entries.get(entry.identity.key) !== entry) return; + if (entry.session.closed) this.entries.delete(entry.identity.key); + else if (!entry.session.busy) { + entry.idleAt = this.now(); + if (entry.retired || entry.idleAt - entry.createdAt >= this.maxAgeMs) this.remove(entry); + } + this.arm(); + } + + private remove(entry: Entry): void { + if (this.entries.get(entry.identity.key) === entry) this.entries.delete(entry.identity.key); + entry.session.dispose(new Error("codex websocket retained session expired")); + this.arm(); + } + + sweep(): void { + const now = this.now(); + for (const entry of this.entries.values()) { + if (!entry.session.busy && (entry.session.closed || entry.retired + || now - entry.idleAt >= this.idleMs || now - entry.createdAt >= this.maxAgeMs)) this.remove(entry); + } + this.arm(); + } + + private arm(): void { + clearTimeout(this.timer); + this.timer = undefined; + if (!this.entries.size) { + this.detachShutdown?.(); + this.detachShutdown = undefined; + return; + } + let deadline = Infinity; + for (const entry of this.entries.values()) if (!entry.session.busy) { + deadline = Math.min(deadline, entry.idleAt + this.idleMs, entry.createdAt + this.maxAgeMs); + } + if (!Number.isFinite(deadline)) return; + this.timer = setTimeout(() => { this.timer = undefined; this.sweep(); }, Math.max(1, deadline - this.now())); + this.timer.unref?.(); + } + + dispose(): void { + clearTimeout(this.timer); + this.timer = undefined; + this.detachShutdown?.(); + this.detachShutdown = undefined; + const entries = [...this.entries.values()]; + this.entries.clear(); + for (const entry of entries) entry.session.dispose(new DOMException("codex websocket pool shutdown", "AbortError")); + } + + snapshot(): { size: number; active: number; timer: boolean } { + return { size: this.entries.size, active: [...this.entries.values()].filter(entry => entry.session.busy).length, timer: this.timer !== undefined }; + } +} + +export const codexWsPool = new CodexWsPool(); diff --git a/src/server/responses/codex-ws-request.ts b/src/server/responses/codex-ws-request.ts index f065de6e3aa..eaa77323e77 100644 --- a/src/server/responses/codex-ws-request.ts +++ b/src/server/responses/codex-ws-request.ts @@ -32,6 +32,13 @@ function applyLiteMetadata(body: Record, headers: Headers): boo body.client_metadata = { ...(metadata as Record | undefined), [CODEX_RESPONSES_LITE_METADATA_KEY]: lite }; } + for (const name of ["x-codex-turn-state", "x-codex-turn-metadata"]) { + const value = headers.get(name); + const current = body.client_metadata as Record | undefined; + if (value !== null && !Object.hasOwn(current ?? {}, name)) { + body.client_metadata = { ...current, [name]: value }; + } + } return true; } diff --git a/src/server/responses/codex-ws-session.ts b/src/server/responses/codex-ws-session.ts new file mode 100644 index 00000000000..bbf62f81373 --- /dev/null +++ b/src/server/responses/codex-ws-session.ts @@ -0,0 +1,93 @@ +export const MAX_CODEX_WS_SESSION_EXCHANGES = 32; + +/** Owns one physical socket; request listeners belong to the exchange, not this object. */ +export class CodexWsSession { + readonly socket: WebSocket; + opened = false; + closed = false; + busy = false; + private owner?: (reason: Error) => void; + private readonly completedIds = new Set(); + + constructor(url: string, headers: Record, readonly retainable = false, + private readonly changed: () => void = () => {}) { + this.socket = new WebSocket(url, { headers } as unknown as string[]); + this.socket.addEventListener("open", this.onOpen); + this.socket.addEventListener("message", this.onIdleMessage); + this.socket.addEventListener("close", this.onClose); + this.socket.addEventListener("error", this.onIdleError); + } + + get reused(): boolean { return this.completedIds.size > 0; } + hasCompleted(id: string): boolean { return this.completedIds.has(id); } + + reserve(): boolean { + if (this.closed || this.busy || (this.opened && this.socket.readyState !== undefined && this.socket.readyState !== 1)) return false; + this.busy = true; + const socket = this.socket as WebSocket & { ref?: () => void }; + try { socket.ref?.(); } catch { /* optional keepalive hint */ } + return true; + } + + bindOwner(owner: (reason: Error) => void): () => void { + if (!this.busy || this.closed || this.owner) throw new Error("codex websocket lease is unavailable"); + this.owner = owner; + return () => { if (this.owner === owner) this.owner = undefined; }; + } + + release(completedId: string | null): void { + this.owner = undefined; + if (this.closed) return; + if (!this.retainable || !completedId || !this.opened + || (this.socket.readyState !== undefined && this.socket.readyState !== 1)) { + this.dispose(); + return; + } + this.completedIds.add(completedId); + if (this.completedIds.size >= MAX_CODEX_WS_SESSION_EXCHANGES) { + this.dispose(); + return; + } + this.busy = false; + const socket = this.socket as WebSocket & { unref?: () => void }; + try { socket.unref?.(); } catch { /* optional hint; shutdown/expiry still owns cleanup */ } + this.changed(); + } + + dispose(reason = new Error("codex websocket session disposed")): void { + if (this.closed) return; + this.closed = true; + const owner = this.owner; + this.owner = undefined; + this.detach(); + try { owner?.(reason); } finally { + this.busy = false; + this.completedIds.clear(); + try { this.socket.close(); } catch { /* already closing */ } + if (this.retainable) { + try { (this.socket as WebSocket & { terminate?: () => void }).terminate?.(); } catch { /* already closed */ } + } + this.changed(); + } + } + + private onOpen = (): void => { this.opened = true; }; + private onIdleMessage = (): void => { + if (!this.busy) this.dispose(new Error("codex websocket received unsolicited idle data")); + }; + private onIdleError = (): void => { if (!this.busy) this.dispose(); }; + private onClose = (): void => { + this.closed = true; + this.busy = false; + this.completedIds.clear(); + this.detach(); + this.changed(); + // The active exchange's close listener retains pre-send fallback semantics. + }; + private detach(): void { + this.socket.removeEventListener("open", this.onOpen); + this.socket.removeEventListener("message", this.onIdleMessage); + this.socket.removeEventListener("close", this.onClose); + this.socket.removeEventListener("error", this.onIdleError); + } +} diff --git a/src/server/responses/codex-ws-wire.ts b/src/server/responses/codex-ws-wire.ts new file mode 100644 index 00000000000..996ea6ec55e --- /dev/null +++ b/src/server/responses/codex-ws-wire.ts @@ -0,0 +1,144 @@ +import { MAX_CLIENT_SSE_FRAME_BYTES } from "../sse-frame-buffer"; +// If the 101 never arrives (network black hole), give SSE a chance well before +// the caller's connect timeout (default 200s) would fire. +export const UPGRADE_DEADLINE_MS = 10_000; +export const CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS = 30_000; +// Keep the push-based WS transport inside the same memory envelope as the +// bounded SSE relays that consume this response. Unlike fetch response bodies, +// a WebSocket cannot be paused when a ReadableStream applies backpressure, so +// an upstream that outruns the consumer must be disconnected. +export const MAX_CODEX_WS_FRAME_BYTES = MAX_CLIENT_SSE_FRAME_BYTES; +export const MAX_CODEX_WS_QUEUE_BYTES = 8 * 1024 * 1024; +// The backend drops any inbound message of 16 MiB or more: it closes the socket +// (1009) without a Responses terminal event, which reaches clients as a bare +// 502 upstream_server_error. Measured against the live endpoint 2026-08-23: +// 16,777,000 B completed, 16,777,300 B closed in ~1s, every time. The same +// request body succeeds over HTTP SSE, so the ceiling belongs to this transport +// alone (see #2426). A full-replay thread reaches it with ~11 pasted +// screenshots, and then never recovers, because each retry resends the frame. +export const MAX_CODEX_WS_CREATE_FRAME_BYTES = 16 * 1024 * 1024; +// Bun frames the payload it is handed, so the send-side budget is the JSON text +// itself, and nothing is appended between the check and the send. The margin is +// a conservative cushion, not a computed requirement: it covers RFC 6455 frame +// overhead in case the backend counts it (14 bytes at this payload size — an +// 8-byte extended length plus a 4-byte client mask, leaving ~65.5 KiB spare), +// and it leaves room for a future caller that appends to the frame. +const CODEX_WS_CREATE_FRAME_MARGIN_BYTES = 64 * 1024; +export const CODEX_WS_CREATE_FRAME_LIMIT_BYTES = + MAX_CODEX_WS_CREATE_FRAME_BYTES - CODEX_WS_CREATE_FRAME_MARGIN_BYTES; +/** Close code the backend uses for an oversized message (RFC 6455 "message too big"). */ +const WS_CLOSE_MESSAGE_TOO_BIG = 1009; + +const codexWsUpstreamResponses = new WeakSet(); +const quotaObservedResponses = new WeakSet(); + +/** Quota arrived directly at its captured account; do not replay old HTTP prelude headers. */ +export function isCodexWsQuotaObservedResponse(response: Response): boolean { + return quotaObservedResponses.has(response); +} + +/** True only for a successful Codex WebSocket upgrade, never an HTTP fallback. */ +export function isCodexWsUpstreamResponse(response: Response): boolean { + return codexWsUpstreamResponses.has(response); +} + + +export function markCodexWsResponse(response: Response, observed: boolean): void { + codexWsUpstreamResponses.add(response); + if (observed) quotaObservedResponses.add(response); +} + +const CLOSED_BEFORE_TERMINAL = "codex websocket closed before a Responses terminal event"; + +export type ResponsesWsRelayEvent = { + type: string; + text: string; + payload: Record; +}; + +/** + * Responses WebSocket uses `response.done` as its terminal event, while the + * SSE Responses surface uses status-specific terminal events. Normalize the + * WS-only discriminator before relaying so the existing SSE consumers can + * settle the turn and the socket close cannot be mistaken for a drop. Unknown + * or missing status values fail closed instead of being reported as success. + */ +export function normalizeResponsesWsRelayEvent(text: string): ResponsesWsRelayEvent | null { + let payload: unknown; + try { + payload = JSON.parse(text); + } catch { + return null; + } + if (!payload || typeof payload !== "object" || Array.isArray(payload)) return null; + const record = payload as Record; + if (typeof record.type !== "string") return null; + // A native error may be pretty-printed JSON. SSE prefixes one data line; + // embedded physical newlines would otherwise truncate the JSON for readers. + if (record.type !== "response.done") return { + type: record.type, text: /[\r\n]/.test(text) ? JSON.stringify(record) : text, payload: record, + }; + + const response = record.response; + const status = response && typeof response === "object" && !Array.isArray(response) + ? (response as Record).status + : undefined; + const type = status === "completed" + ? "response.completed" + : status === "failed" + ? "response.failed" + : status === "incomplete" || status === "cancelled" + ? "response.incomplete" + : "response.failed"; + const normalizedRecord: Record = { ...record, type }; + if (type === "response.failed" && status !== "failed") { + normalizedRecord.response = response && typeof response === "object" && !Array.isArray(response) + ? { ...(response as Record), status: "failed" } + : { status: "failed" }; + } + return { type, text: JSON.stringify(normalizedRecord), payload: normalizedRecord }; +} + +/** + * The close code is the only thing that separates "the backend refused this + * payload" from "the network dropped", and both used to reach the caller as the + * same bare 502. Naming the oversized case here puts that distinction in the + * message the client receives. + * + * It does NOT reach the request log as a typed code. The eager relay turns any + * stream error into a generic `upstream_reset` synthetic terminal + * (`relay.ts`, `relay-eager.ts`) without feeding that frame back through the + * inspector, so `/api/logs` keeps neither this message nor a specific code — + * only `streamAborted`. Machine-readable typing would mean changing the error + * taxonomy, which is deliberately out of scope for this transport fix. + */ +export function closedBeforeTerminalMessage(event: unknown): string { + const detail = event as { code?: unknown; reason?: unknown } | null | undefined; + const code = typeof detail?.code === "number" ? detail.code : null; + const reason = typeof detail?.reason === "string" ? detail.reason.trim() : ""; + if (code === null) return CLOSED_BEFORE_TERMINAL; + const suffix = reason ? ` ${code} ${reason}` : ` ${code}`; + if (code === WS_CLOSE_MESSAGE_TOO_BIG) { + return `codex websocket rejected the request frame as too large (close${suffix});` + + ` requests at or above ${MAX_CODEX_WS_CREATE_FRAME_BYTES} bytes must use the HTTP SSE transport`; + } + return `${CLOSED_BEFORE_TERMINAL} (close${suffix})`; +} + +/** + * True when the `response.create` frame is at or above the backend's inbound + * message ceiling, so this turn must take the HTTP SSE path instead. + * + * Sizing a 16 MiB string should not cost a 16 MiB copy. UTF-8 never encodes + * below one byte per UTF-16 code unit and never above three, so both tails are + * settled from the string length alone; only the narrow band between them pays + * for a real byte count, and `Buffer.byteLength` measures without allocating. + */ +export function codexWsCreateFrameExceedsLimit( + frameText: string, + limitBytes: number = CODEX_WS_CREATE_FRAME_LIMIT_BYTES, +): boolean { + if (frameText.length >= limitBytes) return true; + if (frameText.length * 3 < limitBytes) return false; + return Buffer.byteLength(frameText, "utf8") >= limitBytes; +} diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index da1ca0a3d94..e9773d02a34 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -12,10 +12,17 @@ // returned event frames as an SSE byte stream, so every downstream consumer // (passthrough relay, adapter parsers, usage sniffing) is unchanged. -import { MAX_CLIENT_SSE_FRAME_BYTES } from "../sse-frame-buffer"; import { compareBunVersions } from "../../lib/bun-stream-caps"; -import { CodexWsMetadata, type CodexWsQuotaObserver } from "./codex-ws-metadata"; +import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { CODEX_RESPONSES_HTTP_URL, CODEX_RESPONSES_WS_URL, prepareCodexHttpInit, prepareCodexWsRequest } from "./codex-ws-request"; +import { codexWsExchange } from "./codex-ws-exchange"; +import { CodexWsSession } from "./codex-ws-session"; +import { codexWsPool, codexWsReuseIdentity } from "./codex-ws-pool"; +import { codexWsCreateFrameExceedsLimit } from "./codex-ws-wire"; +export { CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, MAX_CODEX_WS_QUEUE_BYTES, + MAX_CODEX_WS_CREATE_FRAME_BYTES, CODEX_WS_CREATE_FRAME_LIMIT_BYTES, codexWsCreateFrameExceedsLimit, + isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./codex-ws-wire"; +export const MIN_BOUNDED_CODEX_WS_BUN_VERSION = "1.4.0"; /** * Dial URL for a request URL. The canonical ChatGPT backend keeps its constant; @@ -46,37 +53,6 @@ function isResponsesWebsocketEligibleUrl(url: string): boolean { return parsed.protocol === "https:" && parsed.pathname.endsWith("/responses"); } -// If the 101 never arrives (network black hole), give SSE a chance well before -// the caller's connect timeout (default 200s) would fire. -const UPGRADE_DEADLINE_MS = 10_000; -export const CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS = 30_000; -// Keep the push-based WS transport inside the same memory envelope as the -// bounded SSE relays that consume this response. Unlike fetch response bodies, -// a WebSocket cannot be paused when a ReadableStream applies backpressure, so -// an upstream that outruns the consumer must be disconnected. -export const MAX_CODEX_WS_FRAME_BYTES = MAX_CLIENT_SSE_FRAME_BYTES; -export const MAX_CODEX_WS_QUEUE_BYTES = 8 * 1024 * 1024; -export const MIN_BOUNDED_CODEX_WS_BUN_VERSION = "1.4.0"; -// The backend drops any inbound message of 16 MiB or more: it closes the socket -// (1009) without a Responses terminal event, which reaches clients as a bare -// 502 upstream_server_error. Measured against the live endpoint 2026-08-23: -// 16,777,000 B completed, 16,777,300 B closed in ~1s, every time. The same -// request body succeeds over HTTP SSE, so the ceiling belongs to this transport -// alone (see #2426). A full-replay thread reaches it with ~11 pasted -// screenshots, and then never recovers, because each retry resends the frame. -export const MAX_CODEX_WS_CREATE_FRAME_BYTES = 16 * 1024 * 1024; -// Bun frames the payload it is handed, so the send-side budget is the JSON text -// itself, and nothing is appended between the check and the send. The margin is -// a conservative cushion, not a computed requirement: it covers RFC 6455 frame -// overhead in case the backend counts it (14 bytes at this payload size — an -// 8-byte extended length plus a 4-byte client mask, leaving ~65.5 KiB spare), -// and it leaves room for a future caller that appends to the frame. -const CODEX_WS_CREATE_FRAME_MARGIN_BYTES = 64 * 1024; -export const CODEX_WS_CREATE_FRAME_LIMIT_BYTES = - MAX_CODEX_WS_CREATE_FRAME_BYTES - CODEX_WS_CREATE_FRAME_MARGIN_BYTES; -/** Close code the backend uses for an oversized message (RFC 6455 "message too big"). */ -const WS_CLOSE_MESSAGE_TOO_BIG = 1009; - export type BunRuntimeIdentity = { version: string; versionWithSha: string; @@ -84,19 +60,6 @@ export type BunRuntimeIdentity = { export type BunRuntimeGateInput = string | BunRuntimeIdentity; -const codexWsUpstreamResponses = new WeakSet(); -const quotaObservedResponses = new WeakSet(); - -/** Quota arrived directly at its captured account; do not replay old HTTP prelude headers. */ -export function isCodexWsQuotaObservedResponse(response: Response): boolean { - return quotaObservedResponses.has(response); -} - -/** True only for a successful Codex WebSocket upgrade, never an HTTP fallback. */ -export function isCodexWsUpstreamResponse(response: Response): boolean { - return codexWsUpstreamResponses.has(response); -} - export function currentBunRuntimeIdentity(): BunRuntimeIdentity { return { version: Bun.version, @@ -158,97 +121,6 @@ export function shouldUseCodexWsUpstream( } } -const CLOSED_BEFORE_TERMINAL = "codex websocket closed before a Responses terminal event"; - -type ResponsesWsRelayEvent = { - type: string; - text: string; - payload: Record; -}; - -/** - * Responses WebSocket uses `response.done` as its terminal event, while the - * SSE Responses surface uses status-specific terminal events. Normalize the - * WS-only discriminator before relaying so the existing SSE consumers can - * settle the turn and the socket close cannot be mistaken for a drop. Unknown - * or missing status values fail closed instead of being reported as success. - */ -function normalizeResponsesWsRelayEvent(text: string): ResponsesWsRelayEvent | null { - let payload: unknown; - try { - payload = JSON.parse(text); - } catch { - return null; - } - if (!payload || typeof payload !== "object" || Array.isArray(payload)) return null; - const record = payload as Record; - if (typeof record.type !== "string") return null; - if (record.type !== "response.done") return { type: record.type, text, payload: record }; - - const response = record.response; - const status = response && typeof response === "object" && !Array.isArray(response) - ? (response as Record).status - : undefined; - const type = status === "completed" - ? "response.completed" - : status === "failed" - ? "response.failed" - : status === "incomplete" || status === "cancelled" - ? "response.incomplete" - : "response.failed"; - const normalizedRecord: Record = { ...record, type }; - if (type === "response.failed" && status !== "failed") { - normalizedRecord.response = response && typeof response === "object" && !Array.isArray(response) - ? { ...(response as Record), status: "failed" } - : { status: "failed" }; - } - return { type, text: JSON.stringify(normalizedRecord), payload: normalizedRecord }; -} - -/** - * The close code is the only thing that separates "the backend refused this - * payload" from "the network dropped", and both used to reach the caller as the - * same bare 502. Naming the oversized case here puts that distinction in the - * message the client receives. - * - * It does NOT reach the request log as a typed code. The eager relay turns any - * stream error into a generic `upstream_reset` synthetic terminal - * (`relay.ts`, `relay-eager.ts`) without feeding that frame back through the - * inspector, so `/api/logs` keeps neither this message nor a specific code — - * only `streamAborted`. Machine-readable typing would mean changing the error - * taxonomy, which is deliberately out of scope for this transport fix. - */ -function closedBeforeTerminalMessage(event: unknown): string { - const detail = event as { code?: unknown; reason?: unknown } | null | undefined; - const code = typeof detail?.code === "number" ? detail.code : null; - const reason = typeof detail?.reason === "string" ? detail.reason.trim() : ""; - if (code === null) return CLOSED_BEFORE_TERMINAL; - const suffix = reason ? ` ${code} ${reason}` : ` ${code}`; - if (code === WS_CLOSE_MESSAGE_TOO_BIG) { - return `codex websocket rejected the request frame as too large (close${suffix});` - + ` requests at or above ${MAX_CODEX_WS_CREATE_FRAME_BYTES} bytes must use the HTTP SSE transport`; - } - return `${CLOSED_BEFORE_TERMINAL} (close${suffix})`; -} - -/** - * True when the `response.create` frame is at or above the backend's inbound - * message ceiling, so this turn must take the HTTP SSE path instead. - * - * Sizing a 16 MiB string should not cost a 16 MiB copy. UTF-8 never encodes - * below one byte per UTF-16 code unit and never above three, so both tails are - * settled from the string length alone; only the narrow band between them pays - * for a real byte count, and `Buffer.byteLength` measures without allocating. - */ -export function codexWsCreateFrameExceedsLimit( - frameText: string, - limitBytes: number = CODEX_WS_CREATE_FRAME_LIMIT_BYTES, -): boolean { - if (frameText.length >= limitBytes) return true; - if (frameText.length * 3 < limitBytes) return false; - return Buffer.byteLength(frameText, "utf8") >= limitBytes; -} - export function codexWsUpstreamFetch( url: string, init: RequestInit, @@ -290,225 +162,17 @@ export function codexWsUpstreamFetch( } catch (error) { return Promise.reject(error); } - return new Promise((resolve, reject) => { - let ws: WebSocket; - try { - // Bun accepts per-handshake headers; the DOM lib types only list protocol arrays. - ws = new WebSocket(wsUpstreamUrlFor(url), { headers } as unknown as string[]); - } catch { - resolve(sseFallback(url, init)); - return; + let session: CodexWsSession; + try { + const identity = codexWsReuseIdentity(url, headers, frameText); + session = (identity ? codexWsPool.acquire(identity, wsUpstreamUrlFor(url), headers) : null) + ?? new CodexWsSession(wsUpstreamUrlFor(url), headers); + if (!session.busy && !session.reserve()) { + session.dispose(); + return sseFallback(url, init); } - - let opened = false; - let settledPreOpen = false; - let sent = false; - let received = false; - let responseCommitted = false; - let terminal = false; - let controller: ReadableStreamDefaultController | null = null; - const encoder = new TextEncoder(); - const metadata = url === CODEX_RESPONSES_HTTP_URL ? new CodexWsMetadata(onQuota) : null; - let preludeTimer: ReturnType | undefined; - const stream = new ReadableStream({ - start(c) { controller = c; }, - cancel() { - terminal = true; - cleanup(); - try { ws.close(); } catch { /* already closing */ } - }, - }, new ByteLengthQueuingStrategy({ highWaterMark: MAX_CODEX_WS_QUEUE_BYTES })); - - const cleanup = () => { - clearTimeout(upgradeTimer); - clearTimeout(preludeTimer); - signal?.removeEventListener("abort", onAbort); - metadata?.finish(); - }; - - const commitResponse = () => { - if (responseCommitted) return; - responseCommitted = true; - clearTimeout(preludeTimer); - const responseHeaders = metadata?.snapshot() ?? new Headers(); - responseHeaders.set("content-type", "text/event-stream; charset=utf-8"); - const response = new Response(stream, { status: 200, headers: responseHeaders }); - metadata?.commit(); - codexWsUpstreamResponses.add(response); - if (metadata && onQuota) quotaObservedResponses.add(response); - resolve(response); - }; - - const failStream = (error: unknown) => { - if (terminal) return; - terminal = true; - // A frame may already be executing upstream. Settle as a body failure, - // never a fetch rejection/5xx that the pre-stream wrapper could resend. - if (sent) commitResponse(); - cleanup(); - try { controller?.error(typeof error === "string" ? new Error(error) : error); } catch { /* stream already done */ } - try { ws.close(); } catch { /* already closing */ } - }; - - const upgradeTimer = setTimeout(() => { - if (opened || settledPreOpen) return; - settledPreOpen = true; - cleanup(); - try { ws.close(); } catch { /* already closing */ } - resolve(sseFallback(url, init)); - }, UPGRADE_DEADLINE_MS); - - const onAbort = () => { - if (!sent) { - if (settledPreOpen) return; - // Settle BEFORE close(): the close handler treats a pre-open close as - // an upgrade rejection and would dial the SSE fallback for a request - // the caller just cancelled. - settledPreOpen = true; - cleanup(); - try { ws.close(); } catch { /* already closing */ } - reject(signal?.reason ?? new DOMException("The operation was aborted.", "AbortError")); - return; - } - failStream(signal?.reason ?? new DOMException("The operation was aborted.", "AbortError")); - }; - signal?.addEventListener("abort", onAbort, { once: true }); - - const onOpen = () => { - if (settledPreOpen) return; - clearTimeout(upgradeTimer); - opened = true; - try { - beforeDispatch?.(new Headers(headers)); - } catch (error) { - // Settle and detach before close: a synchronous close event must not resend over SSE. - settledPreOpen = true; - terminal = true; - cleanup(); - ws.removeEventListener("open", onOpen); - ws.removeEventListener("message", onMessage); - ws.removeEventListener("close", onClose); - ws.removeEventListener("error", onError); - try { ws.close(); } catch { /* already closing */ } - reject(error); - return; - } - sent = true; - try { - ws.send(frameText); - } catch { - if (received || responseCommitted) { - failStream("codex websocket send failed after response activity"); - return; - } - // send() throwing means the frame never left, so no upstream turn - // started and the SSE resend cannot double-generate. Falling back - // (instead of erroring a synthetic 200 body) keeps the pre-stream - // HTTP error/refresh/failover machinery in charge. - settledPreOpen = true; - sent = false; - cleanup(); - try { ws.close(); } catch { /* already closing */ } - resolve(sseFallback(url, init)); - return; - } - if (!metadata) commitResponse(); - else if (!responseCommitted && !terminal) { - preludeTimer = setTimeout(() => failStream("codex websocket response prelude timed out"), CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS); - } - }; - - const onMessage = (event: MessageEvent) => { - if (!controller || terminal) return; - received = true; - const text = typeof event.data === "string" ? event.data : ""; - if (!text) return; - // UTF-8 byte length is always at least the JS string length. Reject this - // cheap lower bound before parsing so an obviously oversized frame does - // not create another large object graph. - if (text.length > MAX_CODEX_WS_FRAME_BYTES) { - failStream("codex websocket frame exceeds the response size limit"); - return; - } - const rawEncodedText = encoder.encode(text); - if (rawEncodedText.byteLength > MAX_CODEX_WS_FRAME_BYTES) { - failStream("codex websocket frame exceeds the response size limit"); - return; - } - const normalized = normalizeResponsesWsRelayEvent(text); - if (!normalized) return; - const { type } = normalized; - let relayText = normalized.text; - let controlFrame = false; - if (metadata) { - try { - const sanitized = metadata.consume(normalized.payload, rawEncodedText.byteLength); - if (sanitized !== null) { - relayText = sanitized; - controlFrame = true; - } - } catch (error) { - failStream(error); - return; - } - } - const encodedText = relayText === text ? rawEncodedText : encoder.encode(relayText); - if (encodedText.byteLength > MAX_CODEX_WS_FRAME_BYTES) { - failStream("codex websocket frame exceeds the response size limit"); - return; - } - if (!controlFrame && !type.startsWith("response.") && type !== "error") return; - if (!controlFrame) commitResponse(); - const prefix = encoder.encode(`event: ${type}\ndata: `); - const suffix = encoder.encode("\n\n"); - const frameBytes = prefix.byteLength + encodedText.byteLength + suffix.byteLength; - if (frameBytes > MAX_CLIENT_SSE_FRAME_BYTES) { - failStream("codex websocket frame exceeds the response size limit"); - return; - } - const availableBytes = controller.desiredSize ?? 0; - if (frameBytes > availableBytes) { - failStream("codex websocket response exceeded the buffered queue limit"); - return; - } - const sseFrame = new Uint8Array(frameBytes); - sseFrame.set(prefix); - sseFrame.set(encodedText, prefix.byteLength); - sseFrame.set(suffix, prefix.byteLength + encodedText.byteLength); - try { - controller.enqueue(sseFrame); - } catch { - failStream("codex websocket response stream closed while enqueueing"); - return; - } - if (type === "response.completed" || type === "response.failed" || type === "response.incomplete" || type === "error") { - terminal = true; - cleanup(); - try { controller.close(); } catch { /* already closed */ } - try { ws.close(); } catch { /* already closing */ } - } - }; - - const onClose = (event: unknown) => { - cleanup(); - if (!opened) { - if (settledPreOpen) return; - settledPreOpen = true; - // Upgrade rejected (401/403/429/5xx). Retry over plain SSE so the real - // HTTP status reaches the existing refresh/rotation handlers. No turn - // started upstream, so the resend cannot double-generate. - resolve(sseFallback(url, init)); - return; - } - if (sent && !terminal) failStream(closedBeforeTerminalMessage(event)); - }; - - const onError = () => { - /* Bun always follows error with close; the close handler settles. */ - }; - ws.addEventListener("open", onOpen); - ws.addEventListener("message", onMessage); - ws.addEventListener("close", onClose); - ws.addEventListener("error", onError); - }); + } catch { + return sseFallback(url, init); + } + return codexWsExchange({ session, url, init, prepared, sseFallback, onQuota, beforeDispatch }); } diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index 200ffb58413..c545d41d8c1 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -414,6 +414,22 @@ so HTTP fallback cannot duplicate that inference. A standalone no-response exchange has a 30-second prelude deadline in addition to the upgrade deadline. These are transport-fidelity guarantees, not a provider-billing guarantee. +Eligible complete-input creates can retain a canonical upstream socket within +one selected account, credential, thread and turn. Model/tier and immutable +handshake headers must also match. Turn-state and turn-metadata headers are +projected into their same-name per-frame metadata slots; explicit body values win. +The pool retains at most 32 sockets, expires idle sockets after 30 seconds, and +retires a socket after five minutes or 32 successful exchanges (after active work +finishes). Cancellation, errors, idle unsolicited frames and shutdown dispose it. +A busy key uses a separate one-shot connection rather than interleaving requests. + +This is connection reuse, not native incremental-input synthesis: complete HTTP +inputs are never trimmed and no previous response id is invented. Explicit +continuation IDs, named lanes, warmup and background requests remain outside this +pool. A fresh credential-dispatch guard runs before every warm send. Per-exchange +listeners, response/item correlation and metadata ownership detach before release. +No pool timer or shutdown registration exists before eligible traffic activates it. + Translated response request-log tracking and the heartbeat relay also reuse `createSseInspector`. This keeps every client-facing SSE observation path on the same byte-bounded, discard-and-resynchronize frame policy and ensures the diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 5276f08f461..114c699eaf8 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1085,6 +1085,7 @@ "winsw-stop-hardening.test.ts": "windows", "winsw.test.ts": "service", "ws-endpoint.test.ts": "responses", + "ws-upstream-reuse.test.ts": "responses", "ws-upstream.test.ts": "responses", "xai-client.test.ts": "images", "xai-oauth-retry.test.ts": "providers/xai", diff --git a/tests/responses/responses-account-label.test.ts b/tests/responses/responses-account-label.test.ts index b617905426f..b7be5c8d7fb 100644 --- a/tests/responses/responses-account-label.test.ts +++ b/tests/responses/responses-account-label.test.ts @@ -143,6 +143,9 @@ describe("Responses account usage attribution", () => { addEventListener(type: string, listener: (event: unknown) => void) { this.listeners.set(type, [...(this.listeners.get(type) ?? []), listener]); } + removeEventListener(type: string, listener: (event: unknown) => void) { + this.listeners.set(type, (this.listeners.get(type) ?? []).filter(value => value !== listener)); + } emit(type: string, event: unknown) { for (const listener of this.listeners.get(type) ?? []) listener(event); } diff --git a/tests/responses/ws-upstream-reuse.test.ts b/tests/responses/ws-upstream-reuse.test.ts new file mode 100644 index 00000000000..fd0a8fb5a15 --- /dev/null +++ b/tests/responses/ws-upstream-reuse.test.ts @@ -0,0 +1,284 @@ +import { afterEach, beforeEach, expect, test } from "bun:test"; +import { codexWsUpstreamFetch } from "../../src/server/responses/ws-upstream"; +import { runOptionalShutdownHooks } from "../../src/lib/optional-shutdown-hooks"; +import { CodexWsPool, codexWsPool } from "../../src/server/responses/codex-ws-pool"; +import { prepareCodexWsRequest } from "../../src/server/responses/codex-ws-request"; + +const URL = "https://chatgpt.com/backend-api/codex/responses"; +const realWebSocket = globalThis.WebSocket; +let sequence = 0; + +class Socket extends EventTarget { + static all: Socket[] = []; + static onSend: (socket: Socket, frame: Record) => void = (socket) => socket.complete(); + readyState = 0; + frames: Record[] = []; + constructor(readonly url: string) { + super(); + Socket.all.push(this); + queueMicrotask(() => { if (this.readyState === 0) { this.readyState = 1; this.dispatchEvent(new Event("open")); } }); + } + send(text: string) { + const frame = JSON.parse(text); + this.frames.push(frame); + Socket.onSend(this, frame); + } + emit(payload: Record) { + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify(payload) })); + } + complete() { + const id = `response-${++sequence}`; + queueMicrotask(() => { + this.emit({ type: "response.created", response: { id } }); + this.emit({ type: "response.completed", response: { id, status: "completed", output: [] } }); + }); + } + close() { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } + ref() {} + unref() {} +} + +function init(input = "first", signal?: AbortSignal): RequestInit { + return { method: "POST", signal, headers: { + authorization: "Bearer fixture-token", "chatgpt-account-id": "fixture-account", "thread-id": "fixture-thread", + }, body: JSON.stringify({ model: "fixture-model", stream: true, input, + client_metadata: { thread_id: "fixture-thread", turn_id: "fixture-turn" } }) }; +} + +const fallback = (async () => { throw new Error("unexpected HTTP fallback"); }) as typeof fetch; +const request = (options = init(), guard?: (headers: Headers) => void) => + codexWsUpstreamFetch(URL, options, fallback, "1.4.0", undefined, guard); +const drain = async (options = init()) => (await request(options)).text(); +function bodyWith(fields: Record) { + const options = init(); + options.body = JSON.stringify({ ...JSON.parse(options.body as string), ...fields }); + return options; +} +beforeEach(() => { globalThis.WebSocket = Socket as unknown as typeof WebSocket; }); + +afterEach(() => { + runOptionalShutdownHooks(); + for (const socket of Socket.all) socket.close(); + Socket.all = []; + Socket.onSend = socket => socket.complete(); + sequence = 0; + globalThis.WebSocket = realWebSocket; +}); + +test("same account/thread/turn reuses one socket without trimming either HTTP input", async () => { + globalThis.WebSocket = Socket as unknown as typeof WebSocket; + await (await codexWsUpstreamFetch(URL, init("first full input"), fallback, "1.4.0")).text(); + await (await codexWsUpstreamFetch(URL, init("second full input"), fallback, "1.4.0")).text(); + expect(Socket.all).toHaveLength(1); + expect(Socket.all[0]!.frames.map(frame => frame.input)).toEqual(["first full input", "second full input"]); + expect(Socket.all[0]!.frames.every(frame => !Object.hasOwn(frame, "previous_response_id"))).toBe(true); +}); + +test.each(["authorization", "chatgpt-account-id", "originator", "x-client-request-id", "x-custom-policy"])( + "changed selected %s cannot reuse an immutable handshake", async name => { + await drain(); + const options = init(); + const headers = new Headers(options.headers); + headers.set(name, name === "authorization" ? "Bearer rotated-token" : "different"); + await drain({ ...options, headers }); + expect(Socket.all).toHaveLength(2); + }); + +test.each([ + { model: "another-model" }, { service_tier: "priority" }, + { client_metadata: { thread_id: "other-thread", turn_id: "fixture-turn" } }, + { client_metadata: { thread_id: "fixture-thread", turn_id: "other-turn" } }, +])("model, tier or native scope changes redial: %j", async fields => { + await drain(); + await drain(bodyWith(fields)); + expect(Socket.all).toHaveLength(2); +}); + +test.each([ + { client_metadata: {} }, { client_metadata: { session_id: "shared", turn_id: "turn" }, }, + { client_metadata: { thread_id: "fixture-thread", turn_id: "" } }, + { previous_response_id: "server-owned-id" }, { stream_id: "main" }, + { generate: false }, { background: true }, +])("ineligible requests stay one-shot: %j", async fields => { + const options = bodyWith(fields); + // session-only fixture must not accidentally inherit the explicit header thread. + if (Object.hasOwn((fields.client_metadata ?? {}) as object, "session_id")) { + const headers = new Headers(options.headers); headers.delete("thread-id"); options.headers = headers; + } + await drain(options); await drain(options); + expect(Socket.all).toHaveLength(2); + expect(codexWsPool.snapshot()).toEqual({ size: 0, active: 0, timer: false }); +}); + +test("mutable turn headers are projected per frame; explicit body values win", async () => { + for (const state of ["state-a", "state-b"]) { + const options = init(); const headers = new Headers(options.headers); + headers.set("x-codex-turn-state", state); + headers.set("x-codex-turn-metadata", JSON.stringify({ turn: state })); + await drain({ ...options, headers }); + } + expect(Socket.all).toHaveLength(1); + expect(Socket.all[0]!.frames.map(frame => (frame.client_metadata as Record)["x-codex-turn-state"])) + .toEqual(["state-a", "state-b"]); + expect((Socket.all[0]!.frames[1]!.client_metadata as Record)["x-codex-turn-metadata"]) + .toBe('{"turn":"state-b"}'); + const options = bodyWith({ client_metadata: { "x-codex-turn-state": "body-state" } }); + const headers = new Headers(options.headers); headers.set("x-codex-turn-state", "header-state"); + const prepared = prepareCodexWsRequest(URL, { ...options, headers })!; + expect(JSON.parse(prepared.frameText).client_metadata["x-codex-turn-state"]).toBe("body-state"); +}); + +test("fresh warm dispatch guard refusal never sends or falls back", async () => { + await drain(); + let checks = 0; + await expect(request(init(), () => { if (++checks === 2) throw new Error("revoked"); })).rejects.toThrow("revoked"); + expect(checks).toBe(2); + expect(Socket.all).toHaveLength(1); + expect(Socket.all[0]!.frames).toHaveLength(1); + expect(codexWsPool.snapshot().size).toBe(0); +}); + +test("busy identity gets an independent one-shot; old abort cannot kill successor", async () => { + const old = new AbortController(); + await drain(init("A", old.signal)); + Socket.onSend = socket => queueMicrotask(() => socket.emit({ type: "response.created", response: { id: `active-${Socket.all.indexOf(socket)}` } })); + const b = await request(init("B")); + const c = await request(init("C")); + expect(Socket.all).toHaveLength(2); + expect(Socket.all[0]!.frames.map(frame => frame.input)).toEqual(["A", "B"]); + old.abort(); + expect(Socket.all[0]!.readyState).toBe(1); + for (const [index, socket] of Socket.all.entries()) socket.emit({ type: "response.completed", response: { id: `active-${index}`, status: "completed" } }); + await b.text(); await c.text(); + expect(Socket.all[0]!.readyState).toBe(1); + expect(Socket.all[1]!.readyState).toBe(3); +}); + +test("overlapping A to changed-header B to A keeps retired busy sockets tracked until release", async () => { + Socket.onSend = () => {}; + const changed = init("B"); + const headers = new Headers(changed.headers); + headers.set("x-custom-policy", "B"); + const pending = [request(init("A")), request({ ...changed, headers }), request(init("A-again"))]; + await Promise.resolve(); + expect(Socket.all).toHaveLength(3); + expect(codexWsPool.snapshot()).toEqual({ size: 2, active: 2, timer: false }); + expect(Socket.all.map(socket => socket.frames.map(frame => frame.input))) + .toEqual([["A"], ["B"], ["A-again"]]); + Socket.all[0]!.complete(); + await (await pending[0]!).text(); + expect(Socket.all[0]!.readyState).toBe(3); + expect(codexWsPool.snapshot()).toEqual({ size: 1, active: 1, timer: false }); + Socket.all[1]!.complete(); Socket.all[2]!.complete(); + await Promise.all(pending.slice(1).map(async result => (await result).text())); + expect(Socket.all.every(socket => socket.readyState === 3)).toBe(true); + expect(codexWsPool.snapshot()).toEqual({ size: 0, active: 0, timer: false }); +}); + +test.each(["abort", "error", "close", "shutdown", "stale-item", "stale-response", "named-lane"])( + "warm %s fails its body without a resend", async reason => { + await drain(); + Socket.onSend = socket => queueMicrotask(() => socket.emit({ type: "response.created", response: { id: "new-response" } })); + const abort = new AbortController(); + const response = await request(init("B", abort.signal)); + const socket = Socket.all[0]!; + if (reason === "abort") abort.abort(); + if (reason === "error") socket.dispatchEvent(new Event("error")); + if (reason === "close") socket.close(); + if (reason === "shutdown") runOptionalShutdownHooks(); + if (reason === "stale-item") socket.emit({ type: "response.output_text.delta", item_id: "old-item", delta: "MUST NOT RELAY" }); + if (reason === "stale-response") socket.emit({ type: "response.completed", response: { id: "response-1", status: "completed" } }); + if (reason === "named-lane") socket.emit({ type: "response.output_text.delta", stream_id: "other", delta: "MUST NOT RELAY" }); + await expect(response.text()).rejects.toThrow(); + expect(Socket.all).toHaveLength(1); + expect(socket.frames).toHaveLength(2); + expect(codexWsPool.snapshot()).toEqual({ size: 0, active: 0, timer: false }); + }); + +test("idle unsolicited data retires the socket before another request", async () => { + await drain(); + Socket.all[0]!.emit({ type: "response.created", response: { id: "unsolicited" } }); + await drain(); + expect(Socket.all).toHaveLength(2); +}); + +test("uncorrelatable legacy response remains usable but never retained", async () => { + Socket.onSend = socket => queueMicrotask(() => socket.emit({ type: "response.completed", response: { status: "completed" } })); + expect(await drain()).toContain("response.completed"); + expect(await drain()).toContain("response.completed"); + expect(Socket.all).toHaveLength(2); + expect(codexWsPool.snapshot().timer).toBe(false); +}); + +test("bounded pool expires idle state, preserves active work, and drains on shutdown", async () => { + let now = 0; + const pool = new CodexWsPool({ now: () => now, idleMs: 30_000, maxAgeMs: 300_000, maxSessions: 2 }); + try { + expect(pool.snapshot()).toEqual({ size: 0, active: 0, timer: false }); + const a = pool.acquire({ key: "a", scope: "a" }, "wss://fixture", {})!; + await Promise.resolve(); a.release("a-response"); + now = 29_999; pool.sweep(); expect(a.closed).toBe(false); + now = 30_000; pool.sweep(); expect(a.closed).toBe(true); + const b = pool.acquire({ key: "b", scope: "b" }, "wss://fixture", {})!; + await Promise.resolve(); + now = 330_000; pool.sweep(); expect(b.closed).toBe(false); + b.release("b-response"); expect(b.closed).toBe(true); + const c = pool.acquire({ key: "c", scope: "c" }, "wss://fixture", {})!; + const d = pool.acquire({ key: "d", scope: "d" }, "wss://fixture", {})!; + expect(pool.acquire({ key: "e", scope: "e" }, "wss://fixture", {})).toBeNull(); + await Promise.resolve(); c.release("c-response"); + const e = pool.acquire({ key: "e", scope: "e" }, "wss://fixture", {})!; + expect(c.closed).toBe(true); expect(d.closed).toBe(false); + pool.dispose(); expect(d.closed).toBe(true); expect(e.closed).toBe(true); + expect(pool.snapshot()).toEqual({ size: 0, active: 0, timer: false }); + } finally { pool.dispose(); } +}); + +test("shutdown before open rejects as cancellation, not fallback", async () => { + const response = request(); + runOptionalShutdownHooks(); + await expect(response).rejects.toMatchObject({ name: "AbortError" }); + expect(Socket.all[0]!.frames).toHaveLength(0); +}); + +test("quota prelude and callbacks belong to each warm exchange, not its predecessor", async () => { + let turn = 0; + Socket.onSend = socket => queueMicrotask(() => { + const id = `quota-${++turn}`; + socket.emit({ type: "codex.rate_limits", rate_limits: { primary: { used_percent: turn, window_minutes: 300 } } }); + socket.emit({ type: "response.created", response: { id } }); + socket.emit({ type: "response.completed", response: { id, status: "completed" } }); + }); + const observed: string[][] = [[], []]; + for (let index = 0; index < 2; index++) { + const response = await codexWsUpstreamFetch(URL, init(), fallback, "1.4.0", headers => { + observed[index]!.push(headers.get("x-codex-primary-used-percent")!); + }); + expect(response.headers.get("x-codex-primary-used-percent")).toBe(String(index + 1)); + await response.text(); + } + expect(Socket.all).toHaveLength(1); + expect(observed).toEqual([["1"], ["2"]]); +}); + +test("retirement bounds remembered response IDs and keeps all full requests intact", async () => { + for (let index = 0; index < 33; index++) await drain(init(`full-${index}`)); + expect(Socket.all).toHaveLength(2); + expect(Socket.all[0]!.frames).toHaveLength(32); + expect(Socket.all[0]!.readyState).toBe(3); + expect(Socket.all[1]!.frames[0]!.input).toBe("full-32"); +}); + +test("a Lite mode change retires the old handshake", async () => { + for (const lite of ["true", "false"]) { + const options = init(); const headers = new Headers(options.headers); + headers.set("x-openai-internal-codex-responses-lite", lite); + await drain({ ...options, headers }); + } + expect(Socket.all).toHaveLength(2); + expect(Socket.all[0]!.readyState).toBe(3); +}); diff --git a/tests/responses/ws-upstream.test.ts b/tests/responses/ws-upstream.test.ts index 0c9e62e6984..fd09513070f 100644 --- a/tests/responses/ws-upstream.test.ts +++ b/tests/responses/ws-upstream.test.ts @@ -188,6 +188,10 @@ class FakeWebSocket { for (const listener of this.listeners.get(type) ?? []) listener(event); } + removeEventListener(type: string, listener: Listener) { + this.listeners.set(type, (this.listeners.get(type) ?? []).filter(value => value !== listener)); + } + send(data: string) { this.sent.push(data); } @@ -581,6 +585,24 @@ describe("codexWsUpstreamFetch", () => { expect(FakeWebSocket.instances[0].closed).toBe(true); }); + test.each(["error", "response.completed"])("multiline upstream %s JSON remains one valid SSE data value", async type => { + const payload = type === "error" + ? { type, status: 400, error: { type: "invalid_request_error", message: "fixture refusal" } } + : { type, response: { id: "pretty-response", status: "completed", output: [] } }; + installFake(ws => { + ws.emit("open", {}); + ws.emit("message", { data: JSON.stringify(payload, null, 2) }); + }); + const response = await codexWsUpstreamFetch(CODEX_URL, streamingInit(), (async () => { + throw new Error("a sent multiline response cannot fall back"); + }) as typeof fetch); + const text = await response.text(); + const data = text.split("\n").filter(line => line.startsWith("data: ")); + expect(data).toHaveLength(1); + expect(JSON.parse(data[0]!.slice(6))).toEqual(payload); + expect(FakeWebSocket.instances[0]!.closed).toBe(true); + }); + test("normalizes the Responses WebSocket response.done terminal to SSE", async () => { installFake(ws => { ws.emit("open", {});