diff --git a/codex-rs/core/src/guardian/tests.rs b/codex-rs/core/src/guardian/tests.rs index c98c78348e..eb8721a0d2 100644 --- a/codex-rs/core/src/guardian/tests.rs +++ b/codex-rs/core/src/guardian/tests.rs @@ -2233,40 +2233,48 @@ async fn run_eager_compaction_review( Ok(analytics_result) } -#[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn guardian_recovers_from_eager_compaction_failure_and_timeout() -> anyhow::Result<()> { - const EAGER_COMPACTION_WAIT_TIMEOUT: Duration = Duration::from_millis(/*millis*/ 50); +#[derive(Clone, Copy)] +enum EagerCompactionRecoveryCase { + Failure, + Timeout, +} + +async fn run_eager_compaction_recovery_case( + recovery_case: EagerCompactionRecoveryCase, +) -> anyhow::Result<()> { + const FAILURE_WAIT_TIMEOUT: Duration = Duration::from_secs(/*secs*/ 3_600); + const TIMEOUT_WAIT_TIMEOUT: Duration = Duration::from_millis(/*millis*/ 50); const FIRST_REVIEW_TOTAL_TOKENS: i64 = 500_000; - let (failed_compaction_tx, failed_compaction_rx) = tokio::sync::oneshot::channel(); let (compaction_tx, compaction_rx) = tokio::sync::oneshot::channel(); + let compaction_body = match recovery_case { + EagerCompactionRecoveryCase::Failure => sse(vec![ + ev_response_created("resp-eager-compact-failed"), + ev_assistant_message("msg-eager-compact-partial", "partial summary"), + ]), + EagerCompactionRecoveryCase::Timeout => { + sse(vec![ev_response_created("resp-eager-compact-timeout")]) + } + }; let (server, _) = start_streaming_sse_server(vec![ guardian_review_sse_with_body_tokens( "1", FIRST_REVIEW_TOTAL_TOKENS - EAGER_COMPACTION_GUARDIAN_PREFILL_TOKENS, ), vec![StreamingSseChunk { - gate: Some(failed_compaction_rx), - body: sse(vec![ - ev_response_created("resp-eager-compact-failed"), - ev_assistant_message("msg-eager-compact-partial", "partial summary"), - ]), + gate: Some(compaction_rx), + body: compaction_body, }], guardian_review_sse("2", ev_completed), - guardian_review_sse_with_body_tokens( - "3", - FIRST_REVIEW_TOTAL_TOKENS - EAGER_COMPACTION_GUARDIAN_PREFILL_TOKENS, - ), - vec![StreamingSseChunk { - gate: Some(compaction_rx), - body: sse(vec![ev_response_created("resp-eager-compact-timeout")]), - }], - guardian_review_sse("4", ev_completed), ]) .await; let (mut session, mut turn) = guardian_test_session_and_turn_with_base_url(server.uri()).await; + let eager_compaction_wait_timeout = match recovery_case { + EagerCompactionRecoveryCase::Failure => FAILURE_WAIT_TIMEOUT, + EagerCompactionRecoveryCase::Timeout => TIMEOUT_WAIT_TIMEOUT, + }; Arc::get_mut(&mut session) .expect("guardian parent session should be uniquely owned") - .guardian_review_session = GuardianReviewSessionManager::new(EAGER_COMPACTION_WAIT_TIMEOUT); + .guardian_review_session = GuardianReviewSessionManager::new(eager_compaction_wait_timeout); configure_guardian_eager_compaction_test(&mut turn); seed_guardian_parent_history(&session, &turn).await; @@ -2282,46 +2290,39 @@ async fn guardian_recovers_from_eager_compaction_failure_and_timeout() -> anyhow server.wait_for_request_count(/*count*/ 2), ) .await?; - failed_compaction_tx - .send(()) - .expect("failed compaction response gate should still be open"); - let second_metadata = run_eager_compaction_review( + let mut compaction_tx = Some(compaction_tx); + if matches!(recovery_case, EagerCompactionRecoveryCase::Failure) { + let compaction_tx = compaction_tx.take().expect("compaction response gate"); + compaction_tx + .send(()) + .expect("failed compaction response gate should still be open"); + } + let metadata = run_eager_compaction_review( &session, &turn, - Some("Retry after eager compaction failure."), + Some("Retry after eager compaction maintenance."), ) .await?; + if let Some(compaction_tx) = compaction_tx { + let _ = compaction_tx.send(()); + } assert_eq!( tokio::fs::read(&discarded_rollout_path).await?, committed_rollout ); - run_eager_compaction_review(&session, &turn, Some("Start a replacement guardian trunk.")) - .await?; - tokio::time::timeout( - EAGER_COMPACTION_TEST_TIMEOUT, - server.wait_for_request_count(/*count*/ 5), - ) - .await?; - let fourth_metadata = run_eager_compaction_review( - &session, - &turn, - Some("Retry after eager compaction timeout."), - ) - .await?; assert!(matches!( - ( - second_metadata.guardian_session_kind, - fourth_metadata.guardian_session_kind - ), - ( - Some(GuardianReviewSessionKind::EphemeralForked), - Some(GuardianReviewSessionKind::EphemeralForked) - ) + metadata.guardian_session_kind, + Some(GuardianReviewSessionKind::EphemeralForked) )); - let _ = compaction_tx.send(()); Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn guardian_recovers_from_eager_compaction_failure_and_timeout() -> anyhow::Result<()> { + run_eager_compaction_recovery_case(EagerCompactionRecoveryCase::Failure).await?; + run_eager_compaction_recovery_case(EagerCompactionRecoveryCase::Timeout).await +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn guardian_reused_trunk_ignores_stale_prior_turn_completion() -> anyhow::Result<()> { skip_if_no_network!(Ok(()));