From 66dfadbe6680a5aa2dad0dc80f3b1092d7a59258 Mon Sep 17 00:00:00 2001 From: felixxia-oai Date: Wed, 16 Sep 2026 16:44:59 +0000 Subject: [PATCH] Replace guardian review result tuples with named structs (#45985) ## What changed Introduce `ReviewTurnResult` and `ReviewSessionResult` to carry review outcomes, session disposition, and completion or analytics data. Propagate `SessionDisposition` directly through review execution and pooling, replacing boolean reuse flags and positional tuple access while preserving existing behavior. Update existing review-session tests to assert the named fields and explicit session dispositions. GitOrigin-RevId: 93eea8b482cd592167956444042becc1ef5f2eb5 --- codex-rs/core/src/guardian/review_session.rs | 106 ++++++++++++------ .../core/src/guardian/review_session_setup.rs | 18 +-- .../core/src/guardian/review_session_tests.rs | 73 ++++++++---- .../ext/guardian-reviewer/src/execution.rs | 81 ++++++++----- codex-rs/ext/guardian-reviewer/src/lib.rs | 2 + codex-rs/ext/guardian-reviewer/src/pool.rs | 25 +++-- 6 files changed, 192 insertions(+), 113 deletions(-) diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index 88cc7acf84..cfb6b1aa47 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -26,6 +26,8 @@ use codex_extension_api::Instructions; use codex_guardian_reviewer::ConversationCheckpoint; use codex_guardian_reviewer::ConversationState; use codex_guardian_reviewer::ReviewModel; +use codex_guardian_reviewer::ReviewSessionResult; +use codex_guardian_reviewer::SessionDisposition; use codex_history::InitialHistory; use codex_history::RolloutItem; use codex_protocol::ThreadId; @@ -315,11 +317,7 @@ async fn run_review_on_session( params: &GuardianReviewSessionParams, guardian_session_kind: GuardianReviewSessionKind, deadline: tokio::time::Instant, -) -> ( - GuardianReviewSessionOutcome, - bool, - GuardianReviewAnalyticsResult, -) { +) -> ReviewSessionResult { let review_model = ¶ms.review_model; let model_info = params .parent_session @@ -367,17 +365,23 @@ async fn run_review_on_session( { Ok(Ok(())) => {} Ok(Err(error)) => { - return ( - GuardianReviewSessionOutcome::SessionFailed { + return ReviewSessionResult { + outcome: GuardianReviewSessionOutcome::SessionFailed { error, error_info: None, retry_at: None, }, - false, - analytics_result, - ); + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; + } + Err(outcome) => { + return ReviewSessionResult { + outcome, + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; } - Err(outcome) => return (outcome, false, analytics_result), } if params.spawn_config.features.enabled(Feature::TokenBudget) @@ -399,19 +403,25 @@ async fn run_review_on_session( let compact_turn_id = match compact_submission { Ok(Ok(turn_id)) => turn_id, Ok(Err(error)) => { - return ( - GuardianReviewSessionOutcome::SessionFailed { + return ReviewSessionResult { + outcome: GuardianReviewSessionOutcome::SessionFailed { error: error.into(), error_info: None, retry_at: None, }, - false, - analytics_result, - ); + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; + } + Err(outcome) => { + return ReviewSessionResult { + outcome, + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; } - Err(outcome) => return (outcome, false, analytics_result), }; - let (outcome, keep_review_session, _) = wait_for_guardian_review( + let result = wait_for_guardian_review( review_session, &compact_turn_id, deadline, @@ -419,8 +429,15 @@ async fn run_review_on_session( &mut analytics_result, ) .await; - if !matches!(outcome, GuardianReviewSessionOutcome::Completed(Ok(_))) { - return (outcome, keep_review_session, analytics_result); + if !matches!( + result.outcome, + GuardianReviewSessionOutcome::Completed(Ok(_)) + ) { + return ReviewSessionResult { + outcome: result.outcome, + disposition: result.disposition, + analytics: analytics_result, + }; } if prior_review_count > 0 { @@ -580,16 +597,22 @@ async fn run_review_on_session( .await; let prompt_items = match prompt_items { Ok(prompt_items) => prompt_items, - Err(outcome) => return (outcome, false, analytics_result), + Err(outcome) => { + return ReviewSessionResult { + outcome, + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; + } }; let (prompt_items, items) = match prompt_items { Ok(prompt_items) => prompt_items, Err(err) => { - return ( - GuardianReviewSessionOutcome::PromptBuildFailed(err), - false, - analytics_result, - ); + return ReviewSessionResult { + outcome: GuardianReviewSessionOutcome::PromptBuildFailed(err), + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; } }; let transcript_cursor = prompt_items.transcript_cursor; @@ -662,7 +685,11 @@ async fn run_review_on_session( .services .thread_extension_data .remove::(); - return (outcome, false, analytics_result); + return ReviewSessionResult { + outcome, + disposition: SessionDisposition::Discard, + analytics: analytics_result, + }; } }; if let Some(response_sequence) = node_repl_evidence_admission { @@ -673,7 +700,7 @@ async fn run_review_on_session( }); } - let outcome = wait_for_guardian_review( + let turn_result = wait_for_guardian_review( review_session, child_turn_id.as_str(), deadline, @@ -686,8 +713,11 @@ async fn run_review_on_session( .services .thread_extension_data .remove::(); - if matches!(outcome.0, GuardianReviewSessionOutcome::Completed(_)) { - if outcome.2 + if matches!( + turn_result.outcome, + GuardianReviewSessionOutcome::Completed(_) + ) { + if turn_result.turn_completed && let Some(total_token_usage) = review_session.session.total_token_usage().await { analytics_result.token_usage = Some(token_usage_delta( @@ -703,7 +733,7 @@ async fn run_review_on_session( .services .thread_extension_data .remove::(); - let result = match outcome.0 { + let result = match turn_result.outcome { GuardianReviewSessionOutcome::SessionFailed { error_info: Some(CodexErrorInfo::ContextWindowExceeded), .. @@ -716,11 +746,15 @@ async fn run_review_on_session( } result => result, }; - ( - result, - outcome.1 && budget_exhausted.is_none(), - analytics_result, - ) + ReviewSessionResult { + outcome: result, + disposition: if budget_exhausted.is_some() { + SessionDisposition::Discard + } else { + turn_result.disposition + }, + analytics: analytics_result, + } } async fn ensure_guardian_followup_reminder(review_session: &GuardianReviewSession) { diff --git a/codex-rs/core/src/guardian/review_session_setup.rs b/codex-rs/core/src/guardian/review_session_setup.rs index 41ead25c63..c87b8cb622 100644 --- a/codex-rs/core/src/guardian/review_session_setup.rs +++ b/codex-rs/core/src/guardian/review_session_setup.rs @@ -4,7 +4,6 @@ use super::*; use codex_guardian_reviewer::ReviewerPool; use codex_guardian_reviewer::ReviewerRequest; -use codex_guardian_reviewer::SessionDisposition; pub struct PreparedGuardianContext { parent: Arc, @@ -182,25 +181,16 @@ impl ReviewerRequest for PreparedReview { &self, session: &GuardianReviewSession, kind: GuardianReviewSessionKind, - ) -> ( - GuardianReviewSessionOutcome, - SessionDisposition, - GuardianReviewAnalyticsResult, - ) { - let (outcome, keep_session, analytics) = Box::pin(run_review_on_session( + ) -> ReviewSessionResult { + let result = Box::pin(run_review_on_session( session, &self.params, kind, self.params.deadline, )) .await; - record_failed_review(&session.session, &self.params, &outcome).await; - let disposition = if keep_session { - SessionDisposition::Reusable - } else { - SessionDisposition::Discard - }; - (outcome, disposition, analytics) + record_failed_review(&session.session, &self.params, &result.outcome).await; + result } } diff --git a/codex-rs/core/src/guardian/review_session_tests.rs b/codex-rs/core/src/guardian/review_session_tests.rs index 3cf37785b6..b981a8677b 100644 --- a/codex-rs/core/src/guardian/review_session_tests.rs +++ b/codex-rs/core/src/guardian/review_session_tests.rs @@ -842,14 +842,17 @@ async fn run_review_on_reused_session_waits_for_submitted_turn() { .await .expect("queue submitted turn completion"); - let (outcome, keep_review_session, analytics_result) = - review.await.expect("review task should complete"); + let ReviewSessionResult { + outcome, + disposition, + analytics: analytics_result, + } = review.await.expect("review task should complete"); let GuardianReviewSessionOutcome::Completed(Ok(last_agent_message)) = outcome else { panic!("expected submitted turn completion"); }; assert_eq!(last_agent_message.as_deref(), Some("fresh")); assert_eq!(analytics_result.time_to_first_token_ms, Some(42)); - assert!(keep_review_session); + assert_eq!(disposition, SessionDisposition::Reusable); } #[tokio::test] @@ -916,7 +919,11 @@ async fn wait_for_guardian_review_ignores_prior_turn_completion() { .expect("queue current turn completion"); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_secs(1), @@ -930,8 +937,8 @@ async fn wait_for_guardian_review_ignores_prior_turn_completion() { }; assert_eq!(last_agent_message.as_deref(), Some("fresh")); assert_eq!(analytics_result.time_to_first_token_ms, Some(42)); - assert!(keep_review_session); - assert!(capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(turn_completed); } #[tokio::test] @@ -958,7 +965,11 @@ async fn wait_for_guardian_review_ignores_prior_turn_errors() { .expect("queue current turn completion"); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_secs(1), @@ -972,8 +983,8 @@ async fn wait_for_guardian_review_ignores_prior_turn_errors() { }; assert_eq!(last_agent_message, None); assert_eq!(analytics_result.time_to_first_token_ms, Some(42)); - assert!(keep_review_session); - assert!(capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(turn_completed); } #[tokio::test] @@ -1000,7 +1011,11 @@ async fn wait_for_guardian_review_preserves_structured_session_error() { .expect("queue current turn completion"); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_secs(1), @@ -1017,8 +1032,8 @@ async fn wait_for_guardian_review_preserves_structured_session_error() { }; assert_eq!(error.to_string(), "temporary failure"); assert_eq!(error_info, Some(CodexErrorInfo::ServerOverloaded)); - assert!(keep_review_session); - assert!(capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(turn_completed); } #[tokio::test] @@ -1034,7 +1049,11 @@ async fn wait_for_guardian_review_ignores_prior_turn_aborts() { .expect("queue current turn completion"); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_secs(1), @@ -1048,8 +1067,8 @@ async fn wait_for_guardian_review_ignores_prior_turn_aborts() { }; assert_eq!(last_agent_message.as_deref(), Some("fresh")); assert_eq!(analytics_result.time_to_first_token_ms, Some(42)); - assert!(keep_review_session); - assert!(capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(turn_completed); } #[tokio::test] @@ -1070,7 +1089,11 @@ async fn wait_for_guardian_review_timeout_drains_expected_turn_after_stale_termi }); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_millis(10), @@ -1083,8 +1106,8 @@ async fn wait_for_guardian_review_timeout_drains_expected_turn_after_stale_termi .await .expect("interrupt response task should complete"); assert!(matches!(outcome, GuardianReviewSessionOutcome::TimedOut)); - assert!(keep_review_session); - assert!(!capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(!turn_completed); } #[tokio::test] @@ -1107,7 +1130,11 @@ async fn wait_for_guardian_review_cancel_drains_expected_turn_after_stale_termin external_cancel.cancel(); let mut analytics_result = GuardianReviewAnalyticsResult::without_session(); - let (outcome, keep_review_session, capture_token_usage) = wait_for_guardian_review( + let codex_guardian_reviewer::ReviewTurnResult { + outcome, + disposition, + turn_completed, + } = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now() + Duration::from_secs(1), @@ -1120,8 +1147,8 @@ async fn wait_for_guardian_review_cancel_drains_expected_turn_after_stale_termin .await .expect("interrupt response task should complete"); assert!(matches!(outcome, GuardianReviewSessionOutcome::Aborted)); - assert!(keep_review_session); - assert!(!capture_token_usage); + assert_eq!(disposition, SessionDisposition::Reusable); + assert!(!turn_completed); } #[tokio::test] @@ -1138,7 +1165,7 @@ async fn interrupt_and_drain_turn_ignores_prior_turn_completion() { let cancellation = CancellationToken::new(); cancellation.cancel(); - let (_, reusable, _) = wait_for_guardian_review( + let result = wait_for_guardian_review( &review_session, "current-turn", tokio::time::Instant::now(), @@ -1146,7 +1173,7 @@ async fn interrupt_and_drain_turn_ignores_prior_turn_completion() { &mut GuardianReviewAnalyticsResult::without_session(), ) .await; - assert!(reusable); + assert_eq!(result.disposition, SessionDisposition::Reusable); assert!(review_session.io.rx_event.try_recv().is_err()); } diff --git a/codex-rs/ext/guardian-reviewer/src/execution.rs b/codex-rs/ext/guardian-reviewer/src/execution.rs index 1c2bbec8d2..36fcdadd2e 100644 --- a/codex-rs/ext/guardian-reviewer/src/execution.rs +++ b/codex-rs/ext/guardian-reviewer/src/execution.rs @@ -15,6 +15,7 @@ use tokio::time::Instant; use tokio_util::sync::CancellationToken; use crate::GuardianReviewSessionOutcome; +use crate::SessionDisposition; const GUARDIAN_INTERRUPT_DRAIN_TIMEOUT: Duration = Duration::from_secs(5); @@ -54,13 +55,21 @@ pub async fn start_review_turn( } } +/// A submitted turn's terminal outcome and the state left behind for reuse and accounting. +pub struct ReviewTurnResult { + pub outcome: GuardianReviewSessionOutcome, + pub disposition: SessionDisposition, + /// A matching TurnComplete event supplied authoritative usage for this turn. + pub turn_completed: bool, +} + pub async fn wait_for_guardian_review( runtime: &impl ReviewerRuntime, expected_turn_id: &str, deadline: tokio::time::Instant, external_cancel: Option<&CancellationToken>, analytics_result: &mut GuardianReviewAnalyticsResult, -) -> (GuardianReviewSessionOutcome, bool, bool) { +) -> ReviewTurnResult { let timeout = tokio::time::sleep_until(deadline); tokio::pin!(timeout); let mut last_error: Option = None; @@ -68,13 +77,16 @@ pub async fn wait_for_guardian_review( loop { tokio::select! { _ = &mut timeout => { - let keep_review_session = interrupt_and_drain_turn( - runtime, - expected_turn_id, - ) - .await - .is_ok(); - return (GuardianReviewSessionOutcome::TimedOut, keep_review_session, false); + let disposition = if interrupt_and_drain_turn(runtime, expected_turn_id).await.is_ok() { + SessionDisposition::Reusable + } else { + SessionDisposition::Discard + }; + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::TimedOut, + disposition, + turn_completed: false, + }; } _ = async { if let Some(cancel_token) = external_cancel { @@ -83,13 +95,16 @@ pub async fn wait_for_guardian_review( std::future::pending::<()>().await; } } => { - let keep_review_session = interrupt_and_drain_turn( - runtime, - expected_turn_id, - ) - .await - .is_ok(); - return (GuardianReviewSessionOutcome::Aborted, keep_review_session, false); + let disposition = if interrupt_and_drain_turn(runtime, expected_turn_id).await.is_ok() { + SessionDisposition::Reusable + } else { + SessionDisposition::Discard + }; + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::Aborted, + disposition, + turn_completed: false, + }; } event = runtime.next_event() => { match event { @@ -105,36 +120,40 @@ pub async fn wait_for_guardian_review( if turn_complete.last_agent_message.is_none() && let Some(error) = last_error { - return ( - GuardianReviewSessionOutcome::SessionFailed { + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::SessionFailed { error: anyhow!(error.message), error_info: error.codex_error_info, retry_at: runtime.retry_at(expected_turn_id), }, - true, - true, - ); + disposition: SessionDisposition::Reusable, + turn_completed: true, + }; } - return ( - GuardianReviewSessionOutcome::Completed(Ok(turn_complete.last_agent_message)), - true, - true, - ); + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::Completed(Ok(turn_complete.last_agent_message)), + disposition: SessionDisposition::Reusable, + turn_completed: true, + }; } EventMsg::Error(error) => { last_error = Some(error); } EventMsg::TurnAborted(_) => { - return (GuardianReviewSessionOutcome::Aborted, true, false); + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::Aborted, + disposition: SessionDisposition::Reusable, + turn_completed: false, + }; } _ => {} }, Err(err) => { - return ( - GuardianReviewSessionOutcome::Completed(Err(err)), - false, - false, - ); + return ReviewTurnResult { + outcome: GuardianReviewSessionOutcome::Completed(Err(err)), + disposition: SessionDisposition::Discard, + turn_completed: false, + }; } } } diff --git a/codex-rs/ext/guardian-reviewer/src/lib.rs b/codex-rs/ext/guardian-reviewer/src/lib.rs index 86347f863e..864d225d1c 100644 --- a/codex-rs/ext/guardian-reviewer/src/lib.rs +++ b/codex-rs/ext/guardian-reviewer/src/lib.rs @@ -41,6 +41,7 @@ pub const MAX_REVIEW_ATTEMPTS: i64 = 3; pub const REVIEW_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(90); pub use deadline::run_before_review_deadline; +pub use pool::ReviewSessionResult; pub use pool::ReviewerPool; pub use pool::ReviewerRequest; pub use pool::ReviewerSession; @@ -53,6 +54,7 @@ pub use completion::ReviewCompletion; pub use completion::complete_review; pub use completion::guardian_timeout_message; +pub use execution::ReviewTurnResult; pub use execution::ReviewerRuntime; pub use execution::start_review_turn; pub use execution::wait_for_guardian_review; diff --git a/codex-rs/ext/guardian-reviewer/src/pool.rs b/codex-rs/ext/guardian-reviewer/src/pool.rs index fc9971d44e..750514ba7e 100644 --- a/codex-rs/ext/guardian-reviewer/src/pool.rs +++ b/codex-rs/ext/guardian-reviewer/src/pool.rs @@ -54,13 +54,14 @@ pub trait ReviewerRequest: Send + Sync { &self, session: &Self::Session, kind: GuardianReviewSessionKind, - ) -> impl Future< - Output = ( - GuardianReviewSessionOutcome, - SessionDisposition, - GuardianReviewAnalyticsResult, - ), - > + Send; + ) -> impl Future + Send; +} + +/// Result of one review, including whether its session can serve the next request. +pub struct ReviewSessionResult { + pub outcome: GuardianReviewSessionOutcome, + pub disposition: SessionDisposition, + pub analytics: GuardianReviewAnalyticsResult, } /// Whether the host drained the session sufficiently for another review to use it. @@ -259,7 +260,11 @@ impl ReviewerPool { }; // Dropping a review before it drains its turn must not leave a reusable agent. let review_lifetime = trunk.cancellation.clone().drop_guard(); - let (outcome, disposition, analytics) = request.run(&trunk.session, kind).await; + let ReviewSessionResult { + outcome, + disposition, + analytics, + } = request.run(&trunk.session, kind).await; if disposition == SessionDisposition::Reusable && matches!(outcome, GuardianReviewSessionOutcome::Completed(_)) { @@ -316,7 +321,9 @@ impl ReviewerPool { } Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()), }; - let (outcome, _, analytics) = request + let ReviewSessionResult { + outcome, analytics, .. + } = request .run(&session, GuardianReviewSessionKind::EphemeralForked) .await; (outcome, analytics)