diff --git a/codex-rs/code-mode/src/remote_session.rs b/codex-rs/code-mode/src/remote_session.rs index d942b377c..af7d27845 100644 --- a/codex-rs/code-mode/src/remote_session.rs +++ b/codex-rs/code-mode/src/remote_session.rs @@ -1,3 +1,5 @@ +use std::ffi::OsString; +use std::io; use std::path::PathBuf; use std::sync::Arc; use std::sync::Mutex as StdMutex; @@ -472,7 +474,17 @@ impl CodeModeSession for ProcessOwnedCodeModeSession { } fn default_host_program() -> PathBuf { - if let Some(path) = std::env::var_os(CODE_MODE_HOST_PATH_ENV) { + resolve_host_program( + std::env::var_os(CODE_MODE_HOST_PATH_ENV), + std::env::current_exe(), + ) +} + +fn resolve_host_program( + override_path: Option, + current_exe: io::Result, +) -> PathBuf { + if let Some(path) = override_path { return PathBuf::from(path); } let executable_name = if cfg!(windows) { @@ -480,13 +492,10 @@ fn default_host_program() -> PathBuf { } else { "codex-code-mode-host" }; - if let Ok(current_exe) = std::env::current_exe() + if let Ok(current_exe) = current_exe && let Some(parent) = current_exe.parent() { - let sibling = parent.join(executable_name); - if sibling.is_file() { - return sibling; - } + return parent.join(executable_name); } PathBuf::from(executable_name) } diff --git a/codex-rs/code-mode/src/remote_session_tests.rs b/codex-rs/code-mode/src/remote_session_tests.rs index d97120755..03d75e523 100644 --- a/codex-rs/code-mode/src/remote_session_tests.rs +++ b/codex-rs/code-mode/src/remote_session_tests.rs @@ -1,9 +1,12 @@ +use std::io; +use std::path::PathBuf; use std::sync::Arc; use codex_code_mode_protocol::CodeModeSessionProvider; use super::ProcessOwnedCodeModeSession; use super::ProcessOwnedCodeModeSessionProvider; +use super::resolve_host_program; use crate::NoopCodeModeSessionDelegate; #[test] @@ -16,6 +19,54 @@ fn provider_reuses_its_live_process_host() { assert!(Arc::ptr_eq(&first, &second)); } +#[test] +fn host_program_override_takes_precedence() { + assert_eq!( + resolve_host_program( + Some("custom-code-mode-host".into()), + Ok(PathBuf::from("/opt/codex/bin/codex")), + ), + PathBuf::from("custom-code-mode-host") + ); +} + +#[test] +fn host_program_is_next_to_the_main_executable_even_when_missing() { + let executable_name = if cfg!(windows) { + "codex-code-mode-host.exe" + } else { + "codex-code-mode-host" + }; + + assert_eq!( + resolve_host_program( + /*override_path*/ None, + Ok(PathBuf::from("/opt/codex/bin/codex")), + ), + PathBuf::from("/opt/codex/bin").join(executable_name) + ); +} + +#[test] +fn host_program_falls_back_to_its_name_when_main_executable_is_unknown() { + let executable_name = if cfg!(windows) { + "codex-code-mode-host.exe" + } else { + "codex-code-mode-host" + }; + + assert_eq!( + resolve_host_program( + /*override_path*/ None, + Err(io::Error::new( + io::ErrorKind::NotFound, + "missing executable" + )), + ), + PathBuf::from(executable_name) + ); +} + #[tokio::test] async fn provider_reports_host_spawn_failure() { let provider = ProcessOwnedCodeModeSessionProvider::with_host_program( diff --git a/codex-rs/core/config.schema.json b/codex-rs/core/config.schema.json index 512c050a2..ab787ba34 100644 --- a/codex-rs/core/config.schema.json +++ b/codex-rs/core/config.schema.json @@ -448,6 +448,9 @@ "code_mode": { "$ref": "#/definitions/FeatureToml_for_CodeModeConfigToml" }, + "code_mode_host": { + "type": "boolean" + }, "code_mode_only": { "type": "boolean" }, @@ -4811,6 +4814,9 @@ "code_mode": { "$ref": "#/definitions/FeatureToml_for_CodeModeConfigToml" }, + "code_mode_host": { + "type": "boolean" + }, "code_mode_only": { "type": "boolean" }, diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index 8874aeae7..d25af5793 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -105,6 +105,7 @@ pub(crate) async fn run_codex_thread_interactive( skills_service: Arc::clone(&parent_session.services.skills_service), plugins_manager: Arc::clone(&parent_session.services.plugins_manager), mcp_manager: Arc::clone(&parent_session.services.mcp_manager), + code_mode_session_provider: parent_session.services.code_mode_service.session_provider(), extensions: Arc::clone(&parent_session.services.extensions), conversation_history, session_source: SessionSource::SubAgent(subagent_source.clone()), diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index 000825754..6a21a6be6 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -416,6 +416,7 @@ pub(crate) struct CodexSpawnArgs { pub(crate) skills_service: Arc, pub(crate) plugins_manager: Arc, pub(crate) mcp_manager: Arc, + pub(crate) code_mode_session_provider: Arc, pub(crate) extensions: Arc>, pub(crate) conversation_history: InitialHistory, pub(crate) session_source: SessionSource, @@ -506,6 +507,7 @@ impl Codex { skills_service, plugins_manager, mcp_manager, + code_mode_session_provider, extensions, conversation_history, session_source, @@ -680,6 +682,7 @@ impl Codex { skills_service, plugins_manager, mcp_manager.clone(), + code_mode_session_provider, extensions, thread_extension_init, supports_openai_form_elicitation, diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index ed07352b1..a9cc89db3 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -489,6 +489,7 @@ impl Session { skills_service: Arc, plugins_manager: Arc, mcp_manager: Arc, + code_mode_session_provider: Arc, extensions: Arc>, mut thread_extension_init: ExtensionDataInit, supports_openai_form_elicitation: bool, @@ -1112,7 +1113,9 @@ impl Session { session_configuration.parent_thread_id, ), ), - code_mode_service: crate::tools::code_mode::CodeModeService::new(), + code_mode_service: crate::tools::code_mode::CodeModeService::new(Arc::clone( + &code_mode_session_provider, + )), tool_search_handler_cache: Default::default(), turn_environments: Arc::clone(&turn_environments), }; diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index 2bc63b404..09fdcf224 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -5225,6 +5225,7 @@ async fn session_new_fails_when_zsh_fork_enabled_without_packaged_zsh() { skills_service, plugins_manager, mcp_manager, + Arc::new(codex_code_mode::InProcessCodeModeSessionProvider), Arc::new(codex_extension_api::ExtensionRegistryBuilder::new().build()), codex_extension_api::ExtensionDataInit::default(), /*supports_openai_form_elicitation*/ false, @@ -5438,7 +5439,9 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { /*item_ids_enabled*/ config.features.enabled(Feature::ItemIds), /*attestation_provider*/ None, ), - code_mode_service: crate::tools::code_mode::CodeModeService::new(), + code_mode_service: crate::tools::code_mode::CodeModeService::new(Arc::new( + codex_code_mode::InProcessCodeModeSessionProvider, + )), tool_search_handler_cache: Default::default(), turn_environments: Arc::clone(&turn_environments), }; @@ -5602,6 +5605,7 @@ async fn make_session_with_config_and_rx( skills_service, plugins_manager, mcp_manager, + Arc::new(codex_code_mode::InProcessCodeModeSessionProvider), Arc::new(codex_extension_api::ExtensionRegistryBuilder::new().build()), codex_extension_api::ExtensionDataInit::default(), /*supports_openai_form_elicitation*/ false, @@ -5708,6 +5712,7 @@ async fn make_session_with_history_source_and_agent_control_and_rx( skills_service, plugins_manager, mcp_manager, + Arc::new(codex_code_mode::InProcessCodeModeSessionProvider), Arc::new(codex_extension_api::ExtensionRegistryBuilder::new().build()), codex_extension_api::ExtensionDataInit::default(), /*supports_openai_form_elicitation*/ false, @@ -7513,7 +7518,9 @@ where /*item_ids_enabled*/ config.features.enabled(Feature::ItemIds), /*attestation_provider*/ None, ), - code_mode_service: crate::tools::code_mode::CodeModeService::new(), + code_mode_service: crate::tools::code_mode::CodeModeService::new(Arc::new( + codex_code_mode::InProcessCodeModeSessionProvider, + )), tool_search_handler_cache: Default::default(), turn_environments: Arc::clone(&turn_environments), }; @@ -9178,6 +9185,7 @@ async fn attach_in_memory_thread_store( originator: "test_originator".to_string(), base_instructions: BaseInstructions::default(), dynamic_tools: Vec::new(), + selected_capability_roots: Vec::new(), multi_agent_version: None, initial_window_id: Uuid::now_v7().to_string(), metadata: ThreadPersistenceMetadata { diff --git a/codex-rs/core/src/session/tests/guardian_tests.rs b/codex-rs/core/src/session/tests/guardian_tests.rs index 4f88c9ed6..dc17af5a1 100644 --- a/codex-rs/core/src/session/tests/guardian_tests.rs +++ b/codex-rs/core/src/session/tests/guardian_tests.rs @@ -727,6 +727,7 @@ async fn guardian_subagent_does_not_inherit_parent_exec_policy_rules() { skills_service, plugins_manager, mcp_manager, + code_mode_session_provider: Arc::new(codex_code_mode::InProcessCodeModeSessionProvider), extensions: codex_extension_api::empty_extension_registry(), conversation_history: InitialHistory::New, session_source: SessionSource::SubAgent(SubAgentSource::Other( diff --git a/codex-rs/core/src/thread_manager.rs b/codex-rs/core/src/thread_manager.rs index b606bb3f7..2ca21c15a 100644 --- a/codex-rs/core/src/thread_manager.rs +++ b/codex-rs/core/src/thread_manager.rs @@ -21,6 +21,9 @@ use codex_agent_graph_store::LocalAgentGraphStore; use codex_analytics::AnalyticsEventsClient; use codex_app_server_protocol::ThreadHistoryBuilder; use codex_app_server_protocol::TurnStatus; +use codex_code_mode::CodeModeSessionProvider; +use codex_code_mode::InProcessCodeModeSessionProvider; +use codex_code_mode::ProcessOwnedCodeModeSessionProvider; use codex_core_plugins::PluginsManager; use codex_exec_server::EnvironmentManager; use codex_extension_api::ExtensionDataInit; @@ -240,6 +243,7 @@ pub(crate) struct ThreadManagerState { skills_service: Arc, plugins_manager: Arc, mcp_manager: Arc, + code_mode_session_provider: Arc, extensions: Arc>, user_instructions_provider: Arc, thread_store: Arc, @@ -336,6 +340,11 @@ impl ThreadManager { skills_service, plugins_manager, mcp_manager, + code_mode_session_provider: if config.features.enabled(Feature::CodeModeHost) { + Arc::new(ProcessOwnedCodeModeSessionProvider::default()) + } else { + Arc::new(InProcessCodeModeSessionProvider) + }, extensions, user_instructions_provider, thread_store, @@ -441,6 +450,7 @@ impl ThreadManager { skills_service, plugins_manager, mcp_manager, + code_mode_session_provider: Arc::new(InProcessCodeModeSessionProvider), extensions: empty_extension_registry(), user_instructions_provider: Arc::new( crate::test_support::EmptyUserInstructionsProvider, @@ -1562,6 +1572,7 @@ impl ThreadManagerState { skills_service: Arc::clone(&self.skills_service), plugins_manager: Arc::clone(&self.plugins_manager), mcp_manager: Arc::clone(&self.mcp_manager), + code_mode_session_provider: Arc::clone(&self.code_mode_session_provider), extensions: Arc::clone(&self.extensions), conversation_history: initial_history, session_source, diff --git a/codex-rs/core/src/thread_manager_tests.rs b/codex-rs/core/src/thread_manager_tests.rs index 5148f1625..b33d95ac3 100644 --- a/codex-rs/core/src/thread_manager_tests.rs +++ b/codex-rs/core/src/thread_manager_tests.rs @@ -395,6 +395,64 @@ async fn shutdown_all_threads_bounded_submits_shutdown_to_every_thread() { assert!(manager.list_thread_ids().await.is_empty()); } +#[tokio::test] +async fn code_mode_session_provider_is_shared_across_threads() { + let temp_dir = tempdir().expect("tempdir"); + let mut config = test_config().await; + config.codex_home = temp_dir.path().join("codex-home").abs(); + config.cwd = config.codex_home.abs(); + std::fs::create_dir_all(&config.codex_home).expect("create codex home"); + + let manager = ThreadManager::with_models_provider_and_home_for_tests( + CodexAuth::from_api_key("dummy"), + config.model_provider.clone(), + config.codex_home.to_path_buf(), + Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), + ); + let first = manager + .start_thread(config.clone()) + .await + .expect("start first thread"); + let second = manager + .start_thread(config) + .await + .expect("start second thread"); + + let first_provider = first + .thread + .codex + .session + .services + .code_mode_service + .session_provider(); + let second_provider = second + .thread + .codex + .session + .services + .code_mode_service + .session_provider(); + assert!(Arc::ptr_eq(&first_provider, &second_provider)); + assert!(Arc::ptr_eq( + &first_provider, + &manager.state.code_mode_session_provider + )); + + let mut completed = vec![first.thread_id, second.thread_id]; + completed.sort_by_key(std::string::ToString::to_string); + let report = manager + .shutdown_all_threads_bounded(Duration::from_secs(10)) + .await; + assert_eq!( + report, + ThreadShutdownReport { + completed, + submit_failed: Vec::new(), + timed_out: Vec::new(), + } + ); +} + #[tokio::test] async fn start_thread_keeps_internal_threads_hidden_from_normal_lookups() { let temp_dir = tempdir().expect("tempdir"); diff --git a/codex-rs/core/src/tools/code_mode/mod.rs b/codex-rs/core/src/tools/code_mode/mod.rs index c28eef072..ad746a1de 100644 --- a/codex-rs/core/src/tools/code_mode/mod.rs +++ b/codex-rs/core/src/tools/code_mode/mod.rs @@ -6,16 +6,19 @@ mod wait_handler; pub(crate) mod wait_spec; use std::sync::Arc; +use std::sync::atomic::AtomicBool; +use std::sync::atomic::Ordering; use std::time::Duration; use codex_code_mode::CellId; use codex_code_mode::CodeModeNestedToolCall; use codex_code_mode::CodeModeSession; +use codex_code_mode::CodeModeSessionProvider; use codex_code_mode::CodeModeToolKind; -use codex_code_mode::InProcessCodeModeSession; use codex_code_mode::RuntimeResponse; use codex_protocol::models::FunctionCallOutputContentItem; use serde_json::Value as JsonValue; +use tokio::sync::OnceCell; use tokio_util::sync::CancellationToken; use crate::function_tool::FunctionCallError; @@ -61,46 +64,62 @@ pub(crate) struct ExecContext { } pub(crate) struct CodeModeService { - session: Option>, + session: OnceCell>, + session_provider: Arc, dispatch_broker: Arc, + shutting_down: AtomicBool, } impl CodeModeService { - pub(crate) fn new() -> Self { + pub(crate) fn new(session_provider: Arc) -> Self { let dispatch_broker = Arc::new(CodeModeDispatchBroker::new()); Self { - session: Some(Arc::new(InProcessCodeModeSession::with_delegate( - dispatch_broker.clone(), - ))), + session: OnceCell::new(), + session_provider, dispatch_broker, + shutting_down: AtomicBool::new(false), } } + pub(crate) fn session_provider(&self) -> Arc { + Arc::clone(&self.session_provider) + } + pub(crate) async fn execute( &self, request: codex_code_mode::ExecuteRequest, ) -> Result { - self.session()?.execute(request).await + self.session().await?.execute(request).await } pub(crate) async fn wait( &self, request: codex_code_mode::WaitRequest, ) -> Result { - self.session()?.wait(request).await + self.session().await?.wait(request).await } pub(crate) async fn terminate( &self, cell_id: CellId, ) -> Result { - self.session()?.terminate(cell_id).await + self.session().await?.terminate(cell_id).await } pub(crate) async fn shutdown(&self) -> Result<(), String> { - match &self.session { - Some(session) => session.shutdown().await, - None => Ok(()), + self.shutting_down.store(true, Ordering::Release); + // Join any initialization already in progress without initializing an unused service. + match self + .session + .get_or_try_init(|| async { + Err::, String>( + "code mode session is shutting down".to_string(), + ) + }) + .await + { + Ok(session) => session.shutdown().await, + Err(_) => Ok(()), } } @@ -121,9 +140,7 @@ impl CodeModeService { ) -> Option { let turn = &step_context.turn; let tool_mode = effective_tool_mode(turn); - if !matches!(tool_mode, ToolMode::CodeMode | ToolMode::CodeModeOnly) - || self.session.is_none() - { + if !matches!(tool_mode, ToolMode::CodeMode | ToolMode::CodeModeOnly) { return None; } @@ -137,10 +154,27 @@ impl CodeModeService { ) } - fn session(&self) -> Result<&Arc, String> { + async fn session(&self) -> Result, String> { + if self.shutting_down.load(Ordering::Acquire) { + return Err("code mode session is shutting down".to_string()); + } self.session - .as_ref() - .ok_or_else(|| "code mode is unavailable".to_string()) + .get_or_try_init(|| async { + if self.shutting_down.load(Ordering::Acquire) { + return Err("code mode session is shutting down".to_string()); + } + let session = self + .session_provider + .create_session(self.dispatch_broker.clone()) + .await?; + if self.shutting_down.load(Ordering::Acquire) { + let _ = session.shutdown().await; + return Err("code mode session is shutting down".to_string()); + } + Ok(session) + }) + .await + .map(Arc::clone) } } @@ -325,10 +359,15 @@ fn build_freeform_tool_payload( #[cfg(test)] mod tests { + use std::sync::Arc; + + use super::CodeModeService; use super::build_nested_tool_payload; use super::truncate_code_mode_result; use crate::tools::context::ToolPayload; use codex_code_mode::CodeModeToolKind; + use codex_code_mode::ExecuteRequest; + use codex_code_mode::ProcessOwnedCodeModeSessionProvider; use codex_protocol::models::FunctionCallOutputContentItem; use codex_tools::ToolName; use serde_json::json; @@ -385,4 +424,28 @@ mod tests { }] ); } + + #[tokio::test] + async fn missing_process_host_is_reported_without_failing_service_creation() { + let service = CodeModeService::new(Arc::new( + ProcessOwnedCodeModeSessionProvider::with_host_program( + "codex-code-mode-host-does-not-exist".into(), + ), + )); + + let error = service + .execute(ExecuteRequest { + tool_call_id: "call-1".to_string(), + enabled_tools: Vec::new(), + source: "text('unreachable')".to_string(), + yield_time_ms: None, + max_output_tokens: None, + }) + .await + .err() + .expect("missing host should reject execution"); + + assert!(error.contains("failed to spawn code-mode host")); + service.shutdown().await.expect("shutdown unused service"); + } } diff --git a/codex-rs/core/tests/suite/code_mode.rs b/codex-rs/core/tests/suite/code_mode.rs index 92b6c9d47..26c041d8f 100644 --- a/codex-rs/core/tests/suite/code_mode.rs +++ b/codex-rs/core/tests/suite/code_mode.rs @@ -219,6 +219,32 @@ async fn run_code_mode_turn_with_model_and_config( Ok((test, second_mock)) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn missing_process_host_returns_a_tool_error() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = responses::start_mock_server().await; + let (_test, follow_up_mock) = + run_code_mode_turn_with_config(&server, "Run code mode", "text('unreachable')", |config| { + config + .features + .enable(Feature::CodeModeHost) + .expect("code mode host should be enabled"); + }) + .await?; + + let output = follow_up_mock + .single_request() + .custom_tool_call_output("call-1"); + assert!( + output["output"] + .as_str() + .is_some_and(|output| output.contains("failed to spawn code-mode host")) + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn code_mode_can_call_standalone_web_search() -> Result<()> { assert_code_mode_standalone_web_search(WebSearchMode::Live, serde_json::json!(true)).await diff --git a/codex-rs/features/src/lib.rs b/codex-rs/features/src/lib.rs index 3bce1a7d1..5e9fbe528 100644 --- a/codex-rs/features/src/lib.rs +++ b/codex-rs/features/src/lib.rs @@ -92,6 +92,8 @@ pub enum Feature { // Experimental /// Enable JavaScript code mode backed by the in-process V8 runtime. CodeMode, + /// Run JavaScript code mode in the standalone host process. + CodeModeHost, /// Restrict model-visible tools to code mode entrypoints (`exec`, `wait`). CodeModeOnly, /// Use the single unified PTY-backed exec tool. @@ -850,6 +852,12 @@ pub const FEATURES: &[FeatureSpec] = &[ stage: Stage::UnderDevelopment, default_enabled: false, }, + FeatureSpec { + id: Feature::CodeModeHost, + key: "code_mode_host", + stage: Stage::UnderDevelopment, + default_enabled: false, + }, FeatureSpec { id: Feature::CodeModeOnly, key: "code_mode_only",