use super::mcp_refresh::McpRefreshInvalidationGuard; use super::*; use crate::environment_selection::combine_selected_capability_roots; use codex_exec_server::ExecutorCapabilityDiscoveryCache; use codex_exec_server::ExecutorCapabilityDiscoverySnapshot; use codex_exec_server::MAX_SELECTED_CAPABILITY_ROOTS; use codex_exec_server::ResolvedSelectedCapabilityRoot; use codex_mcp::CODEX_APPS_MCP_SERVER_NAME; use codex_mcp::ElicitationReviewRequest; use codex_mcp::ElicitationReviewer; use codex_mcp::ElicitationReviewerHandle; use codex_mcp::MCP_TOOL_CODEX_APPS_META_KEY; use codex_protocol::capabilities::CapabilityRootLocation; use codex_protocol::capabilities::SelectedCapabilityRoot; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::mcp::is_node_repl_backed_server; use codex_protocol::mcp_approval_meta::APPROVAL_KIND_KEY as MCP_ELICITATION_APPROVAL_KIND_KEY; use codex_protocol::mcp_approval_meta::APPROVAL_KIND_MCP_TOOL_CALL as MCP_ELICITATION_APPROVAL_KIND_MCP_TOOL_CALL; use codex_protocol::mcp_approval_meta::APPROVAL_KIND_TOOL_SUGGESTION as MCP_ELICITATION_APPROVAL_KIND_TOOL_SUGGESTION; use codex_protocol::mcp_approval_meta::APPROVALS_REVIEWER_KEY as MCP_ELICITATION_APPROVALS_REVIEWER_KEY; use codex_protocol::mcp_approval_meta::CONNECTOR_DESCRIPTION_KEY as MCP_ELICITATION_CONNECTOR_DESCRIPTION_KEY; use codex_protocol::mcp_approval_meta::CONNECTOR_ID_KEY as MCP_ELICITATION_CONNECTOR_ID_KEY; use codex_protocol::mcp_approval_meta::CONNECTOR_NAME_KEY as MCP_ELICITATION_CONNECTOR_NAME_KEY; use codex_protocol::mcp_approval_meta::PERSIST_KEY as MCP_ELICITATION_PERSIST_KEY; use codex_protocol::mcp_approval_meta::REQUEST_TYPE_APPROVAL_REQUEST as MCP_ELICITATION_REQUEST_TYPE_APPROVAL_REQUEST; use codex_protocol::mcp_approval_meta::REQUEST_TYPE_KEY as MCP_ELICITATION_REQUEST_TYPE_KEY; use codex_protocol::mcp_approval_meta::SENSITIVE_ACTION_KEY as MCP_ELICITATION_SENSITIVE_ACTION_KEY; use codex_protocol::mcp_approval_meta::STRICT_AUTO_REVIEW_KEY as MCP_ELICITATION_STRICT_AUTO_REVIEW_KEY; use codex_protocol::mcp_approval_meta::TOOL_DESCRIPTION_KEY as MCP_ELICITATION_TOOL_DESCRIPTION_KEY; use codex_protocol::mcp_approval_meta::TOOL_NAME_KEY as MCP_ELICITATION_TOOL_NAME_KEY; use codex_protocol::mcp_approval_meta::TOOL_PARAMS_KEY as MCP_ELICITATION_TOOL_PARAMS_KEY; use codex_protocol::mcp_approval_meta::TOOL_TITLE_KEY as MCP_ELICITATION_TOOL_TITLE_KEY; use codex_protocol::openai_models::ModelInfo; use codex_rmcp_client::Elicitation; use rmcp::model::ElicitationAction; use rmcp::model::RequestMetaObject; use serde_json::Map; const MCP_ELICITATION_DECLINE_MESSAGE_KEY: &str = "message"; const TOOL_SUGGESTION_ACTION_INSTALL: &str = "install"; const TOOL_SUGGESTION_ACTION_KEY: &str = "suggest_type"; const TOOL_SUGGESTION_TOOL_ID_KEY: &str = "tool_id"; const TOOL_SUGGESTION_TOOL_TYPE_KEY: &str = "tool_type"; #[derive(Debug, PartialEq)] enum GuardianElicitationReview { NotRequested, Decline(&'static str), ApprovalRequest(Box), } struct GuardianMcpElicitationReviewer { session: std::sync::Weak, } pub(crate) struct McpServerElicitationOutcome { pub(crate) response: Option, pub(crate) sent: bool, } #[derive(Debug, PartialEq, Eq)] struct PluginInstallElicitationTelemetryMetadata { tool_type: String, tool_id: String, tool_name: String, } impl GuardianMcpElicitationReviewer { fn new(session: &Arc) -> Self { Self { session: Arc::downgrade(session), } } } impl ElicitationReviewer for GuardianMcpElicitationReviewer { fn review( &self, request: ElicitationReviewRequest, ) -> BoxFuture<'static, anyhow::Result>> { let session = self.session.clone(); Box::pin(async move { let Some(session) = session.upgrade() else { return Ok(None); }; review_guardian_mcp_elicitation(session, request).await }) } } impl Session { pub(crate) async fn runtime_mcp_config(&self, config: &Config) -> McpConfig { self.runtime_mcp_config_and_context(config).await.0 } pub(crate) async fn runtime_mcp_config_and_context( &self, config: &Config, ) -> (McpConfig, McpRuntimeContext) { let originator = self.originator().await; let (session_source, host_fallback_cwd, disabled_plugin_ids) = { let state = self.state.lock().await; ( state.session_configuration.session_source.clone(), state.session_configuration.cwd().clone(), state.active_disabled_plugin_ids.clone(), ) }; let environments = self.services.turn_environments.snapshot().await; let selected_capability_roots = self .resolve_selected_capability_roots_for_step(&environments) .await; let ready_selected_capability_roots = Self::ready_selected_capability_roots(&selected_capability_roots); let executor_capability_discovery = self .executor_capability_discovery_for_step( config, &ready_selected_capability_roots, &environments, ) .await; let mcp_projection = self .services .mcp_manager .runtime_config_for_step( config, &self.services.mcp_thread_init, &self.services.thread_extension_data, McpThreadIdentity { session_source: &session_source, originator: &originator, disabled_plugin_ids: &disabled_plugin_ids, environments: McpEnvironmentScope::Live(&self.services.turn_environments), }, &ready_selected_capability_roots, executor_capability_discovery.as_deref(), ) .await; let mcp_config = self .project_selected_environment_mcp_servers(config, &environments, mcp_projection) .await .config; let local_process_cwd = environments .local_environment_cwd() .map(|cwd| cwd.to_path_buf()) .unwrap_or_else(|| host_fallback_cwd.to_path_buf()); let runtime_context = McpRuntimeContext::new( self.services.turn_environments.environment_manager(), local_process_cwd, ) .with_selected_environments( environments .turn_environments() .map(|environment| { ( environment.selection.environment_id.clone(), Arc::clone(&environment.environment), ) }) .collect(), ); (mcp_config, runtime_context) } pub(crate) async fn runtime_mcp_servers( &self, config: &Config, ) -> HashMap { codex_mcp::configured_mcp_servers(&self.runtime_mcp_config(config).await) } /// Publishes changed MCP state, waiting for any refresh already in progress. #[tracing::instrument(name = "mcp.runtime.refresh_if_dirty", skip_all)] pub(crate) async fn refresh_mcp_if_dirty(self: &Arc) { let Ok(_refresh) = self.mcp_refresh.acquire().await else { error!("MCP runtime refresh semaphore closed"); return; }; loop { let environments = self.services.turn_environments.snapshot().await; let ready_environments = environments .turn_environments() .map(|environment| { ( environment.selection.environment_id.clone(), Arc::clone(&environment.environment), ) }) .collect(); // Attachment resolution can finish after the owner's configuration callback. if !self .services .mcp_runtime .current_environments_match(&ready_environments) { self.mark_mcp_runtime_dirty(); } let auth = self.services.auth_manager.auth_cached(); if !self .services .mcp_runtime .current_auth_matches(auth.as_ref()) { self.mark_mcp_runtime_dirty(); } if !self.mcp_refresh.claim() { return; } let mut refresh_invalidation = McpRefreshInvalidationGuard { refresh: &self.mcp_refresh, published: false, }; let auth = self.services.auth_manager.auth().await; let desired = self.latest_mcp_desired_state(auth).await; let selected_capability_roots = self .resolve_selected_capability_roots_for_step(&desired.environments) .await; let ready_selected_capability_roots = Self::ready_selected_capability_roots(&selected_capability_roots); let executor_capability_discovery = self .executor_capability_discovery_for_step( &desired.config, &ready_selected_capability_roots, &desired.environments, ) .await; let mcp_projection = self .services .mcp_manager .runtime_config_for_step( &desired.config, &self.services.mcp_thread_init, &self.services.thread_extension_data, McpThreadIdentity { session_source: &desired.session_source, originator: &desired.originator, disabled_plugin_ids: &desired.disabled_plugin_ids, environments: McpEnvironmentScope::Live(&self.services.turn_environments), }, &ready_selected_capability_roots, executor_capability_discovery.as_deref(), ) .await; self.publish_mcp_runtime( &desired, mcp_projection, &ready_selected_capability_roots, Some(self.mcp_elicitation_reviewer()), ) .await; refresh_invalidation.published = true; } } /// Refreshes Apps tools on the published thread runtime and returns that client's snapshot. pub(crate) async fn refresh_codex_apps_tools( self: &Arc, ) -> anyhow::Result { // Reconcile unchanged config so failed or closed clients can be replaced. self.mark_mcp_runtime_dirty(); self.refresh_mcp_if_dirty().await; let _refresh = self .mcp_refresh .acquire() .await .map_err(|_| anyhow::anyhow!("MCP runtime refresh semaphore closed"))?; self.services.mcp_runtime.refresh_codex_apps_tools().await } /// Reconnects the runtime so refreshed Apps tools belong to their new exact client. pub(crate) async fn hard_refresh_latest_codex_apps_tools( self: &Arc, ) -> anyhow::Result> { self.refresh_mcp_if_dirty().await; let _refresh = self .mcp_refresh .acquire() .await .map_err(|_| anyhow::anyhow!("MCP runtime refresh semaphore closed"))?; let auth = self.services.auth_manager.auth().await; let desired = self.latest_mcp_desired_state(auth).await; let selected_capability_roots = self .resolve_selected_capability_roots_for_step(&desired.environments) .await; let ready_selected_capability_roots = Self::ready_selected_capability_roots(&selected_capability_roots); let executor_capability_discovery = self .executor_capability_discovery_for_step( &desired.config, &ready_selected_capability_roots, &desired.environments, ) .await; let mcp_projection = self .services .mcp_manager .runtime_config_for_step( &desired.config, &self.services.mcp_thread_init, &self.services.thread_extension_data, McpThreadIdentity { session_source: &desired.session_source, originator: &desired.originator, disabled_plugin_ids: &desired.disabled_plugin_ids, environments: McpEnvironmentScope::Live(&self.services.turn_environments), }, &ready_selected_capability_roots, executor_capability_discovery.as_deref(), ) .await; let mcp_projection = self .project_selected_environment_mcp_servers( &desired.config, &desired.environments, mcp_projection, ) .await; let selected_plugins = mcp_projection.selected_plugins.clone(); let input = self.build_mcp_runtime_input( &desired, mcp_projection, &ready_selected_capability_roots, Some(self.mcp_elicitation_reviewer()), ); anyhow::ensure!( input.mcp_servers.contains_key(CODEX_APPS_MCP_SERVER_NAME), "unknown MCP server '{CODEX_APPS_MCP_SERVER_NAME}'" ); let refreshed = self.services.mcp_runtime.replace_fresh(input).await; self.services.thread_extension_data.insert(selected_plugins); refreshed } pub(super) fn mark_mcp_runtime_dirty(&self) { self.mcp_refresh.invalidate(); } #[tracing::instrument(name = "mcp.runtime.resolve_for_step", skip_all)] pub(crate) async fn mcp_runtime_for_step( self: &Arc, turn_context: &TurnContext, selected_capability_roots: &[ResolvedSelectedCapabilityRoot], required_servers: &[String], required_plugins: &HashSet, ) -> Arc { let ready_selected_capability_roots = Self::ready_selected_capability_roots(selected_capability_roots); if self .services .mcp_runtime .current_ready_selected_capability_roots() != ready_selected_capability_roots { self.mark_mcp_runtime_dirty(); } let recovered_oauth_servers = self .services .mcp_runtime .updated_oauth_credentials_after_auth_failure() .await; if !recovered_oauth_servers.is_empty() && let Ok(_refresh) = self.mcp_refresh.acquire().await && self .services .mcp_runtime .has_authentication_failed_servers(&recovered_oauth_servers) .await { self.mark_mcp_runtime_dirty(); } self.refresh_mcp_if_dirty().await; let required_servers = required_servers .iter() .chain(&recovered_oauth_servers) .cloned() .collect::>(); if let Some(binding) = self .services .mcp_runtime .current_binding_with_requirements(&required_servers, required_plugins) .await { return binding; } let config = Arc::new(self.runtime_mcp_config(&turn_context.config).await); Arc::new(codex_mcp::McpBinding::empty(config)) } #[tracing::instrument( name = "capability_roots.snapshot_for_step", skip_all, fields(root_count = ready_selected_capability_roots.len()) )] pub(crate) async fn executor_capability_discovery_for_step( &self, config: &Config, ready_selected_capability_roots: &[SelectedCapabilityRoot], environments: &TurnEnvironmentSnapshot, ) -> Option> { // Capability roots can currently be selected independently of turn environments, so a // root may be ready when there is no primary `TurnEnvironment`. Keep using the thread // policy in that case so restricted discovery fails closed below. Once every selected // root belongs to a thread/environment attachment whose `EnvironmentConfig` is installed // before the root becomes ready, discovery can use the root owner's policy and this // fallback can be removed. let restricted_file_system = environments.primary().map_or_else( || { !config .permissions .file_system_sandbox_policy() .has_full_disk_read_access() }, |_| { environments.turn_environments().any(|environment| { !environment .permission_profile() .file_system_sandbox_policy() .has_full_disk_read_access() }) }, ); if !restricted_file_system && !config .features .enabled(Feature::ExecutorCapabilityDiscovery) { return None; } let sandbox_contexts = if restricted_file_system { environments .turn_environments() .map(|environment| { ( environment.selection.environment_id.clone(), environment.sandbox_context(/*additional_permissions*/ None), ) }) .collect::>() } else { HashMap::new() }; let environment_manager = self.services.turn_environments.environment_manager(); let cache = self .services .thread_extension_data .get_or_init(|| ExecutorCapabilityDiscoveryCache::new(environment_manager)); let selected_capability_roots = ready_selected_capability_roots .iter() .filter(|selected_root| { if !restricted_file_system { return true; } let CapabilityRootLocation::Environment { environment_id, .. } = &selected_root.location; if sandbox_contexts.contains_key(environment_id) { return true; } warn!( selected_root = selected_root.id, environment_id, "skipping capability root without a filesystem sandbox context" ); false }) .cloned() .collect::>(); let discovery = cache .snapshot(&selected_capability_roots, &sandbox_contexts) .await; if cache.take_recovered_discovery() { // Root selection is unchanged, but recovered manifests can change MCP servers. self.mark_mcp_runtime_dirty(); } Some(Arc::new(discovery)) } pub(crate) async fn resolve_selected_capability_roots_for_step( &self, environments: &TurnEnvironmentSnapshot, ) -> Vec { let thread_root_count = self.services.selected_capability_roots.len(); let mut root_locations_by_id = HashMap::new(); let mut selected_capability_roots = Vec::new(); let mut ready_environment_root_count = 0; let combined_roots = combine_selected_capability_roots( &self.services.selected_capability_roots, environments.turn_environments().map(|environment| { ( environment.config_origin, environment .config_origin .selected_capability_roots(&environment.environment, environment.config()), ) }), ); for (index, root) in combined_roots.into_iter().enumerate() { if let Some(kept_location) = root_locations_by_id.get(&root.id) { if kept_location != &root.location { tracing::warn!( root_id = root.id, ?kept_location, ignored_location = ?root.location, "ignoring selected capability root with conflicting location" ); } continue; } if index >= thread_root_count { if ready_environment_root_count == MAX_SELECTED_CAPABILITY_ROOTS { tracing::warn!( max_root_count = MAX_SELECTED_CAPABILITY_ROOTS, "ignoring excess selected capability roots from ready environments" ); break; } ready_environment_root_count += 1; } root_locations_by_id.insert(root.id.clone(), root.location.clone()); selected_capability_roots.push(root); } self.services .turn_environments .environment_manager() .resolve_selected_capability_roots( &selected_capability_roots, &environments.captured_environments(), ) .await } pub(crate) fn mcp_elicitation_reviewer(self: &Arc) -> ElicitationReviewerHandle { Arc::clone( self.mcp_elicitation_reviewer_handle .get_or_init(|| Arc::new(GuardianMcpElicitationReviewer::new(self))), ) } pub(crate) fn mcp_elicitation_lifecycle(&self) -> codex_mcp::ElicitationLifecycle { self.mcp_elicitation_lifecycle_handle .get_or_init(|| { let elicitations = self.services.elicitations.clone(); codex_mcp::ElicitationLifecycle::new(move || elicitations.register()) }) .clone() } #[expect( clippy::await_holding_invalid_type, reason = "active turn checks and turn state updates must remain atomic" )] pub async fn request_mcp_server_elicitation( &self, turn_context: &TurnContext, server_name: String, request_id: RequestId, request: ElicitationRequest, ) -> McpServerElicitationOutcome { if self.services.mcp_runtime.elicitations_auto_deny() { return McpServerElicitationOutcome { response: Some(ElicitationResponse { action: codex_rmcp_client::ElicitationAction::Accept, content: Some(serde_json::json!({})), meta: None, }), sent: false, }; } let _elicitation = self.services.elicitations.register(); let (tx_response, rx_response) = oneshot::channel(); let prev_entry = { let mut active = self.active_turn.lock().await; match active.as_mut() { Some(at) => { let mut ts = at.turn_state.lock().await; ts.insert_pending_elicitation( server_name.clone(), request_id.clone(), tx_response, ) } None => None, } }; if prev_entry.is_some() { warn!( "Overwriting existing pending elicitation for server_name: {server_name}, request_id: {request_id}" ); } let id = match request_id { rmcp::model::NumberOrString::String(value) => { codex_protocol::mcp::RequestId::String(value.to_string()) } rmcp::model::NumberOrString::Number(value) => { codex_protocol::mcp::RequestId::Integer(value) } }; let event = EventMsg::ElicitationRequest(ElicitationRequestEvent { turn_id: Some(turn_context.sub_id.clone()), server_name, id, request, }); let plugin_install_telemetry = plugin_install_elicitation_telemetry_metadata(&event); turn_context .turn_metadata_state .mark_user_input_requested_during_turn(); self.send_event(turn_context, event).await; if let Some(plugin_install_telemetry) = plugin_install_telemetry { turn_context .session_telemetry .record_plugin_install_elicitation_sent( plugin_install_telemetry.tool_type.as_str(), plugin_install_telemetry.tool_id.as_str(), plugin_install_telemetry.tool_name.as_str(), ); } McpServerElicitationOutcome { response: rx_response.await.ok(), sent: true, } } #[expect( clippy::await_holding_invalid_type, reason = "active turn checks and manager fallback must stay serialized" )] pub async fn resolve_elicitation( &self, server_name: String, id: RequestId, response: ElicitationResponse, ) -> anyhow::Result<()> { let entry = { let mut active = self.active_turn.lock().await; match active.as_mut() { Some(at) => { let mut ts = at.turn_state.lock().await; ts.remove_pending_elicitation(&server_name, &id) } None => None, } }; if let Some(tx_response) = entry { tx_response .send(response) .map_err(|e| anyhow::anyhow!("failed to send elicitation response: {e:?}"))?; return Ok(()); } self.services .mcp_runtime .resolve_elicitation(server_name, id, response) .await } pub(crate) async fn refresh_mcp_servers_now( &self, turn_context: &TurnContext, refresh_config: &Config, elicitation_reviewer: Option, ) { let Ok(_refresh) = self.mcp_refresh.acquire().await else { error!("MCP runtime refresh semaphore closed"); return; }; let auth = self.services.auth_manager.auth().await; let disabled_plugin_ids = { let mut state = self.state.lock().await; let mut config = (*state.session_configuration.original_config_do_not_use).clone(); config.mcp_servers = refresh_config.mcp_servers.clone(); state.session_configuration.original_config_do_not_use = Arc::new(config); state.active_disabled_plugin_ids.clone() }; let ready_selected_capability_roots = self .services .mcp_runtime .current_ready_selected_capability_roots(); let environments = self.services.turn_environments.snapshot().await; let executor_capability_discovery = self .executor_capability_discovery_for_step( refresh_config, &ready_selected_capability_roots, &environments, ) .await; let mcp_projection = self .services .mcp_manager .runtime_config_for_step( refresh_config, &self.services.mcp_thread_init, &self.services.thread_extension_data, McpThreadIdentity { session_source: &turn_context.session_source, originator: &turn_context.originator, disabled_plugin_ids: &disabled_plugin_ids, environments: McpEnvironmentScope::Live(&self.services.turn_environments), }, &ready_selected_capability_roots, executor_capability_discovery.as_deref(), ) .await; let mut desired = self.latest_mcp_desired_state(auth).await; desired.config = Arc::new(refresh_config.clone()); self.publish_mcp_runtime( &desired, mcp_projection, &ready_selected_capability_roots, elicitation_reviewer, ) .await; } pub(crate) fn ready_selected_capability_roots( selected_capability_roots: &[ResolvedSelectedCapabilityRoot], ) -> Vec { selected_capability_roots .iter() .map(|root| root.selected_root().clone()) .collect() } pub(crate) fn cancel_mcp_startup(&self) { self.services.mcp_runtime.cancel_startup(); } } async fn review_guardian_mcp_elicitation( session: Arc, request: ElicitationReviewRequest, ) -> anyhow::Result> { let Some((turn_context, cancellation_token)) = session.active_turn_context_and_cancellation_token().await else { return Ok(None); }; let Some(mcp_config) = session.services.mcp_runtime.current_config() else { return Ok(None); }; let step_settings = turn_context.current_settings.load_full(); // User approval skips ordinary CUA checks, not separate sensitive requests. let user_cua_execution = step_settings.approvals_reviewer() == ApprovalsReviewer::User && is_node_repl_backed_server(&request.server_name) && request.elicitation.meta().is_some_and(|meta| { metadata_str(meta, MCP_ELICITATION_TOOL_NAME_KEY) == Some("js") && metadata_str(meta, MCP_ELICITATION_CONNECTOR_ID_KEY) == Some("node_repl") && metadata_str(meta, MCP_ELICITATION_APPROVAL_KIND_KEY) == Some(MCP_ELICITATION_APPROVAL_KIND_MCP_TOOL_CALL) && meta.get(MCP_ELICITATION_SENSITIVE_ACTION_KEY) != Some(&Value::Bool(true)) && meta.get("codex_requires_user_input") != Some(&Value::Bool(true)) }); // Full Access skips inference, not the active-turn and cancellation checks. if (user_cua_execution || turn_context.environments.has_full_access( turn_context.approval_policy(), &turn_context .config .permissions .effective_permission_profile(), )) && matches!( &request.elicitation, Elicitation::Mcp(rmcp::model::ElicitRequestParams::FormElicitationParams { requested_schema, .. }) if requested_schema.properties.is_empty() ) { let decision = if cancellation_token.is_cancelled() { ReviewDecision::Abort } else { ReviewDecision::Approved }; return Ok(Some(mcp_elicitation_response_from_guardian_decision( decision, turn_context.model_info(), ))); } // The invocation identifies the tool event, but a nested elicitation can // review a different action and connector than the enclosing JavaScript. let originating_call_id = if is_node_repl_backed_server(&request.server_name) && let Some(call_id) = request .elicitation .meta() .and_then(|meta| meta.get("callId")) .and_then(Value::as_str) && let Some((Some(invocation), _)) = session.mcp_tool_approval_metadata(&request.server_name, call_id) && invocation.server == request.server_name { Some(call_id) } else { None }; let require_synchronous_review = matches!( request .elicitation .meta() .and_then(|meta| meta.get(MCP_ELICITATION_SENSITIVE_ACTION_KEY)), Some(Value::Bool(true)) ); let mut review_context = crate::guardian::GuardianReviewContext::from(&turn_context); let strict_auto_review = matches!( request .elicitation .meta() .and_then(|meta| meta.get(MCP_ELICITATION_STRICT_AUTO_REVIEW_KEY)), Some(Value::Bool(true)) ); let guardian_request = if strict_auto_review { let connector_id = elicitation_connector_id(&request.elicitation); let trusted_guardian_request = if request.server_name == CODEX_APPS_MCP_SERVER_NAME { let Some(call_id) = request .elicitation .meta() .and_then(|meta| meta.get(MCP_TOOL_CODEX_APPS_META_KEY)) .and_then(Value::as_object) .and_then(|meta| meta.get("call_id")) .and_then(Value::as_str) else { return Ok(None); }; let Some((Some(invocation), metadata)) = session.mcp_tool_approval_metadata(&request.server_name, call_id) else { return Ok(None); }; if invocation.server != request.server_name || connector_id != metadata.connector_id.as_deref() || request .elicitation .meta() .and_then(|meta| metadata_str(meta, MCP_ELICITATION_TOOL_NAME_KEY)) != Some(invocation.tool.as_str()) { return Ok(None); } Some( crate::mcp_tool_call::build_guardian_mcp_tool_review_request( call_id, &invocation, Some(&metadata), ), ) } else { None }; if !turn_context .config .features .enabled(Feature::GuardianApproval) || mcp_config .config_layer_stack .requirements() .approvals_reviewer .can_set(&ApprovalsReviewer::AutoReview) .is_err() || crate::connectors::mcp_approvals_reviewer_from_layers( &mcp_config.config_layer_stack, ApprovalsReviewer::AutoReview, Some(turn_context.model_info().slug.as_str()), request.server_name.as_str(), connector_id, /*link_id*/ None, ) != ApprovalsReviewer::AutoReview || request .elicitation .meta() .is_some_and(|meta| meta.contains_key(MCP_ELICITATION_PERSIST_KEY)) { return Ok(None); } let GuardianElicitationReview::ApprovalRequest(guardian_request) = guardian_elicitation_review_request(&request, originating_call_id) else { return Ok(None); }; trusted_guardian_request.unwrap_or(*guardian_request).into() } else { let approval_policy = mcp_config.approval_policy.value(); match approval_policy { AskForApproval::Never => { let Some(permission_profile) = mcp_config.permission_profile_for_server(&request.server_name) else { return Ok(Some(mcp_elicitation_decline_without_message())); }; if codex_mcp::mcp_permission_prompt_is_auto_approved( approval_policy, permission_profile, codex_mcp::McpPermissionPromptAutoApproveContext::default(), ) && matches!( &request.elicitation, Elicitation::Mcp( rmcp::model::ElicitRequestParams::FormElicitationParams { requested_schema, .. } ) if requested_schema.properties.is_empty() ) { return Ok(Some(ElicitationResponse { action: ElicitationAction::Accept, content: Some(serde_json::json!({})), meta: None, })); } return Ok(Some(mcp_elicitation_decline_without_message())); } AskForApproval::Granular(config) if !config.allows_mcp_elicitations() => { return Ok(Some(mcp_elicitation_decline_without_message())); } AskForApproval::OnRequest | AskForApproval::UnlessTrusted | AskForApproval::Granular(_) => {} } let approvals_reviewer = crate::connectors::mcp_approvals_reviewer_from_layers( &mcp_config.config_layer_stack, step_settings .mcp_approvals_reviewer_override .unwrap_or(mcp_config.approvals_reviewer), Some(turn_context.model_info().slug.as_str()), request.server_name.as_str(), elicitation_connector_id(&request.elicitation), /*link_id*/ None, ); review_context.approval_policy = approval_policy; review_context.approvals_reviewer = approvals_reviewer; match guardian_elicitation_review_request(&request, originating_call_id) { GuardianElicitationReview::NotRequested => return Ok(None), GuardianElicitationReview::Decline(reason) => { let Elicitation::Mcp(elicitation) = &request.elicitation else { return Ok(None); }; crate::guardian::ReviewAction { action: serde_json::to_value(elicitation) .map(|elicitation| { serde_json::json!({ "tool": "mcp_elicitation", "server": request.server_name, "request": elicitation, }) }) .map_err(|error| error.to_string()), category: codex_protocol::openai_models::GuardianScope::for_mcp_server( &request.server_name, ), request: Err(reason.to_owned()), } } GuardianElicitationReview::ApprovalRequest(guardian_request) => { (*guardian_request).into() } } }; let declined_reason = guardian_request.request.as_ref().err().cloned(); let decision = crate::guardian::decide_approval( session, review_context, crate::guardian::new_guardian_review_id(), guardian_request, crate::tools::sandboxing::ApprovalRequestReasons { approval: None, retry: None, }, crate::guardian::GuardianReviewOptions { require_guardian: strict_auto_review, plugin_attribution_override: None, approval_request_source: codex_analytics::GuardianApprovalRequestSource::MainTurn, external_cancel: Some(cancellation_token), require_synchronous_review, }, ) .await; Ok(decision.map(|decision| { if let Some(reason) = declined_reason { warn!(server_name = %request.server_name, request_id = %mcp_elicitation_request_id(&request.request_id), reason, "declining Guardian MCP elicitation before review"); return mcp_elicitation_decline_without_message(); } mcp_elicitation_response_from_guardian_decision(decision, turn_context.model_info()) })) } fn guardian_elicitation_review_request( request: &ElicitationReviewRequest, originating_call_id: Option<&str>, ) -> GuardianElicitationReview { let (meta, requested_schema) = match &request.elicitation { Elicitation::Mcp(rmcp::model::ElicitRequestParams::FormElicitationParams { meta, requested_schema, .. }) => (meta, Some(requested_schema)), Elicitation::Mcp(rmcp::model::ElicitRequestParams::UrlElicitationParams { meta, .. }) => { return if meta_requests_approval_request(meta) { GuardianElicitationReview::Decline( "guardian MCP elicitation review only supports form elicitations", ) } else { GuardianElicitationReview::NotRequested }; } Elicitation::Mcp(_) => { return GuardianElicitationReview::Decline( "guardian MCP elicitation review does not support this elicitation mode", ); } Elicitation::OpenAiForm { .. } | Elicitation::OpenAiElicitationForm { .. } | Elicitation::UserVerification { .. } => { return GuardianElicitationReview::NotRequested; } }; let Some(meta) = meta.as_ref().map(|meta| &meta.0.0) else { return GuardianElicitationReview::NotRequested; }; if metadata_str(meta, MCP_ELICITATION_REQUEST_TYPE_KEY) != Some(MCP_ELICITATION_REQUEST_TYPE_APPROVAL_REQUEST) { return GuardianElicitationReview::NotRequested; } if metadata_str(meta, MCP_ELICITATION_APPROVAL_KIND_KEY) != Some(MCP_ELICITATION_APPROVAL_KIND_MCP_TOOL_CALL) { return GuardianElicitationReview::Decline( "guardian MCP elicitation metadata must declare mcp_tool_call approval kind", ); } if requested_schema.is_some_and(|schema| !schema.properties.is_empty()) { return GuardianElicitationReview::Decline( "guardian MCP elicitation review only supports empty form schemas", ); } let Some(tool_name) = metadata_owned_string(meta, MCP_ELICITATION_TOOL_NAME_KEY) else { return GuardianElicitationReview::Decline( "guardian MCP elicitation metadata must include a non-empty tool_name", ); }; let arguments = match meta.get(MCP_ELICITATION_TOOL_PARAMS_KEY) { Some(value @ Value::Object(_)) => Some(value.clone()), Some(_) => { return GuardianElicitationReview::Decline( "guardian MCP elicitation tool_params must be an object", ); } None => Some(Value::Object(Map::new())), }; GuardianElicitationReview::ApprovalRequest(Box::new( crate::guardian::GuardianApprovalRequest::McpToolCall { id: originating_call_id.map(str::to_owned).unwrap_or_else(|| { format!( "mcp_elicitation:{}:{}", request.server_name, mcp_elicitation_request_id(&request.request_id) ) }), server: request.server_name.clone(), tool_name, arguments, connector_id: metadata_owned_string(meta, MCP_ELICITATION_CONNECTOR_ID_KEY), connector_name: metadata_owned_string(meta, MCP_ELICITATION_CONNECTOR_NAME_KEY), connector_description: metadata_owned_string( meta, MCP_ELICITATION_CONNECTOR_DESCRIPTION_KEY, ), connected_account_email: None, tool_title: metadata_owned_string(meta, MCP_ELICITATION_TOOL_TITLE_KEY), tool_description: metadata_owned_string(meta, MCP_ELICITATION_TOOL_DESCRIPTION_KEY), annotations: None, }, )) } fn elicitation_connector_id(elicitation: &Elicitation) -> Option<&str> { elicitation .meta() .and_then(|meta| metadata_str(meta, MCP_ELICITATION_CONNECTOR_ID_KEY)) } fn meta_requests_approval_request(meta: &Option) -> bool { meta.as_ref() .and_then(|meta| metadata_str(&meta.0, MCP_ELICITATION_REQUEST_TYPE_KEY)) == Some(MCP_ELICITATION_REQUEST_TYPE_APPROVAL_REQUEST) } fn metadata_str<'a>(meta: &'a Map, key: &str) -> Option<&'a str> { meta.get(key).and_then(Value::as_str) } fn metadata_owned_string(meta: &Map, key: &str) -> Option { metadata_str(meta, key) .map(str::trim) .filter(|value| !value.is_empty()) .map(ToOwned::to_owned) } fn plugin_install_elicitation_telemetry_metadata( event: &EventMsg, ) -> Option { let EventMsg::ElicitationRequest(ElicitationRequestEvent { request, .. }) = event else { return None; }; let codex_protocol::approvals::ElicitationRequest::Form { meta: Some(Value::Object(meta)), .. } = request else { return None; }; if metadata_str(meta, MCP_ELICITATION_APPROVAL_KIND_KEY) != Some(MCP_ELICITATION_APPROVAL_KIND_TOOL_SUGGESTION) || metadata_str(meta, TOOL_SUGGESTION_ACTION_KEY) != Some(TOOL_SUGGESTION_ACTION_INSTALL) { return None; } Some(PluginInstallElicitationTelemetryMetadata { tool_type: metadata_owned_string(meta, TOOL_SUGGESTION_TOOL_TYPE_KEY)?, tool_id: metadata_owned_string(meta, TOOL_SUGGESTION_TOOL_ID_KEY)?, tool_name: metadata_owned_string(meta, MCP_ELICITATION_TOOL_NAME_KEY)?, }) } fn mcp_elicitation_request_id(id: &RequestId) -> String { match id { rmcp::model::NumberOrString::String(value) => value.to_string(), rmcp::model::NumberOrString::Number(value) => value.to_string(), } } fn mcp_elicitation_response_from_guardian_decision( decision: ReviewDecision, model_info: &ModelInfo, ) -> ElicitationResponse { match decision { ReviewDecision::Approved | ReviewDecision::ApprovedForSession | ReviewDecision::ApprovedMcpPolicyAmendment | ReviewDecision::ApprovedExecpolicyAmendment { .. } | ReviewDecision::NetworkPolicyAmendment { .. } => ElicitationResponse { action: ElicitationAction::Accept, content: Some(serde_json::json!({})), meta: Some(mcp_elicitation_auto_meta()), }, ReviewDecision::Denied { rejection } => mcp_elicitation_decline_with_message(rejection), ReviewDecision::TimedOut => mcp_elicitation_decline_with_message( crate::guardian::guardian_timeout_message(model_info), ), ReviewDecision::Abort => ElicitationResponse { action: ElicitationAction::Cancel, content: None, meta: Some(mcp_elicitation_auto_meta()), }, } } fn mcp_elicitation_decline_with_message(message: String) -> ElicitationResponse { ElicitationResponse { action: ElicitationAction::Decline, content: None, meta: Some(serde_json::json!({ MCP_ELICITATION_DECLINE_MESSAGE_KEY: message, MCP_ELICITATION_APPROVALS_REVIEWER_KEY: ApprovalsReviewer::AutoReview, })), } } fn mcp_elicitation_decline_without_message() -> ElicitationResponse { ElicitationResponse { action: ElicitationAction::Decline, content: None, meta: Some(mcp_elicitation_auto_meta()), } } fn mcp_elicitation_auto_meta() -> serde_json::Value { serde_json::json!({ MCP_ELICITATION_APPROVALS_REVIEWER_KEY: ApprovalsReviewer::AutoReview, }) } #[cfg(test)] #[path = "mcp_tests.rs"] mod tests;