diff --git a/codex-rs/app-server-protocol/schema/json/ClientRequest.json b/codex-rs/app-server-protocol/schema/json/ClientRequest.json index a85d272cd3..77f601bad7 100644 --- a/codex-rs/app-server-protocol/schema/json/ClientRequest.json +++ b/codex-rs/app-server-protocol/schema/json/ClientRequest.json @@ -288,6 +288,14 @@ ], "type": "object" }, + "CodexResponseHandoffMode": { + "enum": [ + "thinking", + "commentary", + "bemTags" + ], + "type": "string" + }, "CollaborationMode": { "description": "Collaboration mode for a Codex session.", "properties": { 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 3bb5973e8b..1f2cd11dff 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 @@ -7416,6 +7416,14 @@ } ] }, + "CodexResponseHandoffMode": { + "enum": [ + "thinking", + "commentary", + "bemTags" + ], + "type": "string" + }, "CollabAgentState": { "properties": { "message": { 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 c5f3a63f27..f77838ea66 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 @@ -3609,6 +3609,14 @@ } ] }, + "CodexResponseHandoffMode": { + "enum": [ + "thinking", + "commentary", + "bemTags" + ], + "type": "string" + }, "CollabAgentState": { "properties": { "message": { diff --git a/codex-rs/app-server-protocol/schema/typescript/CodexResponseHandoffMode.ts b/codex-rs/app-server-protocol/schema/typescript/CodexResponseHandoffMode.ts new file mode 100644 index 0000000000..3eb90dad2c --- /dev/null +++ b/codex-rs/app-server-protocol/schema/typescript/CodexResponseHandoffMode.ts @@ -0,0 +1,5 @@ +// GENERATED CODE! DO NOT MODIFY BY HAND! + +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. + +export type CodexResponseHandoffMode = "thinking" | "commentary" | "bemTags"; diff --git a/codex-rs/app-server-protocol/schema/typescript/index.ts b/codex-rs/app-server-protocol/schema/typescript/index.ts index 77f4af9466..e3e88adbf8 100644 --- a/codex-rs/app-server-protocol/schema/typescript/index.ts +++ b/codex-rs/app-server-protocol/schema/typescript/index.ts @@ -10,6 +10,7 @@ export type { AutoCompactTokenLimitScope } from "./AutoCompactTokenLimitScope"; export type { ClientInfo } from "./ClientInfo"; export type { ClientNotification } from "./ClientNotification"; export type { ClientRequest } from "./ClientRequest"; +export type { CodexResponseHandoffMode } from "./CodexResponseHandoffMode"; export type { CollaborationMode } from "./CollaborationMode"; export type { ContentItem } from "./ContentItem"; export type { ConversationGitInfo } from "./ConversationGitInfo"; diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index 176ab77eff..4c906d134d 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -1764,6 +1764,7 @@ mod tests { use codex_protocol::config_types::MultiAgentMode; use codex_protocol::models::BUILT_IN_PERMISSION_PROFILE_READ_ONLY; use codex_protocol::parse_command::ParsedCommand; + use codex_protocol::protocol::CodexResponseHandoffMode; use codex_protocol::protocol::RealtimeConversationVersion; use codex_protocol::protocol::RealtimeOutputModality; use codex_protocol::protocol::RealtimeVoice; @@ -3390,6 +3391,7 @@ mod tests { flush_transcript_tail_on_session_end: Some(true), codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: Some(CodexResponseHandoffMode::BemTags), thread_id: "thr_123".to_string(), model: Some("realtime-treatment-model".to_string()), output_modality: RealtimeOutputModality::Audio, @@ -3411,6 +3413,7 @@ mod tests { "flushTranscriptTailOnSessionEnd": true, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, + "codexResponseHandoffMode": "bemTags", "model": "realtime-treatment-model", "outputModality": "audio", "includeStartupContext": false, @@ -3435,6 +3438,7 @@ mod tests { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: "thr_123".to_string(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -3456,6 +3460,7 @@ mod tests { "flushTranscriptTailOnSessionEnd": null, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, + "codexResponseHandoffMode": null, "model": null, "outputModality": "audio", "includeStartupContext": null, @@ -3475,6 +3480,7 @@ mod tests { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: "thr_123".to_string(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -3496,6 +3502,7 @@ mod tests { "flushTranscriptTailOnSessionEnd": null, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, + "codexResponseHandoffMode": null, "model": null, "outputModality": "audio", "includeStartupContext": null, @@ -3715,6 +3722,7 @@ mod tests { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: "thr_123".to_string(), model: None, output_modality: RealtimeOutputModality::Audio, 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 f375a5ea2e..9333f9d356 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs @@ -1,3 +1,4 @@ +use codex_protocol::protocol::CodexResponseHandoffMode; use codex_protocol::protocol::ConversationTextRole; use codex_protocol::protocol::RealtimeAudioFrame as CoreRealtimeAudioFrame; use codex_protocol::protocol::RealtimeConversationVersion; @@ -74,12 +75,18 @@ pub struct ThreadRealtimeStartParams { /// TODO: Remove this rollout knob once transcript-tail flushing is always enabled. #[ts(optional = nullable)] pub flush_transcript_tail_on_session_end: Option, + // TODO: Remove this experiment-only delivery path after response-item testing is complete. /// Sends automatic Codex responses as realtime conversation items instead of handoff appends. #[ts(optional = nullable)] pub codex_responses_as_items: Option, + // TODO: Remove this experiment-only prefix with `codex_responses_as_items`. /// Optional prefix added to automatic Codex response items when `codexResponsesAsItems` is true. #[ts(optional = nullable)] pub codex_response_item_prefix: Option, + /// Selects how automatic Codex responses are routed in Frameless Bidi sessions. Omitted values + /// default to `thinking`. Realtime V1 and V2 ignore this setting. + #[ts(optional = nullable)] + pub codex_response_handoff_mode: Option, /// Overrides the configured realtime model for this session only. #[ts(optional = nullable)] pub model: Option, diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 1b6c6ca6c1..0b10de1109 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 `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.*`. V1 and V3 send commentary as raw handoff appends and label final or phase-less agent output with `"Agent Final Message":`; V3 streams these appends as Codex produces text deltas. 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. 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/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.*`. For V3 automatic Codex text, `codexResponseHandoffMode` accepts `"thinking"` (the default; all output uses channel-less thinking appends), `"commentary"` (all output uses the commentary channel), or `"bemTags"` (the raw BEM envelope selects the API channel: BEM `analysis` and `commentary` use `commentary`, while BEM `final` and unparsable output use `speakable`). The BEM envelope remains in the appended text for the frontend model to interpret. V1 and V2 ignore this setting. V3 handoffs do not prepend the legacy `"Agent Final Message"` label. 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. 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 `{}`. @@ -923,11 +923,16 @@ 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. In that default mode, V1 and V3 append commentary without a prefix -and label final or phase-less agent output with `"Agent Final Message":`. V3 -emits each text delta as a `delegation.context.append` instead of waiting for -the completed agent message. Older clients may continue to send the removed -`codexResponseHandoffPrefix` field; the server ignores unknown request fields. +behavior. In V3, automatic handoffs default to +`codexResponseHandoffMode: "thinking"`, which omits the context append `channel` +for every automatic response. Pass `"commentary"` to route every response to +commentary, or `"bemTags"` to route BEM commentary tags to `commentary`, final +tags to `speakable`, and analysis tags to `commentary`. Unparsable BEM output +falls back to `speakable`. BEM routing reads the raw envelope and preserves it +in the appended text for the frontend model. This +setting has no effect on V1 or V2. V3 handoffs never prepend the legacy `"Agent Final Message"` label. Older +clients may continue to send the removed `codexResponseHandoffPrefix` field; the +server ignores unknown request fields. 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/src/request_processors/turn_processor.rs b/codex-rs/app-server/src/request_processors/turn_processor.rs index 31d00cd7bd..1b919775eb 100644 --- a/codex-rs/app-server/src/request_processors/turn_processor.rs +++ b/codex-rs/app-server/src/request_processors/turn_processor.rs @@ -1077,6 +1077,7 @@ impl TurnRequestProcessor { .unwrap_or(false), codex_responses_as_items: params.codex_responses_as_items.unwrap_or(false), codex_response_item_prefix: params.codex_response_item_prefix, + codex_response_handoff_mode: params.codex_response_handoff_mode.unwrap_or_default(), model: params.model, output_modality: params.output_modality, include_startup_context: params.include_startup_context.unwrap_or(true), diff --git a/codex-rs/app-server/tests/suite/v2/experimental_api.rs b/codex-rs/app-server/tests/suite/v2/experimental_api.rs index c9c0bd8dfb..5bd44f9f5c 100644 --- a/codex-rs/app-server/tests/suite/v2/experimental_api.rs +++ b/codex-rs/app-server/tests/suite/v2/experimental_api.rs @@ -93,6 +93,7 @@ async fn realtime_conversation_start_requires_experimental_api_capability() -> R flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: "thr_123".to_string(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -221,6 +222,7 @@ async fn realtime_webrtc_start_requires_experimental_api_capability() -> Result< flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: "thr_123".to_string(), model: None, output_modality: RealtimeOutputModality::Audio, 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 3ac137541e..9c6f787a76 100644 --- a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs +++ b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs @@ -44,6 +44,7 @@ use codex_app_server_protocol::TurnStartedNotification; use codex_app_server_protocol::UserInput as V2UserInput; use codex_features::FEATURES; use codex_features::Feature; +use codex_protocol::protocol::CodexResponseHandoffMode; use codex_protocol::protocol::ConversationTextRole; use codex_protocol::protocol::RealtimeConversationVersion; use codex_protocol::protocol::RealtimeOutputModality; @@ -332,6 +333,7 @@ impl RealtimeE2eHarness { offer_sdp, /*client_managed_handoffs*/ None, /*codex_responses_as_items*/ None, + /*codex_response_handoff_mode*/ None, RealtimeConversationVersion::V1, ) .await @@ -345,6 +347,7 @@ impl RealtimeE2eHarness { offer_sdp, /*client_managed_handoffs*/ None, /*codex_responses_as_items*/ Some(true), + /*codex_response_handoff_mode*/ None, RealtimeConversationVersion::V1, ) .await @@ -355,6 +358,7 @@ impl RealtimeE2eHarness { offer_sdp: &str, client_managed_handoffs: Option, codex_responses_as_items: Option, + codex_response_handoff_mode: Option, version: RealtimeConversationVersion, ) -> Result { // Starts realtime through the public JSON-RPC method, then waits for the same client-visible @@ -368,6 +372,7 @@ impl RealtimeE2eHarness { codex_response_item_prefix: codex_responses_as_items .unwrap_or(false) .then(|| RESPONSE_ITEM_PREFIX.to_string()), + codex_response_handoff_mode, codex_responses_as_items, model: None, output_modality: RealtimeOutputModality::Audio, @@ -428,6 +433,7 @@ impl RealtimeE2eHarness { codex_response_item_prefix: codex_responses_as_items .unwrap_or(false) .then(|| RESPONSE_ITEM_PREFIX.to_string()), + codex_response_handoff_mode: None, codex_responses_as_items, model: None, output_modality: RealtimeOutputModality::Audio, @@ -451,7 +457,10 @@ impl RealtimeE2eHarness { .await } - async fn start_frameless_bidi_realtime(&mut self) -> Result { + async fn start_frameless_bidi_realtime( + &mut self, + codex_response_handoff_mode: Option, + ) -> Result { let start_request_id = self .mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { @@ -459,6 +468,7 @@ impl RealtimeE2eHarness { client_managed_handoffs: None, flush_transcript_tail_on_session_end: None, codex_response_item_prefix: None, + codex_response_handoff_mode, codex_responses_as_items: None, model: None, output_modality: RealtimeOutputModality::Audio, @@ -730,6 +740,7 @@ async fn realtime_conversation_streams_v2_notifications() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id.clone(), model: Some("realtime-treatment-model".to_string()), output_modality: RealtimeOutputModality::Audio, @@ -1024,6 +1035,7 @@ async fn realtime_start_can_skip_startup_context() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id.clone(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -1125,6 +1137,7 @@ async fn realtime_text_output_modality_requests_text_output_and_final_transcript flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id.clone(), model: None, output_modality: RealtimeOutputModality::Text, @@ -1312,6 +1325,7 @@ async fn realtime_conversation_stop_emits_closed_notification() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id.clone(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -1419,6 +1433,7 @@ async fn realtime_webrtc_start_emits_sdp_notification() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_id.clone(), model: None, output_modality: RealtimeOutputModality::Audio, @@ -1612,6 +1627,7 @@ async fn webrtc_v3_start_posts_live_session_and_joins_without_session_update() - "v=offer\r\n", /*client_managed_handoffs*/ None, /*codex_responses_as_items*/ None, + /*codex_response_handoff_mode*/ None, RealtimeConversationVersion::V3, ) .await?; @@ -1742,6 +1758,7 @@ async fn webrtc_v1_client_managed_handoffs_disable_automatic_output() -> Result< "v=offer\r\n", /*client_managed_handoffs*/ Some(true), /*codex_responses_as_items*/ None, + /*codex_response_handoff_mode*/ None, RealtimeConversationVersion::V1, ) .await?; @@ -1799,6 +1816,74 @@ async fn webrtc_v1_client_managed_handoffs_disable_automatic_output() -> Result< Ok(()) } +#[tokio::test] +async fn webrtc_v1_ignores_codex_response_handoff_mode() -> Result<()> { + skip_if_no_network!(Ok(())); + + let mut commentary = responses::ev_assistant_message("msg-commentary", "background progress"); + commentary["item"]["phase"] = json!("commentary"); + let mut final_answer = responses::ev_assistant_message("msg-final", "background complete"); + final_answer["item"]["phase"] = json!("final_answer"); + let mut harness = RealtimeE2eHarness::new( + RealtimeTestVersion::V1, + main_loop_responses(vec![responses::sse(vec![ + responses::ev_response_created("resp-1"), + commentary, + final_answer, + responses::ev_completed("resp-1"), + ])]), + realtime_sideband(vec![realtime_sideband_connection(vec![ + vec![ + session_updated("sess_v1_channel_handoff"), + json!({ + "type": "conversation.handoff.requested", + "handoff_id": "handoff_channel", + "item_id": "item_channel", + "input_transcript": "run the background task" + }), + ], + vec![], + 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_mode*/ Some(CodexResponseHandoffMode::BemTags), + RealtimeConversationVersion::V1, + ) + .await?; + assert_eq!(started.started.version, RealtimeConversationVersion::V1); + let _ = harness + .read_notification::("turn/completed") + .await?; + + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 1).await, + json!({ + "type": "conversation.handoff.append", + "handoff_id": "handoff_channel", + "output_text": "background progress", + }) + ); + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 2).await, + json!({ + "type": "conversation.handoff.append", + "handoff_id": "handoff_channel", + "output_text": "\"Agent Final Message\":\n\nbackground complete", + }) + ); + + harness.shutdown().await; + Ok(()) +} + #[tokio::test] async fn webrtc_v1_handoff_request_delegates_context_and_manual_append_speaks() -> Result<()> { skip_if_no_network!(Ok(())); @@ -2138,114 +2223,123 @@ async fn websocket_v2_assistant_output_without_handoff_reaches_realtime_context( } #[tokio::test] -async fn websocket_v3_frameless_delegation_runs_codex_and_appends_context() -> Result<()> { +async fn websocket_v3_routes_handoffs_by_session_mode() -> 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" - }] - } - }), + for (mode, expected_channels) in [ + (None, [None, None, None, None]), + ( + Some(CodexResponseHandoffMode::Commentary), + [ + Some("commentary"), + Some("commentary"), + Some("commentary"), + Some("commentary"), ], - vec![], - vec![], - vec![], - ])]), - ) - .await?; + ), + ( + Some(CodexResponseHandoffMode::BemTags), + [ + Some("commentary"), + Some("commentary"), + Some("speakable"), + Some("speakable"), + ], + ), + ] { + let analysis_text = "<|start|>assistant<|channel|>analysis<|message|>silent context<|end|>"; + let commentary_text = + "<|start|>assistant<|channel|>commentary<|message|>still working<|end|>"; + let final_text = "<|start|>assistant<|channel|>final<|message|>finished<|end|>"; + let fallback_text = "unparsable BEM output"; + let analysis = responses::ev_assistant_message("msg-analysis", analysis_text); + let commentary = responses::ev_assistant_message("msg-commentary", commentary_text); + let final_answer = responses::ev_assistant_message("msg-final", final_text); + let fallback = responses::ev_assistant_message("msg-fallback", fallback_text); + let mut harness = RealtimeE2eHarness::new( + RealtimeTestVersion::V1, + main_loop_responses(vec![responses::sse(vec![ + responses::ev_response_created("resp-1"), + analysis, + commentary, + final_answer, + fallback, + responses::ev_completed("resp-1"), + ])]), + 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![], + vec![], + vec![], + ])]), + ) + .await?; - let started = harness.start_frameless_bidi_realtime().await?; - assert_eq!(started.version, RealtimeConversationVersion::V3); + let started = harness.start_frameless_bidi_realtime(mode).await?; + assert_eq!(started.version, RealtimeConversationVersion::V3); + let _ = harness + .read_notification::("turn/completed") + .await?; + + for (request_index, (text, channel)) in + [analysis_text, commentary_text, final_text, fallback_text] + .into_iter() + .zip(expected_channels) + .enumerate() + { + let mut expected = json!({ + "type": "delegation.context.append", + "delegation_item_id": "delegation_frameless", + "content": [{ + "type": "input_text", + "text": text + }] + }); + if let Some(channel) = channel { + expected["channel"] = json!(channel); + } + assert_eq!( + harness + .sideband_outbound_request(/*request_index*/ request_index + 1) + .await, + expected + ); + } - 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") - ); + .append_speech(harness.thread_id.clone(), "manual spoken update") + .await?; + assert_eq!( + harness.sideband_outbound_request(/*request_index*/ 5).await, + json!({ + "type": "session.context.append", + "content": [{ + "type": "input_text", + "text": "manual spoken update" + }], + "channel": "speakable" + }) + ); - 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; + harness.shutdown().await; + } Ok(()) } @@ -2913,6 +3007,7 @@ async fn realtime_webrtc_start_surfaces_backend_error() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id, model: None, output_modality: RealtimeOutputModality::Audio, @@ -2982,6 +3077,7 @@ async fn realtime_conversation_requires_feature_flag() -> Result<()> { flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, + codex_response_handoff_mode: None, thread_id: thread_start.thread.id.clone(), model: None, output_modality: RealtimeOutputModality::Audio, diff --git a/codex-rs/codex-api/src/endpoint/mod.rs b/codex-rs/codex-api/src/endpoint/mod.rs index 106c5d73ff..5d01a15fe3 100644 --- a/codex-rs/codex-api/src/endpoint/mod.rs +++ b/codex-rs/codex-api/src/endpoint/mod.rs @@ -15,6 +15,7 @@ pub use memories::MemoriesClient; pub use models::ModelsClient; pub use realtime_call::RealtimeCallClient; pub use realtime_call::RealtimeCallResponse; +pub use realtime_websocket::RealtimeContextAppendChannel; pub use realtime_websocket::RealtimeEventParser; pub use realtime_websocket::RealtimeOutputModality; pub use realtime_websocket::RealtimeSessionConfig; 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 b720f44cee..9d075a89c8 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs @@ -9,6 +9,7 @@ use crate::endpoint::realtime_websocket::methods_frameless_bidi::context_append_ 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::RealtimeContextAppendChannel; use crate::endpoint::realtime_websocket::protocol::RealtimeEvent; use crate::endpoint::realtime_websocket::protocol::RealtimeEventParser; use crate::endpoint::realtime_websocket::protocol::RealtimeOutboundMessage; @@ -211,6 +212,7 @@ pub struct RealtimeWebsocketWriter { stream: Arc, is_closed: Arc, event_parser: RealtimeEventParser, + context_append_channel: Option, } #[derive(Clone)] @@ -281,6 +283,7 @@ impl RealtimeWebsocketConnection { stream: Arc::clone(&stream), is_closed: Arc::clone(&is_closed), event_parser, + context_append_channel: None, }, events: RealtimeWebsocketEvents { rx_message, @@ -294,6 +297,11 @@ impl RealtimeWebsocketConnection { } impl RealtimeWebsocketWriter { + pub fn with_context_append_channel(mut self, channel: RealtimeContextAppendChannel) -> Self { + self.context_append_channel = Some(channel); + self + } + pub async fn send_audio_frame(&self, frame: RealtimeAudioFrame) -> Result<(), ApiError> { let message = match self.event_parser { RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => { @@ -315,6 +323,7 @@ impl RealtimeWebsocketWriter { self.event_parser, text, role, + self.context_append_channel, )) .await } @@ -328,6 +337,7 @@ impl RealtimeWebsocketWriter { self.event_parser, handoff_id, output_text, + self.context_append_channel, )) .await } @@ -341,6 +351,7 @@ impl RealtimeWebsocketWriter { self.event_parser, handoff_id, output_text, + self.context_append_channel, )) .await } @@ -354,6 +365,7 @@ impl RealtimeWebsocketWriter { self.event_parser, call_id, output_text, + self.context_append_channel, )) .await } @@ -413,6 +425,7 @@ impl RealtimeWebsocketWriter { match message { RealtimeOutboundMessage::DelegationContextAppend { delegation_item_id, + channel, content, } => { if let Some(content) = content.first() { @@ -420,17 +433,20 @@ impl RealtimeWebsocketWriter { self.send_json_frame(&frameless_delegation_context_append_message( delegation_item_id.clone(), chunk, + *channel, )) .await?; } return Ok(()); } } - RealtimeOutboundMessage::SessionContextAppend { content } => { + RealtimeOutboundMessage::SessionContextAppend { channel, 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?; + self.send_json_frame(&frameless_session_context_append_message( + chunk, *channel, + )) + .await?; } return Ok(()); } 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 76c9147ed0..e6376be768 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 @@ -10,6 +10,7 @@ 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::RealtimeContextAppendChannel; use crate::endpoint::realtime_websocket::protocol::RealtimeOutboundMessage; use crate::endpoint::realtime_websocket::protocol::RealtimeOutputModality; use crate::endpoint::realtime_websocket::protocol::RealtimeSessionConfig; @@ -40,10 +41,13 @@ pub(super) fn conversation_item_create_message( wire_adapter: RealtimeWireAdapter, text: String, role: ConversationTextRole, + context_append_channel: Option, ) -> RealtimeOutboundMessage { match wire_adapter { RealtimeWireAdapter::V1 => v1_conversation_item_create_message(text, role), - RealtimeWireAdapter::FramelessBidi => frameless_session_context_append_message(text), + RealtimeWireAdapter::FramelessBidi => { + frameless_session_context_append_message(text, context_append_channel) + } RealtimeWireAdapter::RealtimeV2 => v2_conversation_item_create_message(text, role), } } @@ -52,12 +56,15 @@ pub(super) fn conversation_handoff_append_message( wire_adapter: RealtimeWireAdapter, handoff_id: String, output_text: String, + context_append_channel: Option, ) -> 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::FramelessBidi => frameless_delegation_context_append_message( + handoff_id, + output_text, + context_append_channel, + ), RealtimeWireAdapter::RealtimeV2 => { unreachable!("realtime v2 does not send conversation handoff output") } @@ -68,10 +75,13 @@ pub(super) fn standalone_handoff_message( wire_adapter: RealtimeWireAdapter, handoff_id: String, output_text: String, + context_append_channel: Option, ) -> 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::FramelessBidi => { + frameless_session_context_append_message(output_text, context_append_channel) + } RealtimeWireAdapter::RealtimeV2 => { unreachable!("realtime v2 does not send standalone handoff output") } @@ -82,6 +92,7 @@ pub(super) fn conversation_function_call_output_message( wire_adapter: RealtimeWireAdapter, call_id: String, output_text: String, + context_append_channel: Option, ) -> RealtimeOutboundMessage { match wire_adapter { RealtimeWireAdapter::V1 => v1_conversation_handoff_append_message( @@ -90,7 +101,8 @@ pub(super) fn conversation_function_call_output_message( ), RealtimeWireAdapter::FramelessBidi => frameless_delegation_context_append_message( call_id, - format!("{AGENT_FINAL_MESSAGE_PREFIX}{output_text}"), + output_text, + context_append_channel, ), RealtimeWireAdapter::RealtimeV2 => { v2_conversation_function_call_output_message(call_id, output_text) 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 index 8b187960c7..63a2672216 100644 --- 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 @@ -1,6 +1,7 @@ use super::conversation_function_call_output_message; use super::conversation_handoff_append_message; use super::standalone_handoff_message; +use crate::endpoint::realtime_websocket::protocol::RealtimeContextAppendChannel; use crate::endpoint::realtime_websocket::protocol::RealtimeWireAdapter; use pretty_assertions::assert_eq; use serde_json::Value; @@ -8,16 +9,18 @@ use serde_json::json; use serde_json::to_value; #[test] -fn identical_handoff_output_encodes_for_each_bidi_wire_protocol() { +fn context_append_channel_only_encodes_for_frameless_handoff_output() { let legacy = conversation_handoff_append_message( RealtimeWireAdapter::V1, "handoff-123".to_string(), "The result".to_string(), + Some(RealtimeContextAppendChannel::Commentary), ); let frameless = conversation_handoff_append_message( RealtimeWireAdapter::FramelessBidi, "handoff-123".to_string(), "The result".to_string(), + Some(RealtimeContextAppendChannel::Commentary), ); assert_eq!( @@ -33,6 +36,7 @@ fn identical_handoff_output_encodes_for_each_bidi_wire_protocol() { json!({ "type": "delegation.context.append", "delegation_item_id": "handoff-123", + "channel": "commentary", "content": [{"type": "input_text", "text": "The result"}], }) ); @@ -44,11 +48,13 @@ fn standalone_handoff_uses_session_context_for_frameless() { RealtimeWireAdapter::V1, "codex".to_string(), "Speak this".to_string(), + Some(RealtimeContextAppendChannel::Speakable), ); let frameless = standalone_handoff_message( RealtimeWireAdapter::FramelessBidi, "codex".to_string(), "Speak this".to_string(), + Some(RealtimeContextAppendChannel::Speakable), ); assert_eq!( @@ -63,19 +69,20 @@ fn standalone_handoff_uses_session_context_for_frameless() { to_value(frameless).expect("frameless standalone handoff should serialize"), json!({ "type": "session.context.append", + "channel": "speakable", "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()); +fn completed_handoff_only_prefixes_v1_payload_text() { 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(), + Some(RealtimeContextAppendChannel::Speakable), )) .expect("handoff output should serialize"); let text = match wire_adapter { @@ -83,6 +90,23 @@ fn completed_handoff_preserves_legacy_payload_text_in_frameless() { RealtimeWireAdapter::FramelessBidi => &encoded["content"][0]["text"], RealtimeWireAdapter::RealtimeV2 => unreachable!(), }; - assert_eq!(text, &expected_text); + assert_eq!( + text, + &Value::String(match wire_adapter { + RealtimeWireAdapter::V1 => "\"Agent Final Message\":\n\nDone".to_string(), + RealtimeWireAdapter::FramelessBidi => "Done".to_string(), + RealtimeWireAdapter::RealtimeV2 => unreachable!(), + }) + ); + assert_eq!( + encoded.get("channel").cloned(), + match wire_adapter { + RealtimeWireAdapter::V1 => None, + RealtimeWireAdapter::FramelessBidi => { + Some(Value::String("speakable".to_string())) + } + RealtimeWireAdapter::RealtimeV2 => unreachable!(), + } + ); } } 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 index 2b5e01e99f..2f539255d6 100644 --- 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 @@ -1,5 +1,6 @@ use crate::endpoint::realtime_websocket::protocol::FramelessContentType; use crate::endpoint::realtime_websocket::protocol::FramelessInputTextContent; +use crate::endpoint::realtime_websocket::protocol::RealtimeContextAppendChannel; use crate::endpoint::realtime_websocket::protocol::RealtimeOutboundMessage; use crate::endpoint::realtime_websocket::protocol::RealtimeVoice; use serde_json::Value; @@ -10,15 +11,21 @@ const CONTEXT_APPEND_MAX_BYTES: usize = 500; pub(super) fn delegation_context_append_message( delegation_item_id: String, text: String, + channel: Option, ) -> RealtimeOutboundMessage { RealtimeOutboundMessage::DelegationContextAppend { delegation_item_id, + channel, content: input_text_content(text), } } -pub(super) fn session_context_append_message(text: String) -> RealtimeOutboundMessage { +pub(super) fn session_context_append_message( + text: String, + channel: Option, +) -> RealtimeOutboundMessage { RealtimeOutboundMessage::SessionContextAppend { + channel, content: input_text_content(text), } } 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 ac56dedac0..6bb4808fae 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/mod.rs @@ -14,6 +14,7 @@ pub use methods::RealtimeWebsocketConnection; pub use methods::RealtimeWebsocketEvents; pub use methods::RealtimeWebsocketWriter; pub use methods_common::session_update_session_json; +pub use protocol::RealtimeContextAppendChannel; pub use protocol::RealtimeEventParser; pub use protocol::RealtimeOutputModality; pub use protocol::RealtimeSessionConfig; 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 85145aa1a4..a47d980a7b 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs @@ -25,6 +25,14 @@ pub enum RealtimeSessionMode { Transcription, } +/// Selects the semantic stream used for Frameless Bidi context appends. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RealtimeContextAppendChannel { + Speakable, + Commentary, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct RealtimeSessionConfig { pub instructions: String, @@ -51,10 +59,14 @@ pub(super) enum RealtimeOutboundMessage { #[serde(rename = "delegation.context.append")] DelegationContextAppend { delegation_item_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + channel: Option, content: Vec, }, #[serde(rename = "session.context.append")] SessionContextAppend { + #[serde(skip_serializing_if = "Option::is_none")] + channel: Option, content: Vec, }, #[serde(rename = "session.close")] diff --git a/codex-rs/codex-api/src/lib.rs b/codex-rs/codex-api/src/lib.rs index 1d1a4c8c6e..2a549c1502 100644 --- a/codex-rs/codex-api/src/lib.rs +++ b/codex-rs/codex-api/src/lib.rs @@ -52,6 +52,7 @@ pub use crate::endpoint::MemoriesClient; pub use crate::endpoint::ModelsClient; pub use crate::endpoint::RealtimeCallClient; pub use crate::endpoint::RealtimeCallResponse; +pub use crate::endpoint::RealtimeContextAppendChannel; pub use crate::endpoint::RealtimeEventParser; pub use crate::endpoint::RealtimeOutputModality; pub use crate::endpoint::RealtimeSessionConfig; diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index a9a194c83c..19a3fec14f 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -16,6 +16,7 @@ use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use codex_api::ApiError; use codex_api::Provider as ApiProvider; use codex_api::RealtimeAudioFrame; +use codex_api::RealtimeContextAppendChannel; use codex_api::RealtimeEvent; use codex_api::RealtimeEventParser; use codex_api::RealtimeSessionConfig; @@ -36,6 +37,7 @@ use codex_protocol::error::CodexErr; use codex_protocol::error::Result as CodexResult; use codex_protocol::models::MessagePhase; use codex_protocol::protocol::CodexErrorInfo; +use codex_protocol::protocol::CodexResponseHandoffMode; use codex_protocol::protocol::ConversationAudioParams; use codex_protocol::protocol::ConversationSpeechParams; use codex_protocol::protocol::ConversationStartParams; @@ -74,6 +76,11 @@ use tracing::error; use tracing::info; use tracing::warn; +mod bem; + +use self::bem::ChannelParser as BemChannelParser; +use self::bem::message_phase as bem_message_phase; + const AUDIO_IN_QUEUE_CAPACITY: usize = 256; const TEXT_IN_QUEUE_CAPACITY: usize = 64; const HANDOFF_OUT_QUEUE_CAPACITY: usize = 64; @@ -121,15 +128,22 @@ enum RealtimeSessionKind { #[derive(Clone, Debug)] struct RealtimeHandoffState { output_tx: Sender, - last_output_text: Arc>>, + last_output: Arc>>, stream: Arc>, client_managed_handoffs: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, + codex_response_handoff_mode: CodexResponseHandoffMode, session_kind: RealtimeSessionKind, event_parser: RealtimeEventParser, } +#[derive(Clone, Debug)] +struct RealtimeHandoffOutput { + text: String, + phase: Option, +} + #[derive(Debug, Default)] struct RealtimeHandoffStreamState { active_handoff: Option, @@ -140,6 +154,8 @@ struct RealtimeHandoffStreamState { struct RealtimeStreamedItem { handoff_id: String, phase: Option, + bem_channel_parser: Option, + prefix_final_message: bool, sent_bytes: usize, buffered_text: String, tail_text: String, @@ -154,7 +170,10 @@ impl RealtimeStreamedItem { } fn output_prefix(&self) -> &'static str { - if self.sent_bytes == 0 && !matches!(self.phase, Some(MessagePhase::Commentary)) { + if self.prefix_final_message + && self.sent_bytes == 0 + && !matches!(self.phase, Some(MessagePhase::Commentary)) + { AGENT_FINAL_MESSAGE_PREFIX } else { "" @@ -182,6 +201,36 @@ impl RealtimeStreamedItem { if text.is_empty() { return; } + + let text = if let Some(parser) = self.bem_channel_parser.as_mut() { + let Some(text) = parser.push(text) else { + return; + }; + self.phase = parser.phase(); + text + } else { + text.to_string() + }; + self.push_output_text(&text); + } + + fn finish_input(&mut self) { + let Some(parser) = self.bem_channel_parser.as_mut() else { + return; + }; + self.phase = parser.phase(); + let text = parser.finish(); + if self.phase.is_none() && !text.is_empty() { + warn!("BEM output ended before a recognized channel header was received"); + self.phase = Some(MessagePhase::FinalAnswer); + } + self.push_output_text(&text); + } + + fn push_output_text(&mut self, text: &str) { + if text.is_empty() { + return; + } if self.truncated { self.tail_text.push_str(text); self.tail_text = @@ -256,12 +305,35 @@ fn take_last_bytes_at_char_boundary(text: &str, max_bytes: usize) -> &str { #[derive(Debug, PartialEq, Eq)] enum RealtimeOutbound { - StandaloneHandoff { text: String }, - HandoffUpdate { handoff_id: String, text: String }, - HandoffAppend { handoff_id: String, text: String }, - CompletedHandoff { handoff_id: String, text: String }, - ConversationItem { text: String }, - HandoffCompleteAck { handoff_id: String }, + StandaloneHandoff { + text: String, + phase: Option, + }, + StandaloneSpeech { + text: String, + }, + HandoffUpdate { + handoff_id: String, + text: String, + phase: Option, + }, + HandoffAppend { + handoff_id: String, + text: String, + phase: Option, + }, + CompletedHandoff { + handoff_id: String, + text: String, + phase: Option, + }, + ConversationItem { + text: String, + phase: Option, + }, + HandoffCompleteAck { + handoff_id: String, + }, } #[derive(Debug, PartialEq, Eq)] @@ -359,16 +431,18 @@ impl RealtimeHandoffState { client_managed_handoffs: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, + codex_response_handoff_mode: CodexResponseHandoffMode, session_kind: RealtimeSessionKind, event_parser: RealtimeEventParser, ) -> Self { Self { output_tx, - last_output_text: Arc::new(Mutex::new(None)), + last_output: Arc::new(Mutex::new(None)), stream: Arc::new(Mutex::new(RealtimeHandoffStreamState::default())), client_managed_handoffs, codex_responses_as_items, codex_response_item_prefix, + codex_response_handoff_mode, session_kind, event_parser, } @@ -379,6 +453,11 @@ impl RealtimeHandoffState { && !self.client_managed_handoffs && !self.codex_responses_as_items } + + fn routes_handoff_by_bem(&self) -> bool { + self.event_parser == RealtimeEventParser::FramelessBidi + && self.codex_response_handoff_mode == CodexResponseHandoffMode::BemTags + } } #[allow(dead_code)] @@ -400,6 +479,7 @@ struct RealtimeStart { flush_transcript_tail_on_session_end: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, + codex_response_handoff_mode: CodexResponseHandoffMode, realtime_call_api_provider: Option, session_config: RealtimeSessionConfig, model_client: ModelClient, @@ -458,6 +538,7 @@ impl RealtimeConversationManager { flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, + codex_response_handoff_mode, realtime_call_api_provider, session_config, model_client, @@ -486,6 +567,7 @@ impl RealtimeConversationManager { client_managed_handoffs, codex_responses_as_items, codex_response_item_prefix, + codex_response_handoff_mode, session_kind, event_parser, ); @@ -668,28 +750,45 @@ impl RealtimeConversationManager { if handoff.client_managed_handoffs { return Ok(()); } + let phase = if handoff.routes_handoff_by_bem() { + match bem_message_phase(&output_text) { + Some(phase) => Some(phase), + None => { + warn!("BEM output did not contain a recognized channel header"); + Some(MessagePhase::FinalAnswer) + } + } + } else { + phase + }; let is_commentary = matches!(phase, Some(MessagePhase::Commentary)); let active_handoff = handoff.stream.lock().await.active_handoff.clone(); let output = match active_handoff { Some(handoff_id) => { let output_text = realtime_backend_output(output_text, handoff.session_kind); - *handoff.last_output_text.lock().await = Some(output_text.clone()); + *handoff.last_output.lock().await = Some(RealtimeHandoffOutput { + text: output_text.clone(), + phase: phase.clone(), + }); if handoff.codex_responses_as_items { RealtimeOutbound::ConversationItem { text: realtime_backend_item( output_text, handoff.codex_response_item_prefix.as_deref(), ), + phase, } - } else if handoff.session_kind == RealtimeSessionKind::V1 && is_commentary { + } else if handoff.event_parser == RealtimeEventParser::V1 && is_commentary { RealtimeOutbound::HandoffAppend { handoff_id, text: output_text, + phase, } } else { RealtimeOutbound::HandoffUpdate { handoff_id, text: output_text, + phase, } } } @@ -702,14 +801,16 @@ impl RealtimeConversationManager { output_text, handoff.codex_response_item_prefix.as_deref(), ), + phase, } } else { RealtimeOutbound::StandaloneHandoff { - text: if handoff.session_kind == RealtimeSessionKind::V1 && !is_commentary { + text: if handoff.event_parser == RealtimeEventParser::V1 && !is_commentary { format!("{AGENT_FINAL_MESSAGE_PREFIX}{output_text}") } else { output_text }, + phase, } } } @@ -745,7 +846,15 @@ impl RealtimeConversationManager { }; let mut streamed_item = RealtimeStreamedItem { handoff_id, - phase, + phase: if handoff.routes_handoff_by_bem() { + None + } else { + phase + }, + bem_channel_parser: handoff + .routes_handoff_by_bem() + .then(BemChannelParser::default), + prefix_final_message: handoff.event_parser == RealtimeEventParser::V1, sent_bytes: 0, buffered_text: String::new(), tail_text: String::new(), @@ -821,6 +930,7 @@ impl RealtimeConversationManager { let Some(mut streamed_item) = handoff.stream.lock().await.items.remove(item_id) else { return false; }; + streamed_item.finish_input(); let chunk = streamed_item.drain_final_chunk(); let sent_output = streamed_item.sent_bytes > 0; if let Some(text) = chunk { @@ -829,6 +939,7 @@ impl RealtimeConversationManager { .send(RealtimeOutbound::HandoffAppend { handoff_id: streamed_item.handoff_id, text, + phase: streamed_item.phase, }) .await; } @@ -852,7 +963,7 @@ impl RealtimeConversationManager { handoff .output_tx - .send(RealtimeOutbound::StandaloneHandoff { + .send(RealtimeOutbound::StandaloneSpeech { text: realtime_backend_output(text, handoff.session_kind), }) .await @@ -879,7 +990,7 @@ impl RealtimeConversationManager { let Some(handoff_id) = handoff.stream.lock().await.active_handoff.clone() else { return Ok(()); }; - let Some(output_text) = handoff.last_output_text.lock().await.clone() else { + let Some(last_output) = handoff.last_output.lock().await.clone() else { return Ok(()); }; @@ -888,7 +999,8 @@ impl RealtimeConversationManager { } else { RealtimeOutbound::CompletedHandoff { handoff_id, - text: output_text, + text: last_output.text, + phase: last_output.phase, } }; @@ -910,7 +1022,7 @@ impl RealtimeConversationManager { stream.active_handoff = None; stream.items.clear(); } - *handoff.last_output_text.lock().await = None; + *handoff.last_output.lock().await = None; } } @@ -987,6 +1099,7 @@ struct PreparedRealtimeConversationStart { flush_transcript_tail_on_session_end: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, + codex_response_handoff_mode: CodexResponseHandoffMode, realtime_call_api_provider: Option, requested_realtime_session_id: Option, version: RealtimeWsVersion, @@ -1072,6 +1185,7 @@ async fn prepare_realtime_start( flush_transcript_tail_on_session_end: params.flush_transcript_tail_on_session_end, codex_responses_as_items: params.codex_responses_as_items, codex_response_item_prefix: params.codex_response_item_prefix, + codex_response_handoff_mode: params.codex_response_handoff_mode, realtime_call_api_provider, requested_realtime_session_id, version, @@ -1243,6 +1357,7 @@ async fn handle_start_inner( flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, + codex_response_handoff_mode, realtime_call_api_provider, requested_realtime_session_id, version, @@ -1261,6 +1376,7 @@ async fn handle_start_inner( flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, + codex_response_handoff_mode, realtime_call_api_provider, session_config, model_client: sess.services.model_client.clone(), @@ -1700,7 +1816,7 @@ async fn handle_text_input( } async fn flush_streamed_handoff_item(handoff: &RealtimeHandoffState, item_id: &str) { - let (handoff_id, text) = { + let (handoff_id, text, phase) = { let mut stream = handoff.stream.lock().await; let Some(streamed_item) = stream.items.get_mut(item_id) else { return; @@ -1710,11 +1826,19 @@ async fn flush_streamed_handoff_item(handoff: &RealtimeHandoffState, item_id: &s return; }; streamed_item.last_flush_at = Instant::now(); - (streamed_item.handoff_id.clone(), text) + ( + streamed_item.handoff_id.clone(), + text, + streamed_item.phase.clone(), + ) }; let _ = handoff .output_tx - .send(RealtimeOutbound::HandoffAppend { handoff_id, text }) + .send(RealtimeOutbound::HandoffAppend { + handoff_id, + text, + phase, + }) .await; } @@ -1730,6 +1854,26 @@ fn schedule_streamed_handoff_flush( }); } +fn v3_output_writer( + writer: &RealtimeWebsocketWriter, + phase: Option<&MessagePhase>, + handoff_mode: CodexResponseHandoffMode, +) -> RealtimeWebsocketWriter { + let channel = match handoff_mode { + CodexResponseHandoffMode::Thinking => None, + CodexResponseHandoffMode::Commentary => Some(RealtimeContextAppendChannel::Commentary), + CodexResponseHandoffMode::BemTags => match phase { + Some(MessagePhase::FinalAnswer) => Some(RealtimeContextAppendChannel::Speakable), + Some(MessagePhase::Commentary) => Some(RealtimeContextAppendChannel::Commentary), + None => Some(RealtimeContextAppendChannel::Speakable), + }, + }; + match channel { + Some(channel) => writer.clone().with_context_append_channel(channel), + None => writer.clone(), + } +} + async fn handle_handoff_output( handoff_output: Result, writer: &RealtimeWebsocketWriter, @@ -1739,34 +1883,117 @@ async fn handle_handoff_output( response_create_queue: &mut RealtimeResponseCreateQueue, ) -> anyhow::Result<()> { let handoff_output = handoff_output.context("handoff output channel closed")?; - let result = match event_parser { - RealtimeEventParser::V1 | RealtimeEventParser::FramelessBidi => match handoff_output { - RealtimeOutbound::StandaloneHandoff { text } => { + RealtimeEventParser::V1 => match handoff_output { + RealtimeOutbound::StandaloneHandoff { text, phase: _ } => { writer .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text) .await } - RealtimeOutbound::HandoffUpdate { handoff_id, text } - | RealtimeOutbound::CompletedHandoff { handoff_id, text } => { + RealtimeOutbound::StandaloneSpeech { text } => { + writer + .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text) + .await + } + RealtimeOutbound::HandoffUpdate { + handoff_id, + text, + phase: _, + } + | RealtimeOutbound::CompletedHandoff { + handoff_id, + text, + phase: _, + } => { writer .send_conversation_function_call_output(handoff_id, text) .await } - RealtimeOutbound::HandoffAppend { handoff_id, text } => { + RealtimeOutbound::HandoffAppend { + handoff_id, + text, + phase: _, + } => { writer .send_conversation_handoff_append(handoff_id, text) .await } - RealtimeOutbound::ConversationItem { text } => { + RealtimeOutbound::ConversationItem { text, phase: _ } => { writer .send_conversation_item_create(text, ConversationTextRole::Developer) .await } RealtimeOutbound::HandoffCompleteAck { .. } => Ok(()), }, + RealtimeEventParser::FramelessBidi => match handoff_output { + RealtimeOutbound::StandaloneHandoff { text, phase } => { + v3_output_writer( + writer, + phase.as_ref(), + handoff_state.codex_response_handoff_mode, + ) + .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text) + .await + } + RealtimeOutbound::StandaloneSpeech { text } => { + writer + .clone() + .with_context_append_channel(RealtimeContextAppendChannel::Speakable) + .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text) + .await + } + RealtimeOutbound::HandoffUpdate { + handoff_id, + text, + phase, + } => { + v3_output_writer( + writer, + phase.as_ref(), + handoff_state.codex_response_handoff_mode, + ) + .send_conversation_function_call_output(handoff_id, text) + .await + } + RealtimeOutbound::HandoffAppend { + handoff_id, + text, + phase, + } => { + v3_output_writer( + writer, + phase.as_ref(), + handoff_state.codex_response_handoff_mode, + ) + .send_conversation_handoff_append(handoff_id, text) + .await + } + RealtimeOutbound::CompletedHandoff { + handoff_id, + text, + phase, + } => { + v3_output_writer( + writer, + phase.as_ref(), + handoff_state.codex_response_handoff_mode, + ) + .send_conversation_function_call_output(handoff_id, text) + .await + } + RealtimeOutbound::ConversationItem { text, phase } => { + v3_output_writer( + writer, + phase.as_ref(), + handoff_state.codex_response_handoff_mode, + ) + .send_conversation_item_create(text, ConversationTextRole::Developer) + .await + } + RealtimeOutbound::HandoffCompleteAck { .. } => Ok(()), + }, RealtimeEventParser::RealtimeV2 => match handoff_output { - RealtimeOutbound::StandaloneHandoff { text } => { + RealtimeOutbound::StandaloneHandoff { text, phase: _ } => { if let Err(err) = writer .send_conversation_item_create(text, ConversationTextRole::User) .await @@ -1778,8 +2005,28 @@ async fn handle_handoff_output( .await; } } - RealtimeOutbound::HandoffUpdate { handoff_id, text } - | RealtimeOutbound::HandoffAppend { handoff_id, text } => { + RealtimeOutbound::StandaloneSpeech { text } => { + if let Err(err) = writer + .send_conversation_item_create(text, ConversationTextRole::User) + .await + { + Err(err) + } else { + return response_create_queue + .request_create(writer, events_tx, "standalone handoff") + .await; + } + } + RealtimeOutbound::HandoffUpdate { + handoff_id, + text, + phase: _, + } + | RealtimeOutbound::HandoffAppend { + handoff_id, + text, + phase: _, + } => { let active_handoff = handoff_state.stream.lock().await.active_handoff.clone(); match active_handoff { Some(active_handoff) if active_handoff == handoff_id => {} @@ -1795,6 +2042,7 @@ async fn handle_handoff_output( RealtimeOutbound::CompletedHandoff { handoff_id, text: _, + phase: _, } => { if let Err(err) = writer .send_conversation_function_call_output( @@ -1810,7 +2058,7 @@ async fn handle_handoff_output( .await; } } - RealtimeOutbound::ConversationItem { text } => { + RealtimeOutbound::ConversationItem { text, phase: _ } => { writer .send_conversation_item_create(text, ConversationTextRole::Developer) .await diff --git a/codex-rs/core/src/realtime_conversation/bem.rs b/codex-rs/core/src/realtime_conversation/bem.rs new file mode 100644 index 0000000000..c6b6df9822 --- /dev/null +++ b/codex-rs/core/src/realtime_conversation/bem.rs @@ -0,0 +1,49 @@ +use codex_protocol::models::MessagePhase; + +const BEM_ASSISTANT_CHANNEL_PREFIX: &str = "<|start|>assistant<|channel|>"; +const BEM_MESSAGE_MARKER: &str = "<|message|>"; + +pub(super) fn message_phase(text: &str) -> Option { + let channel_and_message = text.strip_prefix(BEM_ASSISTANT_CHANNEL_PREFIX)?; + let (channel, _) = channel_and_message.split_once(BEM_MESSAGE_MARKER)?; + match channel { + "analysis" | "commentary" => Some(MessagePhase::Commentary), + "final" => Some(MessagePhase::FinalAnswer), + _ => None, + } +} + +/// Buffers a streamed BEM message until its channel header is complete. +/// +/// Once the channel is known, the original envelope is released unchanged so +/// the frontend model can distinguish BEM `analysis` from `commentary`. +#[derive(Debug, Default)] +pub(super) struct ChannelParser { + buffered_text: String, + phase: Option, +} + +impl ChannelParser { + pub(super) fn push(&mut self, text: &str) -> Option { + if self.phase.is_some() { + return Some(text.to_string()); + } + + self.buffered_text.push_str(text); + self.phase = message_phase(&self.buffered_text); + self.phase.as_ref()?; + Some(std::mem::take(&mut self.buffered_text)) + } + + pub(super) fn phase(&self) -> Option { + self.phase.clone() + } + + pub(super) fn finish(&mut self) -> String { + std::mem::take(&mut self.buffered_text) + } +} + +#[cfg(test)] +#[path = "bem_tests.rs"] +mod tests; diff --git a/codex-rs/core/src/realtime_conversation/bem_tests.rs b/codex-rs/core/src/realtime_conversation/bem_tests.rs new file mode 100644 index 0000000000..580bf4a828 --- /dev/null +++ b/codex-rs/core/src/realtime_conversation/bem_tests.rs @@ -0,0 +1,42 @@ +use super::ChannelParser; +use super::message_phase; +use codex_protocol::models::MessagePhase; +use pretty_assertions::assert_eq; + +#[test] +fn maps_bem_channels_to_realtime_phases() { + for (channel, expected) in [ + ("analysis", MessagePhase::Commentary), + ("commentary", MessagePhase::Commentary), + ("final", MessagePhase::FinalAnswer), + ] { + assert_eq!( + message_phase(&format!( + "<|start|>assistant<|channel|>{channel}<|message|>text<|end|>" + )), + Some(expected) + ); + } +} + +#[test] +fn buffers_streamed_text_until_the_bem_channel_is_complete() { + let mut parser = ChannelParser::default(); + + assert_eq!(parser.push("<|start|>assistant<|channel|>com"), None); + assert_eq!( + parser.push("mentary<|message|>progress"), + Some("<|start|>assistant<|channel|>commentary<|message|>progress".to_string()) + ); + assert_eq!(parser.phase(), Some(MessagePhase::Commentary)); + assert_eq!(parser.push("<|end|>"), Some("<|end|>".to_string())); +} + +#[test] +fn preserves_unrecognized_output_when_the_stream_finishes() { + let mut parser = ChannelParser::default(); + + assert_eq!(parser.push("plain output"), None); + assert_eq!(parser.finish(), "plain output"); + assert_eq!(parser.phase(), None); +} diff --git a/codex-rs/core/src/realtime_conversation_tests.rs b/codex-rs/core/src/realtime_conversation_tests.rs index f9bf00a7c5..4c69238fc1 100644 --- a/codex-rs/core/src/realtime_conversation_tests.rs +++ b/codex-rs/core/src/realtime_conversation_tests.rs @@ -11,6 +11,7 @@ use crate::context::RealtimeDelegationSource; use async_channel::bounded; use codex_api::RealtimeEventParser; use codex_protocol::models::MessagePhase; +use codex_protocol::protocol::CodexResponseHandoffMode; use codex_protocol::protocol::RealtimeHandoffRequested; use codex_protocol::protocol::RealtimeTranscriptEntry; use pretty_assertions::assert_eq; @@ -151,6 +152,7 @@ async fn clears_active_handoff_explicitly() { /*client_managed_handoffs*/ false, /*codex_responses_as_items*/ false, /*codex_response_item_prefix*/ None, + CodexResponseHandoffMode::Thinking, RealtimeSessionKind::V1, /*event_parser*/ RealtimeEventParser::V1, ); @@ -170,6 +172,8 @@ fn streamed_handoff_preserves_a_bounded_final_tail() { let mut item = RealtimeStreamedItem { handoff_id: "handoff_1".to_string(), phase: Some(MessagePhase::FinalAnswer), + bem_channel_parser: None, + prefix_final_message: true, sent_bytes: 0, buffered_text: String::new(), tail_text: String::new(), @@ -193,6 +197,25 @@ fn streamed_handoff_preserves_a_bounded_final_tail() { assert!(output.ends_with("TAIL")); } +#[test] +fn streamed_v3_handoff_omits_the_final_message_prefix() { + let mut item = RealtimeStreamedItem { + handoff_id: "handoff_1".to_string(), + phase: Some(MessagePhase::FinalAnswer), + bem_channel_parser: None, + prefix_final_message: false, + sent_bytes: 0, + buffered_text: String::new(), + tail_text: String::new(), + truncated: false, + last_flush_at: Instant::now(), + flush_scheduled: false, + }; + item.push_text("done"); + + assert_eq!(item.drain_final_chunk(), Some("done".to_string())); +} + #[test] fn uses_quicksilver_alpha_header_for_realtime_v1() { let headers = realtime_request_headers( diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index 885f8c61f4..b021630bf4 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -212,6 +212,8 @@ async fn start_realtime_conversation(codex: &codex_core::CodexThread) -> Result< flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index 3d1a858620..7cdb200624 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -287,6 +287,8 @@ async fn conversation_start_audio_text_close_round_trip() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -433,6 +435,8 @@ async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -528,6 +532,8 @@ async fn conversation_webrtc_start_posts_generated_session() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: Some("session-override-model".to_string()), output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -715,6 +721,8 @@ async fn conversation_webrtc_start_uses_avas_query() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -813,6 +821,8 @@ async fn conversation_webrtc_default_v1_ignores_configured_v2_voice() -> Result< flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -873,6 +883,8 @@ async fn conversation_webrtc_default_v1_rejects_explicit_v2_voice() -> Result<() flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -943,6 +955,8 @@ async fn conversation_webrtc_start_uses_configured_call_base_url_for_avas() -> R flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1037,6 +1051,8 @@ async fn conversation_webrtc_close_while_sideband_connecting_drops_pending_join( flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1128,6 +1144,8 @@ async fn conversation_webrtc_sideband_connect_failure_closes_with_error() -> Res flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1221,6 +1239,8 @@ async fn conversation_start_uses_openai_env_key_fallback_with_chatgpt_auth() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1312,6 +1332,8 @@ async fn assert_transport_close_tail_flush( flush_transcript_tail_on_session_end, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1428,6 +1450,8 @@ async fn conversation_start_preflight_failure_emits_realtime_error_only() -> Res flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1479,6 +1503,8 @@ async fn conversation_start_connect_failure_emits_realtime_error_only() -> Resul flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1578,6 +1604,8 @@ async fn conversation_second_start_replaces_runtime() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1608,6 +1636,8 @@ async fn conversation_second_start_replaces_runtime() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1709,6 +1739,8 @@ async fn conversation_uses_experimental_realtime_ws_base_url_override() -> Resul flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1778,6 +1810,8 @@ async fn conversation_uses_default_realtime_backend_prompt() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1855,6 +1889,8 @@ async fn conversation_uses_empty_instructions_for_null_or_empty_prompt() -> Resu flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1925,6 +1961,8 @@ async fn conversation_uses_explicit_start_voice() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -1987,6 +2025,8 @@ async fn conversation_uses_configured_realtime_voice() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2037,6 +2077,8 @@ async fn conversation_rejects_voice_for_wrong_realtime_version() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2088,6 +2130,8 @@ async fn conversation_uses_experimental_realtime_ws_backend_prompt_override() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2165,6 +2209,8 @@ async fn conversation_uses_experimental_realtime_ws_startup_context_override() - flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2236,6 +2282,8 @@ async fn conversation_disables_realtime_startup_context_with_empty_override() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2300,6 +2348,8 @@ async fn conversation_start_injects_startup_context_from_thread_history() -> Res flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2419,6 +2469,8 @@ async fn conversation_startup_context_current_thread_selects_many_turns_by_budge flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2531,6 +2583,8 @@ async fn conversation_startup_context_falls_back_to_workspace_map() -> Result<() flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2595,6 +2649,8 @@ async fn conversation_startup_context_is_truncated_and_sent_once_per_start() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2680,6 +2736,8 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2781,6 +2839,8 @@ async fn realtime_v2_noop_tool_call_returns_empty_function_output_without_respon flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2884,6 +2944,8 @@ async fn conversation_mirrors_assistant_message_text_to_realtime_handoff() -> Re flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -2958,22 +3020,20 @@ async fn conversation_mirrors_assistant_message_text_to_realtime_handoff() -> Re async fn conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff() -> Result<()> { skip_if_no_network!(Ok(())); - let initial_commentary_text = "seed "; - let first_commentary_delta = "first "; - let second_commentary_delta = "x".repeat(94); + let initial_commentary_text = "<|start|>assistant<|channel|>com"; + let first_commentary_delta = "mentary<|message|>seed first "; + let second_commentary_delta = format!("{}<|end|>", "x".repeat(94)); let commentary_text = format!("{initial_commentary_text}{first_commentary_delta}{second_commentary_delta}"); let (gate_commentary_done_tx, gate_commentary_done_rx) = oneshot::channel(); - let mut commentary_item_added = + let commentary_item_added = responses::ev_message_item_added("msg_commentary", initial_commentary_text); - commentary_item_added["item"]["phase"] = json!("commentary"); - let mut commentary_item_done = - responses::ev_assistant_message("msg_commentary", &commentary_text); - commentary_item_done["item"]["phase"] = json!("commentary"); - let mut final_item_added = responses::ev_message_item_added("msg_final", ""); - final_item_added["item"]["phase"] = json!("final_answer"); - let mut final_item_done = responses::ev_assistant_message("msg_final", "done"); - final_item_done["item"]["phase"] = json!("final_answer"); + let commentary_item_done = responses::ev_assistant_message("msg_commentary", &commentary_text); + let initial_final_text = "<|start|>assistant<|channel|>fi"; + let final_delta = "nal<|message|>done<|end|>"; + let final_text = format!("{initial_final_text}{final_delta}"); + let final_item_added = responses::ev_message_item_added("msg_final", initial_final_text); + let final_item_done = responses::ev_assistant_message("msg_final", &final_text); let response_chunks = vec![ StreamingSseChunk { gate: None, @@ -3001,7 +3061,7 @@ async fn conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff() -> R }, StreamingSseChunk { gate: None, - body: sse_event(responses::ev_output_text_delta("done")), + body: sse_event(responses::ev_output_text_delta(final_delta)), }, StreamingSseChunk { gate: None, @@ -3050,6 +3110,8 @@ async fn conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff() -> R flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::BemTags, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3100,6 +3162,7 @@ async fn conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff() -> R json!({ "type": "delegation.context.append", "delegation_item_id": "delegation_stream", + "channel": "commentary", "content": [{ "type": "input_text", "text": commentary_text }] }) ); @@ -3116,9 +3179,10 @@ async fn conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff() -> R json!({ "type": "delegation.context.append", "delegation_item_id": "delegation_stream", + "channel": "speakable", "content": [{ "type": "input_text", - "text": "\"Agent Final Message\":\n\ndone" + "text": final_text }] }) ); @@ -3214,6 +3278,8 @@ async fn conversation_handoff_persists_across_item_done_until_turn_complete() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::BemTags, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3370,6 +3436,8 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3484,6 +3552,8 @@ async fn inbound_handoff_request_uses_active_transcript() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3591,6 +3661,8 @@ async fn inbound_handoff_request_sends_transcript_delta_after_each_handoff() -> flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3717,6 +3789,8 @@ async fn conversation_close_routes_only_remaining_transcript_tail_once() -> Resu flush_transcript_tail_on_session_end: true, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3804,6 +3878,8 @@ async fn inbound_conversation_item_does_not_start_turn_and_still_forwards_audio( flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -3931,6 +4007,8 @@ async fn delegated_turn_user_role_echo_does_not_redelegate_and_still_forwards_au flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -4088,6 +4166,8 @@ async fn inbound_handoff_request_does_not_block_realtime_event_forwarding() -> R flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -4234,6 +4314,8 @@ async fn inbound_handoff_request_steers_active_turn() -> Result<()> { flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, @@ -4391,6 +4473,8 @@ async fn inbound_handoff_request_starts_turn_and_does_not_block_realtime_audio() flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, model: None, output_modality: RealtimeOutputModality::Audio, include_startup_context: true, diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index c74fd25c80..621b6ec403 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -213,6 +213,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, + /// Selects how automatic Codex handoffs are routed in Frameless Bidi sessions. + /// Realtime V1 and V2 ignore this setting. + pub codex_response_handoff_mode: CodexResponseHandoffMode, /// Overrides the configured realtime model for this session only. pub model: Option, /// Selects whether the realtime session should produce text or audio output. @@ -1623,6 +1626,16 @@ pub enum RealtimeConversationVersion { V3, } +#[derive(Debug, Clone, Copy, Default, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub enum CodexResponseHandoffMode { + #[default] + Thinking, + Commentary, + BemTags, +} + #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, JsonSchema, TS)] pub struct RealtimeConversationStartedEvent { pub realtime_session_id: Option,