1#[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
85pub 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
109pub 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 #[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 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 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 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 None,
413 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 pub async fn shutdown(&self) {
438 self.startup_cancellation_token.cancel();
439 let clients = self.clients.values().cloned().collect::<Vec<_>>();
440 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 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 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 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 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 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 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 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
735fn 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;