Skip to main content

codex_mcp/
connection_manager.rs

1//! Aggregates MCP server connections for Codex.
2//!
3//! [`McpConnectionManager`] owns the set of running async RMCP clients keyed by
4//! MCP server name. It coordinates startup status events, keeps server origin
5//! metadata, aggregates tools/resources/templates across servers, routes tool
6//! calls to the right client, and exposes the public manager API used by
7//! `codex-core`.
8
9#[path = "connection_manager/required.rs"]
10mod required;
11#[path = "connection_manager/tool_catalog.rs"]
12mod tool_catalog;
13
14use std::collections::HashMap;
15use std::path::PathBuf;
16use std::sync::Arc;
17use std::sync::atomic::Ordering;
18use std::time::Duration;
19
20use crate::binding_clients::McpBindingClients;
21use crate::elicitation::ElicitationRequestManager;
22use crate::elicitation::ElicitationRequestRouter;
23use crate::elicitation::ElicitationReviewerHandle;
24use crate::mcp::CODEX_APPS_MCP_SERVER_NAME;
25use crate::mcp::ToolPluginProvenance;
26use crate::rmcp_client::AsyncManagedClient;
27use crate::rmcp_client::DEFAULT_STARTUP_TIMEOUT;
28use crate::rmcp_client::ManagedClient;
29use crate::rmcp_client::StartupOutcomeError;
30use crate::runtime::McpRuntimeContext;
31use crate::server::EffectiveMcpServer;
32use crate::server::McpServerMetadata;
33use crate::tool_catalog_cache::McpToolCatalogCache;
34use crate::tools::ToolInfo;
35use anyhow::Context;
36use anyhow::Result;
37use anyhow::anyhow;
38use async_channel::Sender;
39use codex_api::SharedAuthProvider;
40use codex_config::Constrained;
41use codex_config::McpServerAuth;
42use codex_config::McpServerConfig;
43use codex_config::McpServerTransportConfig;
44use codex_config::types::AuthKeyringBackendKind;
45use codex_config::types::OAuthCredentialsStoreMode;
46use codex_connectors::ConnectorRuntimeContextKey;
47use codex_connectors::ConnectorRuntimeManager;
48use codex_login::AuthManager;
49use codex_login::CodexAuth;
50use codex_protocol::mcp::CallToolResult;
51use codex_protocol::mcp::McpServerInfo;
52use codex_protocol::models::PermissionProfile;
53use codex_protocol::protocol::AskForApproval;
54use codex_protocol::protocol::Event;
55use codex_protocol::protocol::EventMsg;
56use codex_protocol::protocol::McpStartupCompleteEvent;
57use codex_protocol::protocol::McpStartupFailure;
58use codex_protocol::protocol::McpStartupFailureReason;
59use codex_protocol::protocol::McpStartupStatus;
60use codex_protocol::protocol::McpStartupUpdateEvent;
61use codex_rmcp_client::ElicitationResponse;
62use codex_rmcp_client::McpAuthState;
63use codex_rmcp_client::McpLoginRequirement;
64use codex_rmcp_client::determine_streamable_http_auth_status_from_credentials;
65use rmcp::model::ElicitationCapability;
66use rmcp::model::ListResourceTemplatesResult;
67use rmcp::model::ListResourcesResult;
68use rmcp::model::PaginatedRequestParams;
69use rmcp::model::ReadResourceRequestParams;
70use rmcp::model::ReadResourceResult;
71use rmcp::model::RequestId;
72use rmcp::model::Resource;
73use rmcp::model::ResourceTemplate;
74use serde_json::Value as JsonValue;
75use tokio::sync::Mutex;
76use tokio::sync::RwLock;
77use tokio::task::JoinSet;
78use tokio_util::sync::CancellationToken;
79use tracing::warn;
80
81const MCP_UI_META_KEY: &str = "ui";
82const MCP_UI_VISIBILITY_META_KEY: &str = "visibility";
83const MCP_UI_MODEL_VISIBILITY: &str = "model";
84
85/// Returns whether a tool may be included in model-facing tool declarations.
86///
87/// Tools without visibility metadata remain visible.
88/// Tools with visibility metadata are hidden unless they explicitly include `model`.
89///
90/// <https://github.com/modelcontextprotocol/ext-apps/blob/main/specification/2026-01-26/apps.mdx#resource-discovery>
91pub fn tool_is_model_visible(tool: &ToolInfo) -> bool {
92    let Some(visibility) = tool
93        .tool
94        .meta
95        .as_deref()
96        .and_then(|meta| meta.get(MCP_UI_META_KEY))
97        .and_then(JsonValue::as_object)
98        .and_then(|ui| ui.get(MCP_UI_VISIBILITY_META_KEY))
99        .and_then(JsonValue::as_array)
100    else {
101        return true;
102    };
103
104    visibility
105        .iter()
106        .any(|target| target.as_str() == Some(MCP_UI_MODEL_VISIBILITY))
107}
108
109/// A thin wrapper around a set of running [`RmcpClient`] instances.
110pub struct McpConnectionManager {
111    clients: HashMap<String, AsyncManagedClient>,
112    server_metadata: HashMap<String, McpServerMetadata>,
113    required_servers: Vec<String>,
114    tool_catalog_revision: Arc<RwLock<u64>>,
115    codex_apps_tools_override: RwLock<Option<Vec<ToolInfo>>>,
116    codex_apps_refresh_lock: Mutex<()>,
117    tool_plugin_provenance: Arc<ToolPluginProvenance>,
118    prefix_mcp_tool_names: bool,
119    elicitation_requests: ElicitationRequestManager,
120    startup_cancellation_token: CancellationToken,
121}
122
123impl McpConnectionManager {
124    /// Creates an MCP connection manager. Threadless callers can pass no `tx_event`; startup
125    /// notifications are then skipped and interactive elicitations are declined.
126    #[allow(clippy::too_many_arguments)]
127    pub async fn new(
128        mcp_servers: &HashMap<String, EffectiveMcpServer>,
129        store_mode: OAuthCredentialsStoreMode,
130        keyring_backend_kind: AuthKeyringBackendKind,
131        approval_policy: &Constrained<AskForApproval>,
132        submit_id: String,
133        tx_event: Option<Sender<Event>>,
134        startup_cancellation_token: CancellationToken,
135        initial_permission_profile: PermissionProfile,
136        runtime_context: McpRuntimeContext,
137        codex_home: PathBuf,
138        codex_apps_tools_cache: ConnectorRuntimeManager<ToolInfo>,
139        tool_catalog_cache: McpToolCatalogCache,
140        codex_apps_tools_cache_key: ConnectorRuntimeContextKey,
141        prefix_mcp_tool_names: bool,
142        client_elicitation_capability: ElicitationCapability,
143        supports_openai_form_elicitation: bool,
144        tool_plugin_provenance: ToolPluginProvenance,
145        auth: Option<&CodexAuth>,
146        codex_apps_auth_manager: Option<Arc<AuthManager>>,
147        elicitation_reviewer: Option<ElicitationReviewerHandle>,
148        elicitation_lifecycle: Option<crate::ElicitationLifecycle>,
149        elicitation_router: ElicitationRequestRouter,
150    ) -> Self {
151        let mut required_servers = mcp_servers
152            .iter()
153            .filter(|(_, server)| server.enabled() && server.required())
154            .map(|(name, _)| name.clone())
155            .collect::<Vec<_>>();
156        required_servers.sort();
157        let mut clients = HashMap::new();
158        let mut server_metadata = HashMap::new();
159        let mut join_set = JoinSet::new();
160        let elicitation_requests = ElicitationRequestManager::new(
161            approval_policy.value(),
162            initial_permission_profile,
163            elicitation_reviewer,
164            elicitation_lifecycle,
165            elicitation_router,
166        );
167        let tool_plugin_provenance = Arc::new(tool_plugin_provenance);
168        let startup_submit_id = submit_id.clone();
169        let static_chatgpt_auth_provider = auth
170            .filter(|auth| auth.uses_codex_backend())
171            .map(codex_model_provider::auth_provider_from_auth);
172        let codex_apps_auth_provider = codex_apps_auth_manager.and_then(|auth_manager| {
173            auth.filter(|auth| auth.uses_codex_backend()).map(|auth| {
174                codex_model_provider::auth_provider_from_auth_manager(auth_manager, auth)
175            })
176        });
177        let mcp_servers = mcp_servers.clone();
178        for (server_name, server) in mcp_servers
179            .into_iter()
180            .filter(|(_, server)| server.enabled())
181        {
182            server_metadata.insert(server_name.clone(), McpServerMetadata::from(&server));
183            let cancel_token = startup_cancellation_token.child_token();
184            if let Some(tx_event) = tx_event.as_ref() {
185                let _ = emit_update(
186                    startup_submit_id.as_str(),
187                    tx_event,
188                    McpStartupUpdateEvent {
189                        server: server_name.clone(),
190                        status: McpStartupStatus::Starting,
191                    },
192                )
193                .await;
194            }
195            let configured_config = server.configured_config().cloned();
196            let resolved_environment = configured_config.as_ref().map_or_else(
197                || Ok(None),
198                |config| runtime_context.resolve_server_environment(&server_name, config),
199            );
200            // For built-in Codex Apps, `CODEX_CONNECTORS_TOKEN` is a debug
201            // override: it supplies runtime auth but bypasses the shared tools
202            // cache.
203            let uses_env_bearer_token =
204                configured_config
205                    .as_ref()
206                    .is_some_and(|config| match &config.transport {
207                        McpServerTransportConfig::StreamableHttp {
208                            bearer_token_env_var,
209                            ..
210                        } => bearer_token_env_var.is_some(),
211                        McpServerTransportConfig::Stdio { .. } => false,
212                    });
213            let shares_codex_apps_tools_cache =
214                should_share_codex_apps_tools_cache(&server_name, uses_env_bearer_token);
215            let codex_apps_tools_cache_context = shares_codex_apps_tools_cache.then(|| {
216                codex_apps_tools_cache
217                    .context(codex_home.clone(), codex_apps_tools_cache_key.clone())
218            });
219            // The reserved Codex Apps registration follows the shared
220            // AuthManager across refreshes. In the hosted-plugin path, this
221            // is the ChatGPT /ps/mcp connection. User-configured MCP
222            // registrations keep their existing configured auth path.
223            let chatgpt_auth_provider = if server_name == CODEX_APPS_MCP_SERVER_NAME {
224                codex_apps_auth_provider
225                    .clone()
226                    .or_else(|| static_chatgpt_auth_provider.clone())
227            } else {
228                static_chatgpt_auth_provider.clone()
229            };
230            // If Codex Apps has an env bearer token, that is its auth path. Do
231            // not also attach the ambient CodexAuth provider.
232            let runtime_auth_provider =
233                if server_name == CODEX_APPS_MCP_SERVER_NAME && uses_env_bearer_token {
234                    None
235                } else {
236                    chatgpt_auth_provider_for_server(&server, chatgpt_auth_provider)
237                };
238            let tool_catalog_cache_context = if server_name == CODEX_APPS_MCP_SERVER_NAME {
239                None
240            } else if let Some(config) = configured_config.as_ref()
241                && let Ok(environment) = resolved_environment.as_ref()
242            {
243                tool_catalog_cache.context(
244                    &server_name,
245                    config,
246                    &runtime_context,
247                    environment.as_ref(),
248                    &client_elicitation_capability,
249                    supports_openai_form_elicitation,
250                )
251            } else {
252                None
253            };
254            let has_runtime_auth = runtime_auth_provider.is_some();
255            let async_managed_client = AsyncManagedClient::new(
256                server_name.clone(),
257                startup_submit_id.clone(),
258                server,
259                store_mode,
260                keyring_backend_kind,
261                cancel_token.clone(),
262                tx_event.clone(),
263                elicitation_requests.clone(),
264                codex_apps_tools_cache_context,
265                tool_catalog_cache_context,
266                Arc::clone(&tool_plugin_provenance),
267                runtime_context.clone(),
268                resolved_environment,
269                runtime_auth_provider,
270                client_elicitation_capability.clone(),
271                supports_openai_form_elicitation,
272            );
273            clients.insert(server_name.clone(), async_managed_client.clone());
274            let tx_event = tx_event.clone();
275            let submit_id = startup_submit_id.clone();
276            join_set.spawn(async move {
277                let mut outcome = async_managed_client.client().await;
278                if cancel_token.is_cancelled() {
279                    outcome = Err(StartupOutcomeError::Cancelled);
280                }
281                if let Some(tx_event) = tx_event.as_ref() {
282                    let auth_state = match &outcome {
283                        Err(error) if error.is_authentication_required() && !has_runtime_auth => {
284                            configured_config.as_ref().and_then(|config| {
285                                let McpServerTransportConfig::StreamableHttp {
286                                    url,
287                                    bearer_token_env_var,
288                                    http_headers,
289                                    env_http_headers,
290                                } = &config.transport
291                                else {
292                                    return None;
293                                };
294                                match determine_streamable_http_auth_status_from_credentials(
295                                    &server_name,
296                                    url,
297                                    bearer_token_env_var.as_deref(),
298                                    http_headers.clone(),
299                                    env_http_headers.clone(),
300                                    store_mode,
301                                    keyring_backend_kind,
302                                ) {
303                                    Ok(auth_state) => auth_state,
304                                    Err(error) => {
305                                        warn!(
306                                            "failed to read stored auth status for MCP server `{server_name}`: {error:?}"
307                                        );
308                                        None
309                                    }
310                                }
311                            })
312                        }
313                        Ok(_) | Err(_) => None,
314                    };
315                    if cancel_token.is_cancelled() {
316                        outcome = Err(StartupOutcomeError::Cancelled);
317                    }
318                    let status = match &outcome {
319                        Ok(_) => McpStartupStatus::Ready,
320                        Err(StartupOutcomeError::Cancelled) => McpStartupStatus::Cancelled,
321                        Err(error) => {
322                            let reason = mcp_startup_failure_reason(auth_state, error);
323                            let error_str = mcp_init_error_display(
324                                server_name.as_str(),
325                                configured_config.as_ref(),
326                                error,
327                            );
328                            McpStartupStatus::Failed {
329                                error: error_str,
330                                reason,
331                            }
332                        }
333                    };
334
335                    let _ = emit_update(
336                        submit_id.as_str(),
337                        tx_event,
338                        McpStartupUpdateEvent {
339                            server: server_name.clone(),
340                            status,
341                        },
342                    )
343                    .await;
344                }
345                if cancel_token.is_cancelled() {
346                    outcome = Err(StartupOutcomeError::Cancelled);
347                }
348
349                if matches!(&outcome, Err(StartupOutcomeError::Failed { .. })) {
350                    async_managed_client.reconnect_failed_startup().await;
351                }
352
353                (server_name, outcome)
354            });
355        }
356        let manager = Self {
357            clients,
358            server_metadata,
359            required_servers,
360            tool_catalog_revision: Arc::new(RwLock::new(0)),
361            codex_apps_tools_override: RwLock::new(None),
362            codex_apps_refresh_lock: Mutex::new(()),
363            tool_plugin_provenance,
364            prefix_mcp_tool_names,
365            elicitation_requests: elicitation_requests.clone(),
366            startup_cancellation_token: startup_cancellation_token.clone(),
367        };
368        tokio::spawn(async move {
369            let outcomes = join_set.join_all().await;
370            if let Some(tx_event) = tx_event {
371                let mut summary = McpStartupCompleteEvent::default();
372                for (server_name, outcome) in outcomes {
373                    match outcome {
374                        Ok(_) => summary.ready.push(server_name),
375                        Err(StartupOutcomeError::Cancelled) => summary.cancelled.push(server_name),
376                        Err(StartupOutcomeError::Failed { error, .. }) => {
377                            summary.failed.push(McpStartupFailure {
378                                server: server_name,
379                                error,
380                            })
381                        }
382                    }
383                }
384                let _ = tx_event
385                    .send(Event {
386                        id: startup_submit_id,
387                        msg: EventMsg::McpStartupComplete(summary),
388                    })
389                    .await;
390            }
391        });
392        manager
393    }
394
395    pub fn new_uninitialized_with_permission_profile(
396        approval_policy: &Constrained<AskForApproval>,
397        permission_profile: &PermissionProfile,
398        prefix_mcp_tool_names: bool,
399    ) -> Self {
400        Self {
401            clients: HashMap::new(),
402            server_metadata: HashMap::new(),
403            required_servers: Vec::new(),
404            tool_catalog_revision: Arc::new(RwLock::new(0)),
405            codex_apps_tools_override: RwLock::new(None),
406            codex_apps_refresh_lock: Mutex::new(()),
407            tool_plugin_provenance: Arc::new(ToolPluginProvenance::default()),
408            prefix_mcp_tool_names,
409            elicitation_requests: ElicitationRequestManager::new(
410                approval_policy.value(),
411                permission_profile.clone(),
412                /*reviewer*/ None,
413                /*lifecycle*/ None,
414                ElicitationRequestRouter::default(),
415            ),
416            startup_cancellation_token: CancellationToken::new(),
417        }
418    }
419
420    pub fn empty(prefix_mcp_tool_names: bool) -> Self {
421        Self::new_uninitialized_with_permission_profile(
422            &Constrained::allow_any(AskForApproval::OnRequest),
423            &PermissionProfile::default(),
424            prefix_mcp_tool_names,
425        )
426    }
427
428    pub fn has_servers(&self) -> bool {
429        !self.clients.is_empty()
430    }
431
432    pub(crate) fn contains_server(&self, server_name: &str) -> bool {
433        self.clients.contains_key(server_name)
434    }
435
436    /// Stop all MCP clients owned by this manager and terminate stdio server processes.
437    pub async fn shutdown(&self) {
438        self.startup_cancellation_token.cancel();
439        let clients = self.clients.values().cloned().collect::<Vec<_>>();
440        // Keep cleanup alive if an interrupt cancels the refresh that requested it.
441        let shutdown_task = tokio::spawn(async move {
442            for client in clients {
443                client.shutdown().await;
444            }
445        });
446        if let Err(error) = shutdown_task.await {
447            warn!("MCP client shutdown task failed: {error}");
448        }
449    }
450
451    pub fn server_origin(&self, server_name: &str) -> Option<&str> {
452        self.server_metadata
453            .get(server_name)
454            .and_then(|metadata| metadata.origin.as_ref())
455            .map(super::server::McpServerOrigin::as_str)
456    }
457
458    pub fn server_environment_id(&self, server_name: &str) -> Option<&str> {
459        self.server_metadata
460            .get(server_name)
461            .map(|metadata| metadata.environment_id.as_str())
462    }
463
464    pub fn server_pollutes_memory(&self, server_name: &str) -> bool {
465        self.server_metadata
466            .get(server_name)
467            .is_none_or(|metadata| metadata.pollutes_memory)
468    }
469
470    pub fn plugin_id_for_mcp_server_name(&self, server_name: &str) -> Option<&str> {
471        self.tool_plugin_provenance
472            .plugin_id_for_mcp_server_name(server_name)
473    }
474
475    pub fn is_selected_plugin_mcp_server(&self, server_name: &str) -> bool {
476        self.tool_plugin_provenance
477            .is_selected_plugin_mcp_server(server_name)
478    }
479
480    pub fn tool_approval_mode(
481        &self,
482        server_name: &str,
483        tool_name: &str,
484    ) -> codex_config::AppToolApproval {
485        self.server_metadata
486            .get(server_name)
487            .map(|metadata| metadata.tool_approval_mode(tool_name))
488            .unwrap_or_default()
489    }
490
491    pub fn is_host_owned_codex_apps_server(&self, server_name: &str) -> bool {
492        server_name == CODEX_APPS_MCP_SERVER_NAME && self.server_metadata.contains_key(server_name)
493    }
494
495    pub fn set_approval_policy(&self, approval_policy: &Constrained<AskForApproval>) {
496        if let Ok(mut policy) = self.elicitation_requests.approval_policy.lock() {
497            *policy = approval_policy.value();
498        }
499    }
500
501    pub fn set_permission_profile(&self, permission_profile: PermissionProfile) {
502        if let Ok(mut profile) = self.elicitation_requests.permission_profile.lock() {
503            *profile = permission_profile;
504        }
505    }
506
507    pub fn elicitations_auto_deny(&self) -> bool {
508        self.elicitation_requests.auto_deny()
509    }
510
511    pub fn set_elicitations_auto_deny(&self, auto_deny: bool) {
512        self.elicitation_requests.set_auto_deny(auto_deny);
513    }
514
515    pub fn elicitation_router(&self) -> ElicitationRequestRouter {
516        self.elicitation_requests.router()
517    }
518
519    pub async fn resolve_elicitation(
520        &self,
521        server_name: String,
522        id: RequestId,
523        response: ElicitationResponse,
524    ) -> Result<()> {
525        self.elicitation_requests
526            .resolve(server_name, id, response)
527            .await
528    }
529
530    pub async fn wait_for_server_ready(&self, server_name: &str, timeout: Duration) -> bool {
531        let Some(async_managed_client) = self.clients.get(server_name) else {
532            return false;
533        };
534
535        match tokio::time::timeout(timeout, async_managed_client.client()).await {
536            Ok(Ok(_)) => true,
537            Ok(Err(_)) | Err(_) => false,
538        }
539    }
540
541    /// Returns resources from servers selected by `include_server`. Each key
542    /// is the server name and the value is a vector of resources.
543    pub async fn list_all_resources(
544        &self,
545        include_server: impl Fn(&str) -> bool,
546    ) -> HashMap<String, Vec<Resource>> {
547        self.ready_clients_matching(&include_server)
548            .await
549            .list_all_resources(|_| true)
550            .await
551    }
552
553    /// Returns resource templates from servers selected by `include_server`.
554    /// Each key is the server name and the value is a vector of templates.
555    pub async fn list_all_resource_templates(
556        &self,
557        include_server: impl Fn(&str) -> bool,
558    ) -> HashMap<String, Vec<ResourceTemplate>> {
559        self.ready_clients_matching(&include_server)
560            .await
561            .list_all_resource_templates(|_| true)
562            .await
563    }
564
565    async fn ready_clients_matching(
566        &self,
567        include_server: &impl Fn(&str) -> bool,
568    ) -> McpBindingClients {
569        let mut clients = HashMap::new();
570        for (server, client) in self
571            .clients
572            .iter()
573            .filter(|(server, _)| include_server(server))
574        {
575            if let Ok(client) = client.client().await {
576                clients.insert(server.clone(), Arc::new(client));
577            }
578        }
579        McpBindingClients::new(clients)
580    }
581
582    /// Invoke the tool indicated by the (server, tool) pair.
583    pub async fn call_tool(
584        &self,
585        server: &str,
586        tool: &str,
587        arguments: Option<serde_json::Value>,
588        meta: Option<serde_json::Value>,
589    ) -> Result<CallToolResult> {
590        let client = self.client_by_name(server).await?;
591        if !client.tool_filter.allows(tool) {
592            return Err(anyhow!(
593                "tool '{tool}' is disabled for MCP server '{server}'"
594            ));
595        }
596
597        let result: rmcp::model::CallToolResult = client
598            .client
599            .call_tool(tool.to_string(), arguments, meta, client.tool_timeout)
600            .await
601            .with_context(|| format!("tool call failed for `{server}/{tool}`"))?;
602
603        let content = result
604            .content
605            .into_iter()
606            .map(|content| {
607                serde_json::to_value(content)
608                    .unwrap_or_else(|_| serde_json::Value::String("<content>".to_string()))
609            })
610            .collect();
611
612        Ok(CallToolResult {
613            content,
614            structured_content: result.structured_content,
615            is_error: result.is_error,
616            meta: result.meta.and_then(|meta| serde_json::to_value(meta).ok()),
617        })
618    }
619
620    pub async fn server_supports_sandbox_state_meta_capability(
621        &self,
622        server: &str,
623    ) -> Result<bool> {
624        Ok(self
625            .client_by_name(server)
626            .await?
627            .server_supports_sandbox_state_meta_capability)
628    }
629
630    /// List resources from the specified server.
631    pub async fn list_resources(
632        &self,
633        server: &str,
634        params: Option<PaginatedRequestParams>,
635    ) -> Result<ListResourcesResult> {
636        let managed = self.client_by_name(server).await?;
637        let timeout = managed.tool_timeout;
638
639        managed
640            .client
641            .list_resources(params, timeout)
642            .await
643            .with_context(|| format!("resources/list failed for `{server}`"))
644    }
645
646    /// List resource templates from the specified server.
647    pub async fn list_resource_templates(
648        &self,
649        server: &str,
650        params: Option<PaginatedRequestParams>,
651    ) -> Result<ListResourceTemplatesResult> {
652        let managed = self.client_by_name(server).await?;
653        let client = managed.client.clone();
654        let timeout = managed.tool_timeout;
655
656        client
657            .list_resource_templates(params, timeout)
658            .await
659            .with_context(|| format!("resources/templates/list failed for `{server}`"))
660    }
661
662    /// Read a resource from the specified server.
663    pub async fn read_resource(
664        &self,
665        server: &str,
666        params: ReadResourceRequestParams,
667    ) -> Result<ReadResourceResult> {
668        let managed = self.client_by_name(server).await?;
669        let client = managed.client.clone();
670        let timeout = managed.tool_timeout;
671        let uri = params.uri.clone();
672
673        client
674            .read_resource(params, timeout)
675            .await
676            .with_context(|| format!("resources/read failed for `{server}` ({uri})"))
677    }
678
679    /// Returns presentation metadata from the current connection.
680    /// Codex Apps metadata may come from its existing cache; regular MCP server information is
681    /// connection-specific, so pending regular clients are awaited.
682    pub(crate) async fn list_available_server_infos(&self) -> HashMap<String, McpServerInfo> {
683        let mut server_infos = HashMap::new();
684        for (server_name, client) in &self.clients {
685            if !client.startup_complete.load(Ordering::Acquire)
686                && let Some(server_info) = client.cached_server_info.clone()
687            {
688                server_infos.insert(server_name.clone(), server_info);
689                continue;
690            }
691            match client.client().await {
692                Ok(managed_client) => {
693                    server_infos.insert(server_name.clone(), managed_client.server_info);
694                }
695                Err(_) => {
696                    if let Some(server_info) = client.cached_server_info.clone() {
697                        server_infos.insert(server_name.clone(), server_info);
698                    }
699                }
700            }
701        }
702        server_infos
703    }
704
705    async fn client_by_name(&self, name: &str) -> Result<ManagedClient> {
706        self.clients
707            .get(name)
708            .ok_or_else(|| anyhow!("unknown MCP server '{name}'"))?
709            .client()
710            .await
711            .context("failed to get client")
712    }
713
714    #[cfg(test)]
715    fn new_uninitialized(
716        approval_policy: &Constrained<AskForApproval>,
717        permission_profile: &Constrained<PermissionProfile>,
718        prefix_mcp_tool_names: bool,
719    ) -> Self {
720        Self::new_uninitialized_with_permission_profile(
721            approval_policy,
722            permission_profile.get(),
723            prefix_mcp_tool_names,
724        )
725    }
726}
727
728impl Drop for McpConnectionManager {
729    fn drop(&mut self) {
730        self.startup_cancellation_token.cancel();
731        self.clients.clear();
732    }
733}
734
735/// Makes ChatGPT authentication available to servers that explicitly opt in.
736/// The HTTP transport applies it only when no configured authorization resolves.
737fn chatgpt_auth_provider_for_server(
738    server: &EffectiveMcpServer,
739    chatgpt_auth_provider: Option<SharedAuthProvider>,
740) -> Option<SharedAuthProvider> {
741    if !server
742        .configured_config()
743        .is_some_and(|config| matches!(&config.auth, McpServerAuth::ChatGpt))
744    {
745        return None;
746    }
747    chatgpt_auth_provider
748}
749
750fn should_share_codex_apps_tools_cache(server_name: &str, uses_env_bearer_token: bool) -> bool {
751    server_name == CODEX_APPS_MCP_SERVER_NAME && !uses_env_bearer_token
752}
753
754async fn emit_update(
755    submit_id: &str,
756    tx_event: &Sender<Event>,
757    update: McpStartupUpdateEvent,
758) -> Result<(), async_channel::SendError<Event>> {
759    tx_event
760        .send(Event {
761            id: submit_id.to_string(),
762            msg: EventMsg::McpStartupUpdate(update),
763        })
764        .await
765}
766
767fn mcp_startup_failure_reason(
768    auth_state: Option<McpAuthState>,
769    error: &StartupOutcomeError,
770) -> Option<McpStartupFailureReason> {
771    if !error.is_authentication_required() {
772        return None;
773    }
774
775    match auth_state {
776        Some(McpAuthState::LoggedOut(McpLoginRequirement::Reauthentication)) => {
777            Some(McpStartupFailureReason::ReauthenticationRequired)
778        }
779        Some(
780            McpAuthState::Unsupported
781            | McpAuthState::LoggedOut(McpLoginRequirement::Login)
782            | McpAuthState::BearerToken
783            | McpAuthState::OAuth,
784        )
785        | None => None,
786    }
787}
788
789fn mcp_init_error_display(
790    server_name: &str,
791    config: Option<&McpServerConfig>,
792    err: &StartupOutcomeError,
793) -> String {
794    if let Some(McpServerTransportConfig::StreamableHttp {
795        url,
796        bearer_token_env_var,
797        http_headers,
798        ..
799    }) = config.map(|config| &config.transport)
800        && url == "https://api.githubcopilot.com/mcp/"
801        && bearer_token_env_var.is_none()
802        && http_headers.as_ref().map(HashMap::is_empty).unwrap_or(true)
803    {
804        format!(
805            "GitHub MCP does not support OAuth. Log in by adding a personal access token (https://github.com/settings/personal-access-tokens) to your environment and config.toml:\n[mcp_servers.{server_name}]\nbearer_token_env_var = CODEX_GITHUB_PERSONAL_ACCESS_TOKEN"
806        )
807    } else if is_mcp_client_auth_required_error(err) {
808        format!(
809            "The {server_name} MCP server is not logged in. Run `codex mcp login {server_name}`."
810        )
811    } else if is_mcp_client_startup_timeout_error(err) {
812        let startup_timeout_secs = config
813            .and_then(|config| config.startup_timeout_sec)
814            .unwrap_or(DEFAULT_STARTUP_TIMEOUT)
815            .as_secs();
816        format!(
817            "MCP client for `{server_name}` timed out after {startup_timeout_secs} seconds. Add or adjust `startup_timeout_sec` in your config.toml:\n[mcp_servers.{server_name}]\nstartup_timeout_sec = XX"
818        )
819    } else {
820        format!("MCP client for `{server_name}` failed to start: {err:#}")
821    }
822}
823
824fn is_mcp_client_auth_required_error(error: &StartupOutcomeError) -> bool {
825    match error {
826        StartupOutcomeError::Failed { error, .. } => error.contains("Auth required"),
827        _ => false,
828    }
829}
830
831fn is_mcp_client_startup_timeout_error(error: &StartupOutcomeError) -> bool {
832    match error {
833        StartupOutcomeError::Failed { error, .. } => {
834            error.contains("request timed out")
835                || error.contains("timed out handshaking with MCP server")
836                || error.contains("MCP client startup timed out")
837        }
838        _ => false,
839    }
840}
841
842#[cfg(test)]
843#[path = "connection_manager_tests.rs"]
844mod tests;