diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 3d5df268b2..97570b543a 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -1434,7 +1434,7 @@ Because audio is intentionally separate from `ThreadItem`, clients can opt out o The app-server streams JSON-RPC notifications while a turn is running. Each turn emits `turn/started` when it begins running and ends with `turn/completed` (final `turn` status). Token usage events stream separately via `thread/tokenUsage/updated`. Clients subscribe to the events they care about, rendering each item incrementally as updates arrive. The per-item lifecycle is always: `item/started` → zero or more item-specific deltas → `item/completed`. - `turn/started` — `{ turn }` with the turn id, empty `items`, and `status: "inProgress"`. -- `turn/completed` — `{ turn }` where `turn.status` is `completed`, `interrupted`, or `failed`; failures carry `{ error: { message, codexErrorInfo?, additionalDetails? } }`. +- `turn/completed` — `{ turn }` where `turn.status` is `completed`, `interrupted`, or `failed`; successful turns include their final agent message when available, and failures carry `{ error: { message, codexErrorInfo?, additionalDetails? } }`. - `turn/diff/updated` — `{ threadId, turnId, diff }` represents the up-to-date snapshot of the turn-level unified diff, emitted after every FileChange item. `diff` is the latest aggregated unified diff across every file change in the turn. UIs can render this to show the full "what changed" view without stitching individual `fileChange` items. - `turn/plan/updated` — `{ turnId, explanation?, plan }` whenever the agent shares or changes its plan; each `plan` entry is `{ step, status }` with `status` in `pending`, `inProgress`, or `completed`. - `rawResponse/completed` — internal-only; when `thread/start.experimentalRawEvents` is enabled, emits `{ threadId, turnId, responseId, usage }` once for each upstream Responses API completion. `usage` is the exact upstream usage payload mapped to the app-server token breakdown shape and is `null` when the upstream completion omitted usage. Unlike `thread/tokenUsage/updated`, this notification is not accumulated, estimated, persisted, or replayed. @@ -1443,7 +1443,7 @@ The app-server streams JSON-RPC notifications while a turn is running. Each turn - `model/verification` — `{ threadId, turnId, verifications }` when the backend flags additional account verification, such as `trustedAccessForCyber`. - `turn/moderationMetadata` — experimental; `{ threadId, turnId, metadata }` when a first-party backend supplies turn-scoped moderation metadata for client-side presentation. -Today both notifications carry an empty `items` array even when item events were streamed; rely on `item/*` notifications for the canonical item list until this is fixed. +`turn/started` carries no items. `turn/completed` carries only the final agent message as a summary fallback; continue consuming `item/*` notifications for the full canonical item list. #### Items diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index d9de3e5a40..3775994a69 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -1249,6 +1249,7 @@ async fn handle_turn_plan_update( struct TurnCompletionMetadata { status: TurnStatus, error: Option, + last_agent_message: Option, started_at: Option, completed_at: Option, duration_ms: Option, @@ -1260,12 +1261,16 @@ async fn emit_turn_completed_with_status( turn_completion_metadata: TurnCompletionMetadata, outgoing: &ThreadScopedOutgoingMessageSender, ) { + let (items, items_view) = match turn_completion_metadata.last_agent_message { + Some(item) => (vec![item], TurnItemsView::Summary), + None => (Vec::new(), TurnItemsView::NotLoaded), + }; let notification = TurnCompletedNotification { thread_id: conversation_id.to_string(), turn: Turn { id: event_turn_id, - items: vec![], - items_view: TurnItemsView::NotLoaded, + items, + items_view, error: turn_completion_metadata.error, status: turn_completion_metadata.status, started_at: turn_completion_metadata.started_at, @@ -1448,9 +1453,9 @@ async fn handle_turn_complete( ) { let turn_summary = find_and_remove_turn_summary(conversation_id, thread_state).await; - let (status, error) = match turn_summary.last_error { - Some(error) => (TurnStatus::Failed, Some(error)), - None => (TurnStatus::Completed, None), + let (status, error, last_agent_message) = match turn_summary.last_error { + Some(error) => (TurnStatus::Failed, Some(error), None), + None => (TurnStatus::Completed, None, turn_summary.last_agent_message), }; emit_turn_completed_with_status( @@ -1459,6 +1464,7 @@ async fn handle_turn_complete( TurnCompletionMetadata { status, error, + last_agent_message, started_at: turn_summary.started_at, completed_at: turn_complete_event.completed_at, duration_ms: turn_complete_event.duration_ms, @@ -1483,6 +1489,7 @@ async fn handle_turn_interrupted( TurnCompletionMetadata { status: TurnStatus::Interrupted, error: None, + last_agent_message: None, started_at: turn_summary.started_at, completed_at: turn_aborted_event.completed_at, duration_ms: turn_aborted_event.duration_ms, @@ -2090,6 +2097,8 @@ mod tests { use codex_app_server_protocol::TurnPlanStepStatus; use codex_login::CodexAuth; use codex_protocol::AgentPath; + use codex_protocol::items::AgentMessageContent as CoreAgentMessageContent; + use codex_protocol::items::AgentMessageItem as CoreAgentMessageItem; use codex_protocol::items::DynamicToolCallItem; use codex_protocol::items::DynamicToolCallStatus as CoreDynamicToolCallStatus; use codex_protocol::items::SubAgentActivityItem; @@ -3493,6 +3502,7 @@ mod tests { ThreadId::new(), ); let thread_state = new_thread_state(); + let event = turn_complete_event(&event_turn_id); { let mut state = thread_state.lock().await; state.track_current_turn_event( @@ -3507,14 +3517,48 @@ mod tests { ); state.track_current_turn_event( &event_turn_id, - &EventMsg::TurnComplete(turn_complete_event(&event_turn_id)), + &EventMsg::ItemCompleted(ItemCompletedEvent { + thread_id: conversation_id, + turn_id: event_turn_id.clone(), + item: CoreTurnItem::AgentMessage(CoreAgentMessageItem { + id: "msg-1".to_string(), + content: vec![ + CoreAgentMessageContent::Text { + text: "complete ".to_string(), + }, + CoreAgentMessageContent::Text { + text: "response".to_string(), + }, + ], + phase: None, + memory_citation: None, + }), + completed_at_ms: 0, + }), ); + state.track_current_turn_event( + &event_turn_id, + &EventMsg::ItemCompleted(ItemCompletedEvent { + thread_id: conversation_id, + turn_id: event_turn_id.clone(), + item: CoreTurnItem::AgentMessage(CoreAgentMessageItem { + id: "msg-2".to_string(), + content: vec![CoreAgentMessageContent::Text { + text: " ".to_string(), + }], + phase: None, + memory_citation: None, + }), + completed_at_ms: 0, + }), + ); + state.track_current_turn_event(&event_turn_id, &EventMsg::TurnComplete(event.clone())); } handle_turn_complete( conversation_id, event_turn_id.clone(), - turn_complete_event(&event_turn_id), + event, &outgoing, &thread_state, ) @@ -3525,8 +3569,12 @@ mod tests { ServerNotification::TurnCompleted(n) => { assert_eq!(n.turn.id, event_turn_id); assert_eq!(n.turn.status, TurnStatus::Completed); - assert_eq!(n.turn.items_view, TurnItemsView::NotLoaded); - assert!(n.turn.items.is_empty()); + assert_eq!(n.turn.items_view, TurnItemsView::Summary); + assert!(matches!( + &n.turn.items[..], + [ThreadItem::AgentMessage { id, text, .. }] + if id == "msg-1" && text == "complete response" + )); assert_eq!(n.turn.error, None); assert_eq!(n.turn.started_at, Some(42)); assert_eq!(n.turn.completed_at, Some(TEST_TURN_COMPLETED_AT)); @@ -3539,7 +3587,7 @@ mod tests { } #[tokio::test] - async fn test_handle_turn_interrupted_emits_interrupted_with_error() -> Result<()> { + async fn test_handle_turn_interrupted_emits_interrupted_without_error() -> Result<()> { let conversation_id = ThreadId::new(); let event_turn_id = "interrupt1".to_string(); let thread_state = new_thread_state(); diff --git a/codex-rs/app-server/src/thread_state.rs b/codex-rs/app-server/src/thread_state.rs index 3843172100..bb9b3ecabb 100644 --- a/codex-rs/app-server/src/thread_state.rs +++ b/codex-rs/app-server/src/thread_state.rs @@ -3,6 +3,7 @@ use crate::outgoing_message::ConnectionRequestId; use codex_app_server_protocol::RequestId; use codex_app_server_protocol::ThreadGoal; use codex_app_server_protocol::ThreadHistoryBuilder; +use codex_app_server_protocol::ThreadItem; use codex_app_server_protocol::ThreadSettings; use codex_app_server_protocol::Turn; use codex_app_server_protocol::TurnError; @@ -12,6 +13,9 @@ use codex_file_watcher::WatchRegistration; use codex_protocol::ThreadId; #[cfg(test)] use codex_protocol::config_types::MultiAgentMode; +use codex_protocol::items::AgentMessageContent as CoreAgentMessageContent; +use codex_protocol::items::TurnItem as CoreTurnItem; +use codex_protocol::models::MessagePhase; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::RolloutItem; use codex_rollout::state_db::StateDbHandle; @@ -77,6 +81,7 @@ pub(crate) struct TurnSummary { pub(crate) started_at: Option, pub(crate) command_execution_started: HashSet, pub(crate) last_error: Option, + pub(crate) last_agent_message: Option, } #[derive(Default)] @@ -150,12 +155,22 @@ impl ThreadState { if let EventMsg::TurnStarted(payload) = event { self.turn_summary.started_at = payload.started_at; } - self.current_turn_history.handle_event(event); - if matches!(event, EventMsg::TurnAborted(_) | EventMsg::TurnComplete(_)) - && !self.current_turn_history.has_active_turn() + if let EventMsg::ItemCompleted(payload) = event + && let CoreTurnItem::AgentMessage(item) = &payload.item + && matches!(item.phase, Some(MessagePhase::FinalAnswer) | None) + && item.content.iter().any(|content| { + matches!(content, CoreAgentMessageContent::Text { text } if !text.trim().is_empty()) + }) { + self.turn_summary.last_agent_message = + Some(ThreadItem::from(CoreTurnItem::AgentMessage(item.clone()))); + } + self.current_turn_history.handle_event(event); + if matches!(event, EventMsg::TurnAborted(_) | EventMsg::TurnComplete(_)) { self.last_terminal_turn_id = Some(event_turn_id.to_string()); - self.current_turn_history.reset(); + if !self.current_turn_history.has_active_turn() { + self.current_turn_history.reset(); + } } } diff --git a/codex-rs/app-server/tests/suite/v2/turn_start.rs b/codex-rs/app-server/tests/suite/v2/turn_start.rs index e508ae62fe..ae8cca7f34 100644 --- a/codex-rs/app-server/tests/suite/v2/turn_start.rs +++ b/codex-rs/app-server/tests/suite/v2/turn_start.rs @@ -294,6 +294,11 @@ async fn turn_start_with_empty_input_runs_model_request() -> Result<()> { assert_eq!(completed.thread_id, thread.id); assert_eq!(completed.turn.id, turn.id); assert_eq!(completed.turn.status, TurnStatus::Completed); + assert_eq!(completed.turn.items_view, TurnItemsView::Summary); + assert!(matches!( + &completed.turn.items[..], + [ThreadItem::AgentMessage { text, .. }] if text == "Done" + )); let requests = server .received_requests() @@ -1683,8 +1688,6 @@ async fn turn_start_emits_notifications_and_accepts_model_override() -> Result<( assert_eq!(completed.thread_id, thread.id); assert_eq!(completed.turn.id, turn.id); assert_eq!(completed.turn.status, TurnStatus::Completed); - assert_eq!(completed.turn.items_view, TurnItemsView::NotLoaded); - assert!(completed.turn.items.is_empty()); // Send a second turn that exercises the overrides path: change the model. let turn_req2 = mcp @@ -1735,8 +1738,6 @@ async fn turn_start_emits_notifications_and_accepts_model_override() -> Result<( assert_eq!(completed2.thread_id, thread.id); assert_eq!(completed2.turn.id, turn2.id); assert_eq!(completed2.turn.status, TurnStatus::Completed); - assert_eq!(completed2.turn.items_view, TurnItemsView::NotLoaded); - assert!(completed2.turn.items.is_empty()); Ok(()) } diff --git a/codex-rs/exec/src/lib.rs b/codex-rs/exec/src/lib.rs index 368bcf64a4..5e0a9ca9ea 100644 --- a/codex-rs/exec/src/lib.rs +++ b/codex-rs/exec/src/lib.rs @@ -1397,7 +1397,7 @@ fn should_backfill_turn_completed_items( return false; }; - !thread_ephemeral && payload.turn.items.is_empty() + !thread_ephemeral && payload.turn.items_view != codex_app_server_protocol::TurnItemsView::Full } fn turn_items_for_thread( diff --git a/codex-rs/exec/src/lib_tests.rs b/codex-rs/exec/src/lib_tests.rs index 568d6ad10e..f624752d0d 100644 --- a/codex-rs/exec/src/lib_tests.rs +++ b/codex-rs/exec/src/lib_tests.rs @@ -396,13 +396,13 @@ fn turn_items_for_thread_returns_matching_turn_items() { } #[test] -fn should_backfill_turn_completed_items_skips_ephemeral_threads() { +fn should_backfill_turn_completed_items_backfills_persisted_summaries_only() { let notification = ServerNotification::TurnCompleted(codex_app_server_protocol::TurnCompletedNotification { thread_id: "thread-1".to_string(), turn: codex_app_server_protocol::Turn { id: "turn-1".to_string(), - items_view: codex_app_server_protocol::TurnItemsView::Full, + items_view: codex_app_server_protocol::TurnItemsView::Summary, items: Vec::new(), status: codex_app_server_protocol::TurnStatus::Completed, error: None, @@ -416,6 +416,10 @@ fn should_backfill_turn_completed_items_skips_ephemeral_threads() { /*thread_ephemeral*/ true, ¬ification )); + assert!(should_backfill_turn_completed_items( + /*thread_ephemeral*/ false, + ¬ification + )); } #[test] diff --git a/codex-rs/tui/src/app/snapshots/codex_tui__app__tests__directive_only_completion_removes_streamed_directive.snap b/codex-rs/tui/src/app/snapshots/codex_tui__app__tests__directive_only_completion_removes_streamed_directive.snap new file mode 100644 index 0000000000..84b516998a --- /dev/null +++ b/codex-rs/tui/src/app/snapshots/codex_tui__app__tests__directive_only_completion_removes_streamed_directive.snap @@ -0,0 +1,5 @@ +--- +source: tui/src/app/tests.rs +expression: "rendered.lines.iter().map(rendered_line_text).collect::>().join(\"\\n\")" +--- +before directive diff --git a/codex-rs/tui/src/app/tests.rs b/codex-rs/tui/src/app/tests.rs index 26521c7b9e..5e6ec53818 100644 --- a/codex-rs/tui/src/app/tests.rs +++ b/codex-rs/tui/src/app/tests.rs @@ -5180,6 +5180,42 @@ async fn required_stream_reflow_during_capped_initial_replay_uses_transcript_tai Ok(()) } +#[tokio::test] +async fn directive_only_completion_removes_streamed_directive() -> Result<()> { + let (mut app, _rx, _op_rx) = make_test_app_with_channels().await; + app.config.terminal_resize_reflow.max_rows = TerminalResizeReflowMaxRows::Limit(20); + app.begin_initial_history_replay_buffer(); + app.transcript_cells = vec![ + plain_line_cell("before directive"), + Arc::new(AgentMessageCell::new( + vec![Line::from(r#"::git-stage{cwd="/tmp"}"#)], + /*is_first_line*/ true, + )), + ]; + + let mut tui = crate::tui::test_support::make_test_tui()?; + app.handle_consolidate_agent_message( + &mut tui, + String::new(), + PathBuf::from("/tmp"), + /*inline_visualization_context*/ None, + ConsolidationScrollbackReflow::Required, + /*deferred_history_cell*/ None, + )?; + + let rendered = app.render_transcript_lines_for_reflow(/*width*/ 80); + assert_snapshot!( + "directive_only_completion_removes_streamed_directive", + rendered + .lines + .iter() + .map(rendered_line_text) + .collect::>() + .join("\n") + ); + Ok(()) +} + #[tokio::test] async fn required_stream_reflow_during_capped_initial_replay_survives_transcript_overlay() -> Result<()> { diff --git a/codex-rs/tui/src/chatwidget/protocol.rs b/codex-rs/tui/src/chatwidget/protocol.rs index 9586f5fd7d..79e2c3fe6e 100644 --- a/codex-rs/tui/src/chatwidget/protocol.rs +++ b/codex-rs/tui/src/chatwidget/protocol.rs @@ -243,12 +243,43 @@ impl ChatWidget { self.last_rendered_user_message_display = None; match notification.turn.status { TurnStatus::Completed => { + let last_agent_message = + notification + .turn + .items + .iter() + .rev() + .find_map(|item| match item { + ThreadItem::AgentMessage { + id, + text, + phase: Some(MessagePhase::FinalAnswer) | None, + .. + } => Some((item.clone(), id.clone(), text.clone())), + _ => None, + }); + if let Some((item, id, _)) = &last_agent_message + && self + .transcript + .last_completed_agent_message + .as_ref() + .is_none_or(|(turn_id, item_id)| { + turn_id != ¬ification.turn.id || item_id != id + }) + { + self.handle_thread_item( + item.clone(), + notification.turn.id.clone(), + replay_kind + .map_or(ThreadItemRenderSource::Live, ThreadItemRenderSource::Replay), + ); + } self.last_non_retry_error = None; self.on_task_complete( - /*last_agent_message*/ None, + last_agent_message.map(|(_, _, text)| text), notification.turn.duration_ms, replay_kind.is_some(), - ) + ); } TurnStatus::Interrupted => { self.last_non_retry_error = None; diff --git a/codex-rs/tui/src/chatwidget/replay.rs b/codex-rs/tui/src/chatwidget/replay.rs index 2a16254647..fad6eb8473 100644 --- a/codex-rs/tui/src/chatwidget/replay.rs +++ b/codex-rs/tui/src/chatwidget/replay.rs @@ -118,6 +118,7 @@ impl ChatWidget { } }), }, + &turn_id, from_replay, ); } diff --git a/codex-rs/tui/src/chatwidget/snapshots/codex_tui__chatwidget__tests__live_app_server_turn_completion_repairs_dropped_message_deltas.snap b/codex-rs/tui/src/chatwidget/snapshots/codex_tui__chatwidget__tests__live_app_server_turn_completion_repairs_dropped_message_deltas.snap new file mode 100644 index 0000000000..eab6e17635 --- /dev/null +++ b/codex-rs/tui/src/chatwidget/snapshots/codex_tui__chatwidget__tests__live_app_server_turn_completion_repairs_dropped_message_deltas.snap @@ -0,0 +1,9 @@ +--- +source: tui/src/chatwidget/tests/app_server.rs +expression: source +--- +The transport kept this. +And dropped this. + +- Finding — /tmp/file.rs:1 + Keep ::git-stage{cwd=/tmp} literal. diff --git a/codex-rs/tui/src/chatwidget/streaming.rs b/codex-rs/tui/src/chatwidget/streaming.rs index 249b8c832b..dcccbd027f 100644 --- a/codex-rs/tui/src/chatwidget/streaming.rs +++ b/codex-rs/tui/src/chatwidget/streaming.rs @@ -20,15 +20,27 @@ impl ChatWidget { } pub(super) fn flush_answer_stream_with_separator(&mut self) { + self.flush_answer_stream(/*completed_message*/ None); + } + + fn flush_answer_stream(&mut self, completed_message: Option<&str>) { let had_stream_controller = self.stream_controller.is_some(); if let Some(mut controller) = self.stream_controller.take() { - let scrollback_reflow = if controller.has_live_tail() { + let had_live_tail = controller.has_live_tail(); + self.clear_active_stream_tail(); + let (cell, streamed_source) = controller.finalize(); + let completed_message_differs = completed_message.is_some_and(|completed| { + let Some(streamed) = streamed_source.as_deref() else { + return true; + }; + // Stream finalization supplies one trailing newline when the last delta omitted it. + streamed != completed && streamed.strip_suffix('\n') != Some(completed) + }); + let scrollback_reflow = if had_live_tail || completed_message_differs { crate::app_event::ConsolidationScrollbackReflow::Required } else { crate::app_event::ConsolidationScrollbackReflow::IfResizeReflowRan }; - self.clear_active_stream_tail(); - let (cell, source) = controller.finalize(); // Match newline-committed streaming behavior: once assistant output is ready to be // committed into history, hide the inline status row so transcript content replaces it. if cell.is_some() { @@ -45,9 +57,12 @@ impl ChatWidget { }; // Consolidate the run of streaming AgentMessageCells into a single AgentMarkdownCell // that can re-render from source on resize. + let source = completed_message.map(str::to_owned).or_else(|| { + streamed_source.map(|source| { + parse_assistant_markdown(&source, self.config.cwd.as_path()).visible_markdown + }) + }); if let Some(source) = source { - let source = - parse_assistant_markdown(&source, self.config.cwd.as_path()).visible_markdown; let inline_visualization_context = self.thread_id.and_then(|thread_id| { crate::inline_visualization::InlineVisualizationContext::from_config( &self.config, @@ -110,15 +125,15 @@ impl ChatWidget { } pub(super) fn finalize_completed_assistant_message(&mut self, message: Option<&str>) { - // If we have a stream_controller, the finalized message payload is redundant because the - // visible content has already been accumulated through deltas. if self.stream_controller.is_none() && let Some(message) = message && !message.is_empty() { self.handle_streaming_delta(message.to_string()); } - self.flush_answer_stream_with_separator(); + // Item completion is authoritative. Use it for consolidation so any + // deltas dropped by a saturated transport cannot truncate the transcript. + self.flush_answer_stream(message); self.handle_stream_finished(); self.request_redraw(); } @@ -301,8 +316,10 @@ impl ChatWidget { pub(super) fn on_agent_message_item_completed( &mut self, item: AgentMessageItem, + turn_id: &str, from_replay: bool, ) { + self.transcript.last_completed_agent_message = Some((turn_id.to_string(), item.id.clone())); let mut message = String::new(); for content in &item.content { match content { @@ -310,9 +327,7 @@ impl ChatWidget { } } let parsed = parse_assistant_markdown(&message, self.config.cwd.as_path()); - self.finalize_completed_assistant_message( - (!parsed.visible_markdown.is_empty()).then_some(parsed.visible_markdown.as_str()), - ); + self.finalize_completed_assistant_message(Some(parsed.visible_markdown.as_str())); if matches!(item.phase, Some(MessagePhase::FinalAnswer) | None) && !parsed.visible_markdown.is_empty() { diff --git a/codex-rs/tui/src/chatwidget/tests/app_server.rs b/codex-rs/tui/src/chatwidget/tests/app_server.rs index a0728e017f..859b40f8f3 100644 --- a/codex-rs/tui/src/chatwidget/tests/app_server.rs +++ b/codex-rs/tui/src/chatwidget/tests/app_server.rs @@ -603,17 +603,18 @@ async fn live_app_server_turn_completed_clears_working_status_after_answer_item( .expect("status indicator should be visible"); assert_eq!(status.header(), "Working"); + let item = AppServerThreadItem::AgentMessage { + id: "msg-1".to_string(), + text: "Yes. What do you need?".to_string(), + phase: Some(MessagePhase::FinalAnswer), + memory_citation: None, + }; chat.handle_server_notification( ServerNotification::ItemCompleted(ItemCompletedNotification { thread_id: "thread-1".to_string(), turn_id: "turn-1".to_string(), completed_at_ms: 0, - item: AppServerThreadItem::AgentMessage { - id: "msg-1".to_string(), - text: "Yes. What do you need?".to_string(), - phase: Some(MessagePhase::FinalAnswer), - memory_citation: None, - }, + item: item.clone(), }), /*replay_kind*/ None, ); @@ -628,8 +629,8 @@ async fn live_app_server_turn_completed_clears_working_status_after_answer_item( thread_id: "thread-1".to_string(), turn: AppServerTurn { id: "turn-1".to_string(), - items_view: codex_app_server_protocol::TurnItemsView::Full, - items: Vec::new(), + items_view: codex_app_server_protocol::TurnItemsView::Summary, + items: vec![item], status: AppServerTurnStatus::Completed, error: None, started_at: None, @@ -640,8 +641,13 @@ async fn live_app_server_turn_completed_clears_working_status_after_answer_item( /*replay_kind*/ None, ); + assert!(drain_insert_history(&mut rx).is_empty()); assert!(!chat.bottom_pane.is_task_running()); assert!(chat.bottom_pane.status_widget().is_none()); + assert_eq!( + chat.transcript.last_completed_agent_message, + Some(("turn-1".to_string(), "msg-1".to_string())) + ); } #[tokio::test] @@ -1089,7 +1095,6 @@ async fn live_app_server_failed_turn_does_not_duplicate_error_history() { }), /*replay_kind*/ None, ); - let first_cells = drain_insert_history(&mut rx); assert_eq!(first_cells.len(), 1); assert!(lines_to_single_string(&first_cells[0]).contains("permission denied")); @@ -1153,6 +1158,64 @@ async fn live_app_server_failed_turn_consolidates_streamed_answer() { ); } +#[tokio::test] +async fn live_app_server_turn_completion_repairs_dropped_message_deltas() { + let (mut chat, mut rx, _op_rx) = make_chatwidget_manual(/*model_override*/ None).await; + + handle_turn_started(&mut chat, "turn-1"); + while rx.try_recv().is_ok() {} + handle_agent_message_delta(&mut chat, "The transport kept this.\n"); + chat.run_commit_tick(); + while rx.try_recv().is_ok() {} + + let mut completed_turn = app_server_turn( + "turn-1", + AppServerTurnStatus::Completed, + Some(1_000), + /*error*/ None, + ); + completed_turn.items_view = codex_app_server_protocol::TurnItemsView::Summary; + completed_turn.items = vec![AppServerThreadItem::AgentMessage { + id: "msg-1".to_string(), + text: concat!( + "The transport kept this.\nAnd dropped this.\n\n", + r#"::code-comment{title="Finding" body="Keep ::git-stage{cwd=/tmp} literal." file="/tmp/file.rs"}"#, + ) + .to_string(), + phase: Some(MessagePhase::FinalAnswer), + memory_citation: None, + }]; + chat.handle_server_notification( + ServerNotification::TurnCompleted(TurnCompletedNotification { + thread_id: "thread-1".to_string(), + turn: completed_turn, + }), + /*replay_kind*/ None, + ); + + let consolidations = std::iter::from_fn(|| rx.try_recv().ok()) + .filter_map(|event| match event { + AppEvent::ConsolidateAgentMessage { + source, + scrollback_reflow, + .. + } => { + assert_eq!( + scrollback_reflow, + crate::app_event::ConsolidationScrollbackReflow::Required + ); + assert_chatwidget_snapshot!( + "live_app_server_turn_completion_repairs_dropped_message_deltas", + source, + ); + Some(()) + } + _ => None, + }) + .collect::>(); + assert_eq!(consolidations.len(), 1); +} + #[tokio::test] async fn live_app_server_stream_recovery_restores_previous_status_header() { let (mut chat, mut rx, _op_rx) = make_chatwidget_manual(/*model_override*/ None).await; diff --git a/codex-rs/tui/src/chatwidget/transcript.rs b/codex-rs/tui/src/chatwidget/transcript.rs index ed66d8a717..bacfa2b992 100644 --- a/codex-rs/tui/src/chatwidget/transcript.rs +++ b/codex-rs/tui/src/chatwidget/transcript.rs @@ -9,6 +9,7 @@ pub(super) struct TranscriptState { pub(super) active_cell_revision: u64, /// Raw markdown of the most recently completed agent response. pub(super) last_agent_markdown: Option, + pub(super) last_completed_agent_message: Option<(String, String)>, /// Raw markdown of the most recently completed proposed plan. pub(super) latest_proposed_plan_markdown: Option, /// Whether this turn already produced a copyable response. @@ -56,6 +57,7 @@ impl TranscriptState { pub(super) fn reset_turn_flags(&mut self) { self.saw_copy_source_this_turn = false; + self.last_completed_agent_message = None; self.saw_plan_update_this_turn = false; self.saw_plan_item_this_turn = false; self.had_work_activity = false; diff --git a/codex-rs/tui/src/chatwidget/turn_runtime.rs b/codex-rs/tui/src/chatwidget/turn_runtime.rs index 0c5152cc6d..f479891608 100644 --- a/codex-rs/tui/src/chatwidget/turn_runtime.rs +++ b/codex-rs/tui/src/chatwidget/turn_runtime.rs @@ -107,20 +107,9 @@ impl ChatWidget { from_replay: bool, ) { self.input_queue.submit_pending_steers_after_interrupt = false; - // Use `last_agent_message` from the turn-complete notification as the copy - // source only when no earlier item-level event (AgentMessageItem, plan - // commit, review output) already recorded markdown for this turn. This - // prevents the final summary from overwriting a more specific source. let sanitized_last_agent_message = last_agent_message.as_deref().map(|message| { parse_assistant_markdown(message, self.config.cwd.as_path()).visible_markdown }); - if let Some(message) = sanitized_last_agent_message - .as_ref() - .filter(|message| !message.is_empty()) - && !self.transcript.saw_copy_source_this_turn - { - self.record_agent_markdown(message); - } // For desktop notifications: prefer the notification payload, fall back to // the item-level copy source if present, otherwise send an empty string. let notification_response = sanitized_last_agent_message