From 3c67d5f8be29752facc04cf059491f7139f1a9de Mon Sep 17 00:00:00 2001 From: Guinness Chen Date: Wed, 24 Jun 2026 15:39:19 -0700 Subject: [PATCH] flush trailing realtime transcript tail --- .../src/protocol/common.rs | 7 + .../src/protocol/v2/realtime.rs | 4 + .../src/request_processors/turn_processor.rs | 3 + .../tests/suite/v2/experimental_api.rs | 2 + .../tests/suite/v2/realtime_conversation.rs | 9 + .../endpoint/realtime_websocket/methods.rs | 57 +++++ codex-rs/core/src/realtime_conversation.rs | 88 +++++-- codex-rs/core/src/session/handlers.rs | 2 +- codex-rs/core/tests/suite/compact_remote.rs | 1 + .../core/tests/suite/realtime_conversation.rs | 216 ++++++++++++++++-- codex-rs/protocol/src/protocol.rs | 3 + 11 files changed, 357 insertions(+), 35 deletions(-) diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index d8db99d4ffeb..c4b94d8b3fb3 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -3195,6 +3195,7 @@ mod tests { request_id: RequestId::Integer(9), params: v2::ThreadRealtimeStartParams { client_managed_handoffs: Some(true), + flush_transcript_tail_on_session_end: Some(true), codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: Some("silent context".to_string()), @@ -3216,6 +3217,7 @@ mod tests { "params": { "threadId": "thr_123", "clientManagedHandoffs": true, + "flushTranscriptTailOnSessionEnd": true, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, "codexResponseHandoffPrefix": "silent context", @@ -3240,6 +3242,7 @@ mod tests { request_id: RequestId::Integer(9), params: v2::ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3261,6 +3264,7 @@ mod tests { "params": { "threadId": "thr_123", "clientManagedHandoffs": null, + "flushTranscriptTailOnSessionEnd": null, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, "codexResponseHandoffPrefix": null, @@ -3280,6 +3284,7 @@ mod tests { request_id: RequestId::Integer(9), params: v2::ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3301,6 +3306,7 @@ mod tests { "params": { "threadId": "thr_123", "clientManagedHandoffs": null, + "flushTranscriptTailOnSessionEnd": null, "codexResponsesAsItems": null, "codexResponseItemPrefix": null, "codexResponseHandoffPrefix": null, @@ -3518,6 +3524,7 @@ mod tests { request_id: RequestId::Integer(1), params: v2::ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, 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 4b8b2056ea0b..b78f390bec7c 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/realtime.rs @@ -70,6 +70,10 @@ pub struct ThreadRealtimeStartParams { /// them automatically. Defaults to false. #[ts(optional = nullable)] pub client_managed_handoffs: Option, + /// Routes any transcript tail remaining at session end through Codex. Defaults to false. + /// TODO: Remove this rollout knob once transcript-tail flushing is always enabled. + #[ts(optional = nullable)] + pub flush_transcript_tail_on_session_end: Option, /// Sends automatic Codex responses as realtime conversation items instead of handoff appends. #[ts(optional = nullable)] pub codex_responses_as_items: Option, 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 bb074e457ba4..3ff64cd34e3c 100644 --- a/codex-rs/app-server/src/request_processors/turn_processor.rs +++ b/codex-rs/app-server/src/request_processors/turn_processor.rs @@ -1002,6 +1002,9 @@ impl TurnRequestProcessor { thread.as_ref(), Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: params.client_managed_handoffs.unwrap_or(false), + flush_transcript_tail_on_session_end: params + .flush_transcript_tail_on_session_end + .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_prefix: params.codex_response_handoff_prefix, 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 d9d3469936e7..80fdb0317151 100644 --- a/codex-rs/app-server/tests/suite/v2/experimental_api.rs +++ b/codex-rs/app-server/tests/suite/v2/experimental_api.rs @@ -82,6 +82,7 @@ async fn realtime_conversation_start_requires_experimental_api_capability() -> R let request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -198,6 +199,7 @@ async fn realtime_webrtc_start_requires_experimental_api_capability() -> Result< let request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, 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 26b46c431773..1d8bd90e6da8 100644 --- a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs +++ b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs @@ -349,6 +349,7 @@ impl RealtimeE2eHarness { .mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs, + flush_transcript_tail_on_session_end: None, thread_id: self.thread_id.clone(), codex_response_item_prefix: codex_responses_as_items .unwrap_or(false) @@ -410,6 +411,7 @@ impl RealtimeE2eHarness { .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: codex_responses_as_items .unwrap_or(false) .then(|| RESPONSE_ITEM_PREFIX.to_string()), @@ -672,6 +674,7 @@ async fn realtime_conversation_streams_v2_notifications() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -963,6 +966,7 @@ async fn realtime_start_can_skip_startup_context() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1061,6 +1065,7 @@ async fn realtime_text_output_modality_requests_text_output_and_final_transcript let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1242,6 +1247,7 @@ async fn realtime_conversation_stop_emits_closed_notification() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1346,6 +1352,7 @@ async fn realtime_webrtc_start_emits_sdp_notification() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2701,6 +2708,7 @@ async fn realtime_webrtc_start_surfaces_backend_error() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2767,6 +2775,7 @@ async fn realtime_conversation_requires_feature_flag() -> Result<()> { let start_request_id = mcp .send_thread_realtime_start_request(ThreadRealtimeStartParams { client_managed_handoffs: None, + flush_transcript_tail_on_session_end: None, codex_responses_as_items: None, codex_response_item_prefix: None, codex_response_handoff_prefix: None, 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 499635579f00..0ae443d00acb 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs @@ -391,6 +391,13 @@ impl RealtimeWebsocketWriter { } impl RealtimeWebsocketEvents { + pub async fn take_transcript_tail(&self) -> Vec { + let mut active_transcript = self.active_transcript.lock().await; + let tail = active_transcript.entries[active_transcript.last_handoff_entry_count..].to_vec(); + active_transcript.last_handoff_entry_count = active_transcript.entries.len(); + tail + } + pub async fn next_event(&self) -> Result, ApiError> { if self.is_closed.load(Ordering::SeqCst) { return Ok(None); @@ -962,6 +969,56 @@ mod tests { ); } + #[tokio::test] + async fn takes_only_transcript_after_last_handoff_once() { + let (_tx_message, rx_message) = async_channel::unbounded(); + let events = RealtimeWebsocketEvents { + rx_message, + active_transcript: Arc::new(Mutex::new(ActiveTranscriptState::default())), + event_parser: RealtimeEventParser::V1, + is_closed: Arc::new(AtomicBool::new(false)), + }; + + assert_eq!(events.take_transcript_tail().await, vec![]); + + let mut covered = RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta { + delta: "already handed off".to_string(), + }); + events.update_active_transcript(&mut covered).await; + let mut handoff = RealtimeEvent::HandoffRequested(RealtimeHandoffRequested { + handoff_id: "handoff_1".to_string(), + item_id: "item_1".to_string(), + input_transcript: "already handed off".to_string(), + active_transcript: vec![], + }); + events.update_active_transcript(&mut handoff).await; + assert_eq!( + handoff, + RealtimeEvent::HandoffRequested(RealtimeHandoffRequested { + handoff_id: "handoff_1".to_string(), + item_id: "item_1".to_string(), + input_transcript: "already handed off".to_string(), + active_transcript: vec![RealtimeTranscriptEntry { + role: "user".to_string(), + text: "already handed off".to_string(), + }], + }) + ); + + let mut tail = RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta { + delta: "tail".to_string(), + }); + events.update_active_transcript(&mut tail).await; + assert_eq!( + events.take_transcript_tail().await, + vec![RealtimeTranscriptEntry { + role: "assistant".to_string(), + text: "tail".to_string(), + }] + ); + assert_eq!(events.take_transcript_tail().await, vec![]); + } + #[test] fn parse_input_transcript_delta_event() { let payload = json!({ diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index 0c81133be11b..3915c8190dda 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -48,6 +48,7 @@ use codex_protocol::protocol::RealtimeConversationSdpEvent; use codex_protocol::protocol::RealtimeConversationStartedEvent; use codex_protocol::protocol::RealtimeHandoffRequested; use codex_protocol::protocol::RealtimeOutputModality; +use codex_protocol::protocol::RealtimeTranscriptEntry; use codex_protocol::protocol::RealtimeVoice; use codex_protocol::protocol::RealtimeVoicesList; use http::HeaderMap; @@ -59,6 +60,7 @@ use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; use tokio::sync::Mutex; use tokio::task::JoinHandle; +use tokio_util::sync::CancellationToken; use tracing::debug; use tracing::error; use tracing::info; @@ -80,6 +82,7 @@ const REALTIME_V2_STEER_ACKNOWLEDGEMENT: &str = "This was sent to steer the previous background agent task."; const REALTIME_ACTIVE_RESPONSE_ERROR_PREFIX: &str = "Conversation already has an active response in progress:"; +const REALTIME_SESSION_ENDED_HANDOFF_INSTRUCTION: &str = "The user just ended their realtime session. Here is the remaining handoff/transcript tail. You probably do not have to do anything; acknowledge the handoff unless the transcript itself asks for something."; #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum RealtimeConversationEnd { @@ -89,7 +92,7 @@ enum RealtimeConversationEnd { } enum RealtimeFanoutTaskStop { - Abort, + Await, Detach, } @@ -203,6 +206,9 @@ struct RealtimeInputTask { handoff_state: RealtimeHandoffState, session_kind: RealtimeSessionKind, event_parser: RealtimeEventParser, + flush_transcript_tail_on_session_end: bool, + transcript_tail_tx: Sender, + stop_token: CancellationToken, } struct RealtimeInputChannels { @@ -242,12 +248,14 @@ struct ConversationState { input_task: JoinHandle<()>, fanout_task: Option>, realtime_active: Arc, + stop_token: CancellationToken, } struct RealtimeStart { api_provider: ApiProvider, extra_headers: Option, client_managed_handoffs: bool, + flush_transcript_tail_on_session_end: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, codex_response_handoff_prefix: Option, @@ -260,6 +268,7 @@ struct RealtimeStart { struct RealtimeStartOutput { realtime_active: Arc, events_rx: Receiver, + transcript_tail_rx: Receiver, sdp: Option, } @@ -294,7 +303,7 @@ impl RealtimeConversationManager { guard.take() }; if let Some(state) = previous_state { - stop_conversation_state(state, RealtimeFanoutTaskStop::Abort).await; + stop_conversation_state(state, RealtimeFanoutTaskStop::Await).await; } self.start_inner(start).await @@ -305,6 +314,7 @@ impl RealtimeConversationManager { api_provider, extra_headers, client_managed_handoffs, + flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, codex_response_handoff_prefix, @@ -327,8 +337,10 @@ impl RealtimeConversationManager { async_channel::bounded::(HANDOFF_OUT_QUEUE_CAPACITY); let (events_tx, events_rx) = async_channel::bounded::(OUTPUT_EVENTS_QUEUE_CAPACITY); + let (transcript_tail_tx, transcript_tail_rx) = async_channel::bounded::(1); let realtime_active = Arc::new(AtomicBool::new(true)); + let stop_token = CancellationToken::new(); let handoff = RealtimeHandoffState::new( handoff_output_tx, client_managed_handoffs, @@ -364,6 +376,9 @@ impl RealtimeConversationManager { session_kind, event_parser, realtime_active: Arc::clone(&realtime_active), + flush_transcript_tail_on_session_end, + transcript_tail_tx, + stop_token: stop_token.clone(), }); (task, Some(call.sdp)) } else { @@ -385,6 +400,9 @@ impl RealtimeConversationManager { handoff_state: handoff.clone(), session_kind, event_parser, + flush_transcript_tail_on_session_end, + transcript_tail_tx, + stop_token: stop_token.clone(), }); (task, None) }; @@ -398,10 +416,12 @@ impl RealtimeConversationManager { input_task: task, fanout_task: None, realtime_active: Arc::clone(&realtime_active), + stop_token, }); Ok(RealtimeStartOutput { realtime_active, events_rx, + transcript_tail_rx, sdp, }) } @@ -652,7 +672,7 @@ impl RealtimeConversationManager { }; if let Some(state) = state { - stop_conversation_state(state, RealtimeFanoutTaskStop::Abort).await; + stop_conversation_state(state, RealtimeFanoutTaskStop::Await).await; } Ok(()) } @@ -663,13 +683,12 @@ async fn stop_conversation_state( fanout_task_stop: RealtimeFanoutTaskStop, ) { state.realtime_active.store(false, Ordering::Relaxed); - state.input_task.abort(); + state.stop_token.cancel(); let _ = state.input_task.await; if let Some(fanout_task) = state.fanout_task.take() { match fanout_task_stop { - RealtimeFanoutTaskStop::Abort => { - fanout_task.abort(); + RealtimeFanoutTaskStop::Await => { let _ = fanout_task.await; } RealtimeFanoutTaskStop::Detach => {} @@ -716,6 +735,7 @@ struct PreparedRealtimeConversationStart { api_provider: ApiProvider, extra_headers: Option, client_managed_handoffs: bool, + flush_transcript_tail_on_session_end: bool, codex_responses_as_items: bool, codex_response_item_prefix: Option, codex_response_handoff_prefix: Option, @@ -800,6 +820,7 @@ async fn prepare_realtime_start( api_provider, extra_headers, client_managed_handoffs: params.client_managed_handoffs, + 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_prefix: params.codex_response_handoff_prefix, @@ -966,6 +987,7 @@ async fn handle_start_inner( api_provider, extra_headers, client_managed_handoffs, + flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, codex_response_handoff_prefix, @@ -984,6 +1006,7 @@ async fn handle_start_inner( api_provider, extra_headers, client_managed_handoffs, + flush_transcript_tail_on_session_end, codex_responses_as_items, codex_response_item_prefix, codex_response_handoff_prefix, @@ -1008,6 +1031,7 @@ async fn handle_start_inner( let RealtimeStartOutput { realtime_active, events_rx, + transcript_tail_rx, sdp, } = start_output; if let Some(sdp) = sdp { @@ -1027,10 +1051,8 @@ async fn handle_start_inner( msg, }; let mut end = RealtimeConversationEnd::TransportClosed; + // Drain already-parsed events so a queued handoff is routed before the final tail. while let Ok(event) = events_rx.recv().await { - if !fanout_realtime_active.load(Ordering::Relaxed) { - break; - } match &event { RealtimeEvent::AudioOut(_) => {} _ => { @@ -1054,9 +1076,6 @@ async fn handle_start_inner( let sess_for_routed_text = Arc::clone(&sess_clone); sess_for_routed_text.route_realtime_text_input(text).await; } - if !fanout_realtime_active.load(Ordering::Relaxed) { - break; - } sess_clone .send_event_raw(ev(EventMsg::RealtimeConversationRealtime( RealtimeConversationRealtimeEvent { @@ -1065,6 +1084,9 @@ async fn handle_start_inner( ))) .await; } + if let Ok(text) = transcript_tail_rx.recv().await { + sess_clone.route_realtime_text_input(text).await; + } if fanout_realtime_active.swap(false, Ordering::Relaxed) { match end { RealtimeConversationEnd::TransportClosed => { @@ -1103,8 +1125,11 @@ pub(crate) async fn handle_audio( } fn realtime_transcript_delta_from_handoff(handoff: &RealtimeHandoffRequested) -> Option { - let active_transcript = handoff - .active_transcript + realtime_transcript_delta(&handoff.active_transcript) +} + +fn realtime_transcript_delta(active_transcript: &[RealtimeTranscriptEntry]) -> Option { + let active_transcript = active_transcript .iter() .map(|entry| format!("{role}: {text}", role = entry.role, text = entry.text)) .collect::>() @@ -1254,6 +1279,9 @@ struct RealtimeWebrtcSidebandInputTask { session_kind: RealtimeSessionKind, event_parser: RealtimeEventParser, realtime_active: Arc, + flush_transcript_tail_on_session_end: bool, + transcript_tail_tx: Sender, + stop_token: CancellationToken, } fn spawn_webrtc_sideband_input_task(input: RealtimeWebrtcSidebandInputTask) -> JoinHandle<()> { @@ -1268,6 +1296,9 @@ fn spawn_webrtc_sideband_input_task(input: RealtimeWebrtcSidebandInputTask) -> J session_kind, event_parser, realtime_active, + flush_transcript_tail_on_session_end, + transcript_tail_tx, + stop_token, } = input; tokio::spawn(async move { @@ -1275,15 +1306,15 @@ fn spawn_webrtc_sideband_input_task(input: RealtimeWebrtcSidebandInputTask) -> J return; } - let connection = match client - .connect_webrtc_sideband( + let connection = match tokio::select! { + connection = client.connect_webrtc_sideband( session_config, &call_id, sideband_headers, default_headers(), - ) - .await - { + ) => connection, + _ = stop_token.cancelled() => return, + } { Ok(connection) => connection, Err(err) => { if realtime_active.load(Ordering::Relaxed) { @@ -1311,6 +1342,9 @@ fn spawn_webrtc_sideband_input_task(input: RealtimeWebrtcSidebandInputTask) -> J handoff_state, session_kind, event_parser, + flush_transcript_tail_on_session_end, + transcript_tail_tx, + stop_token, }) .await; }) @@ -1327,6 +1361,9 @@ async fn run_realtime_input_task(input: RealtimeInputTask) { handoff_state, session_kind, event_parser, + flush_transcript_tail_on_session_end, + transcript_tail_tx, + stop_token, } = input; let mut output_audio_state: Option = None; @@ -1334,6 +1371,7 @@ async fn run_realtime_input_task(input: RealtimeInputTask) { loop { let result = tokio::select! { + _ = stop_token.cancelled() => break, // Text input that should be sent into realtime. text = text_rx.recv() => { handle_text_input( @@ -1378,6 +1416,18 @@ async fn run_realtime_input_task(input: RealtimeInputTask) { break; } } + + if flush_transcript_tail_on_session_end + && let Some(transcript_delta) = + realtime_transcript_delta(&events.take_transcript_tail().await) + { + let _ = transcript_tail_tx + .send(wrap_realtime_delegation_input( + REALTIME_SESSION_ENDED_HANDOFF_INSTRUCTION, + Some(&transcript_delta), + )) + .await; + } } async fn handle_text_input( diff --git a/codex-rs/core/src/session/handlers.rs b/codex-rs/core/src/session/handlers.rs index 9733a9f66aec..924750bc34f6 100644 --- a/codex-rs/core/src/session/handlers.rs +++ b/codex-rs/core/src/session/handlers.rs @@ -589,8 +589,8 @@ async fn shutdown_session_runtime(sess: &Arc) { if let Some(startup_prewarm) = sess.take_session_startup_prewarm().await { startup_prewarm.abort().await; } - sess.abort_all_tasks(TurnAbortReason::Interrupted).await; let _ = sess.conversation.shutdown().await; + sess.abort_all_tasks(TurnAbortReason::Interrupted).await; sess.services .unified_exec_manager .terminate_all_processes() diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index d0e69d656944..e22f70b052f6 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -209,6 +209,7 @@ async fn start_realtime_conversation(codex: &codex_core::CodexThread) -> Result< codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index 1a95223bffb0..f4262e8e4021 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -285,6 +285,7 @@ async fn conversation_start_audio_text_close_round_trip() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -431,6 +432,7 @@ async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -526,6 +528,7 @@ async fn conversation_webrtc_start_posts_generated_session() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -713,6 +716,7 @@ async fn conversation_webrtc_start_uses_avas_query() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -811,6 +815,7 @@ async fn conversation_webrtc_default_v1_ignores_configured_v2_voice() -> Result< test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -871,6 +876,7 @@ async fn conversation_webrtc_default_v1_rejects_explicit_v2_voice() -> Result<() test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -941,6 +947,7 @@ async fn conversation_webrtc_start_uses_configured_call_base_url_for_avas() -> R test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1035,6 +1042,7 @@ async fn conversation_webrtc_close_while_sideband_connecting_drops_pending_join( test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1126,6 +1134,7 @@ async fn conversation_webrtc_sideband_connect_failure_closes_with_error() -> Res test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1219,6 +1228,7 @@ async fn conversation_start_uses_openai_env_key_fallback_with_chatgpt_auth() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1271,27 +1281,46 @@ async fn conversation_start_uses_openai_env_key_fallback_with_chatgpt_auth() -> Ok(()) } -#[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn conversation_transport_close_emits_closed_event() -> Result<()> { +async fn assert_transport_close_tail_flush( + flush_transcript_tail_on_session_end: bool, +) -> Result<()> { skip_if_no_network!(Ok(())); - let session_updated = vec![json!({ - "type": "session.updated", - "session": { "id": "sess_1", "instructions": "backend prompt" } - })]; - let server = start_websocket_server(vec![vec![], vec![session_updated]]).await; + let api_server = start_mock_server().await; + let response_mock = responses::mount_sse_once( + &api_server, + responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_assistant_message("msg-1", "ok"), + responses::ev_completed("resp-1"), + ]), + ) + .await; + let realtime_server = start_websocket_server(vec![vec![vec![ + json!({ + "type": "session.updated", + "session": { "id": "sess_1", "instructions": "backend prompt" } + }), + json!({ + "type": "conversation.input_transcript.delta", + "delta": "transport tail" + }), + ]]]) + .await; - let mut builder = test_codex(); - let test = builder.build_with_websocket_server(&server).await?; - assert!( - server - .wait_for_handshakes(/*expected*/ 1, Duration::from_secs(2)) - .await - ); + let mut builder = test_codex().with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.realtime.version = RealtimeWsVersion::V1; + } + }); + let test = builder.build(&api_server).await?; test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1334,8 +1363,27 @@ async fn conversation_transport_close_emits_closed_event() -> Result<()> { }) .await; assert_eq!(closed.reason.as_deref(), Some("transport_closed")); + if flush_transcript_tail_on_session_end { + let deadline = tokio::time::Instant::now() + Duration::from_secs(2); + while response_mock.requests().is_empty() { + assert!(tokio::time::Instant::now() < deadline); + tokio::time::sleep(Duration::from_millis(10)).await; + } + assert!(response_mock.single_request().message_input_texts("user").iter().any(|text| text + == "\n The user just ended their realtime session. Here is the remaining handoff/transcript tail. You probably do not have to do anything; acknowledge the handoff unless the transcript itself asks for something.\n user: transport tail\n")); + } else { + tokio::time::sleep(Duration::from_millis(200)).await; + assert!(response_mock.requests().is_empty()); + } - server.shutdown().await; + realtime_server.shutdown().await; + Ok(()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn conversation_transport_close_tail_flush_is_opt_in() -> Result<()> { + assert_transport_close_tail_flush(/*flush_transcript_tail_on_session_end*/ false).await?; + assert_transport_close_tail_flush(/*flush_transcript_tail_on_session_end*/ true).await?; Ok(()) } @@ -1389,6 +1437,7 @@ async fn conversation_start_preflight_failure_emits_realtime_error_only() -> Res test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1440,6 +1489,7 @@ async fn conversation_start_connect_failure_emits_realtime_error_only() -> Resul test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1539,6 +1589,7 @@ async fn conversation_second_start_replaces_runtime() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1569,6 +1620,7 @@ async fn conversation_second_start_replaces_runtime() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1670,6 +1722,7 @@ async fn conversation_uses_experimental_realtime_ws_base_url_override() -> Resul test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1739,6 +1792,7 @@ async fn conversation_uses_default_realtime_backend_prompt() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1816,6 +1870,7 @@ async fn conversation_uses_empty_instructions_for_null_or_empty_prompt() -> Resu test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1886,6 +1941,7 @@ async fn conversation_uses_explicit_start_voice() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1948,6 +2004,7 @@ async fn conversation_uses_configured_realtime_voice() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -1998,6 +2055,7 @@ async fn conversation_rejects_voice_for_wrong_realtime_version() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2049,6 +2107,7 @@ async fn conversation_uses_experimental_realtime_ws_backend_prompt_override() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2126,6 +2185,7 @@ async fn conversation_uses_experimental_realtime_ws_startup_context_override() - test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2197,6 +2257,7 @@ async fn conversation_disables_realtime_startup_context_with_empty_override() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2261,6 +2322,7 @@ async fn conversation_start_injects_startup_context_from_thread_history() -> Res test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2380,6 +2442,7 @@ async fn conversation_startup_context_current_thread_selects_many_turns_by_budge codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2492,6 +2555,7 @@ async fn conversation_startup_context_falls_back_to_workspace_map() -> Result<() test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2556,6 +2620,7 @@ async fn conversation_startup_context_is_truncated_and_sent_once_per_start() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2641,6 +2706,7 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2742,6 +2808,7 @@ async fn realtime_v2_noop_tool_call_returns_empty_function_output_without_respon test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2845,6 +2912,7 @@ async fn conversation_mirrors_assistant_message_text_to_realtime_handoff() -> Re test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -2984,6 +3052,7 @@ async fn conversation_handoff_persists_across_item_done_until_turn_complete() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: Some(SILENT_CONTEXT_PREFIX.to_string()), @@ -3140,6 +3209,7 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3254,6 +3324,7 @@ async fn inbound_handoff_request_uses_active_transcript() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3361,6 +3432,7 @@ async fn inbound_handoff_request_sends_transcript_delta_after_each_handoff() -> test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3426,6 +3498,115 @@ async fn inbound_handoff_request_sends_transcript_delta_after_each_handoff() -> Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn conversation_close_routes_only_remaining_transcript_tail_once() -> Result<()> { + skip_if_no_network!(Ok(())); + + let api_server = start_mock_server().await; + let response_mock = responses::mount_sse_sequence( + &api_server, + vec![ + responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_assistant_message("msg-1", "first ok"), + responses::ev_completed("resp-1"), + ]), + responses::sse(vec![ + responses::ev_response_created("resp-2"), + responses::ev_assistant_message("msg-2", "tail ok"), + responses::ev_completed("resp-2"), + ]), + ], + ) + .await; + let realtime_server = start_websocket_server(vec![vec![ + vec![ + json!({ + "type": "session.updated", + "session": { "id": "sess_tail", "instructions": "backend prompt" } + }), + json!({ + "type": "conversation.input_transcript.delta", + "delta": "already handed off" + }), + json!({ + "type": "conversation.handoff.requested", + "handoff_id": "handoff_tail", + "item_id": "item_tail", + "input_transcript": "already handed off" + }), + json!({ + "type": "conversation.output_transcript.delta", + "delta": "remaining answer" + }), + json!({ + "type": "conversation.input_transcript.delta", + "delta": "remaining question" + }), + ], + vec![], + ]]) + .await; + let mut builder = test_codex().with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.realtime.version = RealtimeWsVersion::V1; + } + }); + let test = builder.build(&api_server).await?; + + test.codex + .submit(Op::RealtimeConversationStart(ConversationStartParams { + client_managed_handoffs: false, + flush_transcript_tail_on_session_end: true, + codex_responses_as_items: false, + codex_response_item_prefix: None, + codex_response_handoff_prefix: None, + model: None, + output_modality: RealtimeOutputModality::Audio, + include_startup_context: true, + prompt: Some(Some("backend prompt".to_string())), + realtime_session_id: None, + transport: None, + version: None, + voice: None, + })) + .await?; + + wait_for_event(&test.codex, |event| { + matches!(event, EventMsg::TurnComplete(_)) + }) + .await; + test.codex.submit(Op::RealtimeConversationClose).await?; + + let closed = wait_for_event_match(&test.codex, |msg| match msg { + EventMsg::RealtimeConversationClosed(closed) => Some(closed.clone()), + _ => None, + }) + .await; + assert_eq!(closed.reason.as_deref(), Some("requested")); + + let deadline = tokio::time::Instant::now() + Duration::from_secs(2); + while response_mock.requests().len() < 2 { + assert!(tokio::time::Instant::now() < deadline); + tokio::time::sleep(Duration::from_millis(10)).await; + } + + test.codex.submit(Op::RealtimeConversationClose).await?; + tokio::time::sleep(Duration::from_millis(200)).await; + + let requests = response_mock.requests(); + assert_eq!(requests.len(), 2); + assert!(requests[0].message_input_texts("user").iter().any(|text| text + == "\n already handed off\n user: already handed off\n")); + assert!(requests[1].message_input_texts("user").iter().any(|text| text + == "\n The user just ended their realtime session. Here is the remaining handoff/transcript tail. You probably do not have to do anything; acknowledge the handoff unless the transcript itself asks for something.\n assistant: remaining answer\nuser: remaining question\n")); + + realtime_server.shutdown().await; + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn inbound_conversation_item_does_not_start_turn_and_still_forwards_audio() -> Result<()> { skip_if_no_network!(Ok(())); @@ -3466,6 +3647,7 @@ async fn inbound_conversation_item_does_not_start_turn_and_still_forwards_audio( test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3593,6 +3775,7 @@ async fn delegated_turn_user_role_echo_does_not_redelegate_and_still_forwards_au test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3750,6 +3933,7 @@ async fn inbound_handoff_request_does_not_block_realtime_event_forwarding() -> R test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -3896,6 +4080,7 @@ async fn inbound_handoff_request_steers_active_turn() -> Result<()> { test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, @@ -4053,6 +4238,7 @@ async fn inbound_handoff_request_starts_turn_and_does_not_block_realtime_audio() test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: false, + flush_transcript_tail_on_session_end: false, codex_responses_as_items: false, codex_response_item_prefix: None, codex_response_handoff_prefix: None, diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 2c97af10aa01..ca509a5a1729 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -194,6 +194,9 @@ pub struct McpServerRefreshConfig { pub struct ConversationStartParams { /// Whether Codex response handoffs are managed through explicit client append calls. pub client_managed_handoffs: bool, + /// Whether to route any remaining transcript tail through Codex when the session ends. + /// TODO: Remove this rollout knob once transcript-tail flushing is always enabled. + pub flush_transcript_tail_on_session_end: bool, /// Sends automatic Codex responses as realtime conversation items instead of handoff appends. pub codex_responses_as_items: bool, /// Optional prefix added to automatic Codex response items when `codex_responses_as_items` is set.