import type * as NodeOs from "node:os"; import type * as NodeZlib from "node:zlib"; import type { Tool as OpenAITool, ResponseCreateParamsStreaming, ResponseInput, ResponseStreamEvent, } from "openai/resources/responses/responses.js"; type ProcessWithOsBuiltinModule = typeof process & { getBuiltinModule?: (id: "node:os") => typeof NodeOs; }; function loadNodeOs(): typeof NodeOs | null { if (typeof process === "undefined" || !(process.versions?.node || process.versions?.bun)) { return null; } return (process as ProcessWithOsBuiltinModule).getBuiltinModule?.("node:os") ?? null; } // NEVER convert to top-level runtime imports - breaks browser/Vite builds const _os: typeof NodeOs | null = loadNodeOs(); import { clampThinkingLevel } from "../models.ts"; import { registerSessionResourceCleanup } from "../session-resources.ts"; import type { Api, AssistantMessage, Context, Model, ProviderEnv, ProviderHeaders, SimpleStreamOptions, StreamFunction, StreamOptions, Usage, } from "../types.ts"; import { combineAbortSignals } from "../utils/abort-signals.ts"; import { splitDeferredTools } from "../utils/deferred-tools.ts"; import { appendAssistantMessageDiagnostic, createAssistantMessageDiagnostic, formatThrownValue, } from "../utils/diagnostics.ts"; import { formatProviderError, normalizeProviderError } from "../utils/error-body.ts"; import { AssistantMessageEventStream } from "../utils/event-stream.ts"; import { headersToRecord } from "../utils/headers.ts"; import { resolveHttpProxyUrlForTarget } from "../utils/node-http-proxy.ts"; import { uuidv7 } from "../utils/uuid.ts"; import { clampOpenAIPromptCacheKey } from "./openai-prompt-cache.ts"; import { convertResponsesMessages, convertResponsesTools, processResponsesStream } from "./openai-responses-shared.ts"; import { buildBaseOptions } from "./simple-options.ts"; // ============================================================================ // Configuration // ============================================================================ const DEFAULT_CODEX_BASE_URL = "https://chatgpt.com/backend-api"; const JWT_CLAIM_PATH = "https://api.openai.com/auth" as const; const DEFAULT_MAX_RETRIES = 0; const BASE_DELAY_MS = 1000; const DEFAULT_MAX_RETRY_DELAY_MS = 60_000; const DEFAULT_WEBSOCKET_CONNECT_TIMEOUT_MS = 15_000; // The Codex backend accepts zstd-compressed request bodies on the SSE responses // endpoint (the same endpoint the official Codex client compresses against). const REQUEST_COMPRESSION_ZSTD_LEVEL = 3; const CODEX_TOOL_CALL_PROVIDERS = new Set(["openai", "openai-codex", "opencode"]); const WEBSOCKET_MESSAGE_TOO_BIG_CLOSE_CODE = 1009; const WEBSOCKET_CONNECTION_LIMIT_REACHED_CODE = "websocket_connection_limit_reached"; const PREVIOUS_RESPONSE_NOT_FOUND_CODE = "previous_response_not_found"; const CODEX_RESPONSE_STATUSES = new Set([ "completed", "incomplete", "failed", "cancelled", "queued", "in_progress", ]); // ============================================================================ // Types // ============================================================================ export interface OpenAICodexResponsesOptions extends StreamOptions { reasoningEffort?: "none" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max"; reasoningSummary?: "auto" | "concise" | "detailed" | "off" | "on" | null; serviceTier?: ResponseCreateParamsStreaming["service_tier"]; textVerbosity?: "low" | "medium" | "high"; toolChoice?: "auto" | "none" | "required"; } type CodexResponseStatus = "completed" | "incomplete" | "failed" | "cancelled" | "queued" | "in_progress"; interface RequestBody { model: string; store?: boolean; stream?: boolean; instructions?: string; previous_response_id?: string; input?: ResponseInput; tools?: OpenAITool[]; tool_choice?: OpenAICodexResponsesOptions["toolChoice"]; parallel_tool_calls?: boolean; temperature?: number; reasoning?: { effort?: string; summary?: string }; service_tier?: ResponseCreateParamsStreaming["service_tier"]; text?: { verbosity?: string }; include?: string[]; prompt_cache_key?: string; [key: string]: unknown; } // ============================================================================ // Retry Helpers // ============================================================================ function isTerminalRateLimitError(errorText: string): boolean { return /GoUsageLimitError|FreeUsageLimitError|Monthly usage limit reached|available balance|insufficient_quota|out of budget|quota exceeded|billing/i.test( errorText, ); } function isRetryableError(status: number, errorText: string): boolean { if (status === 429 && isTerminalRateLimitError(errorText)) { return false; } if (status === 429 || status === 500 || status === 502 || status === 503 || status === 504) { return true; } return /rate.?limit|overloaded|service.?unavailable|upstream.?connect|connection.?refused/i.test(errorText); } function getRetryAfterDelayMs(headers: Headers): number | undefined { const retryAfterMs = headers.get("retry-after-ms"); if (retryAfterMs !== null) { const millis = Number(retryAfterMs); if (Number.isFinite(millis)) { return Math.max(0, millis); } } const retryAfter = headers.get("retry-after"); if (!retryAfter) { return undefined; } const seconds = Number(retryAfter); if (Number.isFinite(seconds)) { return Math.max(0, seconds * 1000); } const date = Date.parse(retryAfter); if (!Number.isNaN(date)) { return Math.max(0, date - Date.now()); } return undefined; } function capRetryDelayMs(delayMs: number, options?: StreamOptions): number { const maxRetryDelayMs = options?.maxRetryDelayMs ?? DEFAULT_MAX_RETRY_DELAY_MS; return maxRetryDelayMs > 0 ? Math.min(delayMs, maxRetryDelayMs) : delayMs; } function sleep(ms: number, signal?: AbortSignal): Promise { return new Promise((resolve, reject) => { if (signal?.aborted) { reject(new Error("Request was aborted")); return; } const timeout = setTimeout(resolve, ms); signal?.addEventListener("abort", () => { clearTimeout(timeout); reject(new Error("Request was aborted")); }); }); } function normalizeTimeoutMs(value: number | undefined): number | undefined { if (value === undefined) return undefined; if (!Number.isFinite(value) || value < 0) { throw new Error(`Invalid timeoutMs: ${String(value)}`); } return Math.floor(value); } // ============================================================================ // Request Compression // ============================================================================ type ProcessWithBuiltinModule = typeof process & { getBuiltinModule?: (id: "node:zlib") => typeof NodeZlib; }; function loadNodeZlib(): typeof NodeZlib | null { if (typeof process === "undefined" || !(process.versions?.node || process.versions?.bun)) { return null; } return (process as ProcessWithBuiltinModule).getBuiltinModule?.("node:zlib") ?? null; } // Returns the zstd-compressed body bytes, or null when compression is // unavailable (browser/Vite builds). Callers fall back to sending the // uncompressed JSON when this returns null. function compressRequestBodyZstd(bodyJson: string): Uint8Array | null { const zlib = loadNodeZlib(); if (!zlib || typeof zlib.zstdCompressSync !== "function") { return null; } try { const compressed = zlib.zstdCompressSync(bodyJson, { params: { [zlib.constants.ZSTD_c_compressionLevel]: REQUEST_COMPRESSION_ZSTD_LEVEL }, }); return new Uint8Array(compressed.buffer, compressed.byteOffset, compressed.byteLength); } catch { return null; } } // ============================================================================ // Main Stream Function // ============================================================================ export const stream: StreamFunction<"openai-codex-responses", OpenAICodexResponsesOptions> = ( model: Model<"openai-codex-responses">, context: Context, options?: OpenAICodexResponsesOptions, ): AssistantMessageEventStream => { const stream = new AssistantMessageEventStream(); (async () => { const output: AssistantMessage = { role: "assistant", content: [], api: "openai-codex-responses" as Api, provider: model.provider, model: model.id, usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason: "stop", timestamp: Date.now(), }; try { const apiKey = options?.apiKey; if (!apiKey) { throw new Error(`No API key for provider: ${model.provider}`); } const accountId = extractAccountId(apiKey); let body = buildRequestBody(model, context, options); const nextBody = await options?.onPayload?.(body, model); if (nextBody !== undefined) { body = nextBody as RequestBody; } const codexSessionId = clampOpenAIPromptCacheKey(options?.sessionId); const websocketRequestId = codexSessionId || uuidv7(); const sseHeaders = buildSSEHeaders(model.headers, options?.headers, accountId, apiKey, codexSessionId); const websocketHeaders = buildWebSocketHeaders( model.headers, options?.headers, accountId, apiKey, websocketRequestId, ); const bodyJson = JSON.stringify(body); const httpTimeoutMs = normalizeTimeoutMs(options?.timeoutMs); const websocketConnectTimeoutMs = normalizeTimeoutMs(options?.websocketConnectTimeoutMs); const transport = options?.transport || "auto"; let startEmitted = false; const websocketDisabledForSession = transport !== "sse" && isWebSocketSseFallbackActive(options?.sessionId); if (websocketDisabledForSession) { recordWebSocketSseFallback(options?.sessionId); } if (transport !== "sse" && !websocketDisabledForSession) { let websocketStarted = false; let retriedWebSocketConnectionLimit = false; let retriedMissingWebSocketContinuation = false; while (true) { websocketStarted = false; try { await processWebSocketStream( resolveCodexWebSocketUrl(model.baseUrl), body, websocketHeaders, output, stream, model, () => { websocketStarted = true; if (!startEmitted) { startEmitted = true; stream.push({ type: "start", partial: output }); } }, httpTimeoutMs, websocketConnectTimeoutMs, options, ); if (options?.signal?.aborted) { throw new Error("Request was aborted"); } stream.push({ type: "done", reason: output.stopReason as "stop" | "length" | "toolUse", message: output, }); stream.end(); return; } catch (error) { const aborted = options?.signal?.aborted; const connectionLimitBeforeStart = !websocketStarted && isWebSocketConnectionLimitReachedError(error); const previousResponseNotFound = isPreviousResponseNotFoundError(error); if (!aborted && previousResponseNotFound && !retriedMissingWebSocketContinuation) { retriedMissingWebSocketContinuation = true; continue; } if (!aborted && connectionLimitBeforeStart && !retriedWebSocketConnectionLimit) { retriedWebSocketConnectionLimit = true; continue; } if (aborted || (isCodexNonTransportError(error) && !connectionLimitBeforeStart)) { throw error; } appendAssistantMessageDiagnostic( output, createAssistantMessageDiagnostic("provider_transport_failure", error, { configuredTransport: transport, fallbackTransport: websocketStarted ? undefined : "sse", eventsEmitted: websocketStarted, phase: websocketStarted ? "after_message_stream_start" : "before_message_stream_start", requestBytes: new TextEncoder().encode(bodyJson).byteLength, }), ); recordWebSocketFailure(options?.sessionId, error); if (websocketStarted) { throw error; } recordWebSocketSseFallback(options?.sessionId); break; } } } // Compress the request body once for the SSE path. The Codex backend // decodes Content-Encoding: zstd; the WebSocket transport above sends the // uncompressed JSON frame, matching the official Codex client. const compressedBody = compressRequestBodyZstd(bodyJson); if (compressedBody) { sseHeaders.set("content-encoding", "zstd"); } const sseBody: Uint8Array | string = compressedBody ?? bodyJson; // Fetch with retry logic for rate limits and transient errors let response: Response | undefined; let lastError: Error | undefined; const maxRetries = options?.maxRetries ?? DEFAULT_MAX_RETRIES; for (let attempt = 0; attempt <= maxRetries; attempt++) { if (options?.signal?.aborted) { throw new Error("Request was aborted"); } try { const headerTimeoutSignal = httpTimeoutMs !== undefined && httpTimeoutMs > 0 ? AbortSignal.timeout(httpTimeoutMs) : undefined; const combinedSignal = combineAbortSignals([options?.signal, headerTimeoutSignal]); try { response = await fetch(resolveCodexUrl(model.baseUrl), { method: "POST", headers: sseHeaders, body: sseBody, signal: combinedSignal.signal, }); } catch (error) { if (headerTimeoutSignal?.aborted && !options?.signal?.aborted) { throw new Error(`Codex SSE response headers timed out after ${httpTimeoutMs}ms`); } throw error; } finally { combinedSignal.cleanup(); } await options?.onResponse?.( { status: response.status, headers: headersToRecord(response.headers) }, model, ); if (response.ok) { break; } const errorText = await response.text(); if (attempt < maxRetries && isRetryableError(response.status, errorText)) { const retryAfterDelayMs = getRetryAfterDelayMs(response.headers); const delayMs = retryAfterDelayMs === undefined ? BASE_DELAY_MS * 2 ** attempt : response.status === 429 ? capRetryDelayMs(retryAfterDelayMs, options) : retryAfterDelayMs; await sleep(delayMs, options?.signal); continue; } // Parse error for friendly message on final attempt or non-retryable error const fakeResponse = new Response(errorText, { status: response.status, statusText: response.statusText, }); const info = await parseErrorResponse(fakeResponse); throw new Error(info.friendlyMessage || info.message); } catch (error) { if (error instanceof Error) { if (error.name === "AbortError" || error.message === "Request was aborted") { throw new Error("Request was aborted"); } } lastError = error instanceof Error ? error : new Error(String(error)); // Network errors are retryable if (attempt < maxRetries && !lastError.message.includes("usage limit")) { const delayMs = BASE_DELAY_MS * 2 ** attempt; await sleep(delayMs, options?.signal); continue; } throw lastError; } } if (!response?.ok) { throw lastError ?? new Error("Failed after retries"); } if (!response.body) { throw new Error("No response body"); } if (!startEmitted) { startEmitted = true; stream.push({ type: "start", partial: output }); } await processStream(response, output, stream, model, options); if (options?.signal?.aborted) { throw new Error("Request was aborted"); } stream.push({ type: "done", reason: output.stopReason as "stop" | "length" | "toolUse", message: output }); stream.end(); } catch (error) { for (const block of output.content) { // partialJson is only a streaming scratch buffer; never persist it. delete (block as { partialJson?: string }).partialJson; } output.stopReason = options?.signal?.aborted ? "aborted" : "error"; output.errorMessage = formatProviderError(normalizeProviderError(error)); stream.push({ type: "error", reason: output.stopReason, error: output }); stream.end(); } })(); return stream; }; export const streamSimple: StreamFunction<"openai-codex-responses", SimpleStreamOptions> = ( model: Model<"openai-codex-responses">, context: Context, options?: SimpleStreamOptions, ): AssistantMessageEventStream => { const apiKey = options?.apiKey; if (!apiKey) { throw new Error(`No API key for provider: ${model.provider}`); } const base = buildBaseOptions(model, context, options, apiKey); const clampedReasoning = options?.reasoning ? clampThinkingLevel(model, options.reasoning) : undefined; const reasoningEffort = clampedReasoning === "off" ? undefined : clampedReasoning; return stream(model, context, { ...base, reasoningEffort, } satisfies OpenAICodexResponsesOptions); }; // ============================================================================ // Request Building // ============================================================================ function buildRequestBody( model: Model<"openai-codex-responses">, context: Context, options?: OpenAICodexResponsesOptions, ): RequestBody { const toolPlacement = splitDeferredTools(context, model.compat?.supportsToolSearch ?? false); const messages = convertResponsesMessages(model, context, CODEX_TOOL_CALL_PROVIDERS, { includeSystemPrompt: false, deferredTools: toolPlacement.deferred, }); const body: RequestBody = { model: model.id, store: false, stream: true, instructions: context.systemPrompt || "You are a helpful assistant.", input: messages, text: { verbosity: options?.textVerbosity || "low" }, include: ["reasoning.encrypted_content"], prompt_cache_key: clampOpenAIPromptCacheKey(options?.sessionId), tool_choice: options?.toolChoice ?? "auto", parallel_tool_calls: true, }; if (options?.temperature !== undefined) { body.temperature = options.temperature; } if (options?.serviceTier !== undefined) { body.service_tier = options.serviceTier; } if (toolPlacement.immediate.length > 0) { body.tools = convertResponsesTools(toolPlacement.immediate, { strict: null }); } if (options?.reasoningEffort !== undefined) { const effort = options.reasoningEffort === "none" ? (model.thinkingLevelMap?.off ?? "none") : (model.thinkingLevelMap?.[options.reasoningEffort] ?? options.reasoningEffort); if (effort !== null) { body.reasoning = { effort, summary: options.reasoningSummary ?? "auto", }; } } return body; } function getServiceTierCostMultiplier( model: Pick, "id">, serviceTier: ResponseCreateParamsStreaming["service_tier"] | undefined, ): number { switch (serviceTier) { case "flex": return 0.5; case "priority": return model.id === "gpt-5.5" ? 2.5 : 2; default: return 1; } } function applyServiceTierPricing( usage: Usage, serviceTier: ResponseCreateParamsStreaming["service_tier"] | undefined, model: Pick, "id">, ) { const multiplier = getServiceTierCostMultiplier(model, serviceTier); if (multiplier === 1) return; usage.cost.input *= multiplier; usage.cost.output *= multiplier; usage.cost.cacheRead *= multiplier; usage.cost.cacheWrite *= multiplier; usage.cost.total = usage.cost.input + usage.cost.output + usage.cost.cacheRead + usage.cost.cacheWrite; } function resolveCodexServiceTier( responseServiceTier: ResponseCreateParamsStreaming["service_tier"] | undefined, requestServiceTier: ResponseCreateParamsStreaming["service_tier"] | undefined, ): ResponseCreateParamsStreaming["service_tier"] | undefined { if (responseServiceTier === "default" && (requestServiceTier === "flex" || requestServiceTier === "priority")) { return requestServiceTier; } return responseServiceTier ?? requestServiceTier; } function resolveCodexUrl(baseUrl?: string): string { const raw = baseUrl && baseUrl.trim().length > 0 ? baseUrl : DEFAULT_CODEX_BASE_URL; const normalized = raw.replace(/\/+$/, ""); if (normalized.endsWith("/codex/responses")) return normalized; if (normalized.endsWith("/codex")) return `${normalized}/responses`; return `${normalized}/codex/responses`; } function resolveCodexWebSocketUrl(baseUrl?: string): string { const url = new URL(resolveCodexUrl(baseUrl)); if (url.protocol === "https:") url.protocol = "wss:"; if (url.protocol === "http:") url.protocol = "ws:"; return url.toString(); } // ============================================================================ // Response Processing // ============================================================================ async function processStream( response: Response, output: AssistantMessage, stream: AssistantMessageEventStream, model: Model<"openai-codex-responses">, options?: OpenAICodexResponsesOptions, ): Promise { await processResponsesStream(mapCodexEvents(parseSSE(response, options?.signal)), output, stream, model, { serviceTier: options?.serviceTier, resolveServiceTier: resolveCodexServiceTier, applyServiceTierPricing: (usage, serviceTier) => applyServiceTierPricing(usage, serviceTier, model), }); } class CodexApiError extends Error { readonly code?: string; readonly payload?: Record; constructor(message: string, options?: { code?: string; payload?: Record; cause?: unknown }) { super(message); this.name = "CodexApiError"; this.code = options?.code; this.payload = options?.payload; this.cause = options?.cause; } } class CodexProtocolError extends Error { readonly payload?: unknown; constructor(message: string, options?: { payload?: unknown; cause?: unknown }) { super(message); this.name = "CodexProtocolError"; this.payload = options?.payload; this.cause = options?.cause; } } function isCodexNonTransportError(error: unknown): boolean { return error instanceof CodexApiError || error instanceof CodexProtocolError; } function isWebSocketConnectionLimitReachedError(error: unknown): boolean { return error instanceof CodexApiError && error.code === WEBSOCKET_CONNECTION_LIMIT_REACHED_CODE; } function isPreviousResponseNotFoundError(error: unknown): boolean { return error instanceof CodexApiError && error.code === PREVIOUS_RESPONSE_NOT_FOUND_CODE; } function extractCodexEventError(event: Record): { code?: string; message?: string } { const nested = event.error && typeof event.error === "object" ? (event.error as Record) : undefined; return { code: typeof event.code === "string" ? event.code : typeof nested?.code === "string" ? nested.code : undefined, message: typeof event.message === "string" ? event.message : typeof nested?.message === "string" ? nested.message : undefined, }; } async function* mapCodexEvents(events: AsyncIterable>): AsyncGenerator { for await (const event of events) { const type = typeof event.type === "string" ? event.type : undefined; if (!type) continue; if (type === "error") { const { code, message } = extractCodexEventError(event); throw new CodexApiError(`Codex error: ${message || code || JSON.stringify(event)}`, { code, payload: event, }); } if (type === "response.failed") { const response = (event as { response?: { error?: { code?: string; message?: string } } }).response; const code = response?.error?.code; const message = response?.error?.message; throw new CodexApiError(message || "Codex response failed", { code, payload: event }); } if (type === "response.done" || type === "response.completed" || type === "response.incomplete") { const response = (event as { response?: { status?: unknown } }).response; const normalizedResponse = response ? { ...response, status: normalizeCodexStatus(response.status) } : response; yield { ...event, type: "response.completed", response: normalizedResponse } as ResponseStreamEvent; return; } yield event as unknown as ResponseStreamEvent; } } function normalizeCodexStatus(status: unknown): CodexResponseStatus | undefined { if (typeof status !== "string") return undefined; return CODEX_RESPONSE_STATUSES.has(status as CodexResponseStatus) ? (status as CodexResponseStatus) : undefined; } // ============================================================================ // SSE Parsing // ============================================================================ async function* parseSSE(response: Response, signal?: AbortSignal): AsyncGenerator> { if (!response.body) return; const reader = response.body.getReader(); const decoder = new TextDecoder(); let buffer = ""; const onAbort = () => { void reader.cancel().catch(() => {}); }; signal?.addEventListener("abort", onAbort, { once: true }); try { while (true) { if (signal?.aborted) { throw new Error("Request was aborted"); } const { done, value } = await reader.read(); if (signal?.aborted) { throw new Error("Request was aborted"); } if (done) break; buffer += decoder.decode(value, { stream: true }); let idx = buffer.indexOf("\n\n"); while (idx !== -1) { const chunk = buffer.slice(0, idx); buffer = buffer.slice(idx + 2); const dataLines = chunk .split("\n") .filter((l) => l.startsWith("data:")) .map((l) => l.slice(5).trim()); if (dataLines.length > 0) { const data = dataLines.join("\n").trim(); if (data && data !== "[DONE]") { try { yield JSON.parse(data) as Record; } catch (cause) { throw new CodexProtocolError(`Invalid Codex SSE JSON: ${formatThrownValue(cause)}`, { cause, payload: data, }); } } } idx = buffer.indexOf("\n\n"); } } } finally { signal?.removeEventListener("abort", onAbort); try { await reader.cancel(); } catch {} try { reader.releaseLock(); } catch {} } } // ============================================================================ // WebSocket Parsing // ============================================================================ const OPENAI_BETA_RESPONSES_WEBSOCKETS = "responses_websockets=2026-02-06"; const SESSION_WEBSOCKET_CACHE_TTL_MS = 5 * 60 * 1000; const SESSION_WEBSOCKET_MAX_AGE_MS = 55 * 60 * 1000; type WebSocketEventType = "open" | "message" | "error" | "close"; type WebSocketListener = (event: unknown) => void; interface WebSocketLike { close(code?: number, reason?: string): void; send(data: string): void; addEventListener(type: WebSocketEventType, listener: WebSocketListener): void; removeEventListener(type: WebSocketEventType, listener: WebSocketListener): void; } interface CachedWebSocketContinuationState { lastRequestBody: RequestBody; lastResponseId: string; lastResponseItems: ResponseInput; } interface CachedWebSocketConnection { socket: WebSocketLike; busy: boolean; createdAt: number; idleTimer?: ReturnType; continuation?: CachedWebSocketContinuationState; } export interface OpenAICodexWebSocketDebugStats { requests: number; connectionsCreated: number; connectionsReused: number; cachedContextRequests: number; storeTrueRequests: number; fullContextRequests: number; deltaRequests: number; lastInputItems: number; lastDeltaInputItems?: number; lastPreviousResponseId?: string; websocketFailures: number; sseFallbacks: number; websocketFallbackActive?: boolean; lastWebSocketError?: string; } const websocketSessionCache = new Map(); const websocketDebugStats = new Map(); const websocketSseFallbackSessions = new Set(); function getOrCreateWebSocketDebugStats(sessionId: string): OpenAICodexWebSocketDebugStats { let stats = websocketDebugStats.get(sessionId); if (!stats) { stats = { requests: 0, connectionsCreated: 0, connectionsReused: 0, cachedContextRequests: 0, storeTrueRequests: 0, fullContextRequests: 0, deltaRequests: 0, lastInputItems: 0, websocketFailures: 0, sseFallbacks: 0, }; websocketDebugStats.set(sessionId, stats); } return stats; } export function getOpenAICodexWebSocketDebugStats(sessionId: string): OpenAICodexWebSocketDebugStats | undefined { const stats = websocketDebugStats.get(sessionId); return stats ? { ...stats } : undefined; } export function resetOpenAICodexWebSocketDebugStats(sessionId?: string): void { if (sessionId) { websocketDebugStats.delete(sessionId); websocketSseFallbackSessions.delete(sessionId); return; } websocketDebugStats.clear(); websocketSseFallbackSessions.clear(); } export function closeOpenAICodexWebSocketSessions(sessionId?: string): void { const closeEntry = (entry: CachedWebSocketConnection) => { if (entry.idleTimer) clearTimeout(entry.idleTimer); closeWebSocketSilently(entry.socket, 1000, "debug_close"); }; if (sessionId) { const entry = websocketSessionCache.get(sessionId); if (entry) closeEntry(entry); websocketSessionCache.delete(sessionId); return; } for (const entry of websocketSessionCache.values()) { closeEntry(entry); } websocketSessionCache.clear(); } registerSessionResourceCleanup(closeOpenAICodexWebSocketSessions); function isWebSocketSseFallbackActive(sessionId: string | undefined): boolean { return sessionId ? websocketSseFallbackSessions.has(sessionId) : false; } function recordWebSocketSseFallback(sessionId: string | undefined): void { if (!sessionId) return; const stats = getOrCreateWebSocketDebugStats(sessionId); stats.sseFallbacks++; stats.websocketFallbackActive = isWebSocketSseFallbackActive(sessionId); } function recordWebSocketFailure(sessionId: string | undefined, error: unknown): void { if (!sessionId) return; websocketSseFallbackSessions.add(sessionId); const stats = getOrCreateWebSocketDebugStats(sessionId); stats.websocketFailures++; stats.lastWebSocketError = formatThrownValue(error); stats.websocketFallbackActive = true; } type WebSocketConstructor = new ( url: string, protocols?: string | string[] | { headers?: Record }, ) => WebSocketLike; let _cachedWebsocket: WebSocketConstructor | null = null; async function getWebSocketConstructor(env?: ProviderEnv): Promise { if (!env && _cachedWebsocket) return _cachedWebsocket; // bun doesn't respect http proxy envs, ref: https://github.com/oven-sh/bun/issues/15489 // TODO: remove this when bun supports proxy envs in websocket. if (typeof process !== "undefined" && process.versions?.bun) { const WebSocketWithProxy = class extends WebSocket { constructor(url: string | URL, options?: string | string[] | Record) { let _opts: Record = {}; if (Array.isArray(options) || typeof options === "string") { _opts = { protocols: options }; } else { _opts = { ...options }; } const proxyUrl = resolveHttpProxyUrlForTarget( url.toString().replace(/^wss:/, "https:").replace(/^ws:/, "http:"), env, ); super(url, { ..._opts, ...(proxyUrl ? { proxy: proxyUrl.toString() } : {}) } as any); } }; if (!env) { _cachedWebsocket = WebSocketWithProxy; } return WebSocketWithProxy; } const ctor = (globalThis as { WebSocket?: unknown }).WebSocket; if (typeof ctor !== "function") return null; return ctor as unknown as WebSocketConstructor; } class WebSocketCloseError extends Error { readonly code?: number; readonly reason?: string; readonly wasClean?: boolean; constructor(message: string, options?: { code?: number; reason?: string; wasClean?: boolean }) { super(message); this.name = "WebSocketCloseError"; this.code = options?.code; this.reason = options?.reason; this.wasClean = options?.wasClean; } } function getWebSocketReadyState(socket: WebSocketLike): number | undefined { const readyState = (socket as { readyState?: unknown }).readyState; return typeof readyState === "number" ? readyState : undefined; } function isWebSocketReusable(socket: WebSocketLike): boolean { const readyState = getWebSocketReadyState(socket); // If readyState is unavailable, assume the runtime keeps it open/reusable. return readyState === undefined || readyState === 1; } function isWebSocketSessionExpired(entry: CachedWebSocketConnection): boolean { return Date.now() - entry.createdAt >= SESSION_WEBSOCKET_MAX_AGE_MS; } function closeWebSocketSilently(socket: WebSocketLike, code = 1000, reason = "done"): void { try { socket.close(code, reason); } catch {} } function scheduleSessionWebSocketExpiry(sessionId: string, entry: CachedWebSocketConnection): void { if (entry.idleTimer) { clearTimeout(entry.idleTimer); } entry.idleTimer = setTimeout(() => { if (entry.busy) return; closeWebSocketSilently(entry.socket, 1000, "idle_timeout"); websocketSessionCache.delete(sessionId); }, SESSION_WEBSOCKET_CACHE_TTL_MS); } async function connectWebSocket( url: string, headers: Headers, signal?: AbortSignal, connectTimeoutMs = DEFAULT_WEBSOCKET_CONNECT_TIMEOUT_MS, env?: ProviderEnv, ): Promise { const WebSocketCtor = await getWebSocketConstructor(env); if (!WebSocketCtor) { throw new Error("WebSocket transport is not available in this runtime"); } const wsHeaders = headersToRecord(headers); delete wsHeaders["OpenAI-Beta"]; return new Promise((resolve, reject) => { let settled = false; let timeout: ReturnType | undefined; let socket: WebSocketLike; try { socket = new WebSocketCtor(url, { headers: wsHeaders }); } catch (error) { reject(error instanceof Error ? error : new Error(String(error))); return; } const cleanup = () => { if (timeout) { clearTimeout(timeout); timeout = undefined; } socket.removeEventListener("open", onOpen); socket.removeEventListener("error", onError); socket.removeEventListener("close", onClose); signal?.removeEventListener("abort", onAbort); }; const fail = (error: Error, closeReason?: string) => { if (settled) return; settled = true; cleanup(); if (closeReason) { closeWebSocketSilently(socket, 1000, closeReason); } reject(error); }; const onOpen: WebSocketListener = () => { if (settled) return; settled = true; cleanup(); resolve(socket); }; const onError: WebSocketListener = (event) => { fail(extractWebSocketError(event)); }; const onClose: WebSocketListener = (event) => { fail(extractWebSocketCloseError(event)); }; const onAbort = () => { fail(new Error("Request was aborted"), "aborted"); }; socket.addEventListener("open", onOpen); socket.addEventListener("error", onError); socket.addEventListener("close", onClose); signal?.addEventListener("abort", onAbort); if (connectTimeoutMs > 0) { timeout = setTimeout(() => { fail(new Error(`WebSocket connect timeout after ${connectTimeoutMs}ms`), "connect_timeout"); }, connectTimeoutMs); } if (signal?.aborted) { onAbort(); } }); } async function acquireWebSocket( url: string, headers: Headers, sessionId: string | undefined, signal?: AbortSignal, connectTimeoutMs?: number, env?: ProviderEnv, ): Promise<{ socket: WebSocketLike; entry?: CachedWebSocketConnection; reused: boolean; release: (options?: { keep?: boolean }) => void; }> { if (!sessionId) { const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env); return { socket, reused: false, release: () => closeWebSocketSilently(socket), }; } const cached = websocketSessionCache.get(sessionId); if (cached) { if (cached.idleTimer) { clearTimeout(cached.idleTimer); cached.idleTimer = undefined; } if (!cached.busy && isWebSocketSessionExpired(cached)) { closeWebSocketSilently(cached.socket, 1000, "connection_age_limit"); websocketSessionCache.delete(sessionId); } else if (!cached.busy && isWebSocketReusable(cached.socket)) { cached.busy = true; return { socket: cached.socket, entry: cached, reused: true, release: ({ keep } = {}) => { if (!keep || !isWebSocketReusable(cached.socket)) { closeWebSocketSilently(cached.socket); websocketSessionCache.delete(sessionId); return; } cached.busy = false; scheduleSessionWebSocketExpiry(sessionId, cached); }, }; } if (cached.busy) { const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env); return { socket, reused: false, release: () => { closeWebSocketSilently(socket); }, }; } if (!isWebSocketReusable(cached.socket)) { closeWebSocketSilently(cached.socket); websocketSessionCache.delete(sessionId); } } const socket = await connectWebSocket(url, headers, signal, connectTimeoutMs, env); const entry: CachedWebSocketConnection = { socket, busy: true, createdAt: Date.now() }; websocketSessionCache.set(sessionId, entry); return { socket, entry, reused: false, release: ({ keep } = {}) => { if (!keep || !isWebSocketReusable(entry.socket)) { closeWebSocketSilently(entry.socket); if (entry.idleTimer) clearTimeout(entry.idleTimer); if (websocketSessionCache.get(sessionId) === entry) { websocketSessionCache.delete(sessionId); } return; } entry.busy = false; scheduleSessionWebSocketExpiry(sessionId, entry); }, }; } function extractWebSocketError(event: unknown): Error { if (event && typeof event === "object") { const message = "message" in event ? (event as { message?: unknown }).message : undefined; if (typeof message === "string" && message.length > 0) { return new Error(message); } const nestedError = "error" in event ? (event as { error?: unknown }).error : undefined; if (nestedError instanceof Error && nestedError.message.length > 0) { return nestedError; } if (nestedError && typeof nestedError === "object" && "message" in nestedError) { const nestedMessage = (nestedError as { message?: unknown }).message; if (typeof nestedMessage === "string" && nestedMessage.length > 0) { return new Error(nestedMessage); } } } return new Error("WebSocket error"); } function extractWebSocketCloseError(event: unknown): Error { if (event && typeof event === "object") { const code = "code" in event ? (event as { code?: unknown }).code : undefined; const reason = "reason" in event ? (event as { reason?: unknown }).reason : undefined; const wasClean = "wasClean" in event ? (event as { wasClean?: unknown }).wasClean : undefined; const codeText = typeof code === "number" ? ` ${code}` : ""; let reasonText = typeof reason === "string" && reason.length > 0 ? ` ${reason}` : ""; if (!reasonText && code === WEBSOCKET_MESSAGE_TOO_BIG_CLOSE_CODE) { reasonText = " message too big"; } return new WebSocketCloseError(`WebSocket closed${codeText}${reasonText}`.trim(), { code: typeof code === "number" ? code : undefined, reason: typeof reason === "string" && reason.length > 0 ? reason : undefined, wasClean: typeof wasClean === "boolean" ? wasClean : undefined, }); } return new Error("WebSocket closed"); } async function decodeWebSocketData(data: unknown): Promise { if (typeof data === "string") return data; if (data instanceof ArrayBuffer) { return new TextDecoder().decode(new Uint8Array(data)); } if (ArrayBuffer.isView(data)) { const view = data as ArrayBufferView; return new TextDecoder().decode(new Uint8Array(view.buffer, view.byteOffset, view.byteLength)); } if (data && typeof data === "object" && "arrayBuffer" in data) { const blobLike = data as { arrayBuffer: () => Promise }; const arrayBuffer = await blobLike.arrayBuffer(); return new TextDecoder().decode(new Uint8Array(arrayBuffer)); } return null; } async function* parseWebSocket( socket: WebSocketLike, signal?: AbortSignal, idleTimeoutMs?: number, ): AsyncGenerator> { const queue: Record[] = []; let pending: (() => void) | null = null; let done = false; let failed: Error | null = null; let sawCompletion = false; const wake = () => { if (!pending) return; const resolve = pending; pending = null; resolve(); }; const onMessage: WebSocketListener = (event) => { void (async () => { let text: string | null = null; try { if (!event || typeof event !== "object" || !("data" in event)) return; text = await decodeWebSocketData((event as { data?: unknown }).data); if (!text) return; const parsed = JSON.parse(text) as Record; const type = typeof parsed.type === "string" ? parsed.type : ""; if (type === "response.completed" || type === "response.done" || type === "response.incomplete") { sawCompletion = true; done = true; } queue.push(parsed); wake(); } catch (cause) { failed = new CodexProtocolError(`Invalid Codex WebSocket JSON: ${formatThrownValue(cause)}`, { cause, payload: text, }); done = true; wake(); } })(); }; const onError: WebSocketListener = (event) => { failed = extractWebSocketError(event); done = true; wake(); }; const onClose: WebSocketListener = (event) => { if (sawCompletion) { done = true; wake(); return; } if (!failed) { failed = extractWebSocketCloseError(event); } done = true; wake(); }; const onAbort = () => { failed = new Error("Request was aborted"); done = true; wake(); }; socket.addEventListener("message", onMessage); socket.addEventListener("error", onError); socket.addEventListener("close", onClose); signal?.addEventListener("abort", onAbort); try { while (true) { if (signal?.aborted) { throw new Error("Request was aborted"); } if (queue.length > 0) { yield queue.shift()!; continue; } if (done) break; let timeout: ReturnType | undefined; await new Promise((resolve, reject) => { pending = resolve; if (idleTimeoutMs !== undefined && idleTimeoutMs > 0) { timeout = setTimeout(() => { const error = new Error(`WebSocket idle timeout after ${idleTimeoutMs}ms`); failed = error; done = true; pending = null; closeWebSocketSilently(socket, 1000, "idle_timeout"); reject(error); }, idleTimeoutMs); } }).finally(() => { if (timeout) { clearTimeout(timeout); } }); } if (failed) { throw failed; } if (!sawCompletion) { throw new Error("WebSocket stream closed before response.completed"); } } finally { socket.removeEventListener("message", onMessage); socket.removeEventListener("error", onError); socket.removeEventListener("close", onClose); signal?.removeEventListener("abort", onAbort); } } function requestBodyWithoutInput(body: RequestBody): RequestBody { const { input: _input, previous_response_id: _previousResponseId, ...rest } = body; return rest; } function responseInputsEqual(a: ResponseInput | undefined, b: ResponseInput | undefined): boolean { return JSON.stringify(a ?? []) === JSON.stringify(b ?? []); } function requestBodiesMatchExceptInput(a: RequestBody, b: RequestBody): boolean { return JSON.stringify(requestBodyWithoutInput(a)) === JSON.stringify(requestBodyWithoutInput(b)); } function getCachedWebSocketInputDelta( body: RequestBody, continuation: CachedWebSocketContinuationState, ): ResponseInput | undefined { if (!requestBodiesMatchExceptInput(body, continuation.lastRequestBody)) { return undefined; } const currentInput = body.input ?? []; const baseline = [...(continuation.lastRequestBody.input ?? []), ...continuation.lastResponseItems]; if (currentInput.length < baseline.length) { return undefined; } const prefix = currentInput.slice(0, baseline.length); if (!responseInputsEqual(prefix, baseline)) { return undefined; } return currentInput.slice(baseline.length); } function buildCachedWebSocketRequestBody(entry: CachedWebSocketConnection, body: RequestBody): RequestBody { const continuation = entry.continuation; if (!continuation) { return body; } const delta = getCachedWebSocketInputDelta(body, continuation); if (!delta || !continuation.lastResponseId) { entry.continuation = undefined; return body; } return { ...body, previous_response_id: continuation.lastResponseId, input: delta, }; } async function* startWebSocketOutputOnFirstEvent( events: AsyncIterable, onStart: () => void, ): AsyncGenerator { let started = false; for await (const event of events) { if (!started) { started = true; onStart(); } yield event; } } async function processWebSocketStream( url: string, body: RequestBody, headers: Headers, output: AssistantMessage, stream: AssistantMessageEventStream, model: Model<"openai-codex-responses">, onStart: () => void, idleTimeoutMs: number | undefined, websocketConnectTimeoutMs: number | undefined, options?: OpenAICodexResponsesOptions, ): Promise { const { socket, entry, reused, release } = await acquireWebSocket( url, headers, options?.sessionId, options?.signal, websocketConnectTimeoutMs, options?.env, ); let keepConnection = true; const useCachedContext = options?.transport === "websocket-cached" || options?.transport === "auto"; // ChatGPT Codex Responses rejects `store: true` ("Store must be set to false"). // WebSocket continuation still works via connection-scoped previous_response_id state. const fullBody = body; const requestBody = useCachedContext && entry ? buildCachedWebSocketRequestBody(entry, fullBody) : fullBody; const stats = options?.sessionId ? getOrCreateWebSocketDebugStats(options.sessionId) : undefined; if (stats) { stats.requests++; if (reused) stats.connectionsReused++; else stats.connectionsCreated++; if (useCachedContext) stats.cachedContextRequests++; if (requestBody.store === true) stats.storeTrueRequests++; stats.lastInputItems = requestBody.input?.length ?? 0; if (requestBody.previous_response_id) { stats.deltaRequests++; stats.lastDeltaInputItems = requestBody.input?.length ?? 0; stats.lastPreviousResponseId = requestBody.previous_response_id; } else { stats.fullContextRequests++; stats.lastDeltaInputItems = undefined; stats.lastPreviousResponseId = undefined; } } try { socket.send(JSON.stringify({ type: "response.create", ...requestBody })); await processResponsesStream( startWebSocketOutputOnFirstEvent( mapCodexEvents(parseWebSocket(socket, options?.signal, idleTimeoutMs)), onStart, ), output, stream, model, { serviceTier: options?.serviceTier, resolveServiceTier: resolveCodexServiceTier, applyServiceTierPricing: (usage, serviceTier) => applyServiceTierPricing(usage, serviceTier, model), }, ); if (options?.signal?.aborted) { keepConnection = false; } else if (useCachedContext && entry && output.responseId) { const responseItems = convertResponsesMessages(model, { messages: [output] }, CODEX_TOOL_CALL_PROVIDERS, { includeSystemPrompt: false, }).filter((item) => item.type !== "function_call_output"); entry.continuation = { lastRequestBody: fullBody, lastResponseId: output.responseId, lastResponseItems: responseItems, }; } } catch (error) { if (entry) { entry.continuation = undefined; } keepConnection = false; throw error; } finally { release({ keep: keepConnection }); } } // ============================================================================ // Error Handling // ============================================================================ async function parseErrorResponse(response: Response): Promise<{ message: string; friendlyMessage?: string }> { const raw = await response.text(); let message = raw || response.statusText || "Request failed"; let friendlyMessage: string | undefined; try { const parsed = JSON.parse(raw) as { error?: { code?: string; type?: string; message?: string; plan_type?: string; resets_at?: number }; }; const err = parsed?.error; if (err) { const code = err.code || err.type || ""; if (/usage_limit_reached|usage_not_included|rate_limit_exceeded/i.test(code) || response.status === 429) { const plan = err.plan_type ? ` (${err.plan_type.toLowerCase()} plan)` : ""; const mins = err.resets_at ? Math.max(0, Math.round((err.resets_at * 1000 - Date.now()) / 60000)) : undefined; const when = mins !== undefined ? ` Try again in ~${mins} min.` : ""; friendlyMessage = `You have hit your ChatGPT usage limit${plan}.${when}`.trim(); } message = err.message || friendlyMessage || message; } } catch {} return { message, friendlyMessage }; } // ============================================================================ // Auth & Headers // ============================================================================ function extractAccountId(token: string): string { try { const parts = token.split("."); if (parts.length !== 3) throw new Error("Invalid token"); const payload = JSON.parse(atob(parts[1])); const accountId = payload?.[JWT_CLAIM_PATH]?.chatgpt_account_id; if (!accountId) throw new Error("No account ID in token"); return accountId; } catch { throw new Error("Failed to extract accountId from token"); } } function buildBaseCodexHeaders( initHeaders: Record | undefined, additionalHeaders: ProviderHeaders | undefined, accountId: string, token: string, ): Headers { const headers = new Headers(initHeaders); for (const [key, value] of Object.entries(additionalHeaders || {})) { if (value === null) { headers.delete(key); } else { headers.set(key, value); } } headers.set("Authorization", `Bearer ${token}`); headers.set("chatgpt-account-id", accountId); headers.set("originator", "pi"); const userAgent = _os ? `pi (${_os.platform()} ${_os.release()}; ${_os.arch()})` : "pi (browser)"; headers.set("User-Agent", userAgent); return headers; } function buildSSEHeaders( initHeaders: Record | undefined, additionalHeaders: ProviderHeaders | undefined, accountId: string, token: string, sessionId?: string, ): Headers { const headers = buildBaseCodexHeaders(initHeaders, additionalHeaders, accountId, token); headers.set("OpenAI-Beta", "responses=experimental"); headers.set("accept", "text/event-stream"); headers.set("content-type", "application/json"); if (sessionId) { headers.set("session-id", sessionId); headers.set("x-client-request-id", sessionId); } return headers; } function buildWebSocketHeaders( initHeaders: Record | undefined, additionalHeaders: ProviderHeaders | undefined, accountId: string, token: string, requestId: string, ): Headers { const headers = buildBaseCodexHeaders(initHeaders, additionalHeaders, accountId, token); headers.delete("accept"); headers.delete("content-type"); headers.delete("OpenAI-Beta"); headers.delete("openai-beta"); headers.set("OpenAI-Beta", OPENAI_BETA_RESPONSES_WEBSOCKETS); headers.set("x-client-request-id", requestId); headers.set("session-id", requestId); return headers; }