From 044b9ef7643e43a58fadf31958c2766bead02d6a Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Wed, 15 Apr 2026 15:03:50 -0700 Subject: [PATCH 1/4] [codex-analytics] guardian review analytics events emission --- codex-rs/core/src/codex_delegate.rs | 2 + codex-rs/core/src/guardian/review.rs | 401 +++++++++++++++++++++++++-- codex-rs/core/src/guardian/tests.rs | 1 + 3 files changed, 386 insertions(+), 18 deletions(-) diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index 9cd58e044f..db18081785 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/review.rs b/codex-rs/core/src/guardian/review.rs index 18107f7377..fb78156b16 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -1,5 +1,14 @@ use std::sync::Arc; +use std::time::Instant; +use codex_analytics::GuardianApprovalRequestSource; +use codex_analytics::GuardianReviewDecision; +use codex_analytics::GuardianReviewFailureReason; +use codex_analytics::GuardianReviewSessionKind; +use codex_analytics::GuardianReviewTerminalStatus; +use codex_analytics::GuardianReviewedAction; +use codex_analytics::now_unix_seconds; +use codex_features::Feature; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; @@ -10,12 +19,14 @@ use codex_protocol::protocol::GuardianRiskLevel; use codex_protocol::protocol::GuardianUserAuthorization; use codex_protocol::protocol::ReviewDecision; use codex_protocol::protocol::SubAgentSource; +use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::WarningEvent; use tokio_util::sync::CancellationToken; use crate::codex::Session; use crate::codex::TurnContext; +use super::GUARDIAN_REVIEW_TIMEOUT; use super::GUARDIAN_REVIEWER_NAME; use super::GuardianApprovalRequest; use super::GuardianAssessment; @@ -76,10 +87,34 @@ pub(crate) fn guardian_timeout_message() -> String { #[derive(Debug)] pub(super) enum GuardianReviewOutcome { Completed(anyhow::Result), + 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", @@ -89,6 +124,188 @@ fn guardian_risk_level_str(level: GuardianRiskLevel) -> &'static str { } } +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(), + }, + } +} + +struct GuardianReviewAnalyticsContext { + thread_id: String, + turn_id: String, + review_id: String, + target_item_id: Option, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: GuardianReviewedAction, + started_at: u64, + started_instant: Instant, +} + +struct GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision, + terminal_status: GuardianReviewTerminalStatus, + failure_reason: Option, + risk_level: Option, + user_authorization: Option, + outcome: Option, + guardian_thread_id: Option, + guardian_session_kind: Option, + guardian_model: Option, + guardian_reasoning_effort: Option, + had_prior_review_context: Option, + reviewed_action_truncated: bool, + token_usage: Option, + time_to_first_token_ms: Option, + completed_at: u64, +} + +#[derive(Default)] +struct GuardianReviewMetadataFields { + guardian_thread_id: Option, + guardian_session_kind: Option, + guardian_model: Option, + guardian_reasoning_effort: Option, + had_prior_review_context: Option, + reviewed_action_truncated: bool, + token_usage: Option, + time_to_first_token_ms: Option, +} + +impl GuardianReviewAnalyticsResult { + fn from_metadata(metadata: GuardianReviewMetadataFields, completed_at: u64) -> Self { + Self { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: None, + risk_level: None, + user_authorization: None, + outcome: None, + guardian_thread_id: metadata.guardian_thread_id, + guardian_session_kind: metadata.guardian_session_kind, + guardian_model: metadata.guardian_model, + guardian_reasoning_effort: metadata.guardian_reasoning_effort, + had_prior_review_context: metadata.had_prior_review_context, + reviewed_action_truncated: metadata.reviewed_action_truncated, + token_usage: metadata.token_usage, + time_to_first_token_ms: metadata.time_to_first_token_ms, + completed_at, + } + } +} + +impl GuardianReviewAnalyticsContext { + fn track( + &self, + session: &Session, + turn: &TurnContext, + terminal: GuardianReviewAnalyticsResult, + ) { + if !turn.config.features.enabled(Feature::GeneralAnalytics) { + return; + } + let completion_latency_ms = self.started_instant.elapsed().as_millis() as u64; + session + .services + .analytics_events_client + .track_guardian_review(codex_analytics::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: terminal.reviewed_action_truncated, + decision: terminal.decision, + terminal_status: terminal.terminal_status, + failure_reason: terminal.failure_reason, + risk_level: terminal.risk_level, + user_authorization: terminal.user_authorization, + outcome: terminal.outcome, + guardian_thread_id: terminal.guardian_thread_id, + guardian_session_kind: terminal.guardian_session_kind, + guardian_model: terminal.guardian_model, + guardian_reasoning_effort: terminal.guardian_reasoning_effort, + had_prior_review_context: terminal.had_prior_review_context, + review_timeout_ms: GUARDIAN_REVIEW_TIMEOUT.as_millis() as u64, + // TODO(rhan-oai): plumb nested Guardian review session tool-call counts. + tool_call_count: None, + time_to_first_token_ms: terminal.time_to_first_token_ms, + completion_latency_ms: Some(completion_latency_ms), + started_at: self.started_at, + completed_at: Some(terminal.completed_at), + input_tokens: terminal + .token_usage + .as_ref() + .map(|usage| usage.input_tokens), + cached_input_tokens: terminal + .token_usage + .as_ref() + .map(|usage| usage.cached_input_tokens), + output_tokens: terminal + .token_usage + .as_ref() + .map(|usage| usage.output_tokens), + reasoning_output_tokens: terminal + .token_usage + .as_ref() + .map(|usage| usage.reasoning_output_tokens), + total_tokens: terminal + .token_usage + .as_ref() + .map(|usage| usage.total_tokens), + }); + } +} + /// Whether this turn should route `on-request` approval prompts through the /// guardian reviewer instead of surfacing them to the user. ARC may still /// block actions earlier in the flow. @@ -116,11 +333,24 @@ async fn run_guardian_review( review_id: String, request: GuardianApprovalRequest, retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, external_cancel: Option, ) -> ReviewDecision { + let started_at = now_unix_seconds(); + let started_instant = Instant::now(); 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 analytics_context = GuardianReviewAnalyticsContext { + thread_id: session.conversation_id.to_string(), + turn_id: assessment_turn_id.clone(), + review_id: review_id.clone(), + target_item_id: target_item_id.clone(), + approval_request_source, + reviewed_action: guardian_reviewed_action(&request), + started_at, + started_instant, + }; session .send_event( turn.as_ref(), @@ -142,6 +372,19 @@ async fn run_guardian_review( .as_ref() .is_some_and(CancellationToken::is_cancelled) { + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + ..GuardianReviewAnalyticsResult::from_metadata( + GuardianReviewMetadataFields::default(), + now_unix_seconds(), + ) + }, + ); session .send_event( turn.as_ref(), @@ -167,24 +410,97 @@ async fn run_guardian_review( session.clone(), turn.clone(), request, - retry_reason, + retry_reason.clone(), schema, external_cancel, )) .await; + let completed_at = now_unix_seconds(); + let terminal = || { + GuardianReviewAnalyticsResult::from_metadata( + GuardianReviewMetadataFields::default(), + completed_at, + ) + }; 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(Ok(assessment)) => { + let approved = matches!(assessment.outcome, GuardianAssessmentOutcome::Allow); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + 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), + ..terminal() + }, + ); + assessment + } + GuardianReviewOutcome::Completed(Err(err)) => { + let rationale = format!("Automatic approval review failed: {err}"); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(GuardianReviewFailureReason::SessionError), + ..terminal() + }, + ); + GuardianAssessment { + risk_level: GuardianRiskLevel::High, + user_authorization: GuardianUserAuthorization::Unknown, + outcome: GuardianAssessmentOutcome::Deny, + rationale, + } + } + GuardianReviewOutcome::Failed(failure) => { + let rationale = format!("Automatic approval review failed: {}", failure.error()); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(failure.reason()), + ..terminal() + }, + ); + 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(); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::TimedOut, + failure_reason: Some(GuardianReviewFailureReason::Timeout), + ..terminal() + }, + ); session .send_event( turn.as_ref(), @@ -212,6 +528,16 @@ async fn run_guardian_review( return ReviewDecision::TimedOut; } GuardianReviewOutcome::Aborted => { + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsResult { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + ..terminal() + }, + ); session .send_event( turn.as_ref(), @@ -311,6 +637,7 @@ pub(crate) async fn review_approval_request( review_id, request, retry_reason, + GuardianApprovalRequestSource::MainTurn, /*external_cancel*/ None, )) .await @@ -322,16 +649,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 } @@ -360,7 +689,9 @@ pub(super) async fn run_guardian_review_session( 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)); + } }, None => None, }; @@ -410,7 +741,7 @@ 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)), }; match Box::pin( @@ -432,15 +763,49 @@ pub(super) async fn run_guardian_review_session( ) .await { - GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) => { - GuardianReviewOutcome::Completed(parse_guardian_assessment( - last_agent_message.as_deref(), - )) - } + 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(Ok(assessment)), + Err(err) => GuardianReviewOutcome::Failed(GuardianReviewFailure::Parse(err)), + } + } + None => GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(anyhow::anyhow!( + "guardian review completed without an assessment payload" + ))), + }, GuardianReviewSessionOutcome::Completed(Err(err)) => { - GuardianReviewOutcome::Completed(Err(err)) + GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(err)) } 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/tests.rs b/codex-rs/core/src/guardian/tests.rs index 8940753d27..82d4127a2b 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -701,6 +701,7 @@ async fn cancelled_guardian_review_emits_terminal_abort_without_warning() { .to_string(), }, /*retry_reason*/ None, + GuardianApprovalRequestSource::MainTurn, cancel_token, ) .await; From 647ccdbefb3ee8543e463c8aec2c377c27af9bda Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Wed, 15 Apr 2026 16:48:30 -0700 Subject: [PATCH 2/4] [codex-analytics] guardian review thread and token metadata --- codex-rs/core/src/guardian/mod.rs | 2 + codex-rs/core/src/guardian/review.rs | 203 ++++++++++++------- codex-rs/core/src/guardian/review_session.rs | 178 +++++++++++++--- codex-rs/core/src/guardian/tests.rs | 25 ++- 4 files changed, 303 insertions(+), 105 deletions(-) diff --git a/codex-rs/core/src/guardian/mod.rs b/codex-rs/core/src/guardian/mod.rs index cfab943226..f98bc750e8 100644 --- a/codex-rs/core/src/guardian/mod.rs +++ b/codex-rs/core/src/guardian/mod.rs @@ -94,6 +94,8 @@ use prompt::render_guardian_transcript_entries; #[cfg(test)] use review::GuardianReviewOutcome; #[cfg(test)] +use review::GuardianReviewOutcomeKind; +#[cfg(test)] use review::run_guardian_review_session as run_guardian_review_session_for_test; #[cfg(test)] use review_session::build_guardian_review_session_config as build_guardian_review_session_config_for_test; diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index fb78156b16..acab331b41 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -37,6 +37,7 @@ use super::approval_request::guardian_request_target_item_id; use super::approval_request::guardian_request_turn_id; use super::prompt::guardian_output_schema; use super::prompt::parse_guardian_assessment; +use super::review_session::GuardianReviewSessionMetadata; use super::review_session::GuardianReviewSessionOutcome; use super::review_session::GuardianReviewSessionParams; use super::review_session::build_guardian_review_session_config; @@ -85,13 +86,55 @@ pub(crate) fn guardian_timeout_message() -> String { } #[derive(Debug)] -pub(super) enum GuardianReviewOutcome { +pub(super) struct GuardianReviewOutcome { + pub(super) kind: GuardianReviewOutcomeKind, + pub(super) metadata: Option, +} + +#[derive(Debug)] +pub(super) enum GuardianReviewOutcomeKind { Completed(anyhow::Result), Failed(GuardianReviewFailure), TimedOut, Aborted, } +impl GuardianReviewOutcome { + fn completed( + assessment: anyhow::Result, + metadata: Option, + ) -> Self { + Self { + kind: GuardianReviewOutcomeKind::Completed(assessment), + metadata, + } + } + + fn failed( + failure: GuardianReviewFailure, + metadata: Option, + ) -> Self { + Self { + kind: GuardianReviewOutcomeKind::Failed(failure), + metadata, + } + } + + fn timed_out(metadata: Option) -> Self { + Self { + kind: GuardianReviewOutcomeKind::TimedOut, + metadata, + } + } + + fn aborted(metadata: Option) -> Self { + Self { + kind: GuardianReviewOutcomeKind::Aborted, + metadata, + } + } +} + #[derive(Debug)] pub(super) enum GuardianReviewFailure { PromptBuild(anyhow::Error), @@ -208,37 +251,39 @@ struct GuardianReviewAnalyticsResult { completed_at: u64, } -#[derive(Default)] -struct GuardianReviewMetadataFields { - guardian_thread_id: Option, - guardian_session_kind: Option, - guardian_model: Option, - guardian_reasoning_effort: Option, - had_prior_review_context: Option, - reviewed_action_truncated: bool, - token_usage: Option, - time_to_first_token_ms: Option, -} - impl GuardianReviewAnalyticsResult { - fn from_metadata(metadata: GuardianReviewMetadataFields, completed_at: u64) -> Self { - Self { + fn from_session_metadata( + metadata: Option, + completed_at: u64, + ) -> Self { + let mut terminal = Self { decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::FailedClosed, failure_reason: None, risk_level: None, user_authorization: None, outcome: None, - guardian_thread_id: metadata.guardian_thread_id, - guardian_session_kind: metadata.guardian_session_kind, - guardian_model: metadata.guardian_model, - guardian_reasoning_effort: metadata.guardian_reasoning_effort, - had_prior_review_context: metadata.had_prior_review_context, - reviewed_action_truncated: metadata.reviewed_action_truncated, - token_usage: metadata.token_usage, - time_to_first_token_ms: metadata.time_to_first_token_ms, + 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, completed_at, + }; + + if let Some(metadata) = metadata { + terminal.guardian_thread_id = Some(metadata.guardian_thread_id); + terminal.guardian_session_kind = Some(metadata.guardian_session_kind); + terminal.guardian_model = Some(metadata.guardian_model); + terminal.guardian_reasoning_effort = metadata.guardian_reasoning_effort; + terminal.had_prior_review_context = Some(metadata.had_prior_review_context); + terminal.token_usage = metadata.token_usage; } + + terminal } } @@ -379,10 +424,7 @@ async fn run_guardian_review( decision: GuardianReviewDecision::Aborted, terminal_status: GuardianReviewTerminalStatus::Aborted, failure_reason: Some(GuardianReviewFailureReason::Cancelled), - ..GuardianReviewAnalyticsResult::from_metadata( - GuardianReviewMetadataFields::default(), - now_unix_seconds(), - ) + ..GuardianReviewAnalyticsResult::from_session_metadata(None, now_unix_seconds()) }, ); session @@ -406,7 +448,7 @@ 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 GuardianReviewOutcome { kind, metadata } = Box::pin(run_guardian_review_session( session.clone(), turn.clone(), request, @@ -417,14 +459,10 @@ async fn run_guardian_review( .await; let completed_at = now_unix_seconds(); - let terminal = || { - GuardianReviewAnalyticsResult::from_metadata( - GuardianReviewMetadataFields::default(), - completed_at, - ) - }; - let assessment = match outcome { - GuardianReviewOutcome::Completed(Ok(assessment)) => { + let terminal = + |metadata| GuardianReviewAnalyticsResult::from_session_metadata(metadata, completed_at); + let assessment = match kind { + GuardianReviewOutcomeKind::Completed(Ok(assessment)) => { let approved = matches!(assessment.outcome, GuardianAssessmentOutcome::Allow); analytics_context.track( session.as_ref(), @@ -444,12 +482,12 @@ async fn run_guardian_review( risk_level: Some(assessment.risk_level), user_authorization: Some(assessment.user_authorization), outcome: Some(assessment.outcome), - ..terminal() + ..terminal(metadata) }, ); assessment } - GuardianReviewOutcome::Completed(Err(err)) => { + GuardianReviewOutcomeKind::Completed(Err(err)) => { let rationale = format!("Automatic approval review failed: {err}"); analytics_context.track( session.as_ref(), @@ -458,7 +496,7 @@ async fn run_guardian_review( decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::FailedClosed, failure_reason: Some(GuardianReviewFailureReason::SessionError), - ..terminal() + ..terminal(metadata) }, ); GuardianAssessment { @@ -468,7 +506,7 @@ async fn run_guardian_review( rationale, } } - GuardianReviewOutcome::Failed(failure) => { + GuardianReviewOutcomeKind::Failed(failure) => { let rationale = format!("Automatic approval review failed: {}", failure.error()); analytics_context.track( session.as_ref(), @@ -477,7 +515,7 @@ async fn run_guardian_review( decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::FailedClosed, failure_reason: Some(failure.reason()), - ..terminal() + ..terminal(metadata) }, ); GuardianAssessment { @@ -487,7 +525,7 @@ async fn run_guardian_review( rationale, } } - GuardianReviewOutcome::TimedOut => { + GuardianReviewOutcomeKind::TimedOut => { let rationale = "Automatic approval review timed out while evaluating the requested approval." .to_string(); @@ -498,7 +536,7 @@ async fn run_guardian_review( decision: GuardianReviewDecision::Denied, terminal_status: GuardianReviewTerminalStatus::TimedOut, failure_reason: Some(GuardianReviewFailureReason::Timeout), - ..terminal() + ..terminal(metadata) }, ); session @@ -527,7 +565,7 @@ async fn run_guardian_review( .await; return ReviewDecision::TimedOut; } - GuardianReviewOutcome::Aborted => { + GuardianReviewOutcomeKind::Aborted => { analytics_context.track( session.as_ref(), turn.as_ref(), @@ -535,7 +573,7 @@ async fn run_guardian_review( decision: GuardianReviewDecision::Aborted, terminal_status: GuardianReviewTerminalStatus::Aborted, failure_reason: Some(GuardianReviewFailureReason::Cancelled), - ..terminal() + ..terminal(metadata) }, ); session @@ -690,7 +728,10 @@ pub(super) async fn run_guardian_review_session( Some(network_proxy) => match network_proxy.proxy().current_cfg().await { Ok(config) => Some(config), Err(err) => { - return GuardianReviewOutcome::Failed(GuardianReviewFailure::PromptBuild(err)); + return GuardianReviewOutcome::failed( + GuardianReviewFailure::PromptBuild(err), + None, + ); } }, None => None, @@ -741,45 +782,59 @@ pub(super) async fn run_guardian_review_session( ); let guardian_config = match guardian_config { Ok(config) => config, - Err(err) => return GuardianReviewOutcome::Failed(GuardianReviewFailure::PromptBuild(err)), + Err(err) => { + return GuardianReviewOutcome::failed(GuardianReviewFailure::PromptBuild(err), None); + } }; - match Box::pin( - session - .guardian_review_session - .run_review(GuardianReviewSessionParams { - parent_session: Arc::clone(&session), - parent_turn: turn.clone(), - spawn_config: guardian_config, - request, - retry_reason, - schema, - model: guardian_model, - reasoning_effort: guardian_reasoning_effort, - reasoning_summary: turn.reasoning_summary, - personality: turn.personality, - external_cancel, - }), - ) - .await - { + let (session_outcome, session_metadata) = Box::pin(session.guardian_review_session.run_review( + GuardianReviewSessionParams { + parent_session: Arc::clone(&session), + parent_turn: turn.clone(), + spawn_config: guardian_config, + request, + retry_reason, + schema, + model: guardian_model, + reasoning_effort: guardian_reasoning_effort, + reasoning_summary: turn.reasoning_summary, + personality: turn.personality, + external_cancel, + }, + )) + .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(Ok(assessment)), - Err(err) => GuardianReviewOutcome::Failed(GuardianReviewFailure::Parse(err)), + Ok(assessment) => { + GuardianReviewOutcome::completed(Ok(assessment), session_metadata) + } + Err(err) => GuardianReviewOutcome::failed( + GuardianReviewFailure::Parse(err), + session_metadata, + ), } } - None => GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(anyhow::anyhow!( - "guardian review completed without an assessment payload" - ))), + None => GuardianReviewOutcome::failed( + GuardianReviewFailure::Session(anyhow::anyhow!( + "guardian review completed without an assessment payload" + )), + session_metadata, + ), }, GuardianReviewSessionOutcome::Completed(Err(err)) => { - GuardianReviewOutcome::Failed(GuardianReviewFailure::Session(err)) + GuardianReviewOutcome::failed(GuardianReviewFailure::Session(err), session_metadata) } - GuardianReviewSessionOutcome::TimedOut => GuardianReviewOutcome::TimedOut, - GuardianReviewSessionOutcome::Aborted => GuardianReviewOutcome::Aborted, + GuardianReviewSessionOutcome::PromptBuildFailed(err) => { + GuardianReviewOutcome::failed(GuardianReviewFailure::PromptBuild(err), session_metadata) + } + GuardianReviewSessionOutcome::TimedOut => { + GuardianReviewOutcome::timed_out(session_metadata) + } + GuardianReviewSessionOutcome::Aborted => GuardianReviewOutcome::aborted(session_metadata), } } diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index 239754c508..7aaae56e46 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -5,6 +5,7 @@ use std::sync::Arc; use std::time::Duration; use anyhow::anyhow; +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 +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_util::sync::CancellationToken; @@ -58,10 +60,21 @@ const GUARDIAN_FOLLOWUP_REVIEW_REMINDER: &str = concat!( #[derive(Debug)] pub(crate) enum GuardianReviewSessionOutcome { Completed(anyhow::Result>), + PromptBuildFailed(anyhow::Error), TimedOut, Aborted, } +#[derive(Debug, Clone)] +pub(crate) struct GuardianReviewSessionMetadata { + pub(crate) guardian_thread_id: String, + pub(crate) guardian_session_kind: GuardianReviewSessionKind, + pub(crate) guardian_model: String, + pub(crate) guardian_reasoning_effort: Option, + pub(crate) had_prior_review_context: bool, + pub(crate) token_usage: Option, +} + pub(crate) struct GuardianReviewSessionParams { pub(crate) parent_session: Arc, pub(crate) parent_turn: Arc, @@ -101,6 +114,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>, @@ -267,10 +295,14 @@ impl GuardianReviewSessionManager { pub(crate) async fn run_review( &self, params: GuardianReviewSessionParams, - ) -> GuardianReviewSessionOutcome { + ) -> ( + GuardianReviewSessionOutcome, + Option, + ) { 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(), @@ -304,16 +336,17 @@ impl GuardianReviewSessionManager { { Ok(Ok(review_session)) => Arc::new(review_session), Ok(Err(err)) => { - return GuardianReviewSessionOutcome::Completed(Err(err)); + return (GuardianReviewSessionOutcome::PromptBuildFailed(err), None); } - Err(outcome) => return outcome, + Err(outcome) => return (outcome, None), }; state.trunk = Some(Arc::clone(&review_session)); + spawned_trunk = true; } state.trunk.as_ref().cloned() } - Err(outcome) => return outcome, + Err(outcome) => return (outcome, None), }; if let Some(review_session) = stale_trunk_to_shutdown { @@ -321,9 +354,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" + ))), + None, + ); }; if trunk.reuse_key != next_reuse_key { @@ -349,20 +385,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, metadata) = 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, Some(metadata)) } else { if let Some(review_session) = self.remove_trunk_if_current(&trunk).await { review_session.shutdown_in_background(); } - outcome + (outcome, Some(metadata)) } } @@ -459,7 +505,10 @@ impl GuardianReviewSessionManager { reuse_key: GuardianReviewSessionReuseKey, deadline: tokio::time::Instant, fork_snapshot: Option, - ) -> GuardianReviewSessionOutcome { + ) -> ( + GuardianReviewSessionOutcome, + Option, + ) { let spawn_cancel_token = CancellationToken::new(); let mut fork_config = params.spawn_config.clone(); fork_config.ephemeral = true; @@ -478,17 +527,18 @@ 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), None), + Err(outcome) => return (outcome, None), }; 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, _, metadata) = Box::pin(run_review_on_session( review_session.as_ref(), ¶ms, + GuardianReviewSessionKind::EphemeralForked, deadline, )) .await; @@ -496,7 +546,7 @@ impl GuardianReviewSessionManager { cleanup.disarm(); review_session.shutdown_in_background(); } - outcome + (outcome, Some(metadata)) } } @@ -543,8 +593,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, + GuardianReviewSessionMetadata, +) { let (send_followup_reminder, prompt_mode) = { let state = review_session.state.lock().await; @@ -559,6 +614,14 @@ async fn run_review_on_session( (send_followup_reminder, prompt_mode) }; + let mut guardian_metadata = GuardianReviewSessionMetadata { + guardian_thread_id: review_session.codex.session.conversation_id.to_string(), + guardian_session_kind, + guardian_model: params.model.clone(), + guardian_reasoning_effort: params.reasoning_effort.map(|effort| effort.to_string()), + had_prior_review_context: had_prior_review_context(&prompt_mode), + token_usage: None, + }; if send_followup_reminder { append_guardian_followup_reminder(review_session).await; } @@ -583,6 +646,8 @@ async fn run_review_on_session( prompt_mode, ) .await?; + let token_usage_at_review_start = + review_session.codex.session.total_token_usage().await; review_session .codex @@ -602,29 +667,45 @@ async fn run_review_on_session( }) .await?; - Ok::(prompt_items.transcript_cursor) + Ok::<(GuardianTranscriptCursor, Option), 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, guardian_metadata), }; - 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, + guardian_metadata, + ); } }; 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(token_usage_at_review_start) = token_usage_at_review_start + && let Some(total_token_usage) = review_session.codex.session.total_token_usage().await + { + guardian_metadata.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, guardian_metadata) } async fn append_guardian_followup_reminder(review_session: &GuardianReviewSession) { @@ -653,7 +734,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; @@ -662,7 +743,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 { @@ -672,7 +753,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 { @@ -684,18 +765,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); } _ => {} }, @@ -703,6 +786,7 @@ async fn wait_for_guardian_review( return ( GuardianReviewSessionOutcome::Completed(Err(err.into())), false, + false, ); } } @@ -954,4 +1038,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 82d4127a2b..a531fd6be9 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::NetworkDomainPermissionsToml; use crate::config_loader::RequirementSource; use crate::config_loader::Sourced; use crate::test_support; +use codex_analytics::GuardianApprovalRequestSource; use codex_config::config_toml::ConfigToml; use codex_network_proxy::NetworkProxyConfig; use codex_protocol::ThreadId; @@ -917,7 +918,11 @@ async fn guardian_review_request_layout_matches_model_visible_request_snapshot() /*external_cancel*/ None, ) .await; - let GuardianReviewOutcome::Completed(Ok(assessment)) = outcome else { + let GuardianReviewOutcome { + kind: GuardianReviewOutcomeKind::Completed(Ok(assessment)), + .. + } = outcome + else { panic!("expected guardian assessment"); }; assert_eq!(assessment.outcome, GuardianAssessmentOutcome::Allow); @@ -1126,13 +1131,25 @@ async fn guardian_reuses_prompt_cache_key_and_appends_prior_reviews() -> anyhow: ) .await; - let GuardianReviewOutcome::Completed(Ok(first_assessment)) = first_outcome else { + let GuardianReviewOutcome { + kind: GuardianReviewOutcomeKind::Completed(Ok(first_assessment)), + .. + } = first_outcome + else { panic!("expected first guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(second_assessment)) = second_outcome else { + let GuardianReviewOutcome { + kind: GuardianReviewOutcomeKind::Completed(Ok(second_assessment)), + .. + } = second_outcome + else { panic!("expected second guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(third_assessment)) = third_outcome else { + let GuardianReviewOutcome { + kind: GuardianReviewOutcomeKind::Completed(Ok(third_assessment)), + .. + } = third_outcome + else { panic!("expected third guardian assessment"); }; assert_eq!(first_assessment.outcome, GuardianAssessmentOutcome::Allow); From f8d87a8cd8bcf8dcc8728b048cdbf835d0454840 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Wed, 15 Apr 2026 16:48:30 -0700 Subject: [PATCH 3/4] [codex-analytics] guardian review truncation --- .../core/src/guardian/approval_request.rs | 66 +++++++++------ codex-rs/core/src/guardian/mod.rs | 2 +- codex-rs/core/src/guardian/prompt.rs | 25 ++++-- codex-rs/core/src/guardian/review.rs | 1 + codex-rs/core/src/guardian/review_session.rs | 28 ++++--- codex-rs/core/src/guardian/tests.rs | 80 +++++++++++++++++-- 6 files changed, 151 insertions(+), 51 deletions(-) diff --git a/codex-rs/core/src/guardian/approval_request.rs b/codex-rs/core/src/guardian/approval_request.rs index 6d1d3f76af..b3fb7f71dc 100644 --- a/codex-rs/core/src/guardian/approval_request.rs +++ b/codex-rs/core/src/guardian/approval_request.rs @@ -9,7 +9,7 @@ use serde::Serialize; use serde_json::Value; use super::GUARDIAN_MAX_ACTION_STRING_TOKENS; -use super::prompt::guardian_truncate_text; +use super::prompt::guardian_truncate_text_with_metadata; #[derive(Debug, Clone, PartialEq)] pub(crate) enum GuardianApprovalRequest { @@ -167,32 +167,49 @@ fn guardian_command_source_tool_name(source: GuardianCommandSource) -> &'static } } -fn truncate_guardian_action_value(value: Value) -> Value { +fn truncate_guardian_action_value(value: Value) -> (Value, bool) { match value { - Value::String(text) => Value::String(guardian_truncate_text( - &text, - GUARDIAN_MAX_ACTION_STRING_TOKENS, - )), - Value::Array(values) => Value::Array( - values + Value::String(text) => { + let (text, truncated) = + guardian_truncate_text_with_metadata(&text, GUARDIAN_MAX_ACTION_STRING_TOKENS); + (Value::String(text), truncated) + } + Value::Array(values) => { + let mut truncated = false; + let values = values .into_iter() - .map(truncate_guardian_action_value) - .collect::>(), - ), + .map(|value| { + let (value, value_truncated) = truncate_guardian_action_value(value); + truncated |= value_truncated; + value + }) + .collect::>(); + (Value::Array(values), truncated) + } Value::Object(values) => { let mut entries = values.into_iter().collect::>(); entries.sort_by(|(left, _), (right, _)| left.cmp(right)); - Value::Object( - entries - .into_iter() - .map(|(key, value)| (key, truncate_guardian_action_value(value))) - .collect(), - ) + let mut truncated = false; + let values = entries + .into_iter() + .map(|(key, value)| { + let (value, value_truncated) = truncate_guardian_action_value(value); + truncated |= value_truncated; + (key, value) + }) + .collect(); + (Value::Object(values), truncated) } - other => other, + other => (other, false), } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct FormattedGuardianAction { + pub(crate) text: String, + pub(crate) truncated: bool, +} + pub(crate) fn guardian_approval_request_to_json( action: &GuardianApprovalRequest, ) -> serde_json::Result { @@ -382,10 +399,13 @@ pub(crate) fn guardian_request_turn_id<'a>( } } -pub(crate) fn format_guardian_action_pretty( +pub(crate) fn format_guardian_action_pretty_with_truncation( action: &GuardianApprovalRequest, -) -> serde_json::Result { - let mut value = guardian_approval_request_to_json(action)?; - value = truncate_guardian_action_value(value); - serde_json::to_string_pretty(&value) +) -> serde_json::Result { + let value = guardian_approval_request_to_json(action)?; + let (value, truncated) = truncate_guardian_action_value(value); + Ok(FormattedGuardianAction { + text: serde_json::to_string_pretty(&value)?, + truncated, + }) } diff --git a/codex-rs/core/src/guardian/mod.rs b/codex-rs/core/src/guardian/mod.rs index f98bc750e8..98c3bde56d 100644 --- a/codex-rs/core/src/guardian/mod.rs +++ b/codex-rs/core/src/guardian/mod.rs @@ -62,7 +62,7 @@ pub(crate) struct GuardianRejection { } #[cfg(test)] -use approval_request::format_guardian_action_pretty; +use approval_request::format_guardian_action_pretty_with_truncation; #[cfg(test)] use approval_request::guardian_assessment_action; #[cfg(test)] diff --git a/codex-rs/core/src/guardian/prompt.rs b/codex-rs/core/src/guardian/prompt.rs index 761a5ec2fc..6308d2a711 100644 --- a/codex-rs/core/src/guardian/prompt.rs +++ b/codex-rs/core/src/guardian/prompt.rs @@ -19,7 +19,7 @@ use super::GUARDIAN_RECENT_ENTRY_LIMIT; use super::GuardianApprovalRequest; use super::GuardianAssessment; use super::TRUNCATION_TAG; -use super::approval_request::format_guardian_action_pretty; +use super::approval_request::format_guardian_action_pretty_with_truncation; /// Transcript entry retained for guardian review after filtering. #[derive(Debug, PartialEq, Eq)] @@ -56,6 +56,7 @@ impl GuardianTranscriptEntryKind { pub(crate) struct GuardianPromptItems { pub(crate) items: Vec, pub(crate) transcript_cursor: GuardianTranscriptCursor, + pub(crate) reviewed_action_truncated: bool, } /// Points to the end of the transcript that the guardian has already reviewed. @@ -91,7 +92,7 @@ pub(crate) async fn build_guardian_prompt_items( parent_history_version: history.history_version(), transcript_entry_count: transcript_entries.len(), }; - let planned_action_json = format_guardian_action_pretty(&request)?; + let planned_action = format_guardian_action_pretty_with_truncation(&request)?; let prompt_shape = match mode { GuardianPromptMode::Full => GuardianPromptShape::Full, @@ -176,11 +177,12 @@ pub(crate) async fn build_guardian_prompt_items( .to_string(), ); push_text("Planned action JSON:\n".to_string()); - push_text(format!("{planned_action_json}\n")); + push_text(format!("{}\n", planned_action.text)); push_text(">>> APPROVAL REQUEST END\n".to_string()); Ok(GuardianPromptItems { items, transcript_cursor, + reviewed_action_truncated: planned_action.truncated, }) } @@ -420,20 +422,23 @@ pub(crate) fn collect_guardian_transcript_entries( entries } -pub(crate) fn guardian_truncate_text(content: &str, token_cap: usize) -> String { +pub(crate) fn guardian_truncate_text_with_metadata( + content: &str, + token_cap: usize, +) -> (String, bool) { if content.is_empty() { - return String::new(); + return (String::new(), false); } let max_bytes = approx_bytes_for_tokens(token_cap); if content.len() <= max_bytes { - return content.to_string(); + return (content.to_string(), false); } let omitted_tokens = approx_tokens_from_byte_count(content.len().saturating_sub(max_bytes)); let marker = format!("<{TRUNCATION_TAG} omitted_approx_tokens=\"{omitted_tokens}\" />"); if max_bytes <= marker.len() { - return marker; + return (marker, true); } let available_bytes = max_bytes.saturating_sub(marker.len()); @@ -441,7 +446,11 @@ pub(crate) fn guardian_truncate_text(content: &str, token_cap: usize) -> String let suffix_budget = available_bytes.saturating_sub(prefix_budget); let (prefix, suffix) = split_guardian_truncation_bounds(content, prefix_budget, suffix_budget); - format!("{prefix}{marker}{suffix}") + (format!("{prefix}{marker}{suffix}"), true) +} + +pub(crate) fn guardian_truncate_text(content: &str, token_cap: usize) -> String { + guardian_truncate_text_with_metadata(content, token_cap).0 } fn split_guardian_truncation_bounds( diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index acab331b41..4643acbb38 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -280,6 +280,7 @@ impl GuardianReviewAnalyticsResult { terminal.guardian_model = Some(metadata.guardian_model); terminal.guardian_reasoning_effort = metadata.guardian_reasoning_effort; terminal.had_prior_review_context = Some(metadata.had_prior_review_context); + terminal.reviewed_action_truncated = metadata.reviewed_action_truncated; terminal.token_usage = metadata.token_usage; } diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index 7aaae56e46..80ccb75cb0 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -72,6 +72,7 @@ pub(crate) struct GuardianReviewSessionMetadata { pub(crate) guardian_model: String, pub(crate) guardian_reasoning_effort: Option, pub(crate) had_prior_review_context: bool, + pub(crate) reviewed_action_truncated: bool, pub(crate) token_usage: Option, } @@ -620,6 +621,7 @@ async fn run_review_on_session( guardian_model: params.model.clone(), guardian_reasoning_effort: params.reasoning_effort.map(|effort| effort.to_string()), had_prior_review_context: had_prior_review_context(&prompt_mode), + reviewed_action_truncated: false, token_usage: None, }; if send_followup_reminder { @@ -646,6 +648,7 @@ async fn run_review_on_session( prompt_mode, ) .await?; + let reviewed_action_truncated = prompt_items.reviewed_action_truncated; let token_usage_at_review_start = review_session.codex.session.total_token_usage().await; @@ -667,8 +670,9 @@ async fn run_review_on_session( }) .await?; - Ok::<(GuardianTranscriptCursor, Option), anyhow::Error>(( + Ok::<(GuardianTranscriptCursor, bool, Option), anyhow::Error>(( prompt_items.transcript_cursor, + reviewed_action_truncated, token_usage_at_review_start, )) }), @@ -678,16 +682,18 @@ async fn run_review_on_session( Ok(submit_result) => submit_result, Err(outcome) => return (outcome, false, guardian_metadata), }; - let (transcript_cursor, token_usage_at_review_start) = match submit_result { - Ok(submit_result) => submit_result, - Err(err) => { - return ( - GuardianReviewSessionOutcome::PromptBuildFailed(err), - false, - guardian_metadata, - ); - } - }; + let (transcript_cursor, reviewed_action_truncated, token_usage_at_review_start) = + match submit_result { + Ok(submit_result) => submit_result, + Err(err) => { + return ( + GuardianReviewSessionOutcome::PromptBuildFailed(err), + false, + guardian_metadata, + ); + } + }; + guardian_metadata.reviewed_action_truncated = reviewed_action_truncated; let outcome = wait_for_guardian_review(review_session, deadline, params.external_cancel.as_ref()).await; diff --git a/codex-rs/core/src/guardian/tests.rs b/codex-rs/core/src/guardian/tests.rs index a531fd6be9..2aae4a5513 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -585,11 +585,29 @@ fn format_guardian_action_pretty_truncates_large_string_fields() -> serde_json:: patch: patch.clone(), }; - let rendered = format_guardian_action_pretty(&action)?; + let rendered = format_guardian_action_pretty_with_truncation(&action)?; - assert!(rendered.contains("\"tool\": \"apply_patch\"")); - assert!(rendered.contains(" anyhow: let GuardianReviewOutcome { kind: GuardianReviewOutcomeKind::Completed(Ok(first_assessment)), - .. + metadata: first_metadata, } = first_outcome else { panic!("expected first guardian assessment"); }; + let first_metadata = first_metadata.expect("first guardian session metadata"); let GuardianReviewOutcome { kind: GuardianReviewOutcomeKind::Completed(Ok(second_assessment)), - .. + metadata: second_metadata, } = second_outcome else { panic!("expected second guardian assessment"); }; + let second_metadata = second_metadata.expect("second guardian session metadata"); let GuardianReviewOutcome { kind: GuardianReviewOutcomeKind::Completed(Ok(third_assessment)), - .. + metadata: third_metadata, } = third_outcome else { panic!("expected third guardian assessment"); }; + let third_metadata = third_metadata.expect("third guardian session metadata"); assert_eq!(first_assessment.outcome, GuardianAssessmentOutcome::Allow); assert_eq!(second_assessment.outcome, GuardianAssessmentOutcome::Allow); assert_eq!(third_assessment.outcome, GuardianAssessmentOutcome::Allow); + assert!(matches!( + first_metadata.guardian_session_kind, + codex_analytics::GuardianReviewSessionKind::TrunkNew + )); + assert!(matches!( + second_metadata.guardian_session_kind, + codex_analytics::GuardianReviewSessionKind::TrunkReused + )); + assert!(matches!( + third_metadata.guardian_session_kind, + codex_analytics::GuardianReviewSessionKind::TrunkReused + )); + ThreadId::from_string(&first_metadata.guardian_thread_id) + .expect("first guardian thread id should be a valid UUID"); + ThreadId::from_string(&second_metadata.guardian_thread_id) + .expect("second guardian thread id should be a valid UUID"); + ThreadId::from_string(&third_metadata.guardian_thread_id) + .expect("third guardian thread id should be a valid UUID"); + assert!(!first_metadata.had_prior_review_context); + assert!(second_metadata.had_prior_review_context); + assert!(third_metadata.had_prior_review_context); + assert_eq!( + first_metadata.guardian_thread_id, + second_metadata.guardian_thread_id + ); + assert_eq!( + second_metadata.guardian_thread_id, + third_metadata.guardian_thread_id + ); let requests = request_log.requests(); assert_eq!(requests.len(), 3); From bf12018ce290fde8f5a6ce99322a5270e0772d0a Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Wed, 15 Apr 2026 16:48:30 -0700 Subject: [PATCH 4/4] [codex-analytics] guardian review TTFT plumbing and emission --- .../src/protocol/thread_history.rs | 16 +++++++++++++++ .../app-server/src/bespoke_event_handling.rs | 1 + codex-rs/core/src/agent/control_tests.rs | 3 +++ .../src/codex/rollout_reconstruction_tests.rs | 20 +++++++++++++++++++ codex-rs/core/src/codex_tests.rs | 10 ++++++++++ codex-rs/core/src/guardian/review.rs | 1 + codex-rs/core/src/guardian/review_session.rs | 15 ++++++++++++-- codex-rs/core/src/guardian/tests.rs | 4 ++++ codex-rs/core/src/tasks/mod.rs | 5 +++++ .../src/tools/handlers/multi_agents_tests.rs | 3 +++ codex-rs/core/src/turn_timing.rs | 10 ++++++++++ codex-rs/core/tests/suite/resume_warning.rs | 1 + codex-rs/protocol/src/protocol.rs | 4 ++++ codex-rs/tui/src/app/app_server_adapter.rs | 2 ++ .../tui/src/chatwidget/tests/exec_flow.rs | 5 +++++ .../tui/src/chatwidget/tests/plan_mode.rs | 4 ++++ .../tui/src/chatwidget/tests/review_mode.rs | 2 ++ .../src/chatwidget/tests/slash_commands.rs | 5 +++++ .../src/chatwidget/tests/status_and_layout.rs | 3 +++ 19 files changed, 112 insertions(+), 2 deletions(-) diff --git a/codex-rs/app-server-protocol/src/protocol/thread_history.rs b/codex-rs/app-server-protocol/src/protocol/thread_history.rs index 9e94515dd9..a390d33392 100644 --- a/codex-rs/app-server-protocol/src/protocol/thread_history.rs +++ b/codex-rs/app-server-protocol/src/protocol/thread_history.rs @@ -1332,6 +1332,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -1406,6 +1407,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -1729,6 +1731,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2235,6 +2238,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::TurnStarted(TurnStartedEvent { turn_id: "turn-b".into(), @@ -2272,6 +2276,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2324,6 +2329,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::TurnStarted(TurnStartedEvent { turn_id: "turn-b".into(), @@ -2361,6 +2367,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2535,6 +2542,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::TurnStarted(TurnStartedEvent { turn_id: "turn-b".into(), @@ -2553,6 +2561,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::AgentMessage(AgentMessageEvent { message: "still in b".into(), @@ -2564,6 +2573,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2598,6 +2608,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::TurnStarted(TurnStartedEvent { turn_id: "turn-b".into(), @@ -2654,6 +2665,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -2899,6 +2911,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), EventMsg::Error(ErrorEvent { message: "request-level failure".into(), @@ -2958,6 +2971,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -3009,6 +3023,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -3057,6 +3072,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index 4fddee4a11..d5195f5960 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -3042,6 +3042,7 @@ mod tests { last_agent_message: None, completed_at: Some(TEST_TURN_COMPLETED_AT), duration_ms: Some(TEST_TURN_DURATION_MS), + time_to_first_token_ms: None, } } diff --git a/codex-rs/core/src/agent/control_tests.rs b/codex-rs/core/src/agent/control_tests.rs index 6fc74b30e8..2e800ea198 100644 --- a/codex-rs/core/src/agent/control_tests.rs +++ b/codex-rs/core/src/agent/control_tests.rs @@ -275,6 +275,7 @@ async fn on_event_updates_status_from_task_complete() { last_agent_message: Some("done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })); let expected = AgentStatus::Completed(Some("done".to_string())); assert_eq!(status, Some(expected)); @@ -1225,6 +1226,7 @@ async fn multi_agent_v2_completion_ignores_dead_direct_parent() { last_agent_message: Some("done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ) .await; @@ -1311,6 +1313,7 @@ async fn multi_agent_v2_completion_queues_message_for_direct_parent() { last_agent_message: Some("done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ) .await; diff --git a/codex-rs/core/src/codex/rollout_reconstruction_tests.rs b/codex-rs/core/src/codex/rollout_reconstruction_tests.rs index 753244ac2b..f3d62a090c 100644 --- a/codex-rs/core/src/codex/rollout_reconstruction_tests.rs +++ b/codex-rs/core/src/codex/rollout_reconstruction_tests.rs @@ -148,6 +148,7 @@ async fn record_initial_history_resumed_hydrates_previous_turn_settings_from_lif last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), ]; @@ -215,6 +216,7 @@ async fn reconstruct_history_rollback_keeps_history_and_metadata_in_sync_for_com last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -242,6 +244,7 @@ async fn reconstruct_history_rollback_keeps_history_and_metadata_in_sync_for_com last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::ThreadRolledBack( @@ -311,6 +314,7 @@ async fn reconstruct_history_rollback_keeps_history_and_metadata_in_sync_for_inc last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -400,6 +404,7 @@ async fn reconstruct_history_rollback_skips_non_user_turns_for_history_and_metad last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -426,6 +431,7 @@ async fn reconstruct_history_rollback_skips_non_user_turns_for_history_and_metad last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -443,6 +449,7 @@ async fn reconstruct_history_rollback_skips_non_user_turns_for_history_and_metad last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::ThreadRolledBack( @@ -515,6 +522,7 @@ async fn reconstruct_history_rollback_counts_inter_agent_assistant_turns() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -534,6 +542,7 @@ async fn reconstruct_history_rollback_counts_inter_agent_assistant_turns() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::ThreadRolledBack( @@ -601,6 +610,7 @@ async fn reconstruct_history_rollback_clears_history_and_metadata_when_exceeding last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::ThreadRolledBack( @@ -650,6 +660,7 @@ async fn record_initial_history_resumed_rollback_skips_only_user_turns() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), // Standalone task turn (no UserMessage) should not consume rollback skips. @@ -667,6 +678,7 @@ async fn record_initial_history_resumed_rollback_skips_only_user_turns() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::ThreadRolledBack( @@ -720,6 +732,7 @@ async fn record_initial_history_resumed_rollback_drops_incomplete_user_turn_comp last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -876,6 +889,7 @@ async fn reconstruct_history_legacy_compaction_without_replacement_history_clear last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), ]; @@ -945,6 +959,7 @@ async fn record_initial_history_resumed_turn_context_after_compaction_reestablis last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), ]; @@ -1046,6 +1061,7 @@ async fn record_initial_history_resumed_aborted_turn_without_id_clears_active_tu last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -1153,6 +1169,7 @@ async fn record_initial_history_resumed_unmatched_abort_preserves_active_turn_fo last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -1186,6 +1203,7 @@ async fn record_initial_history_resumed_unmatched_abort_preserves_active_turn_fo last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), ]; @@ -1268,6 +1286,7 @@ async fn record_initial_history_resumed_trailing_incomplete_turn_compaction_clea last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( @@ -1418,6 +1437,7 @@ async fn record_initial_history_resumed_replaced_incomplete_compacted_turn_clear last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), RolloutItem::EventMsg(EventMsg::TurnStarted( diff --git a/codex-rs/core/src/codex_tests.rs b/codex-rs/core/src/codex_tests.rs index 6de8909e6d..6e4e870e48 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -1381,6 +1381,7 @@ async fn record_initial_history_forked_hydrates_previous_turn_settings() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }, )), ]; @@ -1566,6 +1567,7 @@ async fn thread_rollback_recomputes_previous_turn_settings_and_reference_context last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), RolloutItem::EventMsg(EventMsg::TurnStarted( codex_protocol::protocol::TurnStartedEvent { @@ -1591,6 +1593,7 @@ async fn thread_rollback_recomputes_previous_turn_settings_and_reference_context last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]) .await; @@ -1668,6 +1671,7 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), RolloutItem::EventMsg(EventMsg::TurnStarted( codex_protocol::protocol::TurnStartedEvent { @@ -1686,6 +1690,7 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), RolloutItem::EventMsg(EventMsg::TurnStarted( codex_protocol::protocol::TurnStartedEvent { @@ -1713,6 +1718,7 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]) .await; @@ -1759,6 +1765,7 @@ async fn thread_rollback_persists_marker_and_replays_cumulatively() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), RolloutItem::EventMsg(EventMsg::TurnStarted( codex_protocol::protocol::TurnStartedEvent { @@ -1782,6 +1789,7 @@ async fn thread_rollback_persists_marker_and_replays_cumulatively() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), RolloutItem::EventMsg(EventMsg::TurnStarted( codex_protocol::protocol::TurnStartedEvent { @@ -1805,6 +1813,7 @@ async fn thread_rollback_persists_marker_and_replays_cumulatively() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]) .await; @@ -4785,6 +4794,7 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() EventMsg::TurnComplete(TurnCompleteEvent { turn_id, last_agent_message: None, + time_to_first_token_ms: None, .. }) if turn_id == tc.sub_id )); diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index 4643acbb38..d3a2631d21 100644 --- a/codex-rs/core/src/guardian/review.rs +++ b/codex-rs/core/src/guardian/review.rs @@ -282,6 +282,7 @@ impl GuardianReviewAnalyticsResult { terminal.had_prior_review_context = Some(metadata.had_prior_review_context); terminal.reviewed_action_truncated = metadata.reviewed_action_truncated; terminal.token_usage = metadata.token_usage; + terminal.time_to_first_token_ms = metadata.time_to_first_token_ms; } terminal diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index 80ccb75cb0..9dbcf18ef0 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -74,6 +74,7 @@ pub(crate) struct GuardianReviewSessionMetadata { pub(crate) had_prior_review_context: bool, pub(crate) reviewed_action_truncated: bool, pub(crate) token_usage: Option, + pub(crate) time_to_first_token_ms: Option, } pub(crate) struct GuardianReviewSessionParams { @@ -623,6 +624,7 @@ async fn run_review_on_session( had_prior_review_context: had_prior_review_context(&prompt_mode), reviewed_action_truncated: false, token_usage: None, + time_to_first_token_ms: None, }; if send_followup_reminder { append_guardian_followup_reminder(review_session).await; @@ -695,8 +697,13 @@ async fn run_review_on_session( }; guardian_metadata.reviewed_action_truncated = reviewed_action_truncated; - let outcome = - wait_for_guardian_review(review_session, deadline, params.external_cancel.as_ref()).await; + let outcome = wait_for_guardian_review( + review_session, + deadline, + params.external_cancel.as_ref(), + &mut guardian_metadata, + ) + .await; if matches!(outcome.0, GuardianReviewSessionOutcome::Completed(_)) { if outcome.2 && let Some(token_usage_at_review_start) = token_usage_at_review_start @@ -740,6 +747,7 @@ async fn wait_for_guardian_review( review_session: &GuardianReviewSession, deadline: tokio::time::Instant, external_cancel: Option<&CancellationToken>, + metadata: &mut GuardianReviewSessionMetadata, ) -> (GuardianReviewSessionOutcome, bool, bool) { let timeout = tokio::time::sleep_until(deadline); tokio::pin!(timeout); @@ -765,6 +773,9 @@ async fn wait_for_guardian_review( match event { Ok(event) => match event.msg { EventMsg::TurnComplete(turn_complete) => { + metadata.time_to_first_token_ms = turn_complete + .time_to_first_token_ms + .and_then(|ms| u64::try_from(ms).ok()); if turn_complete.last_agent_message.is_none() && let Some(error_message) = last_error_message { diff --git a/codex-rs/core/src/guardian/tests.rs b/codex-rs/core/src/guardian/tests.rs index 2aae4a5513..bf0b06b58a 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -958,6 +958,10 @@ async fn guardian_review_request_layout_matches_model_visible_request_snapshot() assert_eq!(metadata.guardian_model, "gpt-5.4"); assert_eq!(metadata.guardian_reasoning_effort.as_deref(), Some("low")); assert!(!metadata.had_prior_review_context); + assert!( + metadata.time_to_first_token_ms.is_some(), + "guardian review metadata should capture TTFT when the nested turn completes" + ); let request = request_log.single_request(); let mut settings = Settings::clone_current(); diff --git a/codex-rs/core/src/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index f17017316e..692a5919a3 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -535,11 +535,16 @@ impl Session { .turn_timing_state .completed_at_and_duration_ms() .await; + let time_to_first_token_ms = turn_context + .turn_timing_state + .time_to_first_token_ms() + .await; let event = EventMsg::TurnComplete(TurnCompleteEvent { turn_id: turn_context.sub_id.clone(), last_agent_message, completed_at, duration_ms, + time_to_first_token_ms, }); self.send_event(turn_context.as_ref(), event).await; diff --git a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs index e35d6a06fb..7c55be087d 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs @@ -1154,6 +1154,7 @@ async fn multi_agent_v2_list_agents_returns_completed_status_and_last_task_messa last_agent_message: Some("done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ) .await; @@ -1631,6 +1632,7 @@ async fn multi_agent_v2_followup_task_completion_notifies_parent_on_every_turn() last_agent_message: Some("first done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ) .await; @@ -1659,6 +1661,7 @@ async fn multi_agent_v2_followup_task_completion_notifies_parent_on_every_turn() last_agent_message: Some("second done".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ) .await; diff --git a/codex-rs/core/src/turn_timing.rs b/codex-rs/core/src/turn_timing.rs index 6c47d3b528..c8692f1976 100644 --- a/codex-rs/core/src/turn_timing.rs +++ b/codex-rs/core/src/turn_timing.rs @@ -74,6 +74,16 @@ impl TurnTimingState { (completed_at, duration_ms) } + pub(crate) async fn time_to_first_token_ms(&self) -> Option { + let state = self.state.lock().await; + let started_at = state.started_at?; + let first_token_at = state.first_token_at?; + Some( + i64::try_from(first_token_at.duration_since(started_at).as_millis()) + .unwrap_or(i64::MAX), + ) + } + pub(crate) async fn record_ttft_for_response_event( &self, event: &ResponseEvent, diff --git a/codex-rs/core/tests/suite/resume_warning.rs b/codex-rs/core/tests/suite/resume_warning.rs index f9ad64c771..70c9014747 100644 --- a/codex-rs/core/tests/suite/resume_warning.rs +++ b/codex-rs/core/tests/suite/resume_warning.rs @@ -69,6 +69,7 @@ fn resume_history( last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ], rollout_path: rollout_path.to_path_buf(), diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 7d20fadeb4..73045ba734 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -2057,6 +2057,10 @@ pub struct TurnCompleteEvent { #[serde(default, skip_serializing_if = "Option::is_none")] #[ts(type = "number | null", optional)] pub duration_ms: Option, + /// Duration between turn start and the first model token in milliseconds, if known. + #[serde(default, skip_serializing_if = "Option::is_none")] + #[ts(type = "number | null", optional)] + pub time_to_first_token_ms: Option, } #[derive(Debug, Clone, Deserialize, Serialize, JsonSchema, TS)] diff --git a/codex-rs/tui/src/app/app_server_adapter.rs b/codex-rs/tui/src/app/app_server_adapter.rs index 09b8ce06cd..2107aef4b7 100644 --- a/codex-rs/tui/src/app/app_server_adapter.rs +++ b/codex-rs/tui/src/app/app_server_adapter.rs @@ -762,6 +762,7 @@ fn append_terminal_turn_events(events: &mut Vec, turn: &Turn, include_fai last_agent_message: None, completed_at: turn.completed_at, duration_ms: turn.duration_ms, + time_to_first_token_ms: None, }), }), TurnStatus::Interrupted => events.push(Event { @@ -793,6 +794,7 @@ fn append_terminal_turn_events(events: &mut Vec, turn: &Turn, include_fai last_agent_message: None, completed_at: turn.completed_at, duration_ms: turn.duration_ms, + time_to_first_token_ms: None, }), }); } diff --git a/codex-rs/tui/src/chatwidget/tests/exec_flow.rs b/codex-rs/tui/src/chatwidget/tests/exec_flow.rs index 53f16a39d1..261fd47f3f 100644 --- a/codex-rs/tui/src/chatwidget/tests/exec_flow.rs +++ b/codex-rs/tui/src/chatwidget/tests/exec_flow.rs @@ -654,6 +654,7 @@ async fn unified_exec_wait_after_final_agent_message_snapshot() { last_agent_message: Some("Final response.".into()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -699,6 +700,7 @@ async fn unified_exec_wait_before_streamed_agent_message_snapshot() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -766,6 +768,7 @@ async fn unified_exec_waiting_multiple_empty_snapshots() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -846,6 +849,7 @@ async fn unified_exec_non_empty_then_empty_snapshots() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -1350,6 +1354,7 @@ async fn turn_complete_keeps_unified_exec_processes() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); diff --git a/codex-rs/tui/src/chatwidget/tests/plan_mode.rs b/codex-rs/tui/src/chatwidget/tests/plan_mode.rs index f3746fdbf0..2cd2790ee6 100644 --- a/codex-rs/tui/src/chatwidget/tests/plan_mode.rs +++ b/codex-rs/tui/src/chatwidget/tests/plan_mode.rs @@ -546,6 +546,7 @@ async fn plan_implementation_popup_skips_replayed_turn_complete() { last_agent_message: Some("Plan details".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })]); let popup = render_bottom_popup(&chat, /*width*/ 80); @@ -572,6 +573,7 @@ async fn plan_implementation_popup_shows_once_when_replay_precedes_live_turn_com last_agent_message: Some("Plan details".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })]); let replay_popup = render_bottom_popup(&chat, /*width*/ 80); assert!( @@ -586,6 +588,7 @@ async fn plan_implementation_popup_shows_once_when_replay_precedes_live_turn_com last_agent_message: Some("Plan details".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -609,6 +612,7 @@ async fn plan_implementation_popup_shows_once_when_replay_precedes_live_turn_com last_agent_message: Some("Plan details".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); let duplicate_popup = render_bottom_popup(&chat, /*width*/ 80); diff --git a/codex-rs/tui/src/chatwidget/tests/review_mode.rs b/codex-rs/tui/src/chatwidget/tests/review_mode.rs index c42ec4fbac..9a05320aa6 100644 --- a/codex-rs/tui/src/chatwidget/tests/review_mode.rs +++ b/codex-rs/tui/src/chatwidget/tests/review_mode.rs @@ -234,6 +234,7 @@ async fn steer_rejection_queues_review_follow_up_before_existing_queued_messages last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -255,6 +256,7 @@ async fn steer_rejection_queues_review_follow_up_before_existing_queued_messages last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); diff --git a/codex-rs/tui/src/chatwidget/tests/slash_commands.rs b/codex-rs/tui/src/chatwidget/tests/slash_commands.rs index 2afb959349..27cd96e388 100644 --- a/codex-rs/tui/src/chatwidget/tests/slash_commands.rs +++ b/codex-rs/tui/src/chatwidget/tests/slash_commands.rs @@ -225,6 +225,7 @@ async fn slash_copy_state_tracks_turn_complete_final_reply() { last_agent_message: Some("Final reply **markdown**".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -257,6 +258,7 @@ async fn slash_copy_state_tracks_plan_item_completion() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -343,6 +345,7 @@ async fn slash_copy_state_is_preserved_during_running_task() { last_agent_message: Some("Previous completed reply".to_string()), completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); chat.on_task_started(); @@ -373,6 +376,7 @@ async fn slash_copy_tracks_replayed_legacy_agent_message_when_turn_complete_omit last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); let _ = drain_insert_history(&mut rx); @@ -410,6 +414,7 @@ async fn slash_copy_uses_agent_message_item_when_turn_complete_omits_final_text( last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); let _ = drain_insert_history(&mut rx); diff --git a/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs b/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs index e59928ff1c..a294e3aadf 100644 --- a/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs +++ b/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs @@ -961,6 +961,7 @@ async fn status_line_branch_refreshes_after_turn_complete() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -1311,6 +1312,7 @@ async fn multiple_agent_messages_in_single_turn_emit_multiple_headers() { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); @@ -2251,6 +2253,7 @@ printf 'fenced within fenced\n' last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), }); for lines in drain_insert_history(&mut rx) {