Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
code-mode: share process host across threads
  • Loading branch information
cconger committed Jun 26, 2026
commit b98668a131fe35f1e4b6bb291214588f44a0b461
1 change: 1 addition & 0 deletions codex-rs/core/src/codex_delegate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
Expand Down
3 changes: 3 additions & 0 deletions codex-rs/core/src/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -506,6 +507,7 @@ impl Codex {
skills_service,
plugins_manager,
mcp_manager,
code_mode_session_provider,
extensions,
conversation_history,
session_source,
Expand Down Expand Up @@ -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,
Expand Down
11 changes: 4 additions & 7 deletions codex-rs/core/src/session/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -1112,13 +1113,9 @@ impl Session {
session_configuration.parent_thread_id,
),
),
code_mode_service: crate::tools::code_mode::CodeModeService::new(
if config.features.enabled(Feature::CodeModeHost) {
Arc::new(codex_code_mode::ProcessOwnedCodeModeSessionProvider::default())
} else {
Arc::new(codex_code_mode::InProcessCodeModeSessionProvider)
},
),
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),
};
Expand Down
3 changes: 3 additions & 0 deletions codex-rs/core/src/session/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -5604,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,
Expand Down Expand Up @@ -5710,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,
Expand Down
1 change: 1 addition & 0 deletions codex-rs/core/src/session/tests/guardian_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
11 changes: 11 additions & 0 deletions codex-rs/core/src/thread_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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>,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
58 changes: 58 additions & 0 deletions codex-rs/core/src/thread_manager_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
76 changes: 34 additions & 42 deletions codex-rs/core/src/tools/code_mode/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,7 @@ use codex_code_mode::CodeModeToolKind;
use codex_code_mode::RuntimeResponse;
use codex_protocol::models::FunctionCallOutputContentItem;
use serde_json::Value as JsonValue;
use tokio::sync::Mutex;
use tokio::sync::Semaphore;
use tokio::sync::OnceCell;
use tokio_util::sync::CancellationToken;

use crate::function_tool::FunctionCallError;
Expand Down Expand Up @@ -65,25 +64,27 @@ pub(crate) struct ExecContext {
}

pub(crate) struct CodeModeService {
session: Mutex<Option<Arc<dyn CodeModeSession>>>,
session: OnceCell<Arc<dyn CodeModeSession>>,
session_provider: Arc<dyn CodeModeSessionProvider>,
dispatch_broker: Arc<CodeModeDispatchBroker>,
session_init_permit: Semaphore,
shutting_down: AtomicBool,
}

impl CodeModeService {
pub(crate) fn new(session_provider: Arc<dyn CodeModeSessionProvider>) -> Self {
let dispatch_broker = Arc::new(CodeModeDispatchBroker::new());
Self {
session: Mutex::new(None),
session: OnceCell::new(),
session_provider,
dispatch_broker,
session_init_permit: Semaphore::new(/*permits*/ 1),
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,
Expand All @@ -107,14 +108,18 @@ impl CodeModeService {

pub(crate) async fn shutdown(&self) -> Result<(), String> {
self.shutting_down.store(true, Ordering::Release);
let _permit = self
.session_init_permit
.acquire()
// 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
.map_err(|_| "code mode session initializer closed".to_string())?;
match self.current_session().await {
Some(session) => session.shutdown().await,
None => Ok(()),
{
Ok(session) => session.shutdown().await,
Err(_) => Ok(()),
}
}

Expand Down Expand Up @@ -153,36 +158,23 @@ impl CodeModeService {
if self.shutting_down.load(Ordering::Acquire) {
return Err("code mode session is shutting down".to_string());
}
if let Some(session) = self.current_session().await {
return Ok(session);
}

let _permit = self
.session_init_permit
.acquire()
self.session
.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_err(|_| "code mode session initializer closed".to_string())?;
if self.shutting_down.load(Ordering::Acquire) {
return Err("code mode session is shutting down".to_string());
}
if let Some(session) = self.current_session().await {
return Ok(session);
}

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());
}
*self.session.lock().await = Some(Arc::clone(&session));
Ok(session)
}

async fn current_session(&self) -> Option<Arc<dyn CodeModeSession>> {
self.session.lock().await.clone()
.map(Arc::clone)
}
}

Expand Down
1 change: 0 additions & 1 deletion codex-rs/core/tests/suite/code_mode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -236,7 +236,6 @@ async fn missing_process_host_returns_a_tool_error() -> Result<()> {
let output = follow_up_mock
.single_request()
.custom_tool_call_output("call-1");
assert_eq!(output["success"], serde_json::json!(false));
assert!(
output["output"]
.as_str()
Expand Down
Loading