diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index fec7c24ec7..dbe04d8ff3 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -10,6 +10,8 @@ use crate::events::CodexHookRunEventRequest; use crate::events::CodexPluginEventRequest; use crate::events::CodexPluginUsedEventRequest; use crate::events::CodexRuntimeMetadata; +use crate::events::CodexToolCallReviewEventParams; +use crate::events::CodexToolCallReviewEventRequest; use crate::events::CodexToolItemEventBase; use crate::events::CodexTurnEventRequest; use crate::events::CommandExecutionFamily; @@ -24,6 +26,10 @@ use crate::events::ThreadInitializedEvent; use crate::events::ThreadInitializedEventParams; use crate::events::ToolItemFinalApprovalOutcome; use crate::events::ToolItemTerminalStatus; +use crate::events::ToolReviewReviewer; +use crate::events::ToolReviewStatus; +use crate::events::ToolReviewToolKind; +use crate::events::ToolReviewTrigger; use crate::events::TrackEventRequest; use crate::events::codex_app_metadata; use crate::events::codex_hook_run_metadata; @@ -62,23 +68,45 @@ use crate::facts::TurnTokenUsageFact; use crate::reducer::AnalyticsReducer; use crate::reducer::normalize_path_for_skill_id; use crate::reducer::skill_id_for_local_skill; +use codex_app_server_protocol::AdditionalNetworkPermissions; use codex_app_server_protocol::ApprovalsReviewer as AppServerApprovalsReviewer; use codex_app_server_protocol::AskForApproval as AppServerAskForApproval; use codex_app_server_protocol::ClientInfo; use codex_app_server_protocol::ClientRequest; use codex_app_server_protocol::ClientResponsePayload; use codex_app_server_protocol::CodexErrorInfo; +use codex_app_server_protocol::CommandExecutionApprovalDecision; +use codex_app_server_protocol::CommandExecutionRequestApprovalParams; +use codex_app_server_protocol::CommandExecutionRequestApprovalResponse; use codex_app_server_protocol::CommandExecutionSource as AppServerCommandExecutionSource; use codex_app_server_protocol::CommandExecutionStatus; +use codex_app_server_protocol::FileChangeApprovalDecision; +use codex_app_server_protocol::FileChangeRequestApprovalParams; +use codex_app_server_protocol::FileChangeRequestApprovalResponse; +use codex_app_server_protocol::GrantedPermissionProfile; +use codex_app_server_protocol::GuardianApprovalReview; +use codex_app_server_protocol::GuardianApprovalReviewAction; +use codex_app_server_protocol::GuardianApprovalReviewStatus; +use codex_app_server_protocol::GuardianCommandSource as AppServerGuardianCommandSource; use codex_app_server_protocol::InitializeCapabilities; use codex_app_server_protocol::InitializeParams; use codex_app_server_protocol::ItemCompletedNotification; +use codex_app_server_protocol::ItemGuardianApprovalReviewCompletedNotification; +use codex_app_server_protocol::ItemGuardianApprovalReviewStartedNotification; use codex_app_server_protocol::ItemStartedNotification; use codex_app_server_protocol::JSONRPCErrorError; +use codex_app_server_protocol::NetworkPolicyAmendment; +use codex_app_server_protocol::NetworkPolicyRuleAction; use codex_app_server_protocol::NonSteerableTurnKind; +use codex_app_server_protocol::PermissionGrantScope; +use codex_app_server_protocol::PermissionsRequestApprovalParams; +use codex_app_server_protocol::PermissionsRequestApprovalResponse; use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::RequestPermissionProfile; use codex_app_server_protocol::SandboxPolicy as AppServerSandboxPolicy; use codex_app_server_protocol::ServerNotification; +use codex_app_server_protocol::ServerRequest; +use codex_app_server_protocol::ServerResponse; use codex_app_server_protocol::SessionSource as AppServerSessionSource; use codex_app_server_protocol::Thread; use codex_app_server_protocol::ThreadArchiveParams; @@ -663,6 +691,171 @@ fn sample_command_execution_item( } } +fn sample_command_approval_request(request_id: i64, approval_id: Option<&str>) -> ServerRequest { + ServerRequest::CommandExecutionRequestApproval { + request_id: RequestId::Integer(request_id), + params: CommandExecutionRequestApprovalParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item_id: "item-1".to_string(), + approval_id: approval_id.map(str::to_string), + reason: None, + network_approval_context: None, + command: Some("echo hi".to_string()), + cwd: None, + command_actions: None, + additional_permissions: None, + proposed_execpolicy_amendment: None, + proposed_network_policy_amendments: None, + available_decisions: None, + }, + } +} + +fn sample_command_approval_response( + request_id: i64, + decision: CommandExecutionApprovalDecision, +) -> ServerResponse { + ServerResponse::CommandExecutionRequestApproval { + request_id: RequestId::Integer(request_id), + response: CommandExecutionRequestApprovalResponse { decision }, + } +} + +fn sample_network_policy_command_approval_request(request_id: i64) -> ServerRequest { + ServerRequest::CommandExecutionRequestApproval { + request_id: RequestId::Integer(request_id), + params: CommandExecutionRequestApprovalParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item_id: "item-1".to_string(), + approval_id: None, + reason: Some("network access requested".to_string()), + network_approval_context: None, + command: None, + cwd: None, + command_actions: None, + additional_permissions: None, + proposed_execpolicy_amendment: None, + proposed_network_policy_amendments: Some(vec![NetworkPolicyAmendment { + host: "example.com".to_string(), + action: NetworkPolicyRuleAction::Allow, + }]), + available_decisions: None, + }, + } +} + +fn sample_file_change_approval_request(request_id: i64) -> ServerRequest { + ServerRequest::FileChangeRequestApproval { + request_id: RequestId::Integer(request_id), + params: FileChangeRequestApprovalParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item_id: "item-1".to_string(), + reason: None, + grant_root: Some(PathBuf::from("/tmp")), + }, + } +} + +fn sample_file_change_approval_response( + request_id: i64, + decision: FileChangeApprovalDecision, +) -> ServerResponse { + ServerResponse::FileChangeRequestApproval { + request_id: RequestId::Integer(request_id), + response: FileChangeRequestApprovalResponse { decision }, + } +} + +fn sample_permissions_approval_request(request_id: i64) -> ServerRequest { + ServerRequest::PermissionsRequestApproval { + request_id: RequestId::Integer(request_id), + params: PermissionsRequestApprovalParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item_id: "item-1".to_string(), + cwd: test_path_buf("/tmp").abs(), + reason: None, + permissions: RequestPermissionProfile { + network: Some(AdditionalNetworkPermissions { + enabled: Some(true), + }), + file_system: None, + }, + }, + } +} + +fn sample_permissions_approval_response(request_id: i64) -> ServerResponse { + ServerResponse::PermissionsRequestApproval { + request_id: RequestId::Integer(request_id), + response: PermissionsRequestApprovalResponse { + permissions: GrantedPermissionProfile { + network: Some(AdditionalNetworkPermissions { + enabled: Some(true), + }), + file_system: None, + }, + scope: PermissionGrantScope::Session, + strict_auto_review: None, + }, + } +} + +fn sample_guardian_review_started( + review_id: &str, + target_item_id: Option<&str>, +) -> ServerNotification { + ServerNotification::ItemGuardianApprovalReviewStarted( + ItemGuardianApprovalReviewStartedNotification { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + review_id: review_id.to_string(), + target_item_id: target_item_id.map(str::to_string), + review: GuardianApprovalReview { + status: GuardianApprovalReviewStatus::InProgress, + risk_level: None, + user_authorization: None, + rationale: None, + }, + action: GuardianApprovalReviewAction::Command { + source: AppServerGuardianCommandSource::Shell, + command: "echo hi".to_string(), + cwd: test_path_buf("/tmp").abs(), + }, + }, + ) +} + +fn sample_guardian_review_completed( + review_id: &str, + target_item_id: Option<&str>, + status: GuardianApprovalReviewStatus, +) -> ServerNotification { + ServerNotification::ItemGuardianApprovalReviewCompleted( + ItemGuardianApprovalReviewCompletedNotification { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + review_id: review_id.to_string(), + target_item_id: target_item_id.map(str::to_string), + decision_source: codex_app_server_protocol::AutoReviewDecisionSource::Agent, + review: GuardianApprovalReview { + status, + risk_level: None, + user_authorization: None, + rationale: None, + }, + action: GuardianApprovalReviewAction::Command { + source: AppServerGuardianCommandSource::Shell, + command: "echo hi".to_string(), + cwd: test_path_buf("/tmp").abs(), + }, + }, + ) +} + fn expected_absolute_path(path: &PathBuf) -> String { std::fs::canonicalize(path) .unwrap_or_else(|_| path.to_path_buf()) @@ -1061,6 +1254,54 @@ fn command_execution_event_serializes_expected_shape() { ); } +#[test] +fn tool_call_review_event_serializes_expected_shape() { + let event = TrackEventRequest::ToolCallReview(CodexToolCallReviewEventRequest { + event_type: "codex_tool_call_review_event", + event_params: CodexToolCallReviewEventParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item_id: None, + review_id: "review-1".to_string(), + thread_source: Some("subagent"), + subagent_source: Some("thread_spawn".to_string()), + parent_thread_id: Some("parent-thread-1".to_string()), + tool_kind: ToolReviewToolKind::NetworkAccess, + tool_name: "network_access".to_string(), + reviewer: ToolReviewReviewer::User, + trigger: ToolReviewTrigger::NetworkRetry, + status: ToolReviewStatus::NetworkPolicyAllow, + created_at: 123, + completed_at: Some(125), + duration_ms: Some(2000), + }, + }); + + let payload = serde_json::to_value(&event).expect("serialize tool review event"); + assert_eq!( + payload, + json!({ + "event_type": "codex_tool_call_review_event", + "event_params": { + "thread_id": "thread-1", + "turn_id": "turn-1", + "item_id": null, + "review_id": "review-1", + "thread_source": "subagent", + "subagent_source": "thread_spawn", + "parent_thread_id": "parent-thread-1", + "tool_kind": "network_access", + "tool_name": "network_access", + "reviewer": "user", + "trigger": "network_retry", + "status": "network_policy_allow", + "created_at": 123, + "completed_at": 125, + "duration_ms": 2000 + } + }) + ); +} #[tokio::test] async fn initialize_caches_client_and_thread_lifecycle_publishes_once_initialized() { let mut reducer = AnalyticsReducer::default(); @@ -1550,6 +1791,337 @@ async fn item_lifecycle_notifications_publish_command_execution_event() { assert_eq!(payload[0]["event_params"]["thread_source"], "user"); } +#[tokio::test] +async fn command_execution_approval_response_publishes_user_review_event() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_command_approval_request( + /*request_id*/ 41, /*approval_id*/ None, + )), + }, + &mut events, + ) + .await; + assert!(events.is_empty()); + + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_command_approval_response( + /*request_id*/ 41, + CommandExecutionApprovalDecision::Accept, + )), + }, + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events).expect("serialize events"); + assert_eq!(payload.as_array().expect("events array").len(), 1); + assert_eq!(payload[0]["event_type"], "codex_tool_call_review_event"); + assert_eq!(payload[0]["event_params"]["thread_id"], "thread-1"); + assert_eq!(payload[0]["event_params"]["turn_id"], "turn-1"); + assert_eq!(payload[0]["event_params"]["item_id"], "item-1"); + assert_eq!(payload[0]["event_params"]["review_id"], "user:41"); + assert_eq!(payload[0]["event_params"]["thread_source"], "user"); + assert_eq!(payload[0]["event_params"]["tool_kind"], "command_execution"); + assert_eq!(payload[0]["event_params"]["tool_name"], "shell"); + assert_eq!(payload[0]["event_params"]["reviewer"], "user"); + assert_eq!(payload[0]["event_params"]["trigger"], "initial"); + assert_eq!(payload[0]["event_params"]["status"], "approved"); + assert!(payload[0]["event_params"]["created_at"].as_u64().is_some()); + assert!( + payload[0]["event_params"]["completed_at"] + .as_u64() + .is_some() + ); + assert!(payload[0]["event_params"]["duration_ms"].as_u64().is_some()); +} + +#[tokio::test] +async fn command_subapproval_maps_to_subcommand_execve_trigger() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_command_approval_request( + /*request_id*/ 42, + Some("approval-1"), + )), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_command_approval_response( + /*request_id*/ 42, + CommandExecutionApprovalDecision::AcceptForSession, + )), + }, + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_params"]["trigger"], "subcommand_execve"); + assert_eq!(payload["event_params"]["status"], "approved_for_session"); +} + +#[tokio::test] +async fn network_policy_amendment_maps_to_network_review_status() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_network_policy_command_approval_request( + /*request_id*/ 43, + )), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_command_approval_response( + /*request_id*/ 43, + CommandExecutionApprovalDecision::ApplyNetworkPolicyAmendment { + network_policy_amendment: NetworkPolicyAmendment { + host: "example.com".to_string(), + action: NetworkPolicyRuleAction::Allow, + }, + }, + )), + }, + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_params"]["trigger"], "network_retry"); + assert_eq!(payload["event_params"]["status"], "network_policy_allow"); +} + +#[tokio::test] +async fn file_change_approval_response_maps_terminal_statuses() { + for (request_id, decision, expected_status) in [ + (51, FileChangeApprovalDecision::Accept, "approved"), + (52, FileChangeApprovalDecision::Decline, "denied"), + (53, FileChangeApprovalDecision::Cancel, "aborted"), + ] { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_file_change_approval_request(request_id)), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_file_change_approval_response(request_id, decision)), + }, + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_params"]["tool_kind"], "file_change"); + assert_eq!(payload["event_params"]["tool_name"], "apply_patch"); + assert_eq!(payload["event_params"]["trigger"], "sandbox_retry"); + assert_eq!(payload["event_params"]["status"], expected_status); + } +} + +#[tokio::test] +async fn permissions_approval_response_maps_to_network_trigger() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_permissions_approval_request(/*request_id*/ 61)), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_permissions_approval_response(/*request_id*/ 61)), + }, + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_params"]["tool_kind"], "permissions"); + assert_eq!(payload["event_params"]["trigger"], "network_retry"); + assert_eq!(payload["event_params"]["status"], "approved_for_session"); +} + +#[tokio::test] +async fn guardian_completed_notification_publishes_review_event_with_thread_metadata() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::Notification(Box::new(sample_guardian_review_started( + "guardian-review-1", + Some("item-1"), + ))), + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::Notification(Box::new(sample_guardian_review_completed( + "guardian-review-1", + Some("item-1"), + GuardianApprovalReviewStatus::Denied, + ))), + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_type"], "codex_tool_call_review_event"); + assert_eq!(payload["event_params"]["review_id"], "guardian-review-1"); + assert_eq!(payload["event_params"]["item_id"], "item-1"); + assert_eq!(payload["event_params"]["thread_source"], "user"); + assert_eq!(payload["event_params"]["tool_kind"], "command_execution"); + assert_eq!(payload["event_params"]["reviewer"], "guardian"); + assert_eq!(payload["event_params"]["status"], "denied"); + assert!(payload["event_params"]["duration_ms"].as_u64().is_some()); +} + +#[tokio::test] +async fn guardian_completed_without_started_uses_null_duration() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::Notification(Box::new(sample_guardian_review_completed( + "guardian-review-2", + /*target_item_id*/ None, + GuardianApprovalReviewStatus::TimedOut, + ))), + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize review event"); + assert_eq!(payload["event_params"]["item_id"], json!(null)); + assert_eq!(payload["event_params"]["status"], "timed_out"); + assert_eq!(payload["event_params"]["duration_ms"], json!(null)); +} + +#[tokio::test] +async fn terminal_reviews_denormalize_counts_onto_tool_item_events() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_tool_review_prerequisites(&mut reducer, &mut events).await; + reducer + .ingest( + AnalyticsFact::ServerRequest { + connection_id: 7, + request: Box::new(sample_command_approval_request( + /*request_id*/ 71, /*approval_id*/ None, + )), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::ServerResponse { + response: Box::new(sample_command_approval_response( + /*request_id*/ 71, + CommandExecutionApprovalDecision::AcceptForSession, + )), + }, + &mut events, + ) + .await; + events.clear(); + + reducer + .ingest( + AnalyticsFact::Notification(Box::new(ServerNotification::ItemStarted( + ItemStartedNotification { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item: sample_command_execution_item( + CommandExecutionStatus::InProgress, + /*exit_code*/ None, + Some(1_000), + /*completed_at_ms*/ None, + /*duration_ms*/ None, + ), + }, + ))), + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::Notification(Box::new(ServerNotification::ItemCompleted( + ItemCompletedNotification { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + item: sample_command_execution_item( + CommandExecutionStatus::Completed, + Some(0), + Some(1_000), + Some(1_042), + Some(42), + ), + }, + ))), + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events[0]).expect("serialize tool item event"); + assert_eq!(payload["event_params"]["review_count"], 1); + assert_eq!(payload["event_params"]["user_review_count"], 1); + assert_eq!(payload["event_params"]["guardian_review_count"], 0); + assert_eq!( + payload["event_params"]["final_approval_outcome"], + "user_approved_for_session" + ); +} + #[test] fn subagent_thread_started_review_serializes_expected_shape() { let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( diff --git a/codex-rs/analytics/src/reducer.rs b/codex-rs/analytics/src/reducer.rs index 99d411b18c..9ebaa44a4e 100644 --- a/codex-rs/analytics/src/reducer.rs +++ b/codex-rs/analytics/src/reducer.rs @@ -19,6 +19,8 @@ use crate::events::CodexMcpToolCallEventRequest; use crate::events::CodexPluginEventRequest; use crate::events::CodexPluginUsedEventRequest; use crate::events::CodexRuntimeMetadata; +use crate::events::CodexToolCallReviewEventParams; +use crate::events::CodexToolCallReviewEventRequest; use crate::events::CodexToolItemEventBase; use crate::events::CodexTurnEventParams; use crate::events::CodexTurnEventRequest; @@ -38,6 +40,10 @@ use crate::events::ThreadInitializedEventParams; use crate::events::ToolItemFailureKind; use crate::events::ToolItemFinalApprovalOutcome; use crate::events::ToolItemTerminalStatus; +use crate::events::ToolReviewReviewer; +use crate::events::ToolReviewStatus; +use crate::events::ToolReviewToolKind; +use crate::events::ToolReviewTrigger; use crate::events::TrackEventRequest; use crate::events::WebSearchActionKind; use crate::events::codex_app_metadata; @@ -77,16 +83,26 @@ use codex_app_server_protocol::CodexErrorInfo; use codex_app_server_protocol::CollabAgentStatus; use codex_app_server_protocol::CollabAgentTool; use codex_app_server_protocol::CollabAgentToolCallStatus; +use codex_app_server_protocol::CommandExecutionApprovalDecision; use codex_app_server_protocol::CommandExecutionSource as AppServerCommandExecutionSource; use codex_app_server_protocol::CommandExecutionStatus; use codex_app_server_protocol::DynamicToolCallOutputContentItem; use codex_app_server_protocol::DynamicToolCallStatus; +use codex_app_server_protocol::FileChangeApprovalDecision; +use codex_app_server_protocol::GuardianApprovalReviewAction; +use codex_app_server_protocol::GuardianApprovalReviewStatus; +use codex_app_server_protocol::GuardianCommandSource as AppServerGuardianCommandSource; use codex_app_server_protocol::InitializeParams; use codex_app_server_protocol::McpToolCallStatus; +use codex_app_server_protocol::NetworkPolicyRuleAction; use codex_app_server_protocol::PatchApplyStatus; use codex_app_server_protocol::PatchChangeKind; +use codex_app_server_protocol::PermissionGrantScope; use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::RequestPermissionProfile; use codex_app_server_protocol::ServerNotification; +use codex_app_server_protocol::ServerRequest; +use codex_app_server_protocol::ServerResponse; use codex_app_server_protocol::ThreadItem; use codex_app_server_protocol::TurnSteerResponse; use codex_app_server_protocol::UserInput; @@ -111,6 +127,9 @@ pub(crate) struct AnalyticsReducer { turns: HashMap, connections: HashMap, threads: HashMap, + tool_review_requests: HashMap, + guardian_review_starts: HashMap, + tool_review_summaries: HashMap, } struct ConnectionState { @@ -194,6 +213,35 @@ enum MissingAnalyticsContext { ThreadMetadata, } +#[derive(Clone)] +struct PendingToolReviewState { + thread_id: String, + turn_id: String, + item_id: Option, + review_id: String, + tool_kind: ToolReviewToolKind, + tool_name: String, + trigger: ToolReviewTrigger, + created_at: u64, + duration_started: bool, + requested_additional_permissions: bool, + requested_network_access: bool, +} + +struct PendingGuardianReviewState { + created_at: u64, +} + +#[derive(Clone, Default)] +struct ToolReviewSummary { + review_count: u64, + guardian_review_count: u64, + user_review_count: u64, + final_approval_outcome: Option, + requested_additional_permissions: bool, + requested_network_access: bool, +} + #[derive(Clone)] struct ThreadMetadataState { thread_source: Option<&'static str>, @@ -311,12 +359,14 @@ impl AnalyticsReducer { self.ingest_notification(*notification, out); } AnalyticsFact::ServerRequest { - connection_id: _connection_id, - request: _request, - } => {} - AnalyticsFact::ServerResponse { - response: _response, - } => {} + connection_id, + request, + } => { + self.ingest_server_request(connection_id, *request); + } + AnalyticsFact::ServerResponse { response } => { + self.ingest_server_response(*response, out); + } AnalyticsFact::Custom(input) => match input { CustomAnalyticsFact::SubAgentThreadStarted(input) => { self.ingest_subagent_thread_started(input, out); @@ -679,6 +729,151 @@ impl AnalyticsReducer { } } + fn ingest_server_request(&mut self, _connection_id: u64, request: ServerRequest) { + match request { + ServerRequest::CommandExecutionRequestApproval { request_id, params } => { + let requested_network_access = params.network_approval_context.is_some() + || params + .proposed_network_policy_amendments + .as_ref() + .is_some_and(|amendments| !amendments.is_empty()) + || params + .additional_permissions + .as_ref() + .and_then(|permissions| permissions.network.as_ref()) + .and_then(|network| network.enabled) + .unwrap_or(false); + let requested_additional_permissions = params.additional_permissions.is_some() + || params.proposed_execpolicy_amendment.is_some(); + let trigger = if params.approval_id.is_some() { + ToolReviewTrigger::SubcommandExecve + } else if requested_network_access { + ToolReviewTrigger::NetworkRetry + } else if requested_additional_permissions { + ToolReviewTrigger::SandboxRetry + } else { + ToolReviewTrigger::Initial + }; + self.tool_review_requests.insert( + request_id.clone(), + PendingToolReviewState { + thread_id: params.thread_id, + turn_id: params.turn_id, + item_id: Some(params.item_id), + review_id: user_review_id(&request_id), + tool_kind: ToolReviewToolKind::CommandExecution, + tool_name: "shell".to_string(), + trigger, + created_at: now_unix_seconds(), + duration_started: true, + requested_additional_permissions, + requested_network_access, + }, + ); + } + ServerRequest::FileChangeRequestApproval { request_id, params } => { + let requested_additional_permissions = params.grant_root.is_some(); + self.tool_review_requests.insert( + request_id.clone(), + PendingToolReviewState { + thread_id: params.thread_id, + turn_id: params.turn_id, + item_id: Some(params.item_id), + review_id: user_review_id(&request_id), + tool_kind: ToolReviewToolKind::FileChange, + tool_name: "apply_patch".to_string(), + trigger: if requested_additional_permissions { + ToolReviewTrigger::SandboxRetry + } else { + ToolReviewTrigger::Initial + }, + created_at: now_unix_seconds(), + duration_started: true, + requested_additional_permissions, + requested_network_access: false, + }, + ); + } + ServerRequest::PermissionsRequestApproval { request_id, params } => { + let requested_network_access = params + .permissions + .network + .as_ref() + .and_then(|network| network.enabled) + .unwrap_or(false); + let requested_additional_permissions = + requested_network_access || params.permissions.file_system.is_some(); + let trigger = if requested_network_access { + ToolReviewTrigger::NetworkRetry + } else if requested_additional_permissions { + ToolReviewTrigger::SandboxRetry + } else { + ToolReviewTrigger::Initial + }; + self.tool_review_requests.insert( + request_id.clone(), + PendingToolReviewState { + thread_id: params.thread_id, + turn_id: params.turn_id, + item_id: Some(params.item_id), + review_id: user_review_id(&request_id), + tool_kind: ToolReviewToolKind::Permissions, + tool_name: "permissions".to_string(), + trigger, + created_at: now_unix_seconds(), + duration_started: true, + requested_additional_permissions, + requested_network_access, + }, + ); + } + _ => {} + } + } + + fn ingest_server_response( + &mut self, + response: ServerResponse, + out: &mut Vec, + ) { + match response { + ServerResponse::CommandExecutionRequestApproval { + request_id, + response, + } => { + let Some(pending_review) = self.tool_review_requests.remove(&request_id) else { + return; + }; + let status = command_execution_review_status(response.decision); + self.emit_tool_review_event(pending_review, ToolReviewReviewer::User, status, out); + } + ServerResponse::FileChangeRequestApproval { + request_id, + response, + } => { + let Some(pending_review) = self.tool_review_requests.remove(&request_id) else { + return; + }; + let status = file_change_review_status(response.decision); + self.emit_tool_review_event(pending_review, ToolReviewReviewer::User, status, out); + } + ServerResponse::PermissionsRequestApproval { + request_id, + response, + } => { + let Some(pending_review) = self.tool_review_requests.remove(&request_id) else { + return; + }; + let status = match response.scope { + PermissionGrantScope::Turn => ToolReviewStatus::Approved, + PermissionGrantScope::Session => ToolReviewStatus::ApprovedForSession, + }; + self.emit_tool_review_event(pending_review, ToolReviewReviewer::User, status, out); + } + _ => {} + } + } + fn ingest_error_response( &mut self, connection_id: u64, @@ -744,15 +939,28 @@ impl AnalyticsReducer { else { return; }; - if let Some(event) = tool_item_event( - ¬ification.thread_id, - ¬ification.turn_id, - ¬ification.item, + if let Some(event) = tool_item_event(ToolItemEventInput { + thread_id: ¬ification.thread_id, + turn_id: ¬ification.turn_id, + item: ¬ification.item, connection_state, thread_metadata, - ) { + review_summary: self.tool_review_summaries.get(item_id), + }) { out.push(event); } + self.tool_review_summaries.remove(item_id); + } + ServerNotification::ItemGuardianApprovalReviewStarted(notification) => { + self.guardian_review_starts.insert( + notification.review_id, + PendingGuardianReviewState { + created_at: now_unix_seconds(), + }, + ); + } + ServerNotification::ItemGuardianApprovalReviewCompleted(notification) => { + self.ingest_guardian_tool_review_completed(notification, out); } ServerNotification::TurnStarted(notification) => { let turn_state = self.turns.entry(notification.turn.id).or_insert(TurnState { @@ -869,6 +1077,47 @@ impl AnalyticsReducer { ))); } + fn ingest_guardian_tool_review_completed( + &mut self, + notification: codex_app_server_protocol::ItemGuardianApprovalReviewCompletedNotification, + out: &mut Vec, + ) { + let completed_at = now_unix_seconds(); + let started = self.guardian_review_starts.remove(¬ification.review_id); + let created_at = started + .as_ref() + .map(|start| start.created_at) + .unwrap_or(completed_at); + let Some(status) = guardian_review_status(notification.review.status) else { + return; + }; + let (tool_kind, tool_name, trigger) = guardian_review_tool_metadata(¬ification.action); + let pending_review = PendingToolReviewState { + thread_id: notification.thread_id, + turn_id: notification.turn_id, + item_id: notification.target_item_id, + review_id: notification.review_id, + tool_kind, + tool_name, + trigger, + created_at, + duration_started: started.is_some(), + requested_additional_permissions: guardian_review_requested_additional_permissions( + ¬ification.action, + ), + requested_network_access: guardian_review_requested_network_access( + ¬ification.action, + ), + }; + self.emit_tool_review_event_at( + pending_review, + ToolReviewReviewer::Guardian, + status, + completed_at, + out, + ); + } + fn ingest_turn_steer_response( &mut self, connection_id: u64, @@ -934,6 +1183,84 @@ impl AnalyticsReducer { })); } + fn emit_tool_review_event( + &mut self, + pending_review: PendingToolReviewState, + reviewer: ToolReviewReviewer, + status: ToolReviewStatus, + out: &mut Vec, + ) { + self.emit_tool_review_event_at(pending_review, reviewer, status, now_unix_seconds(), out); + } + + fn emit_tool_review_event_at( + &mut self, + pending_review: PendingToolReviewState, + reviewer: ToolReviewReviewer, + status: ToolReviewStatus, + completed_at: u64, + out: &mut Vec, + ) { + if let Some(item_id) = pending_review.item_id.as_ref() { + self.record_tool_review_summary(item_id, reviewer, status, &pending_review); + } + let Some(thread_metadata) = self.thread_metadata(&pending_review.thread_id) else { + tracing::warn!( + thread_id = %pending_review.thread_id, + turn_id = %pending_review.turn_id, + review_id = %pending_review.review_id, + "dropping tool review analytics event: missing thread lifecycle metadata" + ); + return; + }; + out.push(TrackEventRequest::ToolCallReview( + CodexToolCallReviewEventRequest { + event_type: "codex_tool_call_review_event", + event_params: CodexToolCallReviewEventParams { + thread_id: pending_review.thread_id, + turn_id: pending_review.turn_id, + item_id: pending_review.item_id, + review_id: pending_review.review_id, + thread_source: thread_metadata.thread_source, + subagent_source: thread_metadata.subagent_source.clone(), + parent_thread_id: thread_metadata.parent_thread_id.clone(), + tool_kind: pending_review.tool_kind, + tool_name: pending_review.tool_name, + reviewer, + trigger: pending_review.trigger, + status, + created_at: pending_review.created_at, + completed_at: Some(completed_at), + duration_ms: pending_review + .duration_started + .then(|| observed_duration_ms(pending_review.created_at, completed_at)) + .flatten(), + }, + }, + )); + } + + fn record_tool_review_summary( + &mut self, + item_id: &str, + reviewer: ToolReviewReviewer, + status: ToolReviewStatus, + pending_review: &PendingToolReviewState, + ) { + let summary = self + .tool_review_summaries + .entry(item_id.to_string()) + .or_default(); + summary.review_count += 1; + match reviewer { + ToolReviewReviewer::Guardian => summary.guardian_review_count += 1, + ToolReviewReviewer::User => summary.user_review_count += 1, + } + summary.final_approval_outcome = Some(tool_item_final_approval_outcome(reviewer, status)); + summary.requested_additional_permissions |= pending_review.requested_additional_permissions; + summary.requested_network_access |= pending_review.requested_network_access; + } + fn maybe_emit_turn_event(&mut self, turn_id: &str, out: &mut Vec) { let Some(turn_state) = self.turns.get(turn_id) else { return; @@ -1057,13 +1384,24 @@ fn tool_item_id(item: &ThreadItem) -> Option<&str> { } } -fn tool_item_event( - thread_id: &str, - turn_id: &str, - item: &ThreadItem, - connection_state: &ConnectionState, - thread_metadata: &ThreadMetadataState, -) -> Option { +struct ToolItemEventInput<'a> { + thread_id: &'a str, + turn_id: &'a str, + item: &'a ThreadItem, + connection_state: &'a ConnectionState, + thread_metadata: &'a ThreadMetadataState, + review_summary: Option<&'a ToolReviewSummary>, +} + +fn tool_item_event(input: ToolItemEventInput<'_>) -> Option { + let ToolItemEventInput { + thread_id, + turn_id, + item, + connection_state, + thread_metadata, + review_summary, + } = input; match item { ThreadItem::CommandExecution { id, @@ -1094,6 +1432,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::CommandExecution( @@ -1136,6 +1475,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::FileChange(CodexFileChangeEventRequest { @@ -1179,6 +1519,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::McpToolCall( @@ -1226,6 +1567,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::DynamicToolCall( @@ -1274,6 +1616,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::CollabAgentToolCall( @@ -1337,6 +1680,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::WebSearch(CodexWebSearchEventRequest { @@ -1377,6 +1721,7 @@ fn tool_item_event( completed_at_ms, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::ImageGeneration( @@ -1406,6 +1751,7 @@ struct ToolItemContext<'a> { completed_at_ms: u64, connection_state: &'a ConnectionState, thread_metadata: &'a ThreadMetadataState, + review_summary: Option<&'a ToolReviewSummary>, } fn tool_item_base( @@ -1417,6 +1763,7 @@ fn tool_item_base( context: ToolItemContext<'_>, ) -> CodexToolItemEventBase { let thread_metadata = context.thread_metadata; + let review_summary = context.review_summary.cloned().unwrap_or_default(); CodexToolItemEventBase { thread_id: thread_id.to_string(), turn_id: turn_id.to_string(), @@ -1431,15 +1778,187 @@ fn tool_item_base( completed_at_ms: context.completed_at_ms, duration_ms: outcome.duration_ms, execution_started: true, - review_count: 0, - guardian_review_count: 0, - user_review_count: 0, - final_approval_outcome: ToolItemFinalApprovalOutcome::Unknown, + review_count: review_summary.review_count, + guardian_review_count: review_summary.guardian_review_count, + user_review_count: review_summary.user_review_count, + final_approval_outcome: review_summary + .final_approval_outcome + .unwrap_or(ToolItemFinalApprovalOutcome::Unknown), terminal_status: outcome.terminal_status, failure_kind: outcome.failure_kind, - requested_additional_permissions: false, - requested_network_access: false, - retry_count: 0, + requested_additional_permissions: review_summary.requested_additional_permissions, + requested_network_access: review_summary.requested_network_access, + retry_count: review_summary.review_count.saturating_sub(1), + } +} + +fn observed_duration_ms(started_at_ms: u64, completed_at_ms: u64) -> Option { + completed_at_ms.checked_sub(started_at_ms) +} + +fn user_review_id(request_id: &RequestId) -> String { + format!("user:{request_id}") +} + +fn command_execution_review_status(decision: CommandExecutionApprovalDecision) -> ToolReviewStatus { + match decision { + CommandExecutionApprovalDecision::Accept => ToolReviewStatus::Approved, + CommandExecutionApprovalDecision::AcceptForSession => ToolReviewStatus::ApprovedForSession, + CommandExecutionApprovalDecision::AcceptWithExecpolicyAmendment { .. } => { + ToolReviewStatus::ApprovedExecpolicyAmendment + } + CommandExecutionApprovalDecision::ApplyNetworkPolicyAmendment { + network_policy_amendment, + } => match network_policy_amendment.action { + NetworkPolicyRuleAction::Allow => ToolReviewStatus::NetworkPolicyAllow, + NetworkPolicyRuleAction::Deny => ToolReviewStatus::NetworkPolicyDeny, + }, + CommandExecutionApprovalDecision::Decline => ToolReviewStatus::Denied, + CommandExecutionApprovalDecision::Cancel => ToolReviewStatus::Aborted, + } +} + +fn file_change_review_status(decision: FileChangeApprovalDecision) -> ToolReviewStatus { + match decision { + FileChangeApprovalDecision::Accept => ToolReviewStatus::Approved, + FileChangeApprovalDecision::AcceptForSession => ToolReviewStatus::ApprovedForSession, + FileChangeApprovalDecision::Decline => ToolReviewStatus::Denied, + FileChangeApprovalDecision::Cancel => ToolReviewStatus::Aborted, + } +} + +fn guardian_review_status(status: GuardianApprovalReviewStatus) -> Option { + match status { + GuardianApprovalReviewStatus::InProgress => None, + GuardianApprovalReviewStatus::Approved => Some(ToolReviewStatus::Approved), + GuardianApprovalReviewStatus::Denied => Some(ToolReviewStatus::Denied), + GuardianApprovalReviewStatus::TimedOut => Some(ToolReviewStatus::TimedOut), + GuardianApprovalReviewStatus::Aborted => Some(ToolReviewStatus::Aborted), + } +} + +fn guardian_review_tool_metadata( + action: &GuardianApprovalReviewAction, +) -> (ToolReviewToolKind, String, ToolReviewTrigger) { + match action { + GuardianApprovalReviewAction::Command { source, .. } => ( + ToolReviewToolKind::CommandExecution, + app_server_guardian_command_tool_name(*source).to_string(), + ToolReviewTrigger::Initial, + ), + GuardianApprovalReviewAction::Execve { source, .. } => ( + ToolReviewToolKind::CommandExecution, + app_server_guardian_command_tool_name(*source).to_string(), + ToolReviewTrigger::SubcommandExecve, + ), + GuardianApprovalReviewAction::ApplyPatch { .. } => ( + ToolReviewToolKind::FileChange, + "apply_patch".to_string(), + ToolReviewTrigger::SandboxRetry, + ), + GuardianApprovalReviewAction::NetworkAccess { .. } => ( + ToolReviewToolKind::NetworkAccess, + "network_access".to_string(), + ToolReviewTrigger::NetworkRetry, + ), + GuardianApprovalReviewAction::RequestPermissions { permissions, .. } => { + let requested_network_access = permissions + .network + .as_ref() + .and_then(|network| network.enabled) + .unwrap_or(false); + let trigger = if requested_network_access { + ToolReviewTrigger::NetworkRetry + } else if permissions.file_system.is_some() { + ToolReviewTrigger::SandboxRetry + } else { + ToolReviewTrigger::Initial + }; + ( + ToolReviewToolKind::Permissions, + "permissions".to_string(), + trigger, + ) + } + GuardianApprovalReviewAction::McpToolCall { tool_name, .. } => ( + ToolReviewToolKind::McpToolCall, + tool_name.clone(), + ToolReviewTrigger::Initial, + ), + } +} + +fn guardian_review_requested_additional_permissions(action: &GuardianApprovalReviewAction) -> bool { + match action { + GuardianApprovalReviewAction::ApplyPatch { .. } + | GuardianApprovalReviewAction::NetworkAccess { .. } => true, + GuardianApprovalReviewAction::RequestPermissions { permissions, .. } => { + guardian_review_request_permissions_network_enabled(permissions) + || permissions.file_system.is_some() + } + GuardianApprovalReviewAction::Command { .. } + | GuardianApprovalReviewAction::Execve { .. } + | GuardianApprovalReviewAction::McpToolCall { .. } => false, + } +} + +fn guardian_review_requested_network_access(action: &GuardianApprovalReviewAction) -> bool { + match action { + GuardianApprovalReviewAction::NetworkAccess { .. } => true, + GuardianApprovalReviewAction::RequestPermissions { permissions, .. } => { + guardian_review_request_permissions_network_enabled(permissions) + } + GuardianApprovalReviewAction::ApplyPatch { .. } + | GuardianApprovalReviewAction::Command { .. } + | GuardianApprovalReviewAction::Execve { .. } + | GuardianApprovalReviewAction::McpToolCall { .. } => false, + } +} + +fn guardian_review_request_permissions_network_enabled( + permissions: &RequestPermissionProfile, +) -> bool { + permissions + .network + .as_ref() + .and_then(|network| network.enabled) + .unwrap_or(false) +} + +fn app_server_guardian_command_tool_name(source: AppServerGuardianCommandSource) -> &'static str { + match source { + AppServerGuardianCommandSource::Shell => "shell", + AppServerGuardianCommandSource::UnifiedExec => "unified_exec", + } +} + +fn tool_item_final_approval_outcome( + reviewer: ToolReviewReviewer, + status: ToolReviewStatus, +) -> ToolItemFinalApprovalOutcome { + match (reviewer, status) { + (ToolReviewReviewer::Guardian, ToolReviewStatus::Approved) => { + ToolItemFinalApprovalOutcome::GuardianApproved + } + ( + ToolReviewReviewer::Guardian, + ToolReviewStatus::Denied | ToolReviewStatus::NetworkPolicyDeny, + ) => ToolItemFinalApprovalOutcome::GuardianDenied, + (ToolReviewReviewer::Guardian, _) => ToolItemFinalApprovalOutcome::GuardianAborted, + ( + ToolReviewReviewer::User, + ToolReviewStatus::Approved + | ToolReviewStatus::ApprovedExecpolicyAmendment + | ToolReviewStatus::NetworkPolicyAllow, + ) => ToolItemFinalApprovalOutcome::UserApproved, + (ToolReviewReviewer::User, ToolReviewStatus::ApprovedForSession) => { + ToolItemFinalApprovalOutcome::UserApprovedForSession + } + ( + ToolReviewReviewer::User, + ToolReviewStatus::Denied | ToolReviewStatus::NetworkPolicyDeny, + ) => ToolItemFinalApprovalOutcome::UserDenied, + (ToolReviewReviewer::User, _) => ToolItemFinalApprovalOutcome::UserAborted, } }