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
This commit is contained in:
felixxia-oai
2026-09-16 16:44:59 +00:00
committed by copyberry
parent 49305d74b4
commit 66dfadbe66
6 changed files with 192 additions and 113 deletions

View File

@@ -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 = &params.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::<super::input_budget::PendingReviewContext>();
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::<super::input_budget::PendingReviewContext>();
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::<super::request_budget::ExhaustedReviewBudget>();
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) {

View File

@@ -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<Session>,
@@ -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
}
}

View File

@@ -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());
}

View File

@@ -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<ErrorEvent> = 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,
};
}
}
}

View File

@@ -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;

View File

@@ -54,13 +54,14 @@ pub trait ReviewerRequest: Send + Sync {
&self,
session: &Self::Session,
kind: GuardianReviewSessionKind,
) -> impl Future<
Output = (
GuardianReviewSessionOutcome,
SessionDisposition,
GuardianReviewAnalyticsResult,
),
> + Send;
) -> impl Future<Output = ReviewSessionResult> + 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<S: ReviewerSession> ReviewerPool<S> {
};
// 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<S: ReviewerSession> ReviewerPool<S> {
}
Err(outcome) => return (outcome, GuardianReviewAnalyticsResult::without_session()),
};
let (outcome, _, analytics) = request
let ReviewSessionResult {
outcome, analytics, ..
} = request
.run(&session, GuardianReviewSessionKind::EphemeralForked)
.await;
(outcome, analytics)