fix(talk): remove unsupported Codex OAuth realtime

This commit is contained in:
Vincent Koc
2026-07-28 21:32:21 +08:00
parent 455e728533
commit 861b9ccce1
11 changed files with 51 additions and 1903 deletions
+3 -2
View File
@@ -51,6 +51,7 @@ Docs: https://docs.openclaw.ai
### Fixes
- **OpenAI Realtime Talk auth:** remove the non-public Codex OAuth realtime fallback and require an OpenAI Platform API key for Talk, Voice Call, and Discord realtime voice, preventing OAuth-only gateways from advertising a browser session that the live service rejects. Fixes #115021.
- **Codex native controls:** stop misclassifying valid thinking/fast runtime controls as provider overrides so Codex routes keep their native controls, while provider-native objects and invalid values stay fail-closed. Thanks @VACInc. (#107588)
- **State snapshot verification:** run SQLite snapshot verification in a separate process so worker-thread file closes no longer drop the Gateway's POSIX WAL locks, eliminating spurious WAL misses and I/O errors. Thanks @VACInc. (#114016)
- **Reply latency with model policies:** reuse one immutable plugin-metadata snapshot per model-selection run instead of repeating plugin discovery, cutting reply delay when a model policy is configured. Thanks @VACInc. (#114117)
@@ -1173,7 +1174,7 @@ The [model catalog](https://docs.openclaw.ai/concepts/models) also reports avail
- Cloudflare 403 challenges on OpenAI or Codex OAuth requests now produce gateway-block guidance instead of incorrectly telling users their authentication failed. [#94440](https://github.com/openclaw/openclaw/pull/94440) Related [#94432](https://github.com/openclaw/openclaw/issues/94432). Thanks @lzyyzznl, @pbm9z95m6z-hue.
- OpenAI authentication errors now point ChatGPT/Codex OAuth users to a model compatible with their existing sign-in instead of recommending an outdated default. [#100579](https://github.com/openclaw/openclaw/pull/100579) Thanks @zhangguiping-xydt.
- For Codex-backed OpenAI models, `/status` now identifies ChatGPT login authentication as `oauth (codex-cli)` instead of incorrectly labeling it as an environment API key. [#91240](https://github.com/openclaw/openclaw/pull/91240) Related [#91099](https://github.com/openclaw/openclaw/issues/91099). Thanks @849261680, @ukstem.
- OpenAI Realtime voice in Talk, Voice Call, and Discord can now use an existing Codex/OpenAI OAuth login when no explicit API key is configured. [#100671](https://github.com/openclaw/openclaw/pull/100671) Thanks @steipete-oai.
- OpenAI Realtime voice in Talk, Voice Call, and Discord was announced with a Codex/OpenAI OAuth fallback in [#100671](https://github.com/openclaw/openclaw/pull/100671). Correction: public Codex OAuth accounts do not have a supported realtime transport, so current builds require an OpenAI Platform API key. Thanks @steipete-oai.
##### Google and Gemini
@@ -6877,7 +6878,7 @@ This audited record covers the complete v2026.5.28..v2026.5.31-beta.4 history: 4
- Agents/compaction: keep contributor diagnostics to a bounded top-three selection without sorting the full history. Thanks @shakkernerd.
- Sessions/UI: avoid full-array sorting while selecting ACPX leases, Google Meet calendar events, and latest chat sessions. Thanks @shakkernerd.
- Plugin SDK: mark direct `deliverOutboundPayloads` and legacy reply-dispatch bridges as deprecated compatibility substrate, enrich `sendDurableMessageBatch` with explicit durable send outcomes, migrate bundled send/turn paths off deprecated APIs, and enforce the split with `check:deprecated-api-usage`.
- OpenAI/Talk: let browser realtime Talk, Gateway relay/Voice Call realtime bridges, and OpenAI realtime transcription use `openai-codex` OAuth when no direct API key is configured, make Google Meet `test_speech` honor `mode: "bidi"`, expose Control UI launch options for provider/model/voice/transport/VAD/reasoning, and update the default OpenAI realtime voice model to `gpt-realtime-2`. Thanks @Solvely-Colin.
- OpenAI/Talk: add browser realtime Talk controls, Google Meet `test_speech` support for `mode: "bidi"`, and the `gpt-realtime-2` default. Correction: the announced `openai-codex` OAuth fallback does not have a supported public realtime transport; Talk, Gateway relay/Voice Call, and realtime transcription require OpenAI Platform credentials. Thanks @Solvely-Colin.
- Telegram: preserve the channel-specific 10-option poll cap in the unified outbound adapter so over-limit polls are rejected before send. (#78762) Thanks @obviyus.
- Telegram/streaming: continue over-limit draft previews in a new message instead of stopping when rendered preview text crosses Telegram's message limit. (#74508) Thanks @anagnorisis2peripeteia.
- Slack: route handled top-level channel turns in implicit-conversation channels to thread-scoped sessions when Slack reply threading is enabled, keeping the root turn and later thread replies on one OpenClaw session. (#78522) Thanks @zeroth-blip.
+5 -15
View File
@@ -97,20 +97,10 @@ Supported keys: `voice` / `voice_id` / `voiceId`, `model` / `model_id` / `modelI
OpenAI browser and iOS WebRTC Talk use Platform credentials in this order:
the configured realtime API key, an `openai` API-key profile, then
`OPENAI_API_KEY`. When none is configured and the bundled Codex runtime is
active, Talk falls back to its logged-in ChatGPT/Codex subscription
automatically. OpenAI OAuth/Codex agent sessions activate that runtime without
an additional Talk auth setting. This experimental fallback supports client-owned WebRTC only;
Gateway relay and backend voice bridges still require OpenAI Platform
credentials.
OpenClaw does not read or copy the Codex OAuth token. The Codex app-server owns
the subscription-authenticated realtime connection and starts an ephemeral,
read-only thread seeded with bounded context from the active agent session.
Codex owns the realtime model, base prompt, and agent handoff on this route;
`talk.realtime.model`, direct provider tools, and Video Talk apply only to the
Platform WebRTC route. Custom Talk instructions and bounded profile context are
added as developer context without replacing Codex's native delegation prompt.
`OPENAI_API_KEY`. ChatGPT/Codex OAuth authenticates the subscription Codex
backend, not the public OpenAI Realtime API, and does not configure Talk,
Voice Call, or Discord realtime voice. Configure a Platform API key even when
agent turns use Codex OAuth.
| Key | Default | Notes |
| ---------------------------------------- | ------------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
@@ -132,7 +122,7 @@ added as developer context without replacing Codex's native delegation prompt.
| `realtime.transport` | - | `webrtc`: client-owned OpenAI WebRTC on iOS and in the browser. `provider-websocket`: browser-owned, stays on Gateway relay on iOS. `gateway-relay`: keeps provider audio on the Gateway; Android uses realtime only with this transport. |
| `realtime.brain` | - | `agent-consult` routes realtime tool calls through Gateway policy; `direct-tools` is legacy direct-tool compatibility; `none` is for transcription/external orchestration. |
| `realtime.consultRouting` | - | `provider-direct` preserves the provider's direct reply when it skips `openclaw_agent_consult`; `force-agent-consult` routes finalized user transcripts through OpenClaw instead. |
| `realtime.instructions` | - | Appends provider-facing system instructions to OpenClaw's built-in realtime prompt. On the Codex OAuth fallback, the text is developer context and Codex's native delegation prompt stays authoritative. |
| `realtime.instructions` | - | Appends provider-facing system instructions to OpenClaw's built-in realtime prompt. |
`talk.catalog` exposes canonical provider ids and registry aliases, each provider's valid modes/transports/brain strategies/realtime audio formats/capability flags, and the runtime-selected readiness result. First-party Talk clients should read that catalog instead of maintaining provider aliases locally; treat an older Gateway that omits group readiness as unverified rather than definitively unconfigured. Streaming transcription providers are discovered through `talk.catalog.transcription`; the current Gateway relay uses the Voice Call streaming provider config until a dedicated Talk transcription config surface ships.
+21 -33
View File
@@ -142,20 +142,20 @@ explicit runtime config.
## OpenClaw feature coverage
| OpenAI capability | OpenClaw surface | Status |
| ------------------------- | --------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------ |
| Chat / Responses | `openai/<model>` model provider | Yes |
| Codex subscription models | `openai/<model>` with OpenAI OAuth | Yes |
| Legacy Codex model refs | old Codex model refs, `codex-cli/<model>` | Repaired by doctor to `openai/<model>` |
| Codex app-server harness | Codex-compatible HTTPS route with runtime unset/`auto`, or explicit `agentRuntime.id: codex` | Yes |
| Server-side web search | Native OpenAI Responses tool | Yes, when web search is enabled and no other provider is pinned |
| Images | `image_generate` | Yes |
| Videos | `video_generate` | Yes |
| Text-to-speech | `tts.provider: "openai"` / `tts` | Yes |
| Batch speech-to-text | `tools.media.audio` / media understanding | Yes |
| Streaming speech-to-text | Voice Call `streaming.provider: "openai"` | Yes |
| Realtime voice | Voice Call `realtime.provider: "openai"` / Control UI Talk `talk.realtime.provider: "openai"` | Yes (Platform API key; experimental Codex OAuth for browser WebRTC Talk) |
| Embeddings | memory embedding provider | Yes |
| OpenAI capability | OpenClaw surface | Status |
| ------------------------- | --------------------------------------------------------------------------------------------- | --------------------------------------------------------------- |
| Chat / Responses | `openai/<model>` model provider | Yes |
| Codex subscription models | `openai/<model>` with OpenAI OAuth | Yes |
| Legacy Codex model refs | old Codex model refs, `codex-cli/<model>` | Repaired by doctor to `openai/<model>` |
| Codex app-server harness | Codex-compatible HTTPS route with runtime unset/`auto`, or explicit `agentRuntime.id: codex` | Yes |
| Server-side web search | Native OpenAI Responses tool | Yes, when web search is enabled and no other provider is pinned |
| Images | `image_generate` | Yes |
| Videos | `video_generate` | Yes |
| Text-to-speech | `tts.provider: "openai"` / `tts` | Yes |
| Batch speech-to-text | `tools.media.audio` / media understanding | Yes |
| Streaming speech-to-text | Voice Call `streaming.provider: "openai"` | Yes |
| Realtime voice | Voice Call `realtime.provider: "openai"` / Control UI Talk `talk.realtime.provider: "openai"` | Yes (Platform API key) |
| Embeddings | memory embedding provider | Yes |
<Note>
OpenAI Realtime voice normally goes through the public **OpenAI Platform
@@ -163,19 +163,11 @@ Realtime API** and requires a Platform API key. Codex OAuth tokens authenticate
the ChatGPT Codex backend instead; they are not interchangeable with Platform
API keys for the public Realtime endpoints.
Control UI and iOS WebRTC Talk can instead use the experimental Codex
app-server route automatically when no Platform credential is configured and
the bundled Codex runtime is active. OpenAI OAuth/Codex agent sessions activate
that runtime without an additional Talk auth setting. Platform auth wins in this order:
configured realtime API key, `openai` API-key profile, then `OPENAI_API_KEY`.
Only when all three are absent does browser WebRTC use the Codex plugin's
logged-in subscription. The OAuth token is never exposed to OpenClaw or the
browser. This fallback is limited to client-owned WebRTC; Voice Call and
Gateway-relay realtime still require Platform credentials. Codex owns the
realtime model, base prompt, and native agent delegation on this route.
OpenClaw adds configured Talk instructions and bounded profile context as
developer context without replacing that prompt. Direct Realtime function
tools, VAD/reasoning tuning, and Video Talk remain Platform-only.
Platform auth is resolved in this order: configured realtime API key, `openai`
API-key profile, then `OPENAI_API_KEY`. ChatGPT/Codex OAuth can still
authenticate agent models and other explicitly supported subscription
surfaces, but it does not configure Talk, Voice Call, Discord realtime voice,
or realtime transcription.
If API-key auth reports missing billing, top up Platform credits at
[platform.openai.com/account/billing](https://platform.openai.com/account/billing)
@@ -952,12 +944,8 @@ compatibility fallback when the shared
against the OpenAI Realtime API when using Platform credentials. The
Gateway mints that client secret with the selected `openai` credential.
Configured realtime keys, API-key profiles, and `OPENAI_API_KEY` use that
path in that order. When none exists and the bundled Codex runtime is
active, browser WebRTC falls back to the logged-in Codex app-server
automatically. Gateway relay and Voice Call backend realtime WebSocket
bridges continue to use Platform credentials. The Codex route keeps
Codex's native realtime prompt, model selection, and agent handoff; it does
not accept the direct Platform model/tool/camera controls.
path in that order. ChatGPT/Codex OAuth is not a Platform Realtime
credential and is not used as a fallback.
Maintainer live verification is available with
`OPENAI_API_KEY=... GEMINI_API_KEY=... node --import tsx scripts/dev/realtime-talk-live-smoke.ts`;
the OpenAI legs verify both the backend WebSocket bridge and the browser
+1 -1
View File
@@ -815,7 +815,7 @@ The [model catalog](/concepts/models) also reports availability and capability m
- ChatGPT OAuth sign-in and refresh now reject unexpectedly large token responses cleanly instead of risking Gateway memory exhaustion. [#99479](https://github.com/openclaw/openclaw/pull/99479) Thanks @pandah97.
- OpenAI-compatible providers routed through the GitHub Copilot BYOK harness now retain bearer authentication, reducing false 401 errors for valid credentials. [#99955](https://github.com/openclaw/openclaw/pull/99955) Thanks @hxy91819.
- A detail-less failure from an OpenAI-compatible provider no longer puts a valid API-key profile into cooldown or forces avoidable fallback traffic. [#100617](https://github.com/openclaw/openclaw/pull/100617) Thanks @fengjikui.
- OpenAI Realtime voice in Talk, Voice Call, and Discord can now use an existing Codex/OpenAI OAuth login when no explicit API key is configured. [#100671](https://github.com/openclaw/openclaw/pull/100671) Thanks @steipete-oai.
- OpenAI Realtime voice in Talk, Voice Call, and Discord was announced with a Codex/OpenAI OAuth fallback in [#100671](https://github.com/openclaw/openclaw/pull/100671). Correction: public Codex OAuth accounts do not have a supported realtime transport, so current builds require an OpenAI Platform API key. Thanks @steipete-oai.
- OpenAI image generation can now be enabled from `models.providers.openai` with a usable API key and custom base URL, without also requiring `OPENAI_API_KEY` or an auth profile. [#100745](https://github.com/openclaw/openclaw/pull/100745) Thanks @amknight.
- Codex users can route both `codex/*` and `openai/*` models through the bundled runtime, and older conversations resume without unnecessary context reprojection. [#105034](https://github.com/openclaw/openclaw/pull/105034)
- The bundled Codex plugin can complete backend requests again after updating its managed app-server runtime, with no model-picker or configuration migration required. [#106098](https://github.com/openclaw/openclaw/pull/106098)
-40
View File
@@ -14,10 +14,6 @@ import type { PluginStateSyncKeyedStore } from "openclaw/plugin-sdk/plugin-state
import { registerCodexCliMetadata } from "./cli-metadata.js";
import { createCodexAppServerAgentHarness } from "./harness.js";
import { buildCodexMediaUnderstandingProvider } from "./media-understanding-provider.js";
import {
CODEX_REALTIME_OFFER_PATH,
configureCodexRealtimeBrowserSession,
} from "./realtime-voice-api.js";
import { readCodexPluginConfig } from "./src/app-server/config.js";
import {
CODEX_APP_SERVER_BINDING_MAX_ENTRIES,
@@ -94,42 +90,6 @@ export default definePluginEntry({
return livePluginConfig;
};
const resolveCurrentPluginConfig = () => resolvePluginConfig(resolveCurrentConfig);
if (api.registrationMode === "full") {
const realtimeBrowserSession = configureCodexRealtimeBrowserSession({
getConfig: resolveCurrentConfig,
getPluginConfig: resolveCurrentPluginConfig,
});
api.registerHttpRoute({
path: CODEX_REALTIME_OFFER_PATH,
auth: "plugin",
match: "exact",
handler: realtimeBrowserSession.handler,
});
api.registerService({
id: "codex-oauth-realtime-browser-session-warmup",
start: () => {
void realtimeBrowserSession.warmup().catch((error: unknown) => {
api.logger.debug?.(
`Codex OAuth realtime warmup unavailable: ${
error instanceof Error ? error.message : String(error)
}`,
);
});
},
});
api.lifecycle.registerRuntimeLifecycle({
id: "codex-oauth-realtime-browser-session",
description: "Release Codex OAuth realtime browser sessions when the plugin stops",
cleanup: async ({ reason }) => {
// Session cleanup must not release the process runtime. Registry
// restart and plugin disable release this registration's lease.
if (reason === "reset" || reason === "delete") {
return;
}
await realtimeBrowserSession.cleanup();
},
});
}
let bindingStateStore: PluginStateSyncKeyedStore<StoredCodexAppServerBinding> | undefined;
const openBindingStateStore = () =>
(bindingStateStore ??= api.runtime.state.openSyncKeyedStore<StoredCodexAppServerBinding>({
@@ -1,54 +0,0 @@
import { afterEach, describe, expect, it } from "vitest";
import { configureCodexRealtimeBrowserSession } from "./realtime-voice-api.js";
const leases: Array<ReturnType<typeof configureCodexRealtimeBrowserSession>> = [];
const createRuntime = () => {
const runtime = configureCodexRealtimeBrowserSession({
getConfig: () => undefined,
getPluginConfig: () => undefined,
});
leases.push(runtime);
return runtime;
};
describe("Codex realtime voice runtime artifact", () => {
afterEach(async () => {
await Promise.all(leases.splice(0).map((runtime) => runtime.cleanup()));
});
it("shares one owner runtime across registration leases", async () => {
const first = createRuntime();
const second = createRuntime();
expect(second).not.toBe(first);
expect(second.broker).toBe(first.broker);
await first.cleanup();
await second.cleanup();
});
it("reads config from the newest active registration", () => {
let firstConfigReads = 0;
let secondConfigReads = 0;
const first = configureCodexRealtimeBrowserSession({
getConfig: () => {
firstConfigReads += 1;
return undefined;
},
getPluginConfig: () => undefined,
});
const second = configureCodexRealtimeBrowserSession({
getConfig: () => {
secondConfigReads += 1;
return undefined;
},
getPluginConfig: () => undefined,
});
leases.push(first, second);
expect(second.broker.isConfigured()).toBe(false);
expect(firstConfigReads).toBe(0);
expect(secondConfigReads).toBe(1);
});
});
-78
View File
@@ -1,78 +0,0 @@
/**
* Bundled Codex realtime integration shared with the bundled OpenAI provider.
*
* This artifact keeps the HTTP offer route and the provider fallback on one
* process-owned runtime without adding a general Plugin SDK registration API.
*/
import { createCodexRealtimeBrowserSessionBroker } from "./src/realtime-browser-session.js";
type CodexRealtimeBrowserSessionRuntime = ReturnType<
typeof createCodexRealtimeBrowserSessionBroker
>;
type CodexRealtimeBrowserSessionParams = Parameters<
typeof createCodexRealtimeBrowserSessionBroker
>[0];
type CodexRealtimeGlobalState = {
version: 1;
runtime?: CodexRealtimeBrowserSessionRuntime;
sources: Map<symbol, CodexRealtimeBrowserSessionParams>;
};
const CODEX_REALTIME_GLOBAL_STATE = Symbol.for("openclaw.codex.realtime-voice.v1");
function getGlobalState(): CodexRealtimeGlobalState {
const root = globalThis as typeof globalThis & {
[CODEX_REALTIME_GLOBAL_STATE]?: CodexRealtimeGlobalState;
};
const state = (root[CODEX_REALTIME_GLOBAL_STATE] ??= {
version: 1,
sources: new Map(),
});
state.sources ??= new Map();
return state;
}
export function configureCodexRealtimeBrowserSession(
params: CodexRealtimeBrowserSessionParams,
): CodexRealtimeBrowserSessionRuntime {
const state = getGlobalState();
const leaseId = Symbol("codex-realtime-registration");
state.sources.set(leaseId, params);
if (state.runtime) {
return createRuntimeLease(state, state.runtime, leaseId);
}
const resolveCurrentSource = (): CodexRealtimeBrowserSessionParams | undefined =>
Array.from(state.sources.values()).at(-1);
const created = createCodexRealtimeBrowserSessionBroker({
getConfig: () => resolveCurrentSource()?.getConfig(),
getPluginConfig: () => resolveCurrentSource()?.getPluginConfig(),
});
state.runtime = created;
return createRuntimeLease(state, created, leaseId);
}
function createRuntimeLease(
state: CodexRealtimeGlobalState,
runtime: CodexRealtimeBrowserSessionRuntime,
leaseId: symbol,
): CodexRealtimeBrowserSessionRuntime {
let released = false;
return {
...runtime,
cleanup: async () => {
if (released) {
return;
}
released = true;
state.sources.delete(leaseId);
if (state.runtime !== runtime || state.sources.size > 0) {
return;
}
state.runtime = undefined;
await runtime.cleanup();
},
};
}
export { CODEX_REALTIME_OFFER_PATH } from "./src/realtime-browser-session.js";
@@ -1,641 +0,0 @@
import { EventEmitter } from "node:events";
import type { IncomingMessage, ServerResponse } from "node:http";
import { Readable } from "node:stream";
import { beforeEach, describe, expect, it, vi } from "vitest";
import type { CodexAppServerClient } from "./app-server/client.js";
import type { CodexServerNotification } from "./app-server/protocol.js";
import {
CODEX_REALTIME_OFFER_PATH,
createCodexRealtimeBrowserSessionBroker,
} from "./realtime-browser-session.js";
const sharedClientMocks = vi.hoisted(() => ({
getClient: vi.fn(),
getSharedClient: vi.fn(),
releaseClient: vi.fn(),
}));
vi.mock("./app-server/shared-client.js", () => ({
getLeasedSharedCodexAppServerClient: sharedClientMocks.getClient,
getSharedCodexAppServerClient: sharedClientMocks.getSharedClient,
releaseLeasedSharedCodexAppServerClient: sharedClientMocks.releaseClient,
}));
function createSdpRequest(token: string, origin?: string): IncomingMessage {
return Object.assign(Readable.from(["v=offer\r\n"]), {
method: "POST",
headers: {
authorization: `Bearer ${token}`,
"content-type": "application/sdp",
...(origin ? { origin } : {}),
},
}) as unknown as IncomingMessage;
}
function createPreflightRequest(origin: string): IncomingMessage {
return Object.assign(Readable.from([]), {
method: "OPTIONS",
headers: {
origin,
"access-control-request-method": "POST",
"access-control-request-headers": "authorization,content-type",
"access-control-request-private-network": "true",
},
}) as unknown as IncomingMessage;
}
function createResponseHarness(options: { autoFinish?: boolean } = {}): {
res: ServerResponse;
end: ReturnType<typeof vi.fn>;
setHeader: ReturnType<typeof vi.fn>;
readBody: () => string;
close: () => void;
} {
let body = "";
const end = vi.fn((value?: string) => {
body = value ?? "";
if (options.autoFinish !== false) {
queueMicrotask(() => res.emit("finish"));
}
});
const setHeader = vi.fn();
const res = Object.assign(new EventEmitter(), {
statusCode: 200,
setHeader,
end,
}) as unknown as ServerResponse;
return {
res,
end,
setHeader,
readBody: () => body,
close: () => {
res.emit("close");
},
};
}
function createFakeClient(
options: {
stallRealtimeStart?: boolean;
genericRealtimeStartAbortError?: boolean;
realtimeStartNotifications?: CodexServerNotification[];
} = {},
): {
client: CodexAppServerClient;
methods: string[];
emitClose: () => void;
emitNotification: (notification: CodexServerNotification) => void;
readRealtimeStartSignal: () => AbortSignal | undefined;
} {
let closeHandler: ((client: CodexAppServerClient) => void) | undefined;
let notificationHandler: ((notification: CodexServerNotification) => void) | undefined;
let realtimeStartSignal: AbortSignal | undefined;
const methods: string[] = [];
const client = {
request: vi.fn(
async (method: string, _params?: unknown, requestOptions?: { signal?: AbortSignal }) => {
methods.push(method);
if (method === "thread/start") {
return {
approvalPolicy: "never",
approvalsReviewer: "user",
cwd: "/tmp/workspace",
model: "gpt-5.4",
modelProvider: "openai",
sandbox: { type: "readOnly" },
thread: {
id: "thread-1",
sessionId: "session-1",
cliVersion: "0.145.0",
createdAt: 1,
updatedAt: 1,
cwd: "/tmp/workspace",
ephemeral: true,
modelProvider: "openai",
preview: "",
source: "appServer",
status: { type: "idle" },
turns: [],
},
};
}
if (method === "thread/realtime/start") {
realtimeStartSignal = requestOptions?.signal;
const notifications = options.realtimeStartNotifications ?? [
{
method: "thread/realtime/sdp",
params: { threadId: "thread-1", sdp: "v=answer\r\n" },
},
];
queueMicrotask(() => {
for (const notification of notifications) {
notificationHandler?.(notification);
}
});
if (options.stallRealtimeStart) {
return await new Promise((_, reject) => {
const signal = requestOptions?.signal;
const rejectAbort = () =>
reject(
options.genericRealtimeStartAbortError
? new Error("request cancelled")
: signal?.reason instanceof Error
? signal.reason
: new Error("realtime start aborted"),
);
signal?.addEventListener("abort", rejectAbort, { once: true });
if (signal?.aborted) {
rejectAbort();
}
});
}
return {};
}
if (method === "thread/realtime/stop" || method === "thread/unsubscribe") {
return {};
}
throw new Error(`Unexpected Codex request: ${method}`);
},
),
addNotificationHandler: vi.fn((handler: (notification: CodexServerNotification) => void) => {
notificationHandler = handler;
return () => {
notificationHandler = undefined;
};
}),
addCloseHandler: vi.fn((handler: (client: CodexAppServerClient) => void) => {
closeHandler = handler;
}),
} as unknown as CodexAppServerClient;
return {
client,
methods,
emitClose: () => closeHandler?.(client),
emitNotification: (notification) => notificationHandler?.(notification),
readRealtimeStartSignal: () => realtimeStartSignal,
};
}
function useFakeClient(fake: ReturnType<typeof createFakeClient>): void {
sharedClientMocks.getClient.mockResolvedValue(fake.client);
sharedClientMocks.getSharedClient.mockResolvedValue(fake.client);
}
describe("Codex OAuth realtime browser session", () => {
beforeEach(() => {
sharedClientMocks.getClient.mockReset();
sharedClientMocks.getSharedClient.mockReset();
sharedClientMocks.releaseClient.mockReset();
});
it("advertises the broker only after subscription warmup succeeds", async () => {
const fake = createFakeClient();
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
expect(realtime.broker.isConfigured()).toBe(false);
await realtime.warmup();
expect(realtime.broker.isConfigured()).toBe(true);
expect(sharedClientMocks.getSharedClient).toHaveBeenCalledWith(
expect.objectContaining({ authRequirement: "subscription" }),
);
await realtime.cleanup();
});
it("handles offer preflights only for configured Control UI origins", async () => {
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({
gateway: {
controlUi: {
allowedOrigins: ["https://Control.Example"],
},
},
}),
getPluginConfig: () => ({}),
});
try {
const accepted = createResponseHarness();
await expect(
realtime.handler(createPreflightRequest("https://control.example"), accepted.res),
).resolves.toBe(true);
expect(accepted.res.statusCode).toBe(204);
expect(accepted.setHeader).toHaveBeenCalledWith(
"Access-Control-Allow-Origin",
"https://control.example",
);
expect(accepted.setHeader).toHaveBeenCalledWith(
"Access-Control-Allow-Methods",
"POST, OPTIONS",
);
expect(accepted.setHeader).toHaveBeenCalledWith(
"Access-Control-Allow-Headers",
"Authorization, Content-Type",
);
expect(accepted.setHeader).toHaveBeenCalledWith(
"Access-Control-Allow-Private-Network",
"true",
);
const rejected = createResponseHarness();
await expect(
realtime.handler(createPreflightRequest("https://untrusted.example"), rejected.res),
).resolves.toBe(true);
expect(rejected.res.statusCode).toBe(403);
expect(rejected.setHeader).not.toHaveBeenCalledWith(
"Access-Control-Allow-Origin",
expect.anything(),
);
expect(sharedClientMocks.getClient).not.toHaveBeenCalled();
} finally {
await realtime.cleanup();
}
});
it("probes again after failed warmup without advertising or reserving the failure", async () => {
const fake = createFakeClient();
const now = vi.spyOn(Date, "now").mockReturnValue(10_000);
sharedClientMocks.getSharedClient.mockRejectedValueOnce(new Error("ChatGPT login required"));
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
await expect(realtime.warmup()).rejects.toThrow("ChatGPT login required");
expect(realtime.broker.isConfigured()).toBe(false);
sharedClientMocks.getSharedClient.mockRejectedValueOnce(new Error("ChatGPT login required"));
await expect(realtime.broker.createBrowserSession({ providerConfig: {} })).rejects.toThrow(
"ChatGPT login required",
);
useFakeClient(fake);
now.mockReturnValue(11_000);
expect(realtime.broker.isConfigured()).toBe(false);
await vi.waitFor(() => {
expect(realtime.broker.isConfigured()).toBe(true);
});
await Promise.all(
Array.from({ length: 8 }, () => realtime.broker.createBrowserSession({ providerConfig: {} })),
);
await realtime.cleanup();
now.mockRestore();
});
it("revalidates after the warmed Codex client closes", async () => {
const first = createFakeClient();
const replacement = createFakeClient();
useFakeClient(first);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
try {
await realtime.warmup();
const session = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (session.transport !== "webrtc") {
throw new Error("Expected Codex browser sessions to use WebRTC");
}
await realtime.handler(createSdpRequest(session.clientSecret), createResponseHarness().res);
sharedClientMocks.getSharedClient.mockResolvedValue(replacement.client);
first.emitClose();
await vi.waitFor(() => {
expect(realtime.broker.isConfigured()).toBe(true);
expect(sharedClientMocks.releaseClient).toHaveBeenCalledWith(first.client);
});
await expect(
realtime.broker.createBrowserSession({ providerConfig: {} }),
).resolves.toMatchObject({
provider: "openai",
transport: "webrtc",
});
expect(sharedClientMocks.getSharedClient).toHaveBeenLastCalledWith(
expect.objectContaining({ authRequirement: "subscription" }),
);
} finally {
await realtime.cleanup();
}
});
it("redeems browser reservations once and invalidates pending ones on cleanup", async () => {
const fake = createFakeClient();
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({
gateway: {
controlUi: {
allowedOrigins: ["https://Control.Example"],
},
},
}),
getPluginConfig: () => ({}),
});
expect(realtime.broker.capabilities).toEqual({
transports: ["webrtc"],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
});
const first = await realtime.broker.createBrowserSession({
providerConfig: {},
instructions: " Keep the same Talk persona. ",
model: " gpt-realtime-2 ",
voice: " Marin ",
initialItems: [
{ role: "user", text: "Earlier question" },
{ role: "assistant", text: "Earlier answer" },
],
});
const second = await realtime.broker.createBrowserSession({ providerConfig: {} });
const cancelled = await realtime.broker.createBrowserSession({ providerConfig: {} });
expect(first).toMatchObject({
provider: "openai",
transport: "webrtc",
offerUrl: CODEX_REALTIME_OFFER_PATH,
voice: "Marin",
clientSecret: expect.stringMatching(/^[A-Za-z0-9_-]{40,}$/),
expiresAt: expect.any(Number),
});
expect(first).not.toHaveProperty("model");
if (first.transport !== "webrtc" || second.transport !== "webrtc") {
throw new Error("Expected Codex browser sessions to use WebRTC");
}
expect(second.clientSecret).not.toBe(first.clientSecret);
await realtime.broker.cancelBrowserSession(cancelled);
try {
const accepted = createResponseHarness();
await expect(
realtime.handler(
createSdpRequest(first.clientSecret, "https://control.example"),
accepted.res,
),
).resolves.toBe(true);
expect(accepted.res.statusCode).toBe(200);
expect(accepted.readBody()).toBe("v=answer\r\n");
expect(accepted.setHeader).toHaveBeenCalledWith(
"Access-Control-Allow-Origin",
"https://control.example",
);
const threadStartParams = (fake.client.request as ReturnType<typeof vi.fn>).mock.calls.find(
([method]) => method === "thread/start",
)?.[1];
expect(threadStartParams).toEqual({
cwd: process.cwd(),
ephemeral: true,
approvalPolicy: "never",
sandbox: "read-only",
config: { "features.realtime_conversation": true },
});
const realtimeStartParams = (fake.client.request as ReturnType<typeof vi.fn>).mock.calls.find(
([method]) => method === "thread/realtime/start",
)?.[1];
expect(realtimeStartParams).toEqual({
threadId: "thread-1",
outputModality: "audio",
transport: { type: "webrtc", sdp: "v=offer\r\n" },
version: "v3",
includeStartupContext: true,
voice: "Marin",
initialItems: [
{ role: "developer", text: "Keep the same Talk persona." },
{ role: "user", text: "Earlier question" },
{ role: "assistant", text: "Earlier answer" },
],
});
expect(realtimeStartParams).not.toHaveProperty("prompt");
expect(realtimeStartParams).not.toHaveProperty("model");
const replayed = createResponseHarness();
await expect(
realtime.handler(createSdpRequest(first.clientSecret), replayed.res),
).resolves.toBe(true);
expect(replayed.res.statusCode).toBe(401);
expect(sharedClientMocks.getClient).toHaveBeenCalledTimes(1);
if (cancelled.transport !== "webrtc") {
throw new Error("Expected cancelled Codex browser session to use WebRTC");
}
const cancelledResponse = createResponseHarness();
await expect(
realtime.handler(createSdpRequest(cancelled.clientSecret), cancelledResponse.res),
).resolves.toBe(true);
expect(cancelledResponse.res.statusCode).toBe(401);
await realtime.cleanup();
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
expect(sharedClientMocks.releaseClient).toHaveBeenCalledWith(fake.client);
const invalidated = createResponseHarness();
await expect(
realtime.handler(createSdpRequest(second.clientSecret), invalidated.res),
).resolves.toBe(true);
expect(invalidated.res.statusCode).toBe(401);
await expect(realtime.broker.createBrowserSession({ providerConfig: {} })).rejects.toThrow(
"Codex OAuth realtime is stopping",
);
} finally {
await realtime.cleanup();
}
});
it("aborts and closes backend startup when the browser offer disconnects", async () => {
const fake = createFakeClient({ stallRealtimeStart: true });
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
const reservation = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (reservation.transport !== "webrtc") {
throw new Error("Expected Codex browser session to use WebRTC");
}
const response = createResponseHarness();
try {
const handling = realtime.handler(createSdpRequest(reservation.clientSecret), response.res);
await vi.waitFor(() => {
expect(fake.readRealtimeStartSignal()).toBeDefined();
});
response.close();
await expect(handling).resolves.toBe(true);
expect(fake.readRealtimeStartSignal()?.aborted).toBe(true);
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
expect(sharedClientMocks.releaseClient).toHaveBeenCalledWith(fake.client);
expect(response.end).not.toHaveBeenCalled();
} finally {
await realtime.cleanup();
}
});
it("returns the Codex startup error when it arrives before the start response", async () => {
const fake = createFakeClient({
stallRealtimeStart: true,
genericRealtimeStartAbortError: true,
realtimeStartNotifications: [
{
method: "thread/realtime/error",
params: { threadId: "thread-1", message: "subscription unavailable" },
},
],
});
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
const reservation = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (reservation.transport !== "webrtc") {
throw new Error("Expected Codex browser session to use WebRTC");
}
const response = createResponseHarness();
try {
await expect(
realtime.handler(createSdpRequest(reservation.clientSecret), response.res),
).resolves.toBe(true);
expect(response.res.statusCode).toBe(502);
expect(response.readBody()).toBe("subscription unavailable");
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
} finally {
await realtime.cleanup();
}
});
it("does not return an SDP answer after Codex closes during startup", async () => {
const fake = createFakeClient({
stallRealtimeStart: true,
genericRealtimeStartAbortError: true,
realtimeStartNotifications: [
{
method: "thread/realtime/sdp",
params: { threadId: "thread-1", sdp: "v=stale-answer\r\n" },
},
{
method: "thread/realtime/closed",
params: { threadId: "thread-1", reason: "backend closed" },
},
],
});
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
const reservation = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (reservation.transport !== "webrtc") {
throw new Error("Expected Codex browser session to use WebRTC");
}
const response = createResponseHarness();
try {
await expect(
realtime.handler(createSdpRequest(reservation.clientSecret), response.res),
).resolves.toBe(true);
expect(response.res.statusCode).toBe(502);
expect(response.readBody()).toBe(
"Codex realtime session closed before returning an SDP answer",
);
expect(response.readBody()).not.toContain("stale-answer");
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
} finally {
await realtime.cleanup();
}
});
it("releases the backend when Codex reports an error after startup", async () => {
const fake = createFakeClient();
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
const reservation = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (reservation.transport !== "webrtc") {
throw new Error("Expected Codex browser session to use WebRTC");
}
const response = createResponseHarness();
try {
await expect(
realtime.handler(createSdpRequest(reservation.clientSecret), response.res),
).resolves.toBe(true);
fake.emitNotification({
method: "thread/realtime/error",
params: { threadId: "thread-1", message: "backend failed" },
});
await vi.waitFor(() => {
expect(sharedClientMocks.releaseClient).toHaveBeenCalledWith(fake.client);
});
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
} finally {
await realtime.cleanup();
}
});
it("closes the backend when the browser disconnects while the SDP answer is flushing", async () => {
const fake = createFakeClient();
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
const reservation = await realtime.broker.createBrowserSession({ providerConfig: {} });
if (reservation.transport !== "webrtc") {
throw new Error("Expected Codex browser session to use WebRTC");
}
const response = createResponseHarness({ autoFinish: false });
try {
const handling = realtime.handler(createSdpRequest(reservation.clientSecret), response.res);
await vi.waitFor(() => {
expect(response.end).toHaveBeenCalledWith("v=answer\r\n");
});
response.close();
await expect(handling).resolves.toBe(true);
expect(fake.methods).toContain("thread/realtime/stop");
expect(fake.methods).toContain("thread/unsubscribe");
expect(sharedClientMocks.releaseClient).toHaveBeenCalledWith(fake.client);
} finally {
await realtime.cleanup();
}
});
it("caps concurrent pending and active browser sessions", async () => {
const fake = createFakeClient();
useFakeClient(fake);
const realtime = createCodexRealtimeBrowserSessionBroker({
getConfig: () => ({}),
getPluginConfig: () => ({}),
});
await Promise.all(
Array.from({ length: 8 }, () => realtime.broker.createBrowserSession({ providerConfig: {} })),
);
await expect(realtime.broker.createBrowserSession({ providerConfig: {} })).rejects.toThrow(
"Too many concurrent Codex OAuth realtime sessions",
);
await realtime.cleanup();
});
});
@@ -1,694 +0,0 @@
// Experimental ChatGPT OAuth browser session broker for Control UI realtime Talk.
import { randomBytes } from "node:crypto";
import type { IncomingMessage, ServerResponse } from "node:http";
import { resolveAgentDir, resolveDefaultAgentId } from "openclaw/plugin-sdk/agent-runtime";
import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts";
import type {
RealtimeVoiceBrowserSession,
RealtimeVoiceBrowserSessionCreateRequest,
RealtimeVoiceProviderCapabilities,
} from "openclaw/plugin-sdk/realtime-voice";
import { readRequestBodyWithLimit } from "openclaw/plugin-sdk/webhook-request-guards";
import {
CODEX_APP_SERVER_UNSUBSCRIBE_TIMEOUT_MS,
unsubscribeCodexThreadBestEffort,
} from "./app-server/attempt-client-cleanup.js";
import type { CodexAppServerClient } from "./app-server/client.js";
import { readCodexPluginConfig, resolveCodexAppServerRuntimeOptions } from "./app-server/config.js";
import { assertCodexThreadStartResponse } from "./app-server/protocol-validators.js";
import type { CodexThreadStartParams } from "./app-server/protocol.js";
import {
getLeasedSharedCodexAppServerClient,
getSharedCodexAppServerClient,
releaseLeasedSharedCodexAppServerClient,
} from "./app-server/shared-client.js";
const CODEX_REALTIME_OFFER_PATH = "/plugins/codex/realtime/calls";
const CODEX_REALTIME_PENDING_TTL_MS = 60_000;
const CODEX_REALTIME_SESSION_TTL_MS = 30 * 60_000;
const CODEX_REALTIME_MAX_SESSIONS = 8;
const CODEX_REALTIME_MAX_SDP_BYTES = 256 * 1024;
const CODEX_REALTIME_START_TIMEOUT_MS = 60_000;
const CODEX_REALTIME_PROBE_COOLDOWN_MS = 1_000;
type CodexRealtimeBrowserSessionCreateRequest = RealtimeVoiceBrowserSessionCreateRequest & {
agentId?: string;
workspaceDir?: string;
initialItems?: Array<{
role: "user" | "assistant";
text: string;
}>;
};
type CodexRealtimeProviderCapabilities = Partial<RealtimeVoiceProviderCapabilities> & {
handlesAgentConsult?: boolean;
};
type PendingOffer = {
expiresAt: number;
request: CodexRealtimeBrowserSessionCreateRequest;
};
type CodexRealtimeBrowserSessionFallback = {
capabilities: CodexRealtimeProviderCapabilities;
isConfigured: () => boolean;
createBrowserSession: (
request: CodexRealtimeBrowserSessionCreateRequest,
) => Promise<RealtimeVoiceBrowserSession>;
cancelBrowserSession: (session: RealtimeVoiceBrowserSession) => Promise<void> | void;
};
type ActiveSession = {
client: CodexAppServerClient;
reservationToken: string;
threadId: string;
timer: NodeJS.Timeout;
disposeNotificationHandler: () => void;
};
type RealtimeNotificationParams = {
threadId?: unknown;
sdp?: unknown;
message?: unknown;
};
type ResponseDeliveryWaiter = {
result: Promise<boolean>;
cancel: () => void;
};
function createResponseDeliveryWaiter(
res: ServerResponse,
onDelivered: () => void,
): ResponseDeliveryWaiter {
let settle!: (delivered: boolean) => void;
const result = new Promise<boolean>((resolve) => {
settle = (delivered) => {
res.removeListener("finish", onFinish);
res.removeListener("close", onClose);
resolve(delivered);
};
});
const onFinish = () => {
// ServerResponse may emit close immediately after finish. Remove the
// disconnect abort synchronously so normal completion keeps WebRTC alive.
onDelivered();
settle(true);
};
const onClose = () => settle(false);
res.once("finish", onFinish);
res.once("close", onClose);
return { result, cancel: () => settle(false) };
}
function respondText(res: ServerResponse, statusCode: number, body: string): void {
res.statusCode = statusCode;
res.setHeader("cache-control", "no-store");
res.setHeader("content-type", "text/plain; charset=utf-8");
res.setHeader("x-content-type-options", "nosniff");
res.end(body);
}
function resolveConfiguredControlUiOrigin(
req: IncomingMessage,
cfg: OpenClawConfig | undefined,
): string | undefined {
const rawOrigin = typeof req.headers.origin === "string" ? req.headers.origin.trim() : "";
if (!rawOrigin) {
return undefined;
}
let origin: string;
try {
const parsed = new URL(rawOrigin);
if (parsed.origin !== rawOrigin || parsed.username || parsed.password) {
return undefined;
}
origin = parsed.origin;
} catch {
return undefined;
}
const allowed = cfg?.gateway?.controlUi?.allowedOrigins ?? [];
return allowed.some((candidate) => {
const normalized = candidate.trim().toLowerCase();
return normalized === "*" || normalized === origin;
})
? origin
: undefined;
}
function applyRealtimeOfferCorsHeaders(
req: IncomingMessage,
res: ServerResponse,
cfg: OpenClawConfig | undefined,
): boolean {
if (!req.headers.origin) {
return true;
}
const origin = resolveConfiguredControlUiOrigin(req, cfg);
if (!origin) {
return false;
}
res.setHeader("Access-Control-Allow-Origin", origin);
res.setHeader("Vary", "Origin");
return true;
}
function readBearerToken(req: IncomingMessage): string | undefined {
const authorization = req.headers.authorization?.trim();
const match = authorization?.match(/^Bearer\s+([^\s]+)$/i);
return match?.[1];
}
function readNotificationParams(value: unknown): RealtimeNotificationParams {
return value && typeof value === "object" ? (value as RealtimeNotificationParams) : {};
}
function buildCodexRealtimeThreadStartParams(params: { cwd: string }): CodexThreadStartParams {
return {
cwd: params.cwd,
ephemeral: true,
approvalPolicy: "never",
sandbox: "read-only",
config: { "features.realtime_conversation": true },
};
}
function buildCodexRealtimeStartParams(params: {
threadId: string;
sdp: string;
developerInstructions?: string;
voice?: string;
initialItems?: CodexRealtimeBrowserSessionCreateRequest["initialItems"];
}): Record<string, unknown> {
const initialItems = [
...(params.developerInstructions
? [{ role: "developer" as const, text: params.developerInstructions }]
: []),
...(params.initialItems ?? []),
];
return {
threadId: params.threadId,
outputModality: "audio",
transport: { type: "webrtc", sdp: params.sdp },
version: "v3",
includeStartupContext: true,
...(params.voice ? { voice: params.voice } : {}),
...(initialItems.length > 0 ? { initialItems } : {}),
};
}
function waitForRealtimeSdpAnswer(
answerPromise: Promise<string>,
signal: AbortSignal,
): Promise<string> {
return new Promise<string>((resolve, reject) => {
let settled = false;
const finish = (result: { answer: string } | { error: Error }) => {
if (settled) {
return;
}
settled = true;
clearTimeout(timeout);
signal.removeEventListener("abort", onAbort);
if ("answer" in result) {
resolve(result.answer);
} else {
reject(result.error);
}
};
const onAbort = () =>
finish({
error:
signal.reason instanceof Error
? signal.reason
: new Error("Codex realtime session stopped during startup"),
});
const timeout = setTimeout(
() => finish({ error: new Error("Codex realtime SDP answer timed out") }),
CODEX_REALTIME_START_TIMEOUT_MS,
);
timeout.unref?.();
signal.addEventListener("abort", onAbort, { once: true });
if (signal.aborted) {
onAbort();
return;
}
void answerPromise.then(
(answer) => finish({ answer }),
(error: unknown) =>
finish({ error: error instanceof Error ? error : new Error("Codex realtime failed") }),
);
});
}
export function createCodexRealtimeBrowserSessionBroker(params: {
getConfig: () => OpenClawConfig | undefined;
getPluginConfig: () => unknown;
}): {
broker: CodexRealtimeBrowserSessionFallback;
handler: (req: IncomingMessage, res: ServerResponse) => Promise<boolean>;
warmup: () => Promise<void>;
cleanup: () => Promise<void>;
} {
const pendingOffers = new Map<string, PendingOffer>();
const reservations = new Set<string>();
const activeSessions = new Set<ActiveSession>();
const inFlightHandlers = new Set<Promise<boolean>>();
const readyClients = new Set<CodexAppServerClient>();
const shutdownController = new AbortController();
let cleanedUp = false;
let probePromise: Promise<void> | undefined;
let nextProbeAt = 0;
const closeSession = async (session: ActiveSession) => {
if (!activeSessions.delete(session)) {
return;
}
clearTimeout(session.timer);
reservations.delete(session.reservationToken);
session.disposeNotificationHandler();
try {
await session.client.request(
"thread/realtime/stop",
{ threadId: session.threadId },
{ timeoutMs: 2_000 },
);
} catch {
// The peer or app-server may already have closed the realtime transport.
}
await unsubscribeCodexThreadBestEffort(session.client, {
threadId: session.threadId,
timeoutMs: CODEX_APP_SERVER_UNSUBSCRIBE_TIMEOUT_MS,
});
releaseLeasedSharedCodexAppServerClient(session.client);
};
const resolveSubscriptionClientOptions = (request: {
cfg?: OpenClawConfig;
agentId?: string;
}) => {
const pluginConfig = readCodexPluginConfig(params.getPluginConfig());
const agentDir =
request.agentId && request.cfg ? resolveAgentDir(request.cfg, request.agentId) : undefined;
const runtime = resolveCodexAppServerRuntimeOptions({
pluginConfig,
config: request.cfg,
agentDir,
});
return {
startOptions: runtime.start,
pluginConfig,
config: request.cfg,
agentDir,
authRequirement: "subscription" as const,
timeoutMs: CODEX_REALTIME_START_TIMEOUT_MS,
};
};
const ensureSubscriptionRuntime = async (request: { cfg?: OpenClawConfig; agentId?: string }) => {
const client = await getSharedCodexAppServerClient({
...resolveSubscriptionClientOptions(request),
abandonSignal: shutdownController.signal,
});
if (!readyClients.has(client)) {
readyClients.add(client);
client.addCloseHandler((closedClient) => {
readyClients.delete(closedClient);
void Promise.allSettled(
[...activeSessions]
.filter((session) => session.client === closedClient)
.map((session) => closeSession(session)),
);
nextProbeAt = 0;
requestProbe();
});
}
};
const runProbe = () => {
if (probePromise) {
return probePromise;
}
const cfg = params.getConfig();
nextProbeAt = Date.now() + CODEX_REALTIME_PROBE_COOLDOWN_MS;
const running = ensureSubscriptionRuntime({
cfg,
...(cfg ? { agentId: resolveDefaultAgentId(cfg) } : {}),
}).finally(() => {
if (probePromise === running) {
probePromise = undefined;
}
});
probePromise = running;
return running;
};
function requestProbe(): void {
if (
cleanedUp ||
shutdownController.signal.aborted ||
probePromise ||
Date.now() < nextProbeAt
) {
return;
}
// Readiness stays false while the async Codex probe runs. This lets auth or
// process recovery become visible without shadowing another healthy provider.
void runProbe().catch(() => undefined);
}
const warmup = async () => {
await runProbe();
};
const prunePendingOffers = () => {
const now = Date.now();
for (const [token, offer] of pendingOffers) {
if (offer.expiresAt <= now) {
pendingOffers.delete(token);
reservations.delete(token);
}
}
};
const broker: CodexRealtimeBrowserSessionFallback = {
capabilities: {
transports: ["webrtc"],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
},
isConfigured: () => {
if (cleanedUp || shutdownController.signal.aborted) {
return false;
}
if (readyClients.size > 0) {
return true;
}
requestProbe();
return false;
},
createBrowserSession: async (request: CodexRealtimeBrowserSessionCreateRequest) => {
if (cleanedUp || shutdownController.signal.aborted) {
throw new Error("Codex OAuth realtime is stopping; restart Gateway and try again");
}
// Revalidate the request's exact agent/runtime before the Gateway persists a
// client voice session. The warmed shared process makes this path cheap.
await ensureSubscriptionRuntime(request);
if (cleanedUp || shutdownController.signal.aborted) {
throw new Error("Codex OAuth realtime is stopping; restart Gateway and try again");
}
prunePendingOffers();
if (reservations.size >= CODEX_REALTIME_MAX_SESSIONS) {
throw new Error("Too many concurrent Codex OAuth realtime sessions; try again in a minute");
}
const token = randomBytes(32).toString("base64url");
const expiresAt = Date.now() + CODEX_REALTIME_PENDING_TTL_MS;
reservations.add(token);
pendingOffers.set(token, { expiresAt, request });
const voice = request.voice?.trim() || undefined;
return {
provider: "openai",
transport: "webrtc",
clientSecret: token,
offerUrl: CODEX_REALTIME_OFFER_PATH,
...(voice ? { voice } : {}),
expiresAt,
};
},
cancelBrowserSession: (session) => {
if (session.transport !== "webrtc") {
return;
}
pendingOffers.delete(session.clientSecret);
reservations.delete(session.clientSecret);
},
};
const handleOffer = async (req: IncomingMessage, res: ServerResponse): Promise<boolean> => {
const corsAllowed = applyRealtimeOfferCorsHeaders(req, res, params.getConfig());
if (req.method === "OPTIONS") {
if (!corsAllowed) {
respondText(res, 403, "Origin not allowed");
return true;
}
res.statusCode = 204;
res.setHeader("cache-control", "no-store");
res.setHeader("Access-Control-Allow-Methods", "POST, OPTIONS");
res.setHeader("Access-Control-Allow-Headers", "Authorization, Content-Type");
res.setHeader(
"Vary",
"Origin, Access-Control-Request-Method, Access-Control-Request-Headers",
);
if (req.headers["access-control-request-private-network"] === "true") {
res.setHeader("Access-Control-Allow-Private-Network", "true");
}
res.setHeader("Access-Control-Max-Age", "600");
res.end();
return true;
}
if (req.method !== "POST") {
respondText(res, 405, "Method not allowed");
return true;
}
if (!req.headers["content-type"]?.toLowerCase().startsWith("application/sdp")) {
respondText(res, 415, "Expected application/sdp");
return true;
}
prunePendingOffers();
const token = readBearerToken(req);
const offer = token ? pendingOffers.get(token) : undefined;
if (!token || !offer || offer.expiresAt <= Date.now()) {
respondText(res, 401, "Invalid or expired realtime session token");
return true;
}
// A browser session token is single-use so captured requests cannot be replayed.
pendingOffers.delete(token);
const requestController = new AbortController();
const abortFromBrowser = () => {
requestController.abort(new Error("Browser realtime offer request closed"));
};
req.once("aborted", abortFromBrowser);
res.once("close", abortFromBrowser);
const detachBrowserAbort = () => {
req.removeListener("aborted", abortFromBrowser);
res.removeListener("close", abortFromBrowser);
};
const lifecycleSignal = AbortSignal.any([shutdownController.signal, requestController.signal]);
let client: CodexAppServerClient | undefined;
let session: ActiveSession | undefined;
let threadId: string | undefined;
let reservationTransferred = false;
let responseDeliveryWaiter: ResponseDeliveryWaiter | undefined;
try {
const sdp = await readRequestBodyWithLimit(req, {
maxBytes: CODEX_REALTIME_MAX_SDP_BYTES,
timeoutMs: 15_000,
});
if (!sdp.trim()) {
respondText(res, 400, "SDP offer is required");
return true;
}
if (lifecycleSignal.aborted) {
throw new Error("Codex realtime session stopped during startup");
}
// Share the agent's normal Codex app-server process. A fresh ephemeral thread
// keeps realtime from replacing a live normal turn on the bound Codex thread.
client = await getLeasedSharedCodexAppServerClient({
...resolveSubscriptionClientOptions(offer.request),
abandonSignal: lifecycleSignal,
});
const started = assertCodexThreadStartResponse(
await client.request(
"thread/start",
buildCodexRealtimeThreadStartParams({
cwd: offer.request.workspaceDir ?? process.cwd(),
}),
{ timeoutMs: CODEX_REALTIME_START_TIMEOUT_MS, signal: lifecycleSignal },
),
);
threadId = started.thread.id;
if (lifecycleSignal.aborted) {
throw new Error("Codex realtime session stopped during startup");
}
let resolveSdp!: (answer: string) => void;
let rejectSdp!: (error: Error) => void;
const answerPromise = new Promise<string>((resolve, reject) => {
resolveSdp = resolve;
rejectSdp = reject;
});
const startupController = new AbortController();
const startupSignal = AbortSignal.any([lifecycleSignal, startupController.signal]);
let startupComplete = false;
const stopStartup = (error: Error) => {
if (!startupController.signal.aborted) {
startupController.abort(error);
}
};
const closeAfterStartup = () => {
// During startup the request handler owns ordered cleanup. Once the SDP
// response is ready, terminal notifications must release the live session.
if (startupComplete && session) {
void closeSession(session);
}
};
const disposeNotificationHandler = client.addNotificationHandler((notification) => {
const notificationParams = readNotificationParams(notification.params);
if (notificationParams.threadId !== threadId) {
return;
}
if (notification.method === "thread/realtime/sdp") {
if (typeof notificationParams.sdp === "string" && notificationParams.sdp.trim()) {
resolveSdp(notificationParams.sdp);
} else {
rejectSdp(new Error("Codex returned an empty realtime SDP answer"));
}
return;
}
if (notification.method === "thread/realtime/error") {
const error = new Error(
typeof notificationParams.message === "string"
? notificationParams.message
: "Codex realtime session failed",
);
rejectSdp(error);
stopStartup(error);
closeAfterStartup();
return;
}
if (notification.method === "thread/realtime/closed") {
const error = new Error("Codex realtime session closed before returning an SDP answer");
rejectSdp(error);
stopStartup(error);
closeAfterStartup();
}
});
const timer = setTimeout(() => {
if (session) {
void closeSession(session);
}
}, CODEX_REALTIME_SESSION_TTL_MS);
timer.unref?.();
session = {
client,
reservationToken: token,
threadId,
timer,
disposeNotificationHandler,
};
reservationTransferred = true;
activeSessions.add(session);
// Observe the answer before starting the request: Codex may emit a terminal
// notification before the matching JSON-RPC response reaches this process.
const answerResultPromise = waitForRealtimeSdpAnswer(answerPromise, startupSignal).then(
(answer) => ({ ok: true as const, answer }),
(error: unknown) => ({
ok: false as const,
error: error instanceof Error ? error : new Error("Codex realtime session failed"),
}),
);
try {
await client.request(
"thread/realtime/start",
buildCodexRealtimeStartParams({
threadId,
sdp,
developerInstructions: offer.request.instructions?.trim() || undefined,
voice: offer.request.voice?.trim() || undefined,
initialItems: offer.request.initialItems,
}),
{ timeoutMs: CODEX_REALTIME_START_TIMEOUT_MS, signal: startupSignal },
);
} catch (error) {
if (startupController.signal.aborted) {
const answerResult = await answerResultPromise;
if (!answerResult.ok) {
throw answerResult.error;
}
throw startupController.signal.reason instanceof Error
? startupController.signal.reason
: new Error("Codex realtime session stopped during startup");
}
stopStartup(
error instanceof Error ? error : new Error("Codex realtime session failed to start"),
);
throw error;
}
const answerResult = await answerResultPromise;
if (!answerResult.ok) {
throw answerResult.error;
}
if (startupSignal.aborted) {
throw startupSignal.reason instanceof Error
? startupSignal.reason
: new Error("Codex realtime session stopped during startup");
}
startupComplete = true;
responseDeliveryWaiter = createResponseDeliveryWaiter(res, detachBrowserAbort);
res.statusCode = 200;
res.setHeader("cache-control", "no-store");
res.setHeader("content-type", "application/sdp");
res.setHeader("x-content-type-options", "nosniff");
res.end(answerResult.answer);
const delivered = await responseDeliveryWaiter.result;
responseDeliveryWaiter = undefined;
if (!delivered || lifecycleSignal.aborted) {
await closeSession(session);
}
return true;
} catch (error) {
if (session) {
await closeSession(session);
} else if (client) {
if (threadId) {
await unsubscribeCodexThreadBestEffort(client, {
threadId,
timeoutMs: CODEX_APP_SERVER_UNSUBSCRIBE_TIMEOUT_MS,
});
}
releaseLeasedSharedCodexAppServerClient(client);
}
if (requestController.signal.aborted) {
return true;
}
const message = error instanceof Error ? error.message : "Codex realtime session failed";
respondText(res, 502, message);
return true;
} finally {
responseDeliveryWaiter?.cancel();
detachBrowserAbort();
if (!reservationTransferred) {
reservations.delete(token);
}
}
};
const trackedHandleOffer = (req: IncomingMessage, res: ServerResponse): Promise<boolean> => {
const handling = handleOffer(req, res);
inFlightHandlers.add(handling);
return handling.finally(() => {
inFlightHandlers.delete(handling);
});
};
const handler = trackedHandleOffer;
const cleanup = async () => {
if (cleanedUp) {
return;
}
cleanedUp = true;
shutdownController.abort();
pendingOffers.clear();
await Promise.all([...activeSessions].map((session) => closeSession(session)));
await Promise.allSettled(inFlightHandlers);
reservations.clear();
};
return { broker, handler, warmup, cleanup };
}
export { CODEX_REALTIME_OFFER_PATH };
+18 -193
View File
@@ -1,38 +1,9 @@
// Openai tests cover realtime voice provider plugin behavior.
import { REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ } from "openclaw/plugin-sdk/realtime-voice";
import type {
RealtimeVoiceBridge,
RealtimeVoiceBrowserSession,
RealtimeVoiceTool,
} from "openclaw/plugin-sdk/realtime-voice";
import type { RealtimeVoiceBridge, RealtimeVoiceTool } from "openclaw/plugin-sdk/realtime-voice";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { buildOpenAIRealtimeVoiceProvider } from "./realtime-voice-provider.js";
const CODEX_REALTIME_GLOBAL_STATE = Symbol.for("openclaw.codex.realtime-voice.v1");
const INTERNAL_REALTIME_VOICE_PROVIDER = Symbol.for("openclaw.internal.realtime-voice-provider.v1");
function readInternalRealtimeVoiceProviderApi(provider: object) {
return Reflect.get(provider, INTERNAL_REALTIME_VOICE_PROVIDER) as {
isBrowserSessionConfigured: (ctx: {
cfg?: object;
providerConfig: Record<string, unknown>;
}) => boolean;
resolveBrowserSessionCapabilities: (ctx: {
cfg?: object;
providerConfig: Record<string, unknown>;
}) => {
handlesAgentConsult?: boolean;
supportsToolCalls?: boolean;
supportsVideoFrames?: boolean;
transports?: string[];
};
cancelBrowserSession: (
request: Record<string, unknown>,
session: RealtimeVoiceBrowserSession,
) => Promise<void>;
};
}
const {
FakeWebSocket,
execFileSyncMock,
@@ -311,7 +282,6 @@ describe("buildOpenAIRealtimeVoiceProvider", () => {
});
afterEach(() => {
Reflect.deleteProperty(globalThis, CODEX_REALTIME_GLOBAL_STATE);
vi.useRealTimers();
vi.unstubAllEnvs();
});
@@ -350,73 +320,6 @@ describe("buildOpenAIRealtimeVoiceProvider", () => {
expect(bridge.supportsToolResultSuppression).toBe(true);
});
it("uses broker-owned capabilities when Codex OAuth is the browser fallback", () => {
const broker = {
capabilities: {
transports: ["webrtc" as const],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
},
isConfigured: () => true,
createBrowserSession: vi.fn(),
cancelBrowserSession: vi.fn(),
};
const provider = buildOpenAIRealtimeVoiceProvider({
resolveCodexRealtimeBrowserSessionFallback: () => broker,
});
const internalApi = readInternalRealtimeVoiceProviderApi(provider);
expect(
internalApi.resolveBrowserSessionCapabilities({
providerConfig: {},
}),
).toMatchObject({
transports: ["webrtc"],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
});
expect(
internalApi.resolveBrowserSessionCapabilities({
providerConfig: { apiKey: "sk-platform" }, // pragma: allowlist secret
}),
).toBe(provider.capabilities);
});
it("discovers the optional Codex OAuth runtime without a Plugin SDK registrar", () => {
const broker = {
capabilities: {
transports: ["webrtc" as const],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
},
isConfigured: () => true,
createBrowserSession: vi.fn(),
cancelBrowserSession: vi.fn(),
};
Reflect.set(globalThis, CODEX_REALTIME_GLOBAL_STATE, {
version: 1,
fallback: broker,
});
const provider = buildOpenAIRealtimeVoiceProvider();
const internalApi = readInternalRealtimeVoiceProviderApi(provider);
expect(provider.isConfigured({ providerConfig: {} })).toBe(false);
expect(internalApi.isBrowserSessionConfigured({ providerConfig: {} })).toBe(true);
expect(
internalApi.resolveBrowserSessionCapabilities({
providerConfig: {},
}),
).toMatchObject({
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
});
});
it("adds OpenClaw attribution headers to native realtime websocket requests", () => {
vi.stubEnv("OPENCLAW_VERSION", "2026.3.22");
const provider = buildOpenAIRealtimeVoiceProvider();
@@ -823,107 +726,29 @@ describe("buildOpenAIRealtimeVoiceProvider", () => {
});
});
it("falls back to Codex OAuth for browser sessions without Platform auth", async () => {
const createBrowserSession = vi.fn(async () => ({
provider: "openai",
transport: "webrtc" as const,
clientSecret: "codex-session-token",
offerUrl: "/plugins/codex/realtime/calls",
}));
const cancelBrowserSession = vi.fn();
const isConfigured = vi.fn(() => true);
const broker = {
capabilities: {
transports: ["webrtc" as const],
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
},
isConfigured,
createBrowserSession,
cancelBrowserSession,
};
const provider = buildOpenAIRealtimeVoiceProvider({
resolveCodexRealtimeBrowserSessionFallback: () => broker,
});
const cfg = { agents: { defaults: {} } } as never;
const request = {
cfg,
providerConfig: {},
model: "gpt-realtime-2",
agentId: "main",
workspaceDir: "/tmp/openclaw-agent-workspace",
initialItems: [],
};
const internalApi = readInternalRealtimeVoiceProviderApi(provider);
expect(provider.isConfigured(request)).toBe(false);
expect(internalApi.isBrowserSessionConfigured(request)).toBe(true);
const session = await provider.createBrowserSession?.(request);
expect(session).toMatchObject({
clientSecret: "codex-session-token",
offerUrl: "/plugins/codex/realtime/calls",
});
if (!session) {
throw new Error("Expected Codex OAuth browser session");
}
await internalApi.cancelBrowserSession(request, session);
expect(isConfigured).toHaveBeenCalledWith();
expect(createBrowserSession).toHaveBeenCalledWith(request);
expect(cancelBrowserSession).toHaveBeenCalledWith(session);
expect(fetchWithSsrFGuardMock).not.toHaveBeenCalled();
});
it("prefers Platform auth over the Codex OAuth browser broker", async () => {
const createBrowserSession = vi.fn();
const broker = {
capabilities: {},
isConfigured: () => true,
createBrowserSession,
cancelBrowserSession: vi.fn(),
};
fetchWithSsrFGuardMock.mockResolvedValueOnce({
response: createJsonResponse({
client_secret: { value: "client-secret-123" },
}),
release: vi.fn(async () => undefined),
});
const provider = buildOpenAIRealtimeVoiceProvider({
resolveCodexRealtimeBrowserSessionFallback: () => broker,
});
await provider.createBrowserSession?.({
providerConfig: { apiKey: "sk-platform" }, // pragma: allowlist secret
});
expect(createBrowserSession).not.toHaveBeenCalled();
expectRecordFields(requireFetchHeaders(), "fetch headers", {
Authorization: "Bearer sk-platform", // pragma: allowlist secret
});
});
it("does not hide an unresolved Platform credential behind Codex OAuth", async () => {
vi.stubEnv("OPENAI_API_KEY", "keychain:openclaw:OPENAI_REALTIME_MISSING_TEST");
execFileSyncMock.mockImplementationOnce(() => {
throw new Error("keychain unavailable");
});
const createBrowserSession = vi.fn();
const broker = {
capabilities: {},
isConfigured: () => true,
createBrowserSession,
cancelBrowserSession: vi.fn(),
};
const provider = buildOpenAIRealtimeVoiceProvider({
resolveCodexRealtimeBrowserSessionFallback: () => broker,
});
it("requires Platform auth for browser sessions", async () => {
const provider = buildOpenAIRealtimeVoiceProvider();
await expect(
provider.createBrowserSession?.({
providerConfig: {},
}),
).rejects.toThrow("OpenAI Realtime voice requires an OpenAI Platform API key");
expect(fetchWithSsrFGuardMock).not.toHaveBeenCalled();
});
it("reports an unresolved Platform credential without trying another auth route", async () => {
vi.stubEnv("OPENAI_API_KEY", "keychain:openclaw:OPENAI_REALTIME_MISSING_TEST");
execFileSyncMock.mockImplementationOnce(() => {
throw new Error("keychain unavailable");
});
const provider = buildOpenAIRealtimeVoiceProvider();
await expect(
provider.createBrowserSession?.({
providerConfig: {},
}),
).rejects.toThrow("OpenAI Realtime voice requires an OpenAI Platform API key");
expect(createBrowserSession).not.toHaveBeenCalled();
});
it("treats OpenAI API-key auth profiles as configured for browser realtime sessions", () => {
+3 -152
View File
@@ -1549,78 +1549,8 @@ function resolveOpenAIRealtimeBrowserOfferHeaders(): Record<string, string> | un
return Object.keys(browserHeaders).length > 0 ? browserHeaders : undefined;
}
type CodexRealtimeBrowserSessionFallback = NonNullable<
ReturnType<typeof readCodexRealtimeBrowserSessionFallback>
>;
type OpenAIInternalRealtimeBrowserSessionCreateRequest =
RealtimeVoiceBrowserSessionCreateRequest & {
agentId: string;
workspaceDir: string;
initialItems: Array<{
role: "user" | "assistant";
text: string;
}>;
};
type OpenAIInternalRealtimeVoiceCapabilities = RealtimeVoiceProviderCapabilities & {
handlesAgentConsult?: boolean;
};
type OpenAIInternalRealtimeVoiceProviderApi = {
isBrowserSessionConfigured: (ctx: {
cfg?: RealtimeVoiceBrowserSessionCreateRequest["cfg"];
providerConfig: RealtimeVoiceProviderConfig;
}) => boolean;
resolveBrowserSessionCapabilities?: (ctx: {
cfg?: RealtimeVoiceBrowserSessionCreateRequest["cfg"];
providerConfig: RealtimeVoiceProviderConfig;
}) => OpenAIInternalRealtimeVoiceCapabilities;
cancelBrowserSession?: (
request: OpenAIInternalRealtimeBrowserSessionCreateRequest,
session: RealtimeVoiceBrowserSession,
) => Promise<void> | void;
};
type CodexRealtimeGlobalState = {
version: 1;
fallback?: {
capabilities: Partial<OpenAIInternalRealtimeVoiceCapabilities>;
isConfigured: () => boolean;
createBrowserSession: (
request: OpenAIInternalRealtimeBrowserSessionCreateRequest,
) => Promise<RealtimeVoiceBrowserSession>;
cancelBrowserSession: (session: RealtimeVoiceBrowserSession) => Promise<void> | void;
};
};
const CODEX_REALTIME_GLOBAL_STATE = Symbol.for("openclaw.codex.realtime-voice.v1");
const INTERNAL_REALTIME_VOICE_PROVIDER = Symbol.for("openclaw.internal.realtime-voice-provider.v1");
function readCodexRealtimeBrowserSessionFallback() {
const state = (
globalThis as typeof globalThis & {
[CODEX_REALTIME_GLOBAL_STATE]?: CodexRealtimeGlobalState;
}
)[CODEX_REALTIME_GLOBAL_STATE];
return state?.version === 1 ? state.fallback : undefined;
}
const codexFallbackBySession = new WeakMap<
RealtimeVoiceBrowserSession,
CodexRealtimeBrowserSessionFallback
>();
function resolveConfiguredCodexRealtimeFallback(
resolveFallback: () => CodexRealtimeBrowserSessionFallback | undefined,
): CodexRealtimeBrowserSessionFallback | undefined {
const fallback = resolveFallback();
return fallback?.isConfigured() === true ? fallback : undefined;
}
async function createOpenAIRealtimeBrowserSession(
req: OpenAIInternalRealtimeBrowserSessionCreateRequest,
resolveCodexFallback: () => CodexRealtimeBrowserSessionFallback | undefined,
req: RealtimeVoiceBrowserSessionCreateRequest,
): Promise<RealtimeVoiceBrowserSession> {
const config = normalizeProviderConfig(req.providerConfig);
if (config.azureEndpoint || config.azureDeployment) {
@@ -1632,21 +1562,6 @@ async function createOpenAIRealtimeBrowserSession(
cfg: req.cfg,
});
if (auth.status === "missing") {
// An authored Platform credential stays authoritative even when it cannot
// be resolved. Falling through would hide a broken key behind OAuth.
if (
!hasOpenAIRealtimePlatformAuthInput({
configuredApiKey: config.apiKey,
cfg: req.cfg,
})
) {
const fallback = resolveConfiguredCodexRealtimeFallback(resolveCodexFallback);
if (fallback) {
const session = await fallback.createBrowserSession(req);
codexFallbackBySession.set(session, fallback);
return session;
}
}
throw new Error(OPENAI_REALTIME_PLATFORM_AUTH_REQUIRED);
}
@@ -1710,22 +1625,7 @@ async function createOpenAIRealtimeBrowserSession(
};
}
async function cancelOpenAIRealtimeBrowserSession(
_req: OpenAIInternalRealtimeBrowserSessionCreateRequest,
session: RealtimeVoiceBrowserSession,
): Promise<void> {
const fallback = codexFallbackBySession.get(session);
codexFallbackBySession.delete(session);
await fallback?.cancelBrowserSession(session);
}
export function buildOpenAIRealtimeVoiceProvider(options?: {
resolveCodexRealtimeBrowserSessionFallback?: () =>
| CodexRealtimeBrowserSessionFallback
| undefined;
}): RealtimeVoiceProviderPlugin {
const resolveCodexFallback =
options?.resolveCodexRealtimeBrowserSessionFallback ?? readCodexRealtimeBrowserSessionFallback;
export function buildOpenAIRealtimeVoiceProvider(): RealtimeVoiceProviderPlugin {
const provider: RealtimeVoiceProviderPlugin = {
id: "openai",
label: "OpenAI Realtime Voice",
@@ -1766,57 +1666,8 @@ export function buildOpenAIRealtimeVoiceProvider(options?: {
azureApiVersion: config.azureApiVersion,
});
},
createBrowserSession: (req) =>
createOpenAIRealtimeBrowserSession(
req as OpenAIInternalRealtimeBrowserSessionCreateRequest,
resolveCodexFallback,
),
createBrowserSession: createOpenAIRealtimeBrowserSession,
};
const internalApi: OpenAIInternalRealtimeVoiceProviderApi = {
isBrowserSessionConfigured: ({ cfg, providerConfig }) => {
const config = normalizeProviderConfig(providerConfig);
if (
config.azureEndpoint ||
config.azureDeployment ||
hasOpenAIRealtimePlatformAuthInput({
configuredApiKey: config.apiKey,
cfg,
})
) {
return false;
}
return resolveConfiguredCodexRealtimeFallback(resolveCodexFallback) !== undefined;
},
resolveBrowserSessionCapabilities: ({ cfg, providerConfig }) => {
const config = normalizeProviderConfig(providerConfig);
if (
config.azureEndpoint ||
config.azureDeployment ||
hasOpenAIRealtimePlatformAuthInput({
configuredApiKey: config.apiKey,
cfg,
})
) {
return OPENAI_REALTIME_CAPABILITIES;
}
const fallback = resolveConfiguredCodexRealtimeFallback(resolveCodexFallback);
if (!fallback) {
return OPENAI_REALTIME_CAPABILITIES;
}
return {
...OPENAI_REALTIME_CAPABILITIES,
handlesAgentConsult: true,
supportsToolCalls: false,
supportsVideoFrames: false,
...fallback.capabilities,
};
},
cancelBrowserSession: cancelOpenAIRealtimeBrowserSession,
};
Object.defineProperty(provider, INTERNAL_REALTIME_VOICE_PROVIDER, {
configurable: true,
value: internalApi,
});
return provider;
}
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */