From 2e1607ee2fa8099a233df7437adee5f16a741905 Mon Sep 17 00:00:00 2001 From: jiayuhuang-openai Date: Wed, 15 Jul 2026 05:54:44 +0000 Subject: [PATCH] Add Frameless Bidi support for realtime conversations (#33261) ## What changed - Add realtime conversation version `v3`, which preserves the V1 Codex Voice behavior while using Frameless Bidi `delegation.*` events. - Translate audio, transcripts, handoffs, session context, and lifecycle events between the app server and the Frameless Bidi wire protocol. - Support `v3` over WebSocket and WebRTC, including the Frameless `/live` endpoint, session configuration, headers, and default model selection. - Update the app-server protocol schemas and documentation for the new version. ## Testing - Add unit coverage for Frameless event parsing, outbound messages, context chunking, URL construction, and call creation. - Add app-server end-to-end coverage for WebSocket delegation and WebRTC session startup. GitOrigin-RevId: 79a3307bc209e1a54582ebd2febc07c909fca016 --- .../schema/json/ClientRequest.json | 3 +- .../schema/json/ServerNotification.json | 3 +- .../codex_app_server_protocol.schemas.json | 3 +- .../codex_app_server_protocol.v2.schemas.json | 3 +- .../v2/ThreadRealtimeStartedNotification.json | 3 +- .../typescript/RealtimeConversationVersion.ts | 2 +- .../src/protocol/common.rs | 4 +- .../src/protocol/v2/realtime.rs | 6 +- codex-rs/app-server/README.md | 17 +- .../tests/suite/v2/realtime_conversation.rs | 285 ++++++++++++++++-- .../codex-api/src/endpoint/realtime_call.rs | 127 ++++++-- .../endpoint/realtime_websocket/methods.rs | 206 +++++++++++-- .../realtime_websocket/methods_common.rs | 136 ++++++--- .../methods_common_tests.rs | 88 ++++++ .../methods_frameless_bidi.rs | 84 ++++++ .../methods_frameless_bidi_tests.rs | 15 + .../src/endpoint/realtime_websocket/mod.rs | 2 + .../endpoint/realtime_websocket/protocol.rs | 33 ++ .../protocol_frameless_bidi.rs | 99 ++++++ .../protocol_frameless_bidi_tests.rs | 64 ++++ codex-rs/core/config.schema.json | 3 +- codex-rs/core/src/realtime_conversation.rs | 42 ++- .../core/src/realtime_conversation_tests.rs | 27 +- codex-rs/protocol/src/protocol.rs | 7 +- 24 files changed, 1116 insertions(+), 146 deletions(-) create mode 100644 codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common_tests.rs create mode 100644 codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi.rs create mode 100644 codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi_tests.rs create mode 100644 codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs create mode 100644 codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi_tests.rs diff --git a/codex-rs/app-server-protocol/schema/json/ClientRequest.json b/codex-rs/app-server-protocol/schema/json/ClientRequest.json index c1c43bd65639..3cfbffe5623d 100644 --- a/codex-rs/app-server-protocol/schema/json/ClientRequest.json +++ b/codex-rs/app-server-protocol/schema/json/ClientRequest.json @@ -2232,7 +2232,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/ServerNotification.json b/codex-rs/app-server-protocol/schema/json/ServerNotification.json index f0f4cbca48d1..abfdc6099216 100644 --- a/codex-rs/app-server-protocol/schema/json/ServerNotification.json +++ b/codex-rs/app-server-protocol/schema/json/ServerNotification.json @@ -3125,7 +3125,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json index ee7778e758f9..7ea4ea8965e1 100644 --- a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json +++ b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json @@ -14973,7 +14973,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json index 8d539a8ff3f7..ea42a95ee042 100644 --- a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json +++ b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json @@ -11330,7 +11330,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/v2/ThreadRealtimeStartedNotification.json b/codex-rs/app-server-protocol/schema/json/v2/ThreadRealtimeStartedNotification.json index 0beb774e7634..f61aa612d351 100644 --- a/codex-rs/app-server-protocol/schema/json/v2/ThreadRealtimeStartedNotification.json +++ b/codex-rs/app-server-protocol/schema/json/v2/ThreadRealtimeStartedNotification.json @@ -4,7 +4,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" } diff --git a/codex-rs/app-server-protocol/schema/typescript/RealtimeConversationVersion.ts b/codex-rs/app-server-protocol/schema/typescript/RealtimeConversationVersion.ts index cedc4bbe5255..81b8d3116e72 100644 --- a/codex-rs/app-server-protocol/schema/typescript/RealtimeConversationVersion.ts +++ b/codex-rs/app-server-protocol/schema/typescript/RealtimeConversationVersion.ts @@ -2,4 +2,4 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. -export type RealtimeConversationVersion = "v1" | "v2"; +export type RealtimeConversationVersion = "v1" | "v2" | "v3"; diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index b74d0f60843a..1602af1d3294 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -3304,7 +3304,7 @@ mod tests { prompt: Some(Some("You are on a call".to_string())), realtime_session_id: Some("sess_456".to_string()), transport: None, - version: Some(RealtimeConversationVersion::V1), + version: Some(RealtimeConversationVersion::V3), voice: Some(RealtimeVoice::Marin), }, }; @@ -3325,7 +3325,7 @@ mod tests { "prompt": "You are on a call", "realtimeSessionId": "sess_456", "transport": null, - "version": "v1", + "version": "v3", "voice": "marin" } }), diff --git a/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs b/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs index b78f390bec7c..dee420dc6596 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs @@ -80,9 +80,9 @@ pub struct ThreadRealtimeStartParams { /// Optional prefix added to automatic Codex response items when `codexResponsesAsItems` is true. #[ts(optional = nullable)] pub codex_response_item_prefix: Option, - /// Optional prefix added to automatic V1 Codex commentary sent with - /// `conversation.handoff.append` when `codexResponsesAsItems` is not true. Final answers are - /// sent without the prefix. + /// Optional prefix added to automatic V1 or V3 Codex commentary sent through the selected + /// Bidi handoff wire event when `codexResponsesAsItems` is not true. Final answers are sent + /// without the prefix. #[ts(optional = nullable)] pub codex_response_handoff_prefix: Option, /// Overrides the configured realtime model for this session only. diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 54b15849d540..35d240b8e168 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -172,7 +172,7 @@ Example with notification opt-out: - `thread/inject_items` — append raw Responses API items to a loaded thread’s model-visible history without starting a user turn; returns `{}` on success. - `turn/steer` — add user input to an already in-flight regular turn without starting a new turn; returns the active `turnId` that accepted the input. `clientUserMessageId` is optional; when supplied, the corresponding `userMessage` item echoes it as `clientId`. Review and manual compaction turns reject `turn/steer`. - `turn/interrupt` — request cancellation of an in-flight turn by `(thread_id, turn_id)`; success is an empty `{}` response and the turn finishes with `status: "interrupted"`. -- `thread/realtime/start` — start a thread-scoped realtime session (experimental); pass `outputModality: "text"` or `outputModality: "audio"` to choose model output, optionally pass `model` and, for websocket transport only, `version` to override configured realtime selection for this session only, and pass `includeStartupContext: false` to omit Codex's generated startup context. By default, automatic Codex text follows the protocol's speakable output path. Pass `clientManagedHandoffs: true` to disable automatic Codex response delivery so only the client's explicit append calls produce handoffs. Pass `codexResponsesAsItems: true` to send automatic Codex responses as realtime conversation items instead, and optionally pass `codexResponseItemPrefix` to prepend experiment instructions to those items. For V1 sessions, pass `codexResponseHandoffPrefix` while item mode is disabled to route automatic Codex commentary through `conversation.handoff.append` with that prefix; final answers remain unprefixed. Returns `{}` and streams `thread/realtime/*` notifications. Omit `transport` for the websocket transport, or pass `{ "type": "webrtc", "sdp": "..." }` to create an AVAS/v1 WebRTC session from a browser-generated SDP offer; the remote answer SDP is emitted as `thread/realtime/sdp`. Explicit `version: "v2"` requests are rejected for WebRTC. +- `thread/realtime/start` — start a thread-scoped realtime session (experimental); pass `outputModality: "text"` or `outputModality: "audio"` to choose model output, optionally pass `model` and `version` to override configured realtime selection for this session only, and pass `includeStartupContext: false` to omit Codex's generated startup context. Version `"v1"` uses legacy Bidi `conversation.handoff.*`, `"v2"` uses the Realtime Voice API, and `"v3"` preserves V1 Codex Voice behavior while using Frameless Bidi `delegation.*`. By default, automatic Codex text follows the protocol's speakable output path. Pass `clientManagedHandoffs: true` to disable automatic Codex response delivery so only the client's explicit append calls produce handoffs. Pass `codexResponsesAsItems: true` to send automatic Codex responses as realtime conversation items instead, and optionally pass `codexResponseItemPrefix` to prepend experiment instructions to those items. For V1 and V3 sessions, pass `codexResponseHandoffPrefix` while item mode is disabled to route automatic Codex commentary through the selected Bidi handoff wire event with that prefix; final answers remain unprefixed. Returns `{}` and streams `thread/realtime/*` notifications. Omit `transport` for the websocket transport, or pass `{ "type": "webrtc", "sdp": "..." }` to create a Bidi WebRTC session from a browser-generated SDP offer; the remote answer SDP is emitted as `thread/realtime/sdp`. Conversation `version: "v2"` requests remain unsupported for WebRTC. - `thread/realtime/appendAudio` — append an input audio chunk to the active realtime session (experimental); returns `{}`. - `thread/realtime/appendText` — append text input to the active realtime session with a required `role` of `user`, `developer`, or `assistant` (experimental); returns `{}`. Older clients that omit `role` default to `user`. - `thread/realtime/appendSpeech` — append text that the realtime model should speak to the user (experimental); returns `{}`. @@ -897,10 +897,9 @@ Omit `prompt` to use Codex's default realtime backend prompt. Send `prompt: null `prompt: ""` when the session should start without that default backend prompt. Clients may also pass `model` on `thread/realtime/start` to select a different realtime session configuration without changing thread or user config. -For websocket transport, clients may pass `version` to select the realtime -protocol version for this session only. WebRTC sessions always use AVAS with -realtime v1; omitting `version` uses v1, and explicitly passing `version: "v2"` -is rejected. +Clients may pass `version` to select the realtime protocol for this session +only. WebRTC uses AVAS and supports legacy Bidi `"v1"` or Frameless Bidi +`"v3"`; Realtime Voice `"v2"` is rejected for WebRTC. Pass `includeStartupContext: false` to skip Codex's startup context for this session while still using the selected backend prompt. Pass `clientManagedHandoffs: true` to suppress automatic Codex response handoffs @@ -911,10 +910,10 @@ Pass `codexResponsesAsItems: true` to inject automatic Codex responses with path. When using that mode, `codexResponseItemPrefix` can prepend short experiment instructions to each automatic Codex response item. Omit `codexResponsesAsItems`, or pass `false`, to preserve the default speakable -behavior. For V1 sessions, `codexResponseHandoffPrefix` instead routes automatic -Codex commentary through `conversation.handoff.append` and prepends the provided -text. Final answers remain unprefixed. Item mode takes precedence when -`codexResponsesAsItems` is true. +behavior. For V1 and V3 sessions, `codexResponseHandoffPrefix` instead routes +automatic Codex commentary through the selected Bidi handoff wire event and +prepends the provided text. Final answers remain unprefixed. Item mode takes +precedence when `codexResponsesAsItems` is true. Call `thread/realtime/appendText` to append app-provided realtime text items, or `thread/realtime/appendSpeech` when the app decides a realtime update should be diff --git a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs index 185631ea2fab..4dbcb6f39410 100644 --- a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs +++ b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs @@ -278,6 +278,16 @@ impl RealtimeE2eHarness { ) .mount(&main_loop_responses_server) .await; + Mock::given(method("POST")) + .and(path("/v1/live")) + .and(call_capture.clone()) + .respond_with( + ResponseTemplate::new(200) + .insert_header("Location", "/v1/live/rtc_e2e") + .set_body_string("v=answer\r\n"), + ) + .mount(&main_loop_responses_server) + .await; let realtime_server = start_websocket_server_with_headers(realtime_sideband.connections).await; @@ -321,8 +331,11 @@ impl RealtimeE2eHarness { async fn start_webrtc_realtime(&mut self, offer_sdp: &str) -> Result { self.start_webrtc_realtime_with_codex_response_routing( - offer_sdp, /*client_managed_handoffs*/ None, - /*codex_responses_as_items*/ None, /*codex_response_handoff_prefix*/ None, + offer_sdp, + /*client_managed_handoffs*/ None, + /*codex_responses_as_items*/ None, + /*codex_response_handoff_prefix*/ None, + RealtimeConversationVersion::V1, ) .await } @@ -336,6 +349,7 @@ impl RealtimeE2eHarness { /*client_managed_handoffs*/ None, /*codex_responses_as_items*/ Some(true), /*codex_response_handoff_prefix*/ None, + RealtimeConversationVersion::V1, ) .await } @@ -346,6 +360,7 @@ impl RealtimeE2eHarness { client_managed_handoffs: Option, codex_responses_as_items: Option, codex_response_handoff_prefix: Option<&str>, + version: RealtimeConversationVersion, ) -> Result { // Starts realtime through the public JSON-RPC method, then waits for the same client-visible // notifications a desktop app needs: started first, SDP answer second. @@ -368,7 +383,7 @@ impl RealtimeE2eHarness { transport: Some(ThreadRealtimeStartTransport::Webrtc { sdp: offer_sdp.to_string(), }), - version: None, + version: Some(version), voice: None, }) .await?; @@ -443,6 +458,38 @@ impl RealtimeE2eHarness { .await } + async fn start_frameless_bidi_realtime(&mut self) -> Result { + let start_request_id = self + .mcp + .send_thread_realtime_start_request(ThreadRealtimeStartParams { + thread_id: self.thread_id.clone(), + client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, + codex_response_item_prefix: None, + codex_response_handoff_prefix: None, + codex_responses_as_items: None, + model: None, + output_modality: RealtimeOutputModality::Audio, + include_startup_context: None, + prompt: Some(Some("backend prompt".to_string())), + realtime_session_id: None, + transport: None, + version: Some(RealtimeConversationVersion::V3), + voice: None, + }) + .await?; + let start_response: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + self.mcp + .read_stream_until_response_message(RequestId::Integer(start_request_id)), + ) + .await??; + let _: ThreadRealtimeStartResponse = to_response(start_response)?; + + self.read_notification::("thread/realtime/started") + .await + } + async fn read_notification(&mut self, method: &str) -> Result { read_notification(&mut self.mcp, method).await } @@ -569,6 +616,13 @@ fn session_updated(realtime_session_id: &str) -> Value { }) } +fn session_started(realtime_session_id: &str) -> Value { + json!({ + "type": "session.started", + "session": { "id": realtime_session_id, "instructions": "backend prompt" } + }) +} + fn v2_background_agent_tool_call(call_id: &str, prompt: &str) -> Value { json!({ "type": "conversation.item.done", @@ -1387,7 +1441,7 @@ async fn realtime_webrtc_start_emits_sdp_notification() -> Result<()> { transport: Some(ThreadRealtimeStartTransport::Webrtc { sdp: "v=offer\r\n".to_string(), }), - version: None, + version: Some(RealtimeConversationVersion::V1), voice: None, }) .await?; @@ -1532,6 +1586,7 @@ async fn webrtc_v1_start_posts_offer_returns_sdp_and_joins_sideband() -> Result< harness.call_capture.single_request(), "v=offer\r\n", v1_session_create_json(), + "/v1/realtime/calls?intent=quicksilver&architecture=avas", )?; let session_update = harness.sideband_outbound_request(/*request_index*/ 0).await; @@ -1554,6 +1609,75 @@ async fn webrtc_v1_start_posts_offer_returns_sdp_and_joins_sideband() -> Result< Ok(()) } +#[tokio::test] +async fn webrtc_v3_start_posts_live_session_and_joins_without_session_update() -> Result<()> { + skip_if_no_network!(Ok(())); + + let mut harness = RealtimeE2eHarness::new( + RealtimeTestVersion::V1, + no_main_loop_responses(), + realtime_sideband(vec![open_realtime_sideband_connection(vec![vec![]])]), + ) + .await?; + + let started = harness + .start_webrtc_realtime_with_codex_response_routing( + "v=offer\r\n", + /*client_managed_handoffs*/ None, + /*codex_responses_as_items*/ None, + /*codex_response_handoff_prefix*/ None, + RealtimeConversationVersion::V3, + ) + .await?; + assert_eq!( + started, + StartedWebrtcRealtime { + started: ThreadRealtimeStartedNotification { + thread_id: harness.thread_id.clone(), + realtime_session_id: Some(harness.thread_id.clone()), + version: RealtimeConversationVersion::V3, + }, + sdp: ThreadRealtimeSdpNotification { + thread_id: harness.thread_id.clone(), + sdp: "v=answer\r\n".to_string(), + }, + } + ); + + assert_call_create_multipart( + harness.call_capture.single_request(), + "v=offer\r\n", + r#"{"audio":{"output":{"voice":"cove"}},"delegation":{"type":"client"},"instructions":"backend prompt\n\nstartup context","model":"gpt-live-1-boulder-alpha"}"#, + "/v1/live", + )?; + assert!( + harness + .realtime_server + .wait_for_handshakes(/*expected*/ 1, DEFAULT_TIMEOUT) + .await, + "Frameless sideband should connect" + ); + assert_eq!( + harness.realtime_server.single_handshake().uri(), + "/v1/live/rtc_e2e" + ); + assert_eq!( + harness + .realtime_server + .single_handshake() + .header("openai-alpha") + .as_deref(), + Some("quicksilver=v2") + ); + assert!( + harness.realtime_server.single_connection().is_empty(), + "Frameless WebRTC sideband must not send a second session.update" + ); + + harness.shutdown().await; + Ok(()) +} + #[tokio::test] async fn webrtc_v1_default_automatic_output_uses_handoff_append() -> Result<()> { skip_if_no_network!(Ok(())); @@ -1633,6 +1757,7 @@ async fn webrtc_v1_client_managed_handoffs_disable_automatic_output() -> Result< /*client_managed_handoffs*/ Some(true), /*codex_responses_as_items*/ None, /*codex_response_handoff_prefix*/ None, + RealtimeConversationVersion::V1, ) .await?; assert_eq!(started.started.version, RealtimeConversationVersion::V1); @@ -1720,6 +1845,7 @@ async fn webrtc_v1_final_automatic_handoff_omits_silent_prefix() -> Result<()> { /*client_managed_handoffs*/ None, /*codex_responses_as_items*/ None, Some(RESPONSE_HANDOFF_PREFIX), + RealtimeConversationVersion::V1, ) .await?; assert_eq!(started.started.version, RealtimeConversationVersion::V1); @@ -1786,6 +1912,7 @@ async fn webrtc_v1_handoff_request_delegates_context_and_manual_append_speaks() harness.call_capture.single_request(), "v=offer\r\n", v1_session_create_json(), + "/v1/realtime/calls?intent=quicksilver&architecture=avas", )?; assert_v1_session_update(&harness.sideband_outbound_request(/*request_index*/ 0).await)?; @@ -2077,6 +2204,118 @@ async fn websocket_v2_assistant_output_without_handoff_reaches_realtime_context( Ok(()) } +#[tokio::test] +async fn websocket_v3_frameless_delegation_runs_codex_and_appends_context() -> Result<()> { + skip_if_no_network!(Ok(())); + + let mut harness = RealtimeE2eHarness::new( + RealtimeTestVersion::V1, + main_loop_responses(vec![create_final_assistant_message_sse_response( + "delegated from frameless", + )?]), + realtime_sideband(vec![realtime_sideband_connection(vec![ + vec![ + session_started("sess_frameless"), + json!({ + "type": "delegation.created", + "offset_ms": 100, + "item": { + "id": "delegation_frameless", + "type": "delegation", + "target": "client", + "content": [{ + "type": "input_text", + "text": "delegate from frameless" + }] + } + }), + ], + vec![], + vec![], + vec![], + ])]), + ) + .await?; + + let started = harness.start_frameless_bidi_realtime().await?; + assert_eq!(started.version, RealtimeConversationVersion::V3); + + let session_update = harness.sideband_outbound_request(/*request_index*/ 0).await; + assert_eq!(session_update["type"], "session.update"); + assert_eq!(session_update["session"]["delegation"]["type"], "client"); + assert_eq!( + harness.realtime_server.single_handshake().uri(), + "/v1/live?model=gpt-live-1-boulder-alpha" + ); + assert_eq!( + harness + .realtime_server + .single_handshake() + .header("openai-alpha") + .as_deref(), + Some("quicksilver=v2") + ); + + let turn_started = harness + .read_notification::("turn/started") + .await?; + assert_eq!(turn_started.thread_id, harness.thread_id); + let turn_completed = harness + .read_notification::("turn/completed") + .await?; + assert_eq!(turn_completed.thread_id, harness.thread_id); + + let requests = harness.main_loop_responses_requests().await?; + assert_eq!(requests.len(), 1); + assert!( + response_request_contains_text(&requests[0], "delegate from frameless"), + "delegated Responses request should contain Frameless input: {}", + requests[0] + ); + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 1).await, + json!({ + "type": "delegation.context.append", + "delegation_item_id": "delegation_frameless", + "content": [{ + "type": "input_text", + "text": "\"Agent Final Message\":\n\ndelegated from frameless" + }] + }) + ); + + harness + .append_speech(harness.thread_id.clone(), "standalone frameless update") + .await?; + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 2).await, + json!({ + "type": "session.context.append", + "content": [{ + "type": "input_text", + "text": "standalone frameless update" + }] + }) + ); + + harness + .append_text(harness.thread_id.clone(), "frameless text context") + .await?; + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 3).await, + json!({ + "type": "session.context.append", + "content": [{ + "type": "input_text", + "text": "frameless text context" + }] + }) + ); + + harness.shutdown().await; + Ok(()) +} + #[tokio::test] async fn websocket_v2_forwards_audio_and_text_between_client_and_sideband() -> Result<()> { skip_if_no_network!(Ok(())); @@ -2751,7 +2990,7 @@ async fn realtime_webrtc_start_surfaces_backend_error() -> Result<()> { transport: Some(ThreadRealtimeStartTransport::Webrtc { sdp: "v=offer\r\n".to_string(), }), - version: None, + version: Some(RealtimeConversationVersion::V1), voice: None, }) .await?; @@ -3066,13 +3305,14 @@ fn assert_v2_session_update(request: &Value) -> Result<()> { fn assert_call_create_multipart( request: WiremockRequest, offer_sdp: &str, - session: &str, + expected_session: &str, + expected_path_and_query: &str, ) -> Result<()> { - assert_eq!(request.url.path(), "/v1/realtime/calls"); - assert_eq!( - request.url.query(), - Some("intent=quicksilver&architecture=avas") - ); + let path_and_query = match request.url.query() { + Some(query) => format!("{}?{query}", request.url.path()), + None => request.url.path().to_string(), + }; + assert_eq!(path_and_query, expected_path_and_query); assert_eq!( request .headers @@ -3081,11 +3321,8 @@ fn assert_call_create_multipart( Some("multipart/form-data; boundary=codex-realtime-call-boundary") ); let body = String::from_utf8(request.body).context("multipart body should be utf-8")?; - let session = normalized_json_string(session)?; - assert_eq!( - body, - format!( - "--codex-realtime-call-boundary\r\n\ + let session_prefix = format!( + "--codex-realtime-call-boundary\r\n\ Content-Disposition: form-data; name=\"sdp\"\r\n\ Content-Type: application/sdp\r\n\ \r\n\ @@ -3093,11 +3330,17 @@ fn assert_call_create_multipart( --codex-realtime-call-boundary\r\n\ Content-Disposition: form-data; name=\"session\"\r\n\ Content-Type: application/json\r\n\ - \r\n\ - {session}\r\n\ - --codex-realtime-call-boundary--\r\n" - ) - ); + \r\n" + ); + let actual_session = body + .strip_prefix(&session_prefix) + .and_then(|body| body.strip_suffix("\r\n--codex-realtime-call-boundary--\r\n")) + .context("multipart body should contain one JSON session part")?; + let actual_session: Value = + serde_json::from_str(actual_session).context("session part should be valid JSON")?; + let expected_session: Value = serde_json::from_str(expected_session) + .context("expected session fixture should be valid JSON")?; + assert_eq!(actual_session, expected_session); Ok(()) } diff --git a/codex-rs/codex-api/src/endpoint/realtime_call.rs b/codex-rs/codex-api/src/endpoint/realtime_call.rs index af706091360f..37b753ee5c77 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_call.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_call.rs @@ -63,6 +63,17 @@ impl RealtimeCallClient { "realtime/calls" } + fn path_for_session(&self, event_parser: RealtimeEventParser) -> &'static str { + if self.uses_backend_request_shape() { + return Self::path(); + } + + match event_parser { + RealtimeEventParser::FramelessBidi => "live", + RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => Self::path(), + } + } + fn uses_backend_request_shape(&self) -> bool { self.session.provider().base_url.contains("/backend-api") } @@ -123,9 +134,11 @@ impl RealtimeCallClient { ) -> Result { trace!(target: "codex_api::realtime_websocket::wire", "realtime call request SDP: {sdp}"); // WebRTC can begin inference as soon as the peer connection comes up, so the initial - // session payload is sent with call creation. The sideband WebSocket still sends its normal - // session.update after it joins. + // session payload is sent with call creation. Legacy sidebands still send session.update + // after joining; Frameless sidebands attach to the session that is already running. validate_avas_session_config(&session_config)?; + let event_parser = session_config.event_parser; + let path = self.path_for_session(event_parser); let mut session = realtime_session_json(session_config)?; if let Some(session) = session.as_object_mut() { session.remove("id"); @@ -139,13 +152,13 @@ impl RealtimeCallClient { .map_err(|err| ApiError::Stream(format!("failed to encode realtime call: {err}")))?; let resp = self .session - .execute_with( - Method::POST, - Self::path(), - extra_headers, - Some(body), - configure_realtime_call_request, - ) + .execute_with(Method::POST, path, extra_headers, Some(body), |request| { + configure_realtime_call_request( + request, + event_parser, + /*uses_backend_request_shape*/ true, + ) + }) .await?; let sdp = decode_sdp_response(resp.body.as_ref())?; let call_id = decode_call_id_from_location(&resp.headers)?; @@ -172,11 +185,15 @@ impl RealtimeCallClient { .session .execute_with( Method::POST, - Self::path(), + path, extra_headers, /*body*/ None, |req| { - configure_realtime_call_request(req); + configure_realtime_call_request( + req, + event_parser, + /*uses_backend_request_shape*/ false, + ); req.headers.insert( CONTENT_TYPE, HeaderValue::from_static(MULTIPART_CONTENT_TYPE), @@ -193,15 +210,23 @@ impl RealtimeCallClient { } } -fn configure_realtime_call_request(request: &mut Request) { - append_query_pair(&mut request.url, "intent", "quicksilver"); - append_query_pair(&mut request.url, "architecture", "avas"); +fn configure_realtime_call_request( + request: &mut Request, + event_parser: RealtimeEventParser, + uses_backend_request_shape: bool, +) { + if event_parser == RealtimeEventParser::V1 + || (uses_backend_request_shape && event_parser == RealtimeEventParser::FramelessBidi) + { + append_query_pair(&mut request.url, "intent", "quicksilver"); + append_query_pair(&mut request.url, "architecture", "avas"); + } } fn validate_avas_session_config(session_config: &RealtimeSessionConfig) -> Result<(), ApiError> { - if session_config.event_parser != RealtimeEventParser::V1 { + if session_config.event_parser == RealtimeEventParser::RealtimeV2 { return Err(ApiError::InvalidRequest { - message: "AVAS realtime calls require realtime v1".to_string(), + message: "AVAS realtime calls require realtime v1 or v3".to_string(), }); } Ok(()) @@ -378,6 +403,13 @@ mod tests { } } + fn frameless_bidi_session_config(session_id: &str) -> RealtimeSessionConfig { + RealtimeSessionConfig { + event_parser: RealtimeEventParser::FramelessBidi, + ..realtime_session_config(session_id) + } + } + #[tokio::test] async fn sends_sdp_offer_as_raw_body() { let transport = CapturingTransport::new(); @@ -520,6 +552,36 @@ mod tests { ); } + #[tokio::test] + async fn sends_frameless_session_call_to_live_without_legacy_query_params() { + let transport = CapturingTransport::with_location("/v1/live/rtc_frameless"); + let client = RealtimeCallClient::new( + transport.clone(), + provider("https://api.openai.com/v1"), + Arc::new(DummyAuth), + ); + + let response = client + .create_with_session( + "v=offer\r\n".to_string(), + frameless_bidi_session_config("sess-api"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.call_id, "rtc_frameless"); + let request = transport.last_request.lock().unwrap().clone().unwrap(); + assert_eq!(request.method, Method::POST); + assert_eq!(request.url, "https://api.openai.com/v1/live"); + let Some(RequestBody::Raw(body)) = request.body else { + panic!("multipart body should be raw"); + }; + let body = std::str::from_utf8(&body).expect("multipart body should be utf-8"); + assert!(body.contains("\"model\":\"gpt-realtime\"")); + assert!(body.contains("\"delegation\":{\"type\":\"client\"}")); + assert!(!body.contains("\"id\":\"sess-api\"")); + } + #[tokio::test] async fn sends_session_call_with_avas_query_params() { let transport = CapturingTransport::new(); @@ -573,7 +635,7 @@ mod tests { assert_eq!( err.to_string(), - "invalid request: AVAS realtime calls require realtime v1" + "invalid request: AVAS realtime calls require realtime v1 or v3" ); assert!(transport.last_request.lock().unwrap().is_none()); } @@ -627,6 +689,37 @@ mod tests { ); } + #[tokio::test] + async fn sends_backend_frameless_session_call_to_realtime_calls() { + let transport = CapturingTransport::with_location("/v1/live/rtc_backend_frameless"); + let client = RealtimeCallClient::new( + transport.clone(), + provider("https://chatgpt.com/backend-api/codex"), + Arc::new(DummyAuth), + ); + + let response = client + .create_with_session( + "v=offer\r\n".to_string(), + frameless_bidi_session_config("sess-backend"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.call_id, "rtc_backend_frameless"); + let request = transport.last_request.lock().unwrap().clone().unwrap(); + assert_eq!(request.method, Method::POST); + assert_eq!( + request.url, + "https://chatgpt.com/backend-api/codex/realtime/calls?intent=quicksilver&architecture=avas" + ); + let Some(RequestBody::Json(body)) = request.body else { + panic!("backend request body should be JSON"); + }; + assert_eq!(body["session"]["delegation"]["type"], "client"); + assert!(body["session"].get("id").is_none()); + } + #[tokio::test] async fn errors_when_location_is_missing() { let transport = CapturingTransport::without_location(); diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs index 7b1a74e4aeb1..b720f44ceed2 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs @@ -1,8 +1,13 @@ use crate::endpoint::realtime_websocket::methods_common::conversation_function_call_output_message; +use crate::endpoint::realtime_websocket::methods_common::conversation_handoff_append_message; use crate::endpoint::realtime_websocket::methods_common::conversation_item_create_message; use crate::endpoint::realtime_websocket::methods_common::normalized_session_mode; -use crate::endpoint::realtime_websocket::methods_common::session_update_session; +use crate::endpoint::realtime_websocket::methods_common::session_update_message; +use crate::endpoint::realtime_websocket::methods_common::standalone_handoff_message; use crate::endpoint::realtime_websocket::methods_common::websocket_intent; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::context_append_chunks; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::delegation_context_append_message as frameless_delegation_context_append_message; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::session_context_append_message as frameless_session_context_append_message; use crate::endpoint::realtime_websocket::protocol::RealtimeAudioFrame; use crate::endpoint::realtime_websocket::protocol::RealtimeEvent; use crate::endpoint::realtime_websocket::protocol::RealtimeEventParser; @@ -25,6 +30,7 @@ use futures::StreamExt; use http::HeaderMap; use http::HeaderValue; use std::collections::HashMap; +use std::collections::VecDeque; use std::sync::Arc; use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; @@ -210,6 +216,7 @@ pub struct RealtimeWebsocketWriter { #[derive(Clone)] pub struct RealtimeWebsocketEvents { rx_message: async_channel::Receiver>, + pending_events: Arc>>, active_transcript: Arc>, event_parser: RealtimeEventParser, is_closed: Arc, @@ -277,6 +284,7 @@ impl RealtimeWebsocketConnection { }, events: RealtimeWebsocketEvents { rx_message, + pending_events: Arc::new(Mutex::new(VecDeque::new())), active_transcript: Arc::new(Mutex::new(ActiveTranscriptState::default())), event_parser, is_closed, @@ -287,8 +295,15 @@ impl RealtimeWebsocketConnection { impl RealtimeWebsocketWriter { pub async fn send_audio_frame(&self, frame: RealtimeAudioFrame) -> Result<(), ApiError> { - self.send_json(&RealtimeOutboundMessage::InputAudioBufferAppend { audio: frame.data }) - .await + let message = match self.event_parser { + RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => { + RealtimeOutboundMessage::InputAudioBufferAppend { audio: frame.data } + } + RealtimeEventParser::FramelessBidi => { + RealtimeOutboundMessage::InputAudioAppend { audio: frame.data } + } + }; + self.send_json(&message).await } pub async fn send_conversation_item_create( @@ -309,10 +324,24 @@ impl RealtimeWebsocketWriter { handoff_id: String, output_text: String, ) -> Result<(), ApiError> { - self.send_json(&RealtimeOutboundMessage::ConversationHandoffAppend { + self.send_json(&conversation_handoff_append_message( + self.event_parser, handoff_id, output_text, - }) + )) + .await + } + + pub async fn send_standalone_handoff( + &self, + handoff_id: String, + output_text: String, + ) -> Result<(), ApiError> { + self.send_json(&standalone_handoff_message( + self.event_parser, + handoff_id, + output_text, + )) .await } @@ -342,21 +371,34 @@ impl RealtimeWebsocketWriter { voice: RealtimeVoice, ) -> Result<(), ApiError> { let session_mode = normalized_session_mode(self.event_parser, session_mode); - let session = session_update_session( + let message = session_update_message( self.event_parser, instructions, session_mode, output_modality, voice, ); - self.send_json(&RealtimeOutboundMessage::SessionUpdate { session }) - .await + self.send_json(&message).await } pub async fn close(&self) -> Result<(), ApiError> { if self.is_closed.swap(true, Ordering::SeqCst) { return Ok(()); } + if self.event_parser == RealtimeEventParser::FramelessBidi { + let payload = + serde_json::to_string(&RealtimeOutboundMessage::SessionClose).map_err(|err| { + ApiError::Stream(format!("failed to encode realtime request: {err}")) + })?; + trace!(target: REALTIME_WIRE_LOG_TARGET, "realtime websocket request: {payload}"); + if let Err(err) = self.stream.send(Message::Text(payload.into())).await + && !matches!(err, WsError::ConnectionClosed | WsError::AlreadyClosed) + { + return Err(ApiError::Stream(format!( + "failed to close frameless realtime session: {err}" + ))); + } + } if let Err(err) = self.stream.close().await && !matches!(err, WsError::ConnectionClosed | WsError::AlreadyClosed) { @@ -368,6 +410,37 @@ impl RealtimeWebsocketWriter { } async fn send_json(&self, message: &RealtimeOutboundMessage) -> Result<(), ApiError> { + match message { + RealtimeOutboundMessage::DelegationContextAppend { + delegation_item_id, + content, + } => { + if let Some(content) = content.first() { + for chunk in context_append_chunks(&content.text) { + self.send_json_frame(&frameless_delegation_context_append_message( + delegation_item_id.clone(), + chunk, + )) + .await?; + } + return Ok(()); + } + } + RealtimeOutboundMessage::SessionContextAppend { content } => { + if let Some(content) = content.first() { + for chunk in context_append_chunks(&content.text) { + self.send_json_frame(&frameless_session_context_append_message(chunk)) + .await?; + } + return Ok(()); + } + } + _ => {} + } + self.send_json_frame(message).await + } + + async fn send_json_frame(&self, message: &RealtimeOutboundMessage) -> Result<(), ApiError> { let payload = serde_json::to_string(message) .map_err(|err| ApiError::Stream(format!("failed to encode realtime request: {err}")))?; debug!(?message, "realtime websocket request"); @@ -403,6 +476,10 @@ impl RealtimeWebsocketEvents { return Ok(None); } + if let Some(event) = self.pending_events.lock().await.pop_front() { + return Ok(Some(event)); + } + loop { let msg = match self.rx_message.recv().await { Ok(Ok(msg)) => msg, @@ -449,6 +526,24 @@ impl RealtimeWebsocketEvents { } } + async fn wait_for_session_started(&self) -> Result<(), ApiError> { + let Some(event) = self.next_event().await? else { + return Err(ApiError::Stream( + "frameless realtime session ended before session.started".to_string(), + )); + }; + match &event { + RealtimeEvent::SessionUpdated { .. } => { + self.pending_events.lock().await.push_back(event); + Ok(()) + } + RealtimeEvent::Error(message) => Err(ApiError::Stream(message.clone())), + _ => Err(ApiError::Stream( + "frameless realtime session received an event before session.started".to_string(), + )), + } + } + async fn update_active_transcript(&self, event: &mut RealtimeEvent) { let mut active_transcript = self.active_transcript.lock().await; match event { @@ -602,8 +697,14 @@ impl RealtimeWebsocketClient { config.event_parser, config.session_mode, )?; - self.connect_realtime_websocket_url(ws_url, config, extra_headers, default_headers) - .await + self.connect_realtime_websocket_url( + ws_url, + config, + extra_headers, + default_headers, + /*initialize_session*/ true, + ) + .await } pub async fn connect_webrtc_sideband( @@ -662,8 +763,14 @@ impl RealtimeWebsocketClient { config.session_mode, call_id, )?; - self.connect_realtime_websocket_url(ws_url, config, extra_headers, default_headers) - .await + self.connect_realtime_websocket_url( + ws_url, + config, + extra_headers, + default_headers, + /*initialize_session*/ false, + ) + .await } async fn connect_realtime_websocket_url( @@ -672,6 +779,7 @@ impl RealtimeWebsocketClient { config: RealtimeSessionConfig, extra_headers: HeaderMap, default_headers: HeaderMap, + initialize_session: bool, ) -> Result { ensure_rustls_crypto_provider(); @@ -708,19 +816,24 @@ impl RealtimeWebsocketClient { let (stream, rx_message) = WsStream::new(stream); let connection = RealtimeWebsocketConnection::new(stream, rx_message, config.event_parser); - debug!( - session_id = config.session_id.as_deref().unwrap_or(""), - "realtime websocket sending session.update" - ); - connection - .writer - .send_session_update( - config.instructions, - config.session_mode, - config.output_modality, - config.voice, - ) - .await?; + if initialize_session || config.event_parser != RealtimeEventParser::FramelessBidi { + debug!( + session_id = config.session_id.as_deref().unwrap_or(""), + "realtime websocket sending session.update" + ); + connection + .writer + .send_session_update( + config.instructions, + config.session_mode, + config.output_modality, + config.voice, + ) + .await?; + } + if initialize_session && config.event_parser == RealtimeEventParser::FramelessBidi { + connection.events.wait_for_session_started().await?; + } Ok(connection) } } @@ -770,7 +883,7 @@ fn websocket_url_from_api_url( let mut url = Url::parse(api_url) .map_err(|err| ApiError::Stream(format!("failed to parse realtime api_url: {err}")))?; - normalize_realtime_path(&mut url); + normalize_realtime_path(&mut url, event_parser); match url.scheme() { "ws" | "wss" => {} @@ -826,11 +939,31 @@ fn websocket_url_from_api_url_for_call( event_parser, session_mode, )?; - url.query_pairs_mut().append_pair("call_id", call_id); + match event_parser { + RealtimeEventParser::FramelessBidi => { + let path = format!("{}/{}", url.path().trim_end_matches('/'), call_id); + url.set_path(&path); + } + RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => { + url.query_pairs_mut().append_pair("call_id", call_id); + } + } Ok(url) } -fn normalize_realtime_path(url: &mut Url) { +fn normalize_realtime_path(url: &mut Url, event_parser: RealtimeEventParser) { + if event_parser == RealtimeEventParser::FramelessBidi { + let path = url.path().to_string(); + if path.is_empty() || path == "/" || path == "/v1" || path == "/v1/" { + url.set_path("/v1/live"); + } else if let Some(prefix) = path.trim_end_matches('/').strip_suffix("/realtime") { + url.set_path(&format!("{prefix}/live")); + } else if path.ends_with("/live/") { + url.set_path(path.trim_end_matches('/')); + } + return; + } + let path = url.path().to_string(); if path.is_empty() || path == "/" { url.set_path("/v1/realtime"); @@ -974,6 +1107,7 @@ mod tests { let (_tx_message, rx_message) = async_channel::unbounded(); let events = RealtimeWebsocketEvents { rx_message, + pending_events: Arc::new(Mutex::new(VecDeque::new())), active_transcript: Arc::new(Mutex::new(ActiveTranscriptState::default())), event_parser: RealtimeEventParser::V1, is_closed: Arc::new(AtomicBool::new(false)), @@ -1546,6 +1680,22 @@ mod tests { ); } + #[test] + fn frameless_websocket_url_rewrites_existing_realtime_path() { + let url = websocket_url_from_api_url( + "wss://example.com/v1/realtime?foo=bar", + /*query_params*/ None, + Some("snapshot"), + RealtimeEventParser::FramelessBidi, + RealtimeSessionMode::Conversational, + ) + .expect("build Frameless websocket url"); + assert_eq!( + url.as_str(), + "wss://example.com/v1/live?foo=bar&model=snapshot" + ); + } + #[test] fn websocket_url_v1_ignores_transcription_mode() { let url = websocket_url_from_api_url( diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common.rs index 131cf27a9447..76c9147ed04c 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common.rs @@ -1,3 +1,7 @@ +use crate::endpoint::realtime_websocket::methods_frameless_bidi::delegation_context_append_message as frameless_delegation_context_append_message; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::session_context_append_message as frameless_session_context_append_message; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::session_json as frameless_session_json; +use crate::endpoint::realtime_websocket::methods_frameless_bidi::session_update_message as frameless_session_update_message; use crate::endpoint::realtime_websocket::methods_v1::conversation_handoff_append_message as v1_conversation_handoff_append_message; use crate::endpoint::realtime_websocket::methods_v1::conversation_item_create_message as v1_conversation_item_create_message; use crate::endpoint::realtime_websocket::methods_v1::session_update_session as v1_session_update_session; @@ -6,13 +10,12 @@ use crate::endpoint::realtime_websocket::methods_v2::conversation_function_call_ use crate::endpoint::realtime_websocket::methods_v2::conversation_item_create_message as v2_conversation_item_create_message; use crate::endpoint::realtime_websocket::methods_v2::session_update_session as v2_session_update_session; use crate::endpoint::realtime_websocket::methods_v2::websocket_intent as v2_websocket_intent; -use crate::endpoint::realtime_websocket::protocol::RealtimeEventParser; use crate::endpoint::realtime_websocket::protocol::RealtimeOutboundMessage; use crate::endpoint::realtime_websocket::protocol::RealtimeOutputModality; use crate::endpoint::realtime_websocket::protocol::RealtimeSessionConfig; use crate::endpoint::realtime_websocket::protocol::RealtimeSessionMode; use crate::endpoint::realtime_websocket::protocol::RealtimeVoice; -use crate::endpoint::realtime_websocket::protocol::SessionUpdateSession; +use crate::endpoint::realtime_websocket::protocol::RealtimeWireAdapter; use codex_protocol::protocol::ConversationTextRole; use serde_json::Result as JsonResult; use serde_json::Value; @@ -22,74 +25,133 @@ pub(super) const REALTIME_AUDIO_SAMPLE_RATE: u32 = 24_000; const AGENT_FINAL_MESSAGE_PREFIX: &str = "\"Agent Final Message\":\n\n"; pub(super) fn normalized_session_mode( - event_parser: RealtimeEventParser, + wire_adapter: RealtimeWireAdapter, session_mode: RealtimeSessionMode, ) -> RealtimeSessionMode { - match event_parser { - RealtimeEventParser::V1 => RealtimeSessionMode::Conversational, - RealtimeEventParser::RealtimeV2 => session_mode, + match wire_adapter { + RealtimeWireAdapter::V1 | RealtimeWireAdapter::FramelessBidi => { + RealtimeSessionMode::Conversational + } + RealtimeWireAdapter::RealtimeV2 => session_mode, } } pub(super) fn conversation_item_create_message( - event_parser: RealtimeEventParser, + wire_adapter: RealtimeWireAdapter, text: String, role: ConversationTextRole, ) -> RealtimeOutboundMessage { - match event_parser { - RealtimeEventParser::V1 => v1_conversation_item_create_message(text, role), - RealtimeEventParser::RealtimeV2 => v2_conversation_item_create_message(text, role), + match wire_adapter { + RealtimeWireAdapter::V1 => v1_conversation_item_create_message(text, role), + RealtimeWireAdapter::FramelessBidi => frameless_session_context_append_message(text), + RealtimeWireAdapter::RealtimeV2 => v2_conversation_item_create_message(text, role), + } +} + +pub(super) fn conversation_handoff_append_message( + wire_adapter: RealtimeWireAdapter, + handoff_id: String, + output_text: String, +) -> RealtimeOutboundMessage { + match wire_adapter { + RealtimeWireAdapter::V1 => v1_conversation_handoff_append_message(handoff_id, output_text), + RealtimeWireAdapter::FramelessBidi => { + frameless_delegation_context_append_message(handoff_id, output_text) + } + RealtimeWireAdapter::RealtimeV2 => { + unreachable!("realtime v2 does not send conversation handoff output") + } + } +} + +pub(super) fn standalone_handoff_message( + wire_adapter: RealtimeWireAdapter, + handoff_id: String, + output_text: String, +) -> RealtimeOutboundMessage { + match wire_adapter { + RealtimeWireAdapter::V1 => v1_conversation_handoff_append_message(handoff_id, output_text), + RealtimeWireAdapter::FramelessBidi => frameless_session_context_append_message(output_text), + RealtimeWireAdapter::RealtimeV2 => { + unreachable!("realtime v2 does not send standalone handoff output") + } } } pub(super) fn conversation_function_call_output_message( - event_parser: RealtimeEventParser, + wire_adapter: RealtimeWireAdapter, call_id: String, output_text: String, ) -> RealtimeOutboundMessage { - match event_parser { - RealtimeEventParser::V1 => v1_conversation_handoff_append_message( + match wire_adapter { + RealtimeWireAdapter::V1 => v1_conversation_handoff_append_message( + call_id, + format!("{AGENT_FINAL_MESSAGE_PREFIX}{output_text}"), + ), + RealtimeWireAdapter::FramelessBidi => frameless_delegation_context_append_message( call_id, format!("{AGENT_FINAL_MESSAGE_PREFIX}{output_text}"), ), - RealtimeEventParser::RealtimeV2 => { + RealtimeWireAdapter::RealtimeV2 => { v2_conversation_function_call_output_message(call_id, output_text) } } } -pub(super) fn session_update_session( - event_parser: RealtimeEventParser, +pub(super) fn session_update_message( + wire_adapter: RealtimeWireAdapter, instructions: String, session_mode: RealtimeSessionMode, output_modality: RealtimeOutputModality, voice: RealtimeVoice, -) -> SessionUpdateSession { - let session_mode = normalized_session_mode(event_parser, session_mode); - match event_parser { - RealtimeEventParser::V1 => v1_session_update_session(instructions, voice), - RealtimeEventParser::RealtimeV2 => { - v2_session_update_session(instructions, session_mode, output_modality, voice) - } +) -> RealtimeOutboundMessage { + let session_mode = normalized_session_mode(wire_adapter, session_mode); + match wire_adapter { + RealtimeWireAdapter::V1 => RealtimeOutboundMessage::SessionUpdate { + session: v1_session_update_session(instructions, voice), + }, + RealtimeWireAdapter::FramelessBidi => frameless_session_update_message(instructions, voice), + RealtimeWireAdapter::RealtimeV2 => RealtimeOutboundMessage::SessionUpdate { + session: v2_session_update_session(instructions, session_mode, output_modality, voice), + }, } } pub fn session_update_session_json(config: RealtimeSessionConfig) -> JsonResult { - let mut session = session_update_session( - config.event_parser, - config.instructions, - config.session_mode, - config.output_modality, - config.voice, - ); - session.id = config.session_id; - session.model = config.model; - to_value(session) + match config.event_parser { + RealtimeWireAdapter::V1 | RealtimeWireAdapter::RealtimeV2 => { + let mut session = match config.event_parser { + RealtimeWireAdapter::V1 => { + v1_session_update_session(config.instructions, config.voice) + } + RealtimeWireAdapter::RealtimeV2 => v2_session_update_session( + config.instructions, + config.session_mode, + config.output_modality, + config.voice, + ), + RealtimeWireAdapter::FramelessBidi => unreachable!(), + }; + session.id = config.session_id; + session.model = config.model; + to_value(session) + } + RealtimeWireAdapter::FramelessBidi => Ok(frameless_session_json( + config.model, + config.instructions, + config.voice, + )), + } } -pub(super) fn websocket_intent(event_parser: RealtimeEventParser) -> Option<&'static str> { - match event_parser { - RealtimeEventParser::V1 => v1_websocket_intent(), - RealtimeEventParser::RealtimeV2 => v2_websocket_intent(), +pub(super) fn websocket_intent(wire_adapter: RealtimeWireAdapter) -> Option<&'static str> { + match wire_adapter { + RealtimeWireAdapter::V1 => v1_websocket_intent(), + RealtimeWireAdapter::FramelessBidi => None, + RealtimeWireAdapter::RealtimeV2 => v2_websocket_intent(), } } + +#[cfg(test)] +#[path = "methods_common_tests.rs"] +mod tests; diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common_tests.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common_tests.rs new file mode 100644 index 000000000000..8b187960c795 --- /dev/null +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_common_tests.rs @@ -0,0 +1,88 @@ +use super::conversation_function_call_output_message; +use super::conversation_handoff_append_message; +use super::standalone_handoff_message; +use crate::endpoint::realtime_websocket::protocol::RealtimeWireAdapter; +use pretty_assertions::assert_eq; +use serde_json::Value; +use serde_json::json; +use serde_json::to_value; + +#[test] +fn identical_handoff_output_encodes_for_each_bidi_wire_protocol() { + let legacy = conversation_handoff_append_message( + RealtimeWireAdapter::V1, + "handoff-123".to_string(), + "The result".to_string(), + ); + let frameless = conversation_handoff_append_message( + RealtimeWireAdapter::FramelessBidi, + "handoff-123".to_string(), + "The result".to_string(), + ); + + assert_eq!( + to_value(legacy).expect("legacy handoff should serialize"), + json!({ + "type": "conversation.handoff.append", + "handoff_id": "handoff-123", + "output_text": "The result", + }) + ); + assert_eq!( + to_value(frameless).expect("frameless handoff should serialize"), + json!({ + "type": "delegation.context.append", + "delegation_item_id": "handoff-123", + "content": [{"type": "input_text", "text": "The result"}], + }) + ); +} + +#[test] +fn standalone_handoff_uses_session_context_for_frameless() { + let legacy = standalone_handoff_message( + RealtimeWireAdapter::V1, + "codex".to_string(), + "Speak this".to_string(), + ); + let frameless = standalone_handoff_message( + RealtimeWireAdapter::FramelessBidi, + "codex".to_string(), + "Speak this".to_string(), + ); + + assert_eq!( + to_value(legacy).expect("legacy standalone handoff should serialize"), + json!({ + "type": "conversation.handoff.append", + "handoff_id": "codex", + "output_text": "Speak this", + }) + ); + assert_eq!( + to_value(frameless).expect("frameless standalone handoff should serialize"), + json!({ + "type": "session.context.append", + "content": [{"type": "input_text", "text": "Speak this"}], + }) + ); +} + +#[test] +fn completed_handoff_preserves_legacy_payload_text_in_frameless() { + let expected_text = Value::String("\"Agent Final Message\":\n\nDone".to_string()); + for wire_adapter in [RealtimeWireAdapter::V1, RealtimeWireAdapter::FramelessBidi] { + let encoded = to_value(conversation_function_call_output_message( + wire_adapter, + "handoff-123".to_string(), + "Done".to_string(), + )) + .expect("handoff output should serialize"); + let text = match wire_adapter { + RealtimeWireAdapter::V1 => &encoded["output_text"], + RealtimeWireAdapter::FramelessBidi => &encoded["content"][0]["text"], + RealtimeWireAdapter::RealtimeV2 => unreachable!(), + }; + assert_eq!(text, &expected_text); + } +} diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi.rs new file mode 100644 index 000000000000..2b5e01e99fae --- /dev/null +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi.rs @@ -0,0 +1,84 @@ +use crate::endpoint::realtime_websocket::protocol::FramelessContentType; +use crate::endpoint::realtime_websocket::protocol::FramelessInputTextContent; +use crate::endpoint::realtime_websocket::protocol::RealtimeOutboundMessage; +use crate::endpoint::realtime_websocket::protocol::RealtimeVoice; +use serde_json::Value; +use serde_json::json; + +const CONTEXT_APPEND_MAX_BYTES: usize = 500; + +pub(super) fn delegation_context_append_message( + delegation_item_id: String, + text: String, +) -> RealtimeOutboundMessage { + RealtimeOutboundMessage::DelegationContextAppend { + delegation_item_id, + content: input_text_content(text), + } +} + +pub(super) fn session_context_append_message(text: String) -> RealtimeOutboundMessage { + RealtimeOutboundMessage::SessionContextAppend { + content: input_text_content(text), + } +} + +pub(super) fn session_update_message( + instructions: String, + voice: RealtimeVoice, +) -> RealtimeOutboundMessage { + RealtimeOutboundMessage::FramelessSessionUpdate { + session: session_json(/*model*/ None, instructions, voice), + } +} + +pub(super) fn session_json( + model: Option, + instructions: String, + voice: RealtimeVoice, +) -> Value { + let mut session = json!({ + "instructions": instructions, + "audio": { + "output": { + "voice": voice, + }, + }, + "delegation": { + "type": "client", + }, + }); + if let Some(model) = model { + session["model"] = Value::String(model); + } + session +} + +fn input_text_content(text: String) -> Vec { + vec![FramelessInputTextContent { + r#type: FramelessContentType::InputText, + text, + }] +} + +pub(super) fn context_append_chunks(text: &str) -> Vec { + if text.len() <= CONTEXT_APPEND_MAX_BYTES { + return vec![text.to_string()]; + } + + let mut chunks = Vec::new(); + let mut start = 0; + while start < text.len() { + let mut end = (start + CONTEXT_APPEND_MAX_BYTES).min(text.len()); + while end > start && !text.is_char_boundary(end) { + end -= 1; + } + chunks.push(text[start..end].to_string()); + start = end; + } + chunks +} + +#[cfg(test)] +#[path = "methods_frameless_bidi_tests.rs"] +mod tests; diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi_tests.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi_tests.rs new file mode 100644 index 000000000000..d3651c65b65f --- /dev/null +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi_tests.rs @@ -0,0 +1,15 @@ +use super::CONTEXT_APPEND_MAX_BYTES; +use super::context_append_chunks; + +#[test] +fn context_append_chunks_preserve_text_within_wire_limit() { + for text in ["a".repeat(1_201), "🙂".repeat(200)] { + let chunks = context_append_chunks(&text); + assert_eq!(chunks.concat(), text); + assert!( + chunks + .iter() + .all(|chunk| chunk.len() <= CONTEXT_APPEND_MAX_BYTES) + ); + } +} diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs index 1fb49b2436f6..ac56dedac048 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs @@ -1,9 +1,11 @@ pub(crate) mod methods; mod methods_common; +mod methods_frameless_bidi; mod methods_v1; mod methods_v2; pub(crate) mod protocol; mod protocol_common; +mod protocol_frameless_bidi; mod protocol_v1; mod protocol_v2; diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs index 48f89b0d33dd..85145aa1a41e 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs @@ -1,3 +1,4 @@ +use crate::endpoint::realtime_websocket::protocol_frameless_bidi::parse_frameless_bidi_event; use crate::endpoint::realtime_websocket::protocol_v1::parse_realtime_event_v1; use crate::endpoint::realtime_websocket::protocol_v2::parse_realtime_event_v2; use codex_protocol::protocol::ConversationTextRole; @@ -12,9 +13,12 @@ use serde_json::Value; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum RealtimeEventParser { V1, + FramelessBidi, RealtimeV2, } +pub type RealtimeWireAdapter = RealtimeEventParser; + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum RealtimeSessionMode { Conversational, @@ -42,14 +46,42 @@ pub(super) enum RealtimeOutboundMessage { handoff_id: String, output_text: String, }, + #[serde(rename = "input_audio.append")] + InputAudioAppend { audio: String }, + #[serde(rename = "delegation.context.append")] + DelegationContextAppend { + delegation_item_id: String, + content: Vec, + }, + #[serde(rename = "session.context.append")] + SessionContextAppend { + content: Vec, + }, + #[serde(rename = "session.close")] + SessionClose, #[serde(rename = "response.create")] ResponseCreate, #[serde(rename = "session.update")] SessionUpdate { session: SessionUpdateSession }, + #[serde(rename = "session.update")] + FramelessSessionUpdate { session: Value }, #[serde(rename = "conversation.item.create")] ConversationItemCreate { item: ConversationItemPayload }, } +#[derive(Debug, Clone, Serialize)] +pub(super) struct FramelessInputTextContent { + #[serde(rename = "type")] + pub(super) r#type: FramelessContentType, + pub(super) text: String, +} + +#[derive(Debug, Clone, Copy, Serialize)] +#[serde(rename_all = "snake_case")] +pub(super) enum FramelessContentType { + InputText, +} + #[derive(Debug, Clone, Serialize)] pub(super) struct SessionUpdateSession { #[serde(skip_serializing_if = "Option::is_none")] @@ -219,6 +251,7 @@ pub(super) fn parse_realtime_event( ) -> Option { match event_parser { RealtimeEventParser::V1 => parse_realtime_event_v1(payload), + RealtimeEventParser::FramelessBidi => parse_frameless_bidi_event(payload), RealtimeEventParser::RealtimeV2 => parse_realtime_event_v2(payload), } } diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs new file mode 100644 index 000000000000..bb2154023b5e --- /dev/null +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs @@ -0,0 +1,99 @@ +use crate::endpoint::realtime_websocket::protocol_common::parse_error_event; +use crate::endpoint::realtime_websocket::protocol_common::parse_realtime_payload; +use crate::endpoint::realtime_websocket::protocol_common::parse_session_updated_event; +use codex_protocol::protocol::RealtimeAudioFrame; +use codex_protocol::protocol::RealtimeEvent; +use codex_protocol::protocol::RealtimeHandoffRequested; +use codex_protocol::protocol::RealtimeTranscriptDelta; +use codex_protocol::protocol::RealtimeTranscriptDone; +use serde_json::Value; +use tracing::debug; + +const DEFAULT_AUDIO_SAMPLE_RATE: u32 = 24_000; +const DEFAULT_AUDIO_CHANNELS: u16 = 1; + +pub(super) fn parse_frameless_bidi_event(payload: &str) -> Option { + let (parsed, message_type) = parse_realtime_payload(payload, "frameless bidi")?; + match message_type.as_str() { + "session.started" | "session.updated" => parse_session_updated_event(&parsed), + "output_audio.delta" => parse_output_audio_delta(&parsed), + "input_transcript.added" => { + parse_transcript_item(&parsed).map(RealtimeEvent::InputTranscriptDelta) + } + "output_transcript.added" => { + parse_transcript_item(&parsed).map(RealtimeEvent::OutputTranscriptDelta) + } + "turn.done" => parse_turn_done(&parsed), + "delegation.created" => parse_delegation_created(&parsed), + "error" => parse_error_event(&parsed), + _ => { + debug!( + "received unsupported frameless bidi event type: {message_type}, data: {payload}" + ); + None + } + } +} + +fn parse_output_audio_delta(parsed: &Value) -> Option { + Some(RealtimeEvent::AudioOut(RealtimeAudioFrame { + data: parsed.get("audio").and_then(Value::as_str)?.to_string(), + sample_rate: DEFAULT_AUDIO_SAMPLE_RATE, + num_channels: DEFAULT_AUDIO_CHANNELS, + samples_per_channel: None, + item_id: None, + })) +} + +fn parse_transcript_item(parsed: &Value) -> Option { + parsed + .get("item") + .and_then(Value::as_object) + .and_then(|item| item.get("text")) + .and_then(Value::as_str) + .map(str::to_string) + .map(|delta| RealtimeTranscriptDelta { delta }) +} + +fn parse_turn_done(parsed: &Value) -> Option { + let turn = parsed.get("turn")?.as_object()?; + let role = turn.get("role").and_then(Value::as_str)?; + let text = turn + .get("transcript") + .and_then(Value::as_str) + .map(str::to_string)?; + let done = RealtimeTranscriptDone { text }; + match role { + "user" => Some(RealtimeEvent::InputTranscriptDone(done)), + "assistant" => Some(RealtimeEvent::OutputTranscriptDone(done)), + _ => None, + } +} + +fn parse_delegation_created(parsed: &Value) -> Option { + let item = parsed.get("item")?.as_object()?; + if item.get("type").and_then(Value::as_str) != Some("delegation") + || item.get("target").and_then(Value::as_str) != Some("client") + { + return None; + } + let item_id = item.get("id").and_then(Value::as_str)?.to_string(); + let input_transcript = item + .get("content") + .and_then(Value::as_array)? + .iter() + .filter(|content| content.get("type").and_then(Value::as_str) == Some("input_text")) + .filter_map(|content| content.get("text").and_then(Value::as_str)) + .collect::(); + + Some(RealtimeEvent::HandoffRequested(RealtimeHandoffRequested { + handoff_id: item_id.clone(), + item_id, + input_transcript, + active_transcript: Vec::new(), + })) +} + +#[cfg(test)] +#[path = "protocol_frameless_bidi_tests.rs"] +mod tests; diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi_tests.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi_tests.rs new file mode 100644 index 000000000000..7b95c0a79a3a --- /dev/null +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi_tests.rs @@ -0,0 +1,64 @@ +use super::parse_frameless_bidi_event; +use crate::endpoint::realtime_websocket::protocol_v1::parse_realtime_event_v1; +use codex_protocol::protocol::RealtimeEvent; +use codex_protocol::protocol::RealtimeHandoffRequested; + +#[test] +fn legacy_and_frameless_delegations_decode_to_the_same_handoff() { + let expected = Some(RealtimeEvent::HandoffRequested(RealtimeHandoffRequested { + handoff_id: "handoff-123".to_string(), + item_id: "handoff-123".to_string(), + input_transcript: "check the weather".to_string(), + active_transcript: Vec::new(), + })); + let legacy = r#"{ + "type": "conversation.handoff.requested", + "handoff_id": "handoff-123", + "item_id": "handoff-123", + "input_transcript": "check the weather" + }"#; + let frameless = r#"{ + "type": "delegation.created", + "offset_ms": 1000, + "item": { + "id": "handoff-123", + "type": "delegation", + "target": "client", + "content": [{"type": "input_text", "text": "check the weather"}] + } + }"#; + + assert_eq!(parse_realtime_event_v1(legacy), expected); + assert_eq!(parse_frameless_bidi_event(frameless), expected); +} + +#[test] +fn frameless_transcript_and_audio_events_reuse_existing_internal_events() { + let input = r#"{ + "type": "input_transcript.added", + "item": {"id": "input-1", "type": "input_transcript", "text": "hello"} + }"#; + let done = r#"{ + "type": "turn.done", + "turn": {"id": "turn-1", "role": "user", "transcript": "hello"} + }"#; + let audio = r#"{ + "type": "output_audio.delta", + "audio": "AAE=", + "start_ms": 0, + "end_ms": 100 + }"#; + + assert!(matches!( + parse_frameless_bidi_event(input), + Some(RealtimeEvent::InputTranscriptDelta(_)) + )); + assert!(matches!( + parse_frameless_bidi_event(done), + Some(RealtimeEvent::InputTranscriptDone(_)) + )); + assert!(matches!( + parse_frameless_bidi_event(audio), + Some(RealtimeEvent::AudioOut(_)) + )); +} diff --git a/codex-rs/core/config.schema.json b/codex-rs/core/config.schema.json index b55512b1691a..89d6bad9b470 100644 --- a/codex-rs/core/config.schema.json +++ b/codex-rs/core/config.schema.json @@ -2640,7 +2640,8 @@ "RealtimeConversationVersion": { "enum": [ "v1", - "v2" + "v2", + "v3" ], "type": "string" }, diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index 3915c8190dda..b1b86b7b622b 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -74,6 +74,7 @@ const REALTIME_STARTUP_CONTEXT_TOKEN_BUDGET: usize = 5_300; const REALTIME_ASSISTANT_OUTPUT_TOKEN_BUDGET: usize = 1_000; const STANDALONE_HANDOFF_ID: &str = "codex"; const DEFAULT_REALTIME_MODEL: &str = "gpt-realtime-1.5"; +const DEFAULT_FRAMELESS_REALTIME_MODEL: &str = "gpt-live-1-boulder-alpha"; pub(crate) const REALTIME_USER_TEXT_PREFIX: &str = "[USER] "; pub(crate) const REALTIME_BACKEND_TEXT_PREFIX: &str = "[BACKEND] "; const REALTIME_V2_HANDOFF_COMPLETE_ACKNOWLEDGEMENT: &str = @@ -325,7 +326,7 @@ impl RealtimeConversationManager { } = start; let event_parser = session_config.event_parser; let session_kind = match event_parser { - RealtimeEventParser::V1 => RealtimeSessionKind::V1, + RealtimeEventParser::V1 | RealtimeEventParser::FramelessBidi => RealtimeSessionKind::V1, RealtimeEventParser::RealtimeV2 => RealtimeSessionKind::V2, }; @@ -796,6 +797,7 @@ async fn prepare_realtime_start( let session_config = build_realtime_session_config(sess, ¶ms, version, configured_voice).await?; let requested_realtime_session_id = session_config.session_id.clone(); + let event_parser = session_config.event_parser; let originator = sess.originator().await; let extra_headers = match transport { ConversationStartTransport::Websocket => { @@ -803,7 +805,7 @@ async fn prepare_realtime_start( realtime_request_headers( requested_realtime_session_id.as_deref(), Some(realtime_api_key.as_str()), - version, + event_parser, originator.as_str(), )? } @@ -811,7 +813,7 @@ async fn prepare_realtime_start( realtime_request_headers( requested_realtime_session_id.as_deref(), /*api_key*/ None, - version, + event_parser, originator.as_str(), )? } @@ -836,9 +838,9 @@ fn validate_avas_webrtc_start( version: RealtimeWsVersion, session_type: RealtimeWsMode, ) -> CodexResult<()> { - if version != RealtimeWsVersion::V1 { + if version == RealtimeWsVersion::V2 { return Err(CodexErr::InvalidRequest( - "AVAS realtime calls require realtime v1".to_string(), + "AVAS realtime calls require realtime v1 or v3".to_string(), )); } if session_type != RealtimeWsMode::Conversational { @@ -883,13 +885,17 @@ pub(crate) async fn build_realtime_session_config( .model .clone() .or_else(|| config.experimental_realtime_ws_model.clone()) - .unwrap_or_else(|| DEFAULT_REALTIME_MODEL.to_string()), + .unwrap_or_else(|| match version { + RealtimeWsVersion::V1 | RealtimeWsVersion::V2 => DEFAULT_REALTIME_MODEL.to_string(), + RealtimeWsVersion::V3 => DEFAULT_FRAMELESS_REALTIME_MODEL.to_string(), + }), ); let event_parser = match version { RealtimeWsVersion::V1 => RealtimeEventParser::V1, RealtimeWsVersion::V2 => RealtimeEventParser::RealtimeV2, + RealtimeWsVersion::V3 => RealtimeEventParser::FramelessBidi, }; - if version == RealtimeWsVersion::V1 + if version != RealtimeWsVersion::V2 && matches!(params.output_modality, RealtimeOutputModality::Text) { return Err(CodexErr::InvalidRequest( @@ -928,7 +934,7 @@ pub(crate) async fn build_realtime_session_config( fn default_realtime_voice(version: RealtimeWsVersion) -> RealtimeVoice { let voices = RealtimeVoicesList::builtin(); match version { - RealtimeWsVersion::V1 => voices.default_v1, + RealtimeWsVersion::V1 | RealtimeWsVersion::V3 => voices.default_v1, RealtimeWsVersion::V2 => voices.default_v2, } } @@ -956,7 +962,7 @@ fn realtime_backend_item(text: String, prefix: Option<&str>) -> String { fn validate_realtime_voice(version: RealtimeWsVersion, voice: RealtimeVoice) -> CodexResult<()> { let voices = RealtimeVoicesList::builtin(); let allowed = match version { - RealtimeWsVersion::V1 => &voices.v1, + RealtimeWsVersion::V1 | RealtimeWsVersion::V3 => &voices.v1, RealtimeWsVersion::V2 => &voices.v2, }; if allowed.contains(&voice) { @@ -966,6 +972,7 @@ fn validate_realtime_voice(version: RealtimeWsVersion, voice: RealtimeVoice) -> let version = match version { RealtimeWsVersion::V1 => "v1", RealtimeWsVersion::V2 => "v2", + RealtimeWsVersion::V3 => "v3", }; let allowed = allowed .iter() @@ -1199,13 +1206,19 @@ fn realtime_api_key(auth: Option<&CodexAuth>, provider: &ModelProviderInfo) -> C fn realtime_request_headers( realtime_session_id: Option<&str>, api_key: Option<&str>, - version: RealtimeWsVersion, + event_parser: RealtimeEventParser, originator: &str, ) -> CodexResult> { let mut headers = HeaderMap::new(); - if version == RealtimeWsVersion::V1 { - headers.insert("openai-alpha", HeaderValue::from_static("quicksilver=v1")); + match event_parser { + RealtimeEventParser::V1 => { + headers.insert("openai-alpha", HeaderValue::from_static("quicksilver=v1")); + } + RealtimeEventParser::FramelessBidi => { + headers.insert("openai-alpha", HeaderValue::from_static("quicksilver=v2")); + } + RealtimeEventParser::RealtimeV2 => {} } if let Some(realtime_session_id) = realtime_session_id @@ -1462,11 +1475,10 @@ async fn handle_handoff_output( let handoff_output = handoff_output.context("handoff output channel closed")?; let result = match event_parser { - RealtimeEventParser::V1 => match handoff_output { + RealtimeEventParser::V1 | RealtimeEventParser::FramelessBidi => match handoff_output { RealtimeOutbound::StandaloneHandoff { text } => { - // TODO(guinness): Use the new client event for standalone handoffs once the API changes are complete. writer - .send_conversation_handoff_append(STANDALONE_HANDOFF_ID.to_string(), text) + .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text) .await } RealtimeOutbound::HandoffUpdate { handoff_id, text } diff --git a/codex-rs/core/src/realtime_conversation_tests.rs b/codex-rs/core/src/realtime_conversation_tests.rs index 51425e4c3fc6..412ce1f8ab6a 100644 --- a/codex-rs/core/src/realtime_conversation_tests.rs +++ b/codex-rs/core/src/realtime_conversation_tests.rs @@ -5,7 +5,7 @@ use super::realtime_request_headers; use super::realtime_text_from_handoff_request; use super::wrap_realtime_delegation_input; use async_channel::bounded; -use codex_config::config_toml::RealtimeWsVersion; +use codex_api::RealtimeEventParser; use codex_protocol::protocol::RealtimeHandoffRequested; use codex_protocol::protocol::RealtimeTranscriptEntry; use pretty_assertions::assert_eq; @@ -152,7 +152,7 @@ fn uses_quicksilver_alpha_header_for_realtime_v1() { let headers = realtime_request_headers( Some("session_1"), Some("sk-test"), - RealtimeWsVersion::V1, + RealtimeEventParser::V1, "codex_work_desktop", ) .expect("headers") @@ -171,7 +171,7 @@ fn omits_quicksilver_alpha_header_for_realtime_v2() { let headers = realtime_request_headers( Some("session_1"), Some("sk-test"), - RealtimeWsVersion::V2, + RealtimeEventParser::RealtimeV2, "codex_work_desktop", ) .expect("headers") @@ -180,6 +180,25 @@ fn omits_quicksilver_alpha_header_for_realtime_v2() { assert!(headers.get("openai-alpha").is_none()); } +#[test] +fn uses_frameless_alpha_header_for_realtime_v3() { + let headers = realtime_request_headers( + Some("session_1"), + Some("sk-test"), + RealtimeEventParser::FramelessBidi, + "codex_work_desktop", + ) + .expect("headers") + .expect("headers"); + + assert_eq!( + headers + .get("openai-alpha") + .and_then(|value| value.to_str().ok()), + Some("quicksilver=v2") + ); +} + #[test] fn realtime_headers_include_only_non_default_originator() { let default_originator = codex_login::default_client::originator(); @@ -190,7 +209,7 @@ fn realtime_headers_include_only_non_default_originator() { let headers = realtime_request_headers( Some("session_1"), Some("sk-test"), - RealtimeWsVersion::V2, + RealtimeEventParser::RealtimeV2, originator, ) .expect("headers") diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 9f1139ea276c..1bb719fa7223 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -211,9 +211,9 @@ pub struct ConversationStartParams { pub codex_responses_as_items: bool, /// Optional prefix added to automatic Codex response items when `codex_responses_as_items` is set. pub codex_response_item_prefix: Option, - /// Optional prefix added to automatic V1 Codex commentary sent with - /// `conversation.handoff.append` when `codex_responses_as_items` is not set. Final answers are - /// sent without the prefix. + /// Optional prefix added to automatic V1 or V3 Codex commentary sent through the selected + /// Bidi handoff wire event when `codex_responses_as_items` is not set. Final answers are sent + /// without the prefix. pub codex_response_handoff_prefix: Option, /// Overrides the configured realtime model for this session only. pub model: Option, @@ -1621,6 +1621,7 @@ pub enum RealtimeConversationVersion { V1, #[default] V2, + V3, } #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, JsonSchema, TS)]