From 6629e0870231fddfd4c0b6abacaad54f2175fc83 Mon Sep 17 00:00:00 2001 From: jif Date: Tue, 14 Jul 2026 14:35:21 +0000 Subject: [PATCH] Add an agent extension runner (#33076) ## What changed - Add `codex-agent-extension` with an `AgentRunner` that starts a resolved agent prompt in a thread forked from its parent. - Propagate the invocation's trace context, select the configured execution environments, submit the initial prompt, and return the spawned thread and turn identifiers. - Reject empty agent prompts and report when the owning thread manager is no longer available. ## Testing - Add an integration test that verifies the agent runs in a forked thread, returns the started turn identifier, completes the turn, and sends the resolved prompt to the model. GitOrigin-RevId: ffb805efabadf758969d6611a6867db6fa9f059a --- codex-rs/Cargo.lock | 12 +++ codex-rs/Cargo.toml | 1 + codex-rs/analytics/src/client_tests.rs | 1 + codex-rs/ext/agent/BUILD.bazel | 6 ++ codex-rs/ext/agent/Cargo.toml | 24 +++++ codex-rs/ext/agent/src/lib.rs | 104 ++++++++++++++++++++++ codex-rs/ext/agent/tests/agent_service.rs | 70 +++++++++++++++ 7 files changed, 218 insertions(+) create mode 100644 codex-rs/ext/agent/BUILD.bazel create mode 100644 codex-rs/ext/agent/Cargo.toml create mode 100644 codex-rs/ext/agent/src/lib.rs create mode 100644 codex-rs/ext/agent/tests/agent_service.rs diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 350fe41fa4ae..e1e13b225d0a 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -1866,6 +1866,18 @@ dependencies = [ "unicode-width 0.2.1", ] +[[package]] +name = "codex-agent-extension" +version = "0.0.0" +dependencies = [ + "anyhow", + "codex-core", + "codex-protocol", + "core_test_support", + "pretty_assertions", + "tokio", +] + [[package]] name = "codex-agent-graph-store" version = "0.0.0" diff --git a/codex-rs/Cargo.toml b/codex-rs/Cargo.toml index 3f0acfd4a7a2..9ee09478b639 100644 --- a/codex-rs/Cargo.toml +++ b/codex-rs/Cargo.toml @@ -48,6 +48,7 @@ members = [ "exec-server-protocol", "exec-server", "execpolicy", + "ext/agent", "ext/connectors", "ext/extension-api", "ext/goal", diff --git a/codex-rs/analytics/src/client_tests.rs b/codex-rs/analytics/src/client_tests.rs index 82b451acfc20..db4cf92d6134 100644 --- a/codex-rs/analytics/src/client_tests.rs +++ b/codex-rs/analytics/src/client_tests.rs @@ -120,6 +120,7 @@ fn sample_mcp_tool_call_event(thread_id: &str, plugin_id: Option<&str>) -> Track event_params: CodexMcpToolCallEventParams { base: CodexToolItemEventBase { thread_id: thread_id.to_string(), + session_id: format!("session-{thread_id}"), turn_id: "turn-1".to_string(), item_id: format!("item-{thread_id}"), app_server_client: CodexAppServerClientMetadata { diff --git a/codex-rs/ext/agent/BUILD.bazel b/codex-rs/ext/agent/BUILD.bazel new file mode 100644 index 000000000000..21793fc95778 --- /dev/null +++ b/codex-rs/ext/agent/BUILD.bazel @@ -0,0 +1,6 @@ +load("//:defs.bzl", "codex_rust_crate") + +codex_rust_crate( + name = "agent", + crate_name = "codex_agent_extension", +) diff --git a/codex-rs/ext/agent/Cargo.toml b/codex-rs/ext/agent/Cargo.toml new file mode 100644 index 000000000000..6c3f10a43319 --- /dev/null +++ b/codex-rs/ext/agent/Cargo.toml @@ -0,0 +1,24 @@ +[package] +edition.workspace = true +license.workspace = true +name = "codex-agent-extension" +version.workspace = true + +[lib] +name = "codex_agent_extension" +path = "src/lib.rs" +doctest = false +test = false + +[lints] +workspace = true + +[dependencies] +codex-core = { workspace = true } +codex-protocol = { workspace = true } + +[dev-dependencies] +anyhow = { workspace = true } +core_test_support = { workspace = true } +pretty_assertions = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/codex-rs/ext/agent/src/lib.rs b/codex-rs/ext/agent/src/lib.rs new file mode 100644 index 000000000000..01e068545ff7 --- /dev/null +++ b/codex-rs/ext/agent/src/lib.rs @@ -0,0 +1,104 @@ +use codex_core::CodexThread; +use codex_core::NewThread; +use codex_core::StartThreadOptions; +use codex_core::ThreadManager; +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; +use std::sync::Weak; + +/// A fully resolved agent invocation. +/// +/// Agent discovery owns rendering `prompt`, including any selected skill +/// references. The runtime only starts that prompt in isolated forked context. +pub struct AgentInvocation { + pub config: Config, + pub prompt: String, + pub parent_trace: Option, +} + +/// A spawned agent whose initial turn has been submitted. +pub struct AgentRun { + pub thread_id: ThreadId, + pub turn_id: String, + pub thread: Arc, +} + +/// Runs resolved agents in threads forked by the owning [`ThreadManager`]. +#[derive(Clone)] +pub struct AgentRunner { + thread_manager: Weak, +} + +impl AgentRunner { + pub fn new(thread_manager: Weak) -> Self { + Self { thread_manager } + } + + /// Starts a resolved agent in a fork of `parent_thread_id`. + pub async fn start( + &self, + parent_thread_id: ThreadId, + invocation: AgentInvocation, + ) -> CodexResult { + let AgentInvocation { + config, + prompt, + parent_trace, + } = invocation; + if prompt.trim().is_empty() { + return Err(CodexErr::InvalidRequest( + "agent prompt must not be empty".to_string(), + )); + } + + let thread_manager = self + .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, + }, + ) + .await?; + let turn_id = thread + .submit_with_trace( + vec![UserInput::Text { + text: prompt, + text_elements: Vec::new(), + }] + .into(), + parent_trace, + ) + .await?; + + Ok(AgentRun { + thread_id, + turn_id, + thread, + }) + } +} diff --git a/codex-rs/ext/agent/tests/agent_service.rs b/codex-rs/ext/agent/tests/agent_service.rs new file mode 100644 index 000000000000..9b36d4a69bd1 --- /dev/null +++ b/codex-rs/ext/agent/tests/agent_service.rs @@ -0,0 +1,70 @@ +use anyhow::Result; +use codex_agent_extension::AgentInvocation; +use codex_agent_extension::AgentRunner; +use codex_protocol::protocol::EventMsg; +use core_test_support::responses; +use core_test_support::skip_if_no_network; +use core_test_support::test_codex::test_codex; +use core_test_support::wait_for_event; +use pretty_assertions::assert_eq; + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn starts_resolved_agent_prompt_in_forked_thread() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = responses::start_mock_server().await; + let response_mock = responses::mount_sse_once( + &server, + responses::sse(vec![ + responses::ev_response_created("agent-response"), + responses::ev_completed("agent-response"), + ]), + ) + .await; + let test = test_codex().build_with_auto_env(&server).await?; + let parent_thread_id = test.session_configured.session_id.into(); + let agent_runner = AgentRunner::new(std::sync::Arc::downgrade(&test.thread_manager)); + + let agent_run = agent_runner + .start( + parent_thread_id, + AgentInvocation { + config: test.config.clone(), + prompt: "Use $example-agent to inspect the current changes.".to_string(), + parent_trace: None, + }, + ) + .await?; + + assert_ne!(agent_run.thread_id, parent_thread_id); + assert_eq!( + agent_run + .thread + .config_snapshot() + .await + .forked_from_thread_id, + Some(parent_thread_id) + ); + let started = wait_for_event(&agent_run.thread, |event| { + matches!(event, EventMsg::TurnStarted(_)) + }) + .await; + let EventMsg::TurnStarted(started) = started else { + unreachable!("event predicate only matches turn started events"); + }; + assert_eq!(started.turn_id, agent_run.turn_id); + wait_for_event(&agent_run.thread, |event| { + matches!(event, EventMsg::TurnComplete(_)) + }) + .await; + + let request = response_mock.single_request(); + assert!( + request + .message_input_texts("user") + .iter() + .any(|text| text == "Use $example-agent to inspect the current changes.") + ); + + Ok(()) +}