From 8a7111f98d2f39de2b4ba1a3836068fc2a317601 Mon Sep 17 00:00:00 2001 From: Ray Morgan Date: Thu, 11 Jun 2026 23:22:00 -0700 Subject: [PATCH] Simplify Guardian trunk ownership --- codex-rs/core/src/guardian/review_session.rs | 71 ++++---------------- 1 file changed, 13 insertions(+), 58 deletions(-) diff --git a/codex-rs/core/src/guardian/review_session.rs b/codex-rs/core/src/guardian/review_session.rs index a5149b0109..345424db2e 100644 --- a/codex-rs/core/src/guardian/review_session.rs +++ b/codex-rs/core/src/guardian/review_session.rs @@ -30,7 +30,6 @@ use codex_protocol::protocol::TokenUsage; use serde_json::Value; use tokio::sync::Mutex; use tokio::sync::Semaphore; -use tokio::sync::SemaphorePermit; use tokio_util::sync::CancellationToken; use tracing::warn; @@ -104,26 +103,9 @@ pub(crate) struct GuardianReviewSessionManager { #[derive(Default)] struct GuardianReviewSessionState { trunk: Option>, - trunk_review_generation: u64, ephemeral_reviews: Vec>, } -impl GuardianReviewSessionState { - fn try_claim_trunk<'a>( - &mut self, - trunk: &'a Arc, - review_generation: u64, - ) -> Option<(SemaphorePermit<'a>, u64)> { - let current_trunk = self.trunk.as_ref()?; - if self.trunk_review_generation != review_generation || !Arc::ptr_eq(current_trunk, trunk) { - return None; - } - let review_guard = trunk.review_lock.try_acquire().ok()?; - self.trunk_review_generation = self.trunk_review_generation.wrapping_add(1); - Some((review_guard, self.trunk_review_generation)) - } -} - struct GuardianReviewSession { codex: Codex, cancel_token: CancellationToken, @@ -418,10 +400,7 @@ impl GuardianReviewSessionManager { spawned_trunk = true; } - state - .trunk - .as_ref() - .map(|trunk| (Arc::clone(trunk), state.trunk_review_generation)) + state.trunk.as_ref().cloned() } Err(outcome) => { return (outcome, GuardianReviewAnalyticsResult::without_session()); @@ -432,7 +411,7 @@ impl GuardianReviewSessionManager { review_session.shutdown_in_background(); } - let Some((trunk, trunk_review_generation)) = trunk_candidate else { + let Some(trunk) = trunk_candidate else { return ( GuardianReviewSessionOutcome::Completed(Err(anyhow!( "guardian review session was not available after spawn" @@ -470,9 +449,7 @@ impl GuardianReviewSessionManager { let eager_compaction_guard = match eager_compaction_guard { Some(guard) if *guard == GuardianEagerCompactionOutcome::Reusable => guard, guard => { - let review_session = self - .remove_trunk_if_current(&trunk, trunk_review_generation) - .await; + let review_session = self.remove_trunk_if_current(&trunk).await; drop(guard); if let Some(review_session) = review_session { review_session.shutdown_in_background(); @@ -487,14 +464,17 @@ impl GuardianReviewSessionManager { } }; // Hold the maintenance latch until review ownership is decided. - let mut state = self.state.lock().await; + let state = self.state.lock().await; let trunk_pointer_is_current = state .trunk .as_ref() .is_some_and(|current| Arc::ptr_eq(current, &trunk)); - let Some((trunk_guard, claimed_review_generation)) = - state.try_claim_trunk(&trunk, trunk_review_generation) - else { + let trunk_guard = if trunk_pointer_is_current { + trunk.review_lock.try_acquire().ok() + } else { + None + }; + let Some(trunk_guard) = trunk_guard else { drop(state); drop(eager_compaction_guard); let fork_snapshot = if trunk_pointer_is_current { @@ -536,9 +516,9 @@ impl GuardianReviewSessionManager { drop(trunk_guard); (outcome, analytics_result) } else { - let review_session = self - .remove_trunk_if_current(&trunk, claimed_review_generation) - .await; + // Remove the trunk before releasing review ownership so another + // review cannot claim it between failure and removal. + let review_session = self.remove_trunk_if_current(&trunk).await; drop(trunk_guard); if let Some(review_session) = review_session { review_session.shutdown_in_background(); @@ -617,14 +597,12 @@ impl GuardianReviewSessionManager { async fn remove_trunk_if_current( &self, trunk: &Arc, - review_generation: u64, ) -> Option> { let mut state = self.state.lock().await; if state .trunk .as_ref() .is_some_and(|current| Arc::ptr_eq(current, trunk)) - && state.trunk_review_generation == review_generation { state.trunk.take() } else { @@ -1231,29 +1209,6 @@ mod tests { ) } - #[tokio::test] - async fn stale_generation_cannot_remove_idle_trunk() { - let (review_session, _, _) = test_review_session().await; - let trunk = Arc::new(review_session); - let manager = GuardianReviewSessionManager::default(); - let review_guard = { - let mut state = manager.state.lock().await; - state.trunk = Some(Arc::clone(&trunk)); - state - .try_claim_trunk(&trunk, /*review_generation*/ 0) - .expect("current trunk should be claimable") - .0 - }; - assert!( - manager - .remove_trunk_if_current(&trunk, /*review_generation*/ 0) - .await - .is_none() - ); - assert!(manager.state.lock().await.trunk.is_some()); - drop(review_guard); - } - fn turn_complete_event( turn_id: &str, last_agent_message: Option<&str>,