[codex-analytics] guardian review analytics events emission

This commit is contained in:
rhan-oai
2026-04-20 23:07:31 -07:00
parent d62421d322
commit 4ded4eb692
7 changed files with 625 additions and 61 deletions

View File

@@ -1,5 +1,6 @@
use crate::events::AppServerRpcTransport;
use crate::events::GuardianReviewEventParams;
use crate::events::GuardianReviewAnalyticsResult;
use crate::events::GuardianReviewTrackContext;
use crate::events::TrackEventRequest;
use crate::events::TrackEventsRequest;
use crate::events::current_runtime_metadata;
@@ -161,9 +162,13 @@ impl AnalyticsEventsClient {
));
}
pub fn track_guardian_review(&self, input: GuardianReviewEventParams) {
pub fn track_guardian_review(
&self,
tracking: &GuardianReviewTrackContext,
result: GuardianReviewAnalyticsResult,
) {
self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::GuardianReview(
Box::new(input),
Box::new(tracking.event_params(result)),
)));
}

View File

@@ -1,3 +1,5 @@
use std::time::Instant;
use crate::facts::AppInvocation;
use crate::facts::CodexCompactionEvent;
use crate::facts::CompactionImplementation;
@@ -16,6 +18,7 @@ use crate::facts::TurnStatus;
use crate::facts::TurnSteerRejectionReason;
use crate::facts::TurnSteerResult;
use crate::facts::TurnSubmissionType;
use crate::now_unix_seconds;
use codex_app_server_protocol::CodexErrorInfo;
use codex_login::default_client::originator;
use codex_plugin::PluginTelemetryMetadata;
@@ -30,6 +33,7 @@ use codex_protocol::protocol::HookEventName;
use codex_protocol::protocol::HookRunStatus;
use codex_protocol::protocol::HookSource;
use codex_protocol::protocol::SubAgentSource;
use codex_protocol::protocol::TokenUsage;
use serde::Serialize;
#[derive(Clone, Copy, Debug, Serialize)]
@@ -200,6 +204,7 @@ pub enum GuardianReviewedAction {
connector_name: Option<String>,
tool_title: Option<String>,
},
RequestPermissions {},
}
#[derive(Clone, Serialize)]
@@ -235,6 +240,142 @@ pub struct GuardianReviewEventParams {
pub total_tokens: Option<i64>,
}
pub struct GuardianReviewTrackContext {
thread_id: String,
turn_id: String,
review_id: String,
target_item_id: Option<String>,
approval_request_source: GuardianApprovalRequestSource,
reviewed_action: GuardianReviewedAction,
review_timeout_ms: u64,
started_at: u64,
started_instant: Instant,
}
impl GuardianReviewTrackContext {
pub fn new(
thread_id: String,
turn_id: String,
review_id: String,
target_item_id: Option<String>,
approval_request_source: GuardianApprovalRequestSource,
reviewed_action: GuardianReviewedAction,
review_timeout_ms: u64,
) -> Self {
Self {
thread_id,
turn_id,
review_id,
target_item_id,
approval_request_source,
reviewed_action,
review_timeout_ms,
started_at: now_unix_seconds(),
started_instant: Instant::now(),
}
}
pub(crate) fn event_params(
&self,
result: GuardianReviewAnalyticsResult,
) -> GuardianReviewEventParams {
GuardianReviewEventParams {
thread_id: self.thread_id.clone(),
turn_id: self.turn_id.clone(),
review_id: self.review_id.clone(),
target_item_id: self.target_item_id.clone(),
approval_request_source: self.approval_request_source,
reviewed_action: self.reviewed_action.clone(),
reviewed_action_truncated: result.reviewed_action_truncated,
decision: result.decision,
terminal_status: result.terminal_status,
failure_reason: result.failure_reason,
risk_level: result.risk_level,
user_authorization: result.user_authorization,
outcome: result.outcome,
guardian_thread_id: result.guardian_thread_id,
guardian_session_kind: result.guardian_session_kind,
guardian_model: result.guardian_model,
guardian_reasoning_effort: result.guardian_reasoning_effort,
had_prior_review_context: result.had_prior_review_context,
review_timeout_ms: self.review_timeout_ms,
// TODO(rhan-oai): plumb nested Guardian review session tool-call counts.
tool_call_count: None,
time_to_first_token_ms: result.time_to_first_token_ms,
completion_latency_ms: Some(self.started_instant.elapsed().as_millis() as u64),
started_at: self.started_at,
completed_at: Some(now_unix_seconds()),
input_tokens: result.token_usage.as_ref().map(|usage| usage.input_tokens),
cached_input_tokens: result
.token_usage
.as_ref()
.map(|usage| usage.cached_input_tokens),
output_tokens: result.token_usage.as_ref().map(|usage| usage.output_tokens),
reasoning_output_tokens: result
.token_usage
.as_ref()
.map(|usage| usage.reasoning_output_tokens),
total_tokens: result.token_usage.as_ref().map(|usage| usage.total_tokens),
}
}
}
#[derive(Debug)]
pub struct GuardianReviewAnalyticsResult {
pub decision: GuardianReviewDecision,
pub terminal_status: GuardianReviewTerminalStatus,
pub failure_reason: Option<GuardianReviewFailureReason>,
pub risk_level: Option<GuardianRiskLevel>,
pub user_authorization: Option<GuardianUserAuthorization>,
pub outcome: Option<GuardianAssessmentOutcome>,
pub guardian_thread_id: Option<String>,
pub guardian_session_kind: Option<GuardianReviewSessionKind>,
pub guardian_model: Option<String>,
pub guardian_reasoning_effort: Option<String>,
pub had_prior_review_context: Option<bool>,
pub reviewed_action_truncated: bool,
pub token_usage: Option<TokenUsage>,
pub time_to_first_token_ms: Option<u64>,
}
impl GuardianReviewAnalyticsResult {
pub fn without_session() -> Self {
Self {
decision: GuardianReviewDecision::Denied,
terminal_status: GuardianReviewTerminalStatus::FailedClosed,
failure_reason: None,
risk_level: None,
user_authorization: None,
outcome: None,
guardian_thread_id: None,
guardian_session_kind: None,
guardian_model: None,
guardian_reasoning_effort: None,
had_prior_review_context: None,
reviewed_action_truncated: false,
token_usage: None,
time_to_first_token_ms: None,
}
}
pub fn from_session(
guardian_thread_id: String,
guardian_session_kind: GuardianReviewSessionKind,
guardian_model: String,
guardian_reasoning_effort: Option<String>,
had_prior_review_context: bool,
) -> Self {
Self {
guardian_thread_id: Some(guardian_thread_id),
guardian_session_kind: Some(guardian_session_kind),
guardian_model: Some(guardian_model),
guardian_reasoning_effort,
had_prior_review_context: Some(had_prior_review_context),
..Self::without_session()
}
}
}
#[derive(Serialize)]
pub(crate) struct GuardianReviewEventPayload {
pub(crate) app_server_client: CodexAppServerClientMetadata,

View File

@@ -9,11 +9,13 @@ use std::time::UNIX_EPOCH;
pub use client::AnalyticsEventsClient;
pub use events::AppServerRpcTransport;
pub use events::GuardianApprovalRequestSource;
pub use events::GuardianReviewAnalyticsResult;
pub use events::GuardianReviewDecision;
pub use events::GuardianReviewEventParams;
pub use events::GuardianReviewFailureReason;
pub use events::GuardianReviewSessionKind;
pub use events::GuardianReviewTerminalStatus;
pub use events::GuardianReviewTrackContext;
pub use events::GuardianReviewedAction;
pub use facts::AnalyticsJsonRpcError;
pub use facts::AppInvocation;

View File

@@ -1,5 +1,6 @@
use std::path::Path;
use codex_analytics::GuardianReviewedAction;
use codex_protocol::approvals::GuardianAssessmentAction;
use codex_protocol::approvals::GuardianCommandSource;
use codex_protocol::approvals::NetworkApprovalProtocol;
@@ -390,6 +391,66 @@ pub(crate) fn guardian_assessment_action(
}
}
pub(crate) fn guardian_reviewed_action(
request: &GuardianApprovalRequest,
) -> GuardianReviewedAction {
match request {
GuardianApprovalRequest::Shell {
sandbox_permissions,
additional_permissions,
..
} => GuardianReviewedAction::Shell {
sandbox_permissions: *sandbox_permissions,
additional_permissions: additional_permissions.clone(),
},
GuardianApprovalRequest::ExecCommand {
sandbox_permissions,
additional_permissions,
tty,
..
} => GuardianReviewedAction::UnifiedExec {
sandbox_permissions: *sandbox_permissions,
additional_permissions: additional_permissions.clone(),
tty: *tty,
},
#[cfg(unix)]
GuardianApprovalRequest::Execve {
source,
program,
additional_permissions,
..
} => GuardianReviewedAction::Execve {
source: *source,
program: program.clone(),
additional_permissions: additional_permissions.clone(),
},
GuardianApprovalRequest::ApplyPatch { .. } => GuardianReviewedAction::ApplyPatch {},
GuardianApprovalRequest::NetworkAccess { protocol, port, .. } => {
GuardianReviewedAction::NetworkAccess {
protocol: *protocol,
port: *port,
}
}
GuardianApprovalRequest::McpToolCall {
server,
tool_name,
connector_id,
connector_name,
tool_title,
..
} => GuardianReviewedAction::McpToolCall {
server: server.clone(),
tool_name: tool_name.clone(),
connector_id: connector_id.clone(),
connector_name: connector_name.clone(),
tool_title: tool_title.clone(),
},
GuardianApprovalRequest::RequestPermissions { .. } => {
GuardianReviewedAction::RequestPermissions {}
}
}
}
pub(crate) fn guardian_request_target_item_id(request: &GuardianApprovalRequest) -> Option<&str> {
match request {
GuardianApprovalRequest::Shell { id, .. }

View File

@@ -1,5 +1,12 @@
use std::sync::Arc;
use codex_analytics::GuardianApprovalRequestSource;
use codex_analytics::GuardianReviewAnalyticsResult;
use codex_analytics::GuardianReviewDecision;
use codex_analytics::GuardianReviewFailureReason;
use codex_analytics::GuardianReviewTerminalStatus;
use codex_analytics::GuardianReviewTrackContext;
use codex_features::Feature;
use codex_protocol::config_types::ApprovalsReviewer;
use codex_protocol::protocol::AskForApproval;
use codex_protocol::protocol::EventMsg;
@@ -17,6 +24,7 @@ use tokio_util::sync::CancellationToken;
use crate::session::session::Session;
use crate::session::turn_context::TurnContext;
use super::GUARDIAN_REVIEW_TIMEOUT;
use super::GUARDIAN_REVIEWER_NAME;
use super::GuardianApprovalRequest;
use super::GuardianAssessment;
@@ -25,6 +33,7 @@ use super::GuardianRejection;
use super::approval_request::guardian_assessment_action;
use super::approval_request::guardian_request_target_item_id;
use super::approval_request::guardian_request_turn_id;
use super::approval_request::guardian_reviewed_action;
use super::prompt::guardian_output_schema;
use super::prompt::parse_guardian_assessment;
use super::review_session::GuardianReviewSessionOutcome;
@@ -76,11 +85,55 @@ pub(crate) fn guardian_timeout_message() -> String {
#[derive(Debug)]
pub(super) enum GuardianReviewOutcome {
Completed(anyhow::Result<GuardianAssessment>),
Completed(GuardianAssessment),
Failed(GuardianReviewError),
TimedOut,
Aborted,
}
#[derive(Debug)]
pub(super) enum GuardianReviewError {
PromptBuild { message: String },
Session { message: String },
Parse { message: String },
}
impl GuardianReviewError {
fn prompt_build(err: anyhow::Error) -> Self {
Self::PromptBuild {
message: err.to_string(),
}
}
fn session(err: anyhow::Error) -> Self {
Self::Session {
message: err.to_string(),
}
}
fn parse(err: anyhow::Error) -> Self {
Self::Parse {
message: err.to_string(),
}
}
fn reason(&self) -> GuardianReviewFailureReason {
match self {
Self::PromptBuild { .. } => GuardianReviewFailureReason::PromptBuildError,
Self::Session { .. } => GuardianReviewFailureReason::SessionError,
Self::Parse { .. } => GuardianReviewFailureReason::ParseError,
}
}
fn message(&self) -> &str {
match self {
Self::PromptBuild { message } | Self::Session { message } | Self::Parse { message } => {
message
}
}
}
}
fn guardian_risk_level_str(level: GuardianRiskLevel) -> &'static str {
match level {
GuardianRiskLevel::Low => "low",
@@ -110,6 +163,21 @@ pub(crate) fn is_guardian_reviewer_source(
)
}
fn track_guardian_review(
session: &Session,
turn: &TurnContext,
tracking: &GuardianReviewTrackContext,
result: GuardianReviewAnalyticsResult,
) {
if !turn.config.features.enabled(Feature::GeneralAnalytics) {
return;
}
session
.services
.analytics_events_client
.track_guardian_review(tracking, result);
}
/// This function always fails closed: timeouts, review-session failures, and
/// parse failures all block execution, but timeouts are still surfaced to the
/// caller as distinct from explicit guardian denials.
@@ -119,11 +187,21 @@ async fn run_guardian_review(
review_id: String,
request: GuardianApprovalRequest,
retry_reason: Option<String>,
approval_request_source: GuardianApprovalRequestSource,
external_cancel: Option<CancellationToken>,
) -> ReviewDecision {
let target_item_id = guardian_request_target_item_id(&request).map(str::to_string);
let assessment_turn_id = guardian_request_turn_id(&request, &turn.sub_id).to_string();
let action_summary = guardian_assessment_action(&request);
let review_tracking = GuardianReviewTrackContext::new(
session.conversation_id.to_string(),
assessment_turn_id.clone(),
review_id.clone(),
target_item_id.clone(),
approval_request_source,
guardian_reviewed_action(&request),
GUARDIAN_REVIEW_TIMEOUT.as_millis() as u64,
);
session
.send_event(
turn.as_ref(),
@@ -145,6 +223,17 @@ async fn run_guardian_review(
.as_ref()
.is_some_and(CancellationToken::is_cancelled)
{
track_guardian_review(
session.as_ref(),
turn.as_ref(),
&review_tracking,
GuardianReviewAnalyticsResult {
decision: GuardianReviewDecision::Aborted,
terminal_status: GuardianReviewTerminalStatus::Aborted,
failure_reason: Some(GuardianReviewFailureReason::Cancelled),
..GuardianReviewAnalyticsResult::without_session()
},
);
session
.send_event(
turn.as_ref(),
@@ -166,28 +255,78 @@ async fn run_guardian_review(
let schema = guardian_output_schema();
let terminal_action = action_summary.clone();
let outcome = Box::pin(run_guardian_review_session(
let (outcome, analytics_result) = Box::pin(run_guardian_review_session(
session.clone(),
turn.clone(),
request,
retry_reason,
retry_reason.clone(),
schema,
external_cancel,
))
.await;
let assessment = match outcome {
GuardianReviewOutcome::Completed(Ok(assessment)) => assessment,
GuardianReviewOutcome::Completed(Err(err)) => GuardianAssessment {
risk_level: GuardianRiskLevel::High,
user_authorization: GuardianUserAuthorization::Unknown,
outcome: GuardianAssessmentOutcome::Deny,
rationale: format!("Automatic approval review failed: {err}"),
},
GuardianReviewOutcome::Completed(assessment) => {
let approved = matches!(assessment.outcome, GuardianAssessmentOutcome::Allow);
track_guardian_review(
session.as_ref(),
turn.as_ref(),
&review_tracking,
GuardianReviewAnalyticsResult {
decision: if approved {
GuardianReviewDecision::Approved
} else {
GuardianReviewDecision::Denied
},
terminal_status: if approved {
GuardianReviewTerminalStatus::Approved
} else {
GuardianReviewTerminalStatus::Denied
},
failure_reason: None,
risk_level: Some(assessment.risk_level),
user_authorization: Some(assessment.user_authorization),
outcome: Some(assessment.outcome),
..analytics_result
},
);
assessment
}
GuardianReviewOutcome::Failed(failure) => {
let rationale = format!("Automatic approval review failed: {}", failure.message());
track_guardian_review(
session.as_ref(),
turn.as_ref(),
&review_tracking,
GuardianReviewAnalyticsResult {
decision: GuardianReviewDecision::Denied,
terminal_status: GuardianReviewTerminalStatus::FailedClosed,
failure_reason: Some(failure.reason()),
..analytics_result
},
);
GuardianAssessment {
risk_level: GuardianRiskLevel::High,
user_authorization: GuardianUserAuthorization::Unknown,
outcome: GuardianAssessmentOutcome::Deny,
rationale,
}
}
GuardianReviewOutcome::TimedOut => {
let rationale =
"Automatic approval review timed out while evaluating the requested approval."
.to_string();
track_guardian_review(
session.as_ref(),
turn.as_ref(),
&review_tracking,
GuardianReviewAnalyticsResult {
decision: GuardianReviewDecision::Denied,
terminal_status: GuardianReviewTerminalStatus::TimedOut,
failure_reason: Some(GuardianReviewFailureReason::Timeout),
..analytics_result
},
);
session
.send_event(
turn.as_ref(),
@@ -215,6 +354,17 @@ async fn run_guardian_review(
return ReviewDecision::TimedOut;
}
GuardianReviewOutcome::Aborted => {
track_guardian_review(
session.as_ref(),
turn.as_ref(),
&review_tracking,
GuardianReviewAnalyticsResult {
decision: GuardianReviewDecision::Aborted,
terminal_status: GuardianReviewTerminalStatus::Aborted,
failure_reason: Some(GuardianReviewFailureReason::Cancelled),
..analytics_result
},
);
session
.send_event(
turn.as_ref(),
@@ -314,6 +464,7 @@ pub(crate) async fn review_approval_request(
review_id,
request,
retry_reason,
GuardianApprovalRequestSource::MainTurn,
/*external_cancel*/ None,
))
.await
@@ -325,16 +476,18 @@ pub(crate) async fn review_approval_request_with_cancel(
review_id: String,
request: GuardianApprovalRequest,
retry_reason: Option<String>,
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
}
@@ -361,6 +514,7 @@ pub(crate) fn spawn_approval_request_review(
review_id,
request,
retry_reason,
GuardianApprovalRequestSource::DelegatedSubagent,
cancel_token,
));
let _ = tx.send(decision);
@@ -389,11 +543,16 @@ pub(super) async fn run_guardian_review_session(
retry_reason: Option<String>,
schema: serde_json::Value,
external_cancel: Option<CancellationToken>,
) -> GuardianReviewOutcome {
) -> (GuardianReviewOutcome, GuardianReviewAnalyticsResult) {
let live_network_config = match session.services.network_proxy.as_ref() {
Some(network_proxy) => match network_proxy.proxy().current_cfg().await {
Ok(config) => Some(config),
Err(err) => return GuardianReviewOutcome::Completed(Err(err)),
Err(err) => {
return (
GuardianReviewOutcome::Failed(GuardianReviewError::prompt_build(err)),
GuardianReviewAnalyticsResult::without_session(),
);
}
},
None => None,
};
@@ -443,10 +602,15 @@ pub(super) async fn run_guardian_review_session(
);
let guardian_config = match guardian_config {
Ok(config) => config,
Err(err) => return GuardianReviewOutcome::Completed(Err(err)),
Err(err) => {
return (
GuardianReviewOutcome::Failed(GuardianReviewError::prompt_build(err)),
GuardianReviewAnalyticsResult::without_session(),
);
}
};
match Box::pin(
let (session_outcome, session_analytics_result) = Box::pin(
session
.guardian_review_session
.run_review(GuardianReviewSessionParams {
@@ -463,17 +627,69 @@ pub(super) async fn run_guardian_review_session(
external_cancel,
}),
)
.await
{
GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) => {
GuardianReviewOutcome::Completed(parse_guardian_assessment(
last_agent_message.as_deref(),
))
.await;
match session_outcome {
GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) => match last_agent_message
{
Some(last_agent_message) => {
match parse_guardian_assessment(Some(&last_agent_message)) {
Ok(assessment) => (
GuardianReviewOutcome::Completed(assessment),
session_analytics_result,
),
Err(err) => (
GuardianReviewOutcome::Failed(GuardianReviewError::parse(err)),
session_analytics_result,
),
}
}
None => (
GuardianReviewOutcome::Failed(GuardianReviewError::session(anyhow::anyhow!(
"guardian review completed without an assessment payload"
))),
session_analytics_result,
),
},
GuardianReviewSessionOutcome::Completed(Err(err)) => (
GuardianReviewOutcome::Failed(GuardianReviewError::session(err)),
session_analytics_result,
),
GuardianReviewSessionOutcome::PromptBuildFailed(err) => (
GuardianReviewOutcome::Failed(GuardianReviewError::prompt_build(err)),
session_analytics_result,
),
GuardianReviewSessionOutcome::TimedOut => {
(GuardianReviewOutcome::TimedOut, session_analytics_result)
}
GuardianReviewSessionOutcome::Completed(Err(err)) => {
GuardianReviewOutcome::Completed(Err(err))
GuardianReviewSessionOutcome::Aborted => {
(GuardianReviewOutcome::Aborted, session_analytics_result)
}
GuardianReviewSessionOutcome::TimedOut => GuardianReviewOutcome::TimedOut,
GuardianReviewSessionOutcome::Aborted => GuardianReviewOutcome::Aborted,
}
}
#[cfg(test)]
mod review_tests {
use super::*;
#[test]
fn guardian_review_error_reason_distinguishes_error_kinds() {
let parse_error = GuardianReviewError::parse(anyhow::anyhow!("bad guardian JSON"));
let prompt_error = GuardianReviewError::prompt_build(anyhow::anyhow!("bad prompt/config"));
let session_error =
GuardianReviewError::session(anyhow::anyhow!("guardian runtime failed"));
assert!(matches!(
parse_error.reason(),
GuardianReviewFailureReason::ParseError
));
assert!(matches!(
prompt_error.reason(),
GuardianReviewFailureReason::PromptBuildError
));
assert!(matches!(
session_error.reason(),
GuardianReviewFailureReason::SessionError
));
}
}

View File

@@ -5,6 +5,8 @@ use std::sync::Arc;
use std::time::Duration;
use anyhow::anyhow;
use codex_analytics::GuardianReviewAnalyticsResult;
use codex_analytics::GuardianReviewSessionKind;
use codex_protocol::config_types::Personality;
use codex_protocol::config_types::ReasoningSummary as ReasoningSummaryConfig;
use codex_protocol::models::DeveloperInstructions;
@@ -17,6 +19,7 @@ use codex_protocol::protocol::Op;
use codex_protocol::protocol::RolloutItem;
use codex_protocol::protocol::SandboxPolicy;
use codex_protocol::protocol::SubAgentSource;
use codex_protocol::protocol::TokenUsage;
use serde_json::Value;
use tokio::sync::Mutex;
use tokio::sync::Semaphore;
@@ -59,6 +62,7 @@ const GUARDIAN_FOLLOWUP_REVIEW_REMINDER: &str = concat!(
#[derive(Debug)]
pub(crate) enum GuardianReviewSessionOutcome {
Completed(anyhow::Result<Option<String>>),
PromptBuildFailed(anyhow::Error),
TimedOut,
Aborted,
}
@@ -102,6 +106,21 @@ struct GuardianReviewState {
last_committed_fork_snapshot: Option<GuardianReviewForkSnapshot>,
}
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<Mutex<GuardianReviewSessionState>>,
review_session: Option<Arc<GuardianReviewSession>>,
@@ -269,13 +288,14 @@ impl GuardianReviewSessionManager {
clippy::await_holding_invalid_type,
reason = "review session selection and trunk spawning must stay serialized"
)]
pub(crate) async fn run_review(
pub(super) async fn run_review(
&self,
params: GuardianReviewSessionParams,
) -> GuardianReviewSessionOutcome {
) -> (GuardianReviewSessionOutcome, GuardianReviewAnalyticsResult) {
let deadline = tokio::time::Instant::now() + GUARDIAN_REVIEW_TIMEOUT;
let next_reuse_key = GuardianReviewSessionReuseKey::from_spawn_config(&params.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(),
@@ -309,16 +329,22 @@ impl GuardianReviewSessionManager {
{
Ok(Ok(review_session)) => Arc::new(review_session),
Ok(Err(err)) => {
return GuardianReviewSessionOutcome::Completed(Err(err));
return (
GuardianReviewSessionOutcome::PromptBuildFailed(err),
GuardianReviewAnalyticsResult::without_session(),
);
}
Err(outcome) => {
return (outcome, GuardianReviewAnalyticsResult::without_session());
}
Err(outcome) => return outcome,
};
state.trunk = Some(Arc::clone(&review_session));
spawned_trunk = true;
}
state.trunk.as_ref().cloned()
}
Err(outcome) => return outcome,
Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()),
};
if let Some(review_session) = stale_trunk_to_shutdown {
@@ -326,9 +352,12 @@ impl GuardianReviewSessionManager {
}
let Some(trunk) = trunk_candidate else {
return GuardianReviewSessionOutcome::Completed(Err(anyhow!(
"guardian review session was not available after spawn"
)));
return (
GuardianReviewSessionOutcome::Completed(Err(anyhow!(
"guardian review session was not available after spawn"
))),
GuardianReviewAnalyticsResult::without_session(),
);
};
if trunk.reuse_key != next_reuse_key {
@@ -354,20 +383,30 @@ impl GuardianReviewSessionManager {
}
};
let (outcome, keep_review_session) =
Box::pin(run_review_on_session(trunk.as_ref(), &params, deadline)).await;
let guardian_session_kind = if spawned_trunk {
GuardianReviewSessionKind::TrunkNew
} else {
GuardianReviewSessionKind::TrunkReused
};
let (outcome, keep_review_session, analytics_result) = Box::pin(run_review_on_session(
trunk.as_ref(),
&params,
guardian_session_kind,
deadline,
))
.await;
if keep_review_session && matches!(outcome, GuardianReviewSessionOutcome::Completed(_)) {
trunk.refresh_last_committed_fork_snapshot().await;
}
drop(trunk_guard);
if keep_review_session {
outcome
(outcome, analytics_result)
} else {
if let Some(review_session) = self.remove_trunk_if_current(&trunk).await {
review_session.shutdown_in_background();
}
outcome
(outcome, analytics_result)
}
}
@@ -464,7 +503,7 @@ impl GuardianReviewSessionManager {
reuse_key: GuardianReviewSessionReuseKey,
deadline: tokio::time::Instant,
fork_snapshot: Option<GuardianReviewForkSnapshot>,
) -> GuardianReviewSessionOutcome {
) -> (GuardianReviewSessionOutcome, GuardianReviewAnalyticsResult) {
let spawn_cancel_token = CancellationToken::new();
let mut fork_config = params.spawn_config.clone();
fork_config.ephemeral = true;
@@ -483,17 +522,23 @@ impl GuardianReviewSessionManager {
.await
{
Ok(Ok(review_session)) => Arc::new(review_session),
Ok(Err(err)) => return GuardianReviewSessionOutcome::Completed(Err(err)),
Err(outcome) => return outcome,
Ok(Err(err)) => {
return (
GuardianReviewSessionOutcome::PromptBuildFailed(err),
GuardianReviewAnalyticsResult::without_session(),
);
}
Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()),
};
self.register_active_ephemeral(Arc::clone(&review_session))
.await;
let mut cleanup =
EphemeralReviewCleanup::new(Arc::clone(&self.state), Arc::clone(&review_session));
let (outcome, _) = Box::pin(run_review_on_session(
let (outcome, _, analytics_result) = Box::pin(run_review_on_session(
review_session.as_ref(),
&params,
GuardianReviewSessionKind::EphemeralForked,
deadline,
))
.await;
@@ -501,7 +546,7 @@ impl GuardianReviewSessionManager {
cleanup.disarm();
review_session.shutdown_in_background();
}
outcome
(outcome, analytics_result)
}
}
@@ -548,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,
GuardianReviewAnalyticsResult,
) {
let (send_followup_reminder, prompt_mode) = {
let state = review_session.state.lock().await;
@@ -564,6 +614,29 @@ async fn run_review_on_session(
(send_followup_reminder, prompt_mode)
};
let model_info = params
.parent_session
.services
.models_manager
.get_model_info(
params.model.as_str(),
&params.spawn_config.to_models_manager_config(),
)
.await;
let guardian_reasoning_effort = if model_info.supports_reasoning_summaries {
params
.reasoning_effort
.or(model_info.default_reasoning_level)
} else {
None
};
let mut analytics_result = GuardianReviewAnalyticsResult::from_session(
review_session.codex.session.conversation_id.to_string(),
guardian_session_kind,
params.model.clone(),
guardian_reasoning_effort.map(|effort| effort.to_string()),
had_prior_review_context(&prompt_mode),
);
if send_followup_reminder {
append_guardian_followup_reminder(review_session).await;
}
@@ -588,6 +661,12 @@ async fn run_review_on_session(
prompt_mode,
)
.await?;
let token_usage_at_review_start = review_session
.codex
.session
.total_token_usage()
.await
.unwrap_or_default();
review_session
.codex
@@ -607,29 +686,44 @@ async fn run_review_on_session(
})
.await?;
Ok::<GuardianTranscriptCursor, anyhow::Error>(prompt_items.transcript_cursor)
Ok::<(GuardianTranscriptCursor, TokenUsage), anyhow::Error>((
prompt_items.transcript_cursor,
token_usage_at_review_start,
))
}),
)
.await;
let submit_result = match submit_result {
Ok(submit_result) => submit_result,
Err(outcome) => return (outcome, false),
Err(outcome) => return (outcome, false, analytics_result),
};
let transcript_cursor = match submit_result {
Ok(transcript_cursor) => transcript_cursor,
let (transcript_cursor, token_usage_at_review_start) = match submit_result {
Ok(submit_result) => submit_result,
Err(err) => {
return (GuardianReviewSessionOutcome::Completed(Err(err)), false);
return (
GuardianReviewSessionOutcome::PromptBuildFailed(err),
false,
analytics_result,
);
}
};
let outcome =
wait_for_guardian_review(review_session, deadline, params.external_cancel.as_ref()).await;
if matches!(outcome.0, GuardianReviewSessionOutcome::Completed(_)) {
if outcome.2
&& let Some(total_token_usage) = review_session.codex.session.total_token_usage().await
{
analytics_result.token_usage = Some(token_usage_delta(
&token_usage_at_review_start,
&total_token_usage,
));
}
let mut state = review_session.state.lock().await;
state.prior_review_count = state.prior_review_count.saturating_add(1);
state.last_reviewed_transcript_cursor = Some(transcript_cursor);
}
outcome
(outcome.0, outcome.1, analytics_result)
}
async fn append_guardian_followup_reminder(review_session: &GuardianReviewSession) {
@@ -658,7 +752,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<String> = None;
@@ -667,7 +761,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 {
@@ -677,7 +771,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 {
@@ -689,18 +783,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);
}
_ => {}
},
@@ -708,6 +804,7 @@ async fn wait_for_guardian_review(
return (
GuardianReviewSessionOutcome::Completed(Err(err.into())),
false,
false,
);
}
}
@@ -1005,4 +1102,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,
}
);
}
}

View File

@@ -15,6 +15,7 @@ use crate::config_loader::Sourced;
use crate::session::session::Session;
use crate::session::turn_context::TurnContext;
use crate::test_support;
use codex_analytics::GuardianApprovalRequestSource;
use codex_config::config_toml::ConfigToml;
use codex_config::types::McpServerConfig;
use codex_exec_server::LOCAL_FS;
@@ -707,6 +708,7 @@ async fn cancelled_guardian_review_emits_terminal_abort_without_warning() {
.to_string(),
},
/*retry_reason*/ None,
GuardianApprovalRequestSource::MainTurn,
cancel_token,
)
.await;
@@ -941,7 +943,7 @@ async fn guardian_review_request_layout_matches_model_visible_request_snapshot()
/*external_cancel*/ None,
)
.await;
let GuardianReviewOutcome::Completed(Ok(assessment)) = outcome else {
let (GuardianReviewOutcome::Completed(assessment), _) = outcome else {
panic!("expected guardian assessment");
};
assert_eq!(assessment.outcome, GuardianAssessmentOutcome::Allow);
@@ -1150,13 +1152,13 @@ async fn guardian_reuses_prompt_cache_key_and_appends_prior_reviews() -> anyhow:
)
.await;
let GuardianReviewOutcome::Completed(Ok(first_assessment)) = first_outcome else {
let (GuardianReviewOutcome::Completed(first_assessment), _) = first_outcome else {
panic!("expected first guardian assessment");
};
let GuardianReviewOutcome::Completed(Ok(second_assessment)) = second_outcome else {
let (GuardianReviewOutcome::Completed(second_assessment), _) = second_outcome else {
panic!("expected second guardian assessment");
};
let GuardianReviewOutcome::Completed(Ok(third_assessment)) = third_outcome else {
let (GuardianReviewOutcome::Completed(third_assessment), _) = third_outcome else {
panic!("expected third guardian assessment");
};
assert_eq!(first_assessment.outcome, GuardianAssessmentOutcome::Allow);