tui: align pending steers with core acceptance (#12868)

## Summary
- submit `Enter` steers immediately while a turn is already running
instead of routing them through `queued_user_messages`
- keep those submitted steers visible in the footer as `pending_steers`
until core records them as a user message or aborts the turn
- reconcile pending steers on `ItemCompleted(UserMessage)`, not
`RawResponseItem`
- emit user-message item lifecycle for leftover pending input at task
finish, then remove the TUI `TurnComplete` fallback
- keep `queued_user_messages` for actual queued drafts, rendered below
pending steers

## Problem
While the assistant was generating, pressing `Enter` could send the
input into `queued_user_messages`. That queue only drains after the turn
ends, so ordinary steers behaved like queued drafts instead of landing
at the next core sampling boundary.

The first version of this fix also used `RawResponseItem` to decide when
a steer had landed. Review feedback was that this is the wrong
abstraction for client behavior.

There was also a late edge case in core: if pending steer input was
accepted after the final sampling decision but before `TurnComplete`,
core would record that user message into history at task finish without
emitting `ItemStarted(UserMessage)` / `ItemCompleted(UserMessage)`. TUI
had a fallback to paper over that gap locally.

## Approach
- `Enter` during an active turn now submits a normal `Op::UserTurn`
immediately
- TUI keeps a local pending-steer preview instead of rendering that user
message into history immediately
- when core records the steer as `ItemCompleted(UserMessage)`, TUI
matches and removes the corresponding pending preview, then renders the
committed user message
- core now emits the same user-message lifecycle when
`on_task_finished(...)` drains leftover pending user input, before
`TurnComplete`
- with that lifecycle gap closed in core, TUI no longer needs to flush
pending steers into history on `TurnComplete`
- if the turn is interrupted, pending steers and queued drafts are both
restored into the composer, with pending steers first

## Notes
- `Tab` still uses the real queued-message path
- `queued_user_messages` and `pending_steers` are separate state with
separate semantics
- the pending-steer matching key is built directly from `UserInput`
- this removes the new TUI dependency on `RawResponseItem`

## Validation
- `just fmt`
- `cargo test -p codex-core
task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input --
--nocapture`
- `cargo test -p codex-tui`
This commit is contained in:
Charley Cunningham
2026-03-03 15:31:52 -08:00
committed by GitHub
Unverified
parent 24a2d0c696
commit 299b8ac445
15 changed files with 1247 additions and 264 deletions
+70 -2
View File
@@ -6590,6 +6590,7 @@ mod tests {
use crate::protocol::TokenCountEvent;
use crate::protocol::TokenUsage;
use crate::protocol::TokenUsageInfo;
use crate::protocol::TurnCompleteEvent;
use crate::protocol::UserMessageEvent;
use crate::rollout::policy::EventPersistenceMode;
use crate::rollout::recorder::RolloutRecorder;
@@ -9307,8 +9308,8 @@ mod tests {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn task_finish_persists_leftover_pending_input() {
let (sess, tc, _rx) = make_session_and_context_with_rx().await;
async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() {
let (sess, tc, rx) = make_session_and_context_with_rx().await;
let input = vec![UserInput::Text {
text: "hello".to_string(),
text_elements: Vec::new(),
@@ -9323,6 +9324,8 @@ mod tests {
)
.await;
while rx.try_recv().is_ok() {}
sess.inject_response_items(vec![ResponseInputItem::Message {
role: "user".to_string(),
content: vec![ContentItem::InputText {
@@ -9348,6 +9351,71 @@ mod tests {
history.raw_items().iter().any(|item| item == &expected),
"expected pending input to be persisted into history on turn completion"
);
let first = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("expected raw response item event")
.expect("channel open");
assert!(matches!(first.msg, EventMsg::RawResponseItem(_)));
let second = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("expected item started event")
.expect("channel open");
assert!(matches!(
second.msg,
EventMsg::ItemStarted(ItemStartedEvent {
item: TurnItem::UserMessage(UserMessageItem { content, .. }),
..
}) if content == vec![UserInput::Text {
text: "late pending input".to_string(),
text_elements: Vec::new(),
}]
));
let third = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("expected item completed event")
.expect("channel open");
assert!(matches!(
third.msg,
EventMsg::ItemCompleted(ItemCompletedEvent {
item: TurnItem::UserMessage(UserMessageItem { content, .. }),
..
}) if content == vec![UserInput::Text {
text: "late pending input".to_string(),
text_elements: Vec::new(),
}]
));
let fourth = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("expected legacy user message event")
.expect("channel open");
assert!(matches!(
fourth.msg,
EventMsg::UserMessage(UserMessageEvent {
message,
images,
text_elements,
local_images,
}) if message == "late pending input"
&& images == Some(Vec::new())
&& text_elements.is_empty()
&& local_images.is_empty()
));
let fifth = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("expected turn complete event")
.expect("channel open");
assert!(matches!(
fifth.msg,
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id,
last_agent_message: None,
}) if turn_id == tc.sub_id
));
}
#[tokio::test]
+21 -2
View File
@@ -22,6 +22,7 @@ use crate::AuthManager;
use crate::codex::Session;
use crate::codex::TurnContext;
use crate::contextual_user_message::TURN_ABORTED_OPEN_TAG;
use crate::event_mapping::parse_turn_item;
use crate::models_manager::manager::ModelsManager;
use crate::protocol::EventMsg;
use crate::protocol::TurnAbortReason;
@@ -30,6 +31,7 @@ use crate::protocol::TurnCompleteEvent;
use crate::state::ActiveTurn;
use crate::state::RunningTask;
use crate::state::TaskKind;
use codex_protocol::items::TurnItem;
use codex_protocol::models::ContentItem;
use codex_protocol::models::ResponseInputItem;
use codex_protocol::models::ResponseItem;
@@ -213,8 +215,25 @@ impl Session {
.into_iter()
.map(ResponseItem::from)
.collect::<Vec<_>>();
self.record_conversation_items(turn_context.as_ref(), &pending_response_items)
.await;
for response_item in pending_response_items {
if let Some(TurnItem::UserMessage(user_message)) = parse_turn_item(&response_item) {
// Keep leftover user input on the same persistence + lifecycle path as the
// normal pre-sampling drain. This helper records the response item once, then
// emits ItemStarted/UserMessage and ItemCompleted/UserMessage for clients.
self.record_user_prompt_and_emit_turn_item(
turn_context.as_ref(),
&user_message.content,
response_item,
)
.await;
} else {
self.record_conversation_items(
turn_context.as_ref(),
std::slice::from_ref(&response_item),
)
.await;
}
}
}
let event = EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_context.sub_id.clone(),
+22 -16
View File
@@ -17,8 +17,8 @@ use std::path::PathBuf;
use crate::app_event::ConnectorsSnapshot;
use crate::app_event_sender::AppEventSender;
use crate::bottom_pane::pending_input_preview::PendingInputPreview;
use crate::bottom_pane::pending_thread_approvals::PendingThreadApprovals;
use crate::bottom_pane::queued_user_messages::QueuedUserMessages;
use crate::bottom_pane::unified_exec_footer::UnifiedExecFooter;
use crate::key_hint;
use crate::key_hint::KeyBinding;
@@ -93,9 +93,9 @@ pub(crate) use skills_toggle_view::SkillsToggleView;
pub(crate) use status_line_setup::StatusLineItem;
pub(crate) use status_line_setup::StatusLineSetupView;
mod paste_burst;
mod pending_input_preview;
mod pending_thread_approvals;
pub mod popup_consts;
mod queued_user_messages;
mod scroll_state;
mod selection_popup_common;
mod textarea;
@@ -172,8 +172,8 @@ pub(crate) struct BottomPane {
/// When a status row exists, this summary is mirrored inline in that row;
/// when no status row exists, it renders as its own footer row.
unified_exec_footer: UnifiedExecFooter,
/// Queued user messages to show above the composer while a turn is running.
queued_user_messages: QueuedUserMessages,
/// Preview of pending steers and queued drafts shown above the composer.
pending_input_preview: PendingInputPreview,
/// Inactive threads with pending approval requests.
pending_thread_approvals: PendingThreadApprovals,
context_window_percent: Option<i64>,
@@ -223,7 +223,7 @@ impl BottomPane {
is_task_running: false,
status: None,
unified_exec_footer: UnifiedExecFooter::new(),
queued_user_messages: QueuedUserMessages::new(),
pending_input_preview: PendingInputPreview::new(),
pending_thread_approvals: PendingThreadApprovals::new(),
esc_backtrack_hint: false,
animations_enabled,
@@ -317,7 +317,7 @@ impl BottomPane {
/// Update the key hint shown next to queued messages so it matches the
/// binding that `ChatWidget` actually listens for.
pub(crate) fn set_queued_message_edit_binding(&mut self, binding: KeyBinding) {
self.queued_user_messages.set_edit_binding(binding);
self.pending_input_preview.set_edit_binding(binding);
self.request_redraw();
}
@@ -774,9 +774,14 @@ impl BottomPane {
true
}
/// Update the queued messages preview shown above the composer.
pub(crate) fn set_queued_user_messages(&mut self, queued: Vec<String>) {
self.queued_user_messages.messages = queued;
/// Update the pending-input preview shown above the composer.
pub(crate) fn set_pending_input_preview(
&mut self,
queued: Vec<String>,
pending_steers: Vec<String>,
) {
self.pending_input_preview.pending_steers = pending_steers;
self.pending_input_preview.queued_messages = queued;
self.request_redraw();
}
@@ -1019,18 +1024,19 @@ impl BottomPane {
flex.push(0, RenderableItem::Borrowed(&self.unified_exec_footer));
}
let has_pending_thread_approvals = !self.pending_thread_approvals.is_empty();
let has_queued_messages = !self.queued_user_messages.messages.is_empty();
let has_pending_input = !self.pending_input_preview.queued_messages.is_empty()
|| !self.pending_input_preview.pending_steers.is_empty();
let has_status_or_footer =
self.status.is_some() || !self.unified_exec_footer.is_empty();
let has_inline_previews = has_pending_thread_approvals || has_queued_messages;
let has_inline_previews = has_pending_thread_approvals || has_pending_input;
if has_inline_previews && has_status_or_footer {
flex.push(0, RenderableItem::Owned("".into()));
}
flex.push(1, RenderableItem::Borrowed(&self.pending_thread_approvals));
if has_pending_thread_approvals && has_queued_messages {
if has_pending_thread_approvals && has_pending_input {
flex.push(0, RenderableItem::Owned("".into()));
}
flex.push(1, RenderableItem::Borrowed(&self.queued_user_messages));
flex.push(1, RenderableItem::Borrowed(&self.pending_input_preview));
if !has_inline_previews && has_status_or_footer {
flex.push(0, RenderableItem::Owned("".into()));
}
@@ -1406,7 +1412,7 @@ mod tests {
StatusDetailsCapitalization::CapitalizeFirst,
STATUS_DETAILS_DEFAULT_MAX_LINES,
);
pane.set_queued_user_messages(vec!["Queued follow-up question".to_string()]);
pane.set_pending_input_preview(vec!["Queued follow-up question".to_string()], Vec::new());
let width = 48;
let height = pane.desired_height(width);
@@ -1433,7 +1439,7 @@ mod tests {
});
pane.set_task_running(true);
pane.set_queued_user_messages(vec!["Queued follow-up question".to_string()]);
pane.set_pending_input_preview(vec!["Queued follow-up question".to_string()], Vec::new());
pane.hide_status_indicator();
let width = 48;
@@ -1461,7 +1467,7 @@ mod tests {
});
pane.set_task_running(true);
pane.set_queued_user_messages(vec!["Queued follow-up question".to_string()]);
pane.set_pending_input_preview(vec!["Queued follow-up question".to_string()], Vec::new());
let width = 48;
let height = pane.desired_height(width);
@@ -10,23 +10,26 @@ use crate::render::renderable::Renderable;
use crate::wrapping::RtOptions;
use crate::wrapping::adaptive_wrap_lines;
/// Widget that displays a list of user messages queued while a turn is in progress.
/// Widget that displays pending steers plus user messages queued while a turn is in progress.
///
/// The widget shows a key hint at the bottom (e.g. "⌥ + ↑ edit") telling the
/// user how to pop the most recent queued message back into the composer.
/// Because some terminals intercept certain modifier-key combinations, the
/// displayed binding is configurable via [`set_edit_binding`](Self::set_edit_binding).
pub(crate) struct QueuedUserMessages {
pub messages: Vec<String>,
/// The widget shows pending steers first, then queued user messages. It only
/// shows the edit hint at the bottom (e.g. "⌥ + ↑ edit") when there are actual
/// queued user messages to pop back into the composer. Because some terminals
/// intercept certain modifier-key combinations, the displayed binding is
/// configurable via [`set_edit_binding`](Self::set_edit_binding).
pub(crate) struct PendingInputPreview {
pub pending_steers: Vec<String>,
pub queued_messages: Vec<String>,
/// Key combination rendered in the hint line. Defaults to Alt+Up but may
/// be overridden for terminals where that chord is unavailable.
edit_binding: key_hint::KeyBinding,
}
impl QueuedUserMessages {
impl PendingInputPreview {
pub(crate) fn new() -> Self {
Self {
messages: Vec::new(),
pending_steers: Vec::new(),
queued_messages: Vec::new(),
edit_binding: key_hint::alt(KeyCode::Up),
}
}
@@ -39,13 +42,31 @@ impl QueuedUserMessages {
}
fn as_renderable(&self, width: u16) -> Box<dyn Renderable> {
if self.messages.is_empty() || width < 4 {
if (self.pending_steers.is_empty() && self.queued_messages.is_empty()) || width < 4 {
return Box::new(());
}
let mut lines = vec![];
for message in &self.messages {
for steer in &self.pending_steers {
let wrapped = adaptive_wrap_lines(
steer
.lines()
.map(|line| format!("pending steer: {line}").dim()),
RtOptions::new(width as usize)
.initial_indent(Line::from(" ! ".dim()))
.subsequent_indent(Line::from(" ")),
);
let len = wrapped.len();
for line in wrapped.into_iter().take(3) {
lines.push(line);
}
if len > 3 {
lines.push(Line::from("".dim()));
}
}
for message in &self.queued_messages {
let wrapped = adaptive_wrap_lines(
message.lines().map(|line| line.dim().italic()),
RtOptions::new(width as usize)
@@ -61,20 +82,22 @@ impl QueuedUserMessages {
}
}
lines.push(
Line::from(vec![
" ".into(),
self.edit_binding.into(),
" edit".into(),
])
.dim(),
);
if !self.queued_messages.is_empty() {
lines.push(
Line::from(vec![
" ".into(),
self.edit_binding.into(),
" edit".into(),
])
.dim(),
);
}
Paragraph::new(lines).into()
}
}
impl Renderable for QueuedUserMessages {
impl Renderable for PendingInputPreview {
fn render(&self, area: Rect, buf: &mut Buffer) {
if area.is_empty() {
return;
@@ -96,21 +119,21 @@ mod tests {
#[test]
fn desired_height_empty() {
let queue = QueuedUserMessages::new();
let queue = PendingInputPreview::new();
assert_eq!(queue.desired_height(40), 0);
}
#[test]
fn desired_height_one_message() {
let mut queue = QueuedUserMessages::new();
queue.messages.push("Hello, world!".to_string());
let mut queue = PendingInputPreview::new();
queue.queued_messages.push("Hello, world!".to_string());
assert_eq!(queue.desired_height(40), 2);
}
#[test]
fn render_one_message() {
let mut queue = QueuedUserMessages::new();
queue.messages.push("Hello, world!".to_string());
let mut queue = PendingInputPreview::new();
queue.queued_messages.push("Hello, world!".to_string());
let width = 40;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
@@ -120,9 +143,11 @@ mod tests {
#[test]
fn render_two_messages() {
let mut queue = QueuedUserMessages::new();
queue.messages.push("Hello, world!".to_string());
queue.messages.push("This is another message".to_string());
let mut queue = PendingInputPreview::new();
queue.queued_messages.push("Hello, world!".to_string());
queue
.queued_messages
.push("This is another message".to_string());
let width = 40;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
@@ -132,11 +157,17 @@ mod tests {
#[test]
fn render_more_than_three_messages() {
let mut queue = QueuedUserMessages::new();
queue.messages.push("Hello, world!".to_string());
queue.messages.push("This is another message".to_string());
queue.messages.push("This is a third message".to_string());
queue.messages.push("This is a fourth message".to_string());
let mut queue = PendingInputPreview::new();
queue.queued_messages.push("Hello, world!".to_string());
queue
.queued_messages
.push("This is another message".to_string());
queue
.queued_messages
.push("This is a third message".to_string());
queue
.queued_messages
.push("This is a fourth message".to_string());
let width = 40;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
@@ -146,11 +177,13 @@ mod tests {
#[test]
fn render_wrapped_message() {
let mut queue = QueuedUserMessages::new();
let mut queue = PendingInputPreview::new();
queue
.messages
.queued_messages
.push("This is a longer message that should be wrapped".to_string());
queue.messages.push("This is another message".to_string());
queue
.queued_messages
.push("This is another message".to_string());
let width = 40;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
@@ -160,9 +193,9 @@ mod tests {
#[test]
fn render_many_line_message() {
let mut queue = QueuedUserMessages::new();
let mut queue = PendingInputPreview::new();
queue
.messages
.queued_messages
.push("This is\na message\nwith many\nlines".to_string());
let width = 40;
let height = queue.desired_height(width);
@@ -173,8 +206,8 @@ mod tests {
#[test]
fn long_url_like_message_does_not_expand_into_wrapped_ellipsis_rows() {
let mut queue = QueuedUserMessages::new();
queue.messages.push(
let mut queue = PendingInputPreview::new();
queue.queued_messages.push(
"example.test/api/v1/projects/alpha-team/releases/2026-02-17/builds/1234567890/artifacts/reports/performance/summary/detail/session_id=abc123def456ghi789"
.to_string(),
);
@@ -202,4 +235,35 @@ mod tests {
"expected no wrapped-ellipsis row for URL-like token, got rows: {rendered_rows:?}"
);
}
#[test]
fn render_one_pending_steer() {
let mut queue = PendingInputPreview::new();
queue.pending_steers.push("Please continue.".to_string());
let width = 48;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
queue.render(Rect::new(0, 0, width, height), &mut buf);
assert_snapshot!("render_one_pending_steer", format!("{buf:?}"));
}
#[test]
fn render_pending_steers_above_queued_messages() {
let mut queue = PendingInputPreview::new();
queue.pending_steers.push("Please continue.".to_string());
queue
.pending_steers
.push("Check the last command output.".to_string());
queue
.queued_messages
.push("Queued follow-up question".to_string());
let width = 52;
let height = queue.desired_height(width);
let mut buf = Buffer::empty(Rect::new(0, 0, width, height));
queue.render(Rect::new(0, 0, width, height), &mut buf);
assert_snapshot!(
"render_pending_steers_above_queued_messages",
format!("{buf:?}")
);
}
}
@@ -0,0 +1,15 @@
---
source: tui/src/bottom_pane/pending_input_preview.rs
assertion_line: 237
expression: "format!(\"{buf:?}\")"
---
Buffer {
area: Rect { x: 0, y: 0, width: 48, height: 1 },
content: [
" ! pending steer: Please continue. ",
],
styles: [
x: 0, y: 0, fg: Reset, bg: Reset, underline: Reset, modifier: DIM,
x: 35, y: 0, fg: Reset, bg: Reset, underline: Reset, modifier: NONE,
]
}
@@ -0,0 +1,25 @@
---
source: tui/src/bottom_pane/pending_input_preview.rs
assertion_line: 252
expression: "format!(\"{buf:?}\")"
---
Buffer {
area: Rect { x: 0, y: 0, width: 52, height: 4 },
content: [
" ! pending steer: Please continue. ",
" ! pending steer: Check the last command output. ",
" ↳ Queued follow-up question ",
" ⌥ + ↑ edit ",
],
styles: [
x: 0, y: 0, fg: Reset, bg: Reset, underline: Reset, modifier: DIM,
x: 35, y: 0, fg: Reset, bg: Reset, underline: Reset, modifier: NONE,
x: 0, y: 1, fg: Reset, bg: Reset, underline: Reset, modifier: DIM,
x: 49, y: 1, fg: Reset, bg: Reset, underline: Reset, modifier: NONE,
x: 0, y: 2, fg: Reset, bg: Reset, underline: Reset, modifier: DIM,
x: 4, y: 2, fg: Reset, bg: Reset, underline: Reset, modifier: DIM | ITALIC,
x: 29, y: 2, fg: Reset, bg: Reset, underline: Reset, modifier: NONE,
x: 0, y: 3, fg: Reset, bg: Reset, underline: Reset, modifier: DIM,
x: 14, y: 3, fg: Reset, bg: Reset, underline: Reset, modifier: NONE,
]
}
+189 -56
View File
@@ -37,6 +37,7 @@ use std::sync::atomic::Ordering;
use std::time::Duration;
use std::time::Instant;
use self::realtime::PendingSteerCompareKey;
use crate::app_event::RealtimeAudioDeviceKind;
#[cfg(all(not(target_os = "linux"), feature = "voice-input"))]
use crate::audio_device::list_realtime_audio_device_names;
@@ -85,6 +86,7 @@ use codex_protocol::config_types::ServiceTier;
use codex_protocol::config_types::Settings;
#[cfg(target_os = "windows")]
use codex_protocol::config_types::WindowsSandboxLevel;
use codex_protocol::items::AgentMessageContent;
use codex_protocol::items::AgentMessageItem;
use codex_protocol::models::MessagePhase;
use codex_protocol::models::local_image_label_text;
@@ -618,6 +620,11 @@ pub(crate) struct ChatWidget {
suppress_session_configured_redraw: bool,
// User messages queued while a turn is in progress
queued_user_messages: VecDeque<UserMessage>,
// Steers already submitted to core but not yet committed into history.
//
// The bottom pane shows these above queued drafts until core records the
// corresponding user message item.
pending_steers: VecDeque<PendingSteer>,
/// Terminal-appropriate keybinding for popping the most-recently queued
/// message back into the composer. Determined once at construction time via
/// [`queued_message_edit_binding_for_terminal`] and propagated to
@@ -650,7 +657,9 @@ pub(crate) struct ChatWidget {
had_work_activity: bool,
// Whether the current turn emitted a plan update.
saw_plan_update_this_turn: bool,
// Whether the current turn emitted a proposed plan item.
// Whether the current turn emitted a proposed plan item that has not been superseded by a
// later steer. This is cleared when the user submits a steer so the plan popup only appears
// if a newer proposed plan arrives afterward.
saw_plan_item_this_turn: bool,
// Incremental buffer for streamed plan content.
plan_delta_buffer: String,
@@ -751,6 +760,11 @@ impl From<&str> for UserMessage {
}
}
struct PendingSteer {
user_message: UserMessage,
compare_key: PendingSteerCompareKey,
}
pub(crate) fn create_initial_user_message(
text: Option<String>,
local_image_paths: Vec<PathBuf>,
@@ -778,6 +792,21 @@ pub(crate) fn create_initial_user_message(
}
}
fn append_text_with_rebased_elements(
target_text: &mut String,
target_text_elements: &mut Vec<TextElement>,
text: &str,
text_elements: impl IntoIterator<Item = TextElement>,
) {
let offset = target_text.len();
target_text.push_str(text);
target_text_elements.extend(text_elements.into_iter().map(|mut element| {
element.byte_range.start += offset;
element.byte_range.end += offset;
element
}));
}
// When merging multiple queued drafts (e.g., after interrupt), each draft starts numbering
// its attachments at [Image #1]. Reassign placeholder labels based on the attachment list so
// the combined local_image_paths order matches the labels, even if placeholders were moved
@@ -1290,17 +1319,24 @@ impl ChatWidget {
self.request_redraw();
}
fn on_agent_message(&mut self, message: String) {
// If we have a stream_controller, then the final agent message is redundant and will be a
// duplicate of what has already been streamed.
if self.stream_controller.is_none() && !message.is_empty() {
self.handle_streaming_delta(message);
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();
self.handle_stream_finished();
self.request_redraw();
}
fn on_agent_message(&mut self, message: String) {
self.finalize_completed_assistant_message(Some(&message));
}
fn on_agent_message_delta(&mut self, delta: String) {
self.handle_streaming_delta(delta);
}
@@ -1483,7 +1519,10 @@ impl ChatWidget {
self.unified_exec_wait_streak = None;
self.request_redraw();
if !from_replay && self.queued_user_messages.is_empty() {
let had_pending_steers = !self.pending_steers.is_empty();
self.refresh_pending_input_preview();
if !from_replay && self.queued_user_messages.is_empty() && !had_pending_steers {
self.maybe_prompt_plan_implementation();
}
// Keep this flag for replayed completion events so a subsequent live TurnComplete can
@@ -1865,30 +1904,31 @@ impl ChatWidget {
if reason == TurnAbortReason::Interrupted {
self.clear_unified_exec_processes();
}
if reason != TurnAbortReason::ReviewEnded {
self.add_to_history(history_cell::new_error_event(
"Conversation interrupted - tell the model what to do differently. Something went wrong? Hit `/feedback` to report the issue.".to_owned(),
));
}
if let Some(combined) = self.drain_queued_messages_for_restore() {
// Core clears pending_input before emitting TurnAborted, so any unacknowledged steers
// still tracked here must be restored locally instead of waiting for a later commit.
if let Some(combined) = self.drain_pending_messages_for_restore() {
self.restore_user_message_to_composer(combined);
self.refresh_queued_user_messages();
}
self.refresh_pending_input_preview();
self.request_redraw();
}
/// Merge queued drafts (plus the current composer state) into a single message for restore.
/// Merge pending steers, queued drafts, and the current composer state into a single message.
///
/// Each queued draft numbers attachments from `[Image #1]`. When we concatenate drafts, we
/// must renumber placeholders in a stable order so the merged attachment list stays aligned
/// with the labels embedded in text. This helper drains the queue, remaps placeholders, and
/// fixes text element byte ranges as content is appended. Returns `None` when there is nothing
/// to restore.
fn drain_queued_messages_for_restore(&mut self) -> Option<UserMessage> {
if self.queued_user_messages.is_empty() {
/// Each pending message numbers attachments from `[Image #1]` relative to its own remote
/// images. When we concatenate multiple messages after interrupt, we must renumber local-image
/// placeholders in a stable order and rebase text element byte ranges so the restored composer
/// state stays aligned with the merged attachment list. Returns `None` when there is nothing to
/// restore.
fn drain_pending_messages_for_restore(&mut self) -> Option<UserMessage> {
if self.pending_steers.is_empty() && self.queued_user_messages.is_empty() {
return None;
}
@@ -1900,7 +1940,12 @@ impl ChatWidget {
mention_bindings: self.bottom_pane.composer_mention_bindings(),
};
let mut to_merge: Vec<UserMessage> = self.queued_user_messages.drain(..).collect();
let mut to_merge: Vec<UserMessage> = self
.pending_steers
.drain(..)
.map(|steer| steer.user_message)
.collect();
to_merge.extend(self.queued_user_messages.drain(..));
if !existing_message.text.is_empty()
|| !existing_message.local_images.is_empty()
|| !existing_message.remote_image_urls.is_empty()
@@ -1915,7 +1960,6 @@ impl ChatWidget {
remote_image_urls: Vec::new(),
mention_bindings: Vec::new(),
};
let mut combined_offset = 0usize;
let total_remote_images = to_merge
.iter()
.map(|message| message.remote_image_urls.len())
@@ -1925,22 +1969,23 @@ impl ChatWidget {
for (idx, message) in to_merge.into_iter().enumerate() {
if idx > 0 {
combined.text.push('\n');
combined_offset += 1;
}
let message = remap_placeholders_for_message(message, &mut next_image_label);
let base = combined_offset;
combined.text.push_str(&message.text);
combined_offset += message.text.len();
combined
.text_elements
.extend(message.text_elements.into_iter().map(|mut elem| {
elem.byte_range.start += base;
elem.byte_range.end += base;
elem
}));
combined.local_images.extend(message.local_images);
combined.remote_image_urls.extend(message.remote_image_urls);
combined.mention_bindings.extend(message.mention_bindings);
let UserMessage {
text,
text_elements,
local_images,
remote_image_urls,
mention_bindings,
} = remap_placeholders_for_message(message, &mut next_image_label);
append_text_with_rebased_elements(
&mut combined.text,
&mut combined.text_elements,
&text,
text_elements,
);
combined.local_images.extend(local_images);
combined.remote_image_urls.extend(remote_image_urls);
combined.mention_bindings.extend(mention_bindings);
}
Some(combined)
@@ -2350,6 +2395,15 @@ impl ChatWidget {
/// returns once stream queues are idle. Final-answer completion (or absent
/// phase for legacy models) clears the flag to preserve historical behavior.
fn on_agent_message_item_completed(&mut self, item: AgentMessageItem) {
let mut message = String::new();
for content in &item.content {
match content {
AgentMessageContent::Text { text } => message.push_str(text),
}
}
self.finalize_completed_assistant_message(
(!message.is_empty()).then_some(message.as_str()),
);
self.pending_status_indicator_restore = match item.phase {
// Models that don't support preambles only output AgentMessageItems on turn completion.
Some(MessagePhase::FinalAnswer) | None => false,
@@ -2915,6 +2969,7 @@ impl ChatWidget {
thread_name: None,
forked_from: None,
queued_user_messages: VecDeque::new(),
pending_steers: VecDeque::new(),
queued_message_edit_binding,
show_welcome_banner: is_first_run,
startup_tooltip_override,
@@ -3099,6 +3154,7 @@ impl ChatWidget {
plan_delta_buffer: String::new(),
plan_item_active: false,
queued_user_messages: VecDeque::new(),
pending_steers: VecDeque::new(),
queued_message_edit_binding,
show_welcome_banner: is_first_run,
startup_tooltip_override,
@@ -3264,6 +3320,7 @@ impl ChatWidget {
thread_name: None,
forked_from: None,
queued_user_messages: VecDeque::new(),
pending_steers: VecDeque::new(),
queued_message_edit_binding,
show_welcome_banner: false,
startup_tooltip_override: None,
@@ -3394,7 +3451,7 @@ impl ChatWidget {
{
if let Some(user_message) = self.queued_user_messages.pop_back() {
self.restore_user_message_to_composer(user_message);
self.refresh_queued_user_messages();
self.refresh_pending_input_preview();
self.request_redraw();
}
return;
@@ -3429,18 +3486,19 @@ impl ChatWidget {
.bottom_pane
.take_recent_submission_mention_bindings(),
};
if user_message.text.is_empty()
&& user_message.local_images.is_empty()
&& user_message.remote_image_urls.is_empty()
{
return;
}
let Some(user_message) =
self.maybe_defer_user_message_for_realtime(user_message)
else {
return;
};
// Submissions during active final-answer streaming can race with turn
// completion and strand the UI in a running state. Queue those inputs instead
// of injecting immediately; `on_task_complete()` drains this FIFO via
// `maybe_send_next_queued_input()`, so no typed prompt is dropped.
let should_submit_now = self.is_session_configured()
&& !self.is_plan_streaming_in_tui()
&& self.stream_controller.is_none();
let should_submit_now =
self.is_session_configured() && !self.is_plan_streaming_in_tui();
if should_submit_now {
// Submitted is emitted when user submits.
// Reset any reasoning header only when we are actually submitting a turn.
@@ -4086,7 +4144,7 @@ impl ChatWidget {
|| self.is_review_mode
{
self.queued_user_messages.push_back(user_message);
self.refresh_queued_user_messages();
self.refresh_pending_input_preview();
} else {
self.submit_user_message(user_message);
}
@@ -4096,7 +4154,12 @@ impl ChatWidget {
if !self.is_session_configured() {
tracing::warn!("cannot submit user message before session is configured; queueing");
self.queued_user_messages.push_front(user_message);
self.refresh_queued_user_messages();
self.refresh_pending_input_preview();
return;
}
if self.is_review_mode {
self.queued_user_messages.push_back(user_message);
self.refresh_pending_input_preview();
return;
}
@@ -4123,6 +4186,7 @@ impl ChatWidget {
return;
}
let render_in_history = !self.agent_turn_running;
let mut items: Vec<UserInput> = Vec::new();
// Special-case: "!cmd" executes a local shell command instead of sending to the model.
@@ -4251,6 +4315,16 @@ impl ChatWidget {
} else {
None
};
let pending_steer = (!render_in_history).then(|| PendingSteer {
user_message: UserMessage {
text: text.clone(),
local_images: local_images.clone(),
remote_image_urls: remote_image_urls.clone(),
text_elements: text_elements.clone(),
mention_bindings: mention_bindings.clone(),
},
compare_key: Self::pending_steer_compare_key_from_items(&items),
});
let personality = self
.config
.personality
@@ -4271,9 +4345,9 @@ impl ChatWidget {
personality,
};
self.codex_op_tx.send(op).unwrap_or_else(|e| {
tracing::error!("failed to send message: {e}");
});
if !self.submit_op(op) {
return;
}
// Persist the text to cross-session message history.
if !text.is_empty() {
@@ -4292,8 +4366,14 @@ impl ChatWidget {
});
}
if let Some(pending_steer) = pending_steer {
self.pending_steers.push_back(pending_steer);
self.saw_plan_item_this_turn = false;
self.refresh_pending_input_preview();
}
// Show replayable user content in conversation history.
if !text.is_empty() {
if render_in_history && !text.is_empty() {
let local_image_paths = local_images
.into_iter()
.map(|img| img.path)
@@ -4311,7 +4391,7 @@ impl ChatWidget {
local_image_paths,
remote_image_urls,
));
} else if !remote_image_urls.is_empty() {
} else if render_in_history && !remote_image_urls.is_empty() {
self.last_rendered_user_message_event =
Some(Self::rendered_user_message_event_from_parts(
String::new(),
@@ -4424,9 +4504,18 @@ impl ChatWidget {
match msg {
EventMsg::SessionConfigured(e) => self.on_session_configured(e),
EventMsg::ThreadNameUpdated(e) => self.on_thread_name_updated(e),
EventMsg::AgentMessage(AgentMessageEvent { message, .. }) => {
EventMsg::AgentMessage(AgentMessageEvent { .. })
if matches!(replay_kind, Some(ReplayKind::ThreadSnapshot))
&& !self.is_review_mode => {}
EventMsg::AgentMessage(AgentMessageEvent { message, .. })
if from_replay || self.is_review_mode =>
{
// TODO(ccunningham): stop relying on legacy AgentMessage in review mode,
// including thread-snapshot replay, and forward
// ItemCompleted(TurnItem::AgentMessage(_)) instead.
self.on_agent_message(message)
}
EventMsg::AgentMessage(AgentMessageEvent { .. }) => {}
EventMsg::AgentMessageDelta(AgentMessageDeltaEvent { delta }) => {
self.on_agent_message_delta(delta)
}
@@ -4481,6 +4570,8 @@ impl ChatWidget {
self.on_interrupted_turn(ev.reason);
}
TurnAbortReason::Replaced => {
self.pending_steers.clear();
self.refresh_pending_input_preview();
self.on_error("Turn aborted: replaced by a new task".to_owned())
}
TurnAbortReason::ReviewEnded => {
@@ -4599,6 +4690,42 @@ impl ChatWidget {
}
EventMsg::ItemCompleted(event) => {
let item = event.item;
if !from_replay && let codex_protocol::items::TurnItem::UserMessage(item) = &item {
let EventMsg::UserMessage(event) = item.as_legacy_event() else {
unreachable!("user message item should convert to a legacy user message");
};
let rendered = Self::rendered_user_message_event_from_event(&event);
let compare_key = Self::pending_steer_compare_key_from_item(item);
if self
.pending_steers
.front()
.is_some_and(|pending| pending.compare_key == compare_key)
{
if let Some(pending) = self.pending_steers.pop_front() {
self.refresh_pending_input_preview();
let pending_event = UserMessageEvent {
message: pending.user_message.text,
images: Some(pending.user_message.remote_image_urls),
local_images: pending
.user_message
.local_images
.into_iter()
.map(|image| image.path)
.collect(),
text_elements: pending.user_message.text_elements,
};
self.on_user_message_event(pending_event);
} else if self.last_rendered_user_message_event.as_ref() != Some(&rendered)
{
tracing::warn!(
"pending steer matched compare key but queue was empty when rendering committed user message"
);
self.on_user_message_event(event);
}
} else if self.last_rendered_user_message_event.as_ref() != Some(&rendered) {
self.on_user_message_event(event);
}
}
if let codex_protocol::items::TurnItem::Plan(plan_item) = &item {
self.on_plan_item_completed(plan_item.text.clone());
}
@@ -4749,17 +4876,23 @@ impl ChatWidget {
self.submit_user_message(user_message);
}
// Update the list to reflect the remaining queued messages (if any).
self.refresh_queued_user_messages();
self.refresh_pending_input_preview();
}
/// Rebuild and update the queued user messages from the current queue.
fn refresh_queued_user_messages(&mut self) {
let messages: Vec<String> = self
/// Rebuild and update the bottom-pane pending-input preview.
fn refresh_pending_input_preview(&mut self) {
let queued_messages: Vec<String> = self
.queued_user_messages
.iter()
.map(|m| m.text.clone())
.collect();
self.bottom_pane.set_queued_user_messages(messages);
let pending_steers: Vec<String> = self
.pending_steers
.iter()
.map(|steer| steer.user_message.text.clone())
.collect();
self.bottom_pane
.set_pending_input_preview(queued_messages, pending_steers);
}
pub(crate) fn set_pending_thread_approvals(&mut self, threads: Vec<String>) {
+82 -4
View File
@@ -49,10 +49,16 @@ impl RealtimeConversationUiState {
#[derive(Clone, Debug, PartialEq)]
pub(super) struct RenderedUserMessageEvent {
message: String,
remote_image_urls: Vec<String>,
local_images: Vec<PathBuf>,
text_elements: Vec<TextElement>,
pub(super) message: String,
pub(super) remote_image_urls: Vec<String>,
pub(super) local_images: Vec<PathBuf>,
pub(super) text_elements: Vec<TextElement>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct PendingSteerCompareKey {
pub(super) message: String,
pub(super) image_count: usize,
}
impl ChatWidget {
@@ -81,6 +87,78 @@ impl ChatWidget {
)
}
/// Build the compare key for a submitted pending steer without invoking the
/// expensive request-serialization path. Pending steers only need to match the
/// committed `ItemCompleted(UserMessage)` emitted after core drains input, which
/// preserves flattened text and total image count but not UI-only text ranges or
/// local image paths.
pub(super) fn pending_steer_compare_key_from_items(
items: &[UserInput],
) -> PendingSteerCompareKey {
let mut message = String::new();
let mut image_count = 0;
for item in items {
match item {
UserInput::Text { text, .. } => message.push_str(text),
UserInput::Image { .. } | UserInput::LocalImage { .. } => image_count += 1,
UserInput::Skill { .. } | UserInput::Mention { .. } => {}
_ => {}
}
}
PendingSteerCompareKey {
message,
image_count,
}
}
pub(super) fn pending_steer_compare_key_from_item(
item: &codex_protocol::items::UserMessageItem,
) -> PendingSteerCompareKey {
Self::pending_steer_compare_key_from_items(&item.content)
}
#[cfg(test)]
pub(super) fn rendered_user_message_event_from_inputs(
items: &[UserInput],
) -> RenderedUserMessageEvent {
let mut message = String::new();
let mut remote_image_urls = Vec::new();
let mut local_images = Vec::new();
let mut text_elements = Vec::new();
for item in items {
match item {
UserInput::Text {
text,
text_elements: current_text_elements,
} => append_text_with_rebased_elements(
&mut message,
&mut text_elements,
text,
current_text_elements.iter().map(|element| {
TextElement::new(
element.byte_range,
element.placeholder(text).map(str::to_string),
)
}),
),
UserInput::Image { image_url } => remote_image_urls.push(image_url.clone()),
UserInput::LocalImage { path } => local_images.push(path.clone()),
UserInput::Skill { .. } | UserInput::Mention { .. } => {}
_ => {}
}
}
Self::rendered_user_message_event_from_parts(
message,
text_elements,
local_images,
remote_image_urls,
)
}
pub(super) fn should_render_realtime_user_message_event(
&self,
event: &UserMessageEvent,
@@ -1,5 +0,0 @@
---
source: tui/src/chatwidget/tests.rs
expression: combined
---
• Here is the result.
File diff suppressed because it is too large Load Diff