From 4e7399c6b9878aa64051aea081278452bb6f8b62 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Wed, 22 Apr 2026 01:02:47 -0700 Subject: [PATCH] [codex-analytics] guardian review analytics events emission (#17693) ## Why Guardian approvals now run as review sessions, but Codex analytics did not have a terminal event for those reviews. That made it hard to measure approval outcomes, failure modes, Guardian session reuse, model metadata, token usage, and timing separately from the parent turn. ## What changed Adds `codex_guardian_review` analytics emission for Guardian approval reviews. The event is emitted from the Guardian review path with review identity, target item id, approval request source, a PII-minimized reviewed-action shape, terminal decision/status, failure reason, Guardian assessment fields, Guardian session metadata, token usage, and timing metadata. The reviewed-action payload intentionally omits high-risk fields such as shell commands, working directories, argv, file paths, network targets/hosts, rationale, retry reason, and permission justifications. It also classifies prompt-build failures separately from Guardian session/runtime failures so fail-closed cases are distinguishable in analytics. ## Verification - Guardian review analytics tests cover terminal success, timeout/cancel/fail-closed paths, session metadata, and token usage plumbing. - `cargo clippy -p codex-core --lib --tests -- -D warnings` --- [//]: # (BEGIN SAPLING FOOTER) Stack created with [Sapling](https://sapling-scm.com). Best reviewed with [ReviewStack](https://reviewstack.dev/openai/codex/pull/17693). * #17696 * #17695 * __->__ #17693 --- codex-rs/analytics/src/client.rs | 11 +- codex-rs/analytics/src/events.rs | 141 +++++++ codex-rs/analytics/src/lib.rs | 2 + codex-rs/core/src/codex_delegate.rs | 4 + .../core/src/guardian/approval_request.rs | 61 +++ codex-rs/core/src/guardian/review.rs | 385 ++++++++++++++---- codex-rs/core/src/guardian/review_session.rs | 253 +++++++++--- codex-rs/core/src/guardian/tests.rs | 10 +- codex-rs/core/src/session/mod.rs | 1 + 9 files changed, 730 insertions(+), 138 deletions(-) diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index 1a4b5defe..a3a20231f 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -1,5 +1,6 @@ use crate::events::AppServerRpcTransport; -use crate::events::GuardianReviewEventParams; +use crate::events::GuardianReviewAnalyticsResult; +use crate::events::GuardianReviewTrackContext; use crate::events::TrackEventRequest; use crate::events::TrackEventsRequest; use crate::events::current_runtime_metadata; @@ -161,9 +162,13 @@ impl AnalyticsEventsClient { )); } - pub fn track_guardian_review(&self, input: GuardianReviewEventParams) { + pub fn track_guardian_review( + &self, + tracking: &GuardianReviewTrackContext, + result: GuardianReviewAnalyticsResult, + ) { self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::GuardianReview( - Box::new(input), + Box::new(tracking.event_params(result)), ))); } diff --git a/codex-rs/analytics/src/events.rs b/codex-rs/analytics/src/events.rs index b542f8268..73f2886f2 100644 --- a/codex-rs/analytics/src/events.rs +++ b/codex-rs/analytics/src/events.rs @@ -1,3 +1,5 @@ +use std::time::Instant; + use crate::facts::AppInvocation; use crate::facts::CodexCompactionEvent; use crate::facts::CompactionImplementation; @@ -16,6 +18,7 @@ use crate::facts::TurnStatus; use crate::facts::TurnSteerRejectionReason; use crate::facts::TurnSteerResult; use crate::facts::TurnSubmissionType; +use crate::now_unix_seconds; use codex_app_server_protocol::CodexErrorInfo; use codex_login::default_client::originator; use codex_plugin::PluginTelemetryMetadata; @@ -30,6 +33,7 @@ use codex_protocol::protocol::HookEventName; use codex_protocol::protocol::HookRunStatus; use codex_protocol::protocol::HookSource; use codex_protocol::protocol::SubAgentSource; +use codex_protocol::protocol::TokenUsage; use serde::Serialize; #[derive(Clone, Copy, Debug, Serialize)] @@ -200,6 +204,7 @@ pub enum GuardianReviewedAction { connector_name: Option, tool_title: Option, }, + RequestPermissions {}, } #[derive(Clone, Serialize)] @@ -235,6 +240,142 @@ pub struct GuardianReviewEventParams { pub total_tokens: Option, } +pub struct GuardianReviewTrackContext { + thread_id: String, + turn_id: String, + review_id: String, + target_item_id: Option, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: GuardianReviewedAction, + review_timeout_ms: u64, + started_at: u64, + started_instant: Instant, +} + +impl GuardianReviewTrackContext { + pub fn new( + thread_id: String, + turn_id: String, + review_id: String, + target_item_id: Option, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: GuardianReviewedAction, + review_timeout_ms: u64, + ) -> Self { + Self { + thread_id, + turn_id, + review_id, + target_item_id, + approval_request_source, + reviewed_action, + review_timeout_ms, + started_at: now_unix_seconds(), + started_instant: Instant::now(), + } + } + + pub(crate) fn event_params( + &self, + result: GuardianReviewAnalyticsResult, + ) -> GuardianReviewEventParams { + GuardianReviewEventParams { + thread_id: self.thread_id.clone(), + turn_id: self.turn_id.clone(), + review_id: self.review_id.clone(), + target_item_id: self.target_item_id.clone(), + approval_request_source: self.approval_request_source, + reviewed_action: self.reviewed_action.clone(), + reviewed_action_truncated: result.reviewed_action_truncated, + decision: result.decision, + terminal_status: result.terminal_status, + failure_reason: result.failure_reason, + risk_level: result.risk_level, + user_authorization: result.user_authorization, + outcome: result.outcome, + guardian_thread_id: result.guardian_thread_id, + guardian_session_kind: result.guardian_session_kind, + guardian_model: result.guardian_model, + guardian_reasoning_effort: result.guardian_reasoning_effort, + had_prior_review_context: result.had_prior_review_context, + review_timeout_ms: self.review_timeout_ms, + // TODO(rhan-oai): plumb nested Guardian review session tool-call counts. + tool_call_count: None, + time_to_first_token_ms: result.time_to_first_token_ms, + completion_latency_ms: Some(self.started_instant.elapsed().as_millis() as u64), + started_at: self.started_at, + completed_at: Some(now_unix_seconds()), + input_tokens: result.token_usage.as_ref().map(|usage| usage.input_tokens), + cached_input_tokens: result + .token_usage + .as_ref() + .map(|usage| usage.cached_input_tokens), + output_tokens: result.token_usage.as_ref().map(|usage| usage.output_tokens), + reasoning_output_tokens: result + .token_usage + .as_ref() + .map(|usage| usage.reasoning_output_tokens), + total_tokens: result.token_usage.as_ref().map(|usage| usage.total_tokens), + } + } +} + +#[derive(Debug)] +pub struct GuardianReviewAnalyticsResult { + pub decision: GuardianReviewDecision, + pub terminal_status: GuardianReviewTerminalStatus, + pub failure_reason: Option, + pub risk_level: Option, + pub user_authorization: Option, + pub outcome: Option, + pub guardian_thread_id: Option, + pub guardian_session_kind: Option, + pub guardian_model: Option, + pub guardian_reasoning_effort: Option, + pub had_prior_review_context: Option, + pub reviewed_action_truncated: bool, + pub token_usage: Option, + pub time_to_first_token_ms: Option, +} + +impl GuardianReviewAnalyticsResult { + pub fn without_session() -> Self { + Self { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: None, + risk_level: None, + user_authorization: None, + outcome: None, + guardian_thread_id: None, + guardian_session_kind: None, + guardian_model: None, + guardian_reasoning_effort: None, + had_prior_review_context: None, + reviewed_action_truncated: false, + token_usage: None, + time_to_first_token_ms: None, + } + } + + pub fn from_session( + guardian_thread_id: String, + guardian_session_kind: GuardianReviewSessionKind, + guardian_model: String, + guardian_reasoning_effort: Option, + had_prior_review_context: bool, + ) -> Self { + Self { + guardian_thread_id: Some(guardian_thread_id), + guardian_session_kind: Some(guardian_session_kind), + guardian_model: Some(guardian_model), + guardian_reasoning_effort, + had_prior_review_context: Some(had_prior_review_context), + ..Self::without_session() + } + } +} + #[derive(Serialize)] pub(crate) struct GuardianReviewEventPayload { pub(crate) app_server_client: CodexAppServerClientMetadata, diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index 5c4cdfac7..ed0f1036c 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -9,11 +9,13 @@ use std::time::UNIX_EPOCH; pub use client::AnalyticsEventsClient; pub use events::AppServerRpcTransport; pub use events::GuardianApprovalRequestSource; +pub use events::GuardianReviewAnalyticsResult; pub use events::GuardianReviewDecision; pub use events::GuardianReviewEventParams; pub use events::GuardianReviewFailureReason; pub use events::GuardianReviewSessionKind; pub use events::GuardianReviewTerminalStatus; +pub use events::GuardianReviewTrackContext; pub use events::GuardianReviewedAction; pub use facts::AnalyticsJsonRpcError; pub use facts::AppInvocation; diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index 858154eb7..22d1509fe 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -3,6 +3,7 @@ use std::sync::Arc; use async_channel::Receiver; use async_channel::Sender; +use codex_analytics::GuardianApprovalRequestSource; use codex_async_utils::OrCancelExt; use codex_protocol::protocol::ApplyPatchApprovalRequestEvent; use codex_protocol::protocol::Event; @@ -471,6 +472,7 @@ async fn handle_exec_approval( justification: None, }, reason, + GuardianApprovalRequestSource::DelegatedSubagent, review_cancel.clone(), ); await_approval_with_cancel( @@ -573,6 +575,7 @@ async fn handle_patch_approval( patch, }, reason.clone(), + GuardianApprovalRequestSource::DelegatedSubagent, review_cancel.clone(), ); Some( @@ -689,6 +692,7 @@ async fn maybe_auto_review_mcp_request_user_input( new_guardian_review_id(), build_guardian_mcp_tool_review_request(&event.call_id, &invocation, metadata.as_ref()), /*retry_reason*/ None, + GuardianApprovalRequestSource::DelegatedSubagent, review_cancel.clone(), ); let decision = await_approval_with_cancel( diff --git a/codex-rs/core/src/guardian/approval_request.rs b/codex-rs/core/src/guardian/approval_request.rs index 747837738..471e05442 100644 --- a/codex-rs/core/src/guardian/approval_request.rs +++ b/codex-rs/core/src/guardian/approval_request.rs @@ -1,5 +1,6 @@ use std::path::Path; +use codex_analytics::GuardianReviewedAction; use codex_protocol::approvals::GuardianAssessmentAction; use codex_protocol::approvals::GuardianCommandSource; use codex_protocol::approvals::NetworkApprovalProtocol; @@ -390,6 +391,66 @@ pub(crate) fn guardian_assessment_action( } } +pub(crate) fn guardian_reviewed_action( + request: &GuardianApprovalRequest, +) -> GuardianReviewedAction { + match request { + GuardianApprovalRequest::Shell { + sandbox_permissions, + additional_permissions, + .. + } => GuardianReviewedAction::Shell { + sandbox_permissions: *sandbox_permissions, + additional_permissions: additional_permissions.clone(), + }, + GuardianApprovalRequest::ExecCommand { + sandbox_permissions, + additional_permissions, + tty, + .. + } => GuardianReviewedAction::UnifiedExec { + sandbox_permissions: *sandbox_permissions, + additional_permissions: additional_permissions.clone(), + tty: *tty, + }, + #[cfg(unix)] + GuardianApprovalRequest::Execve { + source, + program, + additional_permissions, + .. + } => GuardianReviewedAction::Execve { + source: *source, + program: program.clone(), + additional_permissions: additional_permissions.clone(), + }, + GuardianApprovalRequest::ApplyPatch { .. } => GuardianReviewedAction::ApplyPatch {}, + GuardianApprovalRequest::NetworkAccess { protocol, port, .. } => { + GuardianReviewedAction::NetworkAccess { + protocol: *protocol, + port: *port, + } + } + GuardianApprovalRequest::McpToolCall { + server, + tool_name, + connector_id, + connector_name, + tool_title, + .. + } => GuardianReviewedAction::McpToolCall { + server: server.clone(), + tool_name: tool_name.clone(), + connector_id: connector_id.clone(), + connector_name: connector_name.clone(), + tool_title: tool_title.clone(), + }, + GuardianApprovalRequest::RequestPermissions { .. } => { + GuardianReviewedAction::RequestPermissions {} + } + } +} + pub(crate) fn guardian_request_target_item_id(request: &GuardianApprovalRequest) -> Option<&str> { match request { GuardianApprovalRequest::Shell { id, .. } diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index 486c079f3..ce01ac970 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -1,5 +1,12 @@ use std::sync::Arc; +use codex_analytics::GuardianApprovalRequestSource; +use codex_analytics::GuardianReviewAnalyticsResult; +use codex_analytics::GuardianReviewDecision; +use codex_analytics::GuardianReviewFailureReason; +use codex_analytics::GuardianReviewTerminalStatus; +use codex_analytics::GuardianReviewTrackContext; +use codex_features::Feature; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; @@ -17,6 +24,7 @@ use tokio_util::sync::CancellationToken; use crate::session::session::Session; use crate::session::turn_context::TurnContext; +use super::GUARDIAN_REVIEW_TIMEOUT; use super::GUARDIAN_REVIEWER_NAME; use super::GuardianApprovalRequest; use super::GuardianAssessment; @@ -25,6 +33,7 @@ use super::GuardianRejection; use super::approval_request::guardian_assessment_action; use super::approval_request::guardian_request_target_item_id; use super::approval_request::guardian_request_turn_id; +use super::approval_request::guardian_reviewed_action; use super::prompt::guardian_output_schema; use super::prompt::parse_guardian_assessment; use super::review_session::GuardianReviewSessionOutcome; @@ -76,9 +85,47 @@ pub(crate) fn guardian_timeout_message() -> String { #[derive(Debug)] pub(super) enum GuardianReviewOutcome { - Completed(anyhow::Result), - TimedOut, - Aborted, + Completed(GuardianAssessment), + Error(GuardianReviewError), +} + +#[derive(Debug)] +pub(super) enum GuardianReviewError { + PromptBuild { message: String }, + Session { message: String }, + Parse { message: String }, + Timeout, + Cancelled, +} + +impl GuardianReviewError { + fn prompt_build(err: anyhow::Error) -> Self { + Self::PromptBuild { + message: err.to_string(), + } + } + + fn session(err: anyhow::Error) -> Self { + Self::Session { + message: err.to_string(), + } + } + + fn parse(err: anyhow::Error) -> Self { + Self::Parse { + message: err.to_string(), + } + } + + fn failure_reason(&self) -> GuardianReviewFailureReason { + match self { + Self::PromptBuild { .. } => GuardianReviewFailureReason::PromptBuildError, + Self::Session { .. } => GuardianReviewFailureReason::SessionError, + Self::Parse { .. } => GuardianReviewFailureReason::ParseError, + Self::Timeout => GuardianReviewFailureReason::Timeout, + Self::Cancelled => GuardianReviewFailureReason::Cancelled, + } + } } fn guardian_risk_level_str(level: GuardianRiskLevel) -> &'static str { @@ -110,6 +157,21 @@ pub(crate) fn is_guardian_reviewer_source( ) } +fn track_guardian_review( + session: &Session, + turn: &TurnContext, + tracking: &GuardianReviewTrackContext, + result: GuardianReviewAnalyticsResult, +) { + if !turn.config.features.enabled(Feature::GeneralAnalytics) { + return; + } + session + .services + .analytics_events_client + .track_guardian_review(tracking, result); +} + /// This function always fails closed: timeouts, review-session failures, and /// parse failures all block execution, but timeouts are still surfaced to the /// caller as distinct from explicit guardian denials. @@ -119,11 +181,21 @@ async fn run_guardian_review( review_id: String, request: GuardianApprovalRequest, retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, external_cancel: Option, ) -> ReviewDecision { let target_item_id = guardian_request_target_item_id(&request).map(str::to_string); let assessment_turn_id = guardian_request_turn_id(&request, &turn.sub_id).to_string(); let action_summary = guardian_assessment_action(&request); + let review_tracking = GuardianReviewTrackContext::new( + session.conversation_id.to_string(), + assessment_turn_id.clone(), + review_id.clone(), + target_item_id.clone(), + approval_request_source, + guardian_reviewed_action(&request), + GUARDIAN_REVIEW_TIMEOUT.as_millis() as u64, + ); session .send_event( turn.as_ref(), @@ -145,6 +217,17 @@ async fn run_guardian_review( .as_ref() .is_some_and(CancellationToken::is_cancelled) { + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + ..GuardianReviewAnalyticsResult::without_session() + }, + ); session .send_event( turn.as_ref(), @@ -166,73 +249,146 @@ async fn run_guardian_review( let schema = guardian_output_schema(); let terminal_action = action_summary.clone(); - let outcome = Box::pin(run_guardian_review_session( + let (outcome, analytics_result) = Box::pin(run_guardian_review_session( session.clone(), turn.clone(), request, - retry_reason, + retry_reason.clone(), schema, external_cancel, )) .await; let assessment = match outcome { - GuardianReviewOutcome::Completed(Ok(assessment)) => assessment, - GuardianReviewOutcome::Completed(Err(err)) => GuardianAssessment { - risk_level: GuardianRiskLevel::High, - user_authorization: GuardianUserAuthorization::Unknown, - outcome: GuardianAssessmentOutcome::Deny, - rationale: format!("Automatic approval review failed: {err}"), + GuardianReviewOutcome::Completed(assessment) => { + let approved = matches!(assessment.outcome, GuardianAssessmentOutcome::Allow); + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: if approved { + GuardianReviewDecision::Approved + } else { + GuardianReviewDecision::Denied + }, + terminal_status: if approved { + GuardianReviewTerminalStatus::Approved + } else { + GuardianReviewTerminalStatus::Denied + }, + failure_reason: None, + risk_level: Some(assessment.risk_level), + user_authorization: Some(assessment.user_authorization), + outcome: Some(assessment.outcome), + ..analytics_result + }, + ); + assessment + } + GuardianReviewOutcome::Error(error) => match error { + GuardianReviewError::Timeout => { + let rationale = + "Automatic approval review timed out while evaluating the requested approval." + .to_string(); + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::TimedOut, + failure_reason: Some(error.failure_reason()), + ..analytics_result + }, + ); + session + .send_event( + turn.as_ref(), + EventMsg::Warning(WarningEvent { + message: rationale.clone(), + }), + ) + .await; + session + .send_event( + turn.as_ref(), + EventMsg::GuardianAssessment(GuardianAssessmentEvent { + id: review_id, + target_item_id, + turn_id: assessment_turn_id, + status: GuardianAssessmentStatus::TimedOut, + risk_level: None, + user_authorization: None, + rationale: Some(rationale), + decision_source: Some(GuardianAssessmentDecisionSource::Agent), + action: terminal_action, + }), + ) + .await; + return ReviewDecision::TimedOut; + } + GuardianReviewError::Cancelled => { + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(error.failure_reason()), + ..analytics_result + }, + ); + session + .send_event( + turn.as_ref(), + EventMsg::GuardianAssessment(GuardianAssessmentEvent { + id: review_id, + target_item_id, + turn_id: assessment_turn_id, + status: GuardianAssessmentStatus::Aborted, + risk_level: None, + user_authorization: None, + rationale: None, + decision_source: Some(GuardianAssessmentDecisionSource::Agent), + action: action_summary, + }), + ) + .await; + return ReviewDecision::Abort; + } + GuardianReviewError::PromptBuild { .. } + | GuardianReviewError::Session { .. } + | GuardianReviewError::Parse { .. } => { + let message = match &error { + GuardianReviewError::PromptBuild { message } + | GuardianReviewError::Session { message } + | GuardianReviewError::Parse { message } => message, + GuardianReviewError::Timeout | GuardianReviewError::Cancelled => { + "guardian review failed" + } + }; + let rationale = format!("Automatic approval review failed: {message}"); + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(error.failure_reason()), + ..analytics_result + }, + ); + GuardianAssessment { + risk_level: GuardianRiskLevel::High, + user_authorization: GuardianUserAuthorization::Unknown, + outcome: GuardianAssessmentOutcome::Deny, + rationale, + } + } }, - GuardianReviewOutcome::TimedOut => { - let rationale = - "Automatic approval review timed out while evaluating the requested approval." - .to_string(); - session - .send_event( - turn.as_ref(), - EventMsg::Warning(WarningEvent { - message: rationale.clone(), - }), - ) - .await; - session - .send_event( - turn.as_ref(), - EventMsg::GuardianAssessment(GuardianAssessmentEvent { - id: review_id, - target_item_id, - turn_id: assessment_turn_id, - status: GuardianAssessmentStatus::TimedOut, - risk_level: None, - user_authorization: None, - rationale: Some(rationale), - decision_source: Some(GuardianAssessmentDecisionSource::Agent), - action: terminal_action, - }), - ) - .await; - return ReviewDecision::TimedOut; - } - GuardianReviewOutcome::Aborted => { - session - .send_event( - turn.as_ref(), - EventMsg::GuardianAssessment(GuardianAssessmentEvent { - id: review_id, - target_item_id, - turn_id: assessment_turn_id, - status: GuardianAssessmentStatus::Aborted, - risk_level: None, - user_authorization: None, - rationale: None, - decision_source: Some(GuardianAssessmentDecisionSource::Agent), - action: action_summary, - }), - ) - .await; - return ReviewDecision::Abort; - } }; let approved = match assessment.outcome { @@ -314,6 +470,7 @@ pub(crate) async fn review_approval_request( review_id, request, retry_reason, + GuardianApprovalRequestSource::MainTurn, /*external_cancel*/ None, )) .await @@ -325,16 +482,18 @@ pub(crate) async fn review_approval_request_with_cancel( review_id: String, request: GuardianApprovalRequest, retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, cancel_token: CancellationToken, ) -> ReviewDecision { - Box::pin(run_guardian_review( + run_guardian_review( Arc::clone(session), Arc::clone(turn), review_id, request, retry_reason, + approval_request_source, Some(cancel_token), - )) + ) .await } @@ -344,6 +503,7 @@ pub(crate) fn spawn_approval_request_review( review_id: String, request: GuardianApprovalRequest, retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, cancel_token: CancellationToken, ) -> oneshot::Receiver { let (tx, rx) = oneshot::channel(); @@ -361,6 +521,7 @@ pub(crate) fn spawn_approval_request_review( review_id, request, retry_reason, + approval_request_source, cancel_token, )); let _ = tx.send(decision); @@ -389,11 +550,16 @@ pub(super) async fn run_guardian_review_session( retry_reason: Option, schema: serde_json::Value, external_cancel: Option, -) -> GuardianReviewOutcome { +) -> (GuardianReviewOutcome, GuardianReviewAnalyticsResult) { let live_network_config = match session.services.network_proxy.as_ref() { Some(network_proxy) => match network_proxy.proxy().current_cfg().await { Ok(config) => Some(config), - Err(err) => return GuardianReviewOutcome::Completed(Err(err)), + Err(err) => { + return ( + GuardianReviewOutcome::Error(GuardianReviewError::prompt_build(err)), + GuardianReviewAnalyticsResult::without_session(), + ); + } }, None => None, }; @@ -443,10 +609,15 @@ pub(super) async fn run_guardian_review_session( ); let guardian_config = match guardian_config { Ok(config) => config, - Err(err) => return GuardianReviewOutcome::Completed(Err(err)), + Err(err) => { + return ( + GuardianReviewOutcome::Error(GuardianReviewError::prompt_build(err)), + GuardianReviewAnalyticsResult::without_session(), + ); + } }; - match Box::pin( + let (session_outcome, session_analytics_result) = Box::pin( session .guardian_review_session .run_review(GuardianReviewSessionParams { @@ -463,17 +634,75 @@ pub(super) async fn run_guardian_review_session( external_cancel, }), ) - .await - { - GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) => { - GuardianReviewOutcome::Completed(parse_guardian_assessment( - last_agent_message.as_deref(), - )) - } - GuardianReviewSessionOutcome::Completed(Err(err)) => { - GuardianReviewOutcome::Completed(Err(err)) - } - GuardianReviewSessionOutcome::TimedOut => GuardianReviewOutcome::TimedOut, - GuardianReviewSessionOutcome::Aborted => GuardianReviewOutcome::Aborted, + .await; + + match session_outcome { + GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) => match last_agent_message + { + Some(last_agent_message) => { + match parse_guardian_assessment(Some(&last_agent_message)) { + Ok(assessment) => ( + GuardianReviewOutcome::Completed(assessment), + session_analytics_result, + ), + Err(err) => ( + GuardianReviewOutcome::Error(GuardianReviewError::parse(err)), + session_analytics_result, + ), + } + } + None => ( + GuardianReviewOutcome::Error(GuardianReviewError::session(anyhow::anyhow!( + "guardian review completed without an assessment payload" + ))), + session_analytics_result, + ), + }, + GuardianReviewSessionOutcome::Completed(Err(err)) => ( + GuardianReviewOutcome::Error(GuardianReviewError::session(err)), + session_analytics_result, + ), + GuardianReviewSessionOutcome::PromptBuildFailed(err) => ( + GuardianReviewOutcome::Error(GuardianReviewError::prompt_build(err)), + session_analytics_result, + ), + GuardianReviewSessionOutcome::SessionFailed(err) => ( + GuardianReviewOutcome::Error(GuardianReviewError::session(err)), + session_analytics_result, + ), + GuardianReviewSessionOutcome::TimedOut => ( + GuardianReviewOutcome::Error(GuardianReviewError::Timeout), + session_analytics_result, + ), + GuardianReviewSessionOutcome::Aborted => ( + GuardianReviewOutcome::Error(GuardianReviewError::Cancelled), + session_analytics_result, + ), + } +} + +#[cfg(test)] +mod review_tests { + use super::*; + + #[test] + fn guardian_review_error_reason_distinguishes_error_kinds() { + let parse_error = GuardianReviewError::parse(anyhow::anyhow!("bad guardian JSON")); + let prompt_error = GuardianReviewError::prompt_build(anyhow::anyhow!("bad prompt/config")); + let session_error = + GuardianReviewError::session(anyhow::anyhow!("guardian runtime failed")); + + assert!(matches!( + parse_error.failure_reason(), + GuardianReviewFailureReason::ParseError + )); + assert!(matches!( + prompt_error.failure_reason(), + GuardianReviewFailureReason::PromptBuildError + )); + assert!(matches!( + session_error.failure_reason(), + GuardianReviewFailureReason::SessionError + )); } } diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index dc5d6e310..e68024dbf 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -5,6 +5,8 @@ use std::sync::Arc; use std::time::Duration; use anyhow::anyhow; +use codex_analytics::GuardianReviewAnalyticsResult; +use codex_analytics::GuardianReviewSessionKind; use codex_protocol::config_types::Personality; use codex_protocol::config_types::ReasoningSummary as ReasoningSummaryConfig; use codex_protocol::models::ResponseItem; @@ -16,6 +18,7 @@ use codex_protocol::protocol::Op; use codex_protocol::protocol::RolloutItem; use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SubAgentSource; +use codex_protocol::protocol::TokenUsage; use serde_json::Value; use tokio::sync::Mutex; use tokio::sync::Semaphore; @@ -52,6 +55,8 @@ const GUARDIAN_INTERRUPT_DRAIN_TIMEOUT: Duration = Duration::from_secs(5); #[derive(Debug)] pub(crate) enum GuardianReviewSessionOutcome { Completed(anyhow::Result>), + PromptBuildFailed(anyhow::Error), + SessionFailed(anyhow::Error), TimedOut, Aborted, } @@ -95,6 +100,21 @@ struct GuardianReviewState { last_committed_fork_snapshot: Option, } +fn had_prior_review_context(prompt_mode: &GuardianPromptMode) -> bool { + matches!(prompt_mode, GuardianPromptMode::Delta { .. }) +} + +fn token_usage_delta(start: &TokenUsage, end: &TokenUsage) -> TokenUsage { + TokenUsage { + input_tokens: (end.input_tokens - start.input_tokens).max(0), + cached_input_tokens: (end.cached_input_tokens - start.cached_input_tokens).max(0), + output_tokens: (end.output_tokens - start.output_tokens).max(0), + reasoning_output_tokens: (end.reasoning_output_tokens - start.reasoning_output_tokens) + .max(0), + total_tokens: (end.total_tokens - start.total_tokens).max(0), + } +} + struct EphemeralReviewCleanup { state: Arc>, review_session: Option>, @@ -262,13 +282,14 @@ impl GuardianReviewSessionManager { clippy::await_holding_invalid_type, reason = "review session selection and trunk spawning must stay serialized" )] - pub(crate) async fn run_review( + pub(super) async fn run_review( &self, params: GuardianReviewSessionParams, - ) -> GuardianReviewSessionOutcome { + ) -> (GuardianReviewSessionOutcome, GuardianReviewAnalyticsResult) { let deadline = tokio::time::Instant::now() + GUARDIAN_REVIEW_TIMEOUT; let next_reuse_key = GuardianReviewSessionReuseKey::from_spawn_config(¶ms.spawn_config); let mut stale_trunk_to_shutdown = None; + let mut spawned_trunk = false; let trunk_candidate = match run_before_review_deadline( deadline, params.external_cancel.as_ref(), @@ -302,16 +323,22 @@ impl GuardianReviewSessionManager { { Ok(Ok(review_session)) => Arc::new(review_session), Ok(Err(err)) => { - return GuardianReviewSessionOutcome::Completed(Err(err)); + return ( + GuardianReviewSessionOutcome::PromptBuildFailed(err), + GuardianReviewAnalyticsResult::without_session(), + ); + } + Err(outcome) => { + return (outcome, GuardianReviewAnalyticsResult::without_session()); } - Err(outcome) => return outcome, }; state.trunk = Some(Arc::clone(&review_session)); + spawned_trunk = true; } state.trunk.as_ref().cloned() } - Err(outcome) => return outcome, + Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()), }; if let Some(review_session) = stale_trunk_to_shutdown { @@ -319,9 +346,12 @@ impl GuardianReviewSessionManager { } let Some(trunk) = trunk_candidate else { - return GuardianReviewSessionOutcome::Completed(Err(anyhow!( - "guardian review session was not available after spawn" - ))); + return ( + GuardianReviewSessionOutcome::Completed(Err(anyhow!( + "guardian review session was not available after spawn" + ))), + GuardianReviewAnalyticsResult::without_session(), + ); }; if trunk.reuse_key != next_reuse_key { @@ -347,20 +377,30 @@ impl GuardianReviewSessionManager { } }; - let (outcome, keep_review_session) = - Box::pin(run_review_on_session(trunk.as_ref(), ¶ms, deadline)).await; + let guardian_session_kind = if spawned_trunk { + GuardianReviewSessionKind::TrunkNew + } else { + GuardianReviewSessionKind::TrunkReused + }; + let (outcome, keep_review_session, analytics_result) = Box::pin(run_review_on_session( + trunk.as_ref(), + ¶ms, + guardian_session_kind, + deadline, + )) + .await; if keep_review_session && matches!(outcome, GuardianReviewSessionOutcome::Completed(_)) { trunk.refresh_last_committed_fork_snapshot().await; } drop(trunk_guard); if keep_review_session { - outcome + (outcome, analytics_result) } else { if let Some(review_session) = self.remove_trunk_if_current(&trunk).await { review_session.shutdown_in_background(); } - outcome + (outcome, analytics_result) } } @@ -457,7 +497,7 @@ impl GuardianReviewSessionManager { reuse_key: GuardianReviewSessionReuseKey, deadline: tokio::time::Instant, fork_snapshot: Option, - ) -> GuardianReviewSessionOutcome { + ) -> (GuardianReviewSessionOutcome, GuardianReviewAnalyticsResult) { let spawn_cancel_token = CancellationToken::new(); let mut fork_config = params.spawn_config.clone(); fork_config.ephemeral = true; @@ -476,17 +516,23 @@ impl GuardianReviewSessionManager { .await { Ok(Ok(review_session)) => Arc::new(review_session), - Ok(Err(err)) => return GuardianReviewSessionOutcome::Completed(Err(err)), - Err(outcome) => return outcome, + Ok(Err(err)) => { + return ( + GuardianReviewSessionOutcome::PromptBuildFailed(err), + GuardianReviewAnalyticsResult::without_session(), + ); + } + Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()), }; self.register_active_ephemeral(Arc::clone(&review_session)) .await; let mut cleanup = EphemeralReviewCleanup::new(Arc::clone(&self.state), Arc::clone(&review_session)); - let (outcome, _) = Box::pin(run_review_on_session( + let (outcome, _, analytics_result) = Box::pin(run_review_on_session( review_session.as_ref(), ¶ms, + GuardianReviewSessionKind::EphemeralForked, deadline, )) .await; @@ -494,7 +540,7 @@ impl GuardianReviewSessionManager { cleanup.disarm(); review_session.shutdown_in_background(); } - outcome + (outcome, analytics_result) } } @@ -541,8 +587,13 @@ async fn spawn_guardian_review_session( async fn run_review_on_session( review_session: &GuardianReviewSession, params: &GuardianReviewSessionParams, + guardian_session_kind: GuardianReviewSessionKind, deadline: tokio::time::Instant, -) -> (GuardianReviewSessionOutcome, bool) { +) -> ( + GuardianReviewSessionOutcome, + bool, + GuardianReviewAnalyticsResult, +) { let (send_followup_reminder, prompt_mode) = { let state = review_session.state.lock().await; @@ -557,11 +608,34 @@ async fn run_review_on_session( (send_followup_reminder, prompt_mode) }; + let model_info = params + .parent_session + .services + .models_manager + .get_model_info( + params.model.as_str(), + ¶ms.spawn_config.to_models_manager_config(), + ) + .await; + let guardian_reasoning_effort = if model_info.supports_reasoning_summaries { + params + .reasoning_effort + .or(model_info.default_reasoning_level) + } else { + None + }; + let mut analytics_result = GuardianReviewAnalyticsResult::from_session( + review_session.codex.session.conversation_id.to_string(), + guardian_session_kind, + params.model.clone(), + guardian_reasoning_effort.map(|effort| effort.to_string()), + had_prior_review_context(&prompt_mode), + ); if send_followup_reminder { append_guardian_followup_reminder(review_session).await; } - let submit_result = run_before_review_deadline( + let prompt_items = run_before_review_deadline( deadline, params.external_cancel.as_ref(), Box::pin(async { @@ -574,56 +648,86 @@ async fn run_review_on_session( ) .await; - let prompt_items = build_guardian_prompt_items( + build_guardian_prompt_items( params.parent_session.as_ref(), params.retry_reason.clone(), params.request.clone(), prompt_mode, ) - .await?; - - review_session - .codex - .submit(Op::UserTurn { - environments: None, - items: prompt_items.items, - cwd: params.parent_turn.cwd.to_path_buf(), - approval_policy: AskForApproval::Never, - approvals_reviewer: None, - sandbox_policy: SandboxPolicy::new_read_only_policy(), - model: params.model.clone(), - effort: params.reasoning_effort, - summary: Some(params.reasoning_summary), - service_tier: None, - final_output_json_schema: Some(params.schema.clone()), - collaboration_mode: None, - personality: params.personality, - }) - .await?; - - Ok::(prompt_items.transcript_cursor) + .await }), ) .await; - let submit_result = match submit_result { - Ok(submit_result) => submit_result, - Err(outcome) => return (outcome, false), + let prompt_items = match prompt_items { + Ok(prompt_items) => prompt_items, + Err(outcome) => return (outcome, false, analytics_result), }; - let transcript_cursor = match submit_result { - Ok(transcript_cursor) => transcript_cursor, + let prompt_items = match prompt_items { + Ok(prompt_items) => prompt_items, Err(err) => { - return (GuardianReviewSessionOutcome::Completed(Err(err)), false); + return ( + GuardianReviewSessionOutcome::PromptBuildFailed(err.into()), + false, + analytics_result, + ); } }; + let transcript_cursor = prompt_items.transcript_cursor; + let token_usage_at_review_start = review_session + .codex + .session + .total_token_usage() + .await + .unwrap_or_default(); + + let submit_result = run_before_review_deadline( + deadline, + params.external_cancel.as_ref(), + Box::pin(review_session.codex.submit(Op::UserTurn { + environments: None, + items: prompt_items.items, + cwd: params.parent_turn.cwd.to_path_buf(), + approval_policy: AskForApproval::Never, + approvals_reviewer: None, + sandbox_policy: SandboxPolicy::new_read_only_policy(), + model: params.model.clone(), + effort: params.reasoning_effort, + summary: Some(params.reasoning_summary), + service_tier: None, + final_output_json_schema: Some(params.schema.clone()), + collaboration_mode: None, + personality: params.personality, + })), + ) + .await; + match submit_result { + Ok(Ok(_)) => {} + Ok(Err(err)) => { + return ( + GuardianReviewSessionOutcome::SessionFailed(err.into()), + false, + analytics_result, + ); + } + Err(outcome) => return (outcome, false, analytics_result), + } let outcome = wait_for_guardian_review(review_session, deadline, params.external_cancel.as_ref()).await; if matches!(outcome.0, GuardianReviewSessionOutcome::Completed(_)) { + if outcome.2 + && let Some(total_token_usage) = review_session.codex.session.total_token_usage().await + { + analytics_result.token_usage = Some(token_usage_delta( + &token_usage_at_review_start, + &total_token_usage, + )); + } let mut state = review_session.state.lock().await; state.prior_review_count = state.prior_review_count.saturating_add(1); state.last_reviewed_transcript_cursor = Some(transcript_cursor); } - outcome + (outcome.0, outcome.1, analytics_result) } async fn append_guardian_followup_reminder(review_session: &GuardianReviewSession) { @@ -651,7 +755,7 @@ async fn wait_for_guardian_review( review_session: &GuardianReviewSession, deadline: tokio::time::Instant, external_cancel: Option<&CancellationToken>, -) -> (GuardianReviewSessionOutcome, bool) { +) -> (GuardianReviewSessionOutcome, bool, bool) { let timeout = tokio::time::sleep_until(deadline); tokio::pin!(timeout); let mut last_error_message: Option = None; @@ -660,7 +764,7 @@ async fn wait_for_guardian_review( tokio::select! { _ = &mut timeout => { let keep_review_session = interrupt_and_drain_turn(&review_session.codex).await.is_ok(); - return (GuardianReviewSessionOutcome::TimedOut, keep_review_session); + return (GuardianReviewSessionOutcome::TimedOut, keep_review_session, false); } _ = async { if let Some(cancel_token) = external_cancel { @@ -670,7 +774,7 @@ async fn wait_for_guardian_review( } } => { let keep_review_session = interrupt_and_drain_turn(&review_session.codex).await.is_ok(); - return (GuardianReviewSessionOutcome::Aborted, keep_review_session); + return (GuardianReviewSessionOutcome::Aborted, keep_review_session, false); } event = review_session.codex.next_event() => { match event { @@ -682,18 +786,20 @@ async fn wait_for_guardian_review( return ( GuardianReviewSessionOutcome::Completed(Err(anyhow!(error_message))), true, + true, ); } return ( GuardianReviewSessionOutcome::Completed(Ok(turn_complete.last_agent_message)), true, + true, ); } EventMsg::Error(error) => { last_error_message = Some(error.message); } EventMsg::TurnAborted(_) => { - return (GuardianReviewSessionOutcome::Aborted, true); + return (GuardianReviewSessionOutcome::Aborted, true, false); } _ => {} }, @@ -701,6 +807,7 @@ async fn wait_for_guardian_review( return ( GuardianReviewSessionOutcome::Completed(Err(err.into())), false, + false, ); } } @@ -999,4 +1106,44 @@ mod tests { assert_eq!(outcome.unwrap(), 42); assert!(!cancel_token.is_cancelled()); } + + #[test] + fn had_prior_review_context_tracks_prompt_mode() { + assert!(!had_prior_review_context(&GuardianPromptMode::Full)); + assert!(had_prior_review_context(&GuardianPromptMode::Delta { + cursor: GuardianTranscriptCursor { + parent_history_version: 7, + transcript_entry_count: 42, + } + })); + } + + #[test] + fn token_usage_delta_never_reports_negative_usage() { + let start = TokenUsage { + input_tokens: 10, + cached_input_tokens: 8, + output_tokens: 6, + reasoning_output_tokens: 4, + total_tokens: 28, + }; + let end = TokenUsage { + input_tokens: 15, + cached_input_tokens: 7, + output_tokens: 10, + reasoning_output_tokens: 2, + total_tokens: 34, + }; + + assert_eq!( + token_usage_delta(&start, &end), + TokenUsage { + input_tokens: 5, + cached_input_tokens: 0, + output_tokens: 4, + reasoning_output_tokens: 0, + total_tokens: 6, + } + ); + } } diff --git a/codex-rs/core/src/guardian/tests.rs b/codex-rs/core/src/guardian/tests.rs index 167cf9a15..9ce265e91 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -15,6 +15,7 @@ use crate::config_loader::Sourced; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::test_support; +use codex_analytics::GuardianApprovalRequestSource; use codex_config::config_toml::ConfigToml; use codex_config::types::McpServerConfig; use codex_exec_server::LOCAL_FS; @@ -707,6 +708,7 @@ async fn cancelled_guardian_review_emits_terminal_abort_without_warning() { .to_string(), }, /*retry_reason*/ None, + GuardianApprovalRequestSource::MainTurn, cancel_token, ) .await; @@ -994,7 +996,7 @@ async fn guardian_review_request_layout_matches_model_visible_request_snapshot() /*external_cancel*/ None, ) .await; - let GuardianReviewOutcome::Completed(Ok(assessment)) = outcome else { + let (GuardianReviewOutcome::Completed(assessment), _) = outcome else { panic!("expected guardian assessment"); }; assert_eq!(assessment.outcome, GuardianAssessmentOutcome::Allow); @@ -1233,13 +1235,13 @@ async fn guardian_reuses_prompt_cache_key_and_appends_prior_reviews() -> anyhow: ) .await; - let GuardianReviewOutcome::Completed(Ok(first_assessment)) = first_outcome else { + let (GuardianReviewOutcome::Completed(first_assessment), _) = first_outcome else { panic!("expected first guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(second_assessment)) = second_outcome else { + let (GuardianReviewOutcome::Completed(second_assessment), _) = second_outcome else { panic!("expected second guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(third_assessment)) = third_outcome else { + let (GuardianReviewOutcome::Completed(third_assessment), _) = third_outcome else { panic!("expected third guardian assessment"); }; assert_eq!(first_assessment.outcome, GuardianAssessmentOutcome::Allow); diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index d45357e3e..2efb59cdf 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -1919,6 +1919,7 @@ impl Session { review_id, request, /*retry_reason*/ None, + codex_analytics::GuardianApprovalRequestSource::MainTurn, cancellation_token.clone(), ); let decision = tokio::select! {