From bb4e510fa7db5c8124deef750a0d75b97f687bd4 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Mon, 13 Apr 2026 14:25:22 -0700 Subject: [PATCH 1/4] [codex-analytics] guardian review analytics events emission --- codex-rs/core/src/codex_delegate.rs | 6 +- codex-rs/core/src/guardian/mod.rs | 11 +- codex-rs/core/src/guardian/review.rs | 500 ++++++++++++++++++++++++++- 3 files changed, 491 insertions(+), 26 deletions(-) diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index 55b3619e11..d7911ff349 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; @@ -40,7 +41,7 @@ use crate::codex::emit_subagent_session_started; use crate::config::Config; use crate::guardian::GuardianApprovalRequest; use crate::guardian::new_guardian_review_id; -use crate::guardian::review_approval_request_with_cancel; +use crate::guardian::review_approval_request_with_source_and_cancel; use crate::guardian::routes_approval_to_guardian; use crate::mcp_tool_call::MCP_TOOL_APPROVAL_ACCEPT; use crate::mcp_tool_call::MCP_TOOL_APPROVAL_ACCEPT_FOR_SESSION; @@ -747,12 +748,13 @@ fn spawn_guardian_review( let _ = tx.send(ReviewDecision::Denied); return; }; - let decision = runtime.block_on(review_approval_request_with_cancel( + let decision = runtime.block_on(review_approval_request_with_source_and_cancel( &session, &turn, review_id, request, retry_reason, + GuardianApprovalRequestSource::DelegatedSubagent, cancel_token, )); let _ = tx.send(decision); diff --git a/codex-rs/core/src/guardian/mod.rs b/codex-rs/core/src/guardian/mod.rs index 67e9a828ee..f5926a0db5 100644 --- a/codex-rs/core/src/guardian/mod.rs +++ b/codex-rs/core/src/guardian/mod.rs @@ -19,6 +19,7 @@ mod review_session; use std::time::Duration; use codex_protocol::protocol::GuardianAssessmentDecisionSource; +use codex_protocol::protocol::GuardianAssessmentOutcome; use serde::Deserialize; use serde::Serialize; @@ -30,7 +31,9 @@ pub(crate) use review::guardian_timeout_message; pub(crate) use review::is_guardian_reviewer_source; pub(crate) use review::new_guardian_review_id; pub(crate) use review::review_approval_request; +#[cfg(test)] pub(crate) use review::review_approval_request_with_cancel; +pub(crate) use review::review_approval_request_with_source_and_cancel; pub(crate) use review::routes_approval_to_guardian; pub(crate) use review_session::GuardianReviewSessionManager; @@ -45,14 +48,6 @@ const GUARDIAN_MAX_ACTION_STRING_TOKENS: usize = 16_000; const GUARDIAN_RECENT_ENTRY_LIMIT: usize = 40; const TRUNCATION_TAG: &str = "truncated"; -/// Final allow/deny outcome returned by the guardian reviewer. -#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq)] -#[serde(rename_all = "lowercase")] -pub(crate) enum GuardianAssessmentOutcome { - Allow, - Deny, -} - /// Structured output contract that the guardian reviewer must satisfy. #[derive(Debug, Clone, Deserialize, Serialize)] pub(crate) struct GuardianAssessment { diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index 861d3576cd..495fba06a1 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,202 @@ fn guardian_risk_level_str(level: GuardianRiskLevel) -> &'static str { } } +fn guardian_reviewed_action(request: &GuardianApprovalRequest) -> GuardianReviewedAction { + match request { + GuardianApprovalRequest::Shell { + command, + cwd, + sandbox_permissions, + additional_permissions, + justification, + .. + } => GuardianReviewedAction::Shell { + command: command.clone(), + command_display: codex_shell_command::parse_command::shlex_join(command), + cwd: cwd.to_string_lossy().into_owned(), + sandbox_permissions: *sandbox_permissions, + additional_permissions: additional_permissions.clone(), + justification: justification.clone(), + }, + GuardianApprovalRequest::ExecCommand { + command, + cwd, + sandbox_permissions, + additional_permissions, + justification, + tty, + .. + } => GuardianReviewedAction::UnifiedExec { + command: command.clone(), + command_display: codex_shell_command::parse_command::shlex_join(command), + cwd: cwd.to_string_lossy().into_owned(), + sandbox_permissions: *sandbox_permissions, + additional_permissions: additional_permissions.clone(), + justification: justification.clone(), + tty: *tty, + }, + #[cfg(unix)] + GuardianApprovalRequest::Execve { + source, + program, + argv, + cwd, + additional_permissions, + .. + } => GuardianReviewedAction::Execve { + source: *source, + program: program.clone(), + argv: argv.clone(), + cwd: cwd.to_string_lossy().into_owned(), + additional_permissions: additional_permissions.clone(), + }, + GuardianApprovalRequest::ApplyPatch { cwd, files, .. } => { + GuardianReviewedAction::ApplyPatch { + cwd: cwd.to_string_lossy().into_owned(), + files: files + .iter() + .map(|file| file.to_string_lossy().into_owned()) + .collect(), + } + } + GuardianApprovalRequest::NetworkAccess { + target, + host, + protocol, + port, + .. + } => GuardianReviewedAction::NetworkAccess { + target: target.clone(), + host: host.clone(), + 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, + retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, + reviewed_action: GuardianReviewedAction, + started_at: u64, + started_instant: Instant, +} + +struct GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision, + terminal_status: GuardianReviewTerminalStatus, + failure_reason: Option, + risk_level: Option, + user_authorization: Option, + outcome: Option, + rationale: 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 GuardianReviewAnalyticsContext { + fn track( + &self, + session: &Session, + turn: &TurnContext, + terminal: GuardianReviewAnalyticsTerminal, + ) { + 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(), + retry_reason: self.retry_reason.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, + rationale: terminal.rationale, + 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 +347,25 @@ 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(), + retry_reason: retry_reason.clone(), + approval_request_source, + reviewed_action: guardian_reviewed_action(&request), + started_at, + started_instant, + }; session .send_event( turn.as_ref(), @@ -142,6 +387,28 @@ async fn run_guardian_review( .as_ref() .is_some_and(CancellationToken::is_cancelled) { + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + risk_level: None, + user_authorization: None, + outcome: None, + rationale: None, + guardian_thread_id: None, + guardian_session_kind: None, + guardian_model: None, + guardian_reasoning_effort: None, + had_prior_review_context: None, + reviewed_action_truncated: false, + token_usage: None, + time_to_first_token_ms: None, + completed_at: now_unix_seconds(), + }, + ); session .send_event( turn.as_ref(), @@ -167,24 +434,142 @@ 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 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 metadata = GuardianReviewMetadataFields::default(); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: if matches!(assessment.outcome, GuardianAssessmentOutcome::Allow) { + GuardianReviewDecision::Approved + } else { + GuardianReviewDecision::Denied + }, + terminal_status: if matches!( + assessment.outcome, + GuardianAssessmentOutcome::Allow + ) { + GuardianReviewTerminalStatus::Approved + } else { + GuardianReviewTerminalStatus::Denied + }, + failure_reason: None, + risk_level: Some(assessment.risk_level), + user_authorization: Some(assessment.user_authorization), + outcome: Some(assessment.outcome), + rationale: Some(assessment.rationale.clone()), + 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, + }, + ); + assessment + } + GuardianReviewOutcome::Completed(Err(err)) => { + let metadata = GuardianReviewMetadataFields::default(); + let rationale = format!("Automatic approval review failed: {err}"); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(GuardianReviewFailureReason::SessionError), + risk_level: None, + user_authorization: None, + outcome: None, + rationale: 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, + }, + ); + GuardianAssessment { + risk_level: GuardianRiskLevel::High, + user_authorization: GuardianUserAuthorization::Unknown, + outcome: GuardianAssessmentOutcome::Deny, + rationale, + } + } + GuardianReviewOutcome::Failed(failure) => { + let metadata = GuardianReviewMetadataFields::default(); + let rationale = format!("Automatic approval review failed: {}", failure.error()); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::FailedClosed, + failure_reason: Some(failure.reason()), + risk_level: None, + user_authorization: None, + outcome: None, + rationale: 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, + }, + ); + GuardianAssessment { + risk_level: GuardianRiskLevel::High, + user_authorization: GuardianUserAuthorization::Unknown, + outcome: GuardianAssessmentOutcome::Deny, + rationale, + } + } GuardianReviewOutcome::TimedOut => { + let metadata = GuardianReviewMetadataFields::default(); let rationale = "Automatic approval review timed out while evaluating the requested approval." .to_string(); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision::Denied, + terminal_status: GuardianReviewTerminalStatus::TimedOut, + failure_reason: Some(GuardianReviewFailureReason::Timeout), + risk_level: None, + user_authorization: None, + outcome: None, + rationale: 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, + }, + ); session .send_event( turn.as_ref(), @@ -212,6 +597,29 @@ async fn run_guardian_review( return ReviewDecision::TimedOut; } GuardianReviewOutcome::Aborted => { + let metadata = GuardianReviewMetadataFields::default(); + analytics_context.track( + session.as_ref(), + turn.as_ref(), + GuardianReviewAnalyticsTerminal { + decision: GuardianReviewDecision::Aborted, + terminal_status: GuardianReviewTerminalStatus::Aborted, + failure_reason: Some(GuardianReviewFailureReason::Cancelled), + risk_level: None, + user_authorization: None, + outcome: None, + rationale: 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, + }, + ); session .send_event( turn.as_ref(), @@ -309,11 +717,34 @@ pub(crate) async fn review_approval_request( review_id, request, retry_reason, + GuardianApprovalRequestSource::MainTurn, /*external_cancel*/ None, ) .await } +pub(crate) async fn review_approval_request_with_source_and_cancel( + session: &Arc, + turn: &Arc, + review_id: String, + request: GuardianApprovalRequest, + retry_reason: Option, + approval_request_source: GuardianApprovalRequestSource, + cancel_token: CancellationToken, +) -> ReviewDecision { + run_guardian_review( + Arc::clone(session), + Arc::clone(turn), + review_id, + request, + retry_reason, + approval_request_source, + Some(cancel_token), + ) + .await +} + +#[cfg(test)] pub(crate) async fn review_approval_request_with_cancel( session: &Arc, turn: &Arc, @@ -328,6 +759,7 @@ pub(crate) async fn review_approval_request_with_cancel( review_id, request, retry_reason, + GuardianApprovalRequestSource::MainTurn, Some(cancel_token), ) .await @@ -358,7 +790,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, }; @@ -408,7 +842,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 session @@ -428,15 +862,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 + )); + } +} From 1caeafafe32ba707ed259841fe6d9cbd9da84216 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Mon, 13 Apr 2026 14:25:26 -0700 Subject: [PATCH 2/4] [codex-analytics] guardian review thread and token metadata --- codex-rs/core/src/guardian/review.rs | 93 ++++++--- codex-rs/core/src/guardian/review_session.rs | 205 ++++++++++++++++--- 2 files changed, 239 insertions(+), 59 deletions(-) diff --git a/codex-rs/core/src/guardian/review.rs b/codex-rs/core/src/guardian/review.rs index 495fba06a1..06e042b815 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; @@ -86,10 +87,13 @@ pub(crate) fn guardian_timeout_message() -> String { #[derive(Debug)] pub(super) enum GuardianReviewOutcome { - Completed(anyhow::Result), - Failed(GuardianReviewFailure), - TimedOut, - Aborted, + Completed( + anyhow::Result, + Option, + ), + Failed(GuardianReviewFailure, Option), + TimedOut(Option), + Aborted(Option), } #[derive(Debug)] @@ -254,6 +258,24 @@ struct GuardianReviewMetadataFields { time_to_first_token_ms: Option, } +fn guardian_review_metadata_fields( + metadata: Option, +) -> GuardianReviewMetadataFields { + match metadata { + Some(metadata) => GuardianReviewMetadataFields { + guardian_thread_id: Some(metadata.guardian_thread_id), + guardian_session_kind: Some(metadata.guardian_session_kind), + guardian_model: Some(metadata.guardian_model), + guardian_reasoning_effort: metadata.guardian_reasoning_effort, + had_prior_review_context: Some(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, + }, + None => GuardianReviewMetadataFields::default(), + } +} + impl GuardianReviewAnalyticsContext { fn track( &self, @@ -442,8 +464,8 @@ async fn run_guardian_review( let completed_at = now_unix_seconds(); let assessment = match outcome { - GuardianReviewOutcome::Completed(Ok(assessment)) => { - let metadata = GuardianReviewMetadataFields::default(); + GuardianReviewOutcome::Completed(Ok(assessment), metadata) => { + let metadata = guardian_review_metadata_fields(metadata); analytics_context.track( session.as_ref(), turn.as_ref(), @@ -479,8 +501,8 @@ async fn run_guardian_review( ); assessment } - GuardianReviewOutcome::Completed(Err(err)) => { - let metadata = GuardianReviewMetadataFields::default(); + GuardianReviewOutcome::Completed(Err(err), metadata) => { + let metadata = guardian_review_metadata_fields(metadata); let rationale = format!("Automatic approval review failed: {err}"); analytics_context.track( session.as_ref(), @@ -511,8 +533,8 @@ async fn run_guardian_review( rationale, } } - GuardianReviewOutcome::Failed(failure) => { - let metadata = GuardianReviewMetadataFields::default(); + GuardianReviewOutcome::Failed(failure, metadata) => { + let metadata = guardian_review_metadata_fields(metadata); let rationale = format!("Automatic approval review failed: {}", failure.error()); analytics_context.track( session.as_ref(), @@ -543,8 +565,8 @@ async fn run_guardian_review( rationale, } } - GuardianReviewOutcome::TimedOut => { - let metadata = GuardianReviewMetadataFields::default(); + GuardianReviewOutcome::TimedOut(metadata) => { + let metadata = guardian_review_metadata_fields(metadata); let rationale = "Automatic approval review timed out while evaluating the requested approval." .to_string(); @@ -596,8 +618,8 @@ async fn run_guardian_review( .await; return ReviewDecision::TimedOut; } - GuardianReviewOutcome::Aborted => { - let metadata = GuardianReviewMetadataFields::default(); + GuardianReviewOutcome::Aborted(metadata) => { + let metadata = guardian_review_metadata_fields(metadata); analytics_context.track( session.as_ref(), turn.as_ref(), @@ -791,7 +813,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, @@ -842,10 +867,12 @@ 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 session + let (session_outcome, session_metadata) = session .guardian_review_session .run_review(GuardianReviewSessionParams { parent_session: Arc::clone(&session), @@ -860,25 +887,37 @@ pub(super) async fn run_guardian_review_session( personality: turn.personality, external_cancel, }) - .await - { + .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::TimedOut(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 f6cd4d02b2..984181ed1f 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; @@ -57,10 +59,23 @@ 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) reviewed_action_truncated: bool, + pub(crate) token_usage: Option, + pub(crate) time_to_first_token_ms: Option, +} + pub(crate) struct GuardianReviewSessionParams { pub(crate) parent_session: Arc, pub(crate) parent_turn: Arc, @@ -100,6 +115,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>, @@ -266,10 +296,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(), @@ -303,16 +337,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 { @@ -320,9 +355,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 { @@ -350,20 +388,25 @@ impl GuardianReviewSessionManager { } }; - let (outcome, keep_review_session) = - 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) = + 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)) } } @@ -460,7 +503,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; @@ -479,20 +525,26 @@ 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, _) = run_review_on_session(review_session.as_ref(), ¶ms, deadline).await; + let (outcome, _, metadata) = run_review_on_session( + review_session.as_ref(), + ¶ms, + GuardianReviewSessionKind::EphemeralForked, + deadline, + ) + .await; if let Some(review_session) = self.take_active_ephemeral(&review_session).await { cleanup.disarm(); review_session.shutdown_in_background(); } - outcome + (outcome, Some(metadata)) } } @@ -539,8 +591,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; @@ -555,6 +612,16 @@ 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), + reviewed_action_truncated: false, + token_usage: None, + time_to_first_token_ms: None, + }; if send_followup_reminder { append_guardian_followup_reminder(review_session).await; } @@ -579,6 +646,9 @@ 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; review_session .codex @@ -598,29 +668,53 @@ async fn run_review_on_session( }) .await?; - Ok::(prompt_items.transcript_cursor) + Ok::<(GuardianTranscriptCursor, bool, Option), anyhow::Error>(( + prompt_items.transcript_cursor, + reviewed_action_truncated, + token_usage_at_review_start, + )) }), ) .await; let submit_result = match submit_result { Ok(submit_result) => submit_result, - Err(outcome) => return (outcome, false), - }; - let transcript_cursor = match submit_result { - Ok(transcript_cursor) => transcript_cursor, - Err(err) => { - return (GuardianReviewSessionOutcome::Completed(Err(err)), false); - } + Err(outcome) => return (outcome, 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; + 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 + && 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) { @@ -649,7 +743,8 @@ async fn wait_for_guardian_review( review_session: &GuardianReviewSession, deadline: tokio::time::Instant, external_cancel: Option<&CancellationToken>, -) -> (GuardianReviewSessionOutcome, bool) { + metadata: &mut GuardianReviewSessionMetadata, +) -> (GuardianReviewSessionOutcome, bool, bool) { let timeout = tokio::time::sleep_until(deadline); tokio::pin!(timeout); let mut last_error_message: Option = None; @@ -658,7 +753,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 { @@ -668,30 +763,35 @@ 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 { 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 { 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); } _ => {} }, @@ -699,6 +799,7 @@ async fn wait_for_guardian_review( return ( GuardianReviewSessionOutcome::Completed(Err(err.into())), false, + false, ); } } @@ -950,4 +1051,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, + } + ); + } } From 449be750aff7aa05beb5a123daf50e32e1dfd48a Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Mon, 13 Apr 2026 14:25:26 -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/tests.rs | 87 +++++++++++++++++-- 4 files changed, 140 insertions(+), 40 deletions(-) diff --git a/codex-rs/core/src/guardian/approval_request.rs b/codex-rs/core/src/guardian/approval_request.rs index c9cc9d9fa4..bba3bc1c4e 100644 --- a/codex-rs/core/src/guardian/approval_request.rs +++ b/codex-rs/core/src/guardian/approval_request.rs @@ -10,7 +10,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 { @@ -168,32 +168,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 { @@ -386,10 +403,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 f5926a0db5..ddd17439ce 100644 --- a/codex-rs/core/src/guardian/mod.rs +++ b/codex-rs/core/src/guardian/mod.rs @@ -64,7 +64,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/tests.rs b/codex-rs/core/src/guardian/tests.rs index b6ba6349f4..26261beba2 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -571,11 +571,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: ) .await; - let GuardianReviewOutcome::Completed(Ok(first_assessment)) = first_outcome else { + let GuardianReviewOutcome::Completed(Ok(first_assessment), first_metadata) = first_outcome + else { panic!("expected first guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(second_assessment)) = second_outcome else { + let first_metadata = first_metadata.expect("first guardian session metadata"); + let GuardianReviewOutcome::Completed(Ok(second_assessment), second_metadata) = second_outcome + else { panic!("expected second guardian assessment"); }; - let GuardianReviewOutcome::Completed(Ok(third_assessment)) = third_outcome else { + let second_metadata = second_metadata.expect("second guardian session metadata"); + let GuardianReviewOutcome::Completed(Ok(third_assessment), 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 dd5d51d9e8943889a5e9121ee228cbf78f6c85cd Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Mon, 13 Apr 2026 14:25:26 -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/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 +++ 16 files changed, 94 insertions(+) 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 d2296c75c5..b93f8b662f 100644 --- a/codex-rs/app-server-protocol/src/protocol/thread_history.rs +++ b/codex-rs/app-server-protocol/src/protocol/thread_history.rs @@ -1330,6 +1330,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -1404,6 +1405,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -1727,6 +1729,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2233,6 +2236,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(), @@ -2270,6 +2274,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2322,6 +2327,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(), @@ -2359,6 +2365,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2533,6 +2540,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(), @@ -2551,6 +2559,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(), @@ -2562,6 +2571,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -2596,6 +2606,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(), @@ -2652,6 +2663,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -2897,6 +2909,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(), @@ -2956,6 +2969,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, }), ]; @@ -3007,6 +3021,7 @@ mod tests { last_agent_message: None, completed_at: None, duration_ms: None, + time_to_first_token_ms: None, })), ]; @@ -3055,6 +3070,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 8d0d40ff94..9313ccbed9 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -2996,6 +2996,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 5648fffd7c..d592bc5f4b 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -1446,6 +1446,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, }, )), ]; @@ -1631,6 +1632,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 { @@ -1656,6 +1658,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; @@ -1733,6 +1736,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 { @@ -1751,6 +1755,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 { @@ -1778,6 +1783,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; @@ -1824,6 +1830,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 { @@ -1847,6 +1854,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 { @@ -1870,6 +1878,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; @@ -4830,6 +4839,7 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() 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/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index b46db8e98f..72c008aa19 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -527,11 +527,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 cb5dea34f0..536b8b000e 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -2041,6 +2041,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 0bd78f353a..e3e32adf7b 100644 --- a/codex-rs/tui/src/app/app_server_adapter.rs +++ b/codex-rs/tui/src/app/app_server_adapter.rs @@ -759,6 +759,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 { @@ -790,6 +791,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 d271fc0076..b23074f791 100644 --- a/codex-rs/tui/src/chatwidget/tests/exec_flow.rs +++ b/codex-rs/tui/src/chatwidget/tests/exec_flow.rs @@ -652,6 +652,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, }), }); @@ -697,6 +698,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, }), }); @@ -764,6 +766,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, }), }); @@ -844,6 +847,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, }), }); @@ -1344,6 +1348,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 083cb8069a..47332bfcf2 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 1c43d2e91d..b4726b7cf9 100644 --- a/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs +++ b/codex-rs/tui/src/chatwidget/tests/status_and_layout.rs @@ -943,6 +943,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, }), }); @@ -1259,6 +1260,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, }), }); @@ -2199,6 +2201,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) {