Keep mirrored user text append-only in realtime

This commit is contained in:
Ahmed Ibrahim
2026-04-12 11:59:57 -07:00
parent b2c7e3f668
commit 582c528fd7
4 changed files with 66 additions and 20 deletions

View File

@@ -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}");
}
}

View File

@@ -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<String>,
user_text_rx: Receiver<RealtimeTextInput>,
handoff_output_rx: Receiver<HandoffOutput>,
audio_rx: Receiver<RealtimeAudioFrame>,
events_tx: Sender<RealtimeEvent>,
@@ -207,7 +219,7 @@ impl RealtimeHandoffState {
#[allow(dead_code)]
struct ConversationState {
audio_tx: Sender<RealtimeAudioFrame>,
user_text_tx: Sender<String>,
user_text_tx: Sender<RealtimeTextInput>,
writer: RealtimeWebsocketWriter,
handoff: RealtimeHandoffState,
input_task: JoinHandle<()>,
@@ -306,7 +318,7 @@ impl RealtimeConversationManager {
let (audio_tx, audio_rx) =
async_channel::bounded::<RealtimeAudioFrame>(AUDIO_IN_QUEUE_CAPACITY);
let (user_text_tx, user_text_rx) =
async_channel::bounded::<String>(USER_TEXT_IN_QUEUE_CAPACITY);
async_channel::bounded::<RealtimeTextInput>(USER_TEXT_IN_QUEUE_CAPACITY);
let (handoff_output_tx, handoff_output_rx) =
async_channel::bounded::<HandoffOutput>(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<String, RecvError>,
async fn handle_realtime_text_input(
input: Result<RealtimeTextInput, RecvError>,
writer: &RealtimeWebsocketWriter,
events_tx: &Sender<RealtimeEvent>,
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(())
}

View File

@@ -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",

View File

@@ -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