From 4d7a5c7c7394b687ebcb67e634528b2b8c5578d9 Mon Sep 17 00:00:00 2001 From: Charlie Marsh Date: Sun, 19 Jul 2026 16:42:26 +0000 Subject: [PATCH] Avoid liveness races when starting side conversations (#34199) ## Why The `thread/started` notification for a newly forked side conversation can arrive after the fork response. Selecting the side thread in that window could incorrectly report that it was unavailable. ## What changed - Seed agent navigation from the side-fork response before selecting the new thread. - Skip redundant liveness and parent-title reads for side threads that already have local state, while preserving liveness checks for uncached agent threads. ## Testing - Cover side-thread selection before `thread/started` is delivered. - Verify uncached threads are still checked and regular forks still resolve their parent title. GitOrigin-RevId: 1f9fb0586db8094bf5a2624bc4e1de06c9a783a1 --- codex-rs/tui/src/app/session_lifecycle.rs | 10 ++- codex-rs/tui/src/app/side.rs | 8 ++ codex-rs/tui/src/app/tests.rs | 98 +++++++++++++++++++++++ codex-rs/tui/src/app_server_session.rs | 55 ++++++++++++- 4 files changed, 165 insertions(+), 6 deletions(-) diff --git a/codex-rs/tui/src/app/session_lifecycle.rs b/codex-rs/tui/src/app/session_lifecycle.rs index 01d73dcfe340..e0eee231ef76 100644 --- a/codex-rs/tui/src/app/session_lifecycle.rs +++ b/codex-rs/tui/src/app/session_lifecycle.rs @@ -408,9 +408,13 @@ impl App { return Ok(()); } - if !self - .refresh_agent_picker_thread_liveness(app_server, thread_id) - .await + // A tracked side thread stays loaded until it is explicitly discarded and already has a + // replay channel, so another liveness read cannot add anything before selection. + if !(self.side_threads.contains_key(&thread_id) + && self.thread_event_channels.contains_key(&thread_id) + || self + .refresh_agent_picker_thread_liveness(app_server, thread_id) + .await) { self.chat_widget .add_error_message(format!("Agent thread {thread_id} is no longer available.")); diff --git a/codex-rs/tui/src/app/side.rs b/codex-rs/tui/src/app/side.rs index 67e9340559ff..50cf8a92134e 100644 --- a/codex-rs/tui/src/app/side.rs +++ b/codex-rs/tui/src/app/side.rs @@ -587,6 +587,14 @@ impl App { } self.side_threads .insert(child_thread_id, SideThreadState::new(parent_thread_id)); + // `thread/started` is delivered after the fork response; seed navigation before + // the first selection without blocking on another app-server read. + self.upsert_agent_picker_thread( + child_thread_id, + /*agent_nickname*/ None, + /*agent_role*/ None, + /*is_closed*/ false, + ); if let Err(err) = app_server .thread_inject_items(child_thread_id, vec![Self::side_boundary_prompt_item()]) .await diff --git a/codex-rs/tui/src/app/tests.rs b/codex-rs/tui/src/app/tests.rs index 270dd9d791e1..48efbba2e08a 100644 --- a/codex-rs/tui/src/app/tests.rs +++ b/codex-rs/tui/src/app/tests.rs @@ -1989,6 +1989,104 @@ async fn refresh_agent_picker_thread_liveness_prunes_closed_metadata_only_thread Ok(()) } +#[tokio::test] +async fn handle_start_side_seeds_navigation_before_thread_started() -> Result<()> { + let (mut app, mut app_event_rx, _op_rx) = make_test_app_with_channels().await; + let config = app.chat_widget.config_ref().clone(); + let parent_thread_id = ThreadId::from_string( + &app_test_support::create_fake_rollout( + config.codex_home.as_path(), + "2025-01-05T12-00-00", + "2025-01-05T12:00:00Z", + "Saved user message", + Some(config.model_provider_id.as_str()), + /*git_info*/ None, + ) + .expect("create source rollout"), + )?; + let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(&config)).await?; + let started = app_server + .resume_thread( + config, + parent_thread_id, + crate::app_server_session::ResumeModelSettings::RestoreFromThread, + ) + .await?; + app.enqueue_primary_thread_session(started.session, started.turns) + .await?; + while app_event_rx.try_recv().is_ok() {} + let mut tui = crate::tui::test_support::make_test_tui()?; + + let control = Box::pin(app.handle_start_side( + &mut tui, + &mut app_server, + parent_thread_id, + /*user_message*/ None, + )) + .await?; + + let side_thread_id = app + .active_thread_id + .expect("side conversation should become active"); + assert!(matches!(control, AppRunControl::Continue)); + assert_ne!(side_thread_id, parent_thread_id); + assert!(app.side_threads.contains_key(&side_thread_id)); + assert!(app.thread_event_channels.contains_key(&side_thread_id)); + assert!( + !app.agent_navigation + .get(&side_thread_id) + .expect("side start should seed navigation before thread/started") + .is_closed + ); + + let mut saw_thread_started = false; + for _ in 0..20 { + let event = time::timeout( + std::time::Duration::from_secs(/*secs*/ 2), + app_server.next_event(), + ) + .await + .expect("app-server should emit an event") + .expect("app-server event stream should remain open"); + if let codex_app_server_client::AppServerEvent::ServerNotification( + ServerNotification::ThreadStarted(notification), + ) = event + && notification.thread.id == side_thread_id.to_string() + { + saw_thread_started = true; + break; + } + } + + assert!(saw_thread_started); + app_server.shutdown().await?; + Ok(()) +} + +#[tokio::test] +async fn select_uncached_agent_thread_still_refreshes_liveness() -> Result<()> { + let mut app = Box::pin(make_test_app()).await; + let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker( + app.chat_widget.config_ref(), + )) + .await?; + let thread_id = ThreadId::new(); + app.agent_navigation.upsert( + thread_id, + Some("Ghost".to_string()), + Some("worker".to_string()), + /*is_closed*/ false, + ); + let mut tui = crate::tui::test_support::make_test_tui()?; + + Box::pin(app.select_agent_thread(&mut tui, &mut app_server, thread_id)).await?; + + assert_eq!(app.active_thread_id, None); + assert_eq!(app.agent_navigation.get(&thread_id), None); + app_server.shutdown().await?; + Ok(()) +} + #[tokio::test] async fn open_agent_picker_prompts_to_enable_multi_agent_when_disabled() -> Result<()> { let (mut app, mut app_event_rx, _op_rx) = Box::pin(make_test_app_with_channels()).await; diff --git a/codex-rs/tui/src/app_server_session.rs b/codex-rs/tui/src/app_server_session.rs index 50d59397cdb5..02ccd9619494 100644 --- a/codex-rs/tui/src/app_server_session.rs +++ b/codex-rs/tui/src/app_server_session.rs @@ -641,9 +641,12 @@ impl AppServerSession { .map_err(|err| { bootstrap_request_error("thread/fork failed during TUI bootstrap", err) })?; - let fork_parent_title = self - .fork_parent_title_from_app_server(response.thread.forked_from_id.as_deref()) - .await; + let fork_parent_title = if presentation == ForkPresentation::SideConversation { + None + } else { + self.fork_parent_title_from_app_server(response.thread.forked_from_id.as_deref()) + .await + }; let mut started = started_thread_from_fork_response(response, &config, self.thread_params_mode()).await?; started.session.fork_parent_title = fork_parent_title; @@ -2444,6 +2447,52 @@ mod tests { Ok(()) } + #[tokio::test] + async fn side_fork_skips_parent_title_lookup_but_normal_ephemeral_fork_keeps_it() -> Result<()> + { + let codex_home = tempfile::tempdir().expect("tempdir"); + let config = build_config(&codex_home).await; + let source_thread_id = ThreadId::from_string( + &create_fake_rollout( + codex_home.path(), + "2025-01-05T12-00-00", + "2025-01-05T12:00:00Z", + "Saved user message", + Some(config.model_provider_id.as_str()), + /*git_info*/ None, + ) + .expect("create source rollout"), + )?; + let mut app_server = crate::start_embedded_app_server_for_picker(&config).await?; + app_server + .resume_thread( + config.clone(), + source_thread_id, + ResumeModelSettings::RestoreFromThread, + ) + .await?; + app_server + .thread_set_name(source_thread_id, "Source thread".to_string()) + .await?; + + let mut ephemeral_config = config; + ephemeral_config.ephemeral = true; + let normal_ephemeral_fork = app_server + .fork_thread(ephemeral_config.clone(), source_thread_id) + .await?; + let side_fork = app_server + .fork_side_thread(ephemeral_config, source_thread_id) + .await?; + + assert_eq!( + normal_ephemeral_fork.session.fork_parent_title.as_deref(), + Some("Source thread") + ); + assert_eq!(side_fork.session.fork_parent_title, None); + app_server.shutdown().await?; + Ok(()) + } + #[tokio::test] async fn config_request_overrides_preserve_implicit_personality_default() { let temp_dir = tempfile::tempdir().expect("tempdir");