fix(scanner): resume checkpoints across cycle and leader changes (#8252)

This commit is contained in:
Jason Kossis
2026-09-29 20:18:20 -04:00
committed by GitHub
parent 50aa482136
commit 56468fd553
3 changed files with 38 additions and 23 deletions
+2 -1
View File
@@ -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()
+13 -5
View File
@@ -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);
@@ -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,