From 410f04c939ed4505617a11bf1c6ee1b627b0070b Mon Sep 17 00:00:00 2001 From: Kosta Milovanovic Date: Tue, 8 Sep 2026 15:11:11 -0400 Subject: [PATCH 1/2] feat(client): add opt-in remote hub voice relay --- .../content/docs/guides/codex-integration.md | 5 + .../src/content/docs/guides/remote-hub.md | 45 ++ scripts/test-layout/layout.json | 1 + .../ocx/references/01_management_surface.md | 18 +- src/cli/capabilities.ts | 15 + src/cli/dispatch.ts | 4 + src/cli/help.ts | 1 + src/cli/registry.ts | 9 + src/cli/voice-relay.ts | 33 ++ src/client/voice-relay.ts | 427 ++++++++++++++++++ tests/clients/client-voice-relay.test.ts | 265 +++++++++++ tests/fixtures/test-layout-expected.json | 1 + 12 files changed, 823 insertions(+), 1 deletion(-) create mode 100644 src/cli/voice-relay.ts create mode 100644 src/client/voice-relay.ts create mode 100644 tests/clients/client-voice-relay.test.ts diff --git a/docs-site/src/content/docs/guides/codex-integration.md b/docs-site/src/content/docs/guides/codex-integration.md index 5b75a21b41..793788252d 100644 --- a/docs-site/src/content/docs/guides/codex-integration.md +++ b/docs-site/src/content/docs/guides/codex-integration.md @@ -49,6 +49,11 @@ current bearer, so the key only keeps the join on the proxy path. It is written `openai_base_url` form, is removed together with it, and a user-owned `experimental_realtime_ws_base_url` is never overwritten. +Authenticated remote clients use a different routing form: the provider-table admission header +does not automatically accompany Codex's dedicated voice transports. Pointing the root voice +URLs directly at a protected hub is therefore not equivalent to the loopback setup above. +See [Remote client voice](/guides/remote-hub/#remote-client-voice) for the opt-in local relay. + ### Voice transport and task handoffs Codex owns the microphone and speaker, WebRTC media negotiation, captions, mute controls, and diff --git a/docs-site/src/content/docs/guides/remote-hub.md b/docs-site/src/content/docs/guides/remote-hub.md index 6a2b2a33bd..7ba842f0a3 100644 --- a/docs-site/src/content/docs/guides/remote-hub.md +++ b/docs-site/src/content/docs/guides/remote-hub.md @@ -297,6 +297,51 @@ both named volumes and a synthetic catalog survive replacement. This check does real provider account, OAuth callback, custom mount migration, or every CPU architecture; perform the authenticated routed-response check above for your deployment. +## Remote client voice + +Codex owns microphone capture, playback, and WebRTC media. The hub forwards call creation and +the realtime control WebSocket, but a remote provider's `x-opencodex-api-key` header is not +automatically included in those dedicated voice transports. The loopback injection described in +[Codex integration](/guides/codex-integration/#config-injection) does not solve remote admission. + +On an already connected client, the opt-in `ocx voice-relay` command provides a local voice-only +transport using that client's existing hub data credential. It is a foreground command, not an +automatically installed service. It does not modify Codex configuration or normal model routing. +Start it with `ocx voice-relay --port 10111` and keep the process running while using voice. +The default port is `10111`; it binds only `127.0.0.1` and does not fall back to another port. +By default it accepts call creation at `POST /v1/live` and `POST /v1/realtime/calls`, followed +by a call-ID WebSocket join. Add `--allow-standalone` for a client that opens a realtime WebSocket +without first creating a WebRTC call. Unrelated API routes and browser-origin requests are not +general-purpose forwarding surfaces. + +The relay exits if its saved connection or credential changes. After disconnect or key rotation, +restore the voice settings or restart the relay against the intended connection. It never repairs +pairing or rotates keys itself. + +Back up the **user-level** Codex config (`$CODEX_HOME/config.toml`, normally +`~/.codex/config.toml`) and set these root keys **before the first TOML table**, using the relay's +reported port. This example assumes port `10111`: + +```toml +experimental_realtime_webrtc_call_base_url = "http://127.0.0.1:10111/v1" +experimental_realtime_ws_base_url = "http://127.0.0.1:10111/v1" +``` + +Do not change `openai_base_url`, `model_provider`, or the generated provider table for this +workaround. Do not put the hub credential in a URL or copy it into these keys. These experimental +settings require a Codex version that supports them; the WebRTC key affects call creation only, +and the WebSocket key affects realtime control only. See the +[upstream config definitions](https://github.com/openai/codex/blob/main/codex-rs/config/src/config_toml.rs). + +Fully quit and reopen Codex when no active work would be interrupted. A listening relay or a +successful hub health check does not prove voice works: start a voice conversation, confirm +microphone input and spoken output, and verify a voice task handoff separately. The relay does +not add voice support to clients that lack it or replace the hub's upstream voice authentication. + +To roll back, restore only the previous values of these two keys (remove them if previously +absent), quit and reopen Codex safely, then stop the foreground relay. Preserve unrelated edits +made since the backup; normal provider routing and hub pairing do not need to be removed. + ## Rollback Inspect existing Serve mappings before changing them. `tailscale serve reset` removes every mapping diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 53a8e65898..57a2a6b746 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -363,6 +363,7 @@ "client-export-modality-enum.test.ts": "clients", "client-fingerprint.test.ts": "clients", "client-hub-relay.test.ts": "clients", + "client-voice-relay.test.ts": "clients", "client-injection-guard.test.ts": "codex-integration", "client-lifecycle-lock.test.ts": "clients", "client-machine-listener.test.ts": "clients", diff --git a/skills/ocx/references/01_management_surface.md b/skills/ocx/references/01_management_surface.md index 512aa3a7e2..02dd2e9fee 100644 --- a/skills/ocx/references/01_management_surface.md +++ b/skills/ocx/references/01_management_surface.md @@ -58,6 +58,22 @@ JSON mode: `envelope`. - Reads /healthz plus local config; drives no management API route. +### `ocx voice-relay` + +Run a loopback-only foreground relay for a connected remote hub's realtime voice routes. + +Drives no management route. + +| Flag | Value | Meaning | +|---|---|---| +| `--port` | number | Loopback port; defaults to 10111. | +| `--allow-standalone` | boolean | Also admit bare /v1/live and /v1/realtime WebSocket sessions. | + +JSON mode: `none`. + +- Requires an intact ocx connect record and owner-matching data credential. +- Does not write Codex configuration or install a service; exits when connection ownership changes. + ### `ocx capabilities` List the declared CLI capabilities and the management routes they drive. @@ -687,6 +703,6 @@ JSON mode: `payload`. ## Counts -- declared capabilities: 37 +- declared capabilities: 38 - of those, state-changing: 16 - head-resolved invocations: 2 diff --git a/src/cli/capabilities.ts b/src/cli/capabilities.ts index 86aa5438df..1600620d43 100644 --- a/src/cli/capabilities.ts +++ b/src/cli/capabilities.ts @@ -155,6 +155,21 @@ export const CAPABILITIES: readonly Capability[] = [ "A rotation left pending by a crash is resumed here — startup and status stop rather than guess which key generation is live.", ], }, + { + command: ["voice-relay"], + summary: "Run a loopback-only foreground relay for a connected remote hub's realtime voice routes.", + routes: [], + flags: [ + { name: "--port", value: "number", summary: "Loopback port; defaults to 10111." }, + { name: "--allow-standalone", value: "boolean", summary: "Also admit bare /v1/live and /v1/realtime WebSocket sessions." }, + ], + mutates: false, + json: "none", + details: [ + "Requires an intact ocx connect record and owner-matching data credential.", + "Does not write Codex configuration or install a service; exits when connection ownership changes.", + ], + }, { command: ["capabilities"], summary: "List the declared CLI capabilities and the management routes they drive.", diff --git a/src/cli/dispatch.ts b/src/cli/dispatch.ts index e85de1c05f..a1af43f329 100644 --- a/src/cli/dispatch.ts +++ b/src/cli/dispatch.ts @@ -438,6 +438,10 @@ const commandRunners: Record = { const { handleDisconnectCommand } = await import("./connect"); return await handleDisconnectCommand(deps.args.slice(1)); }, + "voice-relay": async deps => { + const { handleVoiceRelayCommand } = await import("./voice-relay"); + return await handleVoiceRelayCommand(deps.args.slice(1)); + }, "sync-cache": async deps => { const cacheArgs = deps.args.slice(1); const restartCodex = cacheArgs.includes("--restart-codex"); diff --git a/src/cli/help.ts b/src/cli/help.ts index 43916695b6..12aba44414 100644 --- a/src/cli/help.ts +++ b/src/cli/help.ts @@ -42,6 +42,7 @@ Usage: ocx ensure Ensure the proxy is running and Codex config/cache are current ocx connect Connect this machine to a remote OpenCodex hub (credential via stdin) ocx disconnect Restore local state and clear the hub connection + ocx voice-relay [flags] Foreground loopback relay for connected remote-hub voice ocx sync [--restart-codex] Fetch models from providers and inject into Codex config ocx sync-cache [--restart-codex] Refresh Codex's model cache from the active catalog diff --git a/src/cli/registry.ts b/src/cli/registry.ts index 73bbd68e31..799a72a35d 100644 --- a/src/cli/registry.ts +++ b/src/cli/registry.ts @@ -104,6 +104,15 @@ export const CLI_COMMANDS: CliCommandEntry[] = [ usage: "ocx disconnect [--keep-catalog] [--json]", summary: "Restore local client state offline and clear the remote-hub connection.", }, + { + name: "voice-relay", + usage: "ocx voice-relay [--port ] [--allow-standalone]", + summary: "Relay Codex realtime voice to the connected remote hub from a loopback-only foreground listener.", + details: [ + "Defaults to 127.0.0.1:10111 and exits when the connected hub credential changes.", + "No Codex config is changed; --allow-standalone additionally admits bare realtime WebSocket sessions.", + ], + }, { name: "sync", usage: "ocx sync [--restart-codex] [--restart-desktop-app]", diff --git a/src/cli/voice-relay.ts b/src/cli/voice-relay.ts new file mode 100644 index 0000000000..7d6fb605c8 --- /dev/null +++ b/src/cli/voice-relay.ts @@ -0,0 +1,33 @@ +import { startVoiceRelay, VOICE_RELAY_DEFAULT_PORT } from "../client/voice-relay"; +import { CliUsageError, rejectArgs, runCliAction, takeFlag, takeIntegerOption } from "./runtime-api"; + +export const VOICE_RELAY_USAGE = `Usage: + ocx voice-relay [--port ] [--allow-standalone]`; + +export async function handleVoiceRelayCommand(argv: string[]): Promise { + return runCliAction(async () => { + const args = [...argv]; + const allowStandalone = takeFlag(args, "--allow-standalone"); + const port = takeIntegerOption(args, "--port", { min: 1 }); + if (port !== undefined && port > 65_535) throw new CliUsageError("--port must be between 1 and 65535", VOICE_RELAY_USAGE); + rejectArgs(args, VOICE_RELAY_USAGE, { redactValues: true }); + const relay = startVoiceRelay({ port: port ?? VOICE_RELAY_DEFAULT_PORT, allowStandalone }); + console.log(`Voice relay listening on ${relay.origin}/v1 (Ctrl-C to stop).`); + console.log(`Codex config: experimental_realtime_webrtc_call_base_url = "${relay.origin}/v1"`); + console.log(`Codex config: experimental_realtime_ws_base_url = "${relay.origin}/v1"`); + let stopping = false; + const stop = () => { if (!stopping) { stopping = true; relay.stop(); } }; + process.once("SIGINT", stop); + process.once("SIGTERM", stop); + if (process.platform !== "win32") process.once("SIGHUP", stop); + try { + const reason = await relay.done; + if (reason === "connection_changed") throw new Error("connected hub or credential changed; voice relay stopped"); + } finally { + process.removeListener("SIGINT", stop); + process.removeListener("SIGTERM", stop); + if (process.platform !== "win32") process.removeListener("SIGHUP", stop); + relay.stop(); + } + }); +} diff --git a/src/client/voice-relay.ts b/src/client/voice-relay.ts new file mode 100644 index 0000000000..45a9d9f7b9 --- /dev/null +++ b/src/client/voice-relay.ts @@ -0,0 +1,427 @@ +import type { Server, ServerWebSocket } from "bun"; +import { readBoundedResponseBytes } from "../lib/bounded-body"; +import { clearableDeadline } from "../lib/abort"; +import { + LIVE_CLIENT_PROTOCOL_HEADERS, + parseLiveSidebandTarget, + sanitizeStandaloneRealtimeQuery, + type LiveSidebandTarget, +} from "../server/live"; +import type { OcxClientConnectionConfig } from "../types"; +import { readServiceApiTokenState } from "../lib/service-secrets"; +import { assertNoClientDisconnectPending, readClientConnectionState } from "./state"; + +export const VOICE_RELAY_DEFAULT_PORT = 10_111; +export const VOICE_RELAY_BODY_MAX_BYTES = 16 * 1024 * 1024; +export const VOICE_RELAY_WS_FRAME_MAX_BYTES = 1024 * 1024; +export const VOICE_RELAY_WS_PENDING_MAX_BYTES = 4 * 1024 * 1024; +export const VOICE_RELAY_WS_PENDING_MAX_FRAMES = 64; +export const VOICE_RELAY_WS_BACKPRESSURE_MAX_BYTES = 4 * 1024 * 1024; +export const VOICE_RELAY_CONNECT_TIMEOUT_MS = 15_000; +export const VOICE_RELAY_TOTAL_TIMEOUT_MS = 120_000; +export const VOICE_RELAY_WS_IDLE_SECONDS = 120; +export const VOICE_RELAY_CLOSE_FALLBACK_MS = 2_000; + +const HTTP_REQUEST_HEADERS = [ + "accept", "accept-language", "authorization", "chatgpt-account-id", "content-type", + ...LIVE_CLIENT_PROTOCOL_HEADERS, +] as const; +const WS_REQUEST_HEADERS = ["authorization", "chatgpt-account-id", ...LIVE_CLIENT_PROTOCOL_HEADERS] as const; +const HTTP_RESPONSE_HEADERS = [ + "cache-control", "content-language", "content-type", "location", "openai-processing-ms", "retry-after", +] as const; + +export interface VoiceRelayCredential { + connection: OcxClientConnectionConfig; + token: string; +} + +export interface VoiceRelayOptions { + port?: number; + allowStandalone?: boolean; + credential?: VoiceRelayCredential; + connectionCheck?: (expected: VoiceRelayCredential) => boolean; + fetchImpl?: typeof fetch; + webSocketFactory?: (url: string, headers: Record) => WebSocket; + connectTimeoutMs?: number; + totalTimeoutMs?: number; + closeFallbackMs?: number; + monitorIntervalMs?: number; +} + +export interface VoiceRelayHandle { + port: number; + origin: string; + done: Promise<"stopped" | "connection_changed">; + stop(): void; +} + +interface VoiceRelayWsData { + upstreamUrl: string; + headers: Record; + upstream?: WebSocket; + pending: Array; + pendingBytes: number; + opened: boolean; + closing: boolean; + handshakeTimer?: ReturnType; + closeTimer?: ReturnType; +} + +function jsonError(status: number, error: string): Response { + return Response.json({ error }, { status }); +} + +function frameBytes(value: string | Buffer | ArrayBuffer | ArrayBufferView): number { + if (typeof value === "string") return Buffer.byteLength(value); + return value.byteLength; +} + +function relayHeaders(source: Headers, names: readonly string[]): Headers { + const output = new Headers(); + for (const name of names) { + const value = source.get(name); + if (value !== null) output.set(name, value); + } + return output; +} + +function localAuthorityAllowed(req: Request, port: number): boolean { + let url: URL; + try { url = new URL(req.url); } catch { return false; } + const localHost = url.hostname === "127.0.0.1" || url.hostname.toLowerCase() === "localhost"; + if (!localHost || url.port !== String(port)) return false; + const origin = req.headers.get("origin"); + if (!origin) return true; + try { + const parsed = new URL(origin); + return parsed.protocol === "http:" + && (parsed.hostname === "127.0.0.1" || parsed.hostname.toLowerCase() === "localhost") + && parsed.port === String(port) + && parsed.pathname === "/" && !parsed.search && !parsed.hash; + } catch { + return false; + } +} + +function postRouteAllowed(url: URL, method: string): boolean { + return method === "POST" && (url.pathname === "/v1/live" || url.pathname === "/v1/realtime/calls"); +} + +export function voiceRelayWebSocketRouteAllowed(url: URL, allowStandalone = false): boolean { + const target = parseLiveSidebandTarget(url.pathname, url.searchParams, url.search.slice(1)); + if (!target) return false; + return allowStandalone + || target.style === "frameless-path" + || target.style === "realtime-calls-path" + || target.style === "realtime-query"; +} + +function sanitizedWebSocketPath(url: URL, target: LiveSidebandTarget): string { + const query = sanitizeStandaloneRealtimeQuery(url.search.slice(1)); + if (target.style === "realtime-query") { + const params = new URLSearchParams(query); + // Keep the parser-validated call id authoritative if duplicate values were supplied. + params.delete("call_id"); + params.set("call_id", target.callId); + return `${url.pathname}?${params.toString()}`; + } + return `${url.pathname}${query ? `?${query}` : ""}`; +} + +function upstreamUrl(origin: string, local: URL, websocket: boolean): string { + const target = new URL(`${local.pathname}${local.search}`, `${origin}/`); + if (websocket) target.protocol = target.protocol === "https:" ? "wss:" : "ws:"; + return target.toString(); +} + +function credentialStillOwned(expected: VoiceRelayCredential): boolean { + try { + assertNoClientDisconnectPending(); + const state = readClientConnectionState(); + const token = readServiceApiTokenState(); + return state.kind === "connected" + && JSON.stringify(state.value) === JSON.stringify(expected.connection) + && !state.value.pendingOperation + && token.kind === "present" + && token.fingerprint === expected.connection.tokenFingerprint; + } catch { + return false; + } +} + +export function loadVoiceRelayCredential(): VoiceRelayCredential { + assertNoClientDisconnectPending(); + const state = readClientConnectionState(); + if (state.kind !== "connected" || state.value.pendingOperation) { + throw new Error("voice relay requires a complete, stable 'ocx connect' connection"); + } + const token = readServiceApiTokenState(); + if (token.kind !== "present" || token.fingerprint !== state.value.tokenFingerprint) { + throw new Error(token.kind === "unsafe" + ? "connected voice relay credential is unreadable or unsafe" + : "connected voice relay credential is missing or no longer owned"); + } + return { connection: structuredClone(state.value), token: token.token }; +} + +async function readRequest(req: Request, maxBytes: number, signal: AbortSignal): Promise | Response> { + const declared = req.headers.get("content-length"); + if (declared !== null) { + const bytes = Number(declared); + if (!Number.isSafeInteger(bytes) || bytes < 0 || bytes > maxBytes) return jsonError(413, "voice_relay_request_too_large"); + } + try { + const result = await readBoundedResponseBytes(new Response(req.body), { + maxBytes, + signal, + inactivityTimeoutMs: VOICE_RELAY_CONNECT_TIMEOUT_MS, + }); + return result.oversized ? jsonError(413, "voice_relay_request_too_large") : result.bytes; + } catch { + return jsonError(req.signal.aborted ? 499 : 408, req.signal.aborted ? "voice_relay_client_closed" : "voice_relay_request_timeout"); + } +} + +async function relayHttp(req: Request, url: URL, credential: VoiceRelayCredential, options: VoiceRelayOptions): Promise { + const totalTimeoutMs = options.totalTimeoutMs ?? VOICE_RELAY_TOTAL_TIMEOUT_MS; + const total = AbortSignal.timeout(totalTimeoutMs); + const lifetime = AbortSignal.any([req.signal, total]); + const body = await readRequest(req, VOICE_RELAY_BODY_MAX_BYTES, lifetime); + if (body instanceof Response) return body; + if (!(options.connectionCheck ?? credentialStillOwned)(credential)) return jsonError(409, "voice_relay_connection_changed"); + const headers = relayHeaders(req.headers, HTTP_REQUEST_HEADERS); + headers.set("x-opencodex-api-key", credential.token); + const connect = clearableDeadline(options.connectTimeoutMs ?? VOICE_RELAY_CONNECT_TIMEOUT_MS, lifetime); + let response: Response; + try { + response = await (options.fetchImpl ?? fetch)(upstreamUrl(credential.connection.serverUrl, url, false), { + method: "POST", + headers, + body, + redirect: "manual", + signal: connect.signal, + }); + } catch { + return jsonError(connect.didExpire() ? 504 : 502, connect.didExpire() ? "voice_relay_connect_timeout" : "voice_relay_upstream_unreachable"); + } finally { + connect.clear(); + } + if (response.status >= 300 && response.status < 400) { + try { await response.body?.cancel(); } catch { /* best effort */ } + return jsonError(502, "voice_relay_redirect_refused"); + } + const declared = Number(response.headers.get("content-length") ?? "0"); + if (Number.isFinite(declared) && declared > VOICE_RELAY_BODY_MAX_BYTES) { + try { await response.body?.cancel(); } catch { /* best effort */ } + return jsonError(502, "voice_relay_response_too_large"); + } + try { + const result = await readBoundedResponseBytes(response, { + maxBytes: VOICE_RELAY_BODY_MAX_BYTES, + signal: lifetime, + inactivityTimeoutMs: options.connectTimeoutMs ?? VOICE_RELAY_CONNECT_TIMEOUT_MS, + }); + if (result.oversized) return jsonError(502, "voice_relay_response_too_large"); + return new Response(result.bytes, { status: response.status, statusText: response.statusText, headers: relayHeaders(response.headers, HTTP_RESPONSE_HEADERS) }); + } catch { + return jsonError(504, "voice_relay_response_timeout"); + } +} + +function safeCloseCode(code: number): number { + return code === 1000 || (code >= 1001 && code <= 1014 && code !== 1004 && code !== 1005 && code !== 1006) + || (code >= 3000 && code <= 4999) ? code : 1011; +} + +function safeCloseReason(reason: string): string { + return /^[\x20-\x7e]{0,80}$/.test(reason) ? reason : "peer closed"; +} + +function startWebSocketPeer( + downstream: ServerWebSocket, + options: VoiceRelayOptions, +): void { + const data = downstream.data; + const factory = options.webSocketFactory ?? ((url, headers) => new WebSocket(url, { headers } as unknown as string[])); + const finish = (code = 1000, reason = "") => { + if (data.closing) return; + data.closing = true; + if (data.handshakeTimer) clearTimeout(data.handshakeTimer); + data.pending = []; + data.pendingBytes = 0; + const upstream = data.upstream; + try { if (upstream && upstream.readyState < WebSocket.CLOSING) upstream.close(code, reason); } catch { /* best effort */ } + try { if (downstream.readyState < WebSocket.CLOSING) downstream.close(code, reason); } catch { /* best effort */ } + if (data.closeTimer) clearTimeout(data.closeTimer); + data.closeTimer = setTimeout(() => { + data.closeTimer = undefined; + try { + if (upstream && upstream.readyState !== WebSocket.CLOSED) { + (upstream as WebSocket & { terminate(): void }).terminate(); + } + } catch { /* best effort */ } + try { if (downstream.readyState !== WebSocket.CLOSED) downstream.terminate(); } catch { /* best effort */ } + }, options.closeFallbackMs ?? VOICE_RELAY_CLOSE_FALLBACK_MS); + }; + let upstream: WebSocket; + try { upstream = factory(data.upstreamUrl, data.headers); } + catch { finish(1011, "upstream connect failed"); return; } + data.upstream = upstream; + upstream.binaryType = "arraybuffer"; + data.handshakeTimer = setTimeout(() => finish(1011, "upstream handshake timeout"), options.connectTimeoutMs ?? VOICE_RELAY_CONNECT_TIMEOUT_MS); + upstream.addEventListener("open", () => { + if (data.closing) return; + if (data.handshakeTimer) clearTimeout(data.handshakeTimer); + data.handshakeTimer = undefined; + data.opened = true; + const pending = data.pending; + data.pending = []; + data.pendingBytes = 0; + for (const frame of pending) { + const bytes = frameBytes(frame); + if (upstream.bufferedAmount + bytes > VOICE_RELAY_WS_BACKPRESSURE_MAX_BYTES) { + finish(1013, "upstream backpressure limit"); + return; + } + upstream.send(typeof frame === "string" ? frame : Uint8Array.from(frame)); + } + }, { once: true }); + upstream.addEventListener("message", event => { + if (data.closing) return; + const payload = event.data; + if (typeof payload === "string") { + if (frameBytes(payload) > VOICE_RELAY_WS_FRAME_MAX_BYTES) { finish(1009, "message too large"); return; } + if (downstream.getBufferedAmount() + frameBytes(payload) > VOICE_RELAY_WS_BACKPRESSURE_MAX_BYTES) { finish(1013, "client backpressure limit"); return; } + downstream.send(payload); + return; + } + if (payload instanceof ArrayBuffer) { + if (payload.byteLength > VOICE_RELAY_WS_FRAME_MAX_BYTES) { finish(1009, "message too large"); return; } + if (downstream.getBufferedAmount() + payload.byteLength > VOICE_RELAY_WS_BACKPRESSURE_MAX_BYTES) { finish(1013, "client backpressure limit"); return; } + downstream.send(payload); + return; + } + finish(1003, "unsupported frame"); + }); + upstream.addEventListener("close", event => finish(safeCloseCode(event.code), safeCloseReason(event.reason))); + upstream.addEventListener("error", () => finish(1011, "upstream websocket error"), { once: true }); +} + +export function startVoiceRelay(options: VoiceRelayOptions = {}): VoiceRelayHandle { + const credential = options.credential ?? loadVoiceRelayCredential(); + const port = options.port ?? VOICE_RELAY_DEFAULT_PORT; + if (!Number.isInteger(port) || port < 0 || port > 65_535) throw new Error("voice relay port must be between 0 and 65535"); + const sessions = new Set>(); + let resolveDone!: (reason: "stopped" | "connection_changed") => void; + const done = new Promise<"stopped" | "connection_changed">(resolve => { resolveDone = resolve; }); + let stopped = false; + let monitor: ReturnType | undefined; + let server!: Server; + const stop = (reason: "stopped" | "connection_changed" = "stopped") => { + if (stopped) return; + stopped = true; + if (monitor) clearInterval(monitor); + for (const session of sessions) { + try { session.close(1001, reason); } catch { /* best effort */ } + try { (session.data.upstream as (WebSocket & { terminate(): void }) | undefined)?.terminate(); } catch { /* best effort */ } + if (session.data.handshakeTimer) clearTimeout(session.data.handshakeTimer); + if (session.data.closeTimer) clearTimeout(session.data.closeTimer); + session.data.pending = []; + session.data.pendingBytes = 0; + } + sessions.clear(); + try { server.stop(true); } catch { /* already stopped */ } + resolveDone(reason); + }; + server = Bun.serve({ + port, + hostname: "127.0.0.1", + async fetch(req, bunServer) { + if (!localAuthorityAllowed(req, bunServer.port ?? port)) return jsonError(403, "voice_relay_origin_rejected"); + if (!(options.connectionCheck ?? credentialStillOwned)(credential)) { + queueMicrotask(() => stop("connection_changed")); + return jsonError(409, "voice_relay_connection_changed"); + } + const url = new URL(req.url); + const upgrade = req.headers.get("upgrade")?.toLowerCase() === "websocket"; + const wsTarget = upgrade && req.method === "GET" + ? parseLiveSidebandTarget(url.pathname, url.searchParams, url.search.slice(1)) + : null; + const wsAllowed = wsTarget && (options.allowStandalone + || wsTarget.style === "frameless-path" + || wsTarget.style === "realtime-calls-path" + || wsTarget.style === "realtime-query"); + if (wsAllowed && wsTarget) { + const headers = relayHeaders(req.headers, WS_REQUEST_HEADERS); + headers.set("x-opencodex-api-key", credential.token); + if (bunServer.upgrade(req, { data: { + upstreamUrl: upstreamUrl(credential.connection.serverUrl, new URL(sanitizedWebSocketPath(url, wsTarget), url), true), + headers: Object.fromEntries(headers.entries()), + pending: [], pendingBytes: 0, opened: false, closing: false, + } })) return undefined; + return jsonError(426, "voice_relay_upgrade_failed"); + } + if (postRouteAllowed(url, req.method)) return relayHttp(req, url, credential, options); + return jsonError(404, "voice_relay_route_not_found"); + }, + websocket: { + maxPayloadLength: VOICE_RELAY_WS_FRAME_MAX_BYTES, + idleTimeout: VOICE_RELAY_WS_IDLE_SECONDS, + open(ws) { sessions.add(ws); startWebSocketPeer(ws, options); }, + message(ws, frame) { + const data = ws.data; + if (data.closing) return; + const bytes = frameBytes(frame); + if (bytes > VOICE_RELAY_WS_FRAME_MAX_BYTES) { ws.close(1009, "message too large"); return; } + const upstream = data.upstream; + if (!upstream || !data.opened || upstream.readyState !== WebSocket.OPEN) { + if (data.pending.length >= VOICE_RELAY_WS_PENDING_MAX_FRAMES || data.pendingBytes + bytes > VOICE_RELAY_WS_PENDING_MAX_BYTES) { + ws.close(1009, "pending messages exceeded limit"); + try { upstream?.close(1009, "pending messages exceeded limit"); } catch { /* best effort */ } + data.closing = true; + return; + } + data.pending.push(frame); + data.pendingBytes += bytes; + return; + } + if (upstream.bufferedAmount + bytes > VOICE_RELAY_WS_BACKPRESSURE_MAX_BYTES) { + data.closing = true; + ws.close(1013, "upstream backpressure limit"); + try { upstream.close(1013, "upstream backpressure limit"); } catch { /* best effort */ } + return; + } + upstream.send(typeof frame === "string" ? frame : Uint8Array.from(frame)); + }, + close(ws, code, reason) { + sessions.delete(ws); + const data = ws.data; + data.closing = true; + data.pending = []; + data.pendingBytes = 0; + if (data.handshakeTimer) clearTimeout(data.handshakeTimer); + const upstream = data.upstream; + try { if (upstream && upstream.readyState < WebSocket.CLOSING) upstream.close(safeCloseCode(code), safeCloseReason(reason)); } catch { /* best effort */ } + if (upstream && upstream.readyState !== WebSocket.CLOSED && !data.closeTimer) { + data.closeTimer = setTimeout(() => { + data.closeTimer = undefined; + try { + if (upstream.readyState !== WebSocket.CLOSED) { + (upstream as WebSocket & { terminate(): void }).terminate(); + } + } catch { /* best effort */ } + }, options.closeFallbackMs ?? VOICE_RELAY_CLOSE_FALLBACK_MS); + } + }, + }, + }); + if (!options.credential || options.connectionCheck) { + monitor = setInterval(() => { + if (!(options.connectionCheck ?? credentialStillOwned)(credential)) stop("connection_changed"); + }, options.monitorIntervalMs ?? 1_000); + if (typeof monitor === "object" && "unref" in monitor) monitor.unref(); + } + const boundPort = server.port ?? port; + return { port: boundPort, origin: `http://127.0.0.1:${boundPort}`, done, stop: () => stop() }; +} diff --git a/tests/clients/client-voice-relay.test.ts b/tests/clients/client-voice-relay.test.ts new file mode 100644 index 0000000000..abdc3d0722 --- /dev/null +++ b/tests/clients/client-voice-relay.test.ts @@ -0,0 +1,265 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import type { Server } from "bun"; +import { mkdtempSync } from "node:fs"; +import { createHash } from "node:crypto"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { saveConfig } from "../../src/config"; +import { writeServiceApiTokenFile } from "../../src/lib/service-secrets"; +import { + loadVoiceRelayCredential, + startVoiceRelay, + VOICE_RELAY_BODY_MAX_BYTES, + voiceRelayWebSocketRouteAllowed, + type VoiceRelayCredential, +} from "../../src/client/voice-relay"; +import type { OcxClientConnectionConfig } from "../../src/types"; +import { SERVER_BUDGET_MS } from "../helpers/test-budget"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; + +const servers: Array<{ stop(force?: boolean): void | Promise }> = []; +const previousHome = process.env.OPENCODEX_HOME; + +afterEach(async () => { + for (const server of servers.splice(0).reverse()) await server.stop(true); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; +}); + +function connected(serverUrl: string): VoiceRelayCredential { + const connection: OcxClientConnectionConfig = { + serverUrl, + managementUrl: serverUrl, + managementTransport: "direct", + selectedClients: ["codex"], + tokenEnv: "OPENCODEX_API_AUTH_TOKEN", + apiKeyId: "fixture-key", + tokenFingerprint: "fixture-fingerprint", + protocolVersion: 1, + connectedAt: "2026-09-08T00:00:00.000Z", + catalogFingerprint: "fixture-catalog", + priorCatalog: "", + catalogSyncedAt: "2026-09-08T00:00:00.000Z", + }; + return { connection, token: "ocx_data_fixture-secret" }; +} + +function ws(url: string, headers: Record = {}): WebSocket { + return new WebSocket(url, { headers } as unknown as string[]); +} + +function opened(socket: WebSocket): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("websocket open timeout")), 5_000); + socket.addEventListener("open", () => { clearTimeout(timer); resolve(); }, { once: true }); + socket.addEventListener("error", () => { clearTimeout(timer); reject(new Error("websocket open failed")); }, { once: true }); + }); +} + +function message(socket: WebSocket): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("websocket message timeout")), 5_000); + socket.addEventListener("message", event => { + clearTimeout(timer); + resolve(typeof event.data === "string" ? event.data : new TextDecoder().decode(event.data as ArrayBuffer)); + }, { once: true }); + }); +} + +function rawMessage(socket: WebSocket): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("websocket message timeout")), 5_000); + socket.addEventListener("message", event => { clearTimeout(timer); resolve(event.data as string | ArrayBuffer); }, { once: true }); + }); +} + +describe("remote hub voice relay", () => { + test("parser admits exact call routes and gates standalone sessions", () => { + for (const value of [ + "http://127.0.0.1:10111/v1/live/call_1", + "http://127.0.0.1:10111/v1/realtime/calls/call-2", + "http://127.0.0.1:10111/v1/realtime?call_id=call_3", + ]) expect(voiceRelayWebSocketRouteAllowed(new URL(value))).toBe(true); + for (const value of [ + "http://127.0.0.1:10111/v1/live", + "http://127.0.0.1:10111/v1/realtime?model=gpt-realtime", + "http://127.0.0.1:10111/v1/live/call/extra", + "http://127.0.0.1:10111/v1/realtime/calls", + "http://127.0.0.1:10111/v1/liveevil/call_1", + ]) expect(voiceRelayWebSocketRouteAllowed(new URL(value))).toBe(false); + expect(voiceRelayWebSocketRouteAllowed(new URL("http://127.0.0.1:10111/v1/live?model=gpt"), true)).toBe(true); + expect(voiceRelayWebSocketRouteAllowed(new URL("http://127.0.0.1:10111/v1/realtime?model=gpt"), true)).toBe(true); + }); + + test("fake hub receives exact HTTP route, body, protocol header, and connected data credential", async () => { + const seen: Array<{ path: string; method: string; key: string | null; authorization: string | null; body: string }> = []; + const hub = Bun.serve({ + port: 0, + async fetch(req) { + const url = new URL(req.url); + seen.push({ path: `${url.pathname}${url.search}`, method: req.method, key: req.headers.get("x-opencodex-api-key"), authorization: req.headers.get("authorization"), body: await req.text() }); + return new Response("answer", { status: 201, headers: { "Content-Type": "application/sdp", "Set-Cookie": "private=1" } }); + }, + }); + servers.push(hub); + const relay = startVoiceRelay({ port: 0, credential: connected(hub.url.origin), connectionCheck: () => true }); + servers.push({ stop: () => relay.stop() }); + const response = await fetch(`${relay.origin}/v1/realtime/calls?intent=quicksilver`, { + method: "POST", + headers: { + Host: `127.0.0.1:${relay.port}`, + Origin: relay.origin, + Authorization: "Bearer caller-secret", + "X-OpenCodex-API-Key": "caller-admission", + "OpenAI-Alpha": "quicksilver=v2", + "Content-Type": "application/sdp", + }, + body: "offer", + }); + expect(response.status).toBe(201); + expect(await response.text()).toBe("answer"); + expect(response.headers.get("set-cookie")).toBeNull(); + expect(seen).toEqual([{ + path: "/v1/realtime/calls?intent=quicksilver", method: "POST", + key: "ocx_data_fixture-secret", authorization: "Bearer caller-secret", body: "offer", + }]); + }, SERVER_BUDGET_MS); + + test("wrong methods, prefix variants, Host, and browser Origin fail before hub I/O", async () => { + let calls = 0; + const relay = startVoiceRelay({ + port: 0, + credential: connected("https://hub.example.test"), + connectionCheck: () => true, + fetchImpl: (async () => { calls += 1; return new Response(); }) as typeof fetch, + }); + servers.push({ stop: () => relay.stop() }); + const request = (path: string, init: RequestInit) => fetch(`${relay.origin}${path}`, init); + expect((await request("/v1/live", { method: "GET" })).status).toBe(404); + expect((await request("/v1/liveevil", { method: "POST", body: "x" })).status).toBe(404); + expect((await request("/v1/realtime/calls/id", { method: "POST", body: "x" })).status).toBe(404); + expect((await request("/v1/live", { method: "POST", headers: { Host: "evil.example" }, body: "x" })).status).toBe(403); + expect((await request("/v1/live", { method: "POST", headers: { Origin: "https://evil.example" }, body: "x" })).status).toBe(403); + expect(calls).toBe(0); + }, SERVER_BUDGET_MS); + + test("request limit, redirect refusal, connection deadline, and ownership drift are bounded", async () => { + let calls = 0; + const relay = startVoiceRelay({ + port: 0, + credential: connected("https://hub.example.test"), + connectionCheck: () => true, + connectTimeoutMs: 10, + fetchImpl: (async (_input, init) => { + calls += 1; + if (calls === 1) return new Response(null, { status: 302, headers: { Location: "https://evil.example" } }); + return await new Promise((_resolve, reject) => init?.signal?.addEventListener("abort", () => reject(init.signal?.reason), { once: true })); + }) as typeof fetch, + }); + servers.push({ stop: () => relay.stop() }); + const tooLarge = await fetch(`${relay.origin}/v1/live`, { + method: "POST", + body: new Uint8Array(VOICE_RELAY_BODY_MAX_BYTES + 1), + }); + expect(tooLarge.status).toBe(413); + expect(calls).toBe(0); + expect((await fetch(`${relay.origin}/v1/live`, { method: "POST", body: "x" })).status).toBe(502); + expect((await fetch(`${relay.origin}/v1/live`, { method: "POST", body: "x" })).status).toBe(504); + + let owned = true; + const drift = startVoiceRelay({ port: 0, credential: connected("https://hub.example.test"), connectionCheck: () => owned, monitorIntervalMs: 5 }); + servers.push({ stop: () => drift.stop() }); + owned = false; + expect(await drift.done).toBe("connection_changed"); + }, SERVER_BUDGET_MS); + + test("fake hub WebSocket gets credential, relays frames, rejects Origin, and closes peer", async () => { + let hubKey: string | null = null; + let hubAuthorization: string | null = null; + let hubPath = ""; + let hubClosed = false; + const hub = Bun.serve({ + port: 0, + fetch(req, server) { + hubKey = req.headers.get("x-opencodex-api-key"); + hubAuthorization = req.headers.get("authorization"); + const url = new URL(req.url); + hubPath = `${url.pathname}${url.search}`; + if (server.upgrade(req, { data: {} })) return; + return new Response("upgrade failed", { status: 426 }); + }, + websocket: { + open(socket) { socket.send("hub-ready"); }, + message(socket, value) { + if (typeof value === "string") socket.send(`echo:${value}`); + else socket.send(value); + }, + close() { hubClosed = true; }, + }, + }); + servers.push(hub); + const relay = startVoiceRelay({ port: 0, credential: connected(hub.url.origin), connectionCheck: () => true }); + servers.push({ stop: () => relay.stop() }); + const target = relay.origin.replace(/^http/, "ws") + "/v1/live/call_1?token=drop-me&extension=keep"; + const rejected = ws(target, { Origin: "https://evil.example" }); + await expect(opened(rejected)).rejects.toThrow(); + const socket = ws(target, { Origin: relay.origin, Authorization: "Bearer caller-secret" }); + socket.binaryType = "arraybuffer"; + await opened(socket); + expect(await message(socket)).toBe("hub-ready"); + socket.send("hello"); + expect(await message(socket)).toBe("echo:hello"); + socket.send(Uint8Array.from([0, 127, 128, 255])); + expect(Array.from(new Uint8Array(await rawMessage(socket) as ArrayBuffer))).toEqual([0, 127, 128, 255]); + expect(hubKey).toBe("ocx_data_fixture-secret"); + expect(hubAuthorization).toBe("Bearer caller-secret"); + expect(hubPath).toBe("/v1/live/call_1?extension=keep"); + socket.close(); + for (let i = 0; i < 50 && !hubClosed; i += 1) await Bun.sleep(10); + expect(hubClosed).toBe(true); + }, SERVER_BUDGET_MS); + + test("WebSocket handshake timeout closes a connecting upstream peer", async () => { + class PendingSocket extends EventTarget { + readyState = WebSocket.CONNECTING; + closes = 0; + terminates = 0; + send(): void {} + close(): void { this.closes += 1; this.readyState = WebSocket.CLOSING; } + terminate(): void { this.terminates += 1; this.readyState = WebSocket.CLOSED; } + } + const pending = new PendingSocket(); + const relay = startVoiceRelay({ + port: 0, + credential: connected("https://hub.example.test"), + connectionCheck: () => true, + connectTimeoutMs: 5, + closeFallbackMs: 5, + webSocketFactory: () => pending as unknown as WebSocket, + }); + servers.push({ stop: () => relay.stop() }); + const socket = ws(relay.origin.replace(/^http/, "ws") + "/v1/live/call_1"); + await opened(socket); + await new Promise(resolve => socket.addEventListener("close", () => resolve(), { once: true })); + for (let i = 0; i < 50 && pending.terminates === 0; i += 1) await Bun.sleep(2); + expect(pending.closes).toBeGreaterThan(0); + expect(pending.terminates).toBeGreaterThan(0); + }, SERVER_BUDGET_MS); + + test("default credential loader fails closed for missing and mismatched owner state", () => { + const home = mkdtempSync(join(tmpdir(), "ocx-voice-relay-credential-")); + process.env.OPENCODEX_HOME = home; + try { + expect(() => loadVoiceRelayCredential()).toThrow("requires a complete"); + const fixture = connected("https://hub.example.test").connection; + fixture.tokenFingerprint = createHash("sha256").update("ocx_data_expected-owner").digest("hex"); + fixture.catalogFingerprint = createHash("sha256").update("fixture-catalog").digest("base64url"); + saveConfig({ port: 10100, providers: {}, defaultProvider: "openai", runtimeRole: "client", client: fixture }); + expect(() => loadVoiceRelayCredential()).toThrow("missing or no longer owned"); + writeServiceApiTokenFile("ocx_data_wrong-owner"); + expect(() => loadVoiceRelayCredential()).toThrow("missing or no longer owned"); + } finally { + removeTreeWithRetry(home); + } + }); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index da6be012bc..b537bd2ba4 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -198,6 +198,7 @@ "client-export-modality-enum.test.ts": "clients", "client-fingerprint.test.ts": "clients", "client-hub-relay.test.ts": "clients", + "client-voice-relay.test.ts": "clients", "client-injection-guard.test.ts": "codex-integration", "client-lifecycle-lock.test.ts": "clients", "client-machine-listener.test.ts": "clients", From bbee23d613b6e6131555d5b786364994afdc5739 Mon Sep 17 00:00:00 2001 From: Kosta Milovanovic Date: Tue, 8 Sep 2026 16:14:16 -0400 Subject: [PATCH 2/2] fix(client): strip hub admission bearer in voice relay --- .../src/content/docs/guides/remote-hub.md | 4 + src/client/voice-relay.ts | 21 +++ src/server/auth-cors.ts | 7 +- tests/clients/client-voice-relay.test.ts | 129 +++++++++++++++++- 4 files changed, 157 insertions(+), 4 deletions(-) diff --git a/docs-site/src/content/docs/guides/remote-hub.md b/docs-site/src/content/docs/guides/remote-hub.md index 7ba842f0a3..cada923d4e 100644 --- a/docs-site/src/content/docs/guides/remote-hub.md +++ b/docs-site/src/content/docs/guides/remote-hub.md @@ -314,6 +314,10 @@ by a call-ID WebSocket join. Add `--allow-standalone` for a client that opens a without first creating a WebRTC call. Unrelated API routes and browser-origin requests are not general-purpose forwarding surfaces. +The listener trusts local processes on the client machine. While it is running, another local +process can use the allowed voice routes through the connected hub credential, although the +credential itself is never returned to the caller. Stop the relay when remote voice is not in use. + The relay exits if its saved connection or credential changes. After disconnect or key rotation, restore the voice settings or restart the relay against the intended connection. It never repairs pairing or rotates keys itself. diff --git a/src/client/voice-relay.ts b/src/client/voice-relay.ts index 45a9d9f7b9..d7b6463da3 100644 --- a/src/client/voice-relay.ts +++ b/src/client/voice-relay.ts @@ -1,4 +1,5 @@ import type { Server, ServerWebSocket } from "bun"; +import { timingSafeEqual } from "node:crypto"; import { readBoundedResponseBytes } from "../lib/bounded-body"; import { clearableDeadline } from "../lib/abort"; import { @@ -7,6 +8,7 @@ import { sanitizeStandaloneRealtimeQuery, type LiveSidebandTarget, } from "../server/live"; +import { hasProxyAdmissionSecretShape } from "../server/auth-cors"; import type { OcxClientConnectionConfig } from "../types"; import { readServiceApiTokenState } from "../lib/service-secrets"; import { assertNoClientDisconnectPending, readClientConnectionState } from "./state"; @@ -86,6 +88,23 @@ function relayHeaders(source: Headers, names: readonly string[]): Headers { return output; } +function secretEquals(actual: string, expected: string): boolean { + const actualBytes = Buffer.from(actual); + const expectedBytes = Buffer.from(expected); + return actualBytes.byteLength === expectedBytes.byteLength + && timingSafeEqual(actualBytes, expectedBytes); +} + +function removeRelayAdmissionBearer(headers: Headers, connectedToken: string): void { + const match = headers.get("authorization")?.match(/^Bearer\s+(.+)$/i); + const bearer = match?.[1]?.trim(); + if (!bearer) return; + if (secretEquals(bearer, connectedToken) + || hasProxyAdmissionSecretShape(bearer)) { + headers.delete("authorization"); + } +} + function localAuthorityAllowed(req: Request, port: number): boolean { let url: URL; try { url = new URL(req.url); } catch { return false; } @@ -191,6 +210,7 @@ async function relayHttp(req: Request, url: URL, credential: VoiceRelayCredentia if (body instanceof Response) return body; if (!(options.connectionCheck ?? credentialStillOwned)(credential)) return jsonError(409, "voice_relay_connection_changed"); const headers = relayHeaders(req.headers, HTTP_REQUEST_HEADERS); + removeRelayAdmissionBearer(headers, credential.token); headers.set("x-opencodex-api-key", credential.token); const connect = clearableDeadline(options.connectTimeoutMs ?? VOICE_RELAY_CONNECT_TIMEOUT_MS, lifetime); let response: Response; @@ -354,6 +374,7 @@ export function startVoiceRelay(options: VoiceRelayOptions = {}): VoiceRelayHand || wsTarget.style === "realtime-query"); if (wsAllowed && wsTarget) { const headers = relayHeaders(req.headers, WS_REQUEST_HEADERS); + removeRelayAdmissionBearer(headers, credential.token); headers.set("x-opencodex-api-key", credential.token); if (bunServer.upgrade(req, { data: { upstreamUrl: upstreamUrl(credential.connection.serverUrl, new URL(sanitizedWebSocketPath(url, wsTarget), url), true), diff --git a/src/server/auth-cors.ts b/src/server/auth-cors.ts index 476fd3a4ae..86417fafd0 100644 --- a/src/server/auth-cors.ts +++ b/src/server/auth-cors.ts @@ -446,11 +446,16 @@ export function isManagementAdmissionSecret(token: string): boolean { return !!actual && secretEquals(actual, configuredAdminAuthToken()); } +/** Whether `token` has a minted OpenCodex admission-secret shape. */ +export function hasProxyAdmissionSecretShape(token: string): boolean { + return /^ocx_(?:data|admin|session)_/.test(token) || /^ocx_[0-9a-f]{40}$/.test(token); +} + /** Whether `token` is one of the proxy's own admission secrets and must never reach an upstream. */ export function isProxyAdmissionSecret(token: string, config: OcxConfig): boolean { const actual = token.trim(); if (!actual) return false; - if (/^ocx_(?:data|admin|session)_/.test(actual) || /^ocx_[0-9a-f]{40}$/.test(actual)) return true; + if (hasProxyAdmissionSecretShape(actual)) return true; return isDataPlaneAdmissionSecret(actual, config) || isManagementAdmissionSecret(actual); } diff --git a/tests/clients/client-voice-relay.test.ts b/tests/clients/client-voice-relay.test.ts index abdc3d0722..12008f7e68 100644 --- a/tests/clients/client-voice-relay.test.ts +++ b/tests/clients/client-voice-relay.test.ts @@ -4,8 +4,9 @@ import { mkdtempSync } from "node:fs"; import { createHash } from "node:crypto"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { saveConfig } from "../../src/config"; +import { getDefaultConfig, saveConfig } from "../../src/config"; import { writeServiceApiTokenFile } from "../../src/lib/service-secrets"; +import { ForwardAdmissionCredentialError, validateForwardAdmissionCredential } from "../../src/server/auth-cors"; import { loadVoiceRelayCredential, startVoiceRelay, @@ -19,6 +20,7 @@ import { removeTreeWithRetry } from "../helpers/remove-tree"; const servers: Array<{ stop(force?: boolean): void | Promise }> = []; const previousHome = process.env.OPENCODEX_HOME; +const NATIVE_OAUTH_BEARER = ["Bearer", "native-oauth-access-token"].join(" "); afterEach(async () => { for (const server of servers.splice(0).reverse()) await server.stop(true); @@ -26,7 +28,7 @@ afterEach(async () => { else process.env.OPENCODEX_HOME = previousHome; }); -function connected(serverUrl: string): VoiceRelayCredential { +function connected(serverUrl: string, token = "ocx_data_fixture-secret"): VoiceRelayCredential { const connection: OcxClientConnectionConfig = { serverUrl, managementUrl: serverUrl, @@ -41,7 +43,7 @@ function connected(serverUrl: string): VoiceRelayCredential { priorCatalog: "", catalogSyncedAt: "2026-09-08T00:00:00.000Z", }; - return { connection, token: "ocx_data_fixture-secret" }; + return { connection, token }; } function ws(url: string, headers: Record = {}): WebSocket { @@ -73,6 +75,18 @@ function rawMessage(socket: WebSocket): Promise { }); } +function liveForwardingGuard(headers: Headers, config: ReturnType): Response | null { + try { + validateForwardAdmissionCredential(headers, config); + return null; + } catch (error) { + if (error instanceof ForwardAdmissionCredentialError) { + return Response.json({ error: { type: "authentication_error", message: error.message } }, { status: 401 }); + } + throw error; + } +} + describe("remote hub voice relay", () => { test("parser admits exact call routes and gates standalone sessions", () => { for (const value of [ @@ -125,6 +139,59 @@ describe("remote hub voice relay", () => { }]); }, SERVER_BUDGET_MS); + test("HTTP removes hub admission bearers before the real forwarding guard", async () => { + const connectedToken = "legacy-connected-token"; + const config = getDefaultConfig(); + config.apiKeys = [{ id: "voice-relay", name: "voice relay", key: connectedToken }]; + const seen: Array<{ authorization: string | null; key: string | null; account: string | null }> = []; + const hub = Bun.serve({ + port: 0, + fetch(req) { + const rejected = liveForwardingGuard(req.headers, config); + if (rejected) return rejected; + seen.push({ + authorization: req.headers.get("authorization"), + key: req.headers.get("x-opencodex-api-key"), + account: req.headers.get("chatgpt-account-id"), + }); + return new Response("answer", { status: 201 }); + }, + }); + servers.push(hub); + const relay = startVoiceRelay({ + port: 0, + credential: connected(hub.url.origin, connectedToken), + connectionCheck: () => true, + }); + servers.push({ stop: () => relay.stop() }); + + const admissionBearers = [ + connectedToken, + `ocx_data_${"a".repeat(40)}`, + `ocx_admin_${"b".repeat(40)}`, + `ocx_session_${"c".repeat(40)}`, + `ocx_${"d".repeat(40)}`, + ]; + for (const bearer of admissionBearers) { + const response = await fetch(`${relay.origin}/v1/live`, { + method: "POST", + headers: { Authorization: `Bearer ${bearer}` }, + body: "offer", + }); + expect(response.status).toBe(201); + } + const oauth = await fetch(`${relay.origin}/v1/live`, { + method: "POST", + headers: { Authorization: NATIVE_OAUTH_BEARER, "ChatGPT-Account-ID": "acct-1" }, + body: "offer", + }); + expect(oauth.status).toBe(201); + expect(seen).toEqual([ + ...admissionBearers.map(() => ({ authorization: null, key: connectedToken, account: null })), + { authorization: NATIVE_OAUTH_BEARER, key: connectedToken, account: "acct-1" }, + ]); + }, SERVER_BUDGET_MS); + test("wrong methods, prefix variants, Host, and browser Origin fail before hub I/O", async () => { let calls = 0; const relay = startVoiceRelay({ @@ -219,6 +286,62 @@ describe("remote hub voice relay", () => { expect(hubClosed).toBe(true); }, SERVER_BUDGET_MS); + test("WebSocket removes hub admission bearers before the real forwarding guard", async () => { + const connectedToken = "legacy-connected-token"; + const config = getDefaultConfig(); + config.apiKeys = [{ id: "voice-relay", name: "voice relay", key: connectedToken }]; + const seen: Array<{ authorization: string | null; key: string | null; account: string | null }> = []; + const hub = Bun.serve({ + port: 0, + fetch(req, server) { + const rejected = liveForwardingGuard(req.headers, config); + if (rejected) return rejected; + seen.push({ + authorization: req.headers.get("authorization"), + key: req.headers.get("x-opencodex-api-key"), + account: req.headers.get("chatgpt-account-id"), + }); + if (server.upgrade(req, { data: {} })) return; + return new Response("upgrade failed", { status: 426 }); + }, + websocket: { open(socket) { socket.send("ready"); } }, + }); + servers.push(hub); + const relay = startVoiceRelay({ + port: 0, + credential: connected(hub.url.origin, connectedToken), + connectionCheck: () => true, + }); + servers.push({ stop: () => relay.stop() }); + + const admissionBearers = [ + connectedToken, + `ocx_data_${"a".repeat(40)}`, + `ocx_admin_${"b".repeat(40)}`, + `ocx_session_${"c".repeat(40)}`, + `ocx_${"d".repeat(40)}`, + ]; + for (const bearer of admissionBearers) { + const socket = ws(`${relay.origin.replace(/^http/, "ws")}/v1/live/call_1`, { + Authorization: `Bearer ${bearer}`, + }); + await opened(socket); + expect(await message(socket)).toBe("ready"); + socket.close(); + } + const oauth = ws(`${relay.origin.replace(/^http/, "ws")}/v1/live/call_1`, { + Authorization: NATIVE_OAUTH_BEARER, + "ChatGPT-Account-ID": "acct-1", + }); + await opened(oauth); + expect(await message(oauth)).toBe("ready"); + oauth.close(); + expect(seen).toEqual([ + ...admissionBearers.map(() => ({ authorization: null, key: connectedToken, account: null })), + { authorization: NATIVE_OAUTH_BEARER, key: connectedToken, account: "acct-1" }, + ]); + }, SERVER_BUDGET_MS); + test("WebSocket handshake timeout closes a connecting upstream peer", async () => { class PendingSocket extends EventTarget { readyState = WebSocket.CONNECTING;