Download codex-rs/codex-mcp/src/connection_manager_tests.rs from SaylorTwift/codex: direct link, hf CLI and curl.
- Browser
- Download file 228 kB
-
https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/codex-mcp/src/connection_manager_tests.rs
- Command line
-
hf download hf://SaylorTwift/codex/codex-rs/codex-mcp/src/connection_manager_tests.rs
-
curl -L -o connection_manager_tests.rs https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/codex-mcp/src/connection_manager_tests.rs
228 kB
| use super::*; | |
| use crate::McpBinding; | |
| use crate::client_tool_catalog::ClientToolCatalog; | |
| use crate::elicitation::ElicitationLifecycle; | |
| use crate::elicitation::ElicitationRequestManager; | |
| use crate::elicitation::ElicitationRequestRouter; | |
| use crate::elicitation::ElicitationReviewRequest; | |
| use crate::elicitation::ElicitationReviewer; | |
| use crate::elicitation::elicitation_is_rejected_by_policy; | |
| use crate::mcp::tests::test_elicitation_config; | |
| use crate::rmcp_client::AsyncManagedClient; | |
| use crate::rmcp_client::CODEX_APPS_RECONNECT_INITIAL_BACKOFF; | |
| use crate::rmcp_client::CodexAppsStartupReconnect; | |
| use crate::rmcp_client::ManagedClient; | |
| use crate::rmcp_client::ManagedClientFuture; | |
| use crate::rmcp_client::StartupOutcomeError; | |
| use crate::rmcp_client::list_tools_for_client_uncached; | |
| use crate::runtime::McpRuntimeContext; | |
| use crate::server::EffectiveMcpServer; | |
| use crate::server::McpServerMetadata; | |
| use crate::server::McpServerOrigin; | |
| use crate::tool_catalog_cache::McpToolCatalogCache; | |
| use crate::tools::ToolFilter; | |
| use crate::tools::ToolInfo; | |
| use crate::tools::filter_tools; | |
| use crate::tools::normalize_tools_for_model_with_prefix; | |
| use assert_matches::assert_matches; | |
| use codex_config::AppToolApproval; | |
| use codex_config::Constrained; | |
| use codex_config::McpServerAuth; | |
| use codex_config::McpServerConfig; | |
| use codex_config::McpServerEnvVar; | |
| use codex_config::McpServerToolConfig; | |
| use codex_config::types::AuthKeyringBackendKind; | |
| use codex_config::types::OAuthCredentialsStoreMode; | |
| use codex_connectors::ConnectorRuntimeContext; | |
| use codex_connectors::ConnectorRuntimeContextKey; | |
| use codex_connectors::ConnectorRuntimeFetchSource; | |
| use codex_connectors::ConnectorRuntimeManager; | |
| use codex_exec_server::ExecServerError; | |
| use codex_exec_server::HttpClient; | |
| use codex_exec_server::HttpRequestParams; | |
| use codex_exec_server::HttpRequestResponse; | |
| use codex_exec_server::HttpResponseBodyStream; | |
| use codex_exec_server_test_support::environment_manager_without_environments; | |
| use codex_login::AuthHeaders; | |
| use codex_login::AuthManager; | |
| use codex_login::CodexAuth; | |
| use codex_protocol::ToolName; | |
| use codex_protocol::approvals::ElicitationRequest; | |
| use codex_protocol::mcp::ClientMcpExtensions; | |
| use codex_protocol::mcp::McpServerInfo; | |
| use codex_protocol::models::PermissionProfile; | |
| use codex_protocol::protocol::AskForApproval; | |
| use codex_protocol::protocol::GranularApprovalConfig; | |
| use codex_protocol::protocol::McpStartupFailureReason; | |
| use codex_rmcp_client::ElicitationResponse; | |
| use codex_rmcp_client::InProcessTransportFactory; | |
| use codex_rmcp_client::McpAuthState; | |
| use codex_rmcp_client::McpLoginRequirement; | |
| use codex_rmcp_client::McpOAuthRefreshMode; | |
| use codex_rmcp_client::RmcpClient; | |
| use codex_utils_path_uri::PathUri; | |
| use futures::FutureExt; | |
| use futures::future::BoxFuture; | |
| use pretty_assertions::assert_eq; | |
| use rmcp::ErrorData as McpError; | |
| use rmcp::RoleServer; | |
| use rmcp::ServerHandler; | |
| use rmcp::ServiceExt; | |
| use rmcp::model::ClientCapabilities; | |
| use rmcp::model::ElicitRequestParams; | |
| use rmcp::model::ElicitationAction; | |
| use rmcp::model::ElicitationCapability; | |
| use rmcp::model::Implementation; | |
| use rmcp::model::InitializeRequestParams; | |
| use rmcp::model::JsonObject; | |
| use rmcp::model::ListToolsResult; | |
| use rmcp::model::NumberOrString; | |
| use rmcp::model::PaginatedRequestParams; | |
| use rmcp::model::ProtocolVersion; | |
| use rmcp::model::ServerCapabilities; | |
| use rmcp::model::ServerInfo; | |
| use rmcp::model::Tool; | |
| use rmcp::service::RequestContext; | |
| use std::collections::HashMap; | |
| use std::collections::HashSet; | |
| use std::io; | |
| use std::path::PathBuf; | |
| use std::sync::Arc; | |
| use std::sync::atomic::AtomicBool; | |
| use std::sync::atomic::AtomicUsize; | |
| use std::sync::atomic::Ordering; | |
| use tempfile::tempdir; | |
| use tokio::io::DuplexStream; | |
| use tokio::sync::Notify; | |
| use tokio_util::sync::CancellationToken; | |
| impl McpConnectionSet { | |
| fn new_uninitialized( | |
| approval_policy: &Constrained<AskForApproval>, | |
| permission_profile: &Constrained<PermissionProfile>, | |
| prefix_mcp_tool_names: bool, | |
| ) -> Self { | |
| Self { | |
| servers: HashMap::new(), | |
| event_stream_connection: None, | |
| disabled_servers: Vec::new(), | |
| required_servers: Vec::new(), | |
| optional_startup_deadline: OnceLock::new(), | |
| tool_plugin_context: Arc::new(ToolPluginContext::default()), | |
| prefix_mcp_tool_names, | |
| non_prefixed_mcp_tool_servers: Vec::new(), | |
| elicitation_requests: ElicitationRequestManager::new( | |
| test_elicitation_config( | |
| "server", | |
| approval_policy.value(), | |
| permission_profile.get().clone(), | |
| ), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| ElicitationRequestRouter::default(), | |
| ), | |
| trusted_access: None, | |
| } | |
| } | |
| pub(crate) fn insert_test_client( | |
| &mut self, | |
| name: impl Into<String>, | |
| client: AsyncManagedClient, | |
| ) { | |
| let name = name.into(); | |
| self.servers.insert( | |
| name, | |
| McpServerView { | |
| tool_filter: ToolFilter::default(), | |
| protocol_mode: crate::McpProtocolMode::Legacy, | |
| connection: Arc::new(McpServerConnection { | |
| identity: None, | |
| client, | |
| startup_timeout: DEFAULT_STARTUP_TIMEOUT, | |
| startup_trigger: None, | |
| _diagnostics_guard: LIVE_CONNECTIONS.track(), | |
| }), | |
| metadata: McpServerMetadata { | |
| environment_id: String::new(), | |
| pollutes_memory: true, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| tool_timeout: None, | |
| catalog_item_limit: crate::pagination::MAX_MCP_CATALOG_ITEMS, | |
| }, | |
| ); | |
| } | |
| fn test_client(&self, name: &str) -> &AsyncManagedClient { | |
| &self.servers[name].connection.client | |
| } | |
| fn set_test_server_metadata(&mut self, name: &str, metadata: McpServerMetadata) { | |
| self.servers | |
| .get_mut(name) | |
| .expect("test server exists") | |
| .metadata = metadata; | |
| } | |
| fn shares_test_connection_with(&self, other: &Self, name: &str) -> bool { | |
| let Some(left) = self.servers.get(name) else { | |
| return false; | |
| }; | |
| let Some(right) = other.servers.get(name) else { | |
| return false; | |
| }; | |
| Arc::ptr_eq(&left.connection, &right.connection) | |
| } | |
| } | |
| fn create_test_tool(server_name: &str, tool_name: &str) -> ToolInfo { | |
| ToolInfo { | |
| server_name: server_name.to_string(), | |
| supports_parallel_tool_calls: false, | |
| server_origin: None, | |
| callable_name: tool_name.to_string(), | |
| callable_namespace: server_name.to_string(), | |
| namespace_description: None, | |
| tool: Tool::new( | |
| tool_name.to_string(), | |
| format!("Test tool: {tool_name}"), | |
| Arc::new(JsonObject::default()), | |
| ), | |
| openai_file_input_optional_fields: Default::default(), | |
| connector_id: None, | |
| connector_name: None, | |
| plugin_display_names: Vec::new(), | |
| } | |
| } | |
| fn create_codex_apps_tools_cache_context( | |
| codex_home: PathBuf, | |
| account_id: Option<&str>, | |
| chatgpt_user_id: Option<&str>, | |
| ) -> ConnectorRuntimeContext<ToolInfo> { | |
| ConnectorRuntimeManager::<ToolInfo>::default().context( | |
| codex_home, | |
| ConnectorRuntimeContextKey::personal( | |
| account_id.map(ToOwned::to_owned), | |
| chatgpt_user_id.map(ToOwned::to_owned), | |
| ), | |
| ) | |
| } | |
| fn store_current_tools(cache_context: &ConnectorRuntimeContext<ToolInfo>, tools: Vec<ToolInfo>) { | |
| let _ = cache_context.publish_if_newest_accepted( | |
| cache_context.begin_fetch(ConnectorRuntimeFetchSource::HardRefresh), | |
| &create_test_server_info("Codex Apps"), | |
| tools, | |
| ); | |
| } | |
| async fn capture_binding(manager: &Arc<McpConnectionSet>) -> McpBinding { | |
| let mut config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| config.server_permission_profiles = manager | |
| .servers | |
| .keys() | |
| .map(|name| (name.clone(), PermissionProfile::default())) | |
| .collect(); | |
| manager | |
| .capture_binding_with_metadata( | |
| Arc::new(config), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new(), | |
| ) | |
| .await | |
| } | |
| fn create_test_server_info(title: &str) -> McpServerInfo { | |
| McpServerInfo { | |
| name: "codex-apps".to_string(), | |
| title: Some(title.to_string()), | |
| version: "1.0.0".to_string(), | |
| description: None, | |
| icons: None, | |
| website_url: None, | |
| } | |
| } | |
| struct TestInProcessTransportFactory; | |
| struct PendingHttpClient; | |
| impl HttpClient for PendingHttpClient { | |
| fn http_request( | |
| &self, | |
| _params: HttpRequestParams, | |
| ) -> BoxFuture<'_, Result<HttpRequestResponse, ExecServerError>> { | |
| futures::future::pending().boxed() | |
| } | |
| fn http_request_stream( | |
| &self, | |
| _params: HttpRequestParams, | |
| ) -> BoxFuture<'_, Result<(HttpRequestResponse, HttpResponseBodyStream), ExecServerError>> { | |
| futures::future::pending().boxed() | |
| } | |
| } | |
| impl InProcessTransportFactory for TestInProcessTransportFactory { | |
| fn open(&self) -> BoxFuture<'static, io::Result<DuplexStream>> { | |
| async { | |
| let (client_stream, _server_stream) = tokio::io::duplex(1); | |
| Ok(client_stream) | |
| } | |
| .boxed() | |
| } | |
| } | |
| struct RefreshTestTransportFactory { | |
| tool: Tool, | |
| list_started: Option<Arc<Notify>>, | |
| release_list: Option<Arc<Notify>>, | |
| next_cursor: Option<String>, | |
| list_requests: Arc<AtomicUsize>, | |
| } | |
| impl ServerHandler for RefreshTestTransportFactory { | |
| fn get_info(&self) -> ServerInfo { | |
| ServerInfo::new(ServerCapabilities::builder().enable_tools().build()) | |
| } | |
| async fn list_tools( | |
| &self, | |
| _request: Option<PaginatedRequestParams>, | |
| _context: rmcp::service::RequestContext<rmcp::service::RoleServer>, | |
| ) -> Result<ListToolsResult, McpError> { | |
| self.list_requests | |
| .fetch_add(1, std::sync::atomic::Ordering::SeqCst); | |
| if let Some(list_started) = &self.list_started { | |
| list_started.notify_one(); | |
| } | |
| if let Some(release_list) = &self.release_list { | |
| release_list.notified().await; | |
| } | |
| let mut result = ListToolsResult::with_all_items(vec![self.tool.clone()]); | |
| result.next_cursor = self.next_cursor.clone(); | |
| Ok(result) | |
| } | |
| } | |
| impl InProcessTransportFactory for RefreshTestTransportFactory { | |
| fn open(&self) -> BoxFuture<'static, io::Result<DuplexStream>> { | |
| let server = self.clone(); | |
| async move { | |
| let (client_stream, server_stream) = tokio::io::duplex(4096); | |
| tokio::spawn(async move { | |
| let server = server | |
| .serve(server_stream) | |
| .await | |
| .expect("serve test MCP server"); | |
| server.waiting().await.expect("wait for test MCP server"); | |
| }); | |
| Ok(client_stream) | |
| } | |
| .boxed() | |
| } | |
| } | |
| struct MutableToolsServer { | |
| tools: Arc<tokio::sync::RwLock<Vec<Tool>>>, | |
| block_tool_listing: Arc<AtomicBool>, | |
| } | |
| impl ServerHandler for MutableToolsServer { | |
| fn get_info(&self) -> ServerInfo { | |
| ServerInfo::new(ServerCapabilities::builder().enable_tools().build()) | |
| } | |
| async fn list_tools( | |
| &self, | |
| _request: Option<rmcp::model::PaginatedRequestParams>, | |
| _context: RequestContext<RoleServer>, | |
| ) -> Result<ListToolsResult, rmcp::ErrorData> { | |
| if self.block_tool_listing.load(Ordering::Acquire) { | |
| std::future::pending::<()>().await; | |
| } | |
| Ok(ListToolsResult { | |
| tools: self.tools.read().await.clone(), | |
| ..Default::default() | |
| }) | |
| } | |
| } | |
| struct MutableToolsTransportFactory { | |
| server: MutableToolsServer, | |
| } | |
| impl InProcessTransportFactory for MutableToolsTransportFactory { | |
| fn open(&self) -> BoxFuture<'static, io::Result<DuplexStream>> { | |
| let server = self.server.clone(); | |
| async move { | |
| let (client_stream, server_stream) = tokio::io::duplex(4096); | |
| tokio::spawn(async move { | |
| server | |
| .serve(server_stream) | |
| .await | |
| .expect("serve mutable MCP tools") | |
| .waiting() | |
| .await | |
| .expect("mutable MCP tools server completes"); | |
| }); | |
| Ok(client_stream) | |
| } | |
| .boxed() | |
| } | |
| } | |
| struct DisconnectingToolsTransportFactory { | |
| server: MutableToolsServer, | |
| disconnect: CancellationToken, | |
| } | |
| impl InProcessTransportFactory for DisconnectingToolsTransportFactory { | |
| fn open(&self) -> BoxFuture<'static, io::Result<DuplexStream>> { | |
| let server = self.server.clone(); | |
| let disconnect = self.disconnect.clone(); | |
| async move { | |
| let (client_stream, server_stream) = tokio::io::duplex(4096); | |
| tokio::spawn(async move { | |
| let server = server | |
| .serve(server_stream) | |
| .await | |
| .expect("serve disconnecting MCP tools"); | |
| let cancellation = server.cancellation_token(); | |
| tokio::select! { | |
| () = disconnect.cancelled() => cancellation.cancel(), | |
| result = server.waiting() => { | |
| result.expect("disconnecting MCP server should complete"); | |
| } | |
| } | |
| }); | |
| Ok(client_stream) | |
| } | |
| .boxed() | |
| } | |
| } | |
| async fn legacy_tool_catalog_does_not_follow_pagination_cursor() -> anyhow::Result<()> { | |
| let requests = Arc::new(AtomicUsize::new(0)); | |
| let client = Arc::new( | |
| RmcpClient::new_in_process_client(Arc::new(RefreshTestTransportFactory { | |
| tool: create_test_tool("legacy", "first-page").tool, | |
| list_started: None, | |
| release_list: None, | |
| next_cursor: Some("next-page".to_string()), | |
| list_requests: Arc::clone(&requests), | |
| })) | |
| .await?, | |
| ); | |
| client | |
| .initialize( | |
| InitializeRequestParams::new( | |
| ClientCapabilities::default(), | |
| Implementation::new("codex-test", "0.0.0-test"), | |
| ) | |
| .with_protocol_version(ProtocolVersion::V_2025_06_18), | |
| Some(Duration::from_secs(5)), | |
| Box::new(|_, _| async { Err(anyhow!("unexpected elicitation")) }.boxed()), | |
| ) | |
| .await?; | |
| let tools = list_tools_for_client_uncached( | |
| "legacy", | |
| /*is_codex_apps_mcp_server*/ false, | |
| "test", | |
| &client, | |
| Some(Duration::from_secs(5)), | |
| crate::pagination::MAX_MCP_CATALOG_ITEMS, | |
| /*server_instructions*/ None, | |
| ) | |
| .await?; | |
| assert_eq!(tools.len(), 1); | |
| assert_eq!(tools[0].tool.name.as_ref(), "first-page"); | |
| assert_eq!(requests.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| client.shutdown().await; | |
| Ok(()) | |
| } | |
| async fn create_test_managed_client(tools: Vec<ToolInfo>) -> ManagedClient { | |
| ManagedClient { | |
| _auth_change_notifications: None, | |
| client: Arc::new( | |
| RmcpClient::new_in_process_client(Arc::new(TestInProcessTransportFactory)) | |
| .await | |
| .expect("create in-process RMCP client"), | |
| ), | |
| server_info: create_test_server_info("Ready"), | |
| tool_catalog: Arc::new(ClientToolCatalog::new(tools, /*updates*/ None)), | |
| tool_timeout: None, | |
| server_instructions: None, | |
| server_supports_sandbox_state_meta_capability: false, | |
| codex_apps_tools_cache_context: None, | |
| } | |
| } | |
| async fn prepared_call_timeout_includes_trusted_access_lookup() { | |
| let mut tool = create_test_tool("docs", "access"); | |
| tool.tool.annotations = Some(rmcp::model::ToolAnnotations::new().read_only(true)); | |
| let mut tool_meta = rmcp::model::MetaObject::new(); | |
| tool_meta.insert( | |
| "openai/requestedEntitlements".to_string(), | |
| serde_json::json!(["cyber_trusted_access"]), | |
| ); | |
| tool.tool.meta = Some(tool_meta); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let mut config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| let mut catalog = crate::ResolvedMcpCatalog::builder(); | |
| catalog.register(crate::McpServerRegistration::from_plugin( | |
| "docs".to_string(), | |
| crate::McpPluginAttribution::new("docs@test".to_string(), "Docs".to_string()), | |
| /*plugin_order*/ 0, | |
| serde_json::from_value(serde_json::json!({ "command": "docs" })) | |
| .expect("plugin MCP config"), | |
| )); | |
| config.mcp_server_catalog = catalog.build(); | |
| config | |
| .server_permission_profiles | |
| .insert("docs".to_string(), PermissionProfile::default()); | |
| manager.tool_plugin_context = Arc::new(crate::tool_plugin_context(&config)); | |
| let auth = CodexAuth::create_dummy_chatgpt_auth_for_testing(); | |
| manager.trusted_access = Some(TrustedAccessContext::new( | |
| auth.clone(), | |
| AuthManager::from_auth_for_testing(auth), | |
| "https://chatgpt.com/backend-api".to_string(), | |
| Arc::new(PendingHttpClient), | |
| )); | |
| let manager = Arc::new(manager); | |
| let server_metadata = McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: true, | |
| origin: Some(McpServerOrigin::Stdio), | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }; | |
| let client = Arc::new(create_test_managed_client(vec![tool.clone()]).await); | |
| let catalog_snapshot = client.tool_catalog.read(Arc::new).await; | |
| let prepared = crate::PreparedMcpCall::new( | |
| manager, | |
| client, | |
| Arc::new(config), | |
| catalog_snapshot, | |
| tool, | |
| server_metadata, | |
| Some("docs@test".to_string()), | |
| /*selected_plugin_server*/ false, | |
| ) | |
| .expect("docs should retain its permission profile"); | |
| let started = tokio::time::Instant::now(); | |
| let error = prepared | |
| .call( | |
| Some(serde_json::json!({})), | |
| /*meta*/ None, | |
| Some(Duration::from_secs(1)), | |
| ) | |
| .await | |
| .expect_err("trusted access lookup should consume the call timeout"); | |
| assert_eq!(started.elapsed(), Duration::from_secs(1)); | |
| assert!(format!("{error:#}").contains("timed out awaiting tools/call after 1s")); | |
| } | |
| pub(crate) async fn create_ready_async_managed_client(tools: Vec<ToolInfo>) -> AsyncManagedClient { | |
| AsyncManagedClient { | |
| client: futures::future::ready::<Result<ManagedClient, StartupOutcomeError>>(Ok( | |
| create_test_managed_client(tools).await, | |
| )) | |
| .boxed() | |
| .shared(), | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| } | |
| } | |
| async fn connection_statuses_observe_clients_without_starting_them() { | |
| use codex_protocol::mcp::McpServerConnectionStatus as Status; | |
| let mut manager = McpConnectionSet::empty(/*prefix_mcp_tool_names*/ true); | |
| manager.disabled_servers.push("disabled".to_string()); | |
| let ready = create_ready_async_managed_client(Vec::new()).await; | |
| ready.client().await.expect("ready client"); | |
| manager.insert_test_client("connected", ready); | |
| for (name, error) in [ | |
| ( | |
| "failed", | |
| StartupOutcomeError::Failed { | |
| error: "broken".to_string(), | |
| is_authentication_required: false, | |
| }, | |
| ), | |
| ( | |
| "auth", | |
| StartupOutcomeError::Failed { | |
| error: "login".to_string(), | |
| is_authentication_required: true, | |
| }, | |
| ), | |
| ( | |
| "flattened-auth", | |
| StartupOutcomeError::from(anyhow!("Auth required for server")), | |
| ), | |
| ("cancelled", StartupOutcomeError::Cancelled), | |
| ] { | |
| let mut client = create_ready_async_managed_client(Vec::new()).await; | |
| client.client = futures::future::ready(Err(error)).boxed().shared(); | |
| assert!(client.client().await.is_err()); | |
| manager.insert_test_client(name, client); | |
| } | |
| let mut pending = create_ready_async_managed_client(Vec::new()).await; | |
| pending.client = futures::future::pending().boxed().shared(); | |
| pending.cached_server_info = Some(create_test_server_info("Cached")); | |
| manager.insert_test_client("starting", pending.clone()); | |
| manager.insert_test_client("deferred", pending); | |
| let (trigger, _receiver) = watch::channel(/*init*/ false); | |
| Arc::get_mut(&mut manager.servers.get_mut("deferred").unwrap().connection) | |
| .unwrap() | |
| .startup_trigger = Some(trigger.clone()); | |
| let statuses = tokio::time::timeout( | |
| Duration::from_millis(/*millis*/ 100), | |
| manager.connection_statuses(), | |
| ) | |
| .await | |
| .expect("status must not await startup"); | |
| let mut expected = HashMap::from([ | |
| ("connected".to_string(), Status::Connected), | |
| ("failed".to_string(), Status::Failed), | |
| ("auth".to_string(), Status::AuthenticationRequired), | |
| ("flattened-auth".to_string(), Status::AuthenticationRequired), | |
| ("cancelled".to_string(), Status::Cancelled), | |
| ("starting".to_string(), Status::Starting), | |
| ("deferred".to_string(), Status::NotStarted), | |
| ("disabled".to_string(), Status::Disabled), | |
| ]); | |
| assert_eq!(statuses, expected); | |
| assert!(!*trigger.borrow()); | |
| manager.test_client("connected").cancel_token.cancel(); | |
| expected.insert("connected".to_string(), Status::Cancelled); | |
| assert_eq!(manager.connection_statuses().await, expected); | |
| } | |
| async fn connection_statuses_follow_latest_reconnect_outcome() { | |
| use codex_protocol::mcp::McpServerConnectionStatus as Status; | |
| let recovered = create_test_managed_client(Vec::new()).await; | |
| let attempts = Arc::new(AtomicUsize::new(0)); | |
| let started = Arc::new(Notify::new()); | |
| let release = Arc::new(Notify::new()); | |
| let finished = Arc::new(Notify::new()); | |
| let factory = { | |
| let attempts = Arc::clone(&attempts); | |
| let started = Arc::clone(&started); | |
| let release = Arc::clone(&release); | |
| let finished = Arc::clone(&finished); | |
| Arc::new(move || { | |
| let attempt = attempts.fetch_add(1, Ordering::SeqCst); | |
| let recovered = recovered.clone(); | |
| let started = Arc::clone(&started); | |
| let release = Arc::clone(&release); | |
| let finished = Arc::clone(&finished); | |
| async move { | |
| started.notify_one(); | |
| release.notified().await; | |
| finished.notify_one(); | |
| match attempt { | |
| 0 | 1 => Err(StartupOutcomeError::Failed { | |
| error: "retry failed".to_string(), | |
| is_authentication_required: attempt == 0, | |
| }), | |
| _ => Ok(recovered), | |
| } | |
| } | |
| .boxed() | |
| .shared() | |
| }) | |
| }; | |
| let manager = create_test_manager_with_failed_apps_startup(Vec::new(), factory); | |
| let client = manager.test_client(CODEX_APPS_MCP_SERVER_NAME); | |
| assert!(client.client().await.is_err()); | |
| let expected = |status| HashMap::from([(CODEX_APPS_MCP_SERVER_NAME.to_string(), status)]); | |
| assert_eq!( | |
| manager.connection_statuses().await, | |
| expected(Status::Failed) | |
| ); | |
| for status in [ | |
| Status::AuthenticationRequired, | |
| Status::Failed, | |
| Status::Connected, | |
| ] { | |
| client.reconnect_failed_startup().await; | |
| started.notified().await; | |
| assert_eq!( | |
| manager.connection_statuses().await, | |
| expected(Status::Starting) | |
| ); | |
| release.notify_one(); | |
| finished.notified().await; | |
| assert_eq!(manager.connection_statuses().await, expected(status)); | |
| tokio::time::advance(CODEX_APPS_RECONNECT_INITIAL_BACKOFF * 2).await; | |
| } | |
| assert_eq!(attempts.load(Ordering::SeqCst), 3); | |
| } | |
| fn create_gated_async_managed_client( | |
| client: ManagedClient, | |
| ) -> ( | |
| AsyncManagedClient, | |
| tokio::sync::oneshot::Receiver<()>, | |
| tokio::sync::oneshot::Sender<()>, | |
| ) { | |
| let (started_tx, started_rx) = tokio::sync::oneshot::channel(); | |
| let (release_tx, release_rx) = tokio::sync::oneshot::channel(); | |
| let startup_complete = Arc::new(AtomicBool::new(false)); | |
| let startup_complete_for_client = Arc::clone(&startup_complete); | |
| let client = async move { | |
| started_tx.send(()).expect("signal client startup"); | |
| release_rx.await.expect("release client startup"); | |
| startup_complete_for_client.store(true, std::sync::atomic::Ordering::Release); | |
| Ok(client) | |
| } | |
| .boxed() | |
| .shared(); | |
| ( | |
| AsyncManagedClient { | |
| client, | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete, | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| started_rx, | |
| release_tx, | |
| ) | |
| } | |
| pub(crate) async fn create_test_manager_with_ready_apps_client( | |
| cache_context: ConnectorRuntimeContext<ToolInfo>, | |
| tool_name: &str, | |
| list_started: Option<Arc<Notify>>, | |
| release_list: Option<Arc<Notify>>, | |
| ) -> anyhow::Result<Arc<McpConnectionSet>> { | |
| let tool = create_test_tool(CODEX_APPS_MCP_SERVER_NAME, tool_name); | |
| let client = Arc::new( | |
| RmcpClient::new_in_process_client(Arc::new(RefreshTestTransportFactory { | |
| tool: tool.tool.clone(), | |
| list_started, | |
| release_list, | |
| next_cursor: None, | |
| list_requests: Arc::new(AtomicUsize::new(0)), | |
| })) | |
| .await?, | |
| ); | |
| client | |
| .initialize( | |
| InitializeRequestParams::new( | |
| ClientCapabilities::default(), | |
| Implementation::new("codex-test", "0.0.0-test"), | |
| ) | |
| .with_protocol_version(ProtocolVersion::V_2025_06_18), | |
| Some(Duration::from_secs(5)), | |
| Box::new(|_, _| async { Err(anyhow!("unexpected elicitation")) }.boxed()), | |
| ) | |
| .await?; | |
| let managed_client = ManagedClient { | |
| _auth_change_notifications: None, | |
| client, | |
| server_info: create_test_server_info("Codex Apps"), | |
| tool_catalog: Arc::new(ClientToolCatalog::new(vec![tool], /*updates*/ None)), | |
| tool_timeout: Some(Duration::from_secs(5)), | |
| server_instructions: None, | |
| server_supports_sandbox_state_meta_capability: false, | |
| codex_apps_tools_cache_context: Some(cache_context.clone()), | |
| }; | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: futures::future::ready::<Result<ManagedClient, StartupOutcomeError>>(Ok( | |
| managed_client, | |
| )) | |
| .boxed() | |
| .shared(), | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: Some(create_test_server_info("Codex Apps")), | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| manager.set_test_server_metadata( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: false, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| ); | |
| Ok(Arc::new(manager)) | |
| } | |
| fn create_test_manager_with_failed_apps_startup( | |
| cached_tools: Vec<ToolInfo>, | |
| reconnect_factory: Arc<dyn Fn() -> ManagedClientFuture + Send + Sync>, | |
| ) -> McpConnectionSet { | |
| let client: ManagedClientFuture = futures::future::ready(Err(StartupOutcomeError::Failed { | |
| error: "startup failed".to_string(), | |
| is_authentication_required: false, | |
| })) | |
| .boxed() | |
| .shared(); | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("reconnect-test-account"), | |
| Some("reconnect-test-user"), | |
| ); | |
| store_current_tools(&cache_context, cached_tools); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: Some(Arc::new(CodexAppsStartupReconnect::new(reconnect_factory))), | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| manager | |
| } | |
| fn model_tool_names(tools: &[ToolInfo]) -> HashSet<ToolName> { | |
| tools | |
| .iter() | |
| .map(ToolInfo::canonical_tool_name) | |
| .collect::<HashSet<_>>() | |
| } | |
| fn model_tool_name_len(name: &ToolName) -> usize { | |
| name.namespace | |
| .as_deref() | |
| .map_or(0, |namespace| namespace.len() + "__".len()) | |
| + name.name.len() | |
| } | |
| fn is_code_mode_compatible_tool_name(name: &ToolName) -> bool { | |
| name.namespace | |
| .as_deref() | |
| .into_iter() | |
| .chain(std::iter::once(name.name.as_str())) | |
| .flat_map(str::chars) | |
| .all(|c| c.is_ascii_alphanumeric() || c == '_') | |
| } | |
| fn elicitation_granular_policy_defaults_to_prompting() { | |
| assert!(!elicitation_is_rejected_by_policy( | |
| AskForApproval::OnRequest | |
| )); | |
| assert!(!elicitation_is_rejected_by_policy( | |
| AskForApproval::UnlessTrusted | |
| )); | |
| assert!(elicitation_is_rejected_by_policy(AskForApproval::Granular( | |
| GranularApprovalConfig { | |
| sandbox_approval: true, | |
| rules: true, | |
| skill_approval: true, | |
| request_permissions: true, | |
| mcp_elicitations: false, | |
| } | |
| ))); | |
| } | |
| fn elicitation_granular_policy_respects_never_and_config() { | |
| assert!(elicitation_is_rejected_by_policy(AskForApproval::Never)); | |
| assert!(elicitation_is_rejected_by_policy(AskForApproval::Granular( | |
| GranularApprovalConfig { | |
| sandbox_approval: true, | |
| rules: true, | |
| skill_approval: true, | |
| request_permissions: true, | |
| mcp_elicitations: false, | |
| } | |
| ))); | |
| } | |
| async fn disabled_permissions_auto_accept_elicitation_with_empty_form_schema() { | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config("server", AskForApproval::Never, PermissionProfile::Disabled), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| ElicitationRequestRouter::default(), | |
| ); | |
| let (tx_event, _rx_event) = async_channel::bounded(1); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let response = sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "Confirm?".to_string(), | |
| requested_schema: rmcp::model::ElicitationSchema::builder() | |
| .build() | |
| .expect("schema should build"), | |
| }), | |
| ) | |
| .await | |
| .expect("elicitation should auto accept"); | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Accept, | |
| content: Some(serde_json::json!({})), | |
| meta: None, | |
| } | |
| ); | |
| } | |
| async fn disabled_permissions_do_not_auto_accept_elicitation_with_requested_fields() { | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config("server", AskForApproval::Never, PermissionProfile::Disabled), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| ElicitationRequestRouter::default(), | |
| ); | |
| let (tx_event, _rx_event) = async_channel::bounded(1); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let response = sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "What should I say?".to_string(), | |
| requested_schema: | |
| rmcp::model::ElicitationSchema::builder() | |
| .required_property( | |
| "message", | |
| rmcp::model::PrimitiveSchemaDefinition::String( | |
| rmcp::model::StringSchema::new(), | |
| ), | |
| ) | |
| .build() | |
| .expect("schema should build"), | |
| }), | |
| ) | |
| .await | |
| .expect("elicitation should auto decline"); | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| } | |
| ); | |
| } | |
| fn full_access_form_input_enabled_router() -> ElicitationRequestRouter { | |
| let router = ElicitationRequestRouter::default(); | |
| router.enable_full_access_form_input(); | |
| router | |
| } | |
| fn elicitation_meta(value: serde_json::Value) -> Option<rmcp::model::RequestMetaObject> { | |
| let serde_json::Value::Object(map) = value else { | |
| panic!("elicitation metadata must be an object"); | |
| }; | |
| Some(rmcp::model::RequestMetaObject::from(map)) | |
| } | |
| fn requested_user_input_schema() -> rmcp::model::ElicitationSchema { | |
| rmcp::model::ElicitationSchema::builder() | |
| .required_property( | |
| "message", | |
| rmcp::model::PrimitiveSchemaDefinition::String(rmcp::model::StringSchema::new()), | |
| ) | |
| .build() | |
| .expect("schema should build") | |
| } | |
| struct DecliningElicitationReviewer { | |
| review_count: AtomicUsize, | |
| } | |
| impl ElicitationReviewer for DecliningElicitationReviewer { | |
| fn review( | |
| &self, | |
| _request: ElicitationReviewRequest, | |
| ) -> BoxFuture<'static, anyhow::Result<Option<ElicitationResponse>>> { | |
| self.review_count.fetch_add(1, Ordering::SeqCst); | |
| async { | |
| Ok(Some(ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| })) | |
| } | |
| .boxed() | |
| } | |
| } | |
| async fn assert_elicitation_declined_with_reviewer_calls( | |
| approval_policy: AskForApproval, | |
| server_name: &str, | |
| elicitation: ElicitRequestParams, | |
| expected_reviewer_calls: usize, | |
| ) { | |
| let reviewer = Arc::new(DecliningElicitationReviewer::default()); | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config(server_name, approval_policy, PermissionProfile::Disabled), | |
| Some(reviewer.clone()), | |
| /*lifecycle*/ None, | |
| full_access_form_input_enabled_router(), | |
| ); | |
| let (tx_event, rx_event) = async_channel::bounded(1); | |
| let sender = manager.make_sender( | |
| server_name.to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let response = tokio::select! { | |
| biased; | |
| event = rx_event.recv() => { | |
| panic!("elicitation unexpectedly reached the user: {event:?}"); | |
| } | |
| response = sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp(elicitation), | |
| ) => response.expect("elicitation should be declined"), | |
| }; | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| }, | |
| ); | |
| assert_eq!( | |
| reviewer.review_count.load(Ordering::SeqCst), | |
| expected_reviewer_calls | |
| ); | |
| assert!(rx_event.try_recv().is_err()); | |
| } | |
| async fn assert_requested_user_input_is_declined( | |
| approval_policy: AskForApproval, | |
| permission_profile: PermissionProfile, | |
| router: ElicitationRequestRouter, | |
| ) { | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config("server", approval_policy, permission_profile), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| router, | |
| ); | |
| let (tx_event, rx_event) = async_channel::bounded(1); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let response = tokio::select! { | |
| biased; | |
| event = rx_event.recv() => { | |
| panic!("user-input form unexpectedly reached the user: {event:?}"); | |
| } | |
| response = sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp( | |
| ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "What should I say?".to_string(), | |
| requested_schema: requested_user_input_schema(), | |
| }, | |
| ), | |
| ) => response.expect("restricted user-input request should decline"), | |
| }; | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| }, | |
| ); | |
| assert!(rx_event.try_recv().is_err()); | |
| } | |
| async fn disabled_permissions_do_not_surface_user_input_when_auto_denied() { | |
| let router = full_access_form_input_enabled_router(); | |
| router.set_auto_deny(/*auto_deny*/ true); | |
| assert_requested_user_input_is_declined( | |
| AskForApproval::Never, | |
| PermissionProfile::Disabled, | |
| router, | |
| ) | |
| .await; | |
| } | |
| async fn plugin_tool_suggestion_elicitations_are_declined_before_review() { | |
| assert_elicitation_declined_with_reviewer_calls( | |
| AskForApproval::OnRequest, | |
| "server", | |
| ElicitRequestParams::FormElicitationParams { | |
| meta: elicitation_meta(serde_json::json!({ | |
| "codex_approval_kind": "tool_suggestion", | |
| })), | |
| message: "Install this app?".to_string(), | |
| requested_schema: rmcp::model::ElicitationSchema::builder() | |
| .build() | |
| .expect("schema should build"), | |
| }, | |
| /*expected_reviewer_calls*/ 0, | |
| ) | |
| .await; | |
| } | |
| async fn disabled_permissions_surface_requested_user_input_without_metadata() { | |
| assert_disabled_permissions_surface_requested_user_input(/*meta*/ None).await; | |
| } | |
| async fn disabled_permissions_surface_requested_user_input_with_non_codex_approval_metadata() { | |
| assert_disabled_permissions_surface_requested_user_input(elicitation_meta(serde_json::json!({ | |
| "origin": "https://example.com", | |
| "persist": "always", | |
| }))) | |
| .await; | |
| } | |
| async fn assert_disabled_permissions_surface_requested_user_input( | |
| meta: Option<rmcp::model::RequestMetaObject>, | |
| ) { | |
| let router = full_access_form_input_enabled_router(); | |
| let reviewer = Arc::new(DecliningElicitationReviewer::default()); | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config("server", AskForApproval::Never, PermissionProfile::Disabled), | |
| Some(reviewer.clone()), | |
| /*lifecycle*/ None, | |
| router.clone(), | |
| ); | |
| let (tx_event, rx_event) = async_channel::bounded(1); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let requested_schema = requested_user_input_schema(); | |
| let mut pending = tokio::spawn(sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: meta.clone(), | |
| message: "What should I say?".to_string(), | |
| requested_schema: requested_schema.clone(), | |
| }), | |
| )); | |
| let request = tokio::select! { | |
| event = rx_event.recv() => { | |
| let EventMsg::ElicitationRequest(request) = event.expect("user-input event").msg else { | |
| panic!("expected MCP user-input elicitation"); | |
| }; | |
| request | |
| } | |
| response = &mut pending => { | |
| panic!("user input resolved without reaching the user: {response:?}"); | |
| } | |
| }; | |
| assert_eq!( | |
| request.request, | |
| ElicitationRequest::Form { | |
| meta: meta | |
| .map(serde_json::to_value) | |
| .transpose() | |
| .expect("user-input metadata should serialize"), | |
| message: "What should I say?".to_string(), | |
| requested_schema: serde_json::to_value(requested_schema) | |
| .expect("schema should serialize"), | |
| }, | |
| ); | |
| assert_eq!(request.server_name, "server"); | |
| assert_eq!(reviewer.review_count.load(Ordering::SeqCst), 0); | |
| let codex_protocol::mcp::RequestId::String(request_id) = request.id else { | |
| panic!("expected Codex-owned string request ID"); | |
| }; | |
| let user_response = ElicitationResponse { | |
| action: ElicitationAction::Accept, | |
| content: Some(serde_json::json!({ "message": "The actual user response." })), | |
| meta: None, | |
| }; | |
| router | |
| .resolve( | |
| "server".to_string(), | |
| NumberOrString::String(request_id.into()), | |
| user_response.clone(), | |
| ) | |
| .await | |
| .expect("actual user response should resolve the elicitation"); | |
| assert_eq!( | |
| pending | |
| .await | |
| .expect("user-input task should complete") | |
| .expect("user input should resolve"), | |
| user_response, | |
| ); | |
| } | |
| async fn disabled_permissions_decline_requested_user_input_with_approval_metadata() { | |
| assert_elicitation_declined_with_reviewer_calls( | |
| AskForApproval::Never, | |
| "node_repl", | |
| ElicitRequestParams::FormElicitationParams { | |
| meta: elicitation_meta(serde_json::json!({ | |
| "codex_approval_kind": "mcp_tool_call", | |
| "connector_id": "browser-use", | |
| "tool_name": "access_browser_origin", | |
| })), | |
| message: "Allow Browser Use to access this website?".to_string(), | |
| requested_schema: | |
| rmcp::model::ElicitationSchema::builder() | |
| .required_property( | |
| "confirmation", | |
| rmcp::model::PrimitiveSchemaDefinition::String( | |
| rmcp::model::StringSchema::new(), | |
| ), | |
| ) | |
| .build() | |
| .expect("schema should build"), | |
| }, | |
| /*expected_reviewer_calls*/ 0, | |
| ) | |
| .await; | |
| } | |
| async fn restricted_never_policy_does_not_surface_requested_user_input() { | |
| assert_requested_user_input_is_declined( | |
| AskForApproval::Never, | |
| PermissionProfile::default(), | |
| full_access_form_input_enabled_router(), | |
| ) | |
| .await; | |
| } | |
| async fn granular_policy_does_not_surface_requested_user_input() { | |
| assert_requested_user_input_is_declined( | |
| AskForApproval::Granular(GranularApprovalConfig { | |
| sandbox_approval: true, | |
| rules: true, | |
| skill_approval: true, | |
| request_permissions: true, | |
| mcp_elicitations: false, | |
| }), | |
| PermissionProfile::Disabled, | |
| full_access_form_input_enabled_router(), | |
| ) | |
| .await; | |
| } | |
| async fn on_request_approval_forms_remain_with_the_reviewer() { | |
| assert_elicitation_declined_with_reviewer_calls( | |
| AskForApproval::OnRequest, | |
| "server", | |
| ElicitRequestParams::FormElicitationParams { | |
| meta: elicitation_meta(serde_json::json!({ | |
| "codex_request_type": "approval_request", | |
| "codex_approval_kind": "mcp_tool_call", | |
| "tool_name": "test_tool", | |
| })), | |
| message: "Approve this action?".to_string(), | |
| requested_schema: rmcp::model::ElicitationSchema::builder() | |
| .build() | |
| .expect("schema should build"), | |
| }, | |
| /*expected_reviewer_calls*/ 1, | |
| ) | |
| .await; | |
| } | |
| async fn disabled_permissions_decline_user_input_without_an_event_channel() { | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config("server", AskForApproval::Never, PermissionProfile::Disabled), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| full_access_form_input_enabled_router(), | |
| ); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| /*tx_event*/ None, | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let response = sender( | |
| NumberOrString::Number(1), | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "What should I say?".to_string(), | |
| requested_schema: requested_user_input_schema(), | |
| }), | |
| ) | |
| .await | |
| .expect("headless user-input request should decline"); | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| }, | |
| ); | |
| } | |
| async fn concurrent_authority_updates_never_auto_approve_mixed_policy() { | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::Never, | |
| PermissionProfile::default(), | |
| ), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| ElicitationRequestRouter::default(), | |
| ); | |
| let updating_manager = manager.clone(); | |
| let updater = tokio::spawn(async move { | |
| for _ in 0..1_000 { | |
| assert!(updating_manager.update( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::OnRequest, | |
| PermissionProfile::Disabled | |
| ), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| )); | |
| assert!(updating_manager.update( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::Never, | |
| PermissionProfile::default() | |
| ), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| )); | |
| } | |
| }); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| /*tx_event*/ None, | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let elicitation = | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "Confirm?".to_string(), | |
| requested_schema: rmcp::model::ElicitationSchema::builder() | |
| .build() | |
| .expect("schema should build"), | |
| }); | |
| for _ in 0..1_000 { | |
| let response = sender(NumberOrString::Number(1), elicitation.clone()) | |
| .await | |
| .expect("elicitation should resolve"); | |
| assert_eq!( | |
| response, | |
| ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| } | |
| ); | |
| } | |
| updater.await.expect("authority updates should finish"); | |
| } | |
| async fn shared_elicitation_router_targets_the_exact_pending_request() { | |
| struct Registration(Arc<AtomicUsize>); | |
| impl Drop for Registration { | |
| fn drop(&mut self) { | |
| self.0.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); | |
| } | |
| } | |
| let router = ElicitationRequestRouter::default(); | |
| let outstanding = Arc::new(AtomicUsize::new(0)); | |
| let lifecycle = ElicitationLifecycle::new({ | |
| let outstanding = outstanding.clone(); | |
| move || { | |
| outstanding.fetch_add(1, std::sync::atomic::Ordering::SeqCst); | |
| Registration(outstanding.clone()) | |
| } | |
| }); | |
| let manager_a = ElicitationRequestManager::new( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::OnRequest, | |
| PermissionProfile::default(), | |
| ), | |
| /*reviewer*/ None, | |
| Some(lifecycle.clone()), | |
| router.clone(), | |
| ); | |
| let manager_b = ElicitationRequestManager::new( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::OnRequest, | |
| PermissionProfile::default(), | |
| ), | |
| /*reviewer*/ None, | |
| Some(lifecycle), | |
| router.clone(), | |
| ); | |
| let (tx_event, rx_event) = async_channel::bounded(2); | |
| let sender_a = manager_a.make_sender( | |
| "server".to_string(), | |
| Some(tx_event.clone()), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let sender_b = manager_b.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let elicitation = | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "Which runtime?".to_string(), | |
| requested_schema: | |
| rmcp::model::ElicitationSchema::builder() | |
| .required_property( | |
| "runtime", | |
| rmcp::model::PrimitiveSchemaDefinition::String( | |
| rmcp::model::StringSchema::new(), | |
| ), | |
| ) | |
| .build() | |
| .expect("schema should build"), | |
| }); | |
| let pending_a = tokio::spawn(sender_a(NumberOrString::Number(1), elicitation.clone())); | |
| let EventMsg::ElicitationRequest(request_a) = rx_event.recv().await.expect("request A").msg | |
| else { | |
| panic!("expected elicitation request"); | |
| }; | |
| let pending_b = tokio::spawn(sender_b(NumberOrString::Number(1), elicitation)); | |
| let EventMsg::ElicitationRequest(request_b) = rx_event.recv().await.expect("request B").msg | |
| else { | |
| panic!("expected elicitation request"); | |
| }; | |
| assert_eq!(outstanding.load(std::sync::atomic::Ordering::SeqCst), 2); | |
| let ( | |
| codex_protocol::mcp::RequestId::String(request_a_id), | |
| codex_protocol::mcp::RequestId::String(request_b_id), | |
| ) = (request_a.id, request_b.id) | |
| else { | |
| panic!("expected Codex-owned string request IDs"); | |
| }; | |
| assert_ne!(request_a_id, request_b_id); | |
| let response_a = ElicitationResponse { | |
| action: ElicitationAction::Accept, | |
| content: Some(serde_json::json!({"runtime": "a"})), | |
| meta: None, | |
| }; | |
| router | |
| .resolve( | |
| "server".to_string(), | |
| NumberOrString::String(request_a_id.into()), | |
| response_a.clone(), | |
| ) | |
| .await | |
| .expect("runtime B should route a response to runtime A"); | |
| let response_b = ElicitationResponse { | |
| action: ElicitationAction::Accept, | |
| content: Some(serde_json::json!({"runtime": "b"})), | |
| meta: None, | |
| }; | |
| router | |
| .resolve( | |
| "server".to_string(), | |
| NumberOrString::String(request_b_id.into()), | |
| response_b.clone(), | |
| ) | |
| .await | |
| .expect("runtime A should route a response to runtime B"); | |
| assert_eq!( | |
| pending_a | |
| .await | |
| .expect("request A task") | |
| .expect("request A response"), | |
| response_a | |
| ); | |
| assert_eq!( | |
| pending_b | |
| .await | |
| .expect("request B task") | |
| .expect("request B response"), | |
| response_b | |
| ); | |
| assert_eq!(outstanding.load(std::sync::atomic::Ordering::SeqCst), 0); | |
| } | |
| async fn cancelled_elicitation_is_removed_without_affecting_other_pending_requests() { | |
| let router = ElicitationRequestRouter::default(); | |
| let manager = ElicitationRequestManager::new( | |
| test_elicitation_config( | |
| "server", | |
| AskForApproval::OnRequest, | |
| PermissionProfile::default(), | |
| ), | |
| /*reviewer*/ None, | |
| /*lifecycle*/ None, | |
| router.clone(), | |
| ); | |
| let (tx_event, rx_event) = async_channel::bounded(2); | |
| let sender = manager.make_sender( | |
| "server".to_string(), | |
| Some(tx_event), | |
| &ClientMcpExtensions::default(), | |
| ); | |
| let elicitation = | |
| codex_rmcp_client::Elicitation::Mcp(ElicitRequestParams::FormElicitationParams { | |
| meta: None, | |
| message: "Confirm?".to_string(), | |
| requested_schema: | |
| rmcp::model::ElicitationSchema::builder() | |
| .required_property( | |
| "answer", | |
| rmcp::model::PrimitiveSchemaDefinition::String( | |
| rmcp::model::StringSchema::new(), | |
| ), | |
| ) | |
| .build() | |
| .expect("schema should build"), | |
| }); | |
| let cancelled = tokio::spawn(sender(NumberOrString::Number(1), elicitation.clone())); | |
| let EventMsg::ElicitationRequest(cancelled_request) = | |
| rx_event.recv().await.expect("cancelled request event").msg | |
| else { | |
| panic!("expected elicitation request"); | |
| }; | |
| let pending = tokio::spawn(sender(NumberOrString::Number(2), elicitation)); | |
| let EventMsg::ElicitationRequest(pending_request) = | |
| rx_event.recv().await.expect("pending request event").msg | |
| else { | |
| panic!("expected elicitation request"); | |
| }; | |
| let ( | |
| codex_protocol::mcp::RequestId::String(cancelled_id), | |
| codex_protocol::mcp::RequestId::String(pending_id), | |
| ) = (cancelled_request.id, pending_request.id) | |
| else { | |
| panic!("expected Codex-owned string request IDs"); | |
| }; | |
| cancelled.abort(); | |
| assert!( | |
| cancelled | |
| .await | |
| .expect_err("cancelled request should be aborted") | |
| .is_cancelled() | |
| ); | |
| let response = ElicitationResponse { | |
| action: ElicitationAction::Accept, | |
| content: Some(serde_json::json!({"answer": "yes"})), | |
| meta: None, | |
| }; | |
| let error = router | |
| .resolve( | |
| "server".to_string(), | |
| NumberOrString::String(cancelled_id.into()), | |
| response.clone(), | |
| ) | |
| .await | |
| .expect_err("cancelled request should be removed immediately"); | |
| assert_eq!(error.to_string(), "elicitation request not found"); | |
| router | |
| .resolve( | |
| "server".to_string(), | |
| NumberOrString::String(pending_id.into()), | |
| response.clone(), | |
| ) | |
| .await | |
| .expect("another pending request should remain routable"); | |
| assert_eq!( | |
| pending | |
| .await | |
| .expect("pending request task") | |
| .expect("pending request response"), | |
| response | |
| ); | |
| } | |
| fn test_normalize_tools_short_non_duplicated_names() { | |
| let tools = vec![ | |
| create_test_tool("server1", "tool1"), | |
| create_test_tool("server1", "tool2"), | |
| ]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!( | |
| model_tool_names(&model_tools), | |
| HashSet::from([ | |
| ToolName::namespaced("mcp__server1", "tool1"), | |
| ToolName::namespaced("mcp__server1", "tool2") | |
| ]) | |
| ); | |
| } | |
| fn test_normalize_tools_omits_prefix_only_for_selected_servers() { | |
| let tools = vec![ | |
| create_test_tool("history", "search"), | |
| create_test_tool("notes", "read"), | |
| create_test_tool("calendar", "list"), | |
| ]; | |
| let model_tools = normalize_tools_for_model_with_prefix( | |
| tools, | |
| /*prefix_mcp_tool_names*/ true, | |
| &["history".to_string(), "notes".to_string()], | |
| ); | |
| assert_eq!( | |
| model_tool_names(&model_tools), | |
| HashSet::from([ | |
| ToolName::namespaced("history", "search"), | |
| ToolName::namespaced("notes", "read"), | |
| ToolName::namespaced("mcp__calendar", "list"), | |
| ]) | |
| ); | |
| } | |
| fn test_normalize_tools_selects_raw_server_name() { | |
| let mut tool = create_test_tool("codex_apps", "search"); | |
| tool.callable_namespace = "codex_apps__calendar".to_string(); | |
| let model_tools = normalize_tools_for_model_with_prefix( | |
| vec![tool], | |
| /*prefix_mcp_tool_names*/ true, | |
| &["codex_apps".to_string()], | |
| ); | |
| assert_eq!( | |
| model_tool_names(&model_tools), | |
| HashSet::from([ToolName::namespaced("codex_apps__calendar", "search")]) | |
| ); | |
| } | |
| fn test_normalize_tools_global_feature_omits_prefix_for_every_server() { | |
| let tools = vec![ | |
| create_test_tool("history", "search"), | |
| create_test_tool("calendar", "list"), | |
| ]; | |
| let model_tools = normalize_tools_for_model_with_prefix( | |
| tools, | |
| /*prefix_mcp_tool_names*/ false, | |
| &["history".to_string()], | |
| ); | |
| assert_eq!( | |
| model_tool_names(&model_tools), | |
| HashSet::from([ | |
| ToolName::namespaced("history", "search"), | |
| ToolName::namespaced("calendar", "list"), | |
| ]) | |
| ); | |
| } | |
| fn test_normalize_tools_duplicated_names_skipped() { | |
| let tools = vec![ | |
| create_test_tool("server1", "duplicate_tool"), | |
| create_test_tool("server1", "duplicate_tool"), | |
| ]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| // Only the first tool should remain, the second is skipped | |
| assert_eq!( | |
| model_tool_names(&model_tools), | |
| HashSet::from([ToolName::namespaced("mcp__server1", "duplicate_tool")]) | |
| ); | |
| } | |
| fn test_normalize_tools_respects_responses_api_name_length_boundaries() { | |
| let namespace = "mcp__codex_apps"; | |
| let namespace_len = namespace.len() + "__".len(); | |
| for total_len in [128, 129] { | |
| let tool_name = "a".repeat(total_len - namespace_len); | |
| let model_tools = normalize_tools_for_model_with_prefix( | |
| vec![create_test_tool("codex_apps", &tool_name)], | |
| /*prefix_mcp_tool_names*/ true, | |
| &[], | |
| ); | |
| let model_name = model_tools[0].canonical_tool_name(); | |
| assert_eq!(model_tool_name_len(&model_name), 128); | |
| if total_len == 128 { | |
| assert_eq!(model_name, ToolName::namespaced(namespace, tool_name)); | |
| } else { | |
| assert_ne!(model_name.name, tool_name); | |
| } | |
| } | |
| } | |
| fn test_normalize_tools_long_names_same_server() { | |
| let server_name = "my_server"; | |
| let first_name = "a".repeat(128); | |
| let second_name = "b".repeat(128); | |
| let tools = vec![ | |
| create_test_tool(server_name, &first_name), | |
| create_test_tool(server_name, &second_name), | |
| ]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!(model_tools.len(), 2); | |
| let names = model_tool_names(&model_tools); | |
| assert!(names.iter().all(|name| model_tool_name_len(name) == 128)); | |
| assert!( | |
| names | |
| .iter() | |
| .all(|name| name.namespace.as_deref() == Some("mcp__my_server")) | |
| ); | |
| assert!( | |
| names.iter().all(is_code_mode_compatible_tool_name), | |
| "model-visible names must be code-mode compatible: {names:?}" | |
| ); | |
| } | |
| fn test_normalize_tools_sanitizes_invalid_characters() { | |
| let tools = vec![create_test_tool("server.one", "tool.two-three")]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!(model_tools.len(), 1); | |
| let tool = model_tools.into_iter().next().expect("one tool"); | |
| let model_name = tool.canonical_tool_name(); | |
| assert_eq!( | |
| model_name, | |
| ToolName::namespaced("mcp__server_one", "tool_two_three") | |
| ); | |
| assert_eq!( | |
| ToolName::namespaced(tool.callable_namespace.clone(), tool.callable_name.clone()), | |
| model_name | |
| ); | |
| // The callable parts are sanitized for model-visible tool calls, but the raw | |
| // MCP name is preserved for the actual MCP call. | |
| assert_eq!(tool.server_name, "server.one"); | |
| assert_eq!(tool.callable_namespace, "mcp__server_one"); | |
| assert_eq!(tool.callable_name, "tool_two_three"); | |
| assert_eq!(tool.tool.name, "tool.two-three"); | |
| assert!( | |
| is_code_mode_compatible_tool_name(&model_name), | |
| "model-visible name must be code-mode compatible: {model_name:?}" | |
| ); | |
| } | |
| fn test_normalize_tools_keeps_hyphenated_mcp_tools_callable() { | |
| let tools = vec![create_test_tool("music-studio", "get-strudel-guide")]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!(model_tools.len(), 1); | |
| let tool = model_tools.into_iter().next().expect("one tool"); | |
| assert_eq!( | |
| tool.canonical_tool_name(), | |
| ToolName::namespaced("mcp__music_studio", "get_strudel_guide") | |
| ); | |
| assert_eq!(tool.callable_namespace, "mcp__music_studio"); | |
| assert_eq!(tool.callable_name, "get_strudel_guide"); | |
| assert_eq!(tool.tool.name, "get-strudel-guide"); | |
| } | |
| fn test_normalize_tools_disambiguates_sanitized_namespace_collisions() { | |
| let tools = vec![ | |
| create_test_tool("basic-server", "lookup"), | |
| create_test_tool("basic_server", "query"), | |
| create_test_tool("npm:@scope/package.name", "lookup"), | |
| create_test_tool("npm__scope_package_name", "lookup"), | |
| ]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!(model_tools.len(), 4); | |
| let mut namespaces = model_tools | |
| .iter() | |
| .map(|tool| tool.callable_namespace.as_str()) | |
| .collect::<Vec<_>>(); | |
| namespaces.sort(); | |
| namespaces.dedup(); | |
| assert_eq!(namespaces.len(), 4); | |
| let raw_servers = model_tools | |
| .iter() | |
| .map(|tool| tool.server_name.as_str()) | |
| .collect::<HashSet<_>>(); | |
| assert_eq!( | |
| raw_servers, | |
| HashSet::from([ | |
| "basic-server", | |
| "basic_server", | |
| "npm:@scope/package.name", | |
| "npm__scope_package_name", | |
| ]) | |
| ); | |
| let model_names = model_tool_names(&model_tools); | |
| assert!( | |
| model_names.iter().all(is_code_mode_compatible_tool_name), | |
| "model-visible names must be code-mode compatible: {model_names:?}" | |
| ); | |
| } | |
| fn test_normalize_tools_disambiguates_sanitized_tool_name_collisions() { | |
| let tools = vec![ | |
| create_test_tool("server", "tool-name"), | |
| create_test_tool("server", "tool_name"), | |
| ]; | |
| let model_tools = | |
| normalize_tools_for_model_with_prefix(tools, /*prefix_mcp_tool_names*/ true, &[]); | |
| assert_eq!(model_tools.len(), 2); | |
| let raw_tool_names = model_tools | |
| .iter() | |
| .map(|tool| tool.tool.name.to_string()) | |
| .collect::<HashSet<_>>(); | |
| assert_eq!( | |
| raw_tool_names, | |
| HashSet::from(["tool-name".to_string(), "tool_name".to_string()]) | |
| ); | |
| let callable_tool_names = model_tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<HashSet<_>>(); | |
| assert_eq!(callable_tool_names.len(), 2); | |
| } | |
| fn tool_filter_allows_by_default() { | |
| let filter = ToolFilter::default(); | |
| assert!(filter.allows("any")); | |
| } | |
| fn tool_filter_applies_enabled_list() { | |
| let filter = ToolFilter { | |
| enabled: Some(HashSet::from(["allowed".to_string()])), | |
| disabled: HashSet::new(), | |
| }; | |
| assert!(filter.allows("allowed")); | |
| assert!(!filter.allows("denied")); | |
| } | |
| fn tool_filter_applies_disabled_list() { | |
| let filter = ToolFilter { | |
| enabled: None, | |
| disabled: HashSet::from(["blocked".to_string()]), | |
| }; | |
| assert!(!filter.allows("blocked")); | |
| assert!(filter.allows("open")); | |
| } | |
| fn tool_filter_applies_enabled_then_disabled() { | |
| let filter = ToolFilter { | |
| enabled: Some(HashSet::from(["keep".to_string(), "remove".to_string()])), | |
| disabled: HashSet::from(["remove".to_string()]), | |
| }; | |
| assert!(filter.allows("keep")); | |
| assert!(!filter.allows("remove")); | |
| assert!(!filter.allows("unknown")); | |
| } | |
| fn filter_tools_applies_per_server_filters() { | |
| let server1_tools = vec![ | |
| create_test_tool("server1", "tool_a"), | |
| create_test_tool("server1", "tool_b"), | |
| ]; | |
| let server2_tools = vec![create_test_tool("server2", "tool_a")]; | |
| let server1_filter = ToolFilter { | |
| enabled: Some(HashSet::from(["tool_a".to_string(), "tool_b".to_string()])), | |
| disabled: HashSet::from(["tool_b".to_string()]), | |
| }; | |
| let server2_filter = ToolFilter { | |
| enabled: None, | |
| disabled: HashSet::from(["tool_a".to_string()]), | |
| }; | |
| let filtered: Vec<_> = filter_tools(server1_tools, &server1_filter) | |
| .into_iter() | |
| .chain(filter_tools(server2_tools, &server2_filter)) | |
| .collect(); | |
| assert_eq!(filtered.len(), 1); | |
| assert_eq!(filtered[0].server_name, "server1"); | |
| assert_eq!(filtered[0].callable_name, "tool_a"); | |
| } | |
| fn codex_apps_env_bearer_token_bypasses_shared_tools_cache() { | |
| assert!(!should_share_codex_apps_tools_cache( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| /*uses_env_bearer_token*/ true, | |
| )); | |
| } | |
| async fn hosted_apps_protocol_mode_is_independent_of_generic_mode() -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "url": "http://127.0.0.1:1/ps/mcp" }))?; | |
| for (generic_mode, hosted_mode) in [ | |
| ( | |
| crate::McpProtocolMode::Legacy, | |
| crate::McpProtocolMode::V20260728, | |
| ), | |
| ( | |
| crate::McpProtocolMode::V20260728, | |
| crate::McpProtocolMode::Legacy, | |
| ), | |
| ] { | |
| let mut config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| config.protocol_mode = generic_mode; | |
| config.host_owned_apps_protocol_mode = hosted_mode; | |
| let mut catalog = crate::ResolvedMcpCatalog::builder(); | |
| catalog.register(crate::McpServerRegistration::from_hosted_apps( | |
| "test-host", | |
| /*contribution_order*/ 0, | |
| server_config.clone(), | |
| )); | |
| catalog.register(crate::McpServerRegistration::from_config( | |
| "third_party".to_string(), | |
| server_config.clone(), | |
| )); | |
| config.mcp_server_catalog = catalog.build(); | |
| let startup_cancellation_token = CancellationToken::new(); | |
| startup_cancellation_token.cancel(); | |
| let manager = McpConnectionSet::new( | |
| /*previous*/ None, | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(config), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers: HashMap::from([ | |
| ( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| EffectiveMcpServer::configured(server_config.clone()), | |
| ), | |
| ( | |
| "third_party".to_string(), | |
| EffectiveMcpServer::configured(server_config.clone()), | |
| ), | |
| ]), | |
| submit_id: "protocol-mode-scope".to_string(), | |
| tx_event: None, | |
| startup_cancellation_token, | |
| runtime_context: McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| codex_home.path().to_path_buf(), | |
| ), | |
| codex_apps_tools_cache: ConnectorRuntimeManager::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| assert_eq!( | |
| manager.servers[CODEX_APPS_MCP_SERVER_NAME].protocol_mode, | |
| hosted_mode | |
| ); | |
| assert_eq!(manager.servers["third_party"].protocol_mode, generic_mode); | |
| assert_eq!( | |
| manager | |
| .event_stream_connection | |
| .as_ref() | |
| .expect("hosted Apps event stream") | |
| .protocol_mode, | |
| hosted_mode | |
| ); | |
| } | |
| Ok(()) | |
| } | |
| async fn codex_apps_extension_does_not_share_host_owned_tools_cache() -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let cache_key = ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ); | |
| let codex_apps_tools_cache = ConnectorRuntimeManager::<ToolInfo>::default(); | |
| let cache_context = | |
| codex_apps_tools_cache.context(codex_home.path().to_path_buf(), cache_key.clone()); | |
| store_current_tools( | |
| &cache_context, | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "calendar_create_event", | |
| )], | |
| ); | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "url": "http://127.0.0.1:1" }))?; | |
| for (hosted_mode, extension_mode, expected_mode) in [ | |
| ( | |
| crate::McpProtocolMode::V20260728, | |
| None, | |
| crate::McpProtocolMode::Legacy, | |
| ), | |
| ( | |
| crate::McpProtocolMode::Legacy, | |
| Some(crate::McpProtocolMode::V20260728), | |
| crate::McpProtocolMode::V20260728, | |
| ), | |
| ] { | |
| let mut config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| config.host_owned_apps_protocol_mode = hosted_mode; | |
| let mut registration = crate::McpServerRegistration::from_extension( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| "test-extension", | |
| /*contribution_order*/ 0, | |
| server_config.clone(), | |
| ); | |
| if let Some(mode) = extension_mode { | |
| registration = registration.with_protocol_mode(mode); | |
| } | |
| let mut catalog = crate::ResolvedMcpCatalog::builder(); | |
| catalog.register(registration); | |
| config.mcp_server_catalog = catalog.build(); | |
| let startup_cancellation_token = CancellationToken::new(); | |
| startup_cancellation_token.cancel(); | |
| let manager = McpConnectionSet::new( | |
| /*previous*/ None, | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(config), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers: HashMap::from([( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| EffectiveMcpServer::configured(server_config.clone()), | |
| )]), | |
| submit_id: "cache-ownership-test".to_string(), | |
| tx_event: None, | |
| startup_cancellation_token, | |
| runtime_context: McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| codex_home.path().to_path_buf(), | |
| ), | |
| codex_apps_tools_cache: codex_apps_tools_cache.clone(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: cache_key.clone(), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| let client = manager.test_client(CODEX_APPS_MCP_SERVER_NAME); | |
| assert_eq!( | |
| manager.servers[CODEX_APPS_MCP_SERVER_NAME].protocol_mode, expected_mode, | |
| "an ordinary extension must use its own mode, not the hosted protocol default" | |
| ); | |
| assert!( | |
| client.codex_apps_tools_cache_context.is_none(), | |
| "an extension must not receive the host-owned Apps cache" | |
| ); | |
| assert!( | |
| !client.has_cached_tools(), | |
| "an extension must not expose cached host-owned Apps tools" | |
| ); | |
| } | |
| Ok(()) | |
| } | |
| async fn list_all_tools_uses_shared_codex_apps_cache_while_client_is_pending() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| store_current_tools( | |
| &cache_context, | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "calendar_create_event", | |
| )], | |
| ); | |
| let pending_client = futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let tools = manager.list_all_tools().await; | |
| let tool = tools | |
| .iter() | |
| .find(|tool| { | |
| tool.canonical_tool_name() | |
| == ToolName::namespaced("mcp__codex_apps", "calendar_create_event") | |
| }) | |
| .expect("tool from shared cache"); | |
| assert_eq!(tool.server_name, CODEX_APPS_MCP_SERVER_NAME); | |
| assert_eq!(tool.callable_name, "calendar_create_event"); | |
| } | |
| async fn capture_binding_uses_the_ready_clients_own_tools() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| store_current_tools( | |
| &cache_context, | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "shared_cached_tool", | |
| )], | |
| ); | |
| let mut ready_client = create_test_managed_client(vec![ | |
| create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "client_local_tool"), | |
| create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "client_local_blocked"), | |
| ]) | |
| .await; | |
| let tool_filter = ToolFilter { | |
| enabled: None, | |
| disabled: HashSet::from(["client_local_blocked".to_string()]), | |
| }; | |
| ready_client.codex_apps_tools_cache_context = Some(cache_context.clone()); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: futures::future::ready(Ok(ready_client)).boxed().shared(), | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| manager | |
| .servers | |
| .get_mut(CODEX_APPS_MCP_SERVER_NAME) | |
| .expect("test server exists") | |
| .tool_filter = tool_filter; | |
| manager.set_test_server_metadata( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: false, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| ); | |
| let manager = Arc::new(manager); | |
| assert_eq!( | |
| manager | |
| .list_all_tools() | |
| .await | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["shared_cached_tool"] | |
| ); | |
| let step = capture_binding(&manager).await; | |
| assert_eq!( | |
| step.tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["client_local_tool"] | |
| ); | |
| assert!( | |
| step.prepare_call(CODEX_APPS_MCP_SERVER_NAME, "client_local_tool") | |
| .is_some() | |
| ); | |
| assert!( | |
| step.prepare_call(CODEX_APPS_MCP_SERVER_NAME, "shared_cached_tool") | |
| .is_none() | |
| ); | |
| assert!( | |
| step.prepare_call(CODEX_APPS_MCP_SERVER_NAME, "client_local_blocked") | |
| .is_none() | |
| ); | |
| } | |
| async fn hard_refresh_keeps_client_catalog_local_when_shared_cache_loses_race() -> anyhow::Result<()> | |
| { | |
| let codex_home = tempdir()?; | |
| let shared_cache = ConnectorRuntimeManager::<ToolInfo>::default(); | |
| let cache_key = ConnectorRuntimeContextKey::personal( | |
| Some("shared-account".to_string()), | |
| Some("shared-user".to_string()), | |
| ); | |
| let cache_context_a = shared_cache.context(codex_home.path().to_path_buf(), cache_key.clone()); | |
| let cache_context_b = shared_cache.context(codex_home.path().to_path_buf(), cache_key); | |
| let list_started = Arc::new(Notify::new()); | |
| let release_list = Arc::new(Notify::new()); | |
| let manager_a = create_test_manager_with_ready_apps_client( | |
| cache_context_a.clone(), | |
| "a_only", | |
| Some(Arc::clone(&list_started)), | |
| Some(Arc::clone(&release_list)), | |
| ) | |
| .await?; | |
| let mut manager_b = create_test_manager_with_ready_apps_client( | |
| cache_context_b, | |
| "b_only", | |
| /*list_started*/ None, | |
| /*release_list*/ None, | |
| ) | |
| .await?; | |
| let manager_a_for_refresh = Arc::clone(&manager_a); | |
| let refresh_a = tokio::spawn(async move { | |
| manager_a_for_refresh | |
| .refresh_codex_apps_tools_for_discovery() | |
| .await | |
| }); | |
| list_started.notified().await; | |
| let tools_b = manager_b.refresh_codex_apps_tools_for_discovery().await?; | |
| release_list.notify_one(); | |
| let tools_a = refresh_a.await??; | |
| assert_eq!( | |
| tools_b | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["b_only"] | |
| ); | |
| assert_eq!( | |
| tools_a | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["b_only"] | |
| ); | |
| assert_eq!( | |
| cache_context_a | |
| .current_tools() | |
| .expect("shared cache tools") | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["b_only"] | |
| ); | |
| assert_eq!( | |
| capture_binding(&manager_a) | |
| .await | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["a_only"] | |
| ); | |
| assert_eq!( | |
| capture_binding(&manager_b) | |
| .await | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["b_only"] | |
| ); | |
| let mut config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| config.server_permission_profiles.insert( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| PermissionProfile::default(), | |
| ); | |
| let manager_a_for_refresh = Arc::clone(&manager_a); | |
| let config_for_refresh = config.clone(); | |
| let refresh_a = tokio::spawn(async move { | |
| manager_a_for_refresh | |
| .refresh_codex_apps_client_catalog(&config_for_refresh) | |
| .await | |
| }); | |
| list_started.notified().await; | |
| manager_b.refresh_codex_apps_tools_for_discovery().await?; | |
| release_list.notify_one(); | |
| let snapshot_a = refresh_a.await??; | |
| assert_eq!( | |
| snapshot_a | |
| .tools | |
| .iter() | |
| .map(|tool| tool.tool.name.as_ref()) | |
| .collect::<Vec<_>>(), | |
| vec!["a_only"] | |
| ); | |
| assert_eq!( | |
| snapshot_a.model_visible_tool_names, | |
| HashSet::from(["a_only".to_string()]) | |
| ); | |
| assert_eq!( | |
| cache_context_a | |
| .current_tools() | |
| .expect("shared cache tools") | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["b_only"] | |
| ); | |
| Arc::get_mut(&mut manager_b) | |
| .expect("unshared manager") | |
| .servers | |
| .get_mut(CODEX_APPS_MCP_SERVER_NAME) | |
| .expect("Apps server") | |
| .tool_filter | |
| .disabled | |
| .insert("b_only".to_string()); | |
| let snapshot_b = manager_b.refresh_codex_apps_client_catalog(&config).await?; | |
| assert_eq!( | |
| ( | |
| snapshot_b | |
| .tools | |
| .iter() | |
| .map(|tool| tool.tool.name.as_ref()) | |
| .collect::<Vec<_>>(), | |
| snapshot_b.model_visible_tool_names, | |
| ), | |
| (vec!["b_only"], HashSet::new()) | |
| ); | |
| Ok(()) | |
| } | |
| async fn tool_catalog_cache_sanitizes_tools_and_tracks_environment_generation() { | |
| let cache = McpToolCatalogCache::default(); | |
| let environment_manager = Arc::new(environment_manager_without_environments()); | |
| let replace_environment = |url: &str| { | |
| environment_manager | |
| .upsert_environment( | |
| "remote".to_string(), | |
| url.to_string(), | |
| /*connect_timeout*/ None, | |
| ) | |
| .expect("replace environment"); | |
| }; | |
| replace_environment("ws://127.0.0.1:1"); | |
| let runtime_context = | |
| McpRuntimeContext::new(Arc::clone(&environment_manager), PathBuf::from("/tmp")); | |
| let config: McpServerConfig = serde_json::from_value(serde_json::json!({ | |
| "command": "docs-mcp", | |
| "environment_id": "remote" | |
| })) | |
| .expect("MCP config"); | |
| let resolve_environment = || { | |
| runtime_context | |
| .resolve_server_environment("docs", &config) | |
| .expect("resolve environment") | |
| .expect("remote environment") | |
| }; | |
| let cache_context = |environment: &Arc<codex_exec_server::Environment>| { | |
| cache | |
| .context( | |
| "docs", | |
| &config, | |
| &runtime_context, | |
| Some(environment), | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .expect("cache context") | |
| }; | |
| let first_environment = resolve_environment(); | |
| let first_environment_weak = Arc::downgrade(&first_environment); | |
| let first_context = cache_context(&first_environment); | |
| first_context.publish_if_newest(first_context.begin_fetch(), &[]); | |
| assert!(!first_context.has_tools()); | |
| let mut tool = create_test_tool("docs", "search"); | |
| tool.tool.annotations = Some(rmcp::model::ToolAnnotations::new().read_only(true)); | |
| first_context.publish_if_newest(first_context.begin_fetch(), &[tool]); | |
| assert_eq!( | |
| first_context.current_tools().expect("cached tools")[0] | |
| .tool | |
| .annotations, | |
| None | |
| ); | |
| drop(first_environment); | |
| replace_environment("ws://127.0.0.1:2"); | |
| assert!(first_environment_weak.upgrade().is_none()); | |
| let replacement_environment = resolve_environment(); | |
| assert!(!cache_context(&replacement_environment).has_tools()); | |
| let older = first_context.begin_fetch(); | |
| let newer = first_context.begin_fetch(); | |
| first_context.publish_if_newest(newer, &[create_test_tool("docs", "new")]); | |
| first_context.publish_if_newest(older, &[create_test_tool("docs", "old")]); | |
| assert_eq!( | |
| first_context.current_tools().expect("cached tools")[0].callable_name, | |
| "new" | |
| ); | |
| tokio::time::advance(Duration::from_secs(30 * 60 + 1)).await; | |
| assert!(!first_context.has_tools()); | |
| } | |
| fn tool_catalog_cache_bypasses_remote_sourced_environment_variables() { | |
| let cache = McpToolCatalogCache::default(); | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| PathBuf::from("/tmp"), | |
| ); | |
| let config: McpServerConfig = serde_json::from_value(serde_json::json!({ | |
| "command": "docs-mcp", | |
| "env_vars": [McpServerEnvVar::Config { | |
| name: "DOCS_TOKEN".to_string(), | |
| source: Some("remote".to_string()), | |
| }], | |
| })) | |
| .expect("MCP config"); | |
| assert!( | |
| cache | |
| .context( | |
| "docs", | |
| &config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default() | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .is_none() | |
| ); | |
| } | |
| fn tool_catalog_cache_bypasses_http_headers_helpers() { | |
| let cache = McpToolCatalogCache::default(); | |
| let runtime_context = reusable_server_runtime_context(); | |
| let mut config = reusable_server_config("https://example.com/mcp"); | |
| let identity = reusable_server_identity("docs", &config, &runtime_context); | |
| let context = |config: &McpServerConfig, identity: &McpServerConnectionIdentity| { | |
| cache.context( | |
| "docs", | |
| config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| Some(( | |
| identity, | |
| crate::McpProtocolMode::Legacy, | |
| /*agent_plugin*/ false, | |
| )), | |
| ) | |
| }; | |
| assert!(context(&config, &identity).is_some()); | |
| let McpServerTransportConfig::StreamableHttp { | |
| http_headers_helper, | |
| .. | |
| } = &mut config.transport | |
| else { | |
| unreachable!("expected HTTP transport"); | |
| }; | |
| *http_headers_helper = Some("auth-cli headers".to_string()); | |
| let identity = reusable_server_identity("docs", &config, &runtime_context); | |
| assert!(context(&config, &identity).is_none()); | |
| } | |
| async fn list_available_server_infos_uses_cache_while_client_is_pending() { | |
| let pending_client = futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let server_info = create_test_server_info("Codex Apps"); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: Some(server_info.clone()), | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let timeout_result = tokio::time::timeout( | |
| Duration::from_millis(10), | |
| manager.list_available_server_infos(), | |
| ) | |
| .await; | |
| let server_infos = timeout_result.expect("server info lookup should not block on startup"); | |
| assert_eq!( | |
| server_infos.get(CODEX_APPS_MCP_SERVER_NAME), | |
| Some(&server_info) | |
| ); | |
| } | |
| async fn list_all_tools_accepts_canonical_namespaced_tool_names() { | |
| let managed_client = | |
| create_ready_async_managed_client(vec![create_test_tool("rmcp", "echo")]).await; | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ false, | |
| ); | |
| manager.insert_test_client("rmcp", managed_client); | |
| let tools = manager.list_all_tools().await; | |
| let tool = tools | |
| .iter() | |
| .find(|tool| tool.canonical_tool_name() == ToolName::namespaced("rmcp", "echo")) | |
| .expect("split MCP tool namespace and name should resolve"); | |
| let expected = ("rmcp", "rmcp", "echo", "echo"); | |
| assert_eq!( | |
| ( | |
| tool.server_name.as_str(), | |
| tool.callable_namespace.as_str(), | |
| tool.callable_name.as_str(), | |
| tool.tool.name.as_ref(), | |
| ), | |
| expected | |
| ); | |
| } | |
| async fn capture_binding_exposes_cached_tools_before_startup() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| let mut cached_tool = create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "shared_cached_tool"); | |
| cached_tool.tool.annotations = Some( | |
| rmcp::model::ToolAnnotations::new() | |
| .read_only(true) | |
| .destructive(false) | |
| .open_world(false), | |
| ); | |
| store_current_tools(&cache_context, vec![cached_tool]); | |
| let startup_complete = Arc::new(std::sync::atomic::AtomicBool::new(false)); | |
| let startup_complete_for_client = Arc::clone(&startup_complete); | |
| let (startup_started, wait_for_startup) = tokio::sync::oneshot::channel(); | |
| let (release_startup, startup_released) = tokio::sync::oneshot::channel(); | |
| let pending_client = async move { | |
| startup_started.send(()).expect("signal client startup"); | |
| startup_released.await.expect("release client startup"); | |
| startup_complete_for_client.store(true, std::sync::atomic::Ordering::Release); | |
| Ok(create_test_managed_client(vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "client_local_tool", | |
| )]) | |
| .await) | |
| } | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete, | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| manager.set_test_server_metadata( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: false, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| ); | |
| let manager = Arc::new(manager); | |
| let cached_binding = capture_binding(&manager).await; | |
| assert_eq!( | |
| cached_binding | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["shared_cached_tool"] | |
| ); | |
| assert_eq!( | |
| cached_binding.tools()[0].tool.annotations, | |
| Some( | |
| rmcp::model::ToolAnnotations::new() | |
| .destructive(false) | |
| .open_world(false) | |
| ) | |
| ); | |
| assert!( | |
| cached_binding | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "shared_cached_tool") | |
| .is_none() | |
| ); | |
| let manager_for_startup = Arc::clone(&manager); | |
| let startup = tokio::spawn(async move { | |
| manager_for_startup | |
| .wait_for_server_startup(CODEX_APPS_MCP_SERVER_NAME) | |
| .await | |
| }); | |
| wait_for_startup.await.expect("client startup should begin"); | |
| release_startup.send(()).expect("release client startup"); | |
| assert!(startup.await.expect("startup task")); | |
| let step = capture_binding(&manager).await; | |
| assert_eq!( | |
| step.tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["client_local_tool"] | |
| ); | |
| } | |
| async fn capture_binding_skips_pending_optional_servers_after_configured_shared_startup_grace() { | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let mut plugin_config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| let mut catalog = crate::ResolvedMcpCatalog::builder(); | |
| catalog.register(crate::McpServerRegistration::from_plugin( | |
| "pending-one".to_string(), | |
| crate::McpPluginAttribution::new("optional-plugin".to_string(), "Optional".to_string()), | |
| /*plugin_order*/ 0, | |
| serde_json::from_value(serde_json::json!({ "command": "optional-plugin" })) | |
| .expect("optional plugin MCP config"), | |
| )); | |
| catalog.register(crate::McpServerRegistration::from_selected_plugin( | |
| "pending-selected".to_string(), | |
| crate::McpPluginAttribution::new("selected-plugin".to_string(), "Selected".to_string()), | |
| /*selection_order*/ 0, | |
| serde_json::from_value(serde_json::json!({ "command": "selected-plugin" })) | |
| .expect("selected plugin MCP config"), | |
| )); | |
| plugin_config.mcp_server_catalog = catalog.build(); | |
| plugin_config.optional_mcp_startup_grace = Duration::from_millis(250); | |
| manager.tool_plugin_context = Arc::new(crate::tool_plugin_context(&plugin_config)); | |
| for server_name in ["pending-one", "pending-two", "pending-selected"] { | |
| manager.insert_test_client( | |
| server_name.to_string(), | |
| AsyncManagedClient { | |
| client: futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(), | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| } | |
| let mut required_manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| required_manager.tool_plugin_context = Arc::clone(&manager.tool_plugin_context); | |
| required_manager.insert_test_client( | |
| "pending-selected", | |
| manager.test_client("pending-selected").clone(), | |
| ); | |
| required_manager.required_servers = vec!["pending-selected".to_string()]; | |
| let manager = Arc::new(manager); | |
| assert!( | |
| manager | |
| .stable_catalog_revisions( | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new() | |
| ) | |
| .await | |
| .is_none() | |
| ); | |
| let started = tokio::time::Instant::now(); | |
| let binding = tokio::time::timeout( | |
| Duration::from_millis(500), | |
| manager.capture_binding_with_metadata( | |
| Arc::new(plugin_config), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new(), | |
| ), | |
| ) | |
| .await | |
| .expect("all optional servers should share the configured startup grace"); | |
| assert!(binding.tools().is_empty()); | |
| assert_eq!(started.elapsed(), Duration::from_millis(250)); | |
| let binding = tokio::time::timeout(Duration::from_millis(1), capture_binding(&manager)) | |
| .await | |
| .expect("later bindings must not restart the optional startup grace"); | |
| assert!(binding.tools().is_empty()); | |
| assert!( | |
| tokio::time::timeout( | |
| Duration::from_millis(1), | |
| binding.list_resources("pending-one", /*params*/ None), | |
| ) | |
| .await | |
| .is_err(), | |
| "resources must wait for an omitted server instead of failing immediately" | |
| ); | |
| assert!( | |
| tokio::time::timeout( | |
| Duration::from_millis(1), | |
| binding.list_all_resources(|server| server == "pending-one"), | |
| ) | |
| .await | |
| .is_ok(), | |
| "resource discovery must not wait for an omitted optional server" | |
| ); | |
| for server_name in ["pending-one", "pending-selected"] { | |
| let required_servers = vec![server_name.to_string()]; | |
| let binding = tokio::time::timeout( | |
| Duration::from_millis(1), | |
| manager.capture_binding_with_metadata( | |
| Arc::new(crate::mcp::tests::test_mcp_config(std::env::temp_dir())), | |
| /*plugins_available*/ false, | |
| &required_servers, | |
| /*required_plugins*/ &HashSet::new(), | |
| ), | |
| ) | |
| .await; | |
| assert!(binding.is_err(), "explicitly requested servers must wait"); | |
| } | |
| // A plugin mention must still require startup after the optional grace has elapsed. | |
| for (plugin_id, must_wait) in [ | |
| ("selected-plugin", true), | |
| ("optional-plugin", false), | |
| ("selected-plugin-other", false), | |
| ] { | |
| let required_plugins = HashSet::from([plugin_id.to_string()]); | |
| let binding = tokio::time::timeout( | |
| Duration::from_millis(1), | |
| manager.capture_binding_with_metadata( | |
| Arc::new(crate::mcp::tests::test_mcp_config(std::env::temp_dir())), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| &required_plugins, | |
| ), | |
| ) | |
| .await; | |
| assert_eq!( | |
| binding.is_err(), | |
| must_wait, | |
| "plugin requirement {plugin_id}" | |
| ); | |
| } | |
| assert!( | |
| tokio::time::timeout( | |
| Duration::from_millis(1500), | |
| capture_binding(&Arc::new(required_manager)), | |
| ) | |
| .await | |
| .is_err(), | |
| "configured-required selected plugin servers must wait beyond the optional grace" | |
| ); | |
| } | |
| async fn capture_binding_waits_for_optional_startup_when_shared_grace_is_disabled() { | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let (client, startup_started, release_startup) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("optional", "echo")]).await, | |
| ); | |
| manager.insert_test_client("optional", client); | |
| let mut config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| config.optional_mcp_startup_grace = Duration::ZERO; | |
| config | |
| .server_permission_profiles | |
| .insert("optional".to_string(), PermissionProfile::default()); | |
| let manager = Arc::new(manager); | |
| let mut capture = tokio::spawn(async move { | |
| manager | |
| .capture_binding_with_metadata( | |
| Arc::new(config), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new(), | |
| ) | |
| .await | |
| }); | |
| startup_started.await.expect("client startup should begin"); | |
| assert!( | |
| tokio::time::timeout(Duration::from_millis(1), &mut capture) | |
| .await | |
| .is_err(), | |
| "disabled shared grace should keep waiting for optional startup" | |
| ); | |
| release_startup.send(()).expect("release client startup"); | |
| let binding = capture.await.expect("capture binding task"); | |
| assert!(binding.prepare_call("optional", "echo").is_some()); | |
| } | |
| async fn stable_catalog_revisions_ignore_terminal_optional_server_failures() { | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let ready = create_ready_async_managed_client(vec![create_test_tool("ready", "echo")]).await; | |
| assert!(ready.client().await.is_ok()); | |
| let mut failed = ready.clone(); | |
| manager.insert_test_client("ready", ready); | |
| failed.client = futures::future::ready::<Result<ManagedClient, StartupOutcomeError>>(Err( | |
| StartupOutcomeError::Failed { | |
| error: "optional startup failed".to_string(), | |
| is_authentication_required: false, | |
| }, | |
| )) | |
| .boxed() | |
| .shared(); | |
| assert!(failed.client().await.is_err()); | |
| manager.insert_test_client("failed", failed); | |
| assert!( | |
| manager | |
| .stable_catalog_revisions( | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new() | |
| ) | |
| .await | |
| .is_some() | |
| ); | |
| manager.required_servers.push("failed".to_string()); | |
| assert!( | |
| manager | |
| .stable_catalog_revisions( | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new() | |
| ) | |
| .await | |
| .is_none() | |
| ); | |
| manager.required_servers.clear(); | |
| let binding = capture_binding(&Arc::new(manager)).await; | |
| assert!(binding.prepare_call("ready", "echo").is_some()); | |
| } | |
| async fn capture_binding_shares_optional_startup_grace_across_connection_sets() { | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let cache = McpToolCatalogCache::default(); | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| std::env::temp_dir(), | |
| ); | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "command": "pending-mcp" })) | |
| .expect("pending MCP server configuration"); | |
| let cache_context = cache | |
| .context( | |
| "pending", | |
| &server_config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .expect("shared pending MCP catalog"); | |
| let create_connection_set = || { | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| "pending", | |
| AsyncManagedClient { | |
| client: futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(), | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: Some(cache_context.clone()), | |
| startup_complete: Arc::new(AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| Arc::new(manager) | |
| }; | |
| let first_started = tokio::time::Instant::now(); | |
| let first = tokio::time::timeout( | |
| Duration::from_millis(1500), | |
| capture_binding(&create_connection_set()), | |
| ) | |
| .await | |
| .expect("the first thread should receive the optional startup grace"); | |
| assert!(first.tools().is_empty()); | |
| assert_eq!(first_started.elapsed(), Duration::from_secs(1)); | |
| let second = tokio::time::timeout( | |
| Duration::from_millis(1), | |
| capture_binding(&create_connection_set()), | |
| ) | |
| .await | |
| .expect("the next thread must not restart the same server's startup grace"); | |
| assert!(second.tools().is_empty()); | |
| let mut disabled_config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| disabled_config.optional_mcp_startup_grace = Duration::ZERO; | |
| let disabled_manager = create_connection_set(); | |
| assert!( | |
| tokio::time::timeout( | |
| Duration::from_millis(1), | |
| disabled_manager.capture_binding_with_metadata( | |
| Arc::new(disabled_config), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new(), | |
| ), | |
| ) | |
| .await | |
| .is_err(), | |
| "disabled grace should keep waiting for the pending optional server" | |
| ); | |
| let restored_started = tokio::time::Instant::now(); | |
| let restored = tokio::time::timeout( | |
| Duration::from_millis(1500), | |
| capture_binding(&create_connection_set()), | |
| ) | |
| .await | |
| .expect("restoring the startup grace should create a fresh deadline"); | |
| assert!(restored.tools().is_empty()); | |
| assert_eq!(restored_started.elapsed(), Duration::from_secs(1)); | |
| let mut updated_config = crate::mcp::tests::test_mcp_config(std::env::temp_dir()); | |
| updated_config.optional_mcp_startup_grace = Duration::from_millis(250); | |
| let updated_manager = create_connection_set(); | |
| let updated_started = tokio::time::Instant::now(); | |
| let updated = tokio::time::timeout( | |
| Duration::from_millis(500), | |
| updated_manager.capture_binding_with_metadata( | |
| Arc::new(updated_config), | |
| /*plugins_available*/ false, | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new(), | |
| ), | |
| ) | |
| .await | |
| .expect("a changed startup grace should receive its newly configured deadline"); | |
| assert!(updated.tools().is_empty()); | |
| assert_eq!(updated_started.elapsed(), Duration::from_millis(250)); | |
| cache_context.publish_if_newest( | |
| cache_context.begin_fetch(), | |
| &[create_test_tool("pending", "cached_tool")], | |
| ); | |
| let deadline_after_publication = tokio::time::Instant::now() + Duration::from_secs(1); | |
| assert_eq!( | |
| cache_context.optional_startup_deadline(deadline_after_publication, Duration::from_secs(1)), | |
| deadline_after_publication, | |
| "publishing a catalog must not install a stale startup deadline" | |
| ); | |
| let cached_manager = create_connection_set(); | |
| let cached = tokio::time::timeout(Duration::from_millis(1), capture_binding(&cached_manager)) | |
| .await | |
| .expect("cached tools should be immediately available to later threads"); | |
| assert_eq!( | |
| cached | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_tool"] | |
| ); | |
| tokio::time::advance(Duration::from_secs(30 * 60 + 1)).await; | |
| assert!( | |
| tokio::time::timeout(Duration::from_millis(1), capture_binding(&cached_manager)) | |
| .await | |
| .is_err(), | |
| "an expired catalog should receive a fresh startup grace" | |
| ); | |
| cache_context.disable(); | |
| for _ in 0..2 { | |
| let started = tokio::time::Instant::now(); | |
| let binding = tokio::time::timeout( | |
| Duration::from_millis(1500), | |
| capture_binding(&create_connection_set()), | |
| ) | |
| .await | |
| .expect("non-cacheable servers should keep their per-thread startup grace"); | |
| assert!(binding.tools().is_empty()); | |
| assert_eq!(started.elapsed(), Duration::from_secs(1)); | |
| } | |
| } | |
| async fn capture_binding_uses_cache_published_during_optional_startup() { | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| std::env::temp_dir(), | |
| ); | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "command": "pending-mcp" })) | |
| .expect("server configuration"); | |
| let cache_context = McpToolCatalogCache::default() | |
| .context( | |
| "pending", | |
| &server_config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .expect("shared catalog"); | |
| let (mut client, started, _release) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("pending", "live_tool")]).await, | |
| ); | |
| client.tool_catalog_cache_context = Some(cache_context.clone()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &Constrained::allow_any(AskForApproval::OnRequest), | |
| &Constrained::allow_any(PermissionProfile::default()), | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("pending", client); | |
| let manager = Arc::new(manager); | |
| let binding = capture_binding(&manager); | |
| tokio::pin!(binding); | |
| assert!(futures::poll!(&mut binding).is_pending()); | |
| started | |
| .await | |
| .expect("optional startup began without a cache"); | |
| cache_context.publish_if_newest( | |
| cache_context.begin_fetch(), | |
| &[create_test_tool("pending", "peer_tool")], | |
| ); | |
| tokio::time::advance(Duration::from_secs(/*secs*/ 2)).await; | |
| let binding = binding.await; | |
| assert!( | |
| !manager.servers["pending"] | |
| .connection | |
| .client | |
| .startup_complete | |
| .load(Ordering::Acquire) | |
| ); | |
| assert_eq!( | |
| model_tool_names(binding.tools()), | |
| HashSet::from([ToolName::namespaced("mcp__pending", "peer_tool")]), | |
| ); | |
| } | |
| async fn capture_binding_retains_cached_tools_that_expire_while_waiting_for_another_server() { | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| std::env::temp_dir(), | |
| ); | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "command": "server-a-mcp" })) | |
| .expect("MCP server configuration"); | |
| let cache_context = McpToolCatalogCache::default() | |
| .context( | |
| "server_a", | |
| &server_config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .expect("server A cache context"); | |
| let (mut client_a, _started_a, _release_a) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("server_a", "live_tool")]).await, | |
| ); | |
| client_a.tool_catalog_cache_context = Some(cache_context.clone()); | |
| let (client_b, started_b, release_b) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("server_b", "tool_b")]).await, | |
| ); | |
| cache_context.publish_if_newest( | |
| cache_context.begin_fetch(), | |
| &[create_test_tool("server_a", "cached_tool")], | |
| ); | |
| tokio::time::advance(Duration::from_secs(30 * 60 - 1)).await; | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &Constrained::allow_any(AskForApproval::OnRequest), | |
| &Constrained::allow_any(PermissionProfile::default()), | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("server_a", client_a); | |
| manager.insert_test_client("server_b", client_b); | |
| Arc::get_mut(&mut manager.servers.get_mut("server_a").unwrap().connection) | |
| .expect("unique server A connection") | |
| .startup_trigger = Some(watch::channel(/*init*/ false).0); | |
| manager.required_servers = vec!["server_a".to_string(), "server_b".to_string()]; | |
| let manager = Arc::new(manager); | |
| let binding = capture_binding(&manager); | |
| tokio::pin!(binding); | |
| assert!(futures::poll!(&mut binding).is_pending()); | |
| started_b.await.expect("server B startup should begin"); | |
| assert!(manager.servers["server_a"].connection.startup_is_dormant()); | |
| // Expire A's catalog while capture is waiting for B's required startup. | |
| tokio::time::advance(Duration::from_secs(/*secs*/ 2)).await; | |
| assert!(cache_context.current_tools().is_none()); | |
| release_b.send(()).expect("release server B startup"); | |
| let binding = tokio::time::timeout(Duration::from_secs(/*secs*/ 1), binding) | |
| .await | |
| .expect("cached server A should not need startup"); | |
| assert_eq!( | |
| model_tool_names(binding.tools()), | |
| HashSet::from([ | |
| ToolName::namespaced("mcp__server_a", "cached_tool"), | |
| ToolName::namespaced("mcp__server_b", "tool_b"), | |
| ]) | |
| ); | |
| assert!(manager.servers["server_a"].connection.startup_is_dormant()); | |
| } | |
| async fn capture_binding_omits_cache_disabled_while_waiting_for_another_server() { | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| std::env::temp_dir(), | |
| ); | |
| let server_config: McpServerConfig = | |
| serde_json::from_value(serde_json::json!({ "command": "cached-mcp" })) | |
| .expect("server configuration"); | |
| let cache_context = McpToolCatalogCache::default() | |
| .context( | |
| "cached", | |
| &server_config, | |
| &runtime_context, | |
| /*resolved_environment*/ None, | |
| ( | |
| &ElicitationCapability::default(), | |
| &ClientMcpExtensions::default(), | |
| ), | |
| /*connection_identity*/ None, | |
| ) | |
| .expect("shared catalog"); | |
| cache_context.publish_if_newest( | |
| cache_context.begin_fetch(), | |
| &[create_test_tool("cached", "cached_tool")], | |
| ); | |
| let (mut cached, _started, _release) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("cached", "live_tool")]).await, | |
| ); | |
| cached.tool_catalog_cache_context = Some(cache_context.clone()); | |
| let (waiting, started, release) = create_gated_async_managed_client( | |
| create_test_managed_client(vec![create_test_tool("waiting", "ready_tool")]).await, | |
| ); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &Constrained::allow_any(AskForApproval::OnRequest), | |
| &Constrained::allow_any(PermissionProfile::default()), | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("cached", cached); | |
| manager.insert_test_client("waiting", waiting); | |
| manager.required_servers = vec!["waiting".to_string()]; | |
| let manager = Arc::new(manager); | |
| let binding = capture_binding(&manager); | |
| tokio::pin!(binding); | |
| assert!(futures::poll!(&mut binding).is_pending()); | |
| started.await.expect("required server startup began"); | |
| cache_context.disable(); | |
| release.send(()).expect("release required server startup"); | |
| assert_eq!( | |
| model_tool_names(binding.await.tools()), | |
| HashSet::from([ToolName::namespaced("mcp__waiting", "ready_tool")]), | |
| ); | |
| } | |
| async fn capture_binding_resolves_concurrently_and_rechecks_cached_clients() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| store_current_tools( | |
| &cache_context, | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "shared_cached_tool", | |
| )], | |
| ); | |
| let ready_apps_client = create_test_managed_client(vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "client_local_tool", | |
| )]) | |
| .await; | |
| let (mut apps_client, apps_started, release_apps) = | |
| create_gated_async_managed_client(ready_apps_client); | |
| apps_client.is_codex_apps_mcp_server = true; | |
| apps_client.codex_apps_tools_cache_context = Some(cache_context); | |
| let first_client = | |
| create_test_managed_client(vec![create_test_tool("first", "first_tool")]).await; | |
| let second_client = | |
| create_test_managed_client(vec![create_test_tool("second", "second_tool")]).await; | |
| let (first_client, first_started, release_first) = | |
| create_gated_async_managed_client(first_client); | |
| let (second_client, second_started, release_second) = | |
| create_gated_async_managed_client(second_client); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client(CODEX_APPS_MCP_SERVER_NAME, apps_client); | |
| manager.insert_test_client("first", first_client); | |
| manager.insert_test_client("second", second_client); | |
| let manager = Arc::new(manager); | |
| let manager_for_startup = Arc::clone(&manager); | |
| let startup = tokio::spawn(async move { | |
| manager_for_startup | |
| .wait_for_server_startup(CODEX_APPS_MCP_SERVER_NAME) | |
| .await | |
| }); | |
| tokio::time::timeout(Duration::from_secs(1), apps_started) | |
| .await | |
| .expect("Codex Apps startup should begin") | |
| .expect("signal Codex Apps startup"); | |
| let manager_for_binding = Arc::clone(&manager); | |
| let binding = tokio::spawn(async move { capture_binding(&manager_for_binding).await }); | |
| tokio::time::timeout(Duration::from_secs(1), async { | |
| first_started.await.expect("first server startup"); | |
| second_started.await.expect("second server startup"); | |
| }) | |
| .await | |
| .expect("both uncached servers should start before either is released"); | |
| release_apps.send(()).expect("release Codex Apps startup"); | |
| assert!(startup.await.expect("Codex Apps startup task")); | |
| release_first.send(()).expect("release first server"); | |
| release_second.send(()).expect("release second server"); | |
| let binding = binding.await.expect("binding capture should complete"); | |
| assert_eq!( | |
| binding | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<HashSet<_>>(), | |
| HashSet::from(["client_local_tool", "first_tool", "second_tool"]) | |
| ); | |
| assert!( | |
| binding | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "client_local_tool") | |
| .is_some() | |
| ); | |
| assert!( | |
| binding | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "shared_cached_tool") | |
| .is_none() | |
| ); | |
| assert!(binding.prepare_call("first", "first_tool").is_some()); | |
| assert!(binding.prepare_call("second", "second_tool").is_some()); | |
| } | |
| async fn list_all_tools_applies_legacy_mcp_prefix_by_default() { | |
| let managed_client = | |
| create_ready_async_managed_client(vec![create_test_tool("rmcp", "echo")]).await; | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("rmcp", managed_client); | |
| let tools = manager.list_all_tools().await; | |
| let tool = tools | |
| .iter() | |
| .find(|tool| tool.canonical_tool_name() == ToolName::namespaced("mcp__rmcp", "echo")) | |
| .expect("legacy-prefixed MCP tool name should resolve"); | |
| let expected = ("rmcp", "mcp__rmcp", "echo", "echo"); | |
| assert_eq!( | |
| ( | |
| tool.server_name.as_str(), | |
| tool.callable_namespace.as_str(), | |
| tool.callable_name.as_str(), | |
| tool.tool.name.as_ref(), | |
| ), | |
| expected | |
| ); | |
| } | |
| async fn call_tool_requires_connection_without_waiting_for_startup() { | |
| let client = create_test_managed_client(vec![create_test_tool("docs", "search")]).await; | |
| let (client, startup_started, release_startup) = create_gated_async_managed_client(client); | |
| let startup_client = client.clone(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("docs", client); | |
| let pending_call = tokio::time::timeout( | |
| Duration::from_millis(50), | |
| manager.call_tool( | |
| "docs", | |
| "search", | |
| /*environment_id*/ None, | |
| /*arguments*/ None, | |
| /*meta*/ None, | |
| Some(Duration::from_secs(5)), | |
| /*wait_for_server*/ false, | |
| ), | |
| ) | |
| .await | |
| .expect("ready-only invocation must not wait for pending server startup") | |
| .expect_err("pending server must not accept ready-only calls"); | |
| assert!(pending_call.to_string().contains("not connected")); | |
| let startup = tokio::spawn(async move { startup_client.client().await }); | |
| startup_started.await.expect("server startup should begin"); | |
| release_startup.send(()).expect("release server startup"); | |
| startup | |
| .await | |
| .expect("startup task should finish") | |
| .expect("server startup should succeed"); | |
| let ready_error = manager | |
| .call_tool( | |
| "docs", | |
| "search", | |
| /*environment_id*/ None, | |
| /*arguments*/ None, | |
| /*meta*/ None, | |
| Some(Duration::from_secs(5)), | |
| /*wait_for_server*/ false, | |
| ) | |
| .await | |
| .expect_err("ready server should reach the uninitialized test transport"); | |
| assert!(format!("{ready_error:#}").contains("MCP client not initialized")); | |
| } | |
| async fn connected_call_respects_server_tool_filters() { | |
| let client = create_ready_async_managed_client(vec![create_test_tool("docs", "search")]).await; | |
| client.client().await.expect("server should be ready"); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("docs", client); | |
| manager | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist") | |
| .tool_filter | |
| .disabled | |
| .insert("search".to_string()); | |
| let filtered_error = manager | |
| .call_tool( | |
| "docs", | |
| "search", | |
| /*environment_id*/ None, | |
| /*arguments*/ None, | |
| /*meta*/ None, | |
| Some(Duration::from_secs(5)), | |
| /*wait_for_server*/ false, | |
| ) | |
| .await | |
| .expect_err("disabled tools should not be callable"); | |
| assert!(filtered_error.to_string().contains("disabled")); | |
| } | |
| async fn call_tool_validates_environment_without_waiting_for_ready_connections() { | |
| let client = create_test_managed_client(vec![create_test_tool("docs", "search")]).await; | |
| let (client, _, _) = create_gated_async_managed_client(client); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("docs", client); | |
| manager | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist") | |
| .metadata | |
| .environment_id = "executor-a".to_string(); | |
| let mismatched_environment = manager | |
| .call_tool( | |
| "docs", | |
| "search", | |
| Some("executor-b"), | |
| /*arguments*/ None, | |
| /*meta*/ None, | |
| Some(Duration::from_secs(5)), | |
| /*wait_for_server*/ false, | |
| ) | |
| .await | |
| .expect_err("calls must reject a server from a different environment"); | |
| assert_eq!( | |
| mismatched_environment.to_string(), | |
| "MCP server `docs` is running in environment `executor-a`, expected `executor-b`" | |
| ); | |
| let pending_call = tokio::time::timeout( | |
| Duration::from_millis(50), | |
| manager.call_tool( | |
| "docs", | |
| "search", | |
| Some("executor-a"), | |
| /*arguments*/ None, | |
| /*meta*/ None, | |
| Some(Duration::from_secs(5)), | |
| /*wait_for_server*/ false, | |
| ), | |
| ) | |
| .await | |
| .expect("environment-scoped calls must not wait for pending server startup") | |
| .expect_err("pending server must not accept environment-scoped calls"); | |
| assert!(pending_call.to_string().contains("not connected")); | |
| } | |
| async fn list_all_tools_resolves_server_catalogs_concurrently() { | |
| let first_client = create_test_managed_client(vec![create_test_tool("first", "search")]).await; | |
| let second_client = | |
| create_test_managed_client(vec![create_test_tool("second", "lookup")]).await; | |
| let (first_client, first_started, release_first) = | |
| create_gated_async_managed_client(first_client); | |
| let (second_client, second_started, release_second) = | |
| create_gated_async_managed_client(second_client); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client("first", first_client); | |
| manager.insert_test_client("second", second_client); | |
| let manager = Arc::new(manager); | |
| let manager_for_listing = Arc::clone(&manager); | |
| let listing = tokio::spawn(async move { manager_for_listing.list_all_tools().await }); | |
| tokio::time::timeout(Duration::from_secs(1), async { | |
| first_started.await.expect("first server startup"); | |
| second_started.await.expect("second server startup"); | |
| }) | |
| .await | |
| .expect("both server catalogs should start before either is released"); | |
| release_first.send(()).expect("release first server"); | |
| release_second.send(()).expect("release second server"); | |
| let tools = listing.await.expect("tool listing should complete"); | |
| assert_eq!( | |
| model_tool_names(&tools), | |
| HashSet::from([ | |
| ToolName::namespaced("mcp__first", "search"), | |
| ToolName::namespaced("mcp__second", "lookup"), | |
| ]) | |
| ); | |
| } | |
| async fn list_all_tools_blocks_while_client_is_pending_without_cached_tools() { | |
| let pending_client = futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let timeout_result = | |
| tokio::time::timeout(Duration::from_millis(10), manager.list_all_tools()).await; | |
| assert!(timeout_result.is_err()); | |
| } | |
| async fn cancelling_startup_does_not_disable_a_ready_client() { | |
| let client = create_ready_async_managed_client(vec![create_test_tool("ready", "search")]).await; | |
| client.cancel_token.cancel(); | |
| let managed = client | |
| .client() | |
| .await | |
| .expect("startup cancellation should not disable a ready client"); | |
| assert_eq!( | |
| managed | |
| .tool_catalog | |
| .read(|catalog| model_tool_names(&catalog.tools)) | |
| .await, | |
| HashSet::from([ToolName::namespaced("ready", "search")]) | |
| ); | |
| } | |
| async fn shutdown_cancels_pending_tool_listing() { | |
| let cancel_token = CancellationToken::new(); | |
| let cancel_token_for_startup = cancel_token.clone(); | |
| let (started_tx, started_rx) = tokio::sync::oneshot::channel(); | |
| let pending_client = async move { | |
| let _ = started_tx.send(()); | |
| cancel_token_for_startup.cancelled().await; | |
| Err(StartupOutcomeError::Cancelled) | |
| } | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token, | |
| }, | |
| ); | |
| let manager = Arc::new(manager); | |
| let manager_for_list = Arc::clone(&manager); | |
| let list_task = tokio::spawn(async move { manager_for_list.list_all_tools().await }); | |
| started_rx.await.expect("tool listing should start"); | |
| tokio::time::timeout(Duration::from_secs(1), manager.shutdown()) | |
| .await | |
| .expect("shutdown should cancel speculative tool listing"); | |
| let tools = list_task.await.expect("tool listing task should not panic"); | |
| assert!(tools.is_empty()); | |
| } | |
| async fn shutdown_continues_after_caller_is_aborted() { | |
| let (started_tx, started_rx) = tokio::sync::oneshot::channel(); | |
| let (completed_tx, completed_rx) = tokio::sync::oneshot::channel(); | |
| let release = Arc::new(tokio::sync::Notify::new()); | |
| let release_for_client = Arc::clone(&release); | |
| let blocking_client = async move { | |
| let _ = started_tx.send(()); | |
| release_for_client.notified().await; | |
| let _ = completed_tx.send(()); | |
| Err(StartupOutcomeError::Cancelled) | |
| } | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: blocking_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let manager = Arc::new(manager); | |
| let shutdown_task = tokio::spawn({ | |
| let manager = Arc::clone(&manager); | |
| async move { manager.shutdown().await } | |
| }); | |
| started_rx.await.expect("client shutdown should start"); | |
| shutdown_task.abort(); | |
| let shutdown_error = shutdown_task | |
| .await | |
| .expect_err("caller shutdown task should be aborted"); | |
| assert!(shutdown_error.is_cancelled()); | |
| release.notify_one(); | |
| tokio::time::timeout(Duration::from_secs(1), completed_rx) | |
| .await | |
| .expect("client shutdown should survive caller cancellation") | |
| .expect("client shutdown completion sender should stay alive"); | |
| } | |
| async fn list_all_tools_does_not_block_when_shared_codex_apps_cache_is_empty() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| store_current_tools(&cache_context, Vec::new()); | |
| let pending_client = futures::future::pending::<Result<ManagedClient, StartupOutcomeError>>() | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: pending_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(false)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let timeout_result = | |
| tokio::time::timeout(Duration::from_millis(10), manager.list_all_tools()).await; | |
| let tools = timeout_result.expect("shared empty cache should not block"); | |
| assert!(tools.is_empty()); | |
| } | |
| async fn list_all_tools_uses_shared_codex_apps_cache_when_client_startup_fails() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| Some("account-one"), | |
| Some("user-one"), | |
| ); | |
| store_current_tools( | |
| &cache_context, | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "calendar_create_event", | |
| )], | |
| ); | |
| let server_info = create_test_server_info("Codex Apps"); | |
| let failed_client = futures::future::ready::<Result<ManagedClient, StartupOutcomeError>>(Err( | |
| StartupOutcomeError::Failed { | |
| error: "startup failed".to_string(), | |
| is_authentication_required: false, | |
| }, | |
| )) | |
| .boxed() | |
| .shared(); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let startup_complete = Arc::new(std::sync::atomic::AtomicBool::new(true)); | |
| manager.insert_test_client( | |
| CODEX_APPS_MCP_SERVER_NAME.to_string(), | |
| AsyncManagedClient { | |
| client: failed_client, | |
| is_codex_apps_mcp_server: true, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: Some(server_info.clone()), | |
| codex_apps_tools_cache_context: Some(cache_context), | |
| tool_catalog_cache_context: None, | |
| startup_complete, | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| ); | |
| let tools = manager.list_all_tools().await; | |
| let tool = tools | |
| .iter() | |
| .find(|tool| { | |
| tool.canonical_tool_name() | |
| == ToolName::namespaced("mcp__codex_apps", "calendar_create_event") | |
| }) | |
| .expect("tool from shared cache"); | |
| assert_eq!(tool.server_name, CODEX_APPS_MCP_SERVER_NAME); | |
| assert_eq!(tool.callable_name, "calendar_create_event"); | |
| assert_eq!( | |
| manager | |
| .list_available_server_infos() | |
| .await | |
| .get(CODEX_APPS_MCP_SERVER_NAME), | |
| Some(&server_info) | |
| ); | |
| } | |
| async fn list_all_tools_reconnects_failed_codex_apps_startup_and_reuses_client() { | |
| let recovered_client = create_test_managed_client(vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "drive_search", | |
| )]) | |
| .await; | |
| let attempts = Arc::new(AtomicUsize::new(0)); | |
| let attempts_for_reconnect = Arc::clone(&attempts); | |
| let reconnect_finished = Arc::new(tokio::sync::Notify::new()); | |
| let reconnect_finished_for_factory = Arc::clone(&reconnect_finished); | |
| let reconnect_factory = Arc::new(move || { | |
| attempts_for_reconnect.fetch_add(1, std::sync::atomic::Ordering::SeqCst); | |
| let reconnect_finished = Arc::clone(&reconnect_finished_for_factory); | |
| let recovered_client = recovered_client.clone(); | |
| async move { | |
| reconnect_finished.notify_one(); | |
| Ok(recovered_client) | |
| } | |
| .boxed() | |
| .shared() | |
| }); | |
| let mut manager = create_test_manager_with_failed_apps_startup(Vec::new(), reconnect_factory); | |
| manager | |
| .servers | |
| .get_mut(CODEX_APPS_MCP_SERVER_NAME) | |
| .expect("test server exists") | |
| .metadata = McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: false, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }; | |
| let manager = Arc::new(manager); | |
| assert!( | |
| manager | |
| .stable_catalog_revisions( | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new() | |
| ) | |
| .await | |
| .is_none() | |
| ); | |
| let reconnect_finished_wait = reconnect_finished.notified(); | |
| let tools = manager.list_all_tools().await; | |
| assert!(tools.is_empty()); | |
| reconnect_finished_wait.await; | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["drive_search"] | |
| ); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| assert!( | |
| manager | |
| .stable_catalog_revisions( | |
| /*required_servers*/ &[], | |
| /*required_plugins*/ &HashSet::new() | |
| ) | |
| .await | |
| .is_some() | |
| ); | |
| let step = capture_binding(&manager).await; | |
| let prepared = step | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "drive_search") | |
| .expect("recovered tool should have a prepared call"); | |
| assert!( | |
| !prepared | |
| .server_supports_sandbox_state_meta_capability() | |
| .await | |
| .expect("prepared call should use the recovered client") | |
| ); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["drive_search"] | |
| ); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| } | |
| async fn later_tool_list_retries_after_failed_reconnect_and_keeps_cached_tools() { | |
| let recovered_client = create_test_managed_client(vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "drive_search", | |
| )]) | |
| .await; | |
| let attempts = Arc::new(AtomicUsize::new(0)); | |
| let attempts_for_reconnect = Arc::clone(&attempts); | |
| let reconnect_finished = Arc::new(tokio::sync::Notify::new()); | |
| let reconnect_finished_for_factory = Arc::clone(&reconnect_finished); | |
| let reconnect_factory = Arc::new(move || { | |
| let attempt = attempts_for_reconnect.fetch_add(1, std::sync::atomic::Ordering::SeqCst); | |
| let reconnect_finished = Arc::clone(&reconnect_finished_for_factory); | |
| let recovered_client = recovered_client.clone(); | |
| async move { | |
| let result = if attempt < 2 { | |
| Err(StartupOutcomeError::Failed { | |
| error: "recreated startup failed".to_string(), | |
| is_authentication_required: false, | |
| }) | |
| } else { | |
| Ok(recovered_client) | |
| }; | |
| reconnect_finished.notify_one(); | |
| result | |
| } | |
| .boxed() | |
| .shared() | |
| }); | |
| let manager = create_test_manager_with_failed_apps_startup( | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "cached_drive_search", | |
| )], | |
| reconnect_factory, | |
| ); | |
| let first_reconnect_finished = reconnect_finished.notified(); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| first_reconnect_finished.await; | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| tokio::time::advance(CODEX_APPS_RECONNECT_INITIAL_BACKOFF).await; | |
| let second_reconnect_finished = reconnect_finished.notified(); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| second_reconnect_finished.await; | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 2); | |
| tokio::time::advance(CODEX_APPS_RECONNECT_INITIAL_BACKOFF).await; | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 2); | |
| tokio::time::advance(CODEX_APPS_RECONNECT_INITIAL_BACKOFF).await; | |
| let third_reconnect_finished = reconnect_finished.notified(); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| third_reconnect_finished.await; | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 3); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["drive_search"] | |
| ); | |
| } | |
| async fn tool_lists_do_not_block_and_share_codex_apps_startup_reconnect() { | |
| let recovered_client = create_test_managed_client(vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "drive_search", | |
| )]) | |
| .await; | |
| let attempts = Arc::new(AtomicUsize::new(0)); | |
| let attempts_for_reconnect = Arc::clone(&attempts); | |
| let reconnect_started = Arc::new(tokio::sync::Notify::new()); | |
| let reconnect_started_for_factory = Arc::clone(&reconnect_started); | |
| let release_reconnect = Arc::new(tokio::sync::Notify::new()); | |
| let release_reconnect_for_factory = Arc::clone(&release_reconnect); | |
| let reconnect_factory = Arc::new(move || { | |
| let recovered_client = recovered_client.clone(); | |
| let attempts = Arc::clone(&attempts_for_reconnect); | |
| let reconnect_started = Arc::clone(&reconnect_started_for_factory); | |
| let release_reconnect = Arc::clone(&release_reconnect_for_factory); | |
| async move { | |
| attempts.fetch_add(1, std::sync::atomic::Ordering::SeqCst); | |
| reconnect_started.notify_one(); | |
| release_reconnect.notified().await; | |
| Ok(recovered_client) | |
| } | |
| .boxed() | |
| .shared() | |
| }); | |
| let mut manager = create_test_manager_with_failed_apps_startup( | |
| vec![create_test_tool( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| "cached_drive_search", | |
| )], | |
| reconnect_factory, | |
| ); | |
| manager.set_test_server_metadata( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: false, | |
| origin: None, | |
| supports_parallel_tool_calls: false, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| ); | |
| let manager = Arc::new(manager); | |
| let reconnect_started_wait = reconnect_started.notified(); | |
| let first_tools = tokio::time::timeout(Duration::from_millis(10), manager.list_all_tools()) | |
| .await | |
| .expect("cached tools should not wait for reconnect"); | |
| reconnect_started_wait.await; | |
| let second_tools = tokio::time::timeout(Duration::from_millis(10), manager.list_all_tools()) | |
| .await | |
| .expect("concurrent cached tools should not wait for reconnect"); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| assert_eq!( | |
| first_tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| assert_eq!( | |
| second_tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["cached_drive_search"] | |
| ); | |
| let pending_step = tokio::time::timeout(Duration::from_millis(10), capture_binding(&manager)) | |
| .await | |
| .expect("step capture should not wait for reconnect"); | |
| assert!( | |
| pending_step.tools().is_empty(), | |
| "a model step must not advertise cached tools without an exact ready client" | |
| ); | |
| release_reconnect.notify_one(); | |
| tokio::task::yield_now().await; | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!( | |
| tools | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["drive_search"] | |
| ); | |
| let recovered_step = capture_binding(&manager).await; | |
| assert_eq!( | |
| recovered_step | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.callable_name.as_str()) | |
| .collect::<Vec<_>>(), | |
| vec!["drive_search"] | |
| ); | |
| assert!( | |
| recovered_step | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "drive_search") | |
| .is_some() | |
| ); | |
| assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 1); | |
| } | |
| async fn list_all_tools_adds_server_metadata_to_tools() { | |
| let server_name = "docs"; | |
| let managed_client = | |
| create_ready_async_managed_client(vec![create_test_tool(server_name, "search")]).await; | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| manager.insert_test_client(server_name, managed_client); | |
| manager.set_test_server_metadata( | |
| server_name, | |
| McpServerMetadata { | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| pollutes_memory: true, | |
| origin: Some(McpServerOrigin::StreamableHttp( | |
| "https://docs.example".to_string(), | |
| )), | |
| supports_parallel_tool_calls: true, | |
| default_tools_approval_mode: None, | |
| tool_approval_modes: HashMap::new(), | |
| }, | |
| ); | |
| let tools = manager.list_all_tools().await; | |
| assert_eq!(tools.len(), 1); | |
| let tool = &tools[0]; | |
| assert_eq!(tool.server_name, server_name); | |
| assert!(tool.supports_parallel_tool_calls); | |
| assert_eq!(tool.server_origin.as_deref(), Some("https://docs.example")); | |
| } | |
| fn server_metadata_preserves_tool_approval_policy() { | |
| let mut config = crate::codex_apps_mcp_server_config( | |
| "https://docs.example", | |
| /*apps_mcp_product_sku*/ None, | |
| /*originator*/ None, | |
| ); | |
| config.environment_id = "remote".to_string(); | |
| config.default_tools_approval_mode = Some(AppToolApproval::Prompt); | |
| config.tools.insert( | |
| "search".to_string(), | |
| McpServerToolConfig { | |
| approval_mode: Some(AppToolApproval::Approve), | |
| ..Default::default() | |
| }, | |
| ); | |
| let metadata = McpServerMetadata::from(&EffectiveMcpServer::configured(config)); | |
| assert_eq!(metadata.environment_id, "remote"); | |
| assert_eq!(metadata.tool_approval_mode("read"), AppToolApproval::Prompt); | |
| assert_eq!( | |
| metadata.tool_approval_mode("search"), | |
| AppToolApproval::Approve | |
| ); | |
| } | |
| fn hosted_actor_credentials_are_only_available_to_host_owned_mcp_servers() { | |
| let bootstrap_auth = CodexAuth::create_dummy_chatgpt_auth_for_testing(); | |
| let mut actor_headers = Default::default(); | |
| codex_model_provider::auth_provider_from_auth(&bootstrap_auth) | |
| .add_auth_headers(&mut actor_headers); | |
| actor_headers.insert( | |
| "x-openai-actor-authorization", | |
| "hosted-actor-secret" | |
| .parse() | |
| .expect("valid actor authorization header"), | |
| ); | |
| let hosted_auth = CodexAuth::Headers(AuthHeaders::new(actor_headers)); | |
| let provider = codex_model_provider::auth_provider_from_auth(&hosted_auth); | |
| let mut local_config = crate::codex_apps_mcp_server_config( | |
| "https://chatgpt.com", | |
| /*apps_mcp_product_sku*/ None, | |
| /*originator*/ None, | |
| ); | |
| local_config.auth = McpServerAuth::ChatGpt; | |
| let local_server = EffectiveMcpServer::configured(local_config.clone()); | |
| let local_provider = | |
| chatgpt_auth_provider_for_server(&local_server, Some(Arc::clone(&provider))) | |
| .expect("host-owned Codex Apps must retain hosted authentication"); | |
| assert_eq!( | |
| local_provider | |
| .to_auth_headers() | |
| .get("x-openai-actor-authorization") | |
| .and_then(|value| value.to_str().ok()), | |
| Some("hosted-actor-secret") | |
| ); | |
| let mut remote_config = local_config; | |
| remote_config.environment_id = "customer-executor".to_string(); | |
| let remote_server = EffectiveMcpServer::configured(remote_config); | |
| assert!( | |
| chatgpt_auth_provider_for_server(&remote_server, Some(provider)).is_none(), | |
| "customer-owned executors must never receive hosted actor credentials" | |
| ); | |
| } | |
| async fn executor_owned_chatgpt_mcp_accepts_only_safe_explicit_authorization() -> anyhow::Result<()> | |
| { | |
| let codex_home = tempdir()?; | |
| let environment_manager = Arc::new(environment_manager_without_environments()); | |
| environment_manager.upsert_environment( | |
| "customer-executor".to_string(), | |
| "ws://127.0.0.1:1".to_string(), | |
| /*connect_timeout*/ None, | |
| )?; | |
| let runtime_context = | |
| McpRuntimeContext::new(Arc::clone(&environment_manager), PathBuf::from("/tmp")); | |
| let bootstrap_auth = CodexAuth::create_dummy_chatgpt_auth_for_testing(); | |
| let mut actor_headers = Default::default(); | |
| codex_model_provider::auth_provider_from_auth(&bootstrap_auth) | |
| .add_auth_headers(&mut actor_headers); | |
| actor_headers.insert( | |
| "x-openai-actor-authorization", | |
| "hosted-actor-secret" | |
| .parse() | |
| .expect("valid actor authorization header"), | |
| ); | |
| let hosted_auth = CodexAuth::Headers(AuthHeaders::new(actor_headers)); | |
| let runtime_config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| let cases = [ | |
| ("missing", None, None, None, false), | |
| ("empty", Some(("Authorization", "")), None, None, false), | |
| ( | |
| "whitespace", | |
| Some(("Authorization", " \t ")), | |
| None, | |
| None, | |
| false, | |
| ), | |
| ( | |
| "invalid newline", | |
| Some(("Authorization", "Bearer executor\r\nsecret")), | |
| None, | |
| None, | |
| false, | |
| ), | |
| ( | |
| "invalid NUL", | |
| Some(("Authorization", "Bearer executor\0secret")), | |
| None, | |
| None, | |
| false, | |
| ), | |
| ( | |
| "invalid DEL", | |
| Some(("Authorization", "Bearer executor\u{007f}secret")), | |
| None, | |
| None, | |
| false, | |
| ), | |
| ( | |
| "environment header", | |
| Some(("Authorization", "Bearer executor-secret")), | |
| None, | |
| Some(("aUtHoRiZaTiOn", "CODEX_TEST_HOSTED_SECRET")), | |
| false, | |
| ), | |
| ( | |
| "environment bearer", | |
| Some(("Authorization", "Bearer executor-secret")), | |
| Some("CODEX_TEST_HOSTED_SECRET"), | |
| None, | |
| false, | |
| ), | |
| ( | |
| "mixed-case static header", | |
| Some(("aUtHoRiZaTiOn", "Bearer executor-secret")), | |
| None, | |
| None, | |
| true, | |
| ), | |
| ]; | |
| for (case, static_header, bearer_env_var, env_header, allows_executor_auth) in cases { | |
| let mut server_json = serde_json::json!({ | |
| "url": "https://chatgpt.com/backend-api/ps/mcp", | |
| "auth": "chatgpt", | |
| "environment_id": "customer-executor", | |
| }); | |
| if let Some((name, value)) = static_header { | |
| server_json["http_headers"] = serde_json::json!({ name: value }); | |
| } | |
| if let Some(name) = bearer_env_var { | |
| server_json["bearer_token_env_var"] = serde_json::json!(name); | |
| } | |
| if let Some((name, value)) = env_header { | |
| server_json["env_http_headers"] = serde_json::json!({ name: value }); | |
| } | |
| let server_config = serde_json::from_value::<McpServerConfig>(server_json)?; | |
| let mcp_servers = crate::effective_mcp_servers_from_configured( | |
| HashMap::from([("fake-first-party".to_string(), server_config)]), | |
| &runtime_config, | |
| Some(&hosted_auth), | |
| ); | |
| assert!(matches!( | |
| mcp_servers["fake-first-party"].config().auth, | |
| McpServerAuth::ChatGpt | |
| )); | |
| let remote_server = &mcp_servers["fake-first-party"]; | |
| assert!( | |
| chatgpt_auth_provider_for_server( | |
| remote_server, | |
| Some(codex_model_provider::auth_provider_from_auth(&hosted_auth)), | |
| ) | |
| .is_none(), | |
| "{case}: executor-owned servers must never receive hosted actor credentials" | |
| ); | |
| let resolved_environment = | |
| runtime_context.resolve_server_environment("fake-first-party", remote_server.config()); | |
| let connection_identity = |keyring_backend_kind| { | |
| McpServerConnectionIdentity::new( | |
| "fake-first-party", | |
| remote_server, | |
| /*host_plugin_root*/ None, | |
| OAuthCredentialsStoreMode::File, | |
| keyring_backend_kind, | |
| McpOAuthRefreshMode::Legacy, | |
| &resolved_environment, | |
| &runtime_context, | |
| /*runtime_auth_provider*/ None, | |
| Some(&hosted_auth), | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ) | |
| }; | |
| let direct_keyring_identity = connection_identity(AuthKeyringBackendKind::Direct); | |
| let secrets_keyring_identity = connection_identity(AuthKeyringBackendKind::Secrets); | |
| assert!( | |
| direct_keyring_identity.has_same_connection_config(&secrets_keyring_identity), | |
| "{case}: executor-owned servers must not inspect orchestrator OAuth stores" | |
| ); | |
| assert!( | |
| direct_keyring_identity | |
| .oauth_credentials() | |
| .expect("executor-owned ChatGPT authentication must skip OAuth lookup") | |
| .is_none(), | |
| "{case}: executor-owned servers must not retain hosted OAuth credentials" | |
| ); | |
| let auth_statuses = crate::compute_auth_statuses( | |
| mcp_servers.iter(), | |
| OAuthCredentialsStoreMode::default(), | |
| AuthKeyringBackendKind::default(), | |
| Some(&hosted_auth), | |
| &runtime_context, | |
| ) | |
| .await; | |
| let expected_auth_state = if allows_executor_auth { | |
| McpAuthState::BearerToken | |
| } else { | |
| McpAuthState::Unsupported | |
| }; | |
| assert_eq!( | |
| auth_statuses["fake-first-party"].auth_state, expected_auth_state, | |
| "{case}: auth status must only accept safe executor-owned authorization" | |
| ); | |
| let manager = McpConnectionSet::new( | |
| /*previous*/ None, | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(runtime_config.clone()), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers, | |
| submit_id: "security-test".to_string(), | |
| tx_event: None, | |
| startup_cancellation_token: CancellationToken::new(), | |
| runtime_context: runtime_context.clone(), | |
| codex_apps_tools_cache: ConnectorRuntimeManager::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: Some(hosted_auth.clone()), | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| let error = match manager.test_client("fake-first-party").client().await { | |
| Ok(_) => panic!("{case}: the unreachable fake executor must not connect"), | |
| Err(error) => error, | |
| }; | |
| let StartupOutcomeError::Failed { error, .. } = error else { | |
| panic!("{case}: executor-owned authentication must fail rather than be cancelled"); | |
| }; | |
| if allows_executor_auth { | |
| assert!( | |
| error.contains("127.0.0.1:1"), | |
| "{case}: safe explicit credentials should reach the executor: {error}" | |
| ); | |
| } else { | |
| assert_eq!( | |
| error, | |
| "executor-owned MCP server `fake-first-party` cannot use hosted ChatGPT authentication; configure executor-owned credentials instead", | |
| "{case}: unsafe credentials must fail before contacting the executor" | |
| ); | |
| } | |
| } | |
| Ok(()) | |
| } | |
| async fn no_local_runtime_fails_local_stdio_but_keeps_local_http_server() { | |
| let codex_home = tempdir().expect("tempdir"); | |
| let mcp_servers = HashMap::from([ | |
| ( | |
| "stdio".to_string(), | |
| EffectiveMcpServer::configured(McpServerConfig { | |
| auth: Default::default(), | |
| transport: McpServerTransportConfig::Stdio { | |
| command: "echo".to_string(), | |
| args: Vec::new(), | |
| env: None, | |
| env_vars: Vec::new(), | |
| cwd: None, | |
| }, | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| enabled: true, | |
| required: false, | |
| supports_parallel_tool_calls: false, | |
| omit_tools_from: None, | |
| disabled_reason: None, | |
| startup_timeout_sec: None, | |
| tool_timeout_sec: None, | |
| default_tools_approval_mode: None, | |
| enabled_tools: None, | |
| disabled_tools: None, | |
| scopes: None, | |
| oauth: None, | |
| oauth_resource: None, | |
| tools: HashMap::new(), | |
| }), | |
| ), | |
| ( | |
| "http".to_string(), | |
| EffectiveMcpServer::configured(McpServerConfig { | |
| auth: Default::default(), | |
| transport: McpServerTransportConfig::StreamableHttp { | |
| url: "http://127.0.0.1:1".to_string(), | |
| bearer_token_env_var: None, | |
| http_headers: None, | |
| env_http_headers: None, | |
| http_headers_helper: None, | |
| }, | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| enabled: true, | |
| required: false, | |
| supports_parallel_tool_calls: false, | |
| omit_tools_from: None, | |
| disabled_reason: None, | |
| startup_timeout_sec: None, | |
| tool_timeout_sec: None, | |
| default_tools_approval_mode: None, | |
| enabled_tools: None, | |
| disabled_tools: None, | |
| scopes: None, | |
| oauth: None, | |
| oauth_resource: None, | |
| tools: HashMap::new(), | |
| }), | |
| ), | |
| ]); | |
| let cancel_token = CancellationToken::new(); | |
| let manager = McpConnectionSet::new( | |
| /*previous*/ None, | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(crate::mcp::tests::test_mcp_config( | |
| codex_home.path().to_path_buf(), | |
| )), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers, | |
| submit_id: String::new(), | |
| tx_event: None, | |
| startup_cancellation_token: cancel_token.clone(), | |
| runtime_context: McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| PathBuf::from("/tmp"), | |
| ), | |
| codex_apps_tools_cache: ConnectorRuntimeManager::<ToolInfo>::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| assert!(manager.contains_server("stdio")); | |
| assert!(manager.contains_server("http")); | |
| assert!( | |
| !manager | |
| .wait_for_server_ready("stdio", Duration::from_millis(10)) | |
| .await | |
| ); | |
| let error = match manager.test_client("stdio").client().await { | |
| Ok(_) => panic!("local stdio MCP startup should fail"), | |
| Err(error) => error, | |
| }; | |
| let StartupOutcomeError::Failed { error, .. } = error else { | |
| panic!("local stdio MCP startup should fail rather than be cancelled"); | |
| }; | |
| assert_eq!( | |
| error, | |
| "local stdio MCP server `stdio` requires a local environment" | |
| ); | |
| cancel_token.cancel(); | |
| } | |
| fn elicitation_capability_uses_2025_06_18_shape_for_form_only_support() { | |
| let capability = Some(ElicitationCapability::default()); | |
| assert_eq!( | |
| serde_json::to_value(capability).expect("serialize elicitation capability"), | |
| serde_json::json!({}) | |
| ); | |
| } | |
| fn elicitation_capability_advertises_url_support_when_enabled() { | |
| let capability = Some( | |
| ElicitationCapability::new() | |
| .with_form(rmcp::model::FormElicitationCapability::new()) | |
| .with_url(rmcp::model::UrlElicitationCapability::new()), | |
| ); | |
| assert_eq!( | |
| serde_json::to_value(capability).expect("serialize elicitation capability"), | |
| serde_json::json!({ | |
| "form": {}, | |
| "url": {}, | |
| }) | |
| ); | |
| } | |
| fn mcp_init_error_display_prompts_for_github_pat() { | |
| let server_name = "github"; | |
| let config = McpServerConfig { | |
| auth: Default::default(), | |
| transport: McpServerTransportConfig::StreamableHttp { | |
| url: "https://api.githubcopilot.com/mcp/".to_string(), | |
| bearer_token_env_var: None, | |
| http_headers: None, | |
| env_http_headers: None, | |
| http_headers_helper: None, | |
| }, | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| enabled: true, | |
| required: false, | |
| supports_parallel_tool_calls: false, | |
| omit_tools_from: None, | |
| disabled_reason: None, | |
| startup_timeout_sec: None, | |
| tool_timeout_sec: None, | |
| default_tools_approval_mode: None, | |
| enabled_tools: None, | |
| disabled_tools: None, | |
| scopes: None, | |
| oauth: None, | |
| oauth_resource: None, | |
| tools: HashMap::new(), | |
| }; | |
| let err: StartupOutcomeError = anyhow::anyhow!("OAuth is unsupported").into(); | |
| let display = mcp_init_error_display(server_name, Some(&config), &err, /*reason*/ None); | |
| let expected = format!( | |
| "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" | |
| ); | |
| assert_eq!(expected, display); | |
| } | |
| fn mcp_init_error_display_prompts_for_login_when_auth_required() { | |
| let server_name = "example"; | |
| let expected = format!( | |
| "The {server_name} MCP server is not logged in. Run `codex mcp login {server_name}`." | |
| ); | |
| let executor_config: McpServerConfig = serde_json::from_value(serde_json::json!({ | |
| "url": "https://example.com/mcp", | |
| "environment_id": "executor-1", | |
| })) | |
| .expect("executor MCP configuration should deserialize"); | |
| for error in [ | |
| anyhow::anyhow!("Auth required for server").into(), | |
| StartupOutcomeError::Failed { | |
| error: "OAuth refresh token was rejected: invalid_grant".to_string(), | |
| is_authentication_required: true, | |
| }, | |
| ] { | |
| let display = mcp_init_error_display( | |
| server_name, | |
| /*config*/ None, | |
| &error, | |
| /*reason*/ None, | |
| ); | |
| assert_eq!(expected, display); | |
| let executor_display = mcp_init_error_display( | |
| server_name, | |
| Some(&executor_config), | |
| &error, | |
| /*reason*/ None, | |
| ); | |
| assert_eq!( | |
| format!( | |
| "The {server_name} MCP server is not logged in. Use your client's MCP OAuth sign-in flow." | |
| ), | |
| executor_display | |
| ); | |
| } | |
| } | |
| fn mcp_init_error_display_identifies_oauth_reauthentication() { | |
| let server_name = "example"; | |
| let error = StartupOutcomeError::Failed { | |
| error: "authorization required: Bearer error=\"invalid_token\"".to_string(), | |
| is_authentication_required: true, | |
| }; | |
| let executor_config: McpServerConfig = serde_json::from_value(serde_json::json!({ | |
| "url": "https://example.com/mcp", | |
| "environment_id": "executor-1", | |
| })) | |
| .expect("executor MCP configuration should deserialize"); | |
| for (config, recovery_hint) in [ | |
| (None, "Run `codex mcp login example`."), | |
| ( | |
| Some(&executor_config), | |
| "Use your client's MCP OAuth sign-in flow.", | |
| ), | |
| ] { | |
| assert_eq!( | |
| mcp_init_error_display( | |
| server_name, | |
| config, | |
| &error, | |
| Some(McpStartupFailureReason::ReauthenticationRequired), | |
| ), | |
| format!( | |
| "The {server_name} MCP server requires OAuth reauthentication. {recovery_hint}" | |
| ), | |
| ); | |
| } | |
| } | |
| fn mcp_startup_failure_reason_requires_existing_oauth_and_auth_failure() { | |
| for (auth_state, is_authentication_required, expected) in [ | |
| ( | |
| Some(McpAuthState::LoggedOut( | |
| McpLoginRequirement::Reauthentication, | |
| )), | |
| true, | |
| Some(McpStartupFailureReason::ReauthenticationRequired), | |
| ), | |
| ( | |
| Some(McpAuthState::LoggedOut( | |
| McpLoginRequirement::Reauthentication, | |
| )), | |
| false, | |
| None, | |
| ), | |
| ( | |
| Some(McpAuthState::LoggedOut(McpLoginRequirement::Login)), | |
| true, | |
| None, | |
| ), | |
| (Some(McpAuthState::Unsupported), true, None), | |
| (Some(McpAuthState::BearerToken), true, None), | |
| ( | |
| Some(McpAuthState::OAuth), | |
| true, | |
| Some(McpStartupFailureReason::ReauthenticationRequired), | |
| ), | |
| (Some(McpAuthState::OAuth), false, None), | |
| (None, true, None), | |
| ] { | |
| let error = StartupOutcomeError::Failed { | |
| error: "startup failed".to_string(), | |
| is_authentication_required, | |
| }; | |
| assert_eq!( | |
| mcp_startup_failure_reason(auth_state, &error), | |
| expected, | |
| "auth_state={auth_state:?}, is_authentication_required={is_authentication_required}" | |
| ); | |
| } | |
| } | |
| fn mcp_init_error_display_reports_generic_errors() { | |
| let server_name = "custom"; | |
| let config = McpServerConfig { | |
| auth: Default::default(), | |
| transport: McpServerTransportConfig::StreamableHttp { | |
| url: "https://example.com".to_string(), | |
| bearer_token_env_var: Some("TOKEN".to_string()), | |
| http_headers: None, | |
| env_http_headers: None, | |
| http_headers_helper: None, | |
| }, | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| enabled: true, | |
| required: false, | |
| supports_parallel_tool_calls: false, | |
| omit_tools_from: None, | |
| disabled_reason: None, | |
| startup_timeout_sec: None, | |
| tool_timeout_sec: None, | |
| default_tools_approval_mode: None, | |
| enabled_tools: None, | |
| disabled_tools: None, | |
| scopes: None, | |
| oauth: None, | |
| oauth_resource: None, | |
| tools: HashMap::new(), | |
| }; | |
| let err: StartupOutcomeError = anyhow::anyhow!("boom").into(); | |
| let display = mcp_init_error_display(server_name, Some(&config), &err, /*reason*/ None); | |
| let expected = format!("MCP client for `{server_name}` failed to start: {err:#}"); | |
| assert_eq!(expected, display); | |
| } | |
| fn mcp_init_error_display_quotes_server_names() { | |
| let github_config: McpServerConfig = serde_json::from_value(serde_json::json!({ | |
| "url": "https://api.githubcopilot.com/mcp/", | |
| })) | |
| .expect("GitHub MCP configuration should deserialize"); | |
| let error: StartupOutcomeError = anyhow::anyhow!("request timed out").into(); | |
| let mut displays = Vec::new(); | |
| for server_name in ["npm:@scope/package.name", "server.name"] { | |
| for config in [None, Some(&github_config)] { | |
| displays.push(mcp_init_error_display( | |
| server_name, | |
| config, | |
| &error, | |
| /*reason*/ None, | |
| )); | |
| } | |
| } | |
| insta::assert_snapshot!(displays.join("\n\n")); | |
| } | |
| fn mcp_init_error_display_includes_startup_timeout_hint() { | |
| let server_name = "slow"; | |
| for error in [ | |
| "request timed out", | |
| "MCP client startup timed out after 30s", | |
| ] { | |
| let err: StartupOutcomeError = anyhow::anyhow!(error).into(); | |
| let display = mcp_init_error_display( | |
| server_name, | |
| /*config*/ None, | |
| &err, | |
| /*reason*/ None, | |
| ); | |
| assert_eq!( | |
| "MCP client for `slow` timed out after 30 seconds. Add or adjust `startup_timeout_sec` in your config.toml:\n[mcp_servers.slow]\nstartup_timeout_sec = XX", | |
| display | |
| ); | |
| } | |
| } | |
| fn reusable_server_config(url: &str) -> McpServerConfig { | |
| McpServerConfig { | |
| auth: Default::default(), | |
| transport: McpServerTransportConfig::StreamableHttp { | |
| url: url.to_string(), | |
| bearer_token_env_var: Some("CODEX_MCP_REUSE_TEST_TOKEN".to_string()), | |
| http_headers: None, | |
| env_http_headers: None, | |
| http_headers_helper: None, | |
| }, | |
| environment_id: codex_config::DEFAULT_MCP_SERVER_ENVIRONMENT_ID.to_string(), | |
| enabled: true, | |
| required: false, | |
| supports_parallel_tool_calls: false, | |
| omit_tools_from: None, | |
| disabled_reason: None, | |
| startup_timeout_sec: None, | |
| tool_timeout_sec: None, | |
| default_tools_approval_mode: None, | |
| enabled_tools: None, | |
| disabled_tools: None, | |
| scopes: None, | |
| oauth: None, | |
| oauth_resource: None, | |
| tools: HashMap::new(), | |
| } | |
| } | |
| fn reusable_server_runtime_context() -> McpRuntimeContext { | |
| McpRuntimeContext::new( | |
| Arc::new(environment_manager_without_environments()), | |
| PathBuf::from("/tmp"), | |
| ) | |
| } | |
| fn reusable_server_identity( | |
| server_name: &str, | |
| config: &McpServerConfig, | |
| runtime_context: &McpRuntimeContext, | |
| ) -> McpServerConnectionIdentity { | |
| let server = EffectiveMcpServer::configured(config.clone()); | |
| let resolved_environment = runtime_context.resolve_server_environment(server_name, config); | |
| McpServerConnectionIdentity::new( | |
| server_name, | |
| &server, | |
| /*host_plugin_root*/ None, | |
| OAuthCredentialsStoreMode::default(), | |
| AuthKeyringBackendKind::default(), | |
| McpOAuthRefreshMode::Legacy, | |
| &resolved_environment, | |
| runtime_context, | |
| /*runtime_auth_provider*/ None, | |
| /*auth*/ None, | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ) | |
| } | |
| async fn manager_with_reusable_ready_server( | |
| config: &McpServerConfig, | |
| runtime_context: &McpRuntimeContext, | |
| tools: Vec<ToolInfo>, | |
| ) -> McpConnectionSet { | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut manager = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let server = EffectiveMcpServer::configured(config.clone()); | |
| manager.servers.insert( | |
| "docs".to_string(), | |
| McpServerView { | |
| protocol_mode: crate::McpProtocolMode::Legacy, | |
| connection: Arc::new(McpServerConnection { | |
| identity: Some(reusable_server_identity("docs", config, runtime_context)), | |
| client: create_ready_async_managed_client(tools).await, | |
| startup_timeout: config | |
| .startup_timeout_sec | |
| .unwrap_or(DEFAULT_STARTUP_TIMEOUT), | |
| startup_trigger: None, | |
| _diagnostics_guard: LIVE_CONNECTIONS.track(), | |
| }), | |
| metadata: McpServerMetadata::from(&server), | |
| tool_filter: ToolFilter::from_config(config), | |
| tool_timeout: Some(config.tool_timeout_sec.unwrap_or(DEFAULT_TOOL_TIMEOUT)), | |
| catalog_item_limit: crate::pagination::MAX_MCP_CATALOG_ITEMS, | |
| }, | |
| ); | |
| manager | |
| } | |
| async fn reconcile_reusable_server( | |
| previous: &McpConnectionSet, | |
| config: McpServerConfig, | |
| runtime_context: McpRuntimeContext, | |
| ) -> McpConnectionSet { | |
| let codex_home = tempdir().expect("tempdir"); | |
| reconcile_reusable_server_with_mcp_config( | |
| previous, | |
| "docs", | |
| config, | |
| runtime_context, | |
| crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()), | |
| ) | |
| .await | |
| } | |
| async fn reconcile_reusable_server_with_mcp_config( | |
| previous: &McpConnectionSet, | |
| server_name: &str, | |
| config: McpServerConfig, | |
| runtime_context: McpRuntimeContext, | |
| mcp_config: crate::McpConfig, | |
| ) -> McpConnectionSet { | |
| let (tx_event, _rx_event) = async_channel::unbounded(); | |
| McpConnectionSet::new( | |
| Some(previous), | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(mcp_config), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers: HashMap::from([( | |
| server_name.to_string(), | |
| EffectiveMcpServer::configured(config), | |
| )]), | |
| submit_id: "refresh".to_string(), | |
| tx_event: Some(tx_event), | |
| startup_cancellation_token: CancellationToken::new(), | |
| runtime_context, | |
| codex_apps_tools_cache: ConnectorRuntimeManager::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await | |
| } | |
| async fn apps_catalog_broadcast_preserves_running_calls_and_rejects_stale_calls() | |
| -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| /*account_id*/ None, | |
| /*chatgpt_user_id*/ None, | |
| ) | |
| .with_live_scope("apps".to_string()); | |
| let original = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "original")]; | |
| let updated = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "updated")]; | |
| store_current_tools(&context, original.clone()); | |
| let catalog = ClientToolCatalog::new(original.clone(), context.subscribe()); | |
| let snapshot = catalog.read(Arc::new).await; | |
| store_current_tools(&context, original.clone()); | |
| assert_eq!( | |
| catalog.read(|catalog| catalog.revision).await, | |
| 0, | |
| "an unchanged startup broadcast must not invalidate prepared calls" | |
| ); | |
| let (release, released) = tokio::sync::oneshot::channel::<()>(); | |
| let running = catalog.run_with_snapshot(&snapshot, || async { released.await.unwrap() }); | |
| tokio::pin!(running); | |
| assert!(futures::poll!(&mut running).is_pending()); | |
| store_current_tools(&context, updated.clone()); | |
| let stale = catalog.run_with_snapshot(&snapshot, || async { panic!("stale preparation") }); | |
| tokio::pin!(stale); | |
| assert!( | |
| futures::poll!(&mut stale).is_pending(), | |
| "publication does not wait, but adoption must wait for the running call" | |
| ); | |
| release.send(()).unwrap(); | |
| assert_eq!(running.await, Some(())); | |
| assert_eq!(stale.await, None::<()>); | |
| assert_eq!( | |
| catalog | |
| .read(|catalog| (catalog.revision, catalog.tools.to_vec())) | |
| .await, | |
| (1, updated.clone()) | |
| ); | |
| // A client whose startup finishes late adopts the already-published result. | |
| let late = ClientToolCatalog::new(original, context.subscribe()); | |
| assert_eq!(late.read(|catalog| catalog.tools.to_vec()).await, updated); | |
| Ok(()) | |
| } | |
| async fn apps_catalog_broadcast_restores_equivalent_calls() -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| /*account_id*/ None, | |
| /*chatgpt_user_id*/ None, | |
| ) | |
| .with_live_scope("apps".to_string()); | |
| let original = vec![ | |
| create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "first"), | |
| create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "second"), | |
| ]; | |
| store_current_tools(&context, original.clone()); | |
| let catalog = ClientToolCatalog::new(original.clone(), context.subscribe()); | |
| let snapshot = catalog.read(Arc::new).await; | |
| let mut changed = original.clone(); | |
| changed[0].tool.description = Some("Changed definition".into()); | |
| store_current_tools(&context, changed); | |
| assert_eq!( | |
| catalog | |
| .run_with_snapshot(&snapshot, || async { panic!("changed preparation") }) | |
| .await, | |
| None::<()> | |
| ); | |
| let mut restored = original; | |
| restored.reverse(); | |
| store_current_tools(&context, restored); | |
| assert_eq!( | |
| catalog | |
| .run_with_snapshot(&snapshot, || async { "prepared" }) | |
| .await, | |
| Some("prepared") | |
| ); | |
| catalog | |
| .refresh( | |
| || async { Ok((snapshot.tools.to_vec(), ())) }, | |
| |tools, ()| store_current_tools(&context, tools.to_vec()), | |
| ) | |
| .await?; | |
| assert_eq!( | |
| catalog | |
| .run_with_snapshot(&snapshot, || async { panic!("refreshed preparation") }) | |
| .await, | |
| None::<()> | |
| ); | |
| Ok(()) | |
| } | |
| fn idle_apps_clients_do_not_retain_replaced_tools() { | |
| let context = ConnectorRuntimeManager::<ToolInfo>::new_without_cache() | |
| .context( | |
| PathBuf::from("unused"), | |
| ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| ) | |
| .with_live_scope("apps".into()); | |
| let tool = create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "original"); | |
| let schema = Arc::downgrade(&tool.tool.input_schema); | |
| store_current_tools(&context, vec![tool]); | |
| let clients = [ | |
| ClientToolCatalog::new(Vec::new(), context.subscribe()), | |
| ClientToolCatalog::new(Vec::new(), context.subscribe()), | |
| ]; | |
| store_current_tools( | |
| &context, | |
| vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "updated")], | |
| ); | |
| assert!( | |
| schema.upgrade().is_none(), | |
| "idle clients must not own the old schema" | |
| ); | |
| drop(clients); | |
| } | |
| async fn apps_catalog_broadcast_survives_an_older_local_refresh() -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| /*account_id*/ None, | |
| /*chatgpt_user_id*/ None, | |
| ) | |
| .with_live_scope("apps".to_string()); | |
| store_current_tools(&context, Vec::new()); | |
| let catalog = ClientToolCatalog::new(Vec::new(), context.subscribe()); | |
| let older_ticket = context.begin_fetch(ConnectorRuntimeFetchSource::HardRefresh); | |
| let newer = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "newer")]; | |
| store_current_tools(&context, newer.clone()); | |
| assert_eq!(catalog.read(|catalog| catalog.tools.to_vec()).await, newer); | |
| catalog | |
| .refresh( | |
| || async { Ok((Vec::new(), older_ticket)) }, | |
| |tools, ticket| { | |
| context.publish_if_newest_accepted( | |
| ticket, | |
| &create_test_server_info("Apps"), | |
| tools.to_vec(), | |
| ) | |
| }, | |
| ) | |
| .await?; | |
| assert_eq!(catalog.read(|catalog| catalog.tools.to_vec()).await, newer); | |
| Ok(()) | |
| } | |
| async fn refreshed_catalog_follows_reused_client_without_mutating_old_bindings() | |
| -> anyhow::Result<()> { | |
| let codex_home = tempdir()?; | |
| let cache_context = create_codex_apps_tools_cache_context( | |
| codex_home.path().to_path_buf(), | |
| /*account_id*/ None, | |
| /*chatgpt_user_id*/ None, | |
| ); | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let mut previous = create_test_manager_with_ready_apps_client( | |
| cache_context, | |
| "refreshed", | |
| /*list_started*/ None, | |
| /*release_list*/ None, | |
| ) | |
| .await?; | |
| let manager = Arc::get_mut(&mut previous).expect("unshared manager"); | |
| manager.insert_test_client( | |
| "docs", | |
| create_ready_async_managed_client(vec![create_test_tool("docs", "unrelated")]).await, | |
| ); | |
| let connection = Arc::get_mut( | |
| &mut manager | |
| .servers | |
| .get_mut(CODEX_APPS_MCP_SERVER_NAME) | |
| .expect("Apps server") | |
| .connection, | |
| ) | |
| .expect("unshared Apps connection"); | |
| connection.identity = Some(reusable_server_identity( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| &config, | |
| &runtime_context, | |
| )); | |
| let client = connection.client().await?; | |
| // Seed an earlier catalog before capturing the binding under test. | |
| let startup_tools = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "startup")]; | |
| client | |
| .tool_catalog | |
| .refresh(|| async { Ok((startup_tools, ())) }, |_, ()| {}) | |
| .await?; | |
| let old_binding = capture_binding(&previous).await; | |
| let old_tools = serde_json::to_value(old_binding.tools())?; | |
| let old_call = old_binding | |
| .prepare_call(CODEX_APPS_MCP_SERVER_NAME, "startup") | |
| .expect("startup call"); | |
| let unrelated_call = old_binding | |
| .prepare_call("docs", "unrelated") | |
| .expect("unrelated call"); | |
| previous.refresh_codex_apps_tools_for_discovery().await?; | |
| let republished = Arc::new( | |
| reconcile_reusable_server_with_mcp_config( | |
| &previous, | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| config, | |
| runtime_context, | |
| crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()), | |
| ) | |
| .await, | |
| ); | |
| assert!(previous.shares_test_connection_with(&republished, CODEX_APPS_MCP_SERVER_NAME)); | |
| assert_eq!( | |
| model_tool_names(capture_binding(&republished).await.tools()), | |
| HashSet::from([ToolName::namespaced("mcp__codex_apps", "refreshed")]), | |
| ); | |
| assert_eq!(serde_json::to_value(old_binding.tools())?, old_tools); | |
| let error = old_call | |
| .call_with_preparation(/*requested_timeout*/ None, || async { | |
| panic!("stale call preparation must not run"); | |
| }) | |
| .await | |
| .expect_err("a call from the old catalog must be rejected"); | |
| assert!(error.to_string().contains("catalog changed")); | |
| let error = unrelated_call | |
| .call_with_preparation(/*requested_timeout*/ None, || async { | |
| Err(anyhow!("unrelated preparation reached")) | |
| }) | |
| .await | |
| .expect_err("stop before the unrelated tool executes"); | |
| assert!( | |
| error.to_string().contains("unrelated preparation reached"), | |
| "Apps refresh must not invalidate another client's calls" | |
| ); | |
| Ok(()) | |
| } | |
| async fn reconciliation_reuses_connection_without_relisting_regular_tools() -> anyhow::Result<()> { | |
| let tools = Arc::new(tokio::sync::RwLock::new(vec![Tool::new( | |
| "old_search", | |
| "old search", | |
| Arc::new(JsonObject::default()), | |
| )])); | |
| let block_tool_listing = Arc::new(AtomicBool::new(false)); | |
| let client = Arc::new( | |
| RmcpClient::new_in_process_client(Arc::new(MutableToolsTransportFactory { | |
| server: MutableToolsServer { | |
| tools: Arc::clone(&tools), | |
| block_tool_listing: Arc::clone(&block_tool_listing), | |
| }, | |
| })) | |
| .await?, | |
| ); | |
| let initialize = client | |
| .initialize( | |
| InitializeRequestParams::new( | |
| ClientCapabilities::default(), | |
| Implementation::new("codex-test", "0.0.0-test"), | |
| ) | |
| .with_protocol_version(ProtocolVersion::V_2025_06_18), | |
| /*timeout*/ None, | |
| Box::new(|_, _| { | |
| async { | |
| Ok(ElicitationResponse { | |
| action: ElicitationAction::Decline, | |
| content: None, | |
| meta: None, | |
| }) | |
| } | |
| .boxed() | |
| }), | |
| ) | |
| .await?; | |
| let initial_tools = list_tools_for_client_uncached( | |
| "docs", | |
| /*is_codex_apps_mcp_server*/ false, | |
| /*codex_apps_refresh_trigger*/ "test", | |
| &client, | |
| /*timeout*/ None, | |
| crate::pagination::MAX_MCP_CATALOG_ITEMS, | |
| initialize.instructions.as_deref(), | |
| ) | |
| .await?; | |
| let managed_client = ManagedClient { | |
| _auth_change_notifications: None, | |
| client, | |
| server_info: create_test_server_info("Mutable tools"), | |
| tool_catalog: Arc::new(ClientToolCatalog::new(initial_tools, /*updates*/ None)), | |
| tool_timeout: None, | |
| server_instructions: initialize.instructions, | |
| server_supports_sandbox_state_meta_capability: false, | |
| codex_apps_tools_cache_context: None, | |
| }; | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let approval_policy = Constrained::allow_any(AskForApproval::OnRequest); | |
| let permission_profile = Constrained::allow_any(PermissionProfile::default()); | |
| let mut previous = McpConnectionSet::new_uninitialized( | |
| &approval_policy, | |
| &permission_profile, | |
| /*prefix_mcp_tool_names*/ true, | |
| ); | |
| let server = EffectiveMcpServer::configured(config.clone()); | |
| previous.servers.insert( | |
| "docs".to_string(), | |
| McpServerView { | |
| protocol_mode: crate::McpProtocolMode::Legacy, | |
| connection: Arc::new(McpServerConnection { | |
| identity: Some(reusable_server_identity("docs", &config, &runtime_context)), | |
| client: AsyncManagedClient { | |
| client: futures::future::ready(Ok(managed_client)).boxed().shared(), | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| startup_timeout: config | |
| .startup_timeout_sec | |
| .unwrap_or(DEFAULT_STARTUP_TIMEOUT), | |
| startup_trigger: None, | |
| _diagnostics_guard: LIVE_CONNECTIONS.track(), | |
| }), | |
| metadata: McpServerMetadata::from(&server), | |
| tool_filter: ToolFilter::from_config(&config), | |
| tool_timeout: Some(config.tool_timeout_sec.unwrap_or(DEFAULT_TOOL_TIMEOUT)), | |
| catalog_item_limit: crate::pagination::MAX_MCP_CATALOG_ITEMS, | |
| }, | |
| ); | |
| let previous = Arc::new(previous); | |
| let old_step = capture_binding(&previous).await; | |
| *tools.write().await = vec![Tool::new( | |
| "new_search", | |
| "new search", | |
| Arc::new(JsonObject::default()), | |
| )]; | |
| block_tool_listing.store(true, Ordering::Release); | |
| let reconciled = Arc::new( | |
| tokio::time::timeout( | |
| Duration::from_secs(1), | |
| reconcile_reusable_server(&previous, config, runtime_context), | |
| ) | |
| .await | |
| .expect("connection reuse must not wait for a tool-list request"), | |
| ); | |
| let new_step = capture_binding(&reconciled).await; | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| assert_eq!( | |
| old_step | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.tool.name.to_string()) | |
| .collect::<Vec<_>>(), | |
| vec!["old_search".to_string()] | |
| ); | |
| assert_eq!( | |
| new_step | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.tool.name.to_string()) | |
| .collect::<Vec<_>>(), | |
| vec!["old_search".to_string()] | |
| ); | |
| Ok(()) | |
| } | |
| async fn reconciliation_reuses_an_unchanged_ready_server() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context.clone()).await; | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| assert_eq!( | |
| model_tool_names(&reconciled.list_all_tools().await), | |
| HashSet::from([ToolName::namespaced("mcp__docs", "search")]) | |
| ); | |
| } | |
| async fn reconciliation_reuses_an_unchanged_pending_server_without_waiting() -> anyhow::Result<()> { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| let tools = vec![ | |
| create_test_tool("docs", "search"), | |
| create_test_tool("docs", "write"), | |
| ]; | |
| let mut previous = | |
| manager_with_reusable_ready_server(&config, &runtime_context, tools.clone()).await; | |
| let managed_client = create_test_managed_client(tools).await; | |
| let (pending_client, startup_started, release_startup) = | |
| create_gated_async_managed_client(managed_client); | |
| let startup = tokio::spawn({ | |
| let pending_client = pending_client.clone(); | |
| async move { pending_client.client().await } | |
| }); | |
| startup_started.await?; | |
| let connection = Arc::get_mut( | |
| &mut previous | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist") | |
| .connection, | |
| ) | |
| .expect("test server should have one connection owner"); | |
| connection.client = pending_client; | |
| config.enabled_tools = Some(vec!["search".to_string()]); | |
| config.startup_timeout_sec = Some(DEFAULT_STARTUP_TIMEOUT); | |
| let reconciled = tokio::time::timeout( | |
| Duration::from_millis(100), | |
| reconcile_reusable_server(&previous, config, runtime_context), | |
| ) | |
| .await | |
| .expect("reconciliation must not wait for an unchanged pending MCP server"); | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| release_startup | |
| .send(()) | |
| .map_err(|()| anyhow!("pending startup should still be running"))?; | |
| startup.await??; | |
| assert_eq!( | |
| model_tool_names(&reconciled.list_all_tools().await), | |
| HashSet::from([ToolName::namespaced("mcp__docs", "search")]) | |
| ); | |
| Ok(()) | |
| } | |
| async fn reconciliation_cancels_a_reused_pending_server_when_disabled() -> anyhow::Result<()> { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| let tools = vec![create_test_tool("docs", "search")]; | |
| let mut previous = | |
| manager_with_reusable_ready_server(&config, &runtime_context, tools.clone()).await; | |
| let managed_client = create_test_managed_client(tools).await; | |
| let (pending_client, startup_started, release_startup) = | |
| create_gated_async_managed_client(managed_client); | |
| let cancellation = pending_client.cancel_token.clone(); | |
| let startup = tokio::spawn({ | |
| let pending_client = pending_client.clone(); | |
| async move { pending_client.client().await } | |
| }); | |
| startup_started.await?; | |
| let connection = Arc::get_mut( | |
| &mut previous | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist") | |
| .connection, | |
| ) | |
| .expect("test server should have one connection owner"); | |
| connection.client = pending_client; | |
| let reused = | |
| reconcile_reusable_server(&previous, config.clone(), runtime_context.clone()).await; | |
| assert!(previous.shares_test_connection_with(&reused, "docs")); | |
| config.enabled = false; | |
| let removed = reconcile_reusable_server(&reused, config, runtime_context).await; | |
| assert!(!removed.servers.contains_key("docs")); | |
| drop(previous); | |
| drop(reused); | |
| assert!( | |
| cancellation.is_cancelled(), | |
| "disabling a reused pending MCP server should cancel its obsolete startup" | |
| ); | |
| release_startup | |
| .send(()) | |
| .map_err(|()| anyhow!("pending startup should remain available for test cleanup"))?; | |
| startup.await??; | |
| Ok(()) | |
| } | |
| async fn reconciliation_retries_non_oauth_authentication_failures() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let mut previous = | |
| manager_with_reusable_ready_server(&config, &runtime_context, Vec::new()).await; | |
| let connection = | |
| Arc::get_mut(&mut previous.servers.get_mut("docs").expect("server").connection) | |
| .expect("test server has one connection owner"); | |
| connection.client.client = futures::future::ready(Err(StartupOutcomeError::Failed { | |
| error: "bearer token rejected".to_string(), | |
| is_authentication_required: true, | |
| })) | |
| .boxed() | |
| .shared(); | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context).await; | |
| assert!(!previous.shares_test_connection_with(&reconciled, "docs")); | |
| } | |
| fn connection_identity_uses_effective_authorization_headers() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let missing_env_var = format!("CODEX_TEST_UNSET_MCP_AUTHORIZATION_{}", std::process::id()); | |
| assert!(std::env::var_os(&missing_env_var).is_none()); | |
| for (static_header, environment_header, has_authorization) in [ | |
| (Some("Bearer configured-token"), None, true), | |
| (Some("invalid\nheader"), None, false), | |
| (None, Some("PATH"), true), | |
| (None, Some(missing_env_var.as_str()), false), | |
| (None, None, false), | |
| ] { | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| config.transport = McpServerTransportConfig::StreamableHttp { | |
| url: "http://127.0.0.1:1".to_string(), | |
| bearer_token_env_var: None, | |
| http_headers: static_header | |
| .map(|value| HashMap::from([("aUtHoRiZaTiOn".to_string(), value.to_string())])), | |
| env_http_headers: environment_header | |
| .map(|value| HashMap::from([("aUtHoRiZaTiOn".to_string(), value.to_string())])), | |
| http_headers_helper: None, | |
| }; | |
| let server = EffectiveMcpServer::configured(config); | |
| let identity = |keyring_backend_kind, oauth_refresh_mode| { | |
| McpServerConnectionIdentity::new( | |
| "docs", | |
| &server, | |
| /*host_plugin_root*/ None, | |
| OAuthCredentialsStoreMode::File, | |
| keyring_backend_kind, | |
| oauth_refresh_mode, | |
| &Ok(None), | |
| &runtime_context, | |
| /*runtime_auth_provider*/ None, | |
| /*auth*/ None, | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ) | |
| }; | |
| assert_eq!( | |
| identity(AuthKeyringBackendKind::Direct, McpOAuthRefreshMode::Legacy) | |
| .has_same_connection_config(&identity( | |
| AuthKeyringBackendKind::Secrets, | |
| McpOAuthRefreshMode::Legacy, | |
| )), | |
| has_authorization, | |
| ); | |
| assert_eq!( | |
| identity(AuthKeyringBackendKind::Direct, McpOAuthRefreshMode::Legacy) | |
| .has_same_connection_config(&identity( | |
| AuthKeyringBackendKind::Direct, | |
| McpOAuthRefreshMode::Coordinated, | |
| )), | |
| has_authorization, | |
| ); | |
| } | |
| } | |
| async fn reconciliation_reuses_legacy_stdio_server_with_existing_protocol_marker() { | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), | |
| PathBuf::from("/tmp"), | |
| ); | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| config.transport = McpServerTransportConfig::Stdio { | |
| command: "legacy-server".to_string(), | |
| args: Vec::new(), | |
| env: Some(HashMap::from([( | |
| "CODEX_MCP_PROTOCOL_VERSION".to_string(), | |
| "1999-01-01".to_string(), | |
| )])), | |
| env_vars: Vec::new(), | |
| cwd: None, | |
| }; | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context).await; | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| } | |
| async fn reconciliation_replaces_connection_when_auth_mode_changes() -> anyhow::Result<()> { | |
| let environment_manager = Arc::new(environment_manager_without_environments()); | |
| environment_manager.upsert_environment( | |
| "customer-executor".to_string(), | |
| "ws://127.0.0.1:1".to_string(), | |
| /*connect_timeout*/ None, | |
| )?; | |
| let runtime_context = McpRuntimeContext::new(environment_manager, PathBuf::from("/tmp")); | |
| let codex_home = tempdir()?; | |
| let mcp_config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| let [config, refreshed_config] = [McpServerAuth::OAuth, McpServerAuth::ChatGpt].map(|auth| { | |
| let mut config = reusable_server_config("https://chatgpt.com/backend-api/ps/mcp"); | |
| config.environment_id = "customer-executor".to_string(); | |
| config.auth = auth; | |
| crate::effective_mcp_servers_from_configured( | |
| HashMap::from([("docs".to_string(), config)]), | |
| &mcp_config, | |
| /*auth*/ None, | |
| ) | |
| .remove("docs") | |
| .expect("configured server should survive auth projection") | |
| .config() | |
| .clone() | |
| }); | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let reconciled = reconcile_reusable_server(&previous, refreshed_config, runtime_context).await; | |
| let outcome = reconciled | |
| .servers | |
| .get("docs") | |
| .expect("refreshed server should exist") | |
| .connection | |
| .client() | |
| .await; | |
| assert_matches!( | |
| outcome.err().expect("changed auth mode must be validated"), | |
| StartupOutcomeError::Failed { | |
| error, | |
| is_authentication_required, | |
| } => { | |
| assert_eq!( | |
| (error.as_str(), is_authentication_required), | |
| ( | |
| "executor-owned MCP server `docs` cannot use hosted ChatGPT authentication; configure executor-owned credentials instead", | |
| false, | |
| ) | |
| ); | |
| } | |
| ); | |
| assert_eq!( | |
| model_tool_names(&reconciled.list_all_tools().await), | |
| HashSet::new() | |
| ); | |
| Ok(()) | |
| } | |
| async fn reconciliation_replaces_connection_when_protocol_mode_changes() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let codex_home = tempdir().expect("tempdir"); | |
| let mut mcp_config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| mcp_config.protocol_mode = codex_rmcp_client::McpProtocolMode::V20260728; | |
| let reconciled = McpConnectionSet::new( | |
| Some(&previous), | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(mcp_config), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers: HashMap::from([( | |
| "docs".to_string(), | |
| EffectiveMcpServer::configured(config), | |
| )]), | |
| submit_id: "refresh".to_string(), | |
| tx_event: None, | |
| startup_cancellation_token: CancellationToken::new(), | |
| runtime_context, | |
| codex_apps_tools_cache: ConnectorRuntimeManager::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| assert!(!previous.shares_test_connection_with(&reconciled, "docs")); | |
| } | |
| async fn reconciliation_reuses_legacy_stdio_server_when_modern_protocol_is_enabled() { | |
| let runtime_context = McpRuntimeContext::new( | |
| Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()), | |
| PathBuf::from("/tmp"), | |
| ); | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| config.transport = McpServerTransportConfig::Stdio { | |
| command: "legacy-server".to_string(), | |
| args: Vec::new(), | |
| env: None, | |
| env_vars: Vec::new(), | |
| cwd: None, | |
| }; | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let codex_home = tempdir().expect("tempdir"); | |
| let mut mcp_config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| mcp_config.protocol_mode = codex_rmcp_client::McpProtocolMode::V20260728; | |
| let reconciled = McpConnectionSet::new( | |
| Some(&previous), | |
| McpPublicationGate::already_published(), | |
| McpRuntimeInput { | |
| startup_policy: McpStartupPolicy::Eager, | |
| config: Arc::new(mcp_config), | |
| plugins_available: false, | |
| ready_selected_capability_roots: Vec::new(), | |
| mcp_servers: HashMap::from([( | |
| "docs".to_string(), | |
| EffectiveMcpServer::configured(config), | |
| )]), | |
| submit_id: "refresh".to_string(), | |
| tx_event: None, | |
| startup_cancellation_token: CancellationToken::new(), | |
| runtime_context, | |
| codex_apps_tools_cache: ConnectorRuntimeManager::default(), | |
| tool_catalog_cache: McpToolCatalogCache::default(), | |
| codex_apps_tools_cache_key: ConnectorRuntimeContextKey::personal( | |
| /*account_id*/ None, /*chatgpt_user_id*/ None, | |
| ), | |
| client_mcp_extensions: ClientMcpExtensions::default(), | |
| auth: None, | |
| auth_manager: None, | |
| elicitation_reviewer: None, | |
| elicitation_lifecycle: None, | |
| }, | |
| ElicitationRequestRouter::default(), | |
| ) | |
| .await; | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| } | |
| async fn reconciliation_updates_elicitation_policy_without_restarting_ready_server() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| { | |
| let mut authority = previous | |
| .elicitation_requests | |
| .authority | |
| .lock() | |
| .expect("elicitation authority lock"); | |
| let config = Arc::make_mut( | |
| &mut authority | |
| .as_mut() | |
| .expect("test manager should have permission authority") | |
| .config, | |
| ); | |
| config.approval_policy = Constrained::allow_any(AskForApproval::Never); | |
| config.permission_profile = PermissionProfile::Disabled; | |
| } | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context).await; | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| let authority = reconciled | |
| .elicitation_requests | |
| .authority | |
| .lock() | |
| .expect("elicitation authority lock"); | |
| let config = &authority | |
| .as_ref() | |
| .expect("reconciled manager should have permission authority") | |
| .config; | |
| assert_eq!(config.approval_policy.value(), AskForApproval::OnRequest); | |
| assert_eq!(config.permission_profile, PermissionProfile::default()); | |
| } | |
| async fn reconciliation_reuses_ready_server_when_startup_timeout_changes() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let mut config = reusable_server_config("http://127.0.0.1:1"); | |
| let previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| config.startup_timeout_sec = Some(Duration::from_secs(60)); | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context).await; | |
| assert_eq!( | |
| model_tool_names(&reconciled.list_all_tools().await), | |
| HashSet::from([ToolName::namespaced("mcp__docs", "search")]) | |
| ); | |
| } | |
| async fn reconciliation_replaces_closed_connections() -> anyhow::Result<()> { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let mut previous = manager_with_reusable_ready_server( | |
| &config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let disconnect = CancellationToken::new(); | |
| let client = Arc::new( | |
| RmcpClient::new_in_process_client(Arc::new(DisconnectingToolsTransportFactory { | |
| server: MutableToolsServer { | |
| tools: Arc::new(tokio::sync::RwLock::new(vec![Tool::new( | |
| "search", | |
| "search", | |
| Arc::new(JsonObject::default()), | |
| )])), | |
| block_tool_listing: Arc::new(AtomicBool::new(false)), | |
| }, | |
| disconnect: disconnect.clone(), | |
| })) | |
| .await?, | |
| ); | |
| client | |
| .initialize( | |
| InitializeRequestParams::new( | |
| ClientCapabilities::default(), | |
| Implementation::new("codex-test", "0.0.0-test"), | |
| ) | |
| .with_protocol_version(ProtocolVersion::V_2025_06_18), | |
| /*timeout*/ None, | |
| Box::new(|_, _| async { Err(anyhow!("unexpected elicitation")) }.boxed()), | |
| ) | |
| .await?; | |
| let view = previous | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist"); | |
| let mut connected_client = view.connection.client().await?; | |
| connected_client.client = Arc::clone(&client); | |
| view.connection = Arc::new(McpServerConnection { | |
| identity: Some(reusable_server_identity("docs", &config, &runtime_context)), | |
| client: AsyncManagedClient { | |
| client: futures::future::ready(Ok(connected_client)) | |
| .boxed() | |
| .shared(), | |
| is_codex_apps_mcp_server: false, | |
| server_capabilities: Arc::new(std::sync::Mutex::new(None)), | |
| cached_server_info: None, | |
| codex_apps_tools_cache_context: None, | |
| tool_catalog_cache_context: None, | |
| startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)), | |
| startup_reconnect: None, | |
| cancel_token: CancellationToken::new(), | |
| }, | |
| startup_timeout: config | |
| .startup_timeout_sec | |
| .unwrap_or(DEFAULT_STARTUP_TIMEOUT), | |
| startup_trigger: None, | |
| _diagnostics_guard: LIVE_CONNECTIONS.track(), | |
| }); | |
| assert!(!client.is_closed().await); | |
| disconnect.cancel(); | |
| tokio::time::timeout(Duration::from_secs(2), async { | |
| while !client.is_closed().await { | |
| tokio::task::yield_now().await; | |
| } | |
| }) | |
| .await | |
| .expect("closed MCP transport should be detected"); | |
| let reconciled = reconcile_reusable_server(&previous, config, runtime_context).await; | |
| assert!(!previous.shares_test_connection_with(&reconciled, "docs")); | |
| Ok(()) | |
| } | |
| async fn reconciliation_reconnects_when_connection_identity_changes() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let previous_config = reusable_server_config("http://127.0.0.1:1"); | |
| let previous = manager_with_reusable_ready_server( | |
| &previous_config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let reconciled = reconcile_reusable_server( | |
| &previous, | |
| reusable_server_config("http://127.0.0.1:2"), | |
| runtime_context, | |
| ) | |
| .await; | |
| assert!(!previous.shares_test_connection_with(&reconciled, "docs")); | |
| } | |
| async fn reconciliation_reconnects_when_host_plugin_root_changes() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let server_config = reusable_server_config("http://127.0.0.1:1"); | |
| let original_root = PathUri::parse("file:///plugins/original").expect("valid plugin root URI"); | |
| let replacement_root = | |
| PathUri::parse("file:///plugins/replacement").expect("valid plugin root URI"); | |
| let mut previous = manager_with_reusable_ready_server( | |
| &server_config, | |
| &runtime_context, | |
| vec![create_test_tool("docs", "search")], | |
| ) | |
| .await; | |
| let server = EffectiveMcpServer::configured(server_config.clone()); | |
| let resolved_environment = runtime_context.resolve_server_environment("docs", &server_config); | |
| let original_identity = McpServerConnectionIdentity::new( | |
| "docs", | |
| &server, | |
| Some(&original_root), | |
| OAuthCredentialsStoreMode::default(), | |
| AuthKeyringBackendKind::default(), | |
| McpOAuthRefreshMode::Legacy, | |
| &resolved_environment, | |
| &runtime_context, | |
| /*runtime_auth_provider*/ None, | |
| /*auth*/ None, | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ); | |
| Arc::get_mut( | |
| &mut previous | |
| .servers | |
| .get_mut("docs") | |
| .expect("test server should exist") | |
| .connection, | |
| ) | |
| .expect("test server should have one connection owner") | |
| .identity = Some(original_identity); | |
| let codex_home = tempdir().expect("tempdir"); | |
| let config_for_root = |root| { | |
| let mut config = crate::mcp::tests::test_mcp_config(codex_home.path().to_path_buf()); | |
| let mut catalog = crate::ResolvedMcpCatalog::builder(); | |
| catalog.register(crate::McpServerRegistration::from_plugin( | |
| "docs".to_string(), | |
| crate::McpPluginAttribution::new("docs@test".to_string(), "Docs".to_string()) | |
| .with_host_root(root), | |
| /*plugin_order*/ 0, | |
| server_config.clone(), | |
| )); | |
| config.mcp_server_catalog = catalog.build(); | |
| config | |
| }; | |
| let unchanged = reconcile_reusable_server_with_mcp_config( | |
| &previous, | |
| "docs", | |
| server_config.clone(), | |
| runtime_context.clone(), | |
| config_for_root(original_root), | |
| ) | |
| .await; | |
| assert!(previous.shares_test_connection_with(&unchanged, "docs")); | |
| let replacement_config = config_for_root(replacement_root); | |
| let replacement = reconcile_reusable_server_with_mcp_config( | |
| &unchanged, | |
| "docs", | |
| server_config, | |
| runtime_context, | |
| replacement_config, | |
| ) | |
| .await; | |
| assert!(!unchanged.shares_test_connection_with(&replacement, "docs")); | |
| } | |
| async fn connection_identity_distinguishes_accounts_with_the_same_token() -> anyhow::Result<()> { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let server = EffectiveMcpServer::configured(config); | |
| let access_token = "header.e30.same"; | |
| let previous_auth = CodexAuth::from_external_chatgpt_tokens( | |
| access_token, | |
| "account-a", | |
| /*chatgpt_plan_type*/ None, | |
| )?; | |
| let changed_auth = CodexAuth::from_external_chatgpt_tokens( | |
| access_token, | |
| "account-b", | |
| /*chatgpt_plan_type*/ None, | |
| )?; | |
| let connection_identity = |auth: &CodexAuth| { | |
| let provider = codex_model_provider::auth_provider_from_auth(auth); | |
| McpServerConnectionIdentity::new( | |
| "docs", | |
| &server, | |
| /*host_plugin_root*/ None, | |
| OAuthCredentialsStoreMode::default(), | |
| AuthKeyringBackendKind::default(), | |
| McpOAuthRefreshMode::Legacy, | |
| &Ok(None), | |
| &runtime_context, | |
| Some(&provider), | |
| Some(auth), | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ) | |
| }; | |
| assert_eq!(previous_auth, changed_auth); | |
| assert_eq!(previous_auth.get_token()?, changed_auth.get_token()?); | |
| assert!( | |
| !connection_identity(&previous_auth) | |
| .has_same_connection_config(&connection_identity(&changed_auth)) | |
| ); | |
| Ok(()) | |
| } | |
| async fn connection_identity_distinguishes_agent_account_runtime_and_task() -> anyhow::Result<()> { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let config = reusable_server_config("http://127.0.0.1:1"); | |
| let server = EffectiveMcpServer::configured(config); | |
| let record = codex_login::auth::AgentIdentityAuthRecord { | |
| agent_runtime_id: "agent-a".to_string(), | |
| agent_private_key: "MC4CAQAwBQYDK2VwBCIEIJ7kFBaOujmoz1gvBNEC+BeM2IX87FFB0xmISOZ/XO0c" | |
| .to_string(), | |
| account_id: "account-a".to_string(), | |
| chatgpt_user_id: "user-a".to_string(), | |
| email: Some("agent@example.com".to_string()), | |
| plan_type: codex_protocol::account::PlanType::Plus, | |
| chatgpt_account_is_fedramp: false, | |
| task_id: Some("task-a".to_string()), | |
| }; | |
| let auth_route_config = codex_login::test_support::transport_default_auth_route_config(); | |
| let previous_auth = CodexAuth::AgentIdentity( | |
| codex_login::auth::AgentIdentityAuth::from_record( | |
| record.clone(), | |
| "https://auth.openai.com/api/accounts", | |
| &auth_route_config, | |
| ) | |
| .await?, | |
| ); | |
| let connection_identity = |auth: &CodexAuth| { | |
| let provider = codex_model_provider::auth_provider_from_auth(auth); | |
| McpServerConnectionIdentity::new( | |
| CODEX_APPS_MCP_SERVER_NAME, | |
| &server, | |
| /*host_plugin_root*/ None, | |
| OAuthCredentialsStoreMode::default(), | |
| AuthKeyringBackendKind::default(), | |
| McpOAuthRefreshMode::Legacy, | |
| &Ok(None), | |
| &runtime_context, | |
| Some(&provider), | |
| Some(auth), | |
| /*codex_apps_cache_identity*/ None, | |
| ElicitationCapability::default(), | |
| ClientMcpExtensions::default(), | |
| /*previous_identity*/ None, | |
| ) | |
| }; | |
| let previous_identity = connection_identity(&previous_auth); | |
| for changed_record in [ | |
| codex_login::auth::AgentIdentityAuthRecord { | |
| account_id: "account-b".to_string(), | |
| ..record.clone() | |
| }, | |
| codex_login::auth::AgentIdentityAuthRecord { | |
| chatgpt_user_id: "user-b".to_string(), | |
| ..record.clone() | |
| }, | |
| codex_login::auth::AgentIdentityAuthRecord { | |
| chatgpt_account_is_fedramp: true, | |
| ..record.clone() | |
| }, | |
| codex_login::auth::AgentIdentityAuthRecord { | |
| agent_runtime_id: "agent-b".to_string(), | |
| ..record.clone() | |
| }, | |
| codex_login::auth::AgentIdentityAuthRecord { | |
| task_id: Some("task-b".to_string()), | |
| ..record.clone() | |
| }, | |
| ] { | |
| let changed_auth = CodexAuth::AgentIdentity( | |
| codex_login::auth::AgentIdentityAuth::from_record( | |
| changed_record, | |
| "https://auth.openai.com/api/accounts", | |
| &auth_route_config, | |
| ) | |
| .await?, | |
| ); | |
| assert_eq!(previous_auth, changed_auth); | |
| assert!(!previous_identity.has_same_connection_config(&connection_identity(&changed_auth))); | |
| } | |
| Ok(()) | |
| } | |
| async fn view_only_changes_reuse_connection_and_preserve_the_old_step() { | |
| let runtime_context = reusable_server_runtime_context(); | |
| let mut old_config = reusable_server_config("http://127.0.0.1:1"); | |
| old_config.default_tools_approval_mode = Some(AppToolApproval::Prompt); | |
| let previous = Arc::new( | |
| manager_with_reusable_ready_server( | |
| &old_config, | |
| &runtime_context, | |
| vec![ | |
| create_test_tool("docs", "search"), | |
| create_test_tool("docs", "write"), | |
| ], | |
| ) | |
| .await, | |
| ); | |
| let old_step = capture_binding(&previous).await; | |
| let old_call = old_step | |
| .prepare_call("docs", "search") | |
| .expect("old step should prepare search"); | |
| let mut new_config = old_config; | |
| new_config.enabled_tools = Some(vec!["search".to_string()]); | |
| new_config.default_tools_approval_mode = Some(AppToolApproval::Approve); | |
| let reconciled = | |
| Arc::new(reconcile_reusable_server(previous.as_ref(), new_config, runtime_context).await); | |
| assert!(previous.shares_test_connection_with(&reconciled, "docs")); | |
| let new_step = capture_binding(&reconciled).await; | |
| let new_call = new_step | |
| .prepare_call("docs", "search") | |
| .expect("new step should prepare search"); | |
| drop(previous); | |
| assert_eq!( | |
| old_step | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.tool.name.to_string()) | |
| .collect::<HashSet<_>>(), | |
| HashSet::from(["search".to_string(), "write".to_string()]) | |
| ); | |
| assert_eq!(old_call.tool_approval_mode(), AppToolApproval::Prompt); | |
| assert_eq!( | |
| new_step | |
| .tools() | |
| .iter() | |
| .map(|tool| tool.tool.name.to_string()) | |
| .collect::<Vec<_>>(), | |
| vec!["search".to_string()] | |
| ); | |
| assert_eq!(new_call.tool_approval_mode(), AppToolApproval::Approve); | |
| } | |