From 08ae0fc0cef06a1b57d134019e924692e5be9ffe Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Wed, 22 Jul 2026 19:24:03 +0000 Subject: [PATCH] Consolidate thread startup around `StartThreadOptions` (#34814) ## What changed - Add `StartThreadOptions::new` to provide the standard configuration for a new thread. - Make `ThreadManager::start_thread` the single thread-start entry point and migrate callers from the previous convenience methods. - Derive default environment selections when `environments` is `None`, while preserving explicit selections, including an empty list. GitOrigin-RevId: 8977dc11aed54c5e1215a81eaed2b2cf5fc6087a --- .../app-server/src/bespoke_event_handling.rs | 16 ++- codex-rs/app-server/src/mcp_refresh.rs | 13 ++- .../request_processors/thread_processor.rs | 7 +- .../core/src/agent/control/residency_tests.rs | 5 +- codex-rs/core/src/agent/control_tests.rs | 62 +++------- codex-rs/core/src/prompt_debug.rs | 5 +- codex-rs/core/src/thread_manager.rs | 75 ++++++------- codex-rs/core/src/thread_manager_tests.rs | 106 +++++------------- .../src/tools/handlers/multi_agents_tests.rs | 99 ++++++++-------- codex-rs/core/tests/common/test_codex.rs | 24 +--- codex-rs/core/tests/suite/agents_md.rs | 35 ++---- codex-rs/core/tests/suite/audio_truncation.rs | 6 +- codex-rs/core/tests/suite/client.rs | 3 +- codex-rs/core/tests/suite/code_mode.rs | 17 +-- codex-rs/core/tests/suite/compact_remote.rs | 16 ++- codex-rs/core/tests/suite/git_enrichment.rs | 52 +-------- codex-rs/core/tests/suite/hooks.rs | 16 +-- codex-rs/core/tests/suite/mcp_tool_cache.rs | 3 +- codex-rs/core/tests/suite/search_tool.rs | 11 +- codex-rs/core/tests/suite/sqlite_state.rs | 8 +- .../tests/suite/subagent_notifications.rs | 16 +-- codex-rs/ext/agent/src/lib.rs | 15 +-- codex-rs/mcp-server/src/codex_tool_runner.rs | 6 +- codex-rs/memories/write/src/runtime.rs | 17 +-- codex-rs/thread-manager-sample/src/main.rs | 3 +- 25 files changed, 237 insertions(+), 399 deletions(-) diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index 3775994a6909..dce864366f60 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -2663,7 +2663,9 @@ mod tests { thread_id: conversation_id, thread: conversation, .. - } = thread_manager.start_thread(config.clone()).await?; + } = thread_manager + .start_thread(codex_core::StartThreadOptions::new(config.clone())) + .await?; let thread_state = new_thread_state(); let thread_watch_manager = ThreadWatchManager::new(); let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY); @@ -3241,7 +3243,9 @@ mod tests { thread_id: conversation_id, thread: conversation, .. - } = thread_manager.start_thread(config.clone()).await?; + } = thread_manager + .start_thread(codex_core::StartThreadOptions::new(config.clone())) + .await?; let thread_state = new_thread_state(); { let mut state = thread_state.lock().await; @@ -3329,7 +3333,9 @@ mod tests { thread_id: conversation_id, thread: conversation, .. - } = thread_manager.start_thread(config).await?; + } = thread_manager + .start_thread(codex_core::StartThreadOptions::new(config)) + .await?; let child_thread_id = ThreadId::new(); let child_thread_id_string = child_thread_id.to_string(); let thread_watch_manager = ThreadWatchManager::new(); @@ -3419,7 +3425,9 @@ mod tests { thread_id: conversation_id, thread: conversation, .. - } = thread_manager.start_thread(config).await?; + } = thread_manager + .start_thread(codex_core::StartThreadOptions::new(config)) + .await?; let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY); let outgoing = Arc::new(OutgoingMessageSender::new( tx, diff --git a/codex-rs/app-server/src/mcp_refresh.rs b/codex-rs/app-server/src/mcp_refresh.rs index a45d51056a16..b101a8bd0a01 100644 --- a/codex-rs/app-server/src/mcp_refresh.rs +++ b/codex-rs/app-server/src/mcp_refresh.rs @@ -198,7 +198,10 @@ mod tests { Some(temp_dir.path().join("good")), ) .await?; - let thread = thread_manager.start_thread(thread_config).await?.thread; + let thread = thread_manager + .start_thread(codex_core::StartThreadOptions::new(thread_config)) + .await? + .thread; std::fs::write( temp_dir.path().join(codex_config::CONFIG_TOML_FILE), r#" @@ -337,8 +340,12 @@ enabled = false /*external_time_provider*/ None, ) }); - thread_manager.start_thread(good_config).await?; - thread_manager.start_thread(bad_config).await?; + thread_manager + .start_thread(codex_core::StartThreadOptions::new(good_config)) + .await?; + thread_manager + .start_thread(codex_core::StartThreadOptions::new(bad_config)) + .await?; let loader = Arc::new(CountingThreadConfigLoader { good_cwd: AbsolutePathBuf::try_from(good_cwd)?, diff --git a/codex-rs/app-server/src/request_processors/thread_processor.rs b/codex-rs/app-server/src/request_processors/thread_processor.rs index 802859e2508d..812b42707913 100644 --- a/codex-rs/app-server/src/request_processors/thread_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_processor.rs @@ -1228,8 +1228,7 @@ impl ThreadRequestProcessor { .. } = listener_task_context .thread_manager - .start_thread_with_options(StartThreadOptions { - config, + .start_thread(StartThreadOptions { allow_provider_model_fallback, initial_history: match session_start_source .unwrap_or(codex_app_server_protocol::ThreadStartSource::Startup) @@ -1238,14 +1237,14 @@ impl ThreadRequestProcessor { codex_app_server_protocol::ThreadStartSource::Clear => InitialHistory::Cleared, }, history_mode, - session_source: None, thread_source, dynamic_tools, metrics_service_name: service_name, parent_trace: request_trace, - environments, + environments: Some(environments), thread_extension_init, supports_openai_form_elicitation, + ..StartThreadOptions::new(config) }) .instrument(tracing::info_span!( "app_server.thread_start.create_thread", diff --git a/codex-rs/core/src/agent/control/residency_tests.rs b/codex-rs/core/src/agent/control/residency_tests.rs index 7efd82e92f04..cc51e69da79a 100644 --- a/codex-rs/core/src/agent/control/residency_tests.rs +++ b/codex-rs/core/src/agent/control/residency_tests.rs @@ -1,3 +1,4 @@ +use crate::StartThreadOptions; use crate::ThreadManager; use crate::agent::AgentControl; use crate::codex_thread::CodexThread; @@ -33,7 +34,7 @@ async fn residency_slot_reservation_unloads_oldest_idle_v2_agent() { Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), ); let root = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start root thread"); let control = manager.agent_control(); @@ -79,7 +80,7 @@ async fn interrupted_v2_agent_is_lost_after_residency_eviction() { Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), ); let root = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start root thread"); let control = manager.agent_control(); diff --git a/codex-rs/core/src/agent/control_tests.rs b/codex-rs/core/src/agent/control_tests.rs index 35e62d5d39b0..260b95a483e0 100644 --- a/codex-rs/core/src/agent/control_tests.rs +++ b/codex-rs/core/src/agent/control_tests.rs @@ -36,7 +36,6 @@ use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::CompactedItem; use codex_protocol::protocol::ErrorEvent; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::InterAgentCommunication; use codex_protocol::protocol::ItemCompletedEvent; use codex_protocol::protocol::RolloutItem; @@ -160,7 +159,7 @@ impl AgentControlHarness { async fn start_thread(&self) -> (ThreadId, Arc) { let new_thread = self .manager - .start_thread(self.config.clone()) + .start_thread(StartThreadOptions::new(self.config.clone())) .await .expect("start thread"); (new_thread.thread_id, new_thread.thread) @@ -169,19 +168,10 @@ impl AgentControlHarness { async fn start_paginated_thread(&self) -> (ThreadId, Arc) { let new_thread = self .manager - .start_thread_with_options(StartThreadOptions { - config: self.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, + .start_thread(StartThreadOptions { history_mode: Some(ThreadHistoryMode::Paginated), - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: ExtensionDataInit::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(self.config.clone()) }) .await .expect("start paginated thread"); @@ -1219,7 +1209,7 @@ async fn spawn_agent_can_fork_parent_thread_history_with_sanitized_items() { Some("Child subagent guidance.".to_string()); let new_thread = harness .manager - .start_thread(parent_config.clone()) + .start_thread(StartThreadOptions::new(parent_config.clone())) .await .expect("start parent thread"); let parent_thread_id = new_thread.thread_id; @@ -1445,7 +1435,7 @@ async fn spawn_agent_fork_strips_parent_usage_hints_from_compacted_history() { Some("Child subagent guidance.".to_string()); let new_thread = harness .manager - .start_thread(parent_config) + .start_thread(StartThreadOptions::new(parent_config)) .await .expect("start parent thread"); let parent_thread_id = new_thread.thread_id; @@ -1750,19 +1740,10 @@ async fn spawn_agent_fork_last_n_turns_drops_parent_startup_prefix_when_under_li thread_extension_init.insert(selected_capability_roots.clone()); let parent = harness .manager - .start_thread_with_options(StartThreadOptions { - config: harness.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), + .start_thread(StartThreadOptions { + environments: Some(Vec::new()), thread_extension_init, - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(harness.config.clone()) }) .await .expect("start parent thread"); @@ -1876,7 +1857,7 @@ async fn spawn_agent_fork_last_n_turns_strips_parent_usage_hints() { Some("Child subagent guidance.".to_string()); let new_thread = harness .manager - .start_thread(parent_config) + .start_thread(StartThreadOptions::new(parent_config)) .await .expect("start parent thread"); let parent_thread_id = new_thread.thread_id; @@ -1976,7 +1957,7 @@ async fn spawn_agent_respects_legacy_max_threads_alias() { let control = manager.agent_control(); let _ = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); @@ -2227,7 +2208,7 @@ async fn multi_agent_v2_completion_ignores_dead_direct_parent() { let _ = config.features.enable(Feature::MultiAgentV2); let root = harness .manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("root thread should start"); let root_thread_id = root.thread_id; @@ -2339,7 +2320,7 @@ async fn multi_agent_v2_completion_queues_message_for_direct_parent() { let _ = tester_config.features.enable(Feature::MultiAgentV2); let tester_thread_id = harness .manager - .start_thread(tester_config.clone()) + .start_thread(StartThreadOptions::new(tester_config.clone())) .await .expect("tester thread should start") .thread_id; @@ -2521,19 +2502,10 @@ async fn spawn_thread_subagents_persist_parent_originator_across_new_and_truncat let harness = AgentControlHarness::new().await; let parent = harness .manager - .start_thread_with_options(StartThreadOptions { - config: harness.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), + .start_thread(StartThreadOptions { metrics_service_name: Some("codex_work_desktop".to_string()), - parent_trace: None, - environments: Vec::new(), - thread_extension_init: ExtensionDataInit::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(harness.config.clone()) }) .await .expect("parent thread should start"); @@ -3044,7 +3016,7 @@ async fn list_agent_subtree_thread_ids_finds_live_descendants_of_unloaded_root() ); let control = manager.agent_control(); let parent_thread_id = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("parent should start") .thread_id; diff --git a/codex-rs/core/src/prompt_debug.rs b/codex-rs/core/src/prompt_debug.rs index d8174bd911c2..7670977c9ae9 100644 --- a/codex-rs/core/src/prompt_debug.rs +++ b/codex-rs/core/src/prompt_debug.rs @@ -17,6 +17,7 @@ use crate::session::session::Session; use crate::session::turn::build_prompt; use crate::session::turn::built_tools; use crate::state_db_bridge::StateDbHandle; +use crate::thread_manager::StartThreadOptions; use crate::thread_manager::ThreadManager; use crate::thread_manager::thread_store_from_config; use codex_extension_api::empty_extension_registry; @@ -64,7 +65,9 @@ pub async fn build_prompt_input( /*attestation_provider*/ None, /*external_time_provider*/ None, ); - let thread = thread_manager.start_thread(config).await?; + let thread = thread_manager + .start_thread(StartThreadOptions::new(config)) + .await?; let output = build_prompt_input_from_session(&thread.thread.session, input).await; let shutdown = thread.thread.shutdown_and_wait().await; diff --git a/codex-rs/core/src/thread_manager.rs b/codex-rs/core/src/thread_manager.rs index 451a6d6d56d3..3a4b876dc391 100644 --- a/codex-rs/core/src/thread_manager.rs +++ b/codex-rs/core/src/thread_manager.rs @@ -198,11 +198,30 @@ pub struct StartThreadOptions { pub dynamic_tools: Vec, pub metrics_service_name: Option, pub parent_trace: Option, - pub environments: Vec, + pub environments: Option>, pub thread_extension_init: ExtensionDataInit, pub supports_openai_form_elicitation: bool, } +impl StartThreadOptions { + pub fn new(config: Config) -> Self { + Self { + config, + allow_provider_model_fallback: false, + initial_history: InitialHistory::New, + history_mode: None, + session_source: None, + thread_source: None, + dynamic_tools: Vec::new(), + metrics_service_name: None, + parent_trace: None, + environments: None, + thread_extension_init: ExtensionDataInit::default(), + supports_openai_form_elicitation: false, + } + } +} + fn originator_from_service_name(service_name: Option<&str>) -> Option { let service_name = service_name?.trim(); for originator in [ @@ -662,52 +681,22 @@ impl ThreadManager { Ok(subtree_thread_ids) } - pub async fn start_thread(&self, config: Config) -> CodexResult { - // Box delegated thread-spawn futures so these convenience wrappers do - // not inline the full spawn path into every caller's async state. - Box::pin(self.start_thread_with_tools(config, Vec::new())).await - } - - pub async fn start_thread_with_tools( - &self, - config: Config, - dynamic_tools: Vec, - ) -> CodexResult { - let environments = default_thread_environment_selections( - self.state.environment_manager.as_ref(), - &config.cwd, - &config.workspace_roots, - ); - Box::pin(self.start_thread_with_options(StartThreadOptions { - config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools, - metrics_service_name: None, - parent_trace: None, - environments, - thread_extension_init: ExtensionDataInit::default(), - supports_openai_form_elicitation: false, - })) - .await - } - - pub async fn start_thread_with_options( - &self, - options: StartThreadOptions, - ) -> CodexResult { - self.start_thread_with_options_and_fork_source(options, /*forked_from_thread_id*/ None) - .await + pub async fn start_thread(&self, options: StartThreadOptions) -> CodexResult { + Box::pin(self.start_thread_inner(options, /*forked_from_thread_id*/ None)).await } - async fn start_thread_with_options_and_fork_source( + async fn start_thread_inner( &self, options: StartThreadOptions, forked_from_thread_id: Option, ) -> CodexResult { + let environments = options.environments.unwrap_or_else(|| { + default_thread_environment_selections( + self.state.environment_manager.as_ref(), + &options.config.cwd, + &options.config.workspace_roots, + ) + }); let agent_control = self.agent_control_for_config(&options.config); let (resumed_session_source, resumed_thread_source) = options .initial_history @@ -731,7 +720,7 @@ impl ThreadManager { /*inherited_environments*/ None, /*inherited_exec_policy*/ None, options.parent_trace, - options.environments, + environments, options.thread_extension_init, options.supports_openai_form_elicitation, /*user_shell_override*/ None, @@ -772,7 +761,7 @@ impl ThreadManager { inherited_multi_agent_version, ), ); - self.start_thread_with_options_and_fork_source(options, Some(forked_from_thread_id)) + self.start_thread_inner(options, Some(forked_from_thread_id)) .await } diff --git a/codex-rs/core/src/thread_manager_tests.rs b/codex-rs/core/src/thread_manager_tests.rs index 2c41fd34e290..c08a7ab515bd 100644 --- a/codex-rs/core/src/thread_manager_tests.rs +++ b/codex-rs/core/src/thread_manager_tests.rs @@ -398,12 +398,12 @@ async fn shutdown_all_threads_bounded_submits_shutdown_to_every_thread() { Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), ); let thread_1 = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start first thread") .thread_id; let thread_2 = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start second thread") .thread_id; @@ -435,11 +435,11 @@ async fn code_mode_session_provider_is_shared_across_threads() { Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), ); let first = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start first thread"); let second = manager - .start_thread(config) + .start_thread(StartThreadOptions::new(config)) .await .expect("start second thread"); @@ -491,21 +491,12 @@ async fn start_thread_keeps_internal_threads_hidden_from_normal_lookups() { Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), ); let thread = manager - .start_thread_with_options(StartThreadOptions { - config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, + .start_thread(StartThreadOptions { session_source: Some(SessionSource::Internal( InternalSessionSource::MemoryConsolidation, )), - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(config) }) .await .expect("internal thread should start"); @@ -644,36 +635,19 @@ async fn start_thread_seeds_extension_data_for_mcp_and_lifecycle_contributors() }; let first_thread = manager - .start_thread_with_options(StartThreadOptions { - config: config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), + .start_thread(StartThreadOptions { metrics_service_name: Some("codex_work_desktop".to_string()), - parent_trace: None, - environments: Vec::new(), + environments: Some(Vec::new()), thread_extension_init: selected_root_init("selected-a", "env-a"), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(config.clone()) }) .await .expect("start first thread"); let second_thread = manager - .start_thread_with_options(StartThreadOptions { - config: config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), + .start_thread(StartThreadOptions { + environments: Some(Vec::new()), thread_extension_init: selected_root_init("selected-b", "env-b"), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(config.clone()) }) .await .expect("start second thread"); @@ -783,9 +757,7 @@ async fn selected_capability_roots_round_trip_through_fork() { }, }]; let inherited = manager - .start_thread_with_options(StartThreadOptions { - config, - allow_provider_model_fallback: false, + .start_thread(StartThreadOptions { initial_history: InitialHistory::Forked(vec![RolloutItem::SessionMeta( SessionMetaLine { meta: SessionMeta { @@ -795,15 +767,8 @@ async fn selected_capability_roots_round_trip_through_fork() { git: None, }, )]), - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(config) }) .await .expect("start inherited fork"); @@ -866,19 +831,9 @@ async fn resume_and_fork_do_not_restore_thread_environments_from_rollout() { let mut source_config = config.clone(); source_config.cwd = selected_cwd.clone(); let source = manager - .start_thread_with_options(StartThreadOptions { - config: source_config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: environments.clone(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + .start_thread(StartThreadOptions { + environments: Some(environments.clone()), + ..StartThreadOptions::new(source_config) }) .await .expect("start source thread"); @@ -999,7 +954,7 @@ async fn explicit_installation_id_skips_codex_home_file() { ); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread with explicit installation id"); @@ -1042,7 +997,7 @@ async fn resume_active_thread_from_rollout_returns_running_thread() { ); let source = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start source thread"); source.thread.ensure_rollout_materialized().await; @@ -1104,7 +1059,7 @@ async fn resume_stopped_thread_from_rollout_spawns_new_thread() { ); let source = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start source thread"); source.thread.ensure_rollout_materialized().await; @@ -1173,19 +1128,10 @@ async fn resume_stopped_thread_from_rollout_preserves_thread_source() { ); let source = manager - .start_thread_with_options(StartThreadOptions { - config: config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, + .start_thread(StartThreadOptions { thread_source: Some(ThreadSource::User), - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(config.clone()) }) .await .expect("start source thread"); @@ -1314,7 +1260,7 @@ async fn rollout_path_resume_and_fork_read_history_through_thread_store() { ); let source = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start source thread"); source diff --git a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs index 8ed94d8ec22b..5e5884d27b11 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs @@ -1,4 +1,5 @@ use super::*; +use crate::StartThreadOptions; use crate::ThreadManager; use crate::config::AgentRoleConfig; use crate::config::DEFAULT_AGENT_MAX_DEPTH; @@ -316,7 +317,7 @@ async fn spawn_agent_fork_context_rejects_agent_type_override() { let role_name = install_role_with_model_override(&mut turn).await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -350,7 +351,7 @@ async fn multi_agent_v2_spawn_fork_turns_all_rejects_agent_type_override() { let role_name = install_role_with_model_override(&mut turn).await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -436,7 +437,7 @@ async fn spawn_agent_service_tier_override_validates_the_effective_child_model() let (mut session, turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -541,7 +542,7 @@ async fn spawn_agent_service_tier_inheritance_preserves_supported_or_configured_ turn.config = Arc::new(config); let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -582,7 +583,7 @@ async fn spawn_agent_service_tier_inheritance_preserves_supported_or_configured_ turn.config = Arc::new(config); let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -645,7 +646,7 @@ service_tier = "priority" turn.config = Arc::new(config); let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -718,7 +719,7 @@ service_tier = "turbo" turn.config = Arc::new(config); let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -815,7 +816,7 @@ async fn spawn_agent_full_history_fork_accepts_explicit_service_tier() { .await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -869,7 +870,7 @@ async fn multi_agent_v2_full_history_fork_accepts_explicit_service_tier() { set_turn_config(&mut turn, config); let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -922,7 +923,7 @@ async fn multi_agent_v2_spawn_partial_fork_turns_allows_agent_type_override() { let role_name = install_role_with_model_override(&mut turn).await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1006,7 +1007,7 @@ async fn multi_agent_v2_spawn_requires_task_name() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1040,7 +1041,7 @@ async fn multi_agent_v2_spawn_rejects_legacy_items_field() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1100,7 +1101,7 @@ async fn multi_agent_v2_spawn_returns_path_and_send_message_accepts_relative_pat let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1195,7 +1196,7 @@ async fn multi_agent_v2_spawn_rejects_legacy_fork_context() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1235,7 +1236,7 @@ async fn multi_agent_v2_spawn_rejects_invalid_fork_turns_string() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1275,7 +1276,7 @@ async fn multi_agent_v2_spawn_rejects_zero_fork_turns() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1321,7 +1322,7 @@ async fn multi_agent_v2_send_message_accepts_root_target_from_child() { .expect("test config should allow feature update"); set_turn_config(&mut turn, config); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1397,7 +1398,7 @@ async fn multi_agent_v2_followup_task_rejects_root_target_from_child() { .expect("test config should allow feature update"); set_turn_config(&mut turn, config); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1473,7 +1474,7 @@ async fn multi_agent_v2_list_agents_returns_completed_status() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1561,7 +1562,7 @@ async fn multi_agent_v2_list_agents_filters_by_relative_path_prefix() { let _ = config.features.enable(Feature::MultiAgentV2); set_turn_config(&mut turn, config.clone()); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1642,7 +1643,7 @@ async fn multi_agent_v2_list_agents_omits_closed_agents() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1702,7 +1703,7 @@ async fn multi_agent_v2_list_agents_keeps_interrupted_resident_agents() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1774,7 +1775,7 @@ async fn multi_agent_v2_send_message_rejects_legacy_items_field() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1830,7 +1831,7 @@ async fn multi_agent_v2_send_message_rejects_interrupt_parameter() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -1907,7 +1908,7 @@ async fn multi_agent_v2_followup_task_completion_notifies_parent_on_every_turn() let _ = config.features.enable(Feature::MultiAgentV2); set_turn_config(&mut turn, config); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); // Production spawn_agent calls happen after the parent turn has resolved @@ -2060,7 +2061,7 @@ async fn multi_agent_v2_followup_task_rejects_legacy_items_field() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -2113,7 +2114,7 @@ async fn multi_agent_v2_interrupted_turn_does_not_notify_parent() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -2190,7 +2191,7 @@ async fn multi_agent_v2_spawn_omits_agent_id_when_named() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -2229,7 +2230,7 @@ async fn multi_agent_v2_spawn_surfaces_task_name_validation_errors() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -2449,7 +2450,7 @@ async fn multi_agent_v2_spawn_agent_ignores_configured_max_depth() { .enable(Feature::MultiAgentV2) .expect("test config should allow feature update"); let root = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -2574,7 +2575,7 @@ async fn send_input_interrupts_before_prompt() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -2616,7 +2617,7 @@ async fn send_input_accepts_structured_items() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -2712,7 +2713,7 @@ async fn resume_agent_noops_for_active_agent() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -2918,7 +2919,7 @@ async fn multi_agent_v2_wait_agent_accepts_timeout_only_argument() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3270,7 +3271,7 @@ async fn wait_agent_times_out_when_status_is_not_final() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -3313,7 +3314,7 @@ async fn wait_agent_clamps_short_timeouts_to_minimum() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -3351,7 +3352,7 @@ async fn wait_agent_returns_final_status_without_timeout() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -3401,7 +3402,7 @@ async fn multi_agent_v2_wait_agent_returns_summary_for_mailbox_activity() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3491,7 +3492,7 @@ async fn multi_agent_v2_wait_agent_returns_for_already_queued_mail() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3572,7 +3573,7 @@ async fn multi_agent_v2_wait_agent_wakes_on_any_mailbox_notification() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3663,7 +3664,7 @@ async fn multi_agent_v2_wait_agent_does_not_return_completed_content() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3752,7 +3753,7 @@ async fn multi_agent_v2_interrupt_agent_accepts_task_name_target() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3878,7 +3879,7 @@ async fn multi_agent_v2_interrupt_agent_accepts_unloaded_task_name_target() { Some(state_db.clone()), ); let root = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -3969,7 +3970,7 @@ async fn multi_agent_v2_interrupt_agent_rejects_root_target_and_id() { let (mut session, mut turn) = make_session_and_context().await; let manager = thread_manager(); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -4025,7 +4026,7 @@ async fn multi_agent_v2_interrupt_agent_rejects_self_target_by_id() { .expect("test config should allow feature update"); set_turn_config(&mut turn, config); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -4092,7 +4093,7 @@ async fn multi_agent_v2_interrupt_agent_rejects_self_target_by_task_name() { .expect("test config should allow feature update"); set_turn_config(&mut turn, config); let root = manager - .start_thread((*turn.config).clone()) + .start_thread(StartThreadOptions::new((*turn.config).clone())) .await .expect("root thread should start"); session.services.agent_control = manager.agent_control(); @@ -4155,7 +4156,7 @@ async fn close_agent_submits_shutdown_and_returns_previous_status() { session.services.agent_control = manager.agent_control(); let config = turn.config.as_ref().clone(); let thread = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("start thread"); let agent_id = thread.thread_id; @@ -4216,7 +4217,7 @@ async fn tool_handlers_cascade_close_and_resume_and_keep_explicitly_closed_subtr ); let parent = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("parent thread should start"); let parent_thread_id = parent.thread_id; @@ -4348,7 +4349,7 @@ async fn tool_handlers_cascade_close_and_resume_and_keep_explicitly_closed_subtr ); let operator = manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("operator thread should start"); let operator_session = operator.thread.session.clone(); diff --git a/codex-rs/core/tests/common/test_codex.rs b/codex-rs/core/tests/common/test_codex.rs index d9e00df96a88..e85128e6cd7a 100644 --- a/codex-rs/core/tests/common/test_codex.rs +++ b/codex-rs/core/tests/common/test_codex.rs @@ -41,7 +41,6 @@ use codex_protocol::openai_models::ModelInfo; use codex_protocol::openai_models::ModelsResponse; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::Op; use codex_protocol::protocol::RealtimeConversationVersion as RealtimeWsVersion; use codex_protocol::protocol::SandboxPolicy; @@ -696,24 +695,11 @@ impl TestCodexBuilder { .await? } (None, None) => { - let environments = thread_manager - .default_environment_selections(&config.cwd, &config.workspace_roots); - Box::pin( - thread_manager.start_thread_with_options(StartThreadOptions { - config: config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: self.history_mode, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: self.supports_openai_form_elicitation, - }), - ) + Box::pin(thread_manager.start_thread(StartThreadOptions { + history_mode: self.history_mode, + supports_openai_form_elicitation: self.supports_openai_form_elicitation, + ..StartThreadOptions::new(config.clone()) + })) .await? } }; diff --git a/codex-rs/core/tests/suite/agents_md.rs b/codex-rs/core/tests/suite/agents_md.rs index 65c377e32134..188be19ed101 100644 --- a/codex-rs/core/tests/suite/agents_md.rs +++ b/codex-rs/core/tests/suite/agents_md.rs @@ -8,7 +8,6 @@ use codex_exec_server::REMOTE_ENVIRONMENT_ID; use codex_features::Feature; use codex_home::CodexHomeUserInstructionsProvider; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::Op; use codex_protocol::protocol::RolloutItem; use codex_protocol::protocol::RolloutLine; @@ -484,19 +483,9 @@ async fn loads_user_instructions_without_a_primary_environment() -> Result<()> { let no_environment_thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: test.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + .start_thread(StartThreadOptions { + environments: Some(Vec::new()), + ..StartThreadOptions::new(test.config.clone()) }) .await?; assert_eq!(provider.load_count(), 2); @@ -691,17 +680,8 @@ async fn multi_environment_thread_loads_every_project_and_keeps_creation_snapsho let remote_source = test.config.cwd.join(GLOBAL_AGENTS_FILENAME); let thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: test.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: vec![ + .start_thread(StartThreadOptions { + environments: Some(vec![ TurnEnvironmentSelection { environment_id: REMOTE_ENVIRONMENT_ID.to_string(), cwd: PathUri::from_abs_path(&test.config.cwd), @@ -712,9 +692,8 @@ async fn multi_environment_thread_loads_every_project_and_keeps_creation_snapsho cwd: PathUri::from_host_native_path(local_root.path())?, workspace_roots: vec![PathUri::from_host_native_path(local_root.path())?], }, - ], - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ]), + ..StartThreadOptions::new(test.config.clone()) }) .await?; assert_eq!(provider.load_count(), 2); diff --git a/codex-rs/core/tests/suite/audio_truncation.rs b/codex-rs/core/tests/suite/audio_truncation.rs index 747b65fde340..79732b0230bd 100644 --- a/codex-rs/core/tests/suite/audio_truncation.rs +++ b/codex-rs/core/tests/suite/audio_truncation.rs @@ -1,6 +1,7 @@ use anyhow::Result; use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; +use codex_core::StartThreadOptions; use codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem; use codex_protocol::dynamic_tools::DynamicToolFunctionSpec; use codex_protocol::dynamic_tools::DynamicToolNamespaceSpec; @@ -94,7 +95,10 @@ async fn dynamic_tool_audio_exceeding_the_output_budget_is_omitted() -> Result<( }); let new_thread = base_test .thread_manager - .start_thread_with_tools(base_test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let mut test = base_test; test.codex = new_thread.thread; diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 479f4b97841d..bc0d7a9f79cb 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -4,6 +4,7 @@ use codex_core::ModelClient; use codex_core::NewThread; use codex_core::Prompt; use codex_core::ResponseEvent; +use codex_core::StartThreadOptions; use codex_core::ThreadManager; use codex_core::resolve_installation_id; use codex_core::thread_store_from_config; @@ -1741,7 +1742,7 @@ async fn prefers_apikey_when_config_prefers_apikey_even_with_chatgpt_tokens() { /*external_time_provider*/ None, ); let NewThread { thread: codex, .. } = thread_manager - .start_thread(config.clone()) + .start_thread(StartThreadOptions::new(config.clone())) .await .expect("create new conversation"); diff --git a/codex-rs/core/tests/suite/code_mode.rs b/codex-rs/core/tests/suite/code_mode.rs index 048c15c95e8c..0232c7205106 100644 --- a/codex-rs/core/tests/suite/code_mode.rs +++ b/codex-rs/core/tests/suite/code_mode.rs @@ -5,6 +5,7 @@ use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use codex_config::types::McpServerConfig; use codex_config::types::McpServerTransportConfig; +use codex_core::StartThreadOptions; use codex_core::config::Config; use codex_core::config::CurrentTimeReminderConfig; use codex_extension_api::ExtensionRegistryBuilder; @@ -3628,9 +3629,8 @@ async fn code_mode_can_call_hidden_dynamic_tools() -> Result<()> { let base_test = builder.build(&server).await?; let new_thread = base_test .thread_manager - .start_thread_with_tools( - base_test.config.clone(), - vec![DynamicToolSpec::Namespace(DynamicToolNamespaceSpec { + .start_thread(StartThreadOptions { + dynamic_tools: vec![DynamicToolSpec::Namespace(DynamicToolNamespaceSpec { name: "codex_app".to_string(), description: "Codex app tools.".to_string(), tools: vec![DynamicToolNamespaceTool::Function( @@ -3649,7 +3649,8 @@ async fn code_mode_can_call_hidden_dynamic_tools() -> Result<()> { }, )], })], - ) + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let mut test = base_test; test.codex = new_thread.thread; @@ -3797,9 +3798,8 @@ async fn code_mode_excludes_configured_nested_tool_namespaces() -> Result<()> { let base_test = builder.build(&server).await?; let new_thread = base_test .thread_manager - .start_thread_with_tools( - base_test.config.clone(), - vec![DynamicToolSpec::Namespace(DynamicToolNamespaceSpec { + .start_thread(StartThreadOptions { + dynamic_tools: vec![DynamicToolSpec::Namespace(DynamicToolNamespaceSpec { name: "excluded".to_string(), description: "Excluded tools.".to_string(), tools: vec![DynamicToolNamespaceTool::Function( @@ -3815,7 +3815,8 @@ async fn code_mode_excludes_configured_nested_tool_namespaces() -> Result<()> { }, )], })], - ) + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let mut test = base_test; test.codex = new_thread.thread; diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index 222d5bc4b373..3f0993674513 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -4,6 +4,7 @@ use std::fs; use anyhow::Result; use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; +use codex_core::StartThreadOptions; use codex_core::compact::SUMMARY_PREFIX; use codex_features::Feature; use codex_login::CodexAuth; @@ -1287,7 +1288,10 @@ async fn remote_compact_filters_deferred_dynamic_tools() -> Result<()> { })]; let new_thread = test .thread_manager - .start_thread_with_tools(test.config.clone(), dynamic_tools) + .start_thread(StartThreadOptions { + dynamic_tools, + ..StartThreadOptions::new(test.config.clone()) + }) .await?; test.codex = new_thread.thread; test.session_configured = new_thread.session_configured; @@ -1403,7 +1407,10 @@ async fn remote_compact_does_not_charge_inline_audio_payload_as_text() -> Result }); let new_thread = test .thread_manager - .start_thread_with_tools(test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(test.config.clone()) + }) .await?; test.codex = new_thread.thread; test.session_configured = new_thread.session_configured; @@ -2045,7 +2052,10 @@ async fn remote_compact_trims_tool_search_output_to_empty_tools_array() -> Resul let mut test = builder.build(&server).await?; let new_thread = test .thread_manager - .start_thread_with_tools(test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(test.config.clone()) + }) .await?; test.codex = new_thread.thread; test.session_configured = new_thread.session_configured; diff --git a/codex-rs/core/tests/suite/git_enrichment.rs b/codex-rs/core/tests/suite/git_enrichment.rs index f1b6fb375899..e0e2c27fe687 100644 --- a/codex-rs/core/tests/suite/git_enrichment.rs +++ b/codex-rs/core/tests/suite/git_enrichment.rs @@ -11,7 +11,6 @@ use codex_protocol::config_types::ApprovalsReviewer; #[cfg(not(target_os = "windows"))] use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::Op; use codex_protocol::protocol::ThreadSource; use codex_protocol::user_input::UserInput; @@ -257,24 +256,11 @@ async fn ephemeral_system_thread_prewarm_skips_and_turn_observes_fresh_state() - let mut config = test.config.clone(); config.ephemeral = true; - let environments = test - .thread_manager - .default_environment_selections(&config.cwd, &config.workspace_roots); let system_thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, + .start_thread(StartThreadOptions { thread_source: Some(ThreadSource::Feature("system".to_string())), - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(config) }) .await?; let prewarm = tokio::time::timeout( @@ -386,47 +372,21 @@ async fn concurrent_turns_keep_distinct_worktree_and_repository_metadata() -> Re let mut worktree_config = test.config.clone(); worktree_config.cwd = worktree.abs(); - let worktree_environments = test - .thread_manager - .default_environment_selections(&worktree_config.cwd, &worktree_config.workspace_roots); let worktree_thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: worktree_config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, + .start_thread(StartThreadOptions { thread_source: Some(ThreadSource::User), - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: worktree_environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(worktree_config) }) .await?; let mut other_config = test.config.clone(); other_config.cwd = other_repo.path().to_path_buf().abs(); - let other_environments = test - .thread_manager - .default_environment_selections(&other_config.cwd, &other_config.workspace_roots); let other_thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: other_config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, + .start_thread(StartThreadOptions { thread_source: Some(ThreadSource::User), - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: other_environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(other_config) }) .await?; diff --git a/codex-rs/core/tests/suite/hooks.rs b/codex-rs/core/tests/suite/hooks.rs index 62c2be000bf0..7f7d87b8a4b7 100644 --- a/codex-rs/core/tests/suite/hooks.rs +++ b/codex-rs/core/tests/suite/hooks.rs @@ -19,7 +19,6 @@ use codex_protocol::models::ResponseItem; use codex_protocol::permissions::NetworkSandboxPolicy; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::Op; use codex_protocol::protocol::RolloutItem; use codex_protocol::protocol::RolloutLine; @@ -1415,19 +1414,10 @@ async fn session_end_skips_subagents() -> Result<()> { ] { let subagent = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: test.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, + .start_thread(StartThreadOptions { session_source: Some(SessionSource::SubAgent(source)), - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(test.config.clone()) }) .await?; diff --git a/codex-rs/core/tests/suite/mcp_tool_cache.rs b/codex-rs/core/tests/suite/mcp_tool_cache.rs index 63e7db47ba61..5a5766fe1273 100644 --- a/codex-rs/core/tests/suite/mcp_tool_cache.rs +++ b/codex-rs/core/tests/suite/mcp_tool_cache.rs @@ -3,6 +3,7 @@ use std::time::Duration; use anyhow::Context; use codex_core::NewThread; +use codex_core::StartThreadOptions; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::RemoveOptions; use codex_protocol::models::PermissionProfile; @@ -183,7 +184,7 @@ async fn regular_mcp_definition_cache_preserves_live_session_state() -> anyhow:: .. } = fixture .thread_manager - .start_thread(fixture.config.clone()) + .start_thread(StartThreadOptions::new(fixture.config.clone())) .await?; let second_pid = wait_for_new_pid(fs.as_ref(), &pid_file, Some(&first_pid)).await?; let second_process = process_label(&second_pid); diff --git a/codex-rs/core/tests/suite/search_tool.rs b/codex-rs/core/tests/suite/search_tool.rs index 50ccea0d8764..3143ca6f55f1 100644 --- a/codex-rs/core/tests/suite/search_tool.rs +++ b/codex-rs/core/tests/suite/search_tool.rs @@ -4,6 +4,7 @@ use anyhow::Result; use codex_config::types::McpServerConfig; use codex_config::types::McpServerTransportConfig; +use codex_core::StartThreadOptions; use codex_login::CodexAuth; use codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem; use codex_protocol::dynamic_tools::DynamicToolFunctionSpec; @@ -939,7 +940,10 @@ async fn tool_search_returns_deferred_dynamic_tool_and_routes_follow_up_call() - let base_test = builder.build(&server).await?; let new_thread = base_test .thread_manager - .start_thread_with_tools(base_test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let mut test = base_test; test.codex = new_thread.thread; @@ -1571,7 +1575,10 @@ async fn tool_search_matches_dynamic_tools_by_name_description_namespace_and_sch let base_test = builder.build(&server).await?; let new_thread = base_test .thread_manager - .start_thread_with_tools(base_test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let mut test = base_test; test.codex = new_thread.thread; diff --git a/codex-rs/core/tests/suite/sqlite_state.rs b/codex-rs/core/tests/suite/sqlite_state.rs index e1c055d32108..b29ac524b106 100644 --- a/codex-rs/core/tests/suite/sqlite_state.rs +++ b/codex-rs/core/tests/suite/sqlite_state.rs @@ -1,6 +1,7 @@ use anyhow::Result; use codex_config::types::McpServerConfig; use codex_config::types::McpServerTransportConfig; +use codex_core::StartThreadOptions; use codex_core::config::Config; use codex_extension_api::ExtensionRegistryBuilder; use codex_features::Feature; @@ -154,7 +155,10 @@ async fn resume_restores_dynamic_tools_from_rollout_with_sqlite_enabled() -> Res let base_test = builder.build(&server).await?; let started = base_test .thread_manager - .start_thread_with_tools(base_test.config.clone(), vec![dynamic_tool]) + .start_thread(StartThreadOptions { + dynamic_tools: vec![dynamic_tool], + ..StartThreadOptions::new(base_test.config.clone()) + }) .await?; let rollout_path = started .session_configured @@ -251,7 +255,7 @@ async fn resume_restores_legacy_dynamic_tools_from_rollout_with_sqlite_enabled() let base_test = builder.build(&server).await?; let started = base_test .thread_manager - .start_thread_with_tools(base_test.config.clone(), Vec::new()) + .start_thread(StartThreadOptions::new(base_test.config.clone())) .await?; let rollout_path = started .session_configured diff --git a/codex-rs/core/tests/suite/subagent_notifications.rs b/codex-rs/core/tests/suite/subagent_notifications.rs index 3b2638788f79..4e51ed2cf038 100644 --- a/codex-rs/core/tests/suite/subagent_notifications.rs +++ b/codex-rs/core/tests/suite/subagent_notifications.rs @@ -10,7 +10,6 @@ use codex_protocol::models::PermissionProfile; use codex_protocol::openai_models::ReasoningEffort; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::Op; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; @@ -797,19 +796,10 @@ async fn subagent_stop_replaces_stop_and_skips_internal_subagents() -> Result<() // because the SubagentStop hook above intentionally matches all agent types. let internal_thread = test .thread_manager - .start_thread_with_options(StartThreadOptions { - config: test.config.clone(), - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, + .start_thread(StartThreadOptions { session_source: Some(SessionSource::SubAgent(SubAgentSource::Review)), - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments: Vec::new(), - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + environments: Some(Vec::new()), + ..StartThreadOptions::new(test.config.clone()) }) .await?; diff --git a/codex-rs/ext/agent/src/lib.rs b/codex-rs/ext/agent/src/lib.rs index 01e068545ff7..349fb46c09b7 100644 --- a/codex-rs/ext/agent/src/lib.rs +++ b/codex-rs/ext/agent/src/lib.rs @@ -6,7 +6,6 @@ use codex_core::config::Config; use codex_protocol::ThreadId; use codex_protocol::error::CodexErr; use codex_protocol::error::Result as CodexResult; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::W3cTraceContext; use codex_protocol::user_input::UserInput; use std::sync::Arc; @@ -61,26 +60,14 @@ impl AgentRunner { .thread_manager .upgrade() .ok_or_else(|| CodexErr::UnsupportedOperation("thread manager dropped".to_string()))?; - let environments = - thread_manager.default_environment_selections(&config.cwd, &config.workspace_roots); let NewThread { thread_id, thread, .. } = thread_manager .spawn_subagent( parent_thread_id, StartThreadOptions { - config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, parent_trace: parent_trace.clone(), - environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(config) }, ) .await?; diff --git a/codex-rs/mcp-server/src/codex_tool_runner.rs b/codex-rs/mcp-server/src/codex_tool_runner.rs index 64acffd464dd..4a166226c9c8 100644 --- a/codex-rs/mcp-server/src/codex_tool_runner.rs +++ b/codex-rs/mcp-server/src/codex_tool_runner.rs @@ -11,6 +11,7 @@ use crate::outgoing_message::OutgoingNotificationMeta; use crate::patch_approval::handle_patch_approval_request; use codex_core::CodexThread; use codex_core::NewThread; +use codex_core::StartThreadOptions; use codex_core::ThreadManager; use codex_core::config::Config as CodexConfig; use codex_protocol::ThreadId; @@ -66,7 +67,10 @@ pub async fn run_codex_tool_session( thread_id, thread, session_configured, - } = match thread_manager.start_thread(config.clone()).await { + } = match thread_manager + .start_thread(StartThreadOptions::new(config.clone())) + .await + { Ok(res) => res, Err(e) => { let result = CallToolResult::error(vec![Content::text(format!( diff --git a/codex-rs/memories/write/src/runtime.rs b/codex-rs/memories/write/src/runtime.rs index 4e6f4c8cd9d6..be9b99186cfa 100644 --- a/codex-rs/memories/write/src/runtime.rs +++ b/codex-rs/memories/write/src/runtime.rs @@ -25,7 +25,6 @@ use codex_protocol::ThreadId; use codex_protocol::config_types::ReasoningSummary; use codex_protocol::openai_models::ModelInfo; use codex_protocol::openai_models::ReasoningEffort; -use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::InternalSessionSource; use codex_protocol::protocol::Op; use codex_protocol::protocol::SessionSource; @@ -321,28 +320,16 @@ impl MemoryStartupContext { config: Config, prompt: Vec, ) -> anyhow::Result { - let environments = self - .thread_manager - .default_environment_selections(&config.cwd, &config.workspace_roots); let NewThread { thread_id, thread, .. } = self .thread_manager - .start_thread_with_options(StartThreadOptions { - config, - allow_provider_model_fallback: false, - initial_history: InitialHistory::New, - history_mode: None, + .start_thread(StartThreadOptions { session_source: Some(SessionSource::Internal( InternalSessionSource::MemoryConsolidation, )), thread_source: Some(ThreadSource::MemoryConsolidation), - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments, - thread_extension_init: Default::default(), - supports_openai_form_elicitation: false, + ..StartThreadOptions::new(config) }) .await?; diff --git a/codex-rs/thread-manager-sample/src/main.rs b/codex-rs/thread-manager-sample/src/main.rs index 0ab63a271b4c..9e86a5762591 100644 --- a/codex-rs/thread-manager-sample/src/main.rs +++ b/codex-rs/thread-manager-sample/src/main.rs @@ -45,6 +45,7 @@ use codex_core_api::RealtimeAudioConfig; use codex_core_api::RealtimeConfig; use codex_core_api::SessionPickerViewMode; use codex_core_api::SessionSource; +use codex_core_api::StartThreadOptions; use codex_core_api::TerminalResizeReflowConfig; use codex_core_api::ThreadManager; use codex_core_api::ThreadStoreConfig; @@ -152,7 +153,7 @@ async fn run_main(arg0_paths: Arg0DispatchPaths) -> anyhow::Result<()> { let NewThread { thread_id, thread, .. } = thread_manager - .start_thread(config) + .start_thread(StartThreadOptions::new(config)) .await .context("start Codex thread")?;