diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 8f8cb83eee..97d8c77f98 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -5498,6 +5498,7 @@ dependencies = [ "anyhow", "assert_cmd", "base64 0.22.1", + "codex-analytics", "codex-arg0", "codex-attachment-store", "codex-config", diff --git a/codex-rs/core/src/mcp_tool_call.rs b/codex-rs/core/src/mcp_tool_call.rs index 70ae3e65cb..27e6a77f82 100644 --- a/codex-rs/core/src/mcp_tool_call.rs +++ b/codex-rs/core/src/mcp_tool_call.rs @@ -23,7 +23,9 @@ use crate::tools::sandboxing::ApprovalAction; use crate::tools::sandboxing::ToolError; use crate::turn_metadata::ExecutionMetadata; use codex_analytics::AppInvocation; +use codex_analytics::ElicitationType; use codex_analytics::InvocationType; +use codex_analytics::McpToolCallElicitation; use codex_analytics::build_track_events_context; use codex_api::HostedFileUploadContext; use codex_config::ConfigLayerSource; @@ -41,6 +43,7 @@ use codex_mcp::SandboxState; use codex_mcp::ToolInfo; use codex_mcp::auth_elicitation_completed_result; use codex_mcp::build_auth_elicitation_plan; +use codex_mcp::is_connector_auth_failure_from_tool_result; use codex_mcp::mcp_permission_prompt_is_auto_approved; use codex_protocol::approvals::ElicitationRequest; use codex_protocol::items::McpToolCallError; @@ -450,6 +453,7 @@ async fn handle_approved_mcp_tool_call( let mut tool_input = arguments_value .clone() .unwrap_or_else(|| JsonValue::Object(serde_json::Map::new())); + let mut elicitation_type = None; let result = async { let result = async { let mut result = prepared_call @@ -532,6 +536,8 @@ async fn handle_approved_mcp_tool_call( }) .await .map_err(|error| format!("tool call error: {error:?}"))?; + // Capture trusted server metadata before result callbacks or model-facing rewrites. + elicitation_type = mcp_tool_call_auth_elicitation_type(&server, connector_id, &result); let mcp_tool = McpToolContext::from_prepared_call( &prepared_call, turn_context.config.mcp_servers.get().get(&server), @@ -581,6 +587,9 @@ async fn handle_approved_mcp_tool_call( tracing::warn!("MCP tool call error: {error:?}"); } let duration = start.elapsed(); + if let Some(elicitation_type) = elicitation_type { + track_mcp_tool_call_elicitation(sess, turn_context, call_id, elicitation_type); + } notify_mcp_tool_call_completed( sess, turn_context, @@ -591,7 +600,7 @@ async fn handle_approved_mcp_tool_call( truncate_mcp_tool_result_for_event(&result), ) .await; - maybe_track_codex_app_used(sess, step_context, &server, &metadata).await; + maybe_track_codex_app_used(sess, step_context, &server, &metadata, elicitation_type).await; let outcome = mcp_call_metric_outcome(&result); emit_mcp_call_metrics( @@ -1059,6 +1068,7 @@ async fn maybe_track_codex_app_used( step_context: &StepContext, server: &str, metadata: &McpToolApprovalMetadata, + elicitation_type: Option, ) { if server != CODEX_APPS_MCP_SERVER_NAME { return; @@ -1090,10 +1100,36 @@ async fn maybe_track_codex_app_used( app_name, invocation_type: Some(invocation_type), }, - /*elicitation_type*/ None, + elicitation_type, ); } +fn track_mcp_tool_call_elicitation( + sess: &Session, + turn_context: &TurnContext, + call_id: &str, + elicitation_type: ElicitationType, +) { + sess.services + .analytics_events_client + .track_mcp_tool_call_elicitation(McpToolCallElicitation { + thread_id: sess.thread_id.to_string(), + turn_id: turn_context.sub_id.clone(), + item_id: call_id.to_string(), + elicitation_type, + }); +} + +fn mcp_tool_call_auth_elicitation_type( + server: &str, + connector_id: Option<&str>, + result: &CallToolResult, +) -> Option { + (server == CODEX_APPS_MCP_SERVER_NAME + && is_connector_auth_failure_from_tool_result(result, connector_id)) + .then_some(ElicitationType::AuthOrLink) +} + #[derive(Clone, Copy)] struct McpToolApprovalPolicy { mode: AppToolApproval, @@ -1583,7 +1619,7 @@ pub(crate) async fn request_mcp_tool_user_approval( .as_ref() .map(|rendered_template| rendered_template.question.as_str()), ); - if tool_call_mcp_elicitation_enabled { + let (request_dispatched, decision) = if tool_call_mcp_elicitation_enabled { let link_id = sess .mcp_tool_approval_metadata(server, id) .and_then(|(_, metadata)| metadata.link_id); @@ -1617,27 +1653,40 @@ pub(crate) async fn request_mcp_tool_user_approval( .map(|rendered_template| rendered_template.elicitation_message.as_str()), prompt_options, }); - let decision = parse_mcp_tool_approval_elicitation_response( - sess.request_mcp_server_elicitation(turn_context, server.clone(), request_id, request) - .await - .response, - &question_id, - ); - return normalize_approval_decision_for_mode(decision, *approval_mode); - } - - let args = RequestUserInputArgs { - questions: vec![question], - is_blocking: true, - auto_resolution_ms: None, + let outcome = sess + .request_mcp_server_elicitation(turn_context, server.clone(), request_id, request) + .await; + ( + outcome.sent, + parse_mcp_tool_approval_elicitation_response(outcome.response, &question_id), + ) + } else { + let args = RequestUserInputArgs { + questions: vec![question], + is_blocking: true, + auto_resolution_ms: None, + }; + let response = sess + .request_user_input(turn_context, call_id.to_string(), args) + .await; + ( + true, + parse_mcp_tool_approval_response( + response.map(|accepted| accepted.response), + &question_id, + ), + ) }; - let response = sess - .request_user_input(turn_context, call_id.to_string(), args) - .await; - normalize_approval_decision_for_mode( - parse_mcp_tool_approval_response(response.map(|accepted| accepted.response), &question_id), - *approval_mode, - ) + let decision = normalize_approval_decision_for_mode(decision, *approval_mode); + if request_dispatched + && matches!( + decision, + ReviewDecision::Denied { .. } | ReviewDecision::TimedOut | ReviewDecision::Abort + ) + { + track_mcp_tool_call_elicitation(sess, turn_context, call_id, ElicitationType::Approval); + } + decision } fn session_mcp_tool_approval_key( diff --git a/codex-rs/core/src/mcp_tool_call_tests.rs b/codex-rs/core/src/mcp_tool_call_tests.rs index 2e7fee3950..35d360202d 100644 --- a/codex-rs/core/src/mcp_tool_call_tests.rs +++ b/codex-rs/core/src/mcp_tool_call_tests.rs @@ -15,6 +15,7 @@ use crate::state::ActiveTurn; use crate::test_support::models_manager_with_provider; use crate::tools::hook_names::HookToolName; use crate::turn_metadata::ExecutionMetadata; +use codex_app_server_protocol as app_server_protocol; use codex_config::CONFIG_TOML_FILE; use codex_config::config_toml::ConfigToml; use codex_config::types::AppConfig; @@ -1615,6 +1616,28 @@ fn codex_apps_auth_failure_metadata() -> McpToolApprovalMetadata { ) } +#[test] +fn codex_apps_auth_classification_requires_the_trusted_server_and_connector() { + let result = codex_apps_auth_failure_result(); + + assert_eq!( + mcp_tool_call_auth_elicitation_type( + "untrusted_mcp_server", + Some("connector_calendar"), + &result, + ), + None + ); + assert_eq!( + mcp_tool_call_auth_elicitation_type( + CODEX_APPS_MCP_SERVER_NAME, + Some("connector_drive"), + &result, + ), + None + ); +} + #[tokio::test] async fn codex_apps_auth_elicitation_feature_disabled_returns_original_result() { let (session, mut turn_context, rx_event) = make_session_and_context_with_rx().await; @@ -2153,6 +2176,147 @@ fn accepted_elicitation_without_content_defaults_to_accept() { assert_eq!(response, ReviewDecision::Approved); } +#[tokio::test] +async fn dispatched_mcp_approval_with_closed_response_is_classified_as_approval() { + use wiremock::matchers::method; + use wiremock::matchers::path; + + let server = wiremock::MockServer::start().await; + wiremock::Mock::given(method("POST")) + .and(path("/codex/analytics-events/events")) + .respond_with(wiremock::ResponseTemplate::new(200)) + .mount(&server) + .await; + + let (mut session, turn_context, rx_event) = make_session_and_context_with_rx().await; + let client = codex_analytics::AnalyticsEventsClient::new( + crate::test_support::auth_manager_from_auth( + codex_login::CodexAuth::create_dummy_chatgpt_auth_for_testing(), + ), + server.uri(), + /*analytics_enabled*/ Some(true), + ); + Arc::get_mut(&mut session) + .expect("session should be uniquely owned") + .services + .analytics_events_client = client.clone(); + *session.active_turn.lock().await = Some(ActiveTurn::default()); + + let call_id = "modern-missing-response"; + let thread_id = session.thread_id.to_string(); + let turn_id = turn_context.sub_id.clone(); + client.track_initialize( + /*connection_id*/ 1, + app_server_protocol::InitializeParams::default(), + "test-client".to_string(), + codex_analytics::AppServerRpcTransport::Stdio, + ); + let response = serde_json::from_value(serde_json::json!({ + "thread": { + "id": thread_id, "sessionId": thread_id, "preview": "", "ephemeral": false, + "modelProvider": "openai", "createdAt": 1, "updatedAt": 1, + "status": {"type": "idle"}, "cwd": &turn_context.config.cwd, + "cliVersion": "0.0.0", "source": "exec", "turns": [], + }, + "model": "test-model", "modelProvider": "openai", "cwd": &turn_context.config.cwd, + "approvalPolicy": "on-request", "approvalsReviewer": "user", + "sandbox": {"type": "dangerFullAccess"}, + })) + .expect("thread-start response should deserialize"); + client.track_response( + /*connection_id*/ 1, + app_server_protocol::RequestId::Integer(1), + &app_server_protocol::ClientResponsePayload::ThreadStart(response), + ); + client.track_notification(&app_server_protocol::ServerNotification::TurnStarted( + serde_json::from_value(serde_json::json!({ + "threadId": thread_id, + "turn": {"id": turn_id, "items": [], "status": "inProgress"}, + })) + .expect("turn-started notification should deserialize"), + )); + let mut item = serde_json::json!({ + "type": "mcpToolCall", "id": call_id, "server": "calendar", "tool": "send", + "status": "inProgress", "arguments": {}, + "appContext": {"connectorId": "calendar"}, + }); + client.track_notification(&app_server_protocol::ServerNotification::ItemStarted( + serde_json::from_value(serde_json::json!({ + "threadId": thread_id, "turnId": turn_id, "startedAtMs": 1, "item": item, + })) + .expect("item-started notification should deserialize"), + )); + + let action = ApprovalAction::McpToolCall { + id: call_id.to_string(), + server: "calendar".to_string(), + tool_name: "send".to_string(), + arguments: None, + connector_id: Some("calendar".to_string()), + connector_name: None, + connector_description: None, + connected_account_email: None, + tool_title: None, + tool_description: None, + annotations: None, + hook_tool_name: HookToolName::new("mcp__calendar__send"), + approval_policy: AskForApproval::OnRequest, + reviewer: ApprovalsReviewer::User, + approval_mode: AppToolApproval::Auto, + allow_session_remember: true, + allow_persistent_approval: true, + }; + let approval = tokio::spawn({ + let session = Arc::clone(&session); + let turn_context = Arc::clone(&turn_context); + async move { request_mcp_tool_user_approval(&session, &turn_context, call_id, &action).await } + }); + let event = tokio::time::timeout(std::time::Duration::from_secs(1), rx_event.recv()) + .await + .expect("approval request should be dispatched") + .expect("approval request should be received"); + assert!(matches!(event.msg, EventMsg::ElicitationRequest(_))); + *session.active_turn.lock().await = None; + assert_eq!( + approval.await.expect("approval task should complete"), + ReviewDecision::Abort + ); + + item["status"] = serde_json::json!("failed"); + client.track_notification(&app_server_protocol::ServerNotification::ItemCompleted( + serde_json::from_value(serde_json::json!({ + "threadId": thread_id, "turnId": turn_id, "completedAtMs": 2, "item": item, + })) + .expect("item-completed notification should deserialize"), + )); + client.flush().await; + let events = server + .received_requests() + .await + .expect("analytics requests should be recorded") + .into_iter() + .flat_map(|request| { + serde_json::from_slice::(&request.body) + .expect("analytics request should deserialize")["events"] + .as_array() + .expect("analytics events should be an array") + .clone() + }) + .collect::>(); + let mcp_events = events + .iter() + .filter(|event| { + event["event_type"] == "codex_mcp_tool_call_event" + && event["event_params"]["item_id"] == call_id + }) + .collect::>(); + assert_eq!(mcp_events.len(), 1); + assert_eq!( + mcp_events[0]["event_params"]["elicitation_type"], + "approval" + ); +} + #[tokio::test] async fn persist_codex_app_tool_approval_writes_tool_override() { let tmp = tempdir().expect("tempdir"); diff --git a/codex-rs/core/tests/common/Cargo.toml b/codex-rs/core/tests/common/Cargo.toml index ad4eab6088..5644a06e0a 100644 --- a/codex-rs/core/tests/common/Cargo.toml +++ b/codex-rs/core/tests/common/Cargo.toml @@ -15,6 +15,7 @@ workspace = true anyhow = { workspace = true } assert_cmd = { workspace = true } base64 = { workspace = true } +codex-analytics = { workspace = true } codex-arg0 = { workspace = true } codex-attachment-store = { workspace = true } codex-config = { workspace = true } diff --git a/codex-rs/core/tests/common/test_codex.rs b/codex-rs/core/tests/common/test_codex.rs index 43985f7f30..3c837e96d8 100644 --- a/codex-rs/core/tests/common/test_codex.rs +++ b/codex-rs/core/tests/common/test_codex.rs @@ -13,6 +13,7 @@ use std::time::Duration; use anyhow::Context; use anyhow::Result; use anyhow::anyhow; +use codex_analytics::AnalyticsEventsClient; use codex_attachment_store::AttachmentStore; use codex_config::CloudConfigBundleLoader; use codex_core::CodexThread; @@ -329,6 +330,7 @@ pub fn turn_permission_fields( pub struct TestCodexBuilder { config_mutators: Vec>, auth: CodexAuth, + analytics_events_client: Option, pre_build_hooks: Vec>, workspace_setups: Vec>, home: Option>, @@ -365,6 +367,14 @@ impl TestCodexBuilder { self } + pub fn with_analytics_events_client( + mut self, + analytics_events_client: AnalyticsEventsClient, + ) -> Self { + self.analytics_events_client = Some(analytics_events_client); + self + } + pub fn with_models_manager(mut self, models_manager: SharedModelsManager) -> Self { self.models_manager = Some(models_manager); self @@ -742,7 +752,7 @@ impl TestCodexBuilder { Arc::clone(&environment_manager), Arc::new(extensions.build()), user_instructions_provider, - /*analytics_events_client*/ None, + self.analytics_events_client.clone(), Arc::clone(&self.image_store), Arc::clone(&thread_store), codex_core::local_agent_graph_store_from_state_db(state_db.as_ref()), @@ -1393,6 +1403,7 @@ pub fn test_codex() -> TestCodexBuilder { .expect("test config should allow ShellSnapshot override"); })], auth: CodexAuth::from_api_key("dummy"), + analytics_events_client: None, pre_build_hooks: vec![], workspace_setups: vec![], home: None, diff --git a/codex-rs/core/tests/suite/mcp_auth_elicitation.rs b/codex-rs/core/tests/suite/mcp_auth_elicitation.rs index 590d068a31..c6024968fa 100644 --- a/codex-rs/core/tests/suite/mcp_auth_elicitation.rs +++ b/codex-rs/core/tests/suite/mcp_auth_elicitation.rs @@ -1,10 +1,22 @@ -#![allow(clippy::unwrap_used)] +//! Verify MCP auth prompts and elicitation analytics through actual turns. use anyhow::Result; +use codex_analytics::AnalyticsEventsClient; +use codex_analytics::AppServerRpcTransport; +use codex_app_server_protocol as app; use codex_core::TurnInputRequest; +use codex_core::TurnInputSubmission; use codex_core::config::Constrained; +use codex_extension_api::ExtensionRegistryBuilder; +use codex_extension_api::McpToolResultInput; +use codex_extension_api::ToolLifecycleContributor; +use codex_extension_api::ToolLifecycleFuture; +use codex_features::Feature; +use codex_login::CodexAuth; use codex_mcp::CODEX_APPS_MCP_SERVER_NAME; use codex_protocol::approvals::ElicitationRequest; +use codex_protocol::config_types::ApprovalsReviewer; +use codex_protocol::items::TurnItem; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::ElicitationAction; use codex_protocol::protocol::EventMsg; @@ -14,29 +26,31 @@ use core_test_support::PathExt; use core_test_support::apps_test_server::AppsTestServer; use core_test_support::apps_test_server::SEARCH_CALENDAR_CREATE_TOOL; use core_test_support::apps_test_server::SEARCH_CALENDAR_NAMESPACE; +use core_test_support::apps_test_server::recorded_apps_tool_call_by_call_id; +use core_test_support::apps_test_server::recorded_apps_tool_calls; use core_test_support::apps_test_server::search_capable_apps_builder; -use core_test_support::responses::ev_assistant_message; -use core_test_support::responses::ev_completed; -use core_test_support::responses::ev_function_call_with_namespace; -use core_test_support::responses::ev_response_created; -use core_test_support::responses::mount_sse_sequence; -use core_test_support::responses::sse; -use core_test_support::responses::start_mock_server; +use core_test_support::responses; use core_test_support::skip_if_no_network; -use core_test_support::wait_for_event; use pretty_assertions::assert_eq; use serde_json::Value; use serde_json::json; +use std::sync::Arc; +use std::time::Duration; +use test_case::test_case; use wiremock::Mock; use wiremock::Request; use wiremock::Respond; use wiremock::ResponseTemplate; use wiremock::matchers::body_partial_json; use wiremock::matchers::method; +use wiremock::matchers::path; use wiremock::matchers::path_regex; +const CALL_ID: &str = "calendar-elicitation-call"; +const PRIVATE_SENTINEL: &str = "synthetic-private-value"; + #[derive(Clone, Copy)] -struct AuthFailureResponder; +struct AuthFailureResponder(Scenario); impl Respond for AuthFailureResponder { fn respond(&self, request: &Request) -> ResponseTemplate { @@ -44,13 +58,13 @@ impl Respond for AuthFailureResponder { serde_json::from_slice(&request.body).expect("tools/call request should be valid JSON"); let id = body.get("id").cloned().unwrap_or(Value::Null); - ResponseTemplate::new(/*status*/ 200).set_body_json(json!({ + let mut response = json!({ "jsonrpc": "2.0", "id": id, "result": { "content": [{ "type": "text", - "text": "Connector reauthentication required", + "text": PRIVATE_SENTINEL, }], "isError": true, "_meta": { @@ -67,100 +81,303 @@ impl Respond for AuthFailureResponder { }, }, }, - })) + }); + if self.0 == Scenario::AuthMetadataRemoved { + response["result"]["_meta"] = json!({"_codex_apps": {"connector_auth_failure": { + "is_auth_failure": true, "connector_id": "calendar", + }}}); + } + ResponseTemplate::new(/*status*/ 200).set_body_json(response) } } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum Scenario { + DefaultAuth, + ModernAuth, + AuthMetadataRemoved, + LegacySuccess, + LegacyCancelled, + ModernDeclined, +} + +struct RemoveAuthMetadata; + +impl ToolLifecycleContributor for RemoveAuthMetadata { + fn on_mcp_tool_result<'a>(&'a self, input: McpToolResultInput<'a>) -> ToolLifecycleFuture<'a> { + Box::pin(async move { input.result.meta = None }) + } +} + +#[test_case(Scenario::DefaultAuth; "auth requests elicitation by default")] +#[test_case(Scenario::ModernAuth; "modern accepted auth")] +#[test_case(Scenario::AuthMetadataRemoved; "auth metadata removed by callback")] +#[test_case(Scenario::LegacySuccess; "legacy accepted ordinary")] +#[test_case(Scenario::LegacyCancelled; "legacy cancelled")] +#[test_case(Scenario::ModernDeclined; "modern declined")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn codex_apps_auth_failure_requests_elicitation_by_default() -> Result<()> { +async fn actual_turn_elicitation_analytics(scenario: Scenario) -> Result<()> { skip_if_no_network!(Ok(())); - let server = start_mock_server().await; - let apps_server = AppsTestServer::mount_searchable(&server).await?; + let modern = scenario != Scenario::LegacySuccess && scenario != Scenario::LegacyCancelled; + let expected_approvals = usize::from(scenario != Scenario::DefaultAuth); + let approved = scenario != Scenario::LegacyCancelled && scenario != Scenario::ModernDeclined; + let auth_failure = matches!( + scenario, + Scenario::DefaultAuth | Scenario::ModernAuth | Scenario::AuthMetadataRemoved + ); + let server = responses::start_mock_server().await; + AppsTestServer::mount_searchable(&server).await?; Mock::given(method("POST")) - .and(path_regex("^/api/codex/ps/mcp/?$")) - .and(body_partial_json(json!({ - "method": "tools/call", - "params": { - "name": "calendar_create_event", - }, - }))) - .respond_with(AuthFailureResponder) - .with_priority(/*p*/ 1) + .and(path("/codex/analytics-events/events")) + .respond_with(ResponseTemplate::new(/*status*/ 200)) .mount(&server) .await; - let call_id = "calendar-auth-call"; - let responses = mount_sse_sequence( + if auth_failure { + Mock::given(method("POST")) + .and(path_regex("^/api/codex/ps/mcp/?$")) + .and(body_partial_json(json!({ + "method": "tools/call", + "params": {"name": "calendar_create_event"}, + }))) + .respond_with(AuthFailureResponder(scenario)) + .with_priority(/*p*/ 1) + .mount(&server) + .await; + } + + let arguments = + json!({"title": PRIVATE_SENTINEL, "starts_at": "2026-06-18T12:00:00Z"}).to_string(); + let response_mock = responses::mount_sse_sequence( &server, vec![ - sse(vec![ - ev_response_created("resp-1"), - ev_function_call_with_namespace( - call_id, + responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_function_call_with_namespace( + CALL_ID, SEARCH_CALENDAR_NAMESPACE, SEARCH_CALENDAR_CREATE_TOOL, - &json!({ - "title": "Lunch", - "starts_at": "2026-06-18T12:00:00Z", - }) - .to_string(), + &arguments, ), - ev_completed("resp-1"), - ]), - sse(vec![ - ev_response_created("resp-2"), - ev_assistant_message("msg-1", "done"), - ev_completed("resp-2"), + responses::ev_completed("resp-1"), ]), + responses::sse_completed("resp-2"), ], ) .await; - let mut builder = - search_capable_apps_builder(apps_server.chatgpt_base_url).with_config(|config| { + let client = AnalyticsEventsClient::new( + codex_core::test_support::auth_manager_from_auth( + CodexAuth::create_dummy_chatgpt_auth_for_testing(), + ), + server.uri(), + /*analytics_enabled*/ Some(true), + ); + let mut builder = search_capable_apps_builder(server.uri()) + .with_analytics_events_client(client.clone()) + .with_config(move |config| { config.permissions.approval_policy = Constrained::allow_any(AskForApproval::OnRequest); + config.approvals_reviewer = ApprovalsReviewer::User; + if scenario != Scenario::DefaultAuth { + config + .features + .set_enabled(Feature::ToolCallMcpElicitation, modern) + .expect("approval feature should be configurable"); + } let user_config_path = config.codex_home.join("config.toml").abs(); - let user_config = toml::from_str( - r#" -[apps.calendar] -default_tools_approval_mode = "auto" + let approval_mode = if scenario == Scenario::DefaultAuth { + "auto" + } else { + "prompt" + }; + let user_config = toml::from_str(&format!( + r#"[apps.calendar] +default_tools_approval_mode = "{approval_mode}" +approvals_reviewer = "user" "#, - ) + )) .expect("apps config should parse"); config.config_layer_stack = config .config_layer_stack .with_user_config(&user_config_path, user_config) .expect("apps user config should be valid"); }); - let test = builder.build(&server).await?; + if scenario == Scenario::AuthMetadataRemoved { + let mut extensions = ExtensionRegistryBuilder::new(); + extensions.tool_lifecycle_contributor(Arc::new(RemoveAuthMetadata)); + builder = builder.with_extensions(Arc::new(extensions.build())); + } + let test = builder.build_with_auto_env(&server).await?; + let session = &test.session_configured; + let thread_id = session.thread_id.to_string(); + client.track_initialize( + /*connection_id*/ 1, + app::InitializeParams::default(), + "test-client".to_string(), + AppServerRpcTransport::Stdio, + ); + let thread_response = serde_json::from_value(json!({ + "thread": { + "id": thread_id, "sessionId": session.session_id.to_string(), "preview": "", + "ephemeral": false, "modelProvider": session.model_provider_id, + "createdAt": 1, "updatedAt": 1, "status": {"type": "idle"}, + "cwd": session.cwd, "cliVersion": "0.0.0", "source": "exec", "turns": [], + }, + "model": session.model, "modelProvider": session.model_provider_id, "cwd": session.cwd, + "approvalPolicy": app::AskForApproval::from(session.approval_policy), + "approvalsReviewer": app::ApprovalsReviewer::from(session.approvals_reviewer), + "sandbox": app::SandboxPolicy::from(test.config.legacy_sandbox_policy()), + }))?; + client.track_response( + /*connection_id*/ 1, + app::RequestId::Integer(1), + &app::ClientResponsePayload::ThreadStart(thread_response), + ); - test.codex + let submitted = test + .codex .start_or_steer_turn(TurnInputRequest::user_input(vec![UserInput::Text { text: "Use [$calendar](app://calendar) to create a calendar event.".to_string(), text_elements: Vec::new(), }])) .await?; - - let EventMsg::ElicitationRequest(request) = wait_for_event(&test.codex, |event| { - matches!( - event, - EventMsg::ElicitationRequest(_) | EventMsg::TurnComplete(_) - ) - }) - .await - else { - panic!("default auth elicitation should prompt before completing the turn"); + let TurnInputSubmission::Started { turn_id } = submitted else { + anyhow::bail!("expected a new turn, got {submitted:?}"); }; - assert_eq!(request.server_name, CODEX_APPS_MCP_SERVER_NAME); + let mut approvals_seen = 0; + let mut auth_request = None; + let mut target_lifecycle = [0; 2]; + tokio::time::timeout(Duration::from_secs(15), async { + loop { + let event = test.codex.next_event().await?; + let event_turn_id = event.id; + match event.msg { + EventMsg::TurnStarted(started) => { + assert_eq!(started.turn_id, turn_id); + client.track_notification(&app::ServerNotification::TurnStarted( + serde_json::from_value(json!({ + "threadId": thread_id, + "turn": {"id": started.turn_id, "items": [], + "status": "inProgress", "startedAt": started.started_at}, + }))?, + )); + } + message @ (EventMsg::ItemStarted(_) | EventMsg::ItemCompleted(_)) => { + let (item, item_thread_id, item_turn_id, completed) = match &message { + EventMsg::ItemStarted(event) => { + (&event.item, &event.thread_id, &event.turn_id, false) + } + EventMsg::ItemCompleted(event) => { + (&event.item, &event.thread_id, &event.turn_id, true) + } + _ => unreachable!("item guard guarantees item lifecycle event"), + }; + if let TurnItem::McpToolCall(item) = item + && item.id == CALL_ID + { + assert_eq!(event_turn_id, turn_id); + assert_eq!(item_thread_id, &session.thread_id); + assert_eq!(item_turn_id, &turn_id); + assert_eq!(item.server, CODEX_APPS_MCP_SERVER_NAME); + assert_eq!(item.connector_id.as_deref(), Some("calendar")); + target_lifecycle[usize::from(completed)] += 1; + } + let notification = + app::item_event_to_server_notification(message, &thread_id, &event_turn_id); + client.track_notification(¬ification); + } + EventMsg::RequestUserInput(request) => { + assert!(!modern, "legacy approval requested for {scenario:?}"); + assert!(recorded_apps_tool_calls(&server).await.is_empty()); + assert_eq!(request.turn_id, turn_id); + approvals_seen += 1; + let answer = if approved { "Allow" } else { "Cancel" }; + let response = serde_json::from_value(json!({"answers": { + (request.questions[0].id.clone()): { + "answers": [answer, format!("user_note: {PRIVATE_SENTINEL}")] + } + }}))?; + test.codex + .submit(Op::UserInputAnswer { + id: request.turn_id, + response, + }) + .await?; + } + EventMsg::ElicitationRequest(request) => { + assert_eq!(request.server_name, CODEX_APPS_MCP_SERVER_NAME); + let action = match &request.request { + ElicitationRequest::UserVerification { .. } => { + unreachable!("unexpected verification") + } + ElicitationRequest::Url { .. } => { + assert_eq!(approvals_seen, expected_approvals); + assert!(auth_request.is_none()); + assert_eq!(recorded_apps_tool_calls(&server).await.len(), 1); + assert_eq!( + request.id, + codex_protocol::mcp::RequestId::String(format!( + "codex_apps_auth_{CALL_ID}" + )) + ); + auth_request = Some(request.request.clone()); + ElicitationAction::Accept + } + ElicitationRequest::Form { .. } + | ElicitationRequest::OpenAiForm { .. } + | ElicitationRequest::OpenAiElicitationForm { .. } => { + assert!(modern, "modern approval requested for {scenario:?}"); + assert!(recorded_apps_tool_calls(&server).await.is_empty()); + approvals_seen += 1; + if approved { + ElicitationAction::Accept + } else { + ElicitationAction::Decline + } + } + }; + test.codex + .submit(Op::ResolveElicitation { + server_name: request.server_name, + request_id: request.id, + decision: action, + content: None, + meta: None, + }) + .await?; + } + EventMsg::TurnComplete(completed) => { + assert_eq!(completed.turn_id, turn_id); + client.track_notification(&app::ServerNotification::TurnCompleted( + serde_json::from_value(json!({ + "threadId": thread_id, + "turn": {"id": completed.turn_id, "items": [], "status": "completed", + "startedAt": completed.started_at, + "completedAt": completed.completed_at, + "durationMs": completed.duration_ms}, + }))?, + )); + break; + } + EventMsg::Error(error) => { + anyhow::bail!("unexpected turn error for {scenario:?}: {error:?}"); + } + _ => {} + } + } + Ok::<(), anyhow::Error>(()) + }) + .await??; + assert_eq!( - request.id, - codex_protocol::mcp::RequestId::String(format!("codex_apps_auth_{call_id}")) + (approvals_seen, target_lifecycle), + (expected_approvals, [1, 1]) ); - assert_eq!( - request.request, - ElicitationRequest::Url { + let expected_auth_request = if matches!(scenario, Scenario::DefaultAuth | Scenario::ModernAuth) + { + Some(ElicitationRequest::Url { meta: Some(json!({ "_codex_apps": { "connector_auth_failure": { @@ -179,37 +396,92 @@ default_tools_approval_mode = "auto" message: "Reconnect Calendar on ChatGPT to restore access for this request." .to_string(), url: "https://chatgpt.com/apps/calendar/calendar".to_string(), - elicitation_id: format!("codex_apps_auth_{call_id}"), + elicitation_id: format!("codex_apps_auth_{CALL_ID}"), + }) + } else { + None + }; + assert_eq!(auth_request, expected_auth_request); + let tool_calls = recorded_apps_tool_calls(&server).await; + assert_eq!(tool_calls.len(), usize::from(approved)); + if approved { + let request = recorded_apps_tool_call_by_call_id(&server, CALL_ID).await; + assert_eq!( + request.pointer("/params/name").and_then(Value::as_str), + Some("calendar_create_event") + ); + } + let requests = response_mock.requests(); + assert_eq!(requests.len(), 2); + if matches!(scenario, Scenario::DefaultAuth | Scenario::ModernAuth) { + let output = requests[1].function_call_output(CALL_ID); + let items = output["output"] + .as_array() + .expect("auth elicitation result should contain content items"); + assert_eq!( + &items[1..], + &[json!({ + "type": "input_text", + "text": "Authentication for Calendar was requested and accepted. Retry this tool call now.", + })] + ); + assert!(!output.to_string().contains(PRIVATE_SENTINEL)); + } + + tokio::time::timeout(Duration::from_secs(10), client.flush()).await?; + let mut events = Vec::new(); + for request in server.received_requests().await.unwrap_or_default() { + if request.method != "POST" || request.url.path() != "/codex/analytics-events/events" { + continue; + } + let payload: Value = serde_json::from_slice(&request.body)?; + let batch = payload["events"].as_array().expect("analytics events"); + events.extend(batch.iter().cloned()); + } + let mcp_events = events + .iter() + .filter(|event| { + event["event_type"] == "codex_mcp_tool_call_event" + && event["event_params"]["item_id"] == CALL_ID + }) + .collect::>(); + let app_used_events = events + .iter() + .filter(|event| { + event["event_type"] == "codex_app_used" + && event["event_params"]["connector_id"] == "calendar" + && event["event_params"]["turn_id"] == turn_id + }) + .collect::>(); + assert_eq!(mcp_events.len(), 1); + assert_eq!(app_used_events.len(), usize::from(approved)); + assert_eq!(mcp_events[0]["event_params"]["connector_id"], "calendar"); + assert_eq!(mcp_events[0]["event_params"]["thread_id"], thread_id); + assert_eq!(mcp_events[0]["event_params"]["turn_id"], turn_id); + assert_eq!( + mcp_events[0]["event_params"]["terminal_status"], + if scenario == Scenario::LegacySuccess { + "completed" + } else { + "failed" } ); - - test.codex - .submit(Op::ResolveElicitation { - server_name: request.server_name, - request_id: request.id, - decision: ElicitationAction::Accept, - content: None, - meta: None, - }) - .await?; - wait_for_event(&test.codex, |event| { - matches!(event, EventMsg::TurnComplete(_)) - }) - .await; - - let requests = responses.requests(); - assert_eq!(requests.len(), 2); - let output = requests[1].function_call_output(call_id); - let items = output["output"] - .as_array() - .expect("auth elicitation result should contain content items"); - assert_eq!( - &items[1..], - &[json!({ - "type": "input_text", - "text": "Authentication for Calendar was requested and accepted. Retry this tool call now.", - })] - ); + let expected = match scenario { + Scenario::DefaultAuth | Scenario::ModernAuth | Scenario::AuthMetadataRemoved => { + json!("auth_or_link") + } + Scenario::LegacySuccess => Value::Null, + Scenario::LegacyCancelled | Scenario::ModernDeclined => json!("approval"), + }; + for event in mcp_events.iter().chain(app_used_events.iter()) { + assert_eq!( + event["event_params"].get("elicitation_type"), + Some(&expected) + ); + } + let target_payload = serde_json::to_string(&(mcp_events, app_used_events))?; + assert!(!target_payload.contains(PRIVATE_SENTINEL)); + assert!(!target_payload.contains("https://chatgpt.com/apps/")); Ok(()) }