[codex] wire process-owned code mode host into core (#30142)

## Summary

- add the `code_mode_host` feature flag and select
`ProcessOwnedCodeModeSessionProvider` in `CodeModeService` when enabled
- initialize code-mode sessions lazily so a missing host reports a tool
error without failing thread startup
- resolve `codex-code-mode-host` beside the running Codex binary by
default while preserving `CODEX_CODE_MODE_HOST_PATH` as an override
- add unit and end-to-end coverage for host resolution and graceful
missing-host behavior

## Why

This wires the process-owned session client from #30112 into the core
service behind an opt-in rollout gate. Packaged Codex installations can
place the helper in the same `bin` directory as the main executable
without relying on `PATH`, while development and custom installations
can continue to override the helper path.

## Stack

- Depends on #30112
- Base branch: `cconger/process-owned-session-runtime-4-client`

## Validation

Build `codex` and `codex-code-mode-host`
`CODEX_CODE_MODE_HOST_PATH="$PWD/target/debug/codex-code-mode-host"
./target/debug/codex --enable code_mode_host`
This commit is contained in:
Channing Conger
2026-06-26 00:23:33 -07:00
committed by GitHub
Unverified
parent ab16046c88
commit 7d8906b478
13 changed files with 275 additions and 27 deletions
+15 -6
View File
@@ -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<OsString>,
current_exe: io::Result<PathBuf>,
) -> 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)
}
@@ -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(
+6
View File
@@ -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"
},
+1
View File
@@ -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()),
+3
View File
@@ -416,6 +416,7 @@ pub(crate) struct CodexSpawnArgs {
pub(crate) skills_service: Arc<SkillsService>,
pub(crate) plugins_manager: Arc<PluginsManager>,
pub(crate) mcp_manager: Arc<McpManager>,
pub(crate) code_mode_session_provider: Arc<dyn codex_code_mode::CodeModeSessionProvider>,
pub(crate) extensions: Arc<codex_extension_api::ExtensionRegistry<crate::config::Config>>,
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,
+4 -1
View File
@@ -489,6 +489,7 @@ impl Session {
skills_service: Arc<SkillsService>,
plugins_manager: Arc<PluginsManager>,
mcp_manager: Arc<McpManager>,
code_mode_session_provider: Arc<dyn codex_code_mode::CodeModeSessionProvider>,
extensions: Arc<codex_extension_api::ExtensionRegistry<crate::config::Config>>,
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),
};
+10 -2
View File
@@ -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 {
@@ -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(
+11
View File
@@ -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<SkillsService>,
plugins_manager: Arc<PluginsManager>,
mcp_manager: Arc<McpManager>,
code_mode_session_provider: Arc<dyn CodeModeSessionProvider>,
extensions: Arc<ExtensionRegistry<Config>>,
user_instructions_provider: Arc<dyn UserInstructionsProvider>,
thread_store: Arc<dyn ThreadStore>,
@@ -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,
+58
View File
@@ -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");
+81 -18
View File
@@ -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<Arc<dyn CodeModeSession>>,
session: OnceCell<Arc<dyn CodeModeSession>>,
session_provider: Arc<dyn CodeModeSessionProvider>,
dispatch_broker: Arc<CodeModeDispatchBroker>,
shutting_down: AtomicBool,
}
impl CodeModeService {
pub(crate) fn new() -> Self {
pub(crate) fn new(session_provider: Arc<dyn CodeModeSessionProvider>) -> 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<dyn CodeModeSessionProvider> {
Arc::clone(&self.session_provider)
}
pub(crate) async fn execute(
&self,
request: codex_code_mode::ExecuteRequest,
) -> Result<codex_code_mode::StartedCell, String> {
self.session()?.execute(request).await
self.session().await?.execute(request).await
}
pub(crate) async fn wait(
&self,
request: codex_code_mode::WaitRequest,
) -> Result<codex_code_mode::WaitOutcome, String> {
self.session()?.wait(request).await
self.session().await?.wait(request).await
}
pub(crate) async fn terminate(
&self,
cell_id: CellId,
) -> Result<codex_code_mode::WaitOutcome, String> {
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::<Arc<dyn CodeModeSession>, String>(
"code mode session is shutting down".to_string(),
)
})
.await
{
Ok(session) => session.shutdown().await,
Err(_) => Ok(()),
}
}
@@ -121,9 +140,7 @@ impl CodeModeService {
) -> Option<CodeModeDispatchWorker> {
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<dyn CodeModeSession>, String> {
async fn session(&self) -> Result<Arc<dyn CodeModeSession>, 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");
}
}
+26
View File
@@ -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
+8
View File
@@ -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",