diff --git a/codex-rs/core/src/guardian/metrics.rs b/codex-rs/core/src/guardian/metrics.rs new file mode 100644 index 000000000..9b9d35a52 --- /dev/null +++ b/codex-rs/core/src/guardian/metrics.rs @@ -0,0 +1,418 @@ +use std::time::Duration; + +use codex_analytics::GuardianApprovalRequestSource; +use codex_analytics::GuardianReviewAnalyticsResult; +use codex_analytics::GuardianReviewDecision; +use codex_analytics::GuardianReviewFailureReason; +use codex_analytics::GuardianReviewSessionKind; +use codex_analytics::GuardianReviewTerminalStatus; +use codex_analytics::GuardianReviewedAction; +use codex_otel::GUARDIAN_REVIEW_COUNT_METRIC; +use codex_otel::GUARDIAN_REVIEW_DURATION_METRIC; +use codex_otel::GUARDIAN_REVIEW_TOKEN_USAGE_METRIC; +use codex_otel::GUARDIAN_REVIEW_TTFT_DURATION_METRIC; +use codex_otel::SessionTelemetry; +use codex_otel::sanitize_metric_tag_value; +use codex_protocol::protocol::GuardianAssessmentOutcome; +use codex_protocol::protocol::GuardianRiskLevel; +use codex_protocol::protocol::GuardianUserAuthorization; +use codex_protocol::protocol::TokenUsage; + +pub(crate) fn emit_guardian_review_metrics( + session_telemetry: &SessionTelemetry, + result: &GuardianReviewAnalyticsResult, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: &GuardianReviewedAction, + completion_latency_ms: u64, +) { + let tags = guardian_review_metric_tags(result, approval_request_source, reviewed_action); + let tag_refs: Vec<(&str, &str)> = tags + .iter() + .map(|(key, value)| (*key, value.as_str())) + .collect(); + + session_telemetry.counter(GUARDIAN_REVIEW_COUNT_METRIC, /*inc*/ 1, &tag_refs); + session_telemetry.record_duration( + GUARDIAN_REVIEW_DURATION_METRIC, + Duration::from_millis(completion_latency_ms), + &tag_refs, + ); + + if let Some(time_to_first_token_ms) = result.time_to_first_token_ms { + session_telemetry.record_duration( + GUARDIAN_REVIEW_TTFT_DURATION_METRIC, + Duration::from_millis(time_to_first_token_ms), + &tag_refs, + ); + } + + if let Some(token_usage) = result.token_usage.as_ref() { + emit_guardian_token_usage_histograms(session_telemetry, token_usage, tags); + } +} + +fn emit_guardian_token_usage_histograms( + session_telemetry: &SessionTelemetry, + token_usage: &TokenUsage, + base_tags: Vec<(&'static str, String)>, +) { + for (token_type, value) in [ + ("total", token_usage.total_tokens.max(0)), + ("input", token_usage.input_tokens.max(0)), + ("cached_input", token_usage.cached_input()), + ("non_cached_input", token_usage.non_cached_input()), + ("output", token_usage.output_tokens.max(0)), + ( + "reasoning_output", + token_usage.reasoning_output_tokens.max(0), + ), + ] { + let mut tags = base_tags.clone(); + tags.push(("token_type", token_type.to_string())); + let tag_refs: Vec<(&str, &str)> = tags + .iter() + .map(|(key, value)| (*key, value.as_str())) + .collect(); + session_telemetry.histogram(GUARDIAN_REVIEW_TOKEN_USAGE_METRIC, value, &tag_refs); + } +} + +fn guardian_review_metric_tags( + result: &GuardianReviewAnalyticsResult, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: &GuardianReviewedAction, +) -> Vec<(&'static str, String)> { + vec![ + ("decision", decision_tag(result.decision).to_string()), + ( + "terminal_status", + terminal_status_tag(result.terminal_status).to_string(), + ), + ( + "failure_reason", + failure_reason_tag(result.failure_reason).to_string(), + ), + ( + "approval_request_source", + approval_request_source_tag(approval_request_source).to_string(), + ), + ("action", reviewed_action_tag(reviewed_action).to_string()), + ( + "session_kind", + session_kind_tag(result.guardian_session_kind).to_string(), + ), + ( + "had_prior_review_context", + optional_bool_tag(result.had_prior_review_context).to_string(), + ), + ( + "reviewed_action_truncated", + bool_tag(result.reviewed_action_truncated).to_string(), + ), + ("risk_level", risk_level_tag(result.risk_level).to_string()), + ( + "user_authorization", + user_authorization_tag(result.user_authorization).to_string(), + ), + ("outcome", outcome_tag(result.outcome).to_string()), + ( + "guardian_model", + result + .guardian_model + .as_deref() + .map(sanitize_metric_tag_value) + .unwrap_or_else(|| "none".to_string()), + ), + ( + "guardian_reasoning_effort", + result + .guardian_reasoning_effort + .as_deref() + .map(sanitize_metric_tag_value) + .unwrap_or_else(|| "none".to_string()), + ), + ] +} + +fn decision_tag(decision: GuardianReviewDecision) -> &'static str { + match decision { + GuardianReviewDecision::Approved => "approved", + GuardianReviewDecision::Denied => "denied", + GuardianReviewDecision::Aborted => "aborted", + } +} + +fn terminal_status_tag(status: GuardianReviewTerminalStatus) -> &'static str { + match status { + GuardianReviewTerminalStatus::Approved => "approved", + GuardianReviewTerminalStatus::Denied => "denied", + GuardianReviewTerminalStatus::Aborted => "aborted", + GuardianReviewTerminalStatus::TimedOut => "timed_out", + GuardianReviewTerminalStatus::FailedClosed => "failed_closed", + } +} + +fn failure_reason_tag(reason: Option) -> &'static str { + match reason { + Some(GuardianReviewFailureReason::Timeout) => "timeout", + Some(GuardianReviewFailureReason::Cancelled) => "cancelled", + Some(GuardianReviewFailureReason::PromptBuildError) => "prompt_build_error", + Some(GuardianReviewFailureReason::SessionError) => "session_error", + Some(GuardianReviewFailureReason::ParseError) => "parse_error", + None => "none", + } +} + +fn approval_request_source_tag(source: GuardianApprovalRequestSource) -> &'static str { + match source { + GuardianApprovalRequestSource::MainTurn => "main_turn", + GuardianApprovalRequestSource::DelegatedSubagent => "delegated_subagent", + } +} + +fn reviewed_action_tag(action: &GuardianReviewedAction) -> &'static str { + match action { + GuardianReviewedAction::Shell { .. } => "shell", + GuardianReviewedAction::UnifiedExec { .. } => "unified_exec", + GuardianReviewedAction::Execve { .. } => "execve", + GuardianReviewedAction::ApplyPatch {} => "apply_patch", + GuardianReviewedAction::NetworkAccess { .. } => "network_access", + GuardianReviewedAction::McpToolCall { .. } => "mcp_tool_call", + GuardianReviewedAction::RequestPermissions {} => "request_permissions", + } +} + +fn session_kind_tag(kind: Option) -> &'static str { + match kind { + Some(GuardianReviewSessionKind::TrunkNew) => "trunk_new", + Some(GuardianReviewSessionKind::TrunkReused) => "trunk_reused", + Some(GuardianReviewSessionKind::EphemeralForked) => "ephemeral_forked", + None => "none", + } +} + +fn optional_bool_tag(value: Option) -> &'static str { + match value { + Some(true) => "true", + Some(false) => "false", + None => "unknown", + } +} + +fn bool_tag(value: bool) -> &'static str { + if value { "true" } else { "false" } +} + +fn risk_level_tag(risk_level: Option) -> &'static str { + match risk_level { + Some(GuardianRiskLevel::Low) => "low", + Some(GuardianRiskLevel::Medium) => "medium", + Some(GuardianRiskLevel::High) => "high", + Some(GuardianRiskLevel::Critical) => "critical", + None => "none", + } +} + +fn user_authorization_tag(user_authorization: Option) -> &'static str { + match user_authorization { + Some(GuardianUserAuthorization::Unknown) => "unknown", + Some(GuardianUserAuthorization::Low) => "low", + Some(GuardianUserAuthorization::Medium) => "medium", + Some(GuardianUserAuthorization::High) => "high", + None => "none", + } +} + +fn outcome_tag(outcome: Option) -> &'static str { + match outcome { + Some(GuardianAssessmentOutcome::Allow) => "allow", + Some(GuardianAssessmentOutcome::Deny) => "deny", + None => "none", + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use codex_otel::MetricsClient; + use codex_otel::MetricsConfig; + use codex_protocol::ThreadId; + use codex_protocol::protocol::SessionSource; + use opentelemetry::KeyValue; + use opentelemetry_sdk::metrics::InMemoryMetricExporter; + use opentelemetry_sdk::metrics::data::AggregatedMetrics; + use opentelemetry_sdk::metrics::data::Metric; + use opentelemetry_sdk::metrics::data::MetricData; + use opentelemetry_sdk::metrics::data::ResourceMetrics; + use pretty_assertions::assert_eq; + use std::collections::BTreeMap; + + fn test_session_telemetry() -> SessionTelemetry { + let exporter = InMemoryMetricExporter::default(); + let metrics = MetricsClient::new( + MetricsConfig::in_memory("test", "codex-core", env!("CARGO_PKG_VERSION"), exporter) + .with_runtime_reader(), + ) + .expect("in-memory metrics client"); + SessionTelemetry::new( + ThreadId::new(), + "gpt-5.4", + "gpt-5.4", + /*account_id*/ None, + /*account_email*/ None, + /*auth_mode*/ None, + "test_originator".to_string(), + /*log_user_prompts*/ false, + "tty".to_string(), + SessionSource::Cli, + ) + .with_metrics_without_metadata_tags(metrics) + } + + fn find_metric<'a>(resource_metrics: &'a ResourceMetrics, name: &str) -> &'a Metric { + for scope_metrics in resource_metrics.scope_metrics() { + for metric in scope_metrics.metrics() { + if metric.name() == name { + return metric; + } + } + } + panic!("metric {name} missing"); + } + + fn attributes_to_map<'a>( + attributes: impl Iterator, + ) -> BTreeMap { + attributes + .map(|kv| (kv.key.as_str().to_string(), kv.value.as_str().to_string())) + .collect() + } + + fn counter_point( + resource_metrics: &ResourceMetrics, + name: &str, + ) -> (BTreeMap, u64) { + let metric = find_metric(resource_metrics, name); + match metric.data() { + AggregatedMetrics::U64(data) => match data { + MetricData::Sum(sum) => { + let points: Vec<_> = sum.data_points().collect(); + assert_eq!(points.len(), 1); + let point = points[0]; + (attributes_to_map(point.attributes()), point.value()) + } + _ => panic!("unexpected counter aggregation"), + }, + _ => panic!("unexpected counter data type"), + } + } + + fn histogram_sums(resource_metrics: &ResourceMetrics, name: &str) -> BTreeMap { + let metric = find_metric(resource_metrics, name); + match metric.data() { + AggregatedMetrics::F64(data) => match data { + MetricData::Histogram(histogram) => histogram + .data_points() + .map(|point| { + let attrs = attributes_to_map(point.attributes()); + ( + attrs + .get("token_type") + .cloned() + .unwrap_or_else(|| "sample".to_string()), + point.sum() as u64, + ) + }) + .collect(), + _ => panic!("unexpected histogram aggregation"), + }, + _ => panic!("unexpected histogram data type"), + } + } + + #[test] + fn guardian_review_metrics_record_counts_durations_and_token_usage() { + let session_telemetry = test_session_telemetry(); + let result = GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Approved, + terminal_status: GuardianReviewTerminalStatus::Approved, + risk_level: Some(GuardianRiskLevel::Low), + user_authorization: Some(GuardianUserAuthorization::High), + outcome: Some(GuardianAssessmentOutcome::Allow), + guardian_session_kind: Some(GuardianReviewSessionKind::TrunkReused), + guardian_model: Some("gpt-5.4 guardian".to_string()), + guardian_reasoning_effort: Some("low".to_string()), + had_prior_review_context: Some(true), + reviewed_action_truncated: true, + token_usage: Some(TokenUsage { + input_tokens: 10, + cached_input_tokens: 4, + output_tokens: 3, + reasoning_output_tokens: 2, + total_tokens: 15, + }), + time_to_first_token_ms: Some(123), + ..GuardianReviewAnalyticsResult::without_session() + }; + + emit_guardian_review_metrics( + &session_telemetry, + &result, + GuardianApprovalRequestSource::DelegatedSubagent, + &GuardianReviewedAction::NetworkAccess { + protocol: codex_protocol::approvals::NetworkApprovalProtocol::Https, + port: 443, + }, + /*completion_latency_ms*/ 456, + ); + + let snapshot = session_telemetry + .snapshot_metrics() + .expect("runtime metrics snapshot"); + let (attrs, value) = counter_point(&snapshot, GUARDIAN_REVIEW_COUNT_METRIC); + + assert_eq!(value, 1); + assert_eq!( + attrs, + BTreeMap::from([ + ("action".to_string(), "network_access".to_string()), + ( + "approval_request_source".to_string(), + "delegated_subagent".to_string() + ), + ("decision".to_string(), "approved".to_string()), + ("failure_reason".to_string(), "none".to_string()), + ("guardian_model".to_string(), "gpt-5.4_guardian".to_string()), + ("guardian_reasoning_effort".to_string(), "low".to_string()), + ("had_prior_review_context".to_string(), "true".to_string()), + ("outcome".to_string(), "allow".to_string()), + ("reviewed_action_truncated".to_string(), "true".to_string()), + ("risk_level".to_string(), "low".to_string()), + ("session_kind".to_string(), "trunk_reused".to_string()), + ("terminal_status".to_string(), "approved".to_string()), + ("user_authorization".to_string(), "high".to_string()), + ]) + ); + + assert_eq!( + histogram_sums(&snapshot, GUARDIAN_REVIEW_TOKEN_USAGE_METRIC), + BTreeMap::from([ + ("cached_input".to_string(), 4), + ("input".to_string(), 10), + ("non_cached_input".to_string(), 6), + ("output".to_string(), 3), + ("reasoning_output".to_string(), 2), + ("total".to_string(), 15), + ]) + ); + assert_eq!( + histogram_sums(&snapshot, GUARDIAN_REVIEW_DURATION_METRIC), + BTreeMap::from([("sample".to_string(), 456)]) + ); + assert_eq!( + histogram_sums(&snapshot, GUARDIAN_REVIEW_TTFT_DURATION_METRIC), + BTreeMap::from([("sample".to_string(), 123)]) + ); + } +} diff --git a/codex-rs/core/src/guardian/mod.rs b/codex-rs/core/src/guardian/mod.rs index cd78d9da7..f5c6fe523 100644 --- a/codex-rs/core/src/guardian/mod.rs +++ b/codex-rs/core/src/guardian/mod.rs @@ -12,6 +12,7 @@ //! 4. Apply the guardian's explicit allow/deny outcome. mod approval_request; +mod metrics; mod prompt; mod review; mod review_session; diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index 7df7a9692..db4343448 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -4,6 +4,7 @@ use codex_analytics::GuardianReviewDecision; use codex_analytics::GuardianReviewFailureReason; use codex_analytics::GuardianReviewTerminalStatus; use codex_analytics::GuardianReviewTrackContext; +use codex_analytics::GuardianReviewedAction; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; @@ -36,6 +37,7 @@ 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::metrics::emit_guardian_review_metrics; use super::prompt::guardian_output_schema; use super::prompt::parse_guardian_assessment; use super::review_session::GuardianReviewSessionOutcome; @@ -162,9 +164,18 @@ pub(crate) fn is_guardian_reviewer_source( fn track_guardian_review( session: &Session, tracking: &GuardianReviewTrackContext, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: &GuardianReviewedAction, result: GuardianReviewAnalyticsResult, completed_at_ms: u64, ) { + emit_guardian_review_metrics( + &session.services.session_telemetry, + &result, + approval_request_source, + reviewed_action, + completed_at_ms.saturating_sub(tracking.started_at_ms), + ); session .services .analytics_events_client @@ -244,13 +255,14 @@ async fn run_guardian_review( 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 reviewed_action = guardian_reviewed_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), + reviewed_action.clone(), GUARDIAN_REVIEW_TIMEOUT.as_millis() as u64, ); let started_at_ms = review_tracking.started_at_ms.try_into().unwrap_or_default(); @@ -281,6 +293,8 @@ async fn run_guardian_review( track_guardian_review( session.as_ref(), &review_tracking, + approval_request_source, + &reviewed_action, GuardianReviewAnalyticsResult { decision: GuardianReviewDecision::Aborted, terminal_status: GuardianReviewTerminalStatus::Aborted, @@ -330,6 +344,8 @@ async fn run_guardian_review( track_guardian_review( session.as_ref(), &review_tracking, + approval_request_source, + &reviewed_action, GuardianReviewAnalyticsResult { decision: if approved { GuardianReviewDecision::Approved @@ -361,6 +377,8 @@ async fn run_guardian_review( track_guardian_review( session.as_ref(), &review_tracking, + approval_request_source, + &reviewed_action, GuardianReviewAnalyticsResult { decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::TimedOut, @@ -402,6 +420,8 @@ async fn run_guardian_review( track_guardian_review( session.as_ref(), &review_tracking, + approval_request_source, + &reviewed_action, GuardianReviewAnalyticsResult { decision: GuardianReviewDecision::Aborted, terminal_status: GuardianReviewTerminalStatus::Aborted, @@ -446,6 +466,8 @@ async fn run_guardian_review( track_guardian_review( session.as_ref(), &review_tracking, + approval_request_source, + &reviewed_action, GuardianReviewAnalyticsResult { decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::FailedClosed, diff --git a/codex-rs/otel/src/metrics/names.rs b/codex-rs/otel/src/metrics/names.rs index 817545c87..d2093ed3f 100644 --- a/codex-rs/otel/src/metrics/names.rs +++ b/codex-rs/otel/src/metrics/names.rs @@ -28,6 +28,10 @@ pub const TURN_NETWORK_PROXY_METRIC: &str = "codex.turn.network_proxy"; pub const TURN_MEMORY_METRIC: &str = "codex.turn.memory"; pub const TURN_TOOL_CALL_METRIC: &str = "codex.turn.tool.call"; pub const TURN_TOKEN_USAGE_METRIC: &str = "codex.turn.token_usage"; +pub const GUARDIAN_REVIEW_COUNT_METRIC: &str = "codex.guardian.review"; +pub const GUARDIAN_REVIEW_DURATION_METRIC: &str = "codex.guardian.review.duration_ms"; +pub const GUARDIAN_REVIEW_TTFT_DURATION_METRIC: &str = "codex.guardian.review.ttft.duration_ms"; +pub const GUARDIAN_REVIEW_TOKEN_USAGE_METRIC: &str = "codex.guardian.review.token_usage"; pub const GOAL_CREATED_METRIC: &str = "codex.goal.created"; pub const GOAL_RESUMED_METRIC: &str = "codex.goal.resumed"; pub const GOAL_COMPLETED_METRIC: &str = "codex.goal.completed";