mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 12:35:54 +00:00
fix(scanner): bind resumable scans and cache publication coverage (#7210)
* chore(deps): refresh scanner heal batch dependency baseline Regenerate compatible lockfile selections before the next implementation batch. Cargo upgrade leaves direct requirements unchanged. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(ecstore): remove duplicate local rename implementation Keep the canonical commit module after concurrent storage changes merged. The control-write and rollback changes are already present there. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * chore(deps): refresh profiling dependencies for the next batch Update hotpath and its macro crate to the compatible patch release before the next dependency-ready implementation tasks. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(deps): preserve supported hotpath focus expressions Keep the profiler runtime before its regex-lite compatibility regression. Track the opt-in validation required to remove this constraint in backlog. Refs rustfs/backlog#2302. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): require complete publication coverage Refs rustfs/backlog#2261 and rustfs/backlog#2240. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): retain scoped partial coverage across dirty plans Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): keep stable snapshot rescan behavior Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): verify coverage receipts and scan strength Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * test(scanner): use valid modification times in checkpoint fixtures Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): keep maintenance cycles outside dirty bucket scopes Force complete bucket scope for deep scans and scheduled maintenance while preserving the existing planner for verified ordinary dirty work. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): refresh scope safety independently of idle backoff Inspect maintenance on multi-disk startup and refresh changed or failed evidence even when explicit bitrot configuration disables idle backoff. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): bind bucket cache reuse to scan work requirements Carry stable scan mode and full-maintenance requirements in the existing opaque bucket digest before local and remote cache admission. Different requirements cannot replay a same-cycle Normal cache after root delivery failure; matching requirements remain reusable for the same intent. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): fence set snapshot reuse with the scan work proof Prevent same-cycle set publication from replacing freshly scanned maintenance results with an older Normal aggregate. Recognize uniform completed maintenance baselines when planning later ordinary dirty-bucket work. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * test(scanner): reproduce same-cycle dirty aggregate replay Cover a Normal-to-Normal retry with a new dirty bucket generation after bucket persistence and root delivery failure. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): fence same-cycle caches with full activity coverage Keep structural baseline identity separate from the full activity coverage required by bucket admission and set publication. Require complete set coverage proofs while retaining revision CAS and epoch regression checks. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * test(scanner): supply explicit coverage in publication fixtures Keep the confirmed-empty namespace fixture authoritative under the required coverage contract and qualify the bucket cache metadata test type. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * test(scanner): verify joint checkpoint coverage metadata Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> * fix(scanner): satisfy cache prefix sort lint Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com> Co-authored-by: zhi22915 <qiuzgang@gmail.com> Co-authored-by: Zhengchao An <anzhengchao@gmail.com>
This commit is contained in:
@@ -346,6 +346,9 @@ async fn build_test_scanner() -> (FolderScanner, std::path::PathBuf) {
|
||||
refresh_failed: false,
|
||||
},
|
||||
pending_heals_changed: false,
|
||||
coverage_frontier: None,
|
||||
resume_frontier: None,
|
||||
coverage_gap: false,
|
||||
pending_size_reconciliation_keys: HashSet::new(),
|
||||
pending_size_reconciliation_scopes: HashSet::new(),
|
||||
pending_size_reconciliation_truncated: false,
|
||||
|
||||
@@ -234,6 +234,531 @@ fn checkpoint_fixture_compaction_preserves_aggregate_not_child_enumeration() {
|
||||
);
|
||||
}
|
||||
|
||||
fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) {
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
version: 1,
|
||||
bucket_incarnation: Uuid::from_u128(1),
|
||||
set_layout: DataUsageScanPlanDigest([41; 32]),
|
||||
publication_epoch: 0,
|
||||
tier_registry_generation: 7,
|
||||
scan_mode: HealScanMode::Normal,
|
||||
};
|
||||
let mut cache = DataUsageCache::default();
|
||||
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
|
||||
cache.replace("bucket", "", DataUsageEntry::default());
|
||||
cache.replace(
|
||||
"bucket/static",
|
||||
"bucket",
|
||||
DataUsageEntry {
|
||||
objects: 3,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
cache.info.scan_resume_after = Some("bucket/static".into());
|
||||
cache.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(
|
||||
"bucket/static".into(),
|
||||
DataUsageScanCheckpointReason::Objects,
|
||||
));
|
||||
cache
|
||||
.seal_scan_frontier(Some("bucket/static"))
|
||||
.expect("completed fixture prefix receipt");
|
||||
(cache, identity)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_roundtrip_retains_verified_scope_but_old_reader_rebuilds() {
|
||||
let (cache, identity) = bound_checkpoint();
|
||||
let mut cache = decode_fixture(&cache.marshal_msg().expect("encode bound progress")).expect("read bound progress");
|
||||
let next_plan = DataUsageScanPlanDigest([42; 32]);
|
||||
assert_eq!(
|
||||
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, next_plan, identity),
|
||||
crate::DataUsageCachePrepareOutcome::Reused
|
||||
);
|
||||
assert_eq!(retained(&cache), 3);
|
||||
assert_eq!(cache.info.scan_identity, Some(identity));
|
||||
assert_eq!(
|
||||
cache.info.scan_progress,
|
||||
Some(crate::DataUsageScanProgress {
|
||||
started_plan: PLAN,
|
||||
requested_plan: next_plan
|
||||
})
|
||||
);
|
||||
assert!(cache.info.scan_plan_digest.is_none());
|
||||
let mut old_wire = serde_json::to_value(&cache).expect("map-encoded compatibility fixture");
|
||||
let old_info = old_wire["info"].as_object_mut().expect("cache info is a map");
|
||||
old_info.remove("scan_identity");
|
||||
old_info.remove("scan_progress");
|
||||
old_info.remove("scan_coverage_receipt");
|
||||
let mut old_view: DataUsageCache = serde_json::from_value(old_wire).expect("old writer drops unknown metadata");
|
||||
assert_eq!(
|
||||
old_view.prepare_for_scan("bucket", 11, 7, SOURCE, next_plan, true),
|
||||
crate::DataUsageCachePrepareOutcome::Reset
|
||||
);
|
||||
assert!(old_view.cache.is_empty());
|
||||
assert!(!old_view.info.snapshot_complete);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_optional_coverage_metadata_roundtrips_together() {
|
||||
let (mut cache, _) = bound_checkpoint();
|
||||
cache.info.scan_coverage_digest = Some(DataUsageScanPlanDigest([77; 32]));
|
||||
let encoded = cache.marshal_msg().expect("encode every optional coverage field");
|
||||
let decoded = DataUsageCache::unmarshal(&encoded).expect("decode combined W03/W04 metadata map");
|
||||
assert_eq!(decoded.info.scan_identity, cache.info.scan_identity);
|
||||
assert_eq!(decoded.info.scan_progress, cache.info.scan_progress);
|
||||
assert_eq!(decoded.info.scan_coverage_receipt, cache.info.scan_coverage_receipt);
|
||||
assert_eq!(decoded.info.scan_coverage_digest, cache.info.scan_coverage_digest);
|
||||
assert_eq!(decoded.validated_scan_frontier(), Some("bucket/static"));
|
||||
assert!(!decoded.info.snapshot_complete);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_unchanged_complete_plan_keeps_existing_rescan_policy() {
|
||||
let (mut cache, identity) = bound_checkpoint();
|
||||
cache.info.scan_progress = None;
|
||||
cache.info.scan_plan_digest = Some(PLAN);
|
||||
cache.info.scan_resume_after = None;
|
||||
cache.info.scan_checkpoint = None;
|
||||
cache.info.scan_coverage_receipt = None;
|
||||
cache.info.snapshot_complete = true;
|
||||
assert_eq!(
|
||||
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, identity),
|
||||
crate::DataUsageCachePrepareOutcome::Reused
|
||||
);
|
||||
assert!(
|
||||
cache.info.scan_progress.is_none(),
|
||||
"unchanged complete coverage needs no forced verification sweep"
|
||||
);
|
||||
assert_eq!(cache.info.scan_plan_digest, Some(PLAN));
|
||||
assert_eq!(retained(&cache), 3);
|
||||
let next = DataUsageScanPlanDigest([44; 32]);
|
||||
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, next, identity);
|
||||
assert_eq!(cache.info.scan_progress.expect("changed plan must be verified").started_plan, next);
|
||||
assert!(cache.info.scan_plan_digest.is_none());
|
||||
assert!(!cache.info.snapshot_complete);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() {
|
||||
let (cache, identity) = bound_checkpoint();
|
||||
for next_identity in [
|
||||
crate::DataUsageScanIdentity {
|
||||
bucket_incarnation: Uuid::from_u128(2),
|
||||
..identity
|
||||
},
|
||||
crate::DataUsageScanIdentity {
|
||||
set_layout: DataUsageScanPlanDigest([9; 32]),
|
||||
..identity
|
||||
},
|
||||
crate::DataUsageScanIdentity {
|
||||
publication_epoch: 1,
|
||||
..identity
|
||||
},
|
||||
crate::DataUsageScanIdentity {
|
||||
tier_registry_generation: 8,
|
||||
..identity
|
||||
},
|
||||
crate::DataUsageScanIdentity {
|
||||
scan_mode: HealScanMode::Deep,
|
||||
..identity
|
||||
},
|
||||
] {
|
||||
let mut next = cache.clone();
|
||||
assert_eq!(
|
||||
next.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, next_identity),
|
||||
crate::DataUsageCachePrepareOutcome::Reset
|
||||
);
|
||||
assert!(next.cache.is_empty());
|
||||
assert!(next.info.scan_checkpoint.is_none());
|
||||
assert!(!next.info.snapshot_complete);
|
||||
}
|
||||
for (source, epoch) in [(crate::DataUsageCacheSource::new(1, 0), 7), (SOURCE, 8)] {
|
||||
let mut next = cache.clone();
|
||||
assert_eq!(
|
||||
next.prepare_bucket_checkpoint("bucket", 11, epoch, source, PLAN, identity),
|
||||
crate::DataUsageCachePrepareOutcome::Reset
|
||||
);
|
||||
assert!(next.cache.is_empty());
|
||||
}
|
||||
for (cycle, epoch, expected) in [
|
||||
(10, 7, crate::DataUsageCachePrepareOutcome::RejectedNewerCycle),
|
||||
(11, 6, crate::DataUsageCachePrepareOutcome::RejectedNewerLeader),
|
||||
] {
|
||||
let mut next = cache.clone();
|
||||
assert_eq!(next.prepare_bucket_checkpoint("bucket", cycle, epoch, SOURCE, PLAN, identity), expected);
|
||||
assert_eq!(
|
||||
serde_json::to_value(&next).expect("current cache"),
|
||||
serde_json::to_value(&cache).expect("saved cache")
|
||||
);
|
||||
}
|
||||
for invalid in [
|
||||
crate::DataUsageScanIdentity { version: 2, ..identity },
|
||||
crate::DataUsageScanIdentity {
|
||||
bucket_incarnation: Uuid::nil(),
|
||||
..identity
|
||||
},
|
||||
] {
|
||||
let mut next = cache.clone();
|
||||
crate::scanner_io::current_cache_root_or_prepare_with_generation(
|
||||
&mut next,
|
||||
"bucket",
|
||||
SOURCE,
|
||||
11,
|
||||
7,
|
||||
PLAN,
|
||||
crate::scanner_io::DataUsageCacheReuseOptions {
|
||||
checkpoint_identity: Some(invalid),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
assert!(next.cache.is_empty(), "unsupported identity must not retain coverage");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_corrupt_cursor_restarts_validation_without_claiming_completion() {
|
||||
let (cache, identity) = bound_checkpoint();
|
||||
for resume in ["other/static", "bucket/missing"] {
|
||||
let mut next = cache.clone();
|
||||
next.info.scan_resume_after = Some(resume.into());
|
||||
next.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(resume.into(), DataUsageScanCheckpointReason::Objects));
|
||||
let plan = DataUsageScanPlanDigest([43; 32]);
|
||||
next.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, plan, identity);
|
||||
assert_eq!(retained(&next), 3, "observations may survive an invalid cursor");
|
||||
assert!(next.info.scan_resume_after.is_none());
|
||||
assert!(next.info.scan_checkpoint.is_none());
|
||||
assert_eq!(next.info.scan_progress.expect("new verification sweep").started_plan, plan);
|
||||
assert!(!next.info.snapshot_complete);
|
||||
assert!(next.info.scan_plan_digest.is_none());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_receipt_binds_covered_prefix_not_unvisited_suffix() {
|
||||
let (mut cache, _) = bound_checkpoint();
|
||||
cache.replace(
|
||||
"bucket/z-unvisited",
|
||||
"bucket",
|
||||
DataUsageEntry {
|
||||
objects: 99,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
assert_eq!(cache.validated_scan_frontier(), Some("bucket/static"));
|
||||
let saved = decode_fixture(&cache.marshal_msg().expect("persist receipt")).expect("load receipt");
|
||||
assert_eq!(saved.validated_scan_frontier(), Some("bucket/static"));
|
||||
cache.cache.get_mut("bucket/static").expect("covered prefix").objects = 100;
|
||||
assert!(
|
||||
cache.validated_scan_frontier().is_none(),
|
||||
"altered covered content must invalidate the receipt"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_existing_uncovered_cursor_cannot_skip_to_complete() {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _guard = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for (object, size) in [
|
||||
("a-done/object", 1),
|
||||
("b-pending/first", 7),
|
||||
("b-pending/second", 1),
|
||||
("z-stale/object", 2),
|
||||
] {
|
||||
write_checkpoint_object(&root, object, &[(None, size)]).await;
|
||||
}
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
|
||||
..bound_checkpoint().1
|
||||
};
|
||||
for tamper_receipt_path in [false, true] {
|
||||
let store = FixtureStore::new();
|
||||
let mut cache = DataUsageCache::default();
|
||||
let revisions = cache
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("empty fixture revisions");
|
||||
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
|
||||
cache.replace("bucket", "", DataUsageEntry::default());
|
||||
for (prefix, objects) in [("a-done", 1), ("b-pending", 99), ("z-stale", 99)] {
|
||||
cache.replace(
|
||||
&format!("bucket/{prefix}"),
|
||||
"bucket",
|
||||
DataUsageEntry {
|
||||
objects,
|
||||
size: objects,
|
||||
compacted: true,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
}
|
||||
cache
|
||||
.seal_scan_frontier(Some("bucket/a-done"))
|
||||
.expect("actual completed prefix receipt");
|
||||
cache.info.scan_resume_after = Some("bucket/z-stale".into());
|
||||
cache.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(
|
||||
"bucket/z-stale".into(),
|
||||
DataUsageScanCheckpointReason::Objects,
|
||||
));
|
||||
if tamper_receipt_path {
|
||||
cache.info.scan_coverage_receipt.as_mut().expect("receipt").through = "bucket/z-stale".into();
|
||||
}
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("persist corrupted existing cursor");
|
||||
let mut loaded = store.strict_load().await;
|
||||
let revisions = loaded
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("corrupt cursor CAS revision");
|
||||
assert!(loaded.validated_scan_frontier().is_none());
|
||||
loaded.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
|
||||
assert!(loaded.info.scan_resume_after.is_none());
|
||||
loaded.info.skip_healing = true;
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(2),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let outcome = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget.clone(),
|
||||
vec![scanner.local_disk.clone()],
|
||||
loaded,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("scan must revisit the prefix");
|
||||
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
|
||||
panic!("uncovered suffix must not become complete")
|
||||
};
|
||||
assert_eq!(budget.progress().0, 2);
|
||||
assert_eq!(cache.checked_flatten("bucket/b-pending").expect("revisited prefix").size, 7);
|
||||
assert!(!cache.info.snapshot_complete);
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("persist verified partial");
|
||||
assert!(!store.strict_load().await.info.snapshot_complete);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_failed_child_prevents_receipt_advancing_past_gap() {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _guard = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for object in ["a-good", "b-skipped", "c-later"] {
|
||||
write_checkpoint_object(&root, object, &[(None, 1)]).await;
|
||||
}
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
|
||||
..bound_checkpoint().1
|
||||
};
|
||||
let mut cache = DataUsageCache::default();
|
||||
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
|
||||
cache.info.skip_healing = true;
|
||||
cache.info.failed_objects.insert(
|
||||
root.join("bucket/b-skipped/xl.meta").to_string_lossy().into_owned(),
|
||||
FolderScanner::now_secs(),
|
||||
);
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(2),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let result = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget,
|
||||
vec![scanner.local_disk.clone()],
|
||||
cache,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("scan with a known failed child");
|
||||
let ScannerDiskScanOutcome::Partial(cache) = result else {
|
||||
panic!("skipped failure is not complete")
|
||||
};
|
||||
assert_eq!(cache.validated_scan_frontier(), Some("bucket/a-good"));
|
||||
assert!(!cache.info.failed_objects.is_empty());
|
||||
assert!(!cache.info.snapshot_complete);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_complete_sampling_partial_resumes_with_fixed_budget() {
|
||||
check_complete_sampling_resumption(HealScanMode::Normal).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_normal_partial_reenters_prefix_for_deep_scan() {
|
||||
check_complete_sampling_resumption(HealScanMode::Deep).await;
|
||||
}
|
||||
|
||||
async fn check_complete_sampling_resumption(resume_mode: HealScanMode) {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _guard = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for index in 0..9 {
|
||||
write_checkpoint_object(&root, &format!("prefix/{index:04}"), &[(None, 1)]).await;
|
||||
}
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
|
||||
..bound_checkpoint().1
|
||||
};
|
||||
let mut cache = DataUsageCache::default();
|
||||
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
|
||||
cache.info.skip_healing = true;
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new(&parent, Default::default());
|
||||
// Seed an existing complete baseline; every recovery attempt below is bounded.
|
||||
let baseline = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget,
|
||||
vec![scanner.local_disk.clone()],
|
||||
cache,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("initial complete baseline");
|
||||
let ScannerDiskScanOutcome::Complete(mut cache) = baseline else { panic!("baseline is complete") };
|
||||
let mut deep_current = cache.clone();
|
||||
let deep_identity = crate::DataUsageScanIdentity {
|
||||
scan_mode: HealScanMode::Deep,
|
||||
..identity
|
||||
};
|
||||
let state = crate::scanner_io::current_cache_root_or_prepare_with_generation(
|
||||
&mut deep_current,
|
||||
"bucket",
|
||||
SOURCE,
|
||||
11,
|
||||
7,
|
||||
PLAN,
|
||||
crate::scanner_io::DataUsageCacheReuseOptions {
|
||||
checkpoint_identity: Some(deep_identity),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
assert!(
|
||||
matches!(state, crate::scanner_io::DataUsageCacheScanState::Prepared { .. }),
|
||||
"Normal complete is not Deep Current"
|
||||
);
|
||||
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, identity);
|
||||
assert!(cache.info.scan_progress.is_none(), "complete unchanged baseline uses existing sampling");
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(3),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let outcome = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget,
|
||||
vec![scanner.local_disk.clone()],
|
||||
cache,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("sampling interruption");
|
||||
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
|
||||
panic!("sampling must exhaust the three-object budget")
|
||||
};
|
||||
assert!(cache.info.scan_progress.is_none());
|
||||
assert_eq!(cache.info.scan_plan_digest, Some(PLAN));
|
||||
assert!(cache.info.scan_checkpoint.is_some());
|
||||
let store = FixtureStore::new();
|
||||
let mut loaded = DataUsageCache::default();
|
||||
let revisions = loaded
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("fixture revision");
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("persist sampling partial");
|
||||
if resume_mode == HealScanMode::Deep {
|
||||
write_checkpoint_object(&root, "prefix/0000", &[(None, 7)]).await;
|
||||
}
|
||||
let resumed_identity = crate::DataUsageScanIdentity {
|
||||
scan_mode: resume_mode,
|
||||
..identity
|
||||
};
|
||||
for round in 0..16 {
|
||||
let mut loaded = DataUsageCache::default();
|
||||
let revisions = loaded
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("reload partial each recovery round");
|
||||
loaded.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, resumed_identity);
|
||||
assert!(loaded.info.scan_progress.is_some(), "sampling partial must enter forward validation");
|
||||
loaded.info.skip_healing = true;
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(3),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let result = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(budget.token(), budget, vec![scanner.local_disk.clone()], loaded, None, resume_mode)
|
||||
.await
|
||||
.expect("bounded recovery scan");
|
||||
let (cache, complete) = match result {
|
||||
ScannerDiskScanOutcome::Partial(cache) => (cache, false),
|
||||
ScannerDiskScanOutcome::Complete(cache) => (cache, true),
|
||||
_ => panic!("fixture namespace remains present"),
|
||||
};
|
||||
if round == 0 && resume_mode == HealScanMode::Deep {
|
||||
assert_eq!(cache.find("bucket/prefix/0000").expect("Deep revisits earlier prefix").size, 7);
|
||||
}
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("save bounded recovery");
|
||||
if complete {
|
||||
let root = store.strict_load().await.checked_flatten("bucket").expect("certified root");
|
||||
assert_eq!(root.objects, 9);
|
||||
assert_eq!(root.size, if resume_mode == HealScanMode::Deep { 15 } else { 9 });
|
||||
return;
|
||||
}
|
||||
}
|
||||
panic!("sampling interruption must recover with the unchanged three-object budget");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_save_reload_resume() {
|
||||
@@ -242,23 +767,48 @@ async fn checkpoint_fixture_save_reload_resume() {
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_hot_digest_diagnostic() {
|
||||
async fn checkpoint_fixture_hot_digest_retains_partial_progress() {
|
||||
run_checkpoint_fixture(true).await;
|
||||
}
|
||||
|
||||
async fn write_checkpoint_object(root: &std::path::Path, object: &str, versions: &[(Option<Uuid>, i64)]) {
|
||||
let mut metadata = FileMeta::new();
|
||||
for (index, (version_id, size)) in versions.iter().enumerate() {
|
||||
let mut info = FileInfo::new(object, 4, 2);
|
||||
info.volume = "bucket".into();
|
||||
info.version_id = *version_id;
|
||||
info.versioned = version_id.is_some();
|
||||
info.size = *size;
|
||||
info.mod_time = Some(
|
||||
OffsetDateTime::from_unix_timestamp(1_700_000_000 + i64::try_from(index).expect("fixture index"))
|
||||
.expect("non-sentinel fixture modification time"),
|
||||
);
|
||||
metadata.add_version(info).expect("fixture version");
|
||||
}
|
||||
write_test_object_metadata_bytes(root, "bucket", object, &metadata.marshal_msg().expect("fixture metadata")).await;
|
||||
}
|
||||
|
||||
async fn run_checkpoint_fixture(change_digest: bool) {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _guard = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for index in 0..STATIC_OBJECTS {
|
||||
write_test_object_metadata(&root, "bucket", &format!("static/{index:04}")).await;
|
||||
write_checkpoint_object(&root, &format!("static/{index:04}"), &[(None, 1)]).await;
|
||||
}
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
version: 1,
|
||||
bucket_incarnation: Uuid::from_u128(1),
|
||||
set_layout: DataUsageScanPlanDigest([41; 32]),
|
||||
publication_epoch: 0,
|
||||
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
|
||||
scan_mode: HealScanMode::Normal,
|
||||
};
|
||||
let store = FixtureStore::new();
|
||||
let mut previous = 0;
|
||||
let mut visited = 0;
|
||||
for round in 0..3_u8 {
|
||||
write_test_object_metadata(&root, "bucket", "hot/current").await;
|
||||
write_checkpoint_object(&root, "hot/current", &[(None, 1)]).await;
|
||||
let mut cache = DataUsageCache::default();
|
||||
let revisions = cache
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
@@ -278,9 +828,11 @@ async fn run_checkpoint_fixture(change_digest: bool) {
|
||||
crate::scanner_io::DataUsageCacheReuseOptions {
|
||||
require_source: true,
|
||||
tier_registry_generation: None,
|
||||
checkpoint_identity: Some(identity),
|
||||
},
|
||||
);
|
||||
let prepared = retained(&cache);
|
||||
cache.info.skip_healing = true;
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
@@ -326,13 +878,12 @@ async fn run_checkpoint_fixture(change_digest: bool) {
|
||||
eprintln!(
|
||||
"checkpoint_fixture round={round} hot_digest={change_digest} visited_total={visited} before={previous} prepared={prepared} scanned={scanned} reloaded={reloaded} diagnosis={diagnosis:?}"
|
||||
);
|
||||
if !change_digest || std::env::var_os("RUSTFS_CHECKPOINT_REQUIRE_PROGRESS").is_some() {
|
||||
assert_eq!(
|
||||
diagnosis,
|
||||
CoverageDiagnosis::Progress,
|
||||
"visited growth must produce durable static coverage"
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
diagnosis,
|
||||
CoverageDiagnosis::Progress,
|
||||
"visited growth must produce durable static coverage"
|
||||
);
|
||||
assert!(loaded.info.scan_plan_digest.is_none(), "old readers must rebuild an uncertified sweep");
|
||||
crate::remote_scanner::checkpoint_fixture_partial_return(budget.progress(), budget.entries_visited()).await;
|
||||
previous = reloaded;
|
||||
}
|
||||
@@ -383,28 +934,85 @@ async fn run_checkpoint_fixture(change_digest: bool) {
|
||||
assert!(result.is_err(), "pre-scan cancellation must not produce a complete root");
|
||||
assert_eq!(budget.reason(), None, "parent cancellation is not object budget exhaustion");
|
||||
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new(&parent, Default::default());
|
||||
let result = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget,
|
||||
vec![scanner.local_disk.clone()],
|
||||
loaded,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
store.reject_save.store(false, Ordering::SeqCst);
|
||||
write_checkpoint_object(&root, "static/0000", &[(Some(Uuid::from_u128(2)), 7), (Some(Uuid::from_u128(3)), 3)]).await;
|
||||
tokio::fs::remove_dir_all(root.join("bucket/static/0001"))
|
||||
.await
|
||||
.expect("unbounded scan must complete after durable partial progress");
|
||||
let ScannerDiskScanOutcome::Complete(cache) = result else {
|
||||
panic!("unbounded fixture must produce a complete disk cache");
|
||||
};
|
||||
assert!(cache.info.snapshot_complete);
|
||||
assert!(cache.info.scan_checkpoint.is_none());
|
||||
assert_eq!(
|
||||
cache.checked_flatten("bucket").expect("complete bucket root").objects,
|
||||
usize::try_from(STATIC_OBJECTS + 1).expect("fixture object count fits usize")
|
||||
);
|
||||
.expect("remove previously scanned fixture object");
|
||||
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 {
|
||||
let mut cache = DataUsageCache::default();
|
||||
let revisions = cache
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("reload every bounded round");
|
||||
crate::scanner_io::current_cache_root_or_prepare_with_generation(
|
||||
&mut cache,
|
||||
"bucket",
|
||||
SOURCE,
|
||||
11,
|
||||
7,
|
||||
final_plan,
|
||||
crate::scanner_io::DataUsageCacheReuseOptions {
|
||||
require_source: true,
|
||||
tier_registry_generation: Some(identity.tier_registry_generation),
|
||||
checkpoint_identity: Some(identity),
|
||||
},
|
||||
);
|
||||
cache.info.skip_healing = true;
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(4),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let outcome = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget.clone(),
|
||||
vec![scanner.local_disk.clone()],
|
||||
cache,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("bounded sweep outcome");
|
||||
let (cache, complete) = match outcome {
|
||||
ScannerDiskScanOutcome::Complete(cache) => (cache, true),
|
||||
ScannerDiskScanOutcome::Partial(cache) => (cache, false),
|
||||
ScannerDiskScanOutcome::NamespaceNotFound(_) => panic!("fixture namespace exists"),
|
||||
};
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("save each bounded sweep");
|
||||
let saved = store.strict_load().await;
|
||||
if complete {
|
||||
assert!(
|
||||
saw_mixed_sweep_end,
|
||||
"a clean tail must first finish as partial before a new validation sweep"
|
||||
);
|
||||
assert!(saved.info.snapshot_complete);
|
||||
assert!(saved.info.scan_progress.is_none());
|
||||
assert!(saved.info.scan_checkpoint.is_none());
|
||||
assert_eq!(saved.info.scan_plan_digest, Some(final_plan));
|
||||
let total = saved.checked_flatten("bucket").expect("complete bucket root");
|
||||
assert_eq!((total.objects, total.versions, total.size), (25, 2, 34));
|
||||
assert_eq!(saved.checked_flatten("bucket/static").expect("static subtree").objects, 23);
|
||||
assert_eq!(saved.checked_flatten("bucket/hot").expect("hot subtree").objects, 2);
|
||||
return;
|
||||
}
|
||||
assert!(!saved.info.snapshot_complete);
|
||||
assert!(saved.info.scan_plan_digest.is_none());
|
||||
if !budget.budget_elapsed() {
|
||||
saw_mixed_sweep_end = true;
|
||||
}
|
||||
}
|
||||
panic!("finite stable fixture must converge using the same four-object budget without an unbounded final sweep");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user