From 582c528fd7dca3dfd5dc9b0599ae1a6bd60d36ab Mon Sep 17 00:00:00 2001 From: Ahmed Ibrahim Date: Sun, 12 Apr 2026 11:59:57 -0700 Subject: [PATCH] Keep mirrored user text append-only in realtime --- codex-rs/core/src/codex.rs | 2 +- codex-rs/core/src/realtime_conversation.rs | 63 ++++++++++++++----- .../core/tests/suite/realtime_conversation.rs | 18 +++++- ..._turn_is_sent_to_realtime_when_active.snap | 3 +- 4 files changed, 66 insertions(+), 20 deletions(-) diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 9100850ff2..77577674c5 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -5094,7 +5094,7 @@ mod handlers { if sess.conversation.running_state().await.is_none() { return; } - if let Err(err) = sess.conversation.text_in(text).await { + if let Err(err) = sess.conversation.append_user_text(text).await { debug!("failed to mirror user text to realtime conversation: {err}"); } } diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index d22551a573..5137e5747c 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -113,6 +113,18 @@ enum HandoffOutput { }, } +#[derive(Debug, PartialEq, Eq)] +struct RealtimeTextInput { + text: String, + response_mode: RealtimeTextResponseMode, +} + +#[derive(Debug, PartialEq, Eq)] +enum RealtimeTextResponseMode { + AppendOnly, + RequestResponse, +} + #[derive(Debug, PartialEq, Eq)] struct OutputAudioState { item_id: String, @@ -184,7 +196,7 @@ impl RealtimeResponseCreateQueue { struct RealtimeInputTask { writer: RealtimeWebsocketWriter, events: RealtimeWebsocketEvents, - user_text_rx: Receiver, + user_text_rx: Receiver, handoff_output_rx: Receiver, audio_rx: Receiver, events_tx: Sender, @@ -207,7 +219,7 @@ impl RealtimeHandoffState { #[allow(dead_code)] struct ConversationState { audio_tx: Sender, - user_text_tx: Sender, + user_text_tx: Sender, writer: RealtimeWebsocketWriter, handoff: RealtimeHandoffState, input_task: JoinHandle<()>, @@ -306,7 +318,7 @@ impl RealtimeConversationManager { let (audio_tx, audio_rx) = async_channel::bounded::(AUDIO_IN_QUEUE_CAPACITY); let (user_text_tx, user_text_rx) = - async_channel::bounded::(USER_TEXT_IN_QUEUE_CAPACITY); + async_channel::bounded::(USER_TEXT_IN_QUEUE_CAPACITY); let (handoff_output_tx, handoff_output_rx) = async_channel::bounded::(HANDOFF_OUT_QUEUE_CAPACITY); let (events_tx, events_rx) = @@ -403,6 +415,22 @@ impl RealtimeConversationManager { } pub(crate) async fn text_in(&self, text: String) -> CodexResult<()> { + self.queue_text_input(RealtimeTextInput { + text, + response_mode: RealtimeTextResponseMode::RequestResponse, + }) + .await + } + + pub(crate) async fn append_user_text(&self, text: String) -> CodexResult<()> { + self.queue_text_input(RealtimeTextInput { + text, + response_mode: RealtimeTextResponseMode::AppendOnly, + }) + .await + } + + async fn queue_text_input(&self, input: RealtimeTextInput) -> CodexResult<()> { let sender = { let guard = self.state.lock().await; guard.as_ref().map(|state| state.user_text_tx.clone()) @@ -415,7 +443,7 @@ impl RealtimeConversationManager { }; sender - .send(text) + .send(input) .await .map_err(|_| CodexErr::InvalidRequest("conversation is not running".to_string()))?; Ok(()) @@ -904,7 +932,7 @@ pub(crate) async fn handle_text( sub_id: String, params: ConversationTextParams, ) { - debug!(text = %params.text, "[realtime-text] appending realtime conversation text input"); + debug!(text = %params.text, "[realtime-text] sending realtime conversation text input"); if let Err(err) = sess.conversation.text_in(params.text).await { error!("failed to append realtime text: {err}"); if sess.conversation.running_state().await.is_some() { @@ -939,9 +967,9 @@ fn spawn_realtime_input_task(input: RealtimeInputTask) -> JoinHandle<()> { loop { let result = tokio::select! { - // Text typed by the user that should be sent into realtime. + // Text that should be sent into realtime. user_text = user_text_rx.recv() => { - handle_user_text_input( + handle_realtime_text_input( user_text, &writer, &events_tx, @@ -988,16 +1016,16 @@ fn spawn_realtime_input_task(input: RealtimeInputTask) -> JoinHandle<()> { }) } -async fn handle_user_text_input( - text: Result, +async fn handle_realtime_text_input( + input: Result, writer: &RealtimeWebsocketWriter, events_tx: &Sender, session_kind: RealtimeSessionKind, response_create_queue: &mut RealtimeResponseCreateQueue, ) -> anyhow::Result<()> { - let text = text.context("user text input channel closed")?; + let input = input.context("text input channel closed")?; - if let Err(err) = writer.send_conversation_item_create(text).await { + if let Err(err) = writer.send_conversation_item_create(input.text).await { let mapped_error = map_api_error(err); warn!("failed to send input text: {mapped_error}"); let _ = events_tx @@ -1007,11 +1035,14 @@ async fn handle_user_text_input( } match session_kind { RealtimeSessionKind::V1 => {} - RealtimeSessionKind::V2 => { - response_create_queue - .request_create(writer, events_tx, "text") - .await?; - } + RealtimeSessionKind::V2 => match input.response_mode { + RealtimeTextResponseMode::AppendOnly => {} + RealtimeTextResponseMode::RequestResponse => { + response_create_queue + .request_create(writer, events_tx, "text") + .await?; + } + }, } Ok(()) } diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index e5085687e9..6661f1da53 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -1654,7 +1654,7 @@ async fn conversation_user_text_turn_is_sent_to_realtime_when_active() -> Result 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; + config.experimental_realtime_ws_startup_context = Some(String::new()); } }); let test = builder.build(&api_server).await?; @@ -1708,11 +1708,24 @@ async fn conversation_user_text_turn_is_sent_to_realtime_when_active() -> Result ), (true, Some(user_text.to_string())), ); + let realtime_response_create = timeout(Duration::from_millis(200), async { + wait_for_matching_websocket_request( + &realtime_server, + "unexpected realtime response request for mirrored user text", + |request| request.body_json()["type"].as_str() == Some("response.create"), + ) + .await + }) + .await; + assert!( + realtime_response_create.is_err(), + "mirrored user text should not request a realtime response" + ); let realtime_request_body = realtime_text_request.body_json(); let content = &realtime_request_body["item"]["content"][0]; let snapshot = format!( - "type: {}\nitem.type: {}\nitem.role: {}\ncontent[0].type: {}\ncontent[0].text: {}", + "type: {}\nitem.type: {}\nitem.role: {}\ncontent[0].type: {}\ncontent[0].text: {}\nresponse.create: {}", realtime_request_body["type"].as_str().unwrap_or_default(), realtime_request_body["item"]["type"] .as_str() @@ -1722,6 +1735,7 @@ async fn conversation_user_text_turn_is_sent_to_realtime_when_active() -> Result .unwrap_or_default(), content["type"].as_str().unwrap_or_default(), content["text"].as_str().unwrap_or_default(), + realtime_response_create.is_ok(), ); insta::assert_snapshot!( "conversation_user_text_turn_is_sent_to_realtime_when_active", diff --git a/codex-rs/core/tests/suite/snapshots/all__suite__realtime_conversation__conversation_user_text_turn_is_sent_to_realtime_when_active.snap b/codex-rs/core/tests/suite/snapshots/all__suite__realtime_conversation__conversation_user_text_turn_is_sent_to_realtime_when_active.snap index 070f7ce09b..9cd2a55303 100644 --- a/codex-rs/core/tests/suite/snapshots/all__suite__realtime_conversation__conversation_user_text_turn_is_sent_to_realtime_when_active.snap +++ b/codex-rs/core/tests/suite/snapshots/all__suite__realtime_conversation__conversation_user_text_turn_is_sent_to_realtime_when_active.snap @@ -5,5 +5,6 @@ expression: snapshot type: conversation.item.create item.type: message item.role: user -content[0].type: text +content[0].type: input_text content[0].text: typed follow-up for realtime +response.create: false