Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions codex-rs/tui/src/app/event_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,7 @@ impl App {
self.chat_widget.add_error_message(format!(
"Failed to start TUI session picker: {err}"
));
self.chat_widget.maybe_send_next_queued_input();
return Ok(AppRunControl::Continue);
}
};
Expand Down Expand Up @@ -112,6 +113,7 @@ impl App {
SessionSelection::Fork(_) => {}
}

self.chat_widget.maybe_send_next_queued_input();
// Leaving alt-screen may blank the inline viewport; force a redraw either way.
tui.frame_requester().schedule_frame();
}
Expand Down Expand Up @@ -229,6 +231,7 @@ impl App {
);
}

self.chat_widget.maybe_send_next_queued_input();
tui.frame_requester().schedule_frame();
}
AppEvent::ForkSessionForPromptEdit {
Expand Down
8 changes: 8 additions & 0 deletions codex-rs/tui/src/chatwidget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1770,6 +1770,14 @@ impl ChatWidget {
}

pub(crate) fn prepare_local_op_submission(&mut self, op: &AppCommand) {
if matches!(
op,
AppCommand::Compact
| AppCommand::Review { .. }
| AppCommand::RunUserShellCommand { .. }
) {
self.input_queue.user_turn_pending_start = true;
}
if matches!(op, AppCommand::Interrupt) && self.turn_lifecycle.agent_turn_running {
if let Some(controller) = self.stream_controller.as_mut() {
controller.clear_queue();
Expand Down
20 changes: 14 additions & 6 deletions codex-rs/tui/src/chatwidget/input_flow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,9 @@ impl ChatWidget {
}
let should_submit_now = self.is_session_configured()
&& !self.is_plan_streaming_in_tui()
&& !self.input_queue.suppress_queue_autosend;
&& !self.input_queue.suppress_queue_autosend
&& (!self.input_queue.user_turn_pending_start
|| self.turn_lifecycle.agent_turn_running);
if should_submit_now {
if self.only_user_shell_commands_running()
&& !user_message.text.starts_with('!')
Expand Down Expand Up @@ -108,10 +110,10 @@ impl ChatWidget {
action: QueuedInputAction,
pending_pastes: Vec<(String, String)>,
) {
if !self.is_session_configured()
|| self.is_user_turn_pending_or_running()
|| self.input_queue.suppress_queue_autosend
{
let should_run_now = self.is_session_configured()
&& !self.is_user_turn_pending_or_running()
&& !self.input_queue.suppress_queue_autosend;
if !should_run_now || action != QueuedInputAction::Plain {
self.input_queue
.queued_user_messages
.push_back(QueuedUserMessage {
Expand All @@ -123,6 +125,9 @@ impl ChatWidget {
.queued_user_message_history_records
.push_back(UserMessageHistoryRecord::UserMessageText);
self.refresh_pending_input_preview();
if should_run_now {
self.maybe_send_next_queued_input();
}
} else {
self.submit_user_message(user_message);
}
Expand Down Expand Up @@ -174,7 +179,10 @@ impl ChatWidget {
}

pub(super) fn is_user_turn_pending_or_running(&self) -> bool {
self.input_queue.user_turn_pending_start || self.bottom_pane.is_task_running()
self.input_queue.user_turn_pending_start
|| self.turn_lifecycle.agent_turn_running
|| self.review.is_review_mode
|| (self.bottom_pane.is_task_running() && self.mcp_startup_status.is_none())
}

pub(super) fn only_user_shell_commands_running(&self) -> bool {
Expand Down
3 changes: 3 additions & 0 deletions codex-rs/tui/src/chatwidget/mcp_startup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,9 @@ impl ChatWidget {
self.mcp_startup_pending_next_round.clear();
self.mcp_startup_pending_next_round_saw_starting = false;
self.update_task_running_state();
if self.input_queue.user_turn_pending_start {
self.bottom_pane.set_task_running(/*running*/ true);
}
if self.bottom_pane.is_task_running() && mcp_startup_owned_status {
self.restore_reasoning_status_header();
}
Expand Down
11 changes: 10 additions & 1 deletion codex-rs/tui/src/chatwidget/slash_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -132,7 +132,12 @@ impl ChatWidget {
}

fn slash_command_blocked_by_active_task(&self, cmd: SlashCommand) -> bool {
(!cmd.available_during_task() && self.bottom_pane.is_task_running())
(!cmd.available_during_task()
&& (self.turn_lifecycle.agent_turn_running
|| self.review.is_review_mode
|| (self.bottom_pane.is_task_running()
&& (self.mcp_startup_status.is_none()
|| self.input_queue.user_turn_pending_start))))
|| (cmd == SlashCommand::Resume
&& (self.input_queue.user_turn_pending_start
|| self.turn_lifecycle.agent_turn_running))
Expand Down Expand Up @@ -262,10 +267,14 @@ impl ChatWidget {
if !self.bottom_pane.is_task_running() {
self.bottom_pane.set_task_running(/*running*/ true);
}
self.input_queue.user_turn_pending_start = true;
self.app_event_tx.compact();
}
SlashCommand::Review => {
self.open_review_popup();
if self.mcp_startup_status.is_some() {
self.defer_input_until_settings_applied();
}
}
SlashCommand::Rename => {
self.session_telemetry
Expand Down
154 changes: 154 additions & 0 deletions codex-rs/tui/src/chatwidget/tests/mcp_startup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,160 @@ async fn mcp_startup_complete_does_not_clear_running_task() {
assert_eq!(chat.status_state.current_status.header, "Working");
}

#[tokio::test]
async fn pending_mcp_startup_does_not_block_queued_follow_up() {
let (mut chat, _rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());
handle_turn_started(&mut chat, "turn-1");
chat.queue_user_message("queued follow-up".into());

handle_turn_completed(&mut chat, "turn-1", /*duration_ms*/ None);

assert!(chat.mcp_startup_status.is_some());
assert!(chat.bottom_pane.is_task_running());
assert!(chat.input_queue.queued_user_messages.is_empty());
assert_matches!(next_submit_op(&mut op_rx), Op::UserTurn { items, .. } if items == vec![
UserInput::Text {
text: "queued follow-up".to_string(),
text_elements: Vec::new(),
}
]);

chat.finish_mcp_startup(Vec::new(), Vec::new());

assert!(chat.input_queue.user_turn_pending_start);
}

#[tokio::test]
async fn pending_mcp_startup_dispatches_queued_slash_commands() {
let (mut chat, mut rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());
chat.bottom_pane
.set_composer_text("/resume".to_string(), Vec::new(), Vec::new());

chat.handle_key_event(KeyEvent::new(KeyCode::Tab, KeyModifiers::NONE));

assert_matches!(rx.try_recv(), Ok(AppEvent::OpenResumePicker));
assert_no_submit_op(&mut op_rx);
}

#[tokio::test]
async fn pending_mcp_startup_does_not_reject_queued_compaction() {
let (mut chat, mut rx, _op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());
handle_turn_started(&mut chat, "turn-1");
chat.bottom_pane
.set_composer_text("/compact".to_string(), Vec::new(), Vec::new());
chat.handle_key_event(KeyEvent::new(KeyCode::Tab, KeyModifiers::NONE));

handle_turn_completed(&mut chat, "turn-1", /*duration_ms*/ None);

assert!(
std::iter::from_fn(|| rx.try_recv().ok())
.any(|event| matches!(event, AppEvent::CodexOp(Op::Compact)))
);
}

#[tokio::test]
async fn pending_mcp_startup_does_not_drain_follow_up_before_review_starts() {
let (mut chat, mut rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());
handle_turn_started(&mut chat, "turn-1");
for message in ["/review", "queued follow-up"] {
chat.bottom_pane
.set_composer_text(message.to_string(), Vec::new(), Vec::new());
chat.handle_key_event(KeyEvent::new(KeyCode::Tab, KeyModifiers::NONE));
}
handle_turn_completed(&mut chat, "turn-1", /*duration_ms*/ None);

chat.handle_key_event(KeyEvent::new(KeyCode::Down, KeyModifiers::NONE));
chat.handle_key_event(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE));

assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert!(
std::iter::from_fn(|| rx.try_recv().ok())
.any(|event| { matches!(event, AppEvent::CodexOp(Op::Review { .. })) })
);
assert_no_submit_op(&mut op_rx);
}

#[tokio::test]
async fn pending_mcp_startup_does_not_unblock_external_review() {
let (mut chat, mut rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());

handle_entered_review_mode(&mut chat, "current changes");
chat.queue_user_message("queued follow-up".into());

assert!(chat.review.is_review_mode);
assert!(!chat.turn_lifecycle.agent_turn_running);
assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert_no_submit_op(&mut op_rx);

while rx.try_recv().is_ok() {}
chat.dispatch_command(crate::slash_command::SlashCommand::Fork);
assert_matches!(rx.try_recv(), Ok(AppEvent::InsertHistoryCell(_)));

chat.finish_mcp_startup(Vec::new(), Vec::new());
assert!(chat.bottom_pane.is_task_running());
assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert_no_submit_op(&mut op_rx);
}

#[tokio::test]
async fn pending_mcp_startup_does_not_unblock_foreground_shell() {
let (mut chat, mut rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());

chat.queue_user_message_with_options(
"!echo hi".into(),
QueuedInputAction::RunShell,
Vec::new(),
);

assert_matches!(op_rx.try_recv(), Ok(Op::RunUserShellCommand { command }) if command == "echo hi");
chat.bottom_pane
.set_composer_text("queued follow-up".to_string(), Vec::new(), Vec::new());
chat.handle_key_event(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE));
assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert_no_submit_op(&mut op_rx);

chat.finish_mcp_startup(Vec::new(), Vec::new());
assert!(chat.bottom_pane.is_task_running());
assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert_no_submit_op(&mut op_rx);

while rx.try_recv().is_ok() {}
chat.dispatch_command(crate::slash_command::SlashCommand::Fork);
assert_matches!(rx.try_recv(), Ok(AppEvent::InsertHistoryCell(_)));
}

#[tokio::test]
async fn pending_mcp_startup_does_not_unblock_foreground_compaction() {
let (mut chat, _rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
chat.dispatch_command(crate::slash_command::SlashCommand::Compact);
chat.set_mcp_startup_expected_servers(["slow".to_string()]);
notify_mcp_status(&mut chat, "slow", McpServerStartupState::Starting);
chat.thread_id = Some(ThreadId::new());

chat.queue_user_message("queued follow-up".into());

assert_eq!(chat.input_queue.queued_user_messages.len(), 1);
assert_no_submit_op(&mut op_rx);
}

#[tokio::test]
async fn turn_start_preserves_active_mcp_startup_header() {
let (mut chat, _rx, _op_rx) = make_chatwidget_manual(/*model_override*/ None).await;
Expand Down
2 changes: 2 additions & 0 deletions codex-rs/tui/src/chatwidget/tests/slash_commands.rs
Original file line number Diff line number Diff line change
Expand Up @@ -316,13 +316,15 @@ async fn queued_bang_shell_waits_for_user_shell_completion_before_next_input() {
assert_eq!(next_add_to_history_event(&mut rx), "!echo hi");
assert_eq!(chat.input_queue.queued_user_messages.len(), 1);

handle_turn_started(&mut chat, "turn-2");
let begin = begin_exec_with_source(
&mut chat,
"user-shell-echo",
"echo hi",
ExecCommandSource::UserShell,
);
end_exec(&mut chat, begin, "hi\n", "", /*exit_code*/ 0);
handle_turn_completed(&mut chat, "turn-2", /*duration_ms*/ None);

match next_submit_op(&mut op_rx) {
Op::UserTurn { items, .. } => assert_eq!(
Expand Down
4 changes: 3 additions & 1 deletion codex-rs/tui/src/chatwidget/turn_runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,9 @@ impl ChatWidget {
/// both the agent turn lifecycle and MCP startup lifecycle.
pub(super) fn update_task_running_state(&mut self) {
self.bottom_pane.set_task_running(
self.turn_lifecycle.agent_turn_running || self.mcp_startup_status.is_some(),
self.turn_lifecycle.agent_turn_running
|| self.review.is_review_mode
|| self.mcp_startup_status.is_some(),
);
self.refresh_plan_mode_nudge();
self.refresh_status_surfaces();
Expand Down
Loading