diff --git a/codex-rs/app-server/tests/suite/v2/executor_mcp.rs b/codex-rs/app-server/tests/suite/v2/executor_mcp.rs index fd4921b839..5348bf5b5a 100644 --- a/codex-rs/app-server/tests/suite/v2/executor_mcp.rs +++ b/codex-rs/app-server/tests/suite/v2/executor_mcp.rs @@ -27,6 +27,8 @@ use codex_http_client::HttpClientBuilder; use codex_utils_path_uri::PathUri; use core_test_support::responses; use core_test_support::stdio_server_bin; +use futures::SinkExt; +use futures::StreamExt; use pretty_assertions::assert_eq; use rmcp::handler::server::ServerHandler; use rmcp::model::CallToolRequestParams; @@ -42,6 +44,7 @@ use rmcp::service::RoleServer; use rmcp::transport::StreamableHttpServerConfig; use rmcp::transport::StreamableHttpService; use rmcp::transport::streamable_http_server::session::local::LocalSessionManager; +use serde_json::Value; use serde_json::json; use std::borrow::Cow; use std::collections::BTreeMap; @@ -57,6 +60,9 @@ use tokio::net::TcpListener; use tokio::process::Command; use tokio::sync::mpsc; use tokio::time::timeout; +use tokio_tungstenite::accept_async; +use tokio_tungstenite::connect_async; +use tokio_tungstenite::tungstenite::Message; use url::Url; const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(20); @@ -232,6 +238,179 @@ async fn selected_executor_discovers_browser_mcp_with_executor_only_bearer_token Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn legacy_executor_skips_required_browser_and_keeps_host_owned_mcp() -> Result<()> { + let responses_server = responses::start_mock_server().await; + let codex_home = TempDir::new()?; + let executor_home = TempDir::new()?; + + let http_listener = TcpListener::bind("127.0.0.1:0").await?; + let mcp_url = format!("http://{}/mcp", http_listener.local_addr()?); + let expected_authorization = "Bearer host-only-token"; + let service = StreamableHttpService::new( + || Ok(ExecutorHttpMcpServer), + Arc::new(LocalSessionManager::default()), + StreamableHttpServerConfig::default(), + ); + let router = Router::new() + .nest_service("/mcp", service) + .layer(axum::middleware::from_fn( + move |request: axum::extract::Request, next: axum::middleware::Next| async move { + let authorized = request + .headers() + .get(axum::http::header::AUTHORIZATION) + .and_then(|value| value.to_str().ok()) + .is_some_and(|value| value == expected_authorization); + if !authorized { + return axum::http::StatusCode::UNAUTHORIZED.into_response(); + } + next.run(request).await + }, + )); + let http_server = tokio::spawn(async move { + let _ = axum::serve(http_listener, router).await; + }); + + let root_config = format!( + "[mcp_servers.host_remote]\nurl = \"{mcp_url}\"\nenvironment_id = \"{EXECUTOR_ID}\"\nbearer_token_env_var = \"HOST_MCP_TEST_TOKEN\"\nrequired = true\n" + ); + MockResponsesConfig::new(&responses_server.uri()) + .with_sandbox_mode("danger-full-access") + .with_extra_config(&root_config) + .write(codex_home.path())?; + std::fs::write( + executor_home.path().join("config.toml"), + format!( + "[mcp_servers.node_repl]\nurl = \"http://127.0.0.1:9/mcp\"\nbearer_token_env_var = \"NODE_REPL_AUTH_TOKEN\"\nrequired = true\nstartup_timeout_sec = 1\n\n[mcp_servers.executor_public]\nurl = \"{mcp_url}\"\nhttp_headers = {{ Authorization = \"Bearer host-only-token\" }}\nrequired = true\n" + ), + )?; + + let mut executor = Command::new(codex_utils_cargo_bin::cargo_bin("codex")?) + .args(["exec-server", "--listen", "ws://127.0.0.1:0"]) + .stdout(Stdio::piped()) + .kill_on_drop(true) + .env("CODEX_HOME", executor_home.path()) + .spawn()?; + let stdout = executor.stdout.take().expect("executor stdout is piped"); + let mut lines = BufReader::new(stdout).lines(); + let upstream_url = timeout(DEFAULT_READ_TIMEOUT, lines.next_line()) + .await?? + .expect("executor emits its websocket URL"); + + let proxy_listener = TcpListener::bind("127.0.0.1:0").await?; + let proxy_url = format!("ws://{}", proxy_listener.local_addr()?); + let legacy_executor = tokio::spawn(async move { + let (stream, _) = proxy_listener.accept().await?; + let downstream = accept_async(stream).await?; + let (upstream, _) = connect_async(upstream_url).await?; + let (mut downstream_tx, mut downstream_rx) = downstream.split(); + let (mut upstream_tx, mut upstream_rx) = upstream.split(); + + let requests = async { + while let Some(message) = downstream_rx.next().await { + let mut message = message?; + if let Message::Text(text) = &message { + let mut request: Value = serde_json::from_str(text)?; + if request["method"] == "http/request" + && let Some(headers) = request["params"]["headers"].as_array_mut() + { + for header in headers { + if let Some(header) = header.as_object_mut() { + header.remove("valueEnvVar"); + } + } + message = Message::Text(request.to_string().into()); + } + } + upstream_tx.send(message).await?; + } + Ok::<_, anyhow::Error>(()) + }; + + let responses = async { + while let Some(message) = upstream_rx.next().await { + let mut message = message?; + if let Message::Text(text) = &message { + let mut response: Value = serde_json::from_str(text)?; + if let Some(capabilities) = response + .pointer_mut("/result/capabilities") + .and_then(Value::as_object_mut) + { + capabilities.remove("httpHeaderEnvVars"); + message = Message::Text(response.to_string().into()); + } + } + downstream_tx.send(message).await?; + } + Ok::<_, anyhow::Error>(()) + }; + + tokio::select! { + result = requests => result, + result = responses => result, + } + }); + std::fs::write( + codex_home.path().join("environments.toml"), + format!( + "default = \"{EXECUTOR_ID}\"\ninclude_local = false\n\n[[environments]]\nid = \"{EXECUTOR_ID}\"\nurl = \"{proxy_url}\"\n" + ), + )?; + + let mut app_server = TestAppServer::builder() + .with_codex_home(codex_home.path()) + .with_env_overrides(&[("HOST_MCP_TEST_TOKEN", Some("host-only-token"))]) + .without_auto_env() + .build_initialized_with_timeout(DEFAULT_READ_TIMEOUT) + .await?; + let thread_id = start_thread(&mut app_server, /*selected_capability_roots*/ None).await?; + let servers = mcp_server_statuses(&mut app_server, thread_id.clone()).await?; + assert!(servers.iter().any(|server| server.name == "host_remote")); + assert!( + servers + .iter() + .any(|server| server.name == "executor_public") + ); + assert!(servers.iter().all(|server| server.name != "node_repl")); + + let response = responses::mount_sse_once( + &responses_server, + responses::sse(vec![ + responses::ev_response_created("legacy-executor-turn"), + responses::ev_assistant_message("legacy-executor-message", "Still works"), + responses::ev_completed("legacy-executor-turn"), + ]), + ) + .await; + let request_id = app_server + .send_turn_start_request(TurnStartParams { + thread_id, + input: vec![UserInput::Text { + text: "Check legacy executor compatibility".to_string(), + text_elements: Vec::new(), + }], + ..Default::default() + }) + .await?; + let _: TurnStartResponse = + timeout(DEFAULT_READ_TIMEOUT, app_server.read_response(request_id)).await??; + timeout( + DEFAULT_READ_TIMEOUT, + app_server.read_stream_until_notification_message("turn/completed"), + ) + .await??; + assert!( + response + .single_request() + .tool_by_name("mcp__host_remote", "echo") + .is_some() + ); + + legacy_executor.abort(); + http_server.abort(); + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn guardian_review_does_not_discover_executor_mcp() -> Result<()> { let responses_server = responses::start_mock_server().await; diff --git a/codex-rs/codex-mcp/src/rmcp_client.rs b/codex-rs/codex-mcp/src/rmcp_client.rs index 157acf0c84..1468df5c37 100644 --- a/codex-rs/codex-mcp/src/rmcp_client.rs +++ b/codex-rs/codex-mcp/src/rmcp_client.rs @@ -1139,21 +1139,39 @@ async fn make_rmcp_client( .http_client_for_server(server.config(), resolved_environment.as_ref()) .map_err(|error| StartupOutcomeError::from(anyhow!(error)))?; let http_client = maybe_with_openai_docs_source_attribution(&url, http_client); - let (http_client, resolved_bearer_token) = - if !is_local_environment && let Some(env_var) = bearer_token_env_var.as_ref() { - ( - Arc::new(ExecutorEnvironmentHttpClient { - bearer_token_env_var: env_var.clone(), - http_client, - }) as Arc, - Some(StreamableHttpBearerToken::ProvidedByHttpClient), - ) - } else { - let token = resolve_bearer_token(server_name, bearer_token_env_var.as_deref()) - .map_err(StartupOutcomeError::from)? - .map(StreamableHttpBearerToken::Resolved); - (http_client, token) + let executor_resolves_bearer_token = if !is_local_environment + && bearer_token_env_var.is_some() + { + let Some(environment) = resolved_environment.as_ref() else { + return Err(StartupOutcomeError::from(anyhow!( + "non-local HTTP MCP server `{server_name}` did not resolve an execution environment" + ))); }; + environment + .info() + .await + .map_err(|error| StartupOutcomeError::from(anyhow!(error)))? + .capabilities + .http_header_env_vars + } else { + false + }; + let (http_client, resolved_bearer_token) = if executor_resolves_bearer_token + && let Some(env_var) = bearer_token_env_var.as_ref() + { + ( + Arc::new(ExecutorEnvironmentHttpClient { + bearer_token_env_var: env_var.clone(), + http_client, + }) as Arc, + Some(StreamableHttpBearerToken::ProvidedByHttpClient), + ) + } else { + let token = resolve_bearer_token(server_name, bearer_token_env_var.as_deref()) + .map_err(StartupOutcomeError::from)? + .map(StreamableHttpBearerToken::Resolved); + (http_client, token) + }; let redirect_mode = if server.is_agent_plugin() { StreamableHttpRedirectMode::AgentPluginV1 } else { diff --git a/codex-rs/codex-mcp/src/server.rs b/codex-rs/codex-mcp/src/server.rs index 002e1d5e99..6dea78d27f 100644 --- a/codex-rs/codex-mcp/src/server.rs +++ b/codex-rs/codex-mcp/src/server.rs @@ -335,7 +335,7 @@ fn referenced_environment_variables(config: &McpServerConfig) -> Vec<(String, Op .. } => bearer_token_env_var .iter() - .filter(|_| config.is_local_environment()) + .filter(|name| config.is_local_environment() || std::env::var_os(name).is_some()) .chain(env_http_headers.iter().flat_map(|headers| headers.values())) .cloned() .collect(), diff --git a/codex-rs/codex-mcp/src/server_tests.rs b/codex-rs/codex-mcp/src/server_tests.rs index 87992a8800..7365f75573 100644 --- a/codex-rs/codex-mcp/src/server_tests.rs +++ b/codex-rs/codex-mcp/src/server_tests.rs @@ -18,6 +18,17 @@ fn remote_http_connections_track_host_headers_but_not_executor_bearer_tokens() { vec![("PATH".to_string(), std::env::var_os("PATH"))], ); + let remote_host_bearer: McpServerConfig = serde_json::from_value(serde_json::json!({ + "url": "https://example.com/mcp", + "environment_id": "executor-1", + "bearer_token_env_var": "PATH", + })) + .expect("host-resolved remote MCP configuration should deserialize"); + assert_eq!( + referenced_environment_variables(&remote_host_bearer), + vec![("PATH".to_string(), std::env::var_os("PATH"))], + ); + config.environment_id = DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(); assert_eq!( referenced_environment_variables(&config), diff --git a/codex-rs/exec-server-protocol/src/protocol.rs b/codex-rs/exec-server-protocol/src/protocol.rs index 411a32f46d..c46978f966 100644 --- a/codex-rs/exec-server-protocol/src/protocol.rs +++ b/codex-rs/exec-server-protocol/src/protocol.rs @@ -118,6 +118,9 @@ pub struct EnvironmentCapabilities { /// Whether this executor supports the `environmentConfig/read` request. #[serde(default)] pub environment_config_read: bool, + /// Whether HTTP headers can resolve values from the executor environment. + #[serde(default)] + pub http_header_env_vars: bool, /// Whether filesystem streams can use the requested platform sandbox. #[serde(default)] pub sandboxed_file_streaming: bool, @@ -185,6 +188,7 @@ impl EnvironmentInfo { network_proxy_launch: true, capability_discovery_sandbox: true, environment_config_read: true, + http_header_env_vars: true, sandboxed_file_streaming: true, shell_snapshot_v2: false, }, @@ -963,6 +967,7 @@ mod tests { network_proxy_launch: true, capability_discovery_sandbox: true, environment_config_read: false, + http_header_env_vars: false, sandboxed_file_streaming: false, shell_snapshot_v2: false, } @@ -979,6 +984,7 @@ mod tests { "networkProxyLaunch": false, "capabilityDiscoverySandbox": false, "environmentConfigRead": false, + "httpHeaderEnvVars": false, "sandboxedFileStreaming": false, "shellSnapshotV2": false, }, diff --git a/codex-rs/exec-server/src/environment_config.rs b/codex-rs/exec-server/src/environment_config.rs index d6d32d5e1a..b53208b63a 100644 --- a/codex-rs/exec-server/src/environment_config.rs +++ b/codex-rs/exec-server/src/environment_config.rs @@ -68,20 +68,24 @@ impl Environment { &self, cwd: PathUri, ) -> Result, ExecServerError> { - let response = tokio::time::timeout(std::time::Duration::from_secs(10), async { - if !self.info().await?.capabilities.environment_config_read { - return Ok(None); - } - self.read_environment_config(EnvironmentConfigReadParams { - cwd, - config_paths: vec![vec!["mcp_servers".to_string()]], - requirements_paths: vec![vec!["mcp_servers".to_string()]], + let (response, http_header_env_vars) = + tokio::time::timeout(std::time::Duration::from_secs(10), async { + let capabilities = self.info().await?.capabilities; + if !capabilities.environment_config_read { + return Ok((None, false)); + } + self.read_environment_config(EnvironmentConfigReadParams { + cwd, + config_paths: vec![vec!["mcp_servers".to_string()]], + requirements_paths: vec![vec!["mcp_servers".to_string()]], + }) + .await + .map(|response| (Some(response), capabilities.http_header_env_vars)) }) .await - .map(Some) - }) - .await - .map_err(|_| ExecServerError::Protocol("executor MCP discovery timed out".to_string()))??; + .map_err(|_| { + ExecServerError::Protocol("executor MCP discovery timed out".to_string()) + })??; let Some(response) = response else { return Ok(Vec::new()); }; @@ -107,7 +111,10 @@ impl Environment { .cloned() .unwrap_or_else(|| toml::Value::Table(toml::map::Map::new())); if let Some(servers) = servers.as_table_mut() { - servers.retain(|_, server| server.get("url").is_some()); + servers.retain(|_, server| { + server.get("url").is_some() + && (http_header_env_vars || server.get("bearer_token_env_var").is_none()) + }); } let servers = servers .try_into::>()