//! RMCP client lifecycle for MCP server connections. //! //! This module owns startup of individual RMCP clients: building the transport, //! initializing the server, listing raw tools, applying per-server tool filters, //! and exposing cached Codex Apps tools while a client is still connecting. //! Initialization capabilities survive tool-discovery failures and reset on a new attempt. //! Higher-level aggregation and resource/tool APIs live in //! [`crate::connection_manager`]. #[path = "rmcp_client/status.rs"] mod status; use std::borrow::Cow; use std::collections::BTreeMap; use std::collections::HashMap; use std::env; use std::ffi::OsString; use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; use std::time::Duration; use std::time::Instant; use crate::client_tool_catalog::ClientToolCatalog; use crate::codex_apps::normalize_codex_apps_callable_name; use crate::codex_apps::normalize_codex_apps_callable_namespace; use crate::codex_apps::normalize_codex_apps_tool_title; use crate::codex_apps::prepare_openai_file_params_for_model; use crate::elicitation::ElicitationRequestManager; use crate::executor_environment_http_client::ExecutorEnvironmentHttpClient; use crate::mcp::CODEX_APPS_MCP_SERVER_NAME; use crate::mcp::ToolPluginContext; use crate::openai_docs_source_attribution::maybe_with_openai_docs_source_attribution; use crate::pagination::collect_paginated_with_limit; use crate::runtime::McpRuntimeContext; use crate::runtime::emit_duration; use crate::server::EffectiveMcpServer; use crate::server::has_explicit_http_authorization; use crate::tool_catalog_cache::McpToolCatalogCacheContext; use crate::tool_catalog_cache::McpToolCatalogFetchTicket; use crate::tools::ToolInfo; use anyhow::Result; use anyhow::anyhow; use async_channel::Sender; use codex_api::SharedAuthProvider; use codex_async_utils::CancelErr; use codex_async_utils::OrCancelExt; use codex_config::McpServerAuth; use codex_config::McpServerConfig; use codex_config::McpServerTransportConfig; use codex_config::types::AuthKeyringBackendKind; use codex_config::types::OAuthCredentialsStoreMode; use codex_connectors::ConnectorRuntimeContext; use codex_connectors::ConnectorRuntimeFetchSource; use codex_exec_server::Environment; use codex_login::AuthChangeState; use codex_protocol::mcp::ClientMcpExtensions; use codex_protocol::mcp::McpServerInfo; use codex_protocol::protocol::Event; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::McpStartupStatus; use codex_protocol::protocol::McpStartupUpdateEvent; use codex_rmcp_client::ExecutorStdioServerLauncher; use codex_rmcp_client::LocalStdioServerLauncher; use codex_rmcp_client::McpOAuthRefreshMode; use codex_rmcp_client::McpProtocolMode; use codex_rmcp_client::RmcpClient; use codex_rmcp_client::StdioServerLauncher; use codex_rmcp_client::StreamableHttpBearerToken; use codex_rmcp_client::StreamableHttpRedirectMode; use codex_rmcp_client::ToolWithConnectorId; use codex_rmcp_client::is_authentication_required_error; use futures::future::BoxFuture; use futures::future::FutureExt; use futures::future::Shared; use rmcp::model::ClientCapabilities; use rmcp::model::ElicitationCapability; use rmcp::model::Implementation; use rmcp::model::InitializeRequestParams; use rmcp::model::ProtocolVersion; use rmcp::model::ServerPeerInfo; use rmcp::model::Tool as RmcpTool; use tokio::sync::watch; use tokio::time::Instant as TokioInstant; use tokio_util::sync::CancellationToken; use tokio_util::task::AbortOnDropHandle; use tracing::Instrument; use tracing::instrument; use tracing::warn; /// MCP server capability indicating that Codex should include [`SandboxState`] /// in tool-call request `_meta` under this key. pub const MCP_SANDBOX_STATE_META_CAPABILITY: &str = "codex/sandbox-state-meta"; /// Experimental MCP server capability for development and testing only; production servers should /// not use it. Its `cacheable: false` property disables sharing tool definitions across connections. const MCP_TOOL_CATALOG_CACHE_CAPABILITY: &str = "codex/tool-catalog-cache"; const MCP_TOOL_CATALOG_CACHEABLE_PROPERTY: &str = "cacheable"; pub(crate) const MCP_TOOLS_LIST_DURATION_METRIC: &str = "codex.mcp.tools.list.duration_ms"; pub(crate) const MCP_TOOLS_FETCH_UNCACHED_DURATION_METRIC: &str = "codex.mcp.tools.fetch_uncached.duration_ms"; pub(crate) const CODEX_APPS_REFRESH_DURATION_METRIC: &str = "codex.apps.refresh.duration_ms"; pub(crate) const DEFAULT_STARTUP_TIMEOUT: Duration = Duration::from_secs(30); pub(crate) const DEFAULT_TOOL_TIMEOUT: Duration = Duration::from_secs(300); pub(crate) const CODEX_APPS_RECONNECT_INITIAL_BACKOFF: Duration = Duration::from_secs(1); const CODEX_APPS_RECONNECT_MAX_BACKOFF: Duration = Duration::from_secs(30); const UNTRUSTED_CONNECTOR_META_KEYS: &[&str] = &[ "connector_id", "connector_name", "connector_display_name", "connector_description", "connectorDescription", ]; #[derive(Clone)] pub(crate) struct ManagedClient { pub(crate) _auth_change_notifications: Option>>, pub(crate) client: Arc, pub(crate) server_info: McpServerInfo, pub(crate) tool_catalog: Arc, pub(crate) tool_timeout: Option, pub(crate) server_instructions: Option, pub(crate) server_supports_sandbox_state_meta_capability: bool, pub(crate) codex_apps_tools_cache_context: Option>, } impl ManagedClient { pub(crate) async fn listed_tools(&self) -> Vec { let total_start = Instant::now(); self.tool_catalog .read(|catalog| { // Discovery may use the shared cache until this client is refreshed. // Executable bindings always capture this client's own catalog. if catalog.revision == 0 && let Some(cache_context) = &self.codex_apps_tools_cache_context { let tools = cache_context.current_tools(); emit_duration( MCP_TOOLS_LIST_DURATION_METRIC, total_start.elapsed(), &[("cache", if tools.is_some() { "hit" } else { "miss" })], ); if let Some(tools) = tools { return tools; } } catalog.tools.to_vec() }) .await } } pub(crate) type ManagedClientFuture = Shared>>; #[derive(Default)] struct CodexAppsStartupReconnectState { current_client: Option, last_error: Option, reconnect_in_flight: bool, consecutive_failures: u32, retry_not_before: Option, } #[derive(Clone)] struct CodexAppsStartupStatusContext { submit_id: String, server_name: String, tx_event: Sender, } pub(crate) struct CodexAppsStartupReconnect { factory: Arc ManagedClientFuture + Send + Sync>, state: StdMutex, startup_status_context: Option, } impl CodexAppsStartupReconnect { pub(crate) fn new(factory: Arc ManagedClientFuture + Send + Sync>) -> Self { Self { factory, state: StdMutex::new(CodexAppsStartupReconnectState::default()), startup_status_context: None, } } fn with_startup_status_context( mut self, submit_id: String, server_name: String, tx_event: Option>, ) -> Self { self.startup_status_context = tx_event.map(|tx_event| CodexAppsStartupStatusContext { submit_id, server_name, tx_event, }); self } fn current_client(&self) -> Option { self.state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .current_client .clone() } fn reconnect_in_background(self: &Arc) { { let mut state = self .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); if state.current_client.is_some() || state.reconnect_in_flight { return; } if state .retry_not_before .is_some_and(|retry_not_before| TokioInstant::now() < retry_not_before) { return; } state.reconnect_in_flight = true; } let reconnect = Arc::clone(self); tokio::spawn(async move { let result = (reconnect.factory)().await; let startup_status_context = reconnect.startup_status_context.clone(); let recovered = { let mut state = reconnect .state .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); state.reconnect_in_flight = false; match result { Ok(client) => { state.current_client = Some(client); state.last_error = None; state.consecutive_failures = 0; state.retry_not_before = None; true } Err(error) => { state.last_error = Some(error.clone()); state.consecutive_failures = state.consecutive_failures.saturating_add(1); let retry_after = codex_apps_reconnect_backoff(state.consecutive_failures); state.retry_not_before = Some(TokioInstant::now() + retry_after); warn!( error = %error, retry_after_ms = retry_after.as_millis(), "Apps MCP startup reconnect failed; continuing with cached tools" ); false } } }; if recovered && let Some(context) = startup_status_context { let _ = context .tx_event .send(Event { id: context.submit_id, msg: EventMsg::McpStartupUpdate(McpStartupUpdateEvent { server: context.server_name, status: McpStartupStatus::Ready, }), }) .await; } }); } } fn codex_apps_reconnect_backoff(consecutive_failures: u32) -> Duration { let exponent = consecutive_failures.saturating_sub(1).min(5); CODEX_APPS_RECONNECT_INITIAL_BACKOFF .saturating_mul(1 << exponent) .min(CODEX_APPS_RECONNECT_MAX_BACKOFF) } #[derive(Clone)] struct ManagedClientStartup { server_name: String, server: EffectiveMcpServer, store_mode: OAuthCredentialsStoreMode, keyring_backend_kind: AuthKeyringBackendKind, oauth_refresh_mode: McpOAuthRefreshMode, tx_event: Option>, elicitation_requests: ElicitationRequestManager, codex_apps_tools_cache_context: Option>, tool_catalog_cache_context: Option, runtime_context: McpRuntimeContext, resolved_environment: std::result::Result>, String>, runtime_auth_provider: Option, client_elicitation_capability: ElicitationCapability, client_mcp_extensions: ClientMcpExtensions, auth_changes: Option>, protocol_mode: McpProtocolMode, catalog_item_limit: usize, cancel_token: CancellationToken, startup_complete: Arc, server_capabilities: Arc>>, } impl ManagedClientStartup { fn start(&self) -> ManagedClientFuture { // A new attempt must not expose capabilities from an earlier connection. *self .server_capabilities .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) = None; let Self { server_name, server, store_mode, keyring_backend_kind, oauth_refresh_mode, tx_event, elicitation_requests, codex_apps_tools_cache_context, tool_catalog_cache_context, runtime_context, resolved_environment, runtime_auth_provider, client_elicitation_capability, client_mcp_extensions, auth_changes, protocol_mode, catalog_item_limit, cancel_token, startup_complete, server_capabilities, } = self.clone(); let is_codex_apps_mcp_server = server_name == CODEX_APPS_MCP_SERVER_NAME; let startup_timeout = server .config() .startup_timeout_sec .unwrap_or(DEFAULT_STARTUP_TIMEOUT); let cancel_token_for_fut = cancel_token; async move { let tool_catalog_fetch_ticket = tool_catalog_cache_context .as_ref() .map(McpToolCatalogCacheContext::begin_fetch); let refresh_start = is_codex_apps_mcp_server.then(Instant::now); let outcome = match async { if let Err(error) = validate_mcp_server_name(&server_name) { return Err(error.into()); } let client = match tokio::time::timeout( startup_timeout, make_rmcp_client( &server_name, server.clone(), store_mode, keyring_backend_kind, oauth_refresh_mode, runtime_context, resolved_environment, runtime_auth_provider, protocol_mode, ), ) .await { Ok(result) => Arc::new(result?), Err(_) => { return Err(StartupOutcomeError::from(anyhow!( "MCP client startup timed out after {startup_timeout:?}" ))); } }; start_server_task( server_name.clone(), client, StartServerTaskParams { is_codex_apps_mcp_server, startup_timeout: Some(startup_timeout), tx_event, elicitation_requests, codex_apps_tools_cache_context, tool_catalog_cache_context, tool_catalog_fetch_ticket, client_elicitation_capability, client_mcp_extensions, auth_changes, catalog_item_limit, server_capabilities, }, ) .await } .or_cancel(&cancel_token_for_fut) .await { Ok(result) => result, Err(CancelErr::Cancelled) => Err(StartupOutcomeError::Cancelled), }; // Log once per startup attempt, including discovery without startup notifications. if let Err(StartupOutcomeError::Failed { error, .. }) = &outcome { warn!(server_name, %error, "MCP server startup failed"); } if outcome.is_ok() && let Some(refresh_start) = refresh_start { emit_duration( CODEX_APPS_REFRESH_DURATION_METRIC, refresh_start.elapsed(), &[("path", "legacy"), ("trigger", "initial")], ); } startup_complete.store(true, Ordering::Release); outcome } .in_current_span() .boxed() .shared() } } #[derive(Clone)] pub(crate) struct AsyncManagedClient { pub(crate) client: ManagedClientFuture, pub(crate) is_codex_apps_mcp_server: bool, pub(crate) cached_server_info: Option, /// Retained after initialization even if subsequent tool discovery fails. pub(crate) server_capabilities: Arc>>, pub(crate) codex_apps_tools_cache_context: Option>, pub(crate) tool_catalog_cache_context: Option, pub(crate) startup_complete: Arc, pub(crate) startup_reconnect: Option>, pub(crate) cancel_token: CancellationToken, } impl AsyncManagedClient { // Keep this constructor flat so the startup inputs remain readable at the // single call site instead of introducing a one-off params wrapper. #[instrument(level = "trace", skip_all, fields(server_name = %server_name))] #[allow(clippy::too_many_arguments)] pub(crate) fn new( server_name: String, startup_submit_id: String, server: EffectiveMcpServer, store_mode: OAuthCredentialsStoreMode, keyring_backend_kind: AuthKeyringBackendKind, oauth_refresh_mode: McpOAuthRefreshMode, cancel_token: CancellationToken, tx_event: Option>, elicitation_requests: ElicitationRequestManager, codex_apps_tools_cache_context: Option>, tool_catalog_cache_context: Option, runtime_context: McpRuntimeContext, resolved_environment: std::result::Result>, String>, runtime_auth_provider: Option, client_elicitation_capability: ElicitationCapability, client_mcp_extensions: ClientMcpExtensions, auth_changes: Option>, protocol_mode: McpProtocolMode, catalog_item_limit: usize, ) -> Self { let is_codex_apps_mcp_server = server_name == CODEX_APPS_MCP_SERVER_NAME; let reconnect_server_name = server_name.clone(); let reconnect_tx_event = tx_event.clone(); let cached_server_info = if is_codex_apps_mcp_server { codex_apps_tools_cache_context .as_ref() .and_then(ConnectorRuntimeContext::cached_server_info) } else { None }; let startup_complete = Arc::new(AtomicBool::new(false)); let server_capabilities = Arc::new(StdMutex::new(None)); let startup = Arc::new(ManagedClientStartup { server_name, server, store_mode, keyring_backend_kind, oauth_refresh_mode, tx_event, elicitation_requests, codex_apps_tools_cache_context: codex_apps_tools_cache_context.clone(), tool_catalog_cache_context: tool_catalog_cache_context.clone(), runtime_context, resolved_environment, runtime_auth_provider, client_elicitation_capability, client_mcp_extensions, auth_changes, protocol_mode, catalog_item_limit, cancel_token: cancel_token.clone(), startup_complete: Arc::clone(&startup_complete), server_capabilities: Arc::clone(&server_capabilities), }); let client = startup.start(); let startup_reconnect = is_codex_apps_mcp_server.then(|| { let startup = Arc::clone(&startup); Arc::new( CodexAppsStartupReconnect::new(Arc::new(move || startup.start())) .with_startup_status_context( startup_submit_id, reconnect_server_name, reconnect_tx_event, ), ) }); Self { client, is_codex_apps_mcp_server, cached_server_info, server_capabilities, codex_apps_tools_cache_context, tool_catalog_cache_context, startup_complete, startup_reconnect, cancel_token, } } pub(crate) async fn client(&self) -> Result { if let Some(client) = self .startup_reconnect .as_ref() .and_then(|reconnect| reconnect.current_client()) { return Ok(client); } self.client.clone().await } /// Returns the current ready client, including its tool catalog and metadata, /// without waiting for startup or initiating a reconnection. pub(crate) fn ready_client(&self) -> Option { self.startup_reconnect .as_ref() .and_then(|reconnect| reconnect.current_client()) .or_else(|| { self.client .peek() .and_then(|result| result.as_ref().ok()) .cloned() }) } pub(crate) fn ready_transport(&self) -> Option> { self.ready_client().map(|client| client.client) } pub(crate) async fn reconnect_failed_startup(&self) { let Some(startup_reconnect) = self.startup_reconnect.as_ref() else { return; }; if !self.startup_complete.load(Ordering::Acquire) { return; } if matches!(self.client().await, Err(StartupOutcomeError::Failed { .. })) { startup_reconnect.reconnect_in_background(); } } pub(crate) async fn shutdown(&self) { self.cancel_token.cancel(); match self.client().await { Ok(client) => client.client.shutdown().await, Err(StartupOutcomeError::Cancelled) => {} Err(error) => { warn!("failed to initialize MCP client during shutdown: {error:#}"); } } } pub(crate) fn has_cached_tools(&self) -> bool { self.codex_apps_tools_cache_context .as_ref() .is_some_and(ConnectorRuntimeContext::has_current_tools) || self .tool_catalog_cache_context .as_ref() .is_some_and(McpToolCatalogCacheContext::has_tools) } pub(crate) fn cached_tools(&self) -> Option> { self.cached_tools_or(/*fallback*/ None) } pub(crate) fn cached_tools_or(&self, fallback: Option>) -> Option> { self.codex_apps_tools_cache_context .as_ref() .and_then(ConnectorRuntimeContext::current_tools) .or_else(|| { self.tool_catalog_cache_context .as_ref() .and_then(|cache| cache.current_tools_or(fallback)) }) } pub(crate) async fn listed_tools(&self) -> Result, StartupOutcomeError> { // Plugin provenance is resolved per-session rather than stored in shared cache payloads. if !self.startup_complete.load(Ordering::Acquire) && let Some(startup_tools) = self.cached_tools() { Ok(startup_tools) } else { match self.client().await { Ok(client) => Ok(client.listed_tools().await), Err(error) if self.is_codex_apps_mcp_server => self.cached_tools().ok_or(error), Err(error) => Err(error), } } } } #[derive(Debug, Clone, thiserror::Error)] pub(crate) enum StartupOutcomeError { #[error("MCP startup cancelled")] Cancelled, // We can't store the original error here because anyhow::Error doesn't implement // `Clone`. #[error("MCP startup failed: {error}")] Failed { error: String, is_authentication_required: bool, }, } impl StartupOutcomeError { pub(crate) fn is_authentication_required(&self) -> bool { match self { Self::Cancelled => false, Self::Failed { error, is_authentication_required, } => *is_authentication_required || error.contains("Auth required"), } } } impl From for StartupOutcomeError { fn from(error: anyhow::Error) -> Self { let is_authentication_required = is_authentication_required_error(&error); Self::Failed { error: format!("{error:#}"), is_authentication_required, } } } #[instrument(level = "trace", skip_all, fields(server_name = %server_name))] pub(crate) async fn list_tools_for_client_uncached( server_name: &str, is_codex_apps_mcp_server: bool, codex_apps_refresh_trigger: &'static str, client: &Arc, timeout: Option, catalog_item_limit: usize, server_instructions: Option<&str>, ) -> Result> { let fetch_start = Instant::now(); let protocol_mode = client.protocol_mode(); let tools = collect_paginated_with_limit("tools/list", timeout, catalog_item_limit, |params| { let client = Arc::clone(client); async move { let response = client .list_tools_with_connector_ids(params, timeout) .await?; let next_cursor = match protocol_mode { McpProtocolMode::Legacy => None, McpProtocolMode::V20260728 => response.next_cursor, }; Ok((response.tools, next_cursor)) } }) .await? .into_iter() .map(|tool| { tool_info_from_listed_tool( server_name, is_codex_apps_mcp_server, server_instructions, tool, ) }) .collect(); if is_codex_apps_mcp_server { emit_duration( MCP_TOOLS_FETCH_UNCACHED_DURATION_METRIC, fetch_start.elapsed(), &[("trigger", codex_apps_refresh_trigger)], ); } else { emit_duration( MCP_TOOLS_FETCH_UNCACHED_DURATION_METRIC, fetch_start.elapsed(), &[], ); } Ok(tools) } /// Filters disabled connectors, presents declared Codex Apps file parameters to the model as /// local-path inputs, and adds plugin names to each tool. Plugin membership is resolved by /// connector ID, falling back to the MCP server when absent. pub(crate) fn prepare_codex_apps_tools_for_model( mut tools: Vec, tool_plugin_context: &ToolPluginContext, ) -> Vec { tools.retain(|tool| tool_plugin_context.allows_connector_id(tool.connector_id.as_deref())); for tool in &mut tools { prepare_openai_file_params_for_model(tool); let plugin_names = match tool.connector_id.as_deref() { Some(connector_id) => { tool_plugin_context.plugin_display_names_for_connector_id(connector_id) } None => tool_plugin_context .plugin_display_names_for_mcp_server_name(tool.server_name.as_str()), }; add_plugin_provenance_to_tool(tool, plugin_names); } tools } /// Stores plugin names on the tool and appends a model-visible plugin membership note. fn add_plugin_provenance_to_tool(tool: &mut ToolInfo, plugin_names: &[String]) { tool.plugin_display_names = plugin_names.to_vec(); if plugin_names.is_empty() { return; } let plugin_source_note = if plugin_names.len() == 1 { format!("This tool is part of plugin `{}`.", plugin_names[0]) } else { format!( "This tool is part of plugins {}.", plugin_names .iter() .map(|plugin_name| format!("`{plugin_name}`")) .collect::>() .join(", ") ) }; let description = tool .tool .description .as_deref() .map(str::trim) .unwrap_or(""); let annotated_description = if description.is_empty() { plugin_source_note } else if matches!(description.chars().last(), Some('.' | '!' | '?')) { format!("{description} {plugin_source_note}") } else { format!("{description}. {plugin_source_note}") }; tool.tool.description = Some(Cow::Owned(annotated_description)); } /// Adds server-scoped plugin names to regular MCP tools without changing their input schemas. pub(crate) fn prepare_regular_mcp_tools_for_model( mut tools: Vec, tool_plugin_context: &ToolPluginContext, ) -> Vec { for tool in &mut tools { let plugin_names = tool_plugin_context.plugin_display_names_for_mcp_server_name(tool.server_name.as_str()); add_plugin_provenance_to_tool(tool, plugin_names); } tools } fn tool_info_from_listed_tool( server_name: &str, is_codex_apps_mcp_server: bool, server_instructions: Option<&str>, tool: ToolWithConnectorId, ) -> ToolInfo { if is_codex_apps_mcp_server { codex_apps_tool_info_from_listed_tool(server_name, server_instructions, tool) } else { regular_mcp_tool_info_from_listed_tool(server_name, server_instructions, tool) } } /// Converts a Codex Apps tool by preserving connector fields, removing connector prefixes from /// model-visible names and titles, and using the connector description for its tool namespace. fn codex_apps_tool_info_from_listed_tool( server_name: &str, server_instructions: Option<&str>, tool: ToolWithConnectorId, ) -> ToolInfo { let mut tool_def = tool.tool; let connector_id = tool.connector_id; let connector_name = tool.connector_name; let connector_description = tool.connector_description; let callable_name = normalize_codex_apps_callable_name( &tool_def.name, connector_id.as_deref(), connector_name.as_deref(), ); let callable_namespace = normalize_codex_apps_callable_namespace(server_name, connector_name.as_deref()); if let Some(title) = tool_def.title.as_deref() { let normalized_title = normalize_codex_apps_tool_title(connector_name.as_deref(), title); if tool_def.title.as_deref() != Some(normalized_title.as_str()) { tool_def.title = Some(normalized_title); } } let has_connector_metadata = connector_id.is_some() || connector_name.is_some() || connector_description.is_some(); let namespace_description = if has_connector_metadata { connector_description } else { server_instructions.map(str::to_string) }; ToolInfo { server_name: server_name.to_owned(), supports_parallel_tool_calls: false, server_origin: None, callable_name, callable_namespace, namespace_description, tool: tool_def, openai_file_input_optional_fields: HashMap::new(), connector_id, connector_name, plugin_display_names: Vec::new(), } } /// Converts a regular MCP tool by removing reserved connector metadata, keeping its raw tool name, /// and using the MCP server name and instructions for the model-visible namespace. fn regular_mcp_tool_info_from_listed_tool( server_name: &str, server_instructions: Option<&str>, tool: ToolWithConnectorId, ) -> ToolInfo { let mut tool_def = tool.tool; strip_untrusted_connector_meta(&mut tool_def); ToolInfo { server_name: server_name.to_owned(), supports_parallel_tool_calls: false, server_origin: None, callable_name: tool_def.name.to_string(), callable_namespace: server_name.to_string(), namespace_description: server_instructions.map(str::to_string), tool: tool_def, openai_file_input_optional_fields: HashMap::new(), connector_id: None, connector_name: None, plugin_display_names: Vec::new(), } } fn strip_untrusted_connector_meta(tool: &mut RmcpTool) { if let Some(meta) = tool.meta.as_mut() { meta.retain(|key, _| !is_untrusted_connector_meta_key(key)); } } fn is_untrusted_connector_meta_key(key: &str) -> bool { UNTRUSTED_CONNECTOR_META_KEYS.contains(&key) } fn resolve_bearer_token( server_name: &str, bearer_token_env_var: Option<&str>, ) -> Result> { let Some(env_var) = bearer_token_env_var else { return Ok(None); }; match env::var(env_var) { Ok(value) => { if value.is_empty() { Err(anyhow!( "Environment variable {env_var} for MCP server '{server_name}' is empty" )) } else { Ok(Some(value)) } } Err(env::VarError::NotPresent) => Err(anyhow!( "Environment variable {env_var} for MCP server '{server_name}' is not set" )), Err(env::VarError::NotUnicode(_)) => Err(anyhow!( "Environment variable {env_var} for MCP server '{server_name}' contains invalid Unicode" )), } } fn validate_mcp_server_name(server_name: &str) -> Result<()> { let re = regex_lite::Regex::new(r"^[a-zA-Z0-9_:@/.-]+$")?; if !re.is_match(server_name) { return Err(anyhow!( "Invalid MCP server name '{server_name}': must match pattern {pattern}", pattern = re.as_str() )); } Ok(()) } #[instrument(level = "trace", skip_all, fields(server_name = %server_name))] async fn start_server_task( server_name: String, client: Arc, params: StartServerTaskParams, ) -> Result { let StartServerTaskParams { is_codex_apps_mcp_server, startup_timeout, tx_event, elicitation_requests, codex_apps_tools_cache_context, tool_catalog_cache_context, tool_catalog_fetch_ticket, client_elicitation_capability, client_mcp_extensions, auth_changes, catalog_item_limit, server_capabilities, } = params; let send_elicitation = elicitation_requests.make_sender(server_name.clone(), tx_event, &client_mcp_extensions); let mut params = mcp_initialize_request_params(client_elicitation_capability, client_mcp_extensions); if auth_changes.is_some() { params .capabilities .experimental .get_or_insert_default() .insert( crate::auth_changes::CAPABILITY.to_string(), Default::default(), ); } let requested_capabilities = params.capabilities.clone(); let started_at = Instant::now(); let initialize_result = client .initialize(params, startup_timeout, send_elicitation) .await; record_protocol_discovery_metrics( client.protocol_mode(), is_codex_apps_mcp_server, started_at, &initialize_result, ); let initialize_result = initialize_result.map_err(StartupOutcomeError::from)?; *server_capabilities .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(serde_json::json!(initialize_result.capabilities)); let auth_change_notifications = crate::auth_changes::start( Arc::clone(&client), &initialize_result.capabilities, auth_changes, ) .await .map_err(StartupOutcomeError::from)?; let server_disables_tool_catalog_cache = initialize_result .capabilities .experimental .as_ref() .and_then(|experimental| experimental.get(MCP_TOOL_CATALOG_CACHE_CAPABILITY)) .and_then(|capability| capability.get(MCP_TOOL_CATALOG_CACHEABLE_PROPERTY)) .and_then(serde_json::Value::as_bool) == Some(false); if server_disables_tool_catalog_cache && let Some(cache_context) = tool_catalog_cache_context.as_ref() { cache_context.disable(); } let server_supports_sandbox_state_meta_capability = initialize_result .capabilities .experimental .as_ref() .and_then(|exp| exp.get(MCP_SANDBOX_STATE_META_CAPABILITY)) .is_some(); let codex_apps_tools_cache_context = codex_apps_tools_cache_context.map(|context| { if server_disables_tool_catalog_cache { context.without_live_scope() } else { // Converted tools can inherit these instructions. Server capabilities and the // negotiated protocol can also differ between otherwise identical connections. let mut scope = serde_json::json!([ requested_capabilities, initialize_result.protocol_version, initialize_result.capabilities, initialize_result.instructions, initialize_result.server_info, ]); scope.sort_all_objects(); context.with_live_scope(scope.to_string()) } }); let list_start = Instant::now(); let server_info = mcp_server_info_from_implementation(&server_name, initialize_result.server_info); let fetch_ticket = codex_apps_tools_cache_context .as_ref() .map(|context| context.begin_fetch(ConnectorRuntimeFetchSource::Startup)); let tools = list_tools_for_client_uncached( &server_name, is_codex_apps_mcp_server, /*codex_apps_refresh_trigger*/ "initial", &client, startup_timeout, catalog_item_limit, initialize_result.instructions.as_deref(), ) .await .map_err(StartupOutcomeError::from)?; let client_tools: Arc<[ToolInfo]> = match (codex_apps_tools_cache_context.as_ref(), fetch_ticket) { (Some(context), Some(ticket)) if server_disables_tool_catalog_cache => { context.publish_runtime_if_newest_accepted(ticket, &server_info, tools.clone()); tools.into() } (Some(context), Some(ticket)) => context .publish_runtime_if_newest_accepted(ticket, &server_info, tools) .shared_tools(), (None, None) => tools.into(), _ => unreachable!("Codex Apps fetch ticket requires cache context"), }; if let (Some(cache_context), Some(fetch_ticket)) = ( tool_catalog_cache_context.as_ref(), tool_catalog_fetch_ticket, ) { cache_context.publish_if_newest(fetch_ticket, &client_tools); } if is_codex_apps_mcp_server || tool_catalog_cache_context.is_some() { emit_duration( MCP_TOOLS_LIST_DURATION_METRIC, list_start.elapsed(), &[("cache", "miss")], ); } let managed = ManagedClient { _auth_change_notifications: auth_change_notifications, client: Arc::clone(&client), server_info, tool_catalog: Arc::new(ClientToolCatalog::new( client_tools, codex_apps_tools_cache_context .as_ref() .and_then(ConnectorRuntimeContext::subscribe), )), tool_timeout: None, server_instructions: initialize_result.instructions, server_supports_sandbox_state_meta_capability, codex_apps_tools_cache_context, }; Ok(managed) } fn record_protocol_discovery_metrics( mode: McpProtocolMode, is_codex_apps_mcp_server: bool, started_at: Instant, result: &Result, ) { let Some(metrics) = codex_otel::global() else { return; }; let mode = match mode { McpProtocolMode::Legacy => "legacy", McpProtocolMode::V20260728 => "auto", }; let outcome = match result { Ok(result) if result.protocol_version == ProtocolVersion::V_2026_07_28 => "modern", Ok(_) => "legacy", Err(_) => "failure", }; let mut tags = vec![("mode", mode), ("outcome", outcome)]; if is_codex_apps_mcp_server { tags.push(("server_kind", "openai_codex_apps")); } let _ = metrics.counter("codex.mcp.protocol_discovery", /*inc*/ 1, &tags); let _ = metrics.record_duration( "codex.mcp.protocol_discovery.duration_ms", started_at.elapsed(), &tags, ); } pub(crate) fn mcp_initialize_request_params( client_elicitation_capability: ElicitationCapability, client_mcp_extensions: ClientMcpExtensions, ) -> InitializeRequestParams { let mut capabilities = ClientCapabilities::default(); capabilities.elicitation = Some(client_elicitation_capability); let extensions = client_mcp_extensions .iter() .filter_map(|(id, settings)| { settings .as_object() .cloned() .map(|settings| (id.to_string(), settings)) }) .collect::>(); if !extensions.is_empty() { capabilities.extensions = Some(extensions); } InitializeRequestParams::new( capabilities, Implementation::new("codex-mcp-client", env!("CARGO_PKG_VERSION")).with_title("Codex"), ) .with_protocol_version(ProtocolVersion::V_2025_06_18) } fn mcp_server_info_from_implementation( server_name: &str, server_info: Option, ) -> McpServerInfo { let server_info = server_info.unwrap_or_else(|| Implementation::new(server_name, "")); McpServerInfo { name: server_info.name, title: server_info.title, version: server_info.version, description: server_info.description, icons: server_info.icons.map(|icons| { icons .into_iter() .filter_map(|icon| serde_json::to_value(icon).ok()) .collect() }), website_url: server_info.website_url, } } struct StartServerTaskParams { server_capabilities: Arc>>, is_codex_apps_mcp_server: bool, startup_timeout: Option, // TODO: cancel_token should handle this. tx_event: Option>, elicitation_requests: ElicitationRequestManager, codex_apps_tools_cache_context: Option>, tool_catalog_cache_context: Option, tool_catalog_fetch_ticket: Option, client_elicitation_capability: ElicitationCapability, client_mcp_extensions: ClientMcpExtensions, auth_changes: Option>, catalog_item_limit: usize, } #[allow(clippy::too_many_arguments)] #[instrument(level = "trace", skip_all, fields(server_name = %server_name))] pub(crate) async fn make_rmcp_client( server_name: &str, server: EffectiveMcpServer, store_mode: OAuthCredentialsStoreMode, keyring_backend_kind: AuthKeyringBackendKind, oauth_refresh_mode: McpOAuthRefreshMode, runtime_context: McpRuntimeContext, resolved_environment: std::result::Result>, String>, runtime_auth_provider: Option, protocol_mode: McpProtocolMode, ) -> Result { let config = server.config().clone(); if matches!(config.auth, McpServerAuth::EmaAuth) { return Err(StartupOutcomeError::from(anyhow!( "EMA MCP connections are not enabled in this version" ))); } if matches!(config.auth, McpServerAuth::ChatGpt) && !config.is_local_environment() && !has_explicit_http_authorization(&config) { return Err(StartupOutcomeError::from(anyhow!( "executor-owned MCP server `{server_name}` cannot use hosted ChatGPT authentication; configure executor-owned credentials instead" ))); } let resolved_environment = resolved_environment.map_err(|err| StartupOutcomeError::from(anyhow!(err)))?; let is_local_environment = config.is_local_environment(); let oauth_credential_name = config.oauth_credential_name(server_name); let McpServerConfig { transport, .. } = config; match transport { McpServerTransportConfig::Stdio { command, args, env, env_vars, cwd, } => { let command_os: OsString = command.into(); let args_os: Vec = args.into_iter().map(Into::into).collect(); let env_os = env.map(|env| { env.into_iter() .map(|(key, value)| (key.into(), value.into())) .collect::>() }); let launcher = if is_local_environment { // TODO(starr): Unify local stdio MCP launch with // `ExecutorStdioServerLauncher` once the executor-backed path // preserves `LocalStdioServerLauncher` semantics. Arc::new(LocalStdioServerLauncher::new( runtime_context.local_process_cwd(), )) as Arc } else { let Some(environment) = resolved_environment.as_ref() else { unreachable!( "non-local stdio MCP servers resolve an environment before launch" ); }; Arc::new(ExecutorStdioServerLauncher::new( environment.get_exec_backend(), )) as Arc }; let cwd = cwd.map(codex_utils_path_uri::LegacyAppPathString::into_string); RmcpClient::new_stdio_client_with_protocol_mode( command_os, args_os, env_os, &env_vars, cwd, launcher, protocol_mode, ) .await .map_err(|err| StartupOutcomeError::from(anyhow!(err))) } McpServerTransportConfig::StreamableHttp { url, http_headers, env_http_headers, bearer_token_env_var, http_headers_helper: _, } => { let http_client = runtime_context .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 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 { StreamableHttpRedirectMode::Legacy }; RmcpClient::new_streamable_http_client_with_protocol_mode_and_redirect_mode( oauth_credential_name.as_ref(), &url, resolved_bearer_token, http_headers, env_http_headers, store_mode, keyring_backend_kind, http_client, runtime_auth_provider, protocol_mode, redirect_mode, oauth_refresh_mode, ) .await .map_err(StartupOutcomeError::from) } } } #[cfg(test)] mod tests { use super::*; use codex_protocol::mcp::MCP_APP_UI_EXTENSION_ID; use codex_protocol::mcp::OPENAI_FORM_EXTENSION_ID; use pretty_assertions::assert_eq; use rmcp::model::JsonObject; use rmcp::model::MetaObject; use rmcp::transport::auth::AuthError; #[test] fn startup_outcome_error_identifies_authentication_required() { let error = anyhow::Error::new(AuthError::AuthorizationRequired) .context("failed to initialize MCP server"); let error = StartupOutcomeError::from(error); assert!(error.is_authentication_required()); } #[test] fn missing_server_implementation_uses_configured_server_name() { assert_eq!( mcp_server_info_from_implementation("configured-server", /*server_info*/ None), McpServerInfo { name: "configured-server".to_string(), title: None, version: String::new(), description: None, icons: None, website_url: None, } ); } #[test] fn advertised_server_implementation_takes_precedence_over_configured_name() { assert_eq!( mcp_server_info_from_implementation( "configured-server", Some( Implementation::new("advertised-server", "1.2.3") .with_title("Advertised server") .with_description("Advertised description") .with_website_url("https://example.com"), ), ), McpServerInfo { name: "advertised-server".to_string(), title: Some("Advertised server".to_string()), version: "1.2.3".to_string(), description: Some("Advertised description".to_string()), icons: None, website_url: Some("https://example.com".to_string()), } ); } #[test] fn mcp_initialize_advertises_client_extensions() { let unsupported = mcp_initialize_request_params( ElicitationCapability::default(), ClientMcpExtensions::default(), ); assert_eq!(unsupported.capabilities.extensions, None); let app_ui = serde_json::json!({ "mimeTypes": ["text/html;profile=mcp-app"], "futureField": {"preserved": true}, }); let supported = mcp_initialize_request_params( ElicitationCapability::default(), ClientMcpExtensions::new([ (OPENAI_FORM_EXTENSION_ID.to_string(), serde_json::json!({})), (MCP_APP_UI_EXTENSION_ID.to_string(), app_ui.clone()), ]), ); assert_eq!( supported.capabilities.extensions, Some(BTreeMap::from([ (OPENAI_FORM_EXTENSION_ID.to_string(), JsonObject::new()), ( MCP_APP_UI_EXTENSION_ID.to_string(), app_ui.as_object().cloned().expect("app UI settings"), ), ])) ); } fn tool_with_connector_meta() -> RmcpTool { RmcpTool::new( "capture_file_upload", "test tool", Arc::new(JsonObject::default()), ) .with_meta(MetaObject( serde_json::json!({ "connector_id": "connector_gmail", "connector_name": "Gmail", "connector_display_name": "Gmail", "connector_description": "Mail connector", "connectorDescription": "Mail connector", "connectorFutureField": "future connector metadata", "CONNECTOR_UPPERCASE": "uppercase connector metadata", "openai/fileParams": ["file"], "custom": "kept" }) .as_object() .expect("object") .clone(), )) } #[test] fn custom_mcp_connector_metadata_is_stripped() { let mut tool = tool_with_connector_meta(); strip_untrusted_connector_meta(&mut tool); let meta = tool.meta.as_ref().expect("meta"); for key in [ "connector_id", "connector_name", "connector_display_name", "connector_description", "connectorDescription", ] { assert!(!meta.0.contains_key(key), "{key} should be stripped"); } assert!(meta.0.contains_key("connectorFutureField")); assert!(meta.0.contains_key("CONNECTOR_UPPERCASE")); assert!(meta.0.contains_key("openai/fileParams")); assert_eq!( meta.0.get("custom").and_then(|value| value.as_str()), Some("kept") ); } #[test] fn codex_apps_connector_metadata_is_preserved() { let tool = tool_with_connector_meta(); let expected_tool = tool.clone(); let tool_info = tool_info_from_listed_tool( CODEX_APPS_MCP_SERVER_NAME, /*is_codex_apps_mcp_server*/ true, /*server_instructions*/ None, ToolWithConnectorId { tool, connector_id: Some("connector_gmail".to_string()), connector_name: Some("Gmail".to_string()), connector_description: Some("Mail connector".to_string()), }, ); let expected = ToolInfo { server_name: CODEX_APPS_MCP_SERVER_NAME.to_string(), supports_parallel_tool_calls: false, server_origin: None, callable_name: "capture_file_upload".to_string(), callable_namespace: "codex_apps__gmail".to_string(), namespace_description: Some("Mail connector".to_string()), tool: expected_tool, openai_file_input_optional_fields: HashMap::new(), connector_id: Some("connector_gmail".to_string()), connector_name: Some("Gmail".to_string()), plugin_display_names: Vec::new(), }; assert_eq!( serde_json::to_value(tool_info).expect("serialize actual tool info"), serde_json::to_value(expected).expect("serialize expected tool info") ); } }