diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index 25ab40953f..761812cd1e 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 b542f8268f..86db6331f3 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)] @@ -235,6 +239,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 5c4cdfac79..ed0f1036ca 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 5d67223b33..b30cd18b77 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_exec_server::EnvironmentManager; use codex_protocol::protocol::ApplyPatchApprovalRequestEvent; @@ -753,6 +754,7 @@ fn spawn_guardian_review( review_id, request, retry_reason, + GuardianApprovalRequestSource::DelegatedSubagent, cancel_token, )); let _ = tx.send(decision); diff --git a/codex-rs/core/src/guardian/approval_request.rs b/codex-rs/core/src/guardian/approval_request.rs index 6d1d3f76af..9e81644eff 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; @@ -355,6 +356,63 @@ 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(), + }, + } +} + 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 ea52114a09..661fe275da 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; @@ -16,6 +23,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; @@ -24,6 +32,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; @@ -75,11 +84,35 @@ pub(crate) fn guardian_timeout_message() -> String { #[derive(Debug)] pub(super) enum GuardianReviewOutcome { - Completed(anyhow::Result), + Completed(GuardianAssessment), + Failed(GuardianReviewFailure), TimedOut, Aborted, } +#[derive(Debug)] +pub(super) enum GuardianReviewFailure { + PromptBuild(anyhow::Error), + Session(anyhow::Error), + Parse(anyhow::Error), +} + +impl GuardianReviewFailure { + fn reason(&self) -> GuardianReviewFailureReason { + match self { + Self::PromptBuild(_) => GuardianReviewFailureReason::PromptBuildError, + Self::Session(_) => GuardianReviewFailureReason::SessionError, + Self::Parse(_) => GuardianReviewFailureReason::ParseError, + } + } + + fn error(&self) -> &anyhow::Error { + match self { + Self::PromptBuild(err) | Self::Session(err) | Self::Parse(err) => err, + } + } +} + fn guardian_risk_level_str(level: GuardianRiskLevel) -> &'static str { match level { GuardianRiskLevel::Low => "low", @@ -107,6 +140,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. @@ -116,11 +164,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(), @@ -142,6 +200,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(), @@ -163,28 +232,78 @@ 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::Failed(failure) => { + let rationale = format!("Automatic approval review failed: {}", failure.error()); + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(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(); + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::TimedOut, + failure_reason: Some(GuardianReviewFailureReason::Timeout), + ..analytics_result + }, + ); session .send_event( turn.as_ref(), @@ -212,6 +331,17 @@ async fn run_guardian_review( return ReviewDecision::TimedOut; } GuardianReviewOutcome::Aborted => { + track_guardian_review( + session.as_ref(), + turn.as_ref(), + &review_tracking, + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + ..analytics_result + }, + ); session .send_event( turn.as_ref(), @@ -311,6 +441,7 @@ pub(crate) async fn review_approval_request( review_id, request, retry_reason, + GuardianApprovalRequestSource::MainTurn, /*external_cancel*/ None, )) .await @@ -322,16 +453,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 } @@ -356,11 +489,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::Failed(GuardianReviewFailure::PromptBuild(err)), + GuardianReviewAnalyticsResult::without_session(), + ); + } }, None => None, }; @@ -410,10 +548,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::Failed(GuardianReviewFailure::PromptBuild(err)), + GuardianReviewAnalyticsResult::without_session(), + ); + } }; - match Box::pin( + let (session_outcome, session_analytics_result) = Box::pin( session .guardian_review_session .run_review(GuardianReviewSessionParams { @@ -430,17 +573,70 @@ 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(), - )) + .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::Failed(GuardianReviewFailure::Parse(err)), + session_analytics_result, + ), + } + } + None => ( + GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(anyhow::anyhow!( + "guardian review completed without an assessment payload" + ))), + session_analytics_result, + ), + }, + GuardianReviewSessionOutcome::Completed(Err(err)) => ( + GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(err)), + session_analytics_result, + ), + GuardianReviewSessionOutcome::PromptBuildFailed(err) => ( + GuardianReviewOutcome::Failed(GuardianReviewFailure::PromptBuild(err)), + session_analytics_result, + ), + GuardianReviewSessionOutcome::TimedOut => { + (GuardianReviewOutcome::TimedOut, session_analytics_result) } - GuardianReviewSessionOutcome::Completed(Err(err)) => { - GuardianReviewOutcome::Completed(Err(err)) + GuardianReviewSessionOutcome::Aborted => { + (GuardianReviewOutcome::Aborted, session_analytics_result) } - GuardianReviewSessionOutcome::TimedOut => GuardianReviewOutcome::TimedOut, - GuardianReviewSessionOutcome::Aborted => GuardianReviewOutcome::Aborted, + } +} + +#[cfg(test)] +mod review_tests { + use super::*; + + #[test] + fn guardian_review_failure_reason_distinguishes_failure_kinds() { + let parse_failure = GuardianReviewFailure::Parse(anyhow::anyhow!("bad guardian JSON")); + let prompt_failure = + GuardianReviewFailure::PromptBuild(anyhow::anyhow!("bad prompt/config")); + let session_failure = + GuardianReviewFailure::Session(anyhow::anyhow!("guardian runtime failed")); + + assert!(matches!( + parse_failure.reason(), + GuardianReviewFailureReason::ParseError + )); + assert!(matches!( + prompt_failure.reason(), + GuardianReviewFailureReason::PromptBuildError + )); + assert!(matches!( + session_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 5bb3e9d1b8..bb8a14cc49 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::DeveloperInstructions; @@ -17,6 +19,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; @@ -59,6 +62,7 @@ const GUARDIAN_FOLLOWUP_REVIEW_REMINDER: &str = concat!( #[derive(Debug)] pub(crate) enum GuardianReviewSessionOutcome { Completed(anyhow::Result>), + PromptBuildFailed(anyhow::Error), TimedOut, Aborted, } @@ -102,6 +106,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>, @@ -265,13 +284,14 @@ impl GuardianReviewSessionManager { } } - 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(), @@ -305,16 +325,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 { @@ -322,9 +348,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 { @@ -350,20 +379,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) } } @@ -460,7 +499,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; @@ -479,17 +518,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; @@ -497,7 +542,7 @@ impl GuardianReviewSessionManager { cleanup.disarm(); review_session.shutdown_in_background(); } - outcome + (outcome, analytics_result) } } @@ -544,8 +589,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; @@ -560,6 +610,29 @@ 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; } @@ -584,6 +657,12 @@ async fn run_review_on_session( prompt_mode, ) .await?; + let token_usage_at_review_start = review_session + .codex + .session + .total_token_usage() + .await + .unwrap_or_default(); review_session .codex @@ -603,29 +682,44 @@ async fn run_review_on_session( }) .await?; - Ok::(prompt_items.transcript_cursor) + Ok::<(GuardianTranscriptCursor, TokenUsage), anyhow::Error>(( + prompt_items.transcript_cursor, + token_usage_at_review_start, + )) }), ) .await; let submit_result = match submit_result { Ok(submit_result) => submit_result, - Err(outcome) => return (outcome, false), + Err(outcome) => return (outcome, false, analytics_result), }; - let transcript_cursor = match submit_result { - Ok(transcript_cursor) => transcript_cursor, + let (transcript_cursor, token_usage_at_review_start) = match submit_result { + Ok(submit_result) => submit_result, Err(err) => { - return (GuardianReviewSessionOutcome::Completed(Err(err)), false); + return ( + GuardianReviewSessionOutcome::PromptBuildFailed(err), + 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) { @@ -654,7 +748,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; @@ -663,7 +757,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 { @@ -673,7 +767,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 { @@ -685,18 +779,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); } _ => {} }, @@ -704,6 +800,7 @@ async fn wait_for_guardian_review( return ( GuardianReviewSessionOutcome::Completed(Err(err.into())), false, + false, ); } } @@ -1001,4 +1098,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 098e0f047a..2848a0dc79 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; @@ -706,6 +707,7 @@ async fn cancelled_guardian_review_emits_terminal_abort_without_warning() { .to_string(), }, /*retry_reason*/ None, + GuardianApprovalRequestSource::MainTurn, cancel_token, ) .await; @@ -921,7 +923,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); @@ -1130,13 +1132,14 @@ 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);