diff --git a/codex-rs/tui/src/app/event_dispatch.rs b/codex-rs/tui/src/app/event_dispatch.rs index 8b5ad52bd9..6bcadcf892 100644 --- a/codex-rs/tui/src/app/event_dispatch.rs +++ b/codex-rs/tui/src/app/event_dispatch.rs @@ -35,6 +35,7 @@ impl App { && !matches!( &event, AppEvent::InsertHistoryCell(_) + | AppEvent::CommitRealtimeTranscriptHistory | AppEvent::AgentsOverviewError(_) | AppEvent::ViewAgentsOverviewUnsentPrompt(_) | AppEvent::ResetTranscriptForThreadSwitch @@ -670,6 +671,11 @@ impl App { self.reset_for_thread_switch(tui)?; self.pending_thread_switch_resets -= 1; } + AppEvent::CommitRealtimeTranscriptHistory => { + for cell in self.chat_widget.take_realtime_transcript_history() { + self.insert_history_cell(tui, cell); + } + } AppEvent::InsertHistoryCell(cell) => { self.insert_history_cell(tui, cell); } diff --git a/codex-rs/tui/src/app/tests/realtime_requests.rs b/codex-rs/tui/src/app/tests/realtime_requests.rs index ea1fb9cebb..79e257eab1 100644 --- a/codex-rs/tui/src/app/tests/realtime_requests.rs +++ b/codex-rs/tui/src/app/tests/realtime_requests.rs @@ -5,6 +5,7 @@ use crate::app::tests::session_lifecycle_requests::RealtimeRequestBehavior; use crate::app::tests::session_lifecycle_requests::recorded_params; use crate::app::tests::session_lifecycle_requests::start_recording_realtime_speech_app_server; use crate::app::tests::session_lifecycle_requests::start_recording_remote_app_server; +use crate::chatwidget::commit_realtime_history_events; use crate::chatwidget::tests::make_chatwidget_manual_with_sender; use codex_app_server_protocol::ItemCompletedNotification; use codex_app_server_protocol::ItemStartedNotification; @@ -410,6 +411,7 @@ async fn switching_threads_keeps_the_source_voice_partial_only_on_reattach() { ), /*replay_kind*/ None, ); + commit_realtime_history_events(&mut app.chat_widget, &mut source_events); let rendered = std::iter::from_fn(|| source_events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -553,6 +555,7 @@ async fn queued_voice_caption_after_switch_returns_once_to_its_source_thread() { empty_thread_snapshot(&app, source), /*resume_restored_queue*/ false, ); + commit_realtime_history_events(&mut app.chat_widget, &mut source_events); let rendered = std::iter::from_fn(|| source_events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -625,6 +628,7 @@ async fn replay_reconciles_only_matching_voice_captions_one_for_one() { }, /*resume_restored_queue*/ false, ); + commit_realtime_history_events(&mut app.chat_widget, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -649,7 +653,7 @@ async fn replay_reconciles_only_matching_voice_captions_one_for_one() { } #[tokio::test] -async fn retained_caption_consumes_only_one_matching_answer_fallback_on_reattach() { +async fn retained_caption_consumes_only_one_matching_answer_fallback_on_reattach() -> Result<()> { let (mut app, _initial_events, _ops) = make_test_app_with_channels().await; let source = ThreadId::new(); let (widget, _, mut events, _) = make_chatwidget_manual_with_sender().await; @@ -685,23 +689,26 @@ async fn retained_caption_consumes_only_one_matching_answer_fallback_on_reattach empty_thread_snapshot(&app, source), /*resume_restored_queue*/ false, ); - let rendered = std::iter::from_fn(|| events.try_recv().ok()) - .filter_map(|event| match event { - AppEvent::InsertHistoryCell(cell) => Some( - cell.display_lines(/*width*/ 80) - .into_iter() - .map(|line| line.to_string()) - .collect::>() - .join("\n"), - ), - _ => None, - }) + let (mut app_server, _requests, proxy) = start_recording_remote_app_server(&app.config).await?; + let mut tui = crate::tui::test_support::make_test_tui()?; + while let Ok(event) = events.try_recv() { + Box::pin(app.handle_event(&mut tui, &mut app_server, event)).await?; + app.chat_widget.pre_draw_tick(); + } + let rendered = app + .transcript_cells + .iter() + .flat_map(|cell| cell.display_lines(/*width*/ 80)) + .map(|line| line.to_string()) .collect::>() .join("\n"); assert_eq!(rendered.matches("Same answer").count(), 1); assert_eq!(rendered.matches("Same answer").count(), 1); assert_eq!(rendered.matches("Different answer").count(), 1); assert!(!rendered.contains("[FINAL]")); + app_server.shutdown().await?; + proxy.await??; + Ok(()) } #[tokio::test] @@ -899,6 +906,7 @@ async fn unrendered_buffered_items_do_not_consume_retained_captions() { }, /*resume_restored_queue*/ false, ); + commit_realtime_history_events(&mut app.chat_widget, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -939,6 +947,7 @@ async fn completed_voice_caption_survives_repeated_thread_replacement() { ), /*replay_kind*/ None, ); + commit_realtime_history_events(&mut app.chat_widget, &mut initial_events); let initial = std::iter::from_fn(|| initial_events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -992,6 +1001,7 @@ async fn completed_voice_caption_survives_repeated_thread_replacement() { snapshot.turns.push(typed.clone()); app.replay_thread_snapshot(snapshot, /*resume_restored_queue*/ false); assert!(!app.pending_realtime_transcript_replay.contains_key(&source)); + commit_realtime_history_events(&mut app.chat_widget, &mut source_events); let rendered = std::iter::from_fn(|| source_events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( diff --git a/codex-rs/tui/src/app_event.rs b/codex-rs/tui/src/app_event.rs index 8e00ef9c3c..b1e802455f 100644 --- a/codex-rs/tui/src/app_event.rs +++ b/codex-rs/tui/src/app_event.rs @@ -1060,6 +1060,9 @@ pub(crate) enum AppEvent { InsertHistoryCell(Box), + /// Move visible completed voice captions into history in one app event. + CommitRealtimeTranscriptHistory, + /// Finish buffering initial resume replay after all replay events have been queued. EndInitialHistoryReplayBuffer, diff --git a/codex-rs/tui/src/chatwidget.rs b/codex-rs/tui/src/chatwidget.rs index 4f91c11b20..627e884892 100644 --- a/codex-rs/tui/src/chatwidget.rs +++ b/codex-rs/tui/src/chatwidget.rs @@ -415,6 +415,8 @@ pub(crate) use realtime::realtime_delegation_display_text; pub(crate) use realtime::realtime_delegation_input; #[cfg(test)] pub(crate) use realtime::tests::activate_voice_for_thread; +#[cfg(test)] +pub(crate) use realtime::tests::commit_realtime_history_events; mod reasoning_shortcuts; use self::realtime::RealtimeConversationUiState; mod rendering; @@ -1221,8 +1223,8 @@ impl ChatWidget { } } - fn flush_completed_tool_activity(&mut self) { - if self.transcript.active_cell.as_ref().is_some_and(|cell| { + fn has_completed_tool_activity(&self) -> bool { + self.transcript.active_cell.as_ref().is_some_and(|cell| { cell.as_any() .downcast_ref::() .is_some_and(|cell| !cell.is_active()) @@ -1230,7 +1232,11 @@ impl ChatWidget { .as_any() .downcast_ref::() .is_some_and(|cell| !cell.is_active()) - }) { + }) + } + + fn flush_completed_tool_activity(&mut self) { + if self.has_completed_tool_activity() { self.flush_active_cell(); } } @@ -1240,6 +1246,17 @@ impl ChatWidget { } fn add_boxed_history(&mut self, cell: Box) { + if let Some(active) = self.take_history_insertion_prefix(cell.as_ref()) { + self.app_event_tx.send(AppEvent::InsertHistoryCell(active)); + self.request_pending_usage_output_insertion(); + } + self.app_event_tx.send(AppEvent::InsertHistoryCell(cell)); + } + + fn take_history_insertion_prefix( + &mut self, + cell: &dyn HistoryCell, + ) -> Option> { // Keep the placeholder session header as the active cell until real session info arrives, // so we can merge headers instead of committing a duplicate box to history. let keep_placeholder_header_active = !self.is_session_configured() @@ -1261,24 +1278,15 @@ impl ChatWidget { { // Only break exec grouping if the cell renders visible lines. if !self.has_active_stream_tail() { - self.flush_active_cell(); + return self.transcript.take_active_cell(); } } else if !keep_placeholder_header_active - && self - .transcript - .active_cell - .as_ref() - .is_some_and(|active_cell| { - active_cell.as_any().is::() - || active_cell - .as_any() - .is::() - }) + && self.has_completed_tool_activity() && !cell.transcript_lines(history_width).is_empty() { - self.flush_completed_tool_activity(); + return self.transcript.take_active_cell(); } - self.app_event_tx.send(AppEvent::InsertHistoryCell(cell)); + None } fn enter_review_mode_with_hint(&mut self, hint: String, from_replay: bool) { @@ -1964,11 +1972,16 @@ impl ChatWidget { /// the main viewport updates. pub(crate) fn active_cell_transcript_key(&self) -> Option { let cell = self.transcript.active_cell.as_ref(); - let realtime_cell = self.realtime_conversation.live_transcript_cell.as_ref(); + let mut realtime_cells = self.realtime_conversation.live_transcript_cells(); let token_activity_cell = self.pending_token_activity_output(); let rate_limit_reset_hint = self.pending_rate_limit_reset_hint(); if cell.is_none() - && realtime_cell.is_none() + && self + .realtime_conversation + .live_transcript_cells() + .next() + .is_none() + && self.realtime_conversation.pending_history_cells.is_empty() && token_activity_cell.is_none() && rate_limit_reset_hint.is_none() { @@ -1981,7 +1994,7 @@ impl ChatWidget { .unwrap_or(false), animation_tick: cell .and_then(|cell| cell.transcript_animation_tick()) - .or_else(|| realtime_cell.and_then(|cell| cell.transcript_animation_tick())), + .or_else(|| realtime_cells.find_map(|cell| cell.transcript_animation_tick())), }) } @@ -1999,7 +2012,12 @@ impl ChatWidget { if let Some(cell) = self.transcript.active_cell.as_ref() { lines.extend(cell.transcript_hyperlink_lines(width)); } - if let Some(cell) = self.realtime_conversation.live_transcript_cell.as_ref() { + for cell in self + .realtime_conversation + .pending_history_cells + .iter() + .chain(self.realtime_conversation.live_transcript_cells()) + { let realtime_lines = cell.transcript_hyperlink_lines(width); if !realtime_lines.is_empty() && !lines.is_empty() { lines.push(HyperlinkLine::from("")); diff --git a/codex-rs/tui/src/chatwidget/realtime.rs b/codex-rs/tui/src/chatwidget/realtime.rs index c7330d40d4..4c2e7037ae 100644 --- a/codex-rs/tui/src/chatwidget/realtime.rs +++ b/codex-rs/tui/src/chatwidget/realtime.rs @@ -1,5 +1,6 @@ //! TUI orchestration for an app-server-signaled, locally owned WebRTC voice session. //! Completed captions and both speakers' partials stay bounded across widget replacement. +//! Interleaved speakers retain separate displays so settled caption text never reanimates. mod recording_controls; mod transcript_replay; @@ -157,11 +158,12 @@ pub(super) struct RealtimeConversationUiState { transcript: String, // The other speaker's bounded partial while duplex deltas interleave. interleaved_transcript: Option<(String, String)>, + interleaved_transcript_cell: Option>, transcript_input_generation: Option, assistant_transcript_generation: Option, assistant_caption_started_after_speech_queue: bool, pub(super) live_transcript_cell: Option>, - pending_history_cells: VecDeque>, + pub(super) pending_history_cells: VecDeque>, accepted_transcripts: VecDeque, replay_transcripts: Option>, latest_input_was_voice: bool, @@ -174,6 +176,24 @@ pub(super) struct RealtimeConversationUiState { pending_speech: VecDeque, } +impl RealtimeConversationUiState { + pub(super) fn live_transcript_cells(&self) -> impl Iterator> { + // Keep the user above the reply even when their last packets interleave. + let cells = if self.transcript_role.as_deref() == Some("user") { + [ + &self.live_transcript_cell, + &self.interleaved_transcript_cell, + ] + } else { + [ + &self.interleaved_transcript_cell, + &self.live_transcript_cell, + ] + }; + cells.into_iter().filter_map(Option::as_ref) + } +} + pub(crate) fn realtime_delegation_input(items: &[UserInput]) -> Option<&str> { let [ UserInput::Text { @@ -1021,6 +1041,7 @@ impl ChatWidget { }); }; if let Some((role, text)) = self.realtime_conversation.interleaved_transcript.take() { + self.realtime_conversation.interleaved_transcript_cell = None; retain_partial(role, text); } if let Some(role) = self.realtime_conversation.transcript_role.take() { @@ -1051,6 +1072,7 @@ impl ChatWidget { self.realtime_conversation .pending_history_cells .push_back(cell); + self.bump_active_cell_revision(); } else { // Keep an unfinished caption editable until its late completion or close. self.on_realtime_transcript_delta(record.role.clone(), record.text.clone()); @@ -1237,13 +1259,22 @@ impl ChatWidget { let previous = self.realtime_conversation.transcript_role.take(); let previous_text = std::mem::take(&mut self.realtime_conversation.transcript); let saved = self.realtime_conversation.interleaved_transcript.take(); + let saved_cell = self + .realtime_conversation + .interleaved_transcript_cell + .take(); + let previous_cell = self.realtime_conversation.live_transcript_cell.take(); self.realtime_conversation.transcript = match saved { - Some((saved_role, saved_text)) if saved_role == role => saved_text, + Some((saved_role, saved_text)) if saved_role == role => { + self.realtime_conversation.live_transcript_cell = saved_cell; + saved_text + } _ => String::new(), }; if let Some(previous_role) = previous { self.realtime_conversation.interleaved_transcript = Some((previous_role, previous_text)); + self.realtime_conversation.interleaved_transcript_cell = previous_cell; } self.realtime_conversation.transcript_role = Some(role); } @@ -1307,6 +1338,8 @@ impl ChatWidget { .is_some_and(|(saved_role, _)| saved_role == &role) { self.realtime_conversation.interleaved_transcript = None; + self.realtime_conversation.interleaved_transcript_cell = None; + self.bump_active_cell_revision(); } if text.trim().is_empty() { if let Some(index) = self @@ -1367,6 +1400,7 @@ impl ChatWidget { self.realtime_conversation .pending_history_cells .push_back(cell); + self.bump_active_cell_revision(); self.flush_realtime_transcript_history(); return; } @@ -1470,6 +1504,8 @@ impl ChatWidget { .is_some_and(|(saved_role, _)| saved_role == &role) { self.realtime_conversation.interleaved_transcript = None; + self.realtime_conversation.interleaved_transcript_cell = None; + self.bump_active_cell_revision(); } if text.len() > MAX_TRANSCRIPT_BYTES { let mut end = MAX_TRANSCRIPT_BYTES; @@ -1505,12 +1541,14 @@ impl ChatWidget { self.realtime_conversation .pending_history_cells .push_back(cell); + self.bump_active_cell_revision(); self.flush_realtime_transcript_history(); } } fn finish_realtime_partial_transcripts(&mut self) { if let Some((role, text)) = self.realtime_conversation.interleaved_transcript.take() { + self.realtime_conversation.interleaved_transcript_cell = None; self.on_realtime_transcript_done(role, text); } if let Some(role) = self.realtime_conversation.transcript_role.clone() { @@ -1537,11 +1575,33 @@ impl ChatWidget { { return; } - while let Some(cell) = self.realtime_conversation.pending_history_cells.pop_front() { - self.add_boxed_history(cell); + if !self.realtime_conversation.pending_history_cells.is_empty() { + self.app_event_tx + .send(AppEvent::CommitRealtimeTranscriptHistory); + self.request_redraw(); } } + pub(crate) fn take_realtime_transcript_history(&mut self) -> Vec> { + if self.stream_controller.is_some() + || self.plan_stream_controller.is_some() + || self.pending_stream_consolidations > 0 + { + return Vec::new(); + } + let mut cells = Vec::new(); + while let Some(cell) = self.realtime_conversation.pending_history_cells.pop_front() { + if let Some(active) = self.take_history_insertion_prefix(cell.as_ref()) { + cells.push(active); + } + cells.push(cell); + } + if !cells.is_empty() { + self.bump_active_cell_revision(); + } + cells + } + pub(crate) fn record_realtime_failure(&mut self) { if self.realtime_conversation.phase == RealtimeConversationPhase::Inactive || self.realtime_conversation.failure_recorded @@ -1713,6 +1773,7 @@ impl ChatWidget { self.realtime_conversation .pending_history_cells .push_back(cell); + self.bump_active_cell_revision(); self.realtime_conversation .accepted_transcripts .push_back(RealtimeTranscriptRecord { @@ -1728,7 +1789,11 @@ impl ChatWidget { std::mem::take(&mut self.realtime_conversation.accepted_transcripts); let delegated_reasoning_turns = std::mem::take(&mut self.realtime_conversation.delegated_reasoning_turns); - let had_live_transcript = self.realtime_conversation.live_transcript_cell.is_some(); + let had_live_transcript = self + .realtime_conversation + .live_transcript_cells() + .next() + .is_some(); self.realtime_conversation = RealtimeConversationUiState { attempt_id: self.realtime_conversation.attempt_id, pending_history_cells, diff --git a/codex-rs/tui/src/chatwidget/realtime_tests.rs b/codex-rs/tui/src/chatwidget/realtime_tests.rs index c27108b372..37280b087c 100644 --- a/codex-rs/tui/src/chatwidget/realtime_tests.rs +++ b/codex-rs/tui/src/chatwidget/realtime_tests.rs @@ -32,6 +32,27 @@ use codex_protocol::models::MessagePhase; use futures::future::AbortHandle; use std::collections::VecDeque; +// Model the app's atomic handoff before inspecting its ordinary history events. +pub(crate) fn commit_realtime_history_events( + chat: &mut ChatWidget, + events: &mut tokio::sync::mpsc::UnboundedReceiver, +) { + let mut forwarded = Vec::new(); + while let Ok(event) = events.try_recv() { + match event { + AppEvent::CommitRealtimeTranscriptHistory => forwarded.extend( + chat.take_realtime_transcript_history() + .into_iter() + .map(AppEvent::InsertHistoryCell), + ), + event => forwarded.push(event), + } + } + for event in forwarded { + chat.app_event_tx.send(event); + } +} + fn activate_voice(chat: &mut ChatWidget) -> ThreadId { let thread_id = ThreadId::new(); activate_voice_for_thread(chat, thread_id); diff --git a/codex-rs/tui/src/chatwidget/realtime_tests/caption_replay.rs b/codex-rs/tui/src/chatwidget/realtime_tests/caption_replay.rs index b9ae7e7852..b25084c6e9 100644 --- a/codex-rs/tui/src/chatwidget/realtime_tests/caption_replay.rs +++ b/codex-rs/tui/src/chatwidget/realtime_tests/caption_replay.rs @@ -46,6 +46,7 @@ async fn replay_preserves_typed_updates_before_voice_steers_the_turn() { ReplayKind::ThreadSnapshot, ); chat.flush_answer_stream_with_separator(); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -121,6 +122,7 @@ async fn in_progress_voice_replay_restores_the_late_reasoning_guard() { ); assert!(chat.is_realtime_delegated_reasoning_turn(turn_id)); } + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { let rendered = cell @@ -167,6 +169,7 @@ async fn accepted_voice_answer_with_only_an_old_caption_returns_to_history_on_cl chat.on_realtime_conversation_closed(Some("transport_closed".into())); chat.restore_undelivered_realtime_speech(delivery_id); let mut answers = 0; + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { let rendered = cell @@ -312,6 +315,7 @@ async fn captioned_voice_answer_does_not_duplicate_on_close() { chat.reset_realtime_conversation(); chat.restore_undelivered_realtime_speech(delivery_id); let mut rendered = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { rendered.extend( @@ -362,6 +366,7 @@ async fn restored_partial_caption_accepts_late_completion_without_duplicate_hist "last words" ); assert!(chat.realtime_conversation.accepted_transcripts[0].complete); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( diff --git a/codex-rs/tui/src/chatwidget/realtime_tests/handoffs.rs b/codex-rs/tui/src/chatwidget/realtime_tests/handoffs.rs index fa3e991d14..062a021f12 100644 --- a/codex-rs/tui/src/chatwidget/realtime_tests/handoffs.rs +++ b/codex-rs/tui/src/chatwidget/realtime_tests/handoffs.rs @@ -33,6 +33,7 @@ async fn closing_voice_before_turn_completion_restores_the_delegated_answer_once ); let mut visible_answers = 0; + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { visible_answers += cell @@ -104,6 +105,7 @@ async fn delegated_voice_work_speaks_only_the_completed_final_answer_once() { if spoken_thread_id == thread_id && text.as_str() == "Lucali in Brooklyn" )); assert!(ops.try_recv().is_err()); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { assert!( @@ -152,6 +154,7 @@ async fn delegation_started_before_peer_connection_keeps_its_voice_origin() { )); assert!(ops.try_recv().is_err()); let mut rendered_history = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { let rendered = cell @@ -269,6 +272,7 @@ async fn delegated_item_that_becomes_final_at_turn_completion_is_recoverable() { }; assert_eq!(text.as_str(), "Public answer"); chat.restore_undelivered_realtime_speech(delivery_id); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { for line in cell.display_lines(/*width*/ 80) { @@ -315,6 +319,7 @@ async fn explicit_final_answer_can_explain_private_channel_markers() { ); assert!(ops.try_recv().is_err()); chat.restore_undelivered_realtime_speech(delivery_id); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( diff --git a/codex-rs/tui/src/chatwidget/realtime_tests/lifecycle.rs b/codex-rs/tui/src/chatwidget/realtime_tests/lifecycle.rs index aedbcc2fb1..fe02d4cf6f 100644 --- a/codex-rs/tui/src/chatwidget/realtime_tests/lifecycle.rs +++ b/codex-rs/tui/src/chatwidget/realtime_tests/lifecycle.rs @@ -10,6 +10,7 @@ async fn enabling_voice_on_an_open_thread_snapshots_the_new_thread_notice() { codex_features::Feature::RealtimeConversation, /*enabled*/ true, ); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -33,6 +34,7 @@ async fn voice_cannot_start_in_a_side_conversation() { chat.toggle_realtime_conversation(); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("voice should report that side conversations are unsupported"); }; @@ -129,6 +131,7 @@ async fn audio_failure_cancels_pending_voice_and_reports_the_device_error() { chat.on_realtime_error("speaker stream failed: device disconnected".to_string()); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("voice should report the speaker failure"); }; @@ -209,6 +212,7 @@ async fn voice_becomes_active_only_after_backend_and_current_peer_are_ready() { chat.realtime_conversation.phase, RealtimeConversationPhase::Active ); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("voice start should insert its history banner"); }; @@ -230,6 +234,7 @@ async fn normal_voice_close_renders_the_ended_message() { chat.on_realtime_conversation_closed(Some("transport_closed".into())); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -291,6 +296,7 @@ async fn startup_retry_waits_for_closed_then_uses_a_fresh_attempt_and_preserves_ ); assert!(ops.try_recv().is_err()); let mut rendered = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { rendered.extend( @@ -394,6 +400,7 @@ async fn startup_transport_close_before_peer_timeout_retries_once_and_ignores_ol ops.try_recv().is_err(), "closed backend needs no extra stop" ); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -487,6 +494,7 @@ async fn failure_cleanup_does_not_attribute_stop_to_the_user() { chat.realtime_conversation.phase, RealtimeConversationPhase::Inactive ); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( diff --git a/codex-rs/tui/src/chatwidget/realtime_tests/transcripts.rs b/codex-rs/tui/src/chatwidget/realtime_tests/transcripts.rs index a7c28b0f97..3df707403c 100644 --- a/codex-rs/tui/src/chatwidget/realtime_tests/transcripts.rs +++ b/codex-rs/tui/src/chatwidget/realtime_tests/transcripts.rs @@ -220,7 +220,9 @@ async fn voice_transcripts_stream_in_the_conversation_instead_of_the_footer() { assert!(events.try_recv().is_err()); chat.on_realtime_transcript_done("user".to_string(), "pick a number".to_string()); + commit_realtime_history_events(&mut chat, &mut events); assert!(chat.active_cell_transcript_lines(/*width*/ 80).is_none()); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("the completed voice transcript should become one normal user turn"); }; @@ -246,6 +248,7 @@ async fn unexpected_voice_close_preserves_partial_transcripts_once() { chat.on_realtime_transcript_done(role.to_string(), text.to_string()); let mut appearances = 0; + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { appearances += @@ -273,6 +276,7 @@ async fn interleaved_partial_transcripts_survive_voice_close() { chat.on_realtime_transcript_done("user".into(), "Correction".into()); chat.on_realtime_transcript_done("assistant".into(), "First answer".into()); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -317,6 +321,7 @@ async fn intentional_stop_preserves_both_interleaved_partials_once() { chat.stop_realtime_conversation(); chat.on_realtime_conversation_closed(/*reason*/ None); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -346,6 +351,7 @@ async fn stopping_voice_keeps_the_complete_late_final_caption() { chat.on_realtime_transcript_done("assistant".into(), "First and last".into()); chat.on_realtime_conversation_closed(Some("requested".into())); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -373,6 +379,7 @@ async fn separate_late_finals_with_a_shared_prefix_keep_both_full_captions() { chat.on_realtime_transcript_done("assistant".into(), "Hello".into()); chat.on_realtime_transcript_done("assistant".into(), "Hello again".into()); + commit_realtime_history_events(&mut chat, &mut events); let captions = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -397,6 +404,7 @@ async fn direct_reset_preserves_both_partial_speakers_for_replay() { chat.on_realtime_transcript_delta("assistant".into(), "answer".into()); chat.reset_realtime_conversation(); + commit_realtime_history_events(&mut chat, &mut events); let records = chat.take_realtime_transcript_cells_for_replay(); assert_eq!( @@ -409,6 +417,7 @@ async fn direct_reset_preserves_both_partial_speakers_for_replay() { ("assistant".to_string(), "First answer".to_string()) ] ); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -452,8 +461,10 @@ async fn stopping_voice_preserves_the_live_transcript_once() { chat.realtime_conversation.phase, RealtimeConversationPhase::Inactive ); + commit_realtime_history_events(&mut chat, &mut events); assert!(chat.active_cell_transcript_key().is_none()); let mut rendered = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { assert!(cell.transcript_animation_tick().is_none()); @@ -490,6 +501,7 @@ async fn transcript_completion_waits_for_normal_agent_stream_consolidation() { chat.stop_realtime_conversation(); assert_eq!(chat.realtime_conversation.pending_history_cells.len(), 2); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { assert!( @@ -505,6 +517,7 @@ async fn transcript_completion_waits_for_normal_agent_stream_consolidation() { chat.flush_realtime_transcript_history(); let mut rendered = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { rendered.extend( @@ -547,6 +560,7 @@ async fn transcript_handoff_moves_deferred_repeats_and_partial_once() { chat.note_stream_consolidation_completed(); chat.flush_realtime_transcript_history(); + commit_realtime_history_events(&mut chat, &mut events); let rendered = std::iter::from_fn(|| events.try_recv().ok()) .filter_map(|event| match event { AppEvent::InsertHistoryCell(cell) => Some( @@ -618,6 +632,7 @@ async fn transcripts_are_preserved_while_the_peer_is_connecting() { chat.on_realtime_transcript_done("assistant".to_string(), "Hello there".to_string()); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("voice should preserve transcripts received before the peer connects"); }; @@ -660,6 +675,7 @@ async fn completed_transcript_is_added_to_history_and_clears_the_caption() { chat.on_realtime_transcript_done("user".to_string(), "Can you hear me?".to_string()); let mut rendered = Vec::new(); + commit_realtime_history_events(&mut chat, &mut events); while let Ok(event) = events.try_recv() { if let AppEvent::InsertHistoryCell(cell) = event { rendered.extend( @@ -716,7 +732,9 @@ async fn live_voice_split_flap_animates_without_changing_final_history() { chat.on_realtime_transcript_done("assistant".to_string(), "gate 73".to_string()); + commit_realtime_history_events(&mut chat, &mut events); assert!(chat.active_cell_transcript_key().is_none()); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("a completed voice transcript should commit ordinary history"); }; @@ -762,6 +780,7 @@ async fn spoken_user_transcript_preserves_red_chevron_and_canonical_history() { .any(|line| line.to_string() == "› hello world") ); chat.on_realtime_transcript_done("user".to_string(), " hello world".to_string()); + commit_realtime_history_events(&mut chat, &mut events); let Ok(AppEvent::InsertHistoryCell(cell)) = events.try_recv() else { panic!("completed voice transcript should retain its ordinary history cell"); }; @@ -827,3 +846,125 @@ async fn new_voice_turn_interrupts_uncaptioned_queued_speech() { ); } } + +#[tokio::test] +async fn completed_user_caption_stays_visible_until_history_commit() { + let (mut chat, _sender, mut events, _ops) = make_chatwidget_manual_with_sender().await; + activate_voice(&mut chat); + chat.local_settings.tui.animations = false; + chat.on_realtime_transcript_delta("user".into(), "Keep these words visible".into()); + chat.on_realtime_transcript_done("user".into(), "Keep these words visible.".into()); + + // A scheduled draw must not clear the caption before its queued history event runs. + chat.pre_draw_tick(); + let visible = chat.active_cell_transcript_lines(/*width*/ 80).unwrap(); + insta::assert_snapshot!(visible.iter().map(ToString::to_string).collect::>().join("\n"), @r" + + › Keep these words visible. + "); + assert!(chat.active_cell_transcript_key().is_some()); + let viewport = render_bottom_popup(&chat, /*width*/ 80); + assert!(viewport.contains("Keep these words visible.")); + + // Streaming output can postpone insertion, but must not postpone display. + chat.on_agent_message_delta("An earlier answer".into()); + assert!(chat.take_realtime_transcript_history().is_empty()); + assert!(render_bottom_popup(&chat, /*width*/ 80).contains("Keep these words visible.")); + chat.finalize_completed_assistant_message(Some("An earlier answer")); + chat.note_stream_consolidation_completed(); + commit_realtime_history_events(&mut chat, &mut events); + let history = std::iter::from_fn(|| events.try_recv().ok()) + .filter_map(|event| { + if let AppEvent::InsertHistoryCell(cell) = event { + Some(cell) + } else { + None + } + }) + .flat_map(|cell| cell.display_lines(/*width*/ 80)) + .map(|line| line.to_string()) + .collect::>() + .join("\n"); + assert_eq!(history.matches("Keep these words visible.").count(), 1); + assert!(!render_bottom_popup(&chat, /*width*/ 80).contains("Keep these words visible.")); +} + +#[tokio::test] +async fn animated_interleaved_captions_keep_settled_words_visible() { + let words = "Keep these words visible"; + let mut settled = Vec::new(); + for (first, second) in [("user", "assistant"), ("assistant", "user")] { + let (mut chat, _sender, _events, _ops) = make_chatwidget_manual_with_sender().await; + activate_voice(&mut chat); + chat.local_settings.tui.animations = true; + chat.on_realtime_transcript_delta(first.into(), words.into()); + // Let the real animation settle before the other speaker's packets arrive. + tokio::time::sleep(std::time::Duration::from_millis(/*millis*/ 750)).await; + for (role, delta) in [(second, "Other speaker"), (first, " please"), (second, ".")] { + chat.on_realtime_transcript_delta(role.into(), delta.into()); + assert!(render_bottom_popup(&chat, /*width*/ 80).contains(words)); + let overlay = chat.active_cell_transcript_lines(/*width*/ 80).unwrap(); + assert!(overlay.iter().any(|line| line.to_string().contains(words))); + let key = chat.active_cell_transcript_key().unwrap(); + assert!(key.animation_tick.is_some()); + } + tokio::time::sleep(std::time::Duration::from_millis(/*millis*/ 750)).await; + let overlay = chat.active_cell_transcript_lines(/*width*/ 80).unwrap(); + settled.push(format!( + "{first} first:\n{}", + overlay + .iter() + .map(|line| line.to_string().trim_end().to_string()) + .collect::>() + .join("\n") + )); + chat.on_realtime_transcript_done(first.into(), "Keep these words visible please".into()); + assert!( + render_bottom_popup(&chat, /*width*/ 80).contains("Keep these words visible please") + ); + chat.on_realtime_transcript_done(second.into(), "Other speaker.".into()); + let live_count = chat.realtime_conversation.live_transcript_cells().count(); + assert_eq!(live_count, 0); + let history = chat.take_realtime_transcript_history(); + assert_eq!(history.len(), 2); + assert!(chat.active_cell_transcript_key().is_none()); + } + insta::assert_snapshot!(settled.join("\n"), @" + user first: + + › Keep these words visible please + + + • Other speaker. + assistant first: + + › Other speaker. + + + • Keep these words visible please + "); +} + +#[tokio::test] +async fn empty_interleaved_caption_completion_invalidates_overlay() { + for phase in [ + RealtimeConversationPhase::Active, + RealtimeConversationPhase::Stopping, + ] { + let (mut chat, _sender, _events, _ops) = make_chatwidget_manual_with_sender().await; + activate_voice(&mut chat); + chat.local_settings.tui.animations = false; + chat.on_realtime_transcript_delta("assistant".into(), "Discard this caption".into()); + chat.on_realtime_transcript_delta("user".into(), "Keep this caption".into()); + chat.realtime_conversation.phase = phase; + let previous_key = chat.active_cell_transcript_key().unwrap(); + chat.on_realtime_transcript_done("assistant".into(), String::new()); + let current_key = chat.active_cell_transcript_key().unwrap(); + assert_ne!(previous_key, current_key); + let visible = chat.active_cell_transcript_lines(/*width*/ 80).unwrap(); + assert_eq!( + visible.iter().map(ToString::to_string).collect::>(), + vec!["", "› Keep this caption", ""] + ); + } +} diff --git a/codex-rs/tui/src/chatwidget/rendering.rs b/codex-rs/tui/src/chatwidget/rendering.rs index d24f11411d..7752f51727 100644 --- a/codex-rs/tui/src/chatwidget/rendering.rs +++ b/codex-rs/tui/src/chatwidget/rendering.rs @@ -158,7 +158,12 @@ impl ChatWidget { }; let mut flex = FlexRenderable::new(); flex.push(/*flex*/ 1, active_cell_renderable); - if let Some(cell) = self.realtime_conversation.live_transcript_cell.as_ref() { + for cell in self + .realtime_conversation + .pending_history_cells + .iter() + .chain(self.realtime_conversation.live_transcript_cells()) + { flex.push( /*flex*/ 1, RenderableItem::Owned(Box::new(TranscriptAreaRenderable {