diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 02b833dd7..b6d2da017 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -916,8 +916,9 @@ impl DataUsageCache { // cursor/coverage receipt. Caches without that proof still take the // normal rebuild path; this keeps an old, incomplete writer from // authorizing a new leader to skip namespace coverage. + // Cycle deadlines advance both counters, so a validated frontier from + // an earlier cycle must remain eligible for adoption. let cross_epoch_checkpoint = self.info.leader_epoch < leader_epoch - && self.info.next_cycle == next_cycle && self.info.scan_identity == Some(identity) && (self.validated_scan_frontier().is_some() || self.validated_raw_enumeration_cursor().is_some() diff --git a/crates/scanner/src/data_usage_define/tests.rs b/crates/scanner/src/data_usage_define/tests.rs index 12bbaabba..047d92105 100644 --- a/crates/scanner/src/data_usage_define/tests.rs +++ b/crates/scanner/src/data_usage_define/tests.rs @@ -1371,14 +1371,18 @@ fn prepare_bucket_checkpoint_preserves_only_valid_raw_enumeration_cursor() { cache.prepare_bucket_checkpoint("bucket", 1, 1, source, TEST_PLAN_DIGEST, identity), DataUsageCachePrepareOutcome::Reused ); + assert_eq!( + cache.prepare_bucket_checkpoint("bucket", 2, 2, source, TEST_PLAN_DIGEST, identity), + DataUsageCachePrepareOutcome::Reused + ); assert_eq!(cache.info.scan_raw_enumeration_cursor, Some(cursor)); let invalid = DataUsageRawEnumerationCursor::new("other/raw".to_string(), Some("entry-001".to_string()), 1, [9; 32]); let mut cache = cache_with_raw_cursor(invalid); cache.info.scan_identity = Some(identity); assert_eq!( - cache.prepare_bucket_checkpoint("bucket", 1, 1, source, TEST_PLAN_DIGEST, identity), - DataUsageCachePrepareOutcome::Reused + cache.prepare_bucket_checkpoint("bucket", 2, 2, source, TEST_PLAN_DIGEST, identity), + DataUsageCachePrepareOutcome::Reset ); assert!(cache.info.scan_raw_enumeration_cursor.is_none()); assert!(cache.info.scan_progress.is_some()); @@ -1410,13 +1414,17 @@ fn prepare_bucket_checkpoint_preserves_only_valid_raw_page_index() { cache.prepare_bucket_checkpoint("bucket", 1, 1, source, TEST_PLAN_DIGEST, identity), DataUsageCachePrepareOutcome::Reused ); + assert_eq!( + cache.prepare_bucket_checkpoint("bucket", 2, 2, source, TEST_PLAN_DIGEST, identity), + DataUsageCachePrepareOutcome::Reused + ); assert_eq!(cache.info.scan_raw_enumeration_page_index, Some(page_index)); let invalid = raw_page_index_fixture("other/raw", &["entry-001"], false); cache.info.scan_raw_enumeration_page_index = Some(invalid); assert_eq!( - cache.prepare_bucket_checkpoint("bucket", 1, 1, source, TEST_PLAN_DIGEST, identity), - DataUsageCachePrepareOutcome::Reused + cache.prepare_bucket_checkpoint("bucket", 3, 3, source, TEST_PLAN_DIGEST, identity), + DataUsageCachePrepareOutcome::Reset ); assert!(cache.info.scan_raw_enumeration_page_index.is_none()); assert!(cache.info.scan_progress.is_some()); @@ -2078,7 +2086,7 @@ fn prepare_bucket_checkpoint_migrates_legacy_epoch_bound_receipt() { assert_eq!(cache.validated_scan_frontier(), Some("bucket/a")); assert_eq!( - cache.prepare_bucket_checkpoint("bucket", 8, 2, source, TEST_PLAN_DIGEST, identity), + cache.prepare_bucket_checkpoint("bucket", 9, 2, source, TEST_PLAN_DIGEST, identity), DataUsageCachePrepareOutcome::Reused ); assert_eq!(cache.info.leader_epoch, 2); diff --git a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs index 246b74f87..0fa3eef30 100644 --- a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs +++ b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs @@ -444,7 +444,7 @@ fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() { ] { let mut next = cache.clone(); assert_eq!( - next.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, next_identity), + next.prepare_bucket_checkpoint("bucket", 12, 8, SOURCE, PLAN, next_identity), crate::DataUsageCachePrepareOutcome::Reset ); assert!(next.cache.is_empty()); @@ -453,32 +453,37 @@ fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() { } let mut source_mismatch = cache.clone(); assert_eq!( - source_mismatch.prepare_bucket_checkpoint("bucket", 11, 7, crate::DataUsageCacheSource::new(1, 0), PLAN, identity,), + source_mismatch.prepare_bucket_checkpoint("bucket", 12, 8, crate::DataUsageCacheSource::new(1, 0), PLAN, identity,), crate::DataUsageCachePrepareOutcome::Reset ); assert!(source_mismatch.cache.is_empty()); - let mut handed_off = cache.clone(); - assert_eq!( - handed_off.prepare_bucket_checkpoint("bucket", 11, 8, SOURCE, PLAN, identity), - crate::DataUsageCachePrepareOutcome::Reused - ); - assert_eq!(handed_off.info.leader_epoch, 8); - assert_eq!(handed_off.validated_scan_frontier(), Some("bucket/static")); - assert_eq!(retained(&handed_off), 3); + for cycle in [11, 12] { + let mut handed_off = cache.clone(); + assert_eq!( + handed_off.prepare_bucket_checkpoint("bucket", cycle, 8, SOURCE, PLAN, identity), + crate::DataUsageCachePrepareOutcome::Reused + ); + assert_eq!(handed_off.info.next_cycle, cycle); + assert_eq!(handed_off.info.leader_epoch, 8); + assert_eq!(handed_off.validated_scan_frontier(), Some("bucket/static")); + assert_eq!(retained(&handed_off), 3); + } let mut unverified = cache.clone(); unverified.info.scan_resume_after = None; unverified.info.scan_checkpoint = None; unverified.info.scan_coverage_receipt = None; assert_eq!( - unverified.prepare_bucket_checkpoint("bucket", 11, 8, SOURCE, PLAN, identity), + unverified.prepare_bucket_checkpoint("bucket", 12, 8, SOURCE, PLAN, identity), crate::DataUsageCachePrepareOutcome::Reset, "a cross-epoch cache without a durable frontier proof must rebuild" ); assert!(unverified.cache.is_empty()); for (cycle, epoch, expected) in [ (10, 7, crate::DataUsageCachePrepareOutcome::RejectedNewerCycle), + (10, 8, crate::DataUsageCachePrepareOutcome::RejectedNewerCycle), (11, 6, crate::DataUsageCachePrepareOutcome::RejectedNewerLeader), + (12, 6, crate::DataUsageCachePrepareOutcome::RejectedNewerLeader), ] { let mut next = cache.clone(); assert_eq!(next.prepare_bucket_checkpoint("bucket", cycle, epoch, SOURCE, PLAN, identity), expected); @@ -864,7 +869,7 @@ async fn check_complete_sampling_resumption(resume_mode: HealScanMode) { #[tokio::test] #[serial] -async fn checkpoint_fixture_save_reload_resume() { +async fn checkpoint_fixture_save_reload_resume_across_cycles_and_leaders() { run_checkpoint_fixture(false).await; } @@ -1013,12 +1018,13 @@ async fn run_checkpoint_fixture(change_digest: bool) { assert_eq!(retained(&store.strict_load().await), previous); } let plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, change_digest.then_some(u64::from(round))); + // A timed-out cycle advances both the cycle and the leader epoch. crate::scanner_io::current_cache_root_or_prepare_with_generation( &mut cache, "bucket", SOURCE, - 11, - 7, + 11 + u64::from(round), + 7 + u64::from(round), plan, crate::scanner_io::DataUsageCacheReuseOptions { require_source: true, @@ -1145,7 +1151,7 @@ async fn run_checkpoint_fixture(change_digest: bool) { write_checkpoint_object(&root, "hot/later", &[(None, 1)]).await; let final_plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, Some(3)); let mut saw_mixed_sweep_end = false; - for _ in 0..32 { + for round in 0..32 { let mut cache = DataUsageCache::default(); let revisions = cache .load_with_revisions(store.clone(), CACHE_NAME) @@ -1155,8 +1161,8 @@ async fn run_checkpoint_fixture(change_digest: bool) { &mut cache, "bucket", SOURCE, - 11, - 7, + 14 + round, + 10 + round, final_plan, crate::scanner_io::DataUsageCacheReuseOptions { require_source: true,