diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index f272e9c065..6b3ad4be9b 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::ClientResponse; 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; 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::ThreadItem; @@ -602,6 +630,190 @@ fn sample_command_execution_item( } } +async fn ingest_tool_review_prerequisites( + reducer: &mut AnalyticsReducer, + events: &mut Vec, +) { + reducer + .ingest(sample_initialize_fact(/*connection_id*/ 7), events) + .await; + reducer + .ingest( + AnalyticsFact::ClientResponse { + connection_id: 7, + response: Box::new(sample_thread_start_response( + "thread-1", /*ephemeral*/ false, "gpt-5", + )), + }, + events, + ) + .await; + events.clear(); +} + +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(), + 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, + }, + } +} + +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()) @@ -1051,6 +1263,55 @@ fn command_execution_event_allows_null_thread_denormalization() { assert_eq!(payload["event_params"]["thread_source"], json!(null)); } +#[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(); @@ -1473,6 +1734,337 @@ async fn item_lifecycle_notifications_publish_command_execution_event() { assert_eq!(payload[0]["event_params"]["thread_source"], json!(null)); } +#[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( + notification_fact(sample_guardian_review_started( + "guardian-review-1", + Some("item-1"), + )), + &mut events, + ) + .await; + reducer + .ingest( + notification_fact(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( + notification_fact(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 { + connection_id: 7, + 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, + /*duration_ms*/ None, + ), + })), + }, + &mut events, + ) + .await; + reducer + .ingest( + AnalyticsFact::Notification { + connection_id: 7, + 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(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 8a3659dd3b..fe83f6c42f 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; @@ -39,6 +41,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; @@ -75,16 +81,25 @@ 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; 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::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; @@ -113,6 +128,9 @@ pub(crate) struct AnalyticsReducer { thread_connections: HashMap, thread_metadata: HashMap, tool_items: HashMap, + tool_review_requests: HashMap, + guardian_review_starts: HashMap, + tool_review_summaries: HashMap, } struct ConnectionState { @@ -125,6 +143,35 @@ struct ToolItemState { started_at: u64, } +#[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>, @@ -252,12 +299,14 @@ impl AnalyticsReducer { self.ingest_notification(connection_id, *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); @@ -614,6 +663,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_secs(), + 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_secs(), + 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_secs(), + 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, @@ -692,17 +886,30 @@ impl AnalyticsReducer { return; }; let completed_at = now_unix_secs(); - if let Some(event) = tool_item_event( - ¬ification.thread_id, - ¬ification.turn_id, - ¬ification.item, - started.started_at, + if let Some(event) = tool_item_event(ToolItemEventInput { + thread_id: ¬ification.thread_id, + turn_id: ¬ification.turn_id, + item: ¬ification.item, + started_at: started.started_at, completed_at, connection_state, - self.thread_metadata.get(¬ification.thread_id), - ) { + thread_metadata: self.thread_metadata.get(¬ification.thread_id), + 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_secs(), + }, + ); + } + 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 { @@ -836,6 +1043,50 @@ impl AnalyticsReducer { ))); } + fn ingest_guardian_tool_review_completed( + &mut self, + notification: codex_app_server_protocol::ItemGuardianApprovalReviewCompletedNotification, + out: &mut Vec, + ) { + let completed_at = now_unix_secs(); + 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: matches!( + notification.action, + GuardianApprovalReviewAction::ApplyPatch { .. } + | GuardianApprovalReviewAction::NetworkAccess { .. } + ), + requested_network_access: matches!( + notification.action, + GuardianApprovalReviewAction::NetworkAccess { .. } + ), + }; + self.emit_tool_review_event_at( + pending_review, + ToolReviewReviewer::Guardian, + status, + completed_at, + out, + ); + } + fn ingest_turn_steer_response( &mut self, connection_id: u64, @@ -899,6 +1150,86 @@ 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_secs(), 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 thread_metadata = self + .thread_metadata + .get(&pending_review.thread_id) + .cloned() + .unwrap_or_default(); + 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, + parent_thread_id: thread_metadata.parent_thread_id, + 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(|| { + completed_duration_ms( + /*item_duration_ms*/ None, + 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; @@ -976,15 +1307,29 @@ fn tool_item_id(item: &ThreadItem) -> Option<&str> { } } -fn tool_item_event( - thread_id: &str, - turn_id: &str, - item: &ThreadItem, +struct ToolItemEventInput<'a> { + thread_id: &'a str, + turn_id: &'a str, + item: &'a ThreadItem, started_at: u64, completed_at: u64, - connection_state: &ConnectionState, - thread_metadata: Option<&ThreadMetadataState>, -) -> Option { + connection_state: &'a ConnectionState, + thread_metadata: Option<&'a ThreadMetadataState>, + review_summary: Option<&'a ToolReviewSummary>, +} + +fn tool_item_event(input: ToolItemEventInput<'_>) -> Option { + let ToolItemEventInput { + thread_id, + turn_id, + item, + started_at, + completed_at, + connection_state, + thread_metadata, + review_summary, + } = input; + match item { ThreadItem::CommandExecution { id, @@ -1015,6 +1360,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::CommandExecution( @@ -1052,6 +1398,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::FileChange(CodexFileChangeEventRequest { @@ -1095,6 +1442,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::McpToolCall( @@ -1142,6 +1490,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::DynamicToolCall( @@ -1185,6 +1534,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::CollabAgentToolCall( @@ -1240,6 +1590,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::WebSearch(CodexWebSearchEventRequest { @@ -1275,6 +1626,7 @@ fn tool_item_event( completed_at, connection_state, thread_metadata, + review_summary, }, ); Some(TrackEventRequest::ImageGeneration( @@ -1304,6 +1656,7 @@ struct ToolItemContext<'a> { completed_at: u64, connection_state: &'a ConnectionState, thread_metadata: Option<&'a ThreadMetadataState>, + review_summary: Option<&'a ToolReviewSummary>, } fn tool_item_base( @@ -1315,6 +1668,7 @@ fn tool_item_base( context: ToolItemContext<'_>, ) -> CodexToolItemEventBase { let thread_metadata = context.thread_metadata.cloned().unwrap_or_default(); + let review_summary = context.review_summary.cloned().unwrap_or_default(); CodexToolItemEventBase { thread_id: thread_id.to_string(), turn_id: turn_id.to_string(), @@ -1329,15 +1683,17 @@ fn tool_item_base( completed_at: Some(context.completed_at), duration_ms: outcome.duration_ms, execution_started: true, - review_count: 0, - guardian_review_count: 0, - user_review_count: 0, - final_approval_outcome: ToolItemFinalApprovalOutcome::NotNeeded, + 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::NotNeeded), 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), } } @@ -1357,6 +1713,116 @@ fn observed_completed_duration_ms(started_at: u64, completed_at: u64) -> Option< completed_duration_ms(/*item_duration_ms*/ None, started_at, completed_at) } +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::McpToolCall { tool_name, .. } => ( + ToolReviewToolKind::McpToolCall, + tool_name.clone(), + ToolReviewTrigger::Initial, + ), + } +} + +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, + } +} + fn command_execution_source_kind(source: CommandExecutionSource) -> CommandExecutionSourceKind { match source { CommandExecutionSource::Agent => CommandExecutionSourceKind::Agent,