mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 12:31:36 +00:00
fix(scanner): preserve checkpoint children in large prefixes (#8292)
* fix(scanner): preserve checkpoint children in large prefixes * fix(ecstore): restore strict Clippy compatibility on Rust 1.99 Use try_update without changing atomic ordering or overflow behavior. Keep recursive storage futures boxed once at each frame and remove the redundant async-recursion macro, including its non-recursive SQL planner use. Remove needless closure borrows and orphaned dependency entries. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> (cherry picked from commit537c15277a) * fix(deps): replace yanked yoke-derive release (cherry picked from commitb54b32b8a9) * chore: keep scanner repair separate from CI corrections * test(scanner): verify large-prefix resume and compaction lifecycle Exercise Normal and Deep scans across cycle and leader changes. Validate persisted frontiers, bound replay to interrupted boundary objects, and prove that a completed same-plan sweep restores compaction without losing object or byte totals. Refs: rustfs/backlog#2710 Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com> --------- Co-authored-by: Hauser <housemecn@gmail.com> Co-authored-by: heihutu <heihutu@gmail.com> Co-authored-by: zhi22915 <qiuzgang@gmail.com> Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
@@ -2072,8 +2072,10 @@ impl FolderScanner {
|
||||
break;
|
||||
}
|
||||
|
||||
// If we have many subfolders, compact ourself.
|
||||
// Checkpointed sweeps need child entries to validate their resume
|
||||
// frontier. Keep the emergency compaction limit for large trees.
|
||||
let should_compact = (self.new_cache.info.name != folder.name
|
||||
&& self.old_cache.info.scan_progress.is_none()
|
||||
&& existing_folders.len() + new_folders.len() >= DATA_SCANNER_COMPACT_AT_FOLDERS)
|
||||
|| existing_folders.len() + new_folders.len() >= DATA_SCANNER_FORCE_COMPACT_AT_FOLDERS;
|
||||
|
||||
|
||||
@@ -202,3 +202,94 @@ async fn deep_compacted_budget_preserves_partial_checkpoint() {
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn large_prefix_resumes_after_budget_across_cycles_and_leaders() {
|
||||
check_large_prefix_resumption(HealScanMode::Normal).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn large_prefix_deep_resumes_after_budget_across_cycles_and_leaders() {
|
||||
check_large_prefix_resumption(HealScanMode::Deep).await;
|
||||
}
|
||||
|
||||
async fn check_large_prefix_resumption(mode: HealScanMode) {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _cleanup = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for index in 0..DATA_SCANNER_COMPACT_AT_FOLDERS {
|
||||
write_checkpoint_object(&root, &format!("prefix/{index:04}"), &[(None, 1)]).await;
|
||||
}
|
||||
write_checkpoint_object(&root, "z-last", &[(None, 1)]).await;
|
||||
let identity = crate::DataUsageScanIdentity {
|
||||
scan_mode: mode,
|
||||
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
|
||||
..bound_checkpoint().1
|
||||
};
|
||||
let store = FixtureStore::new();
|
||||
let mut cache = DataUsageCache::default();
|
||||
let mut previous = 0;
|
||||
let mut scanned_objects = 0;
|
||||
for round in 0..8 {
|
||||
cache.prepare_bucket_checkpoint("bucket", 11 + round, 7 + round, SOURCE, PLAN, identity);
|
||||
cache.info.skip_healing = true;
|
||||
let (outcome, budget) = scan(&scanner.local_disk, cache, mode, 1000).await;
|
||||
scanned_objects += budget.progress().0;
|
||||
let (scanned, complete) = match outcome {
|
||||
ScannerDiskScanOutcome::Complete(cache) => (cache, true),
|
||||
ScannerDiskScanOutcome::Partial(cache) => (cache, false),
|
||||
ScannerDiskScanOutcome::NamespaceNotFound(_) => panic!("fixture namespace exists"),
|
||||
};
|
||||
cache = save_reload(&store, &scanned).await;
|
||||
let total = cache.checked_flatten("bucket").expect("saved bucket coverage");
|
||||
assert!(total.objects > previous, "round {round}: large-prefix coverage must advance after reload");
|
||||
if complete {
|
||||
let expected = DATA_SCANNER_COMPACT_AT_FOLDERS + 1;
|
||||
assert_eq!(total.objects, expected);
|
||||
assert_eq!(total.size, expected);
|
||||
assert_eq!(cache.find("bucket/z-last").map(|entry| entry.objects), Some(1));
|
||||
assert!(cache.info.snapshot_complete);
|
||||
let expected_reads = u64::try_from(expected).expect("fixture object count fits u64");
|
||||
// The budget cancels after accounting the last object, before its
|
||||
// child traversal completes. Only that boundary object may replay.
|
||||
assert!(
|
||||
(expected_reads..=expected_reads + round).contains(&scanned_objects),
|
||||
"resumption must not reread objects before the completed frontier"
|
||||
);
|
||||
assert!(cache.info.scan_progress.is_none());
|
||||
assert!(cache.info.scan_checkpoint.is_none());
|
||||
assert!(cache.info.scan_coverage_receipt.is_none());
|
||||
|
||||
// Completing the sweep must restore compaction on the next same-plan scan.
|
||||
let retained_entries = cache.cache.len();
|
||||
assert!(!cache.find(PREFIX).expect("completed prefix").compacted);
|
||||
cache.prepare_bucket_checkpoint("bucket", 12 + round, 7 + round, SOURCE, PLAN, identity);
|
||||
assert!(cache.info.scan_progress.is_none(), "a complete same-plan cache needs no forward sweep");
|
||||
let (outcome, _) = scan(&scanner.local_disk, cache, mode, expected_reads + 1).await;
|
||||
let ScannerDiskScanOutcome::Complete(compacted) = outcome else {
|
||||
panic!("same-plan scan must complete after resuming the large prefix")
|
||||
};
|
||||
let reloaded = save_reload(&store, &compacted).await;
|
||||
assert!(reloaded.find(PREFIX).expect("compacted prefix").compacted);
|
||||
assert!(
|
||||
reloaded.cache.len() < retained_entries,
|
||||
"completed sweeps must release child cache entries"
|
||||
);
|
||||
let total = reloaded.checked_flatten("bucket").expect("compacted bucket coverage");
|
||||
assert_eq!((total.objects, total.size), (expected, expected));
|
||||
assert!(reloaded.info.snapshot_complete);
|
||||
return;
|
||||
}
|
||||
assert_eq!(budget.reason(), Some(ScannerCycleBudgetReason::Objects));
|
||||
assert!(!cache.info.snapshot_complete);
|
||||
let frontier = cache
|
||||
.validated_scan_frontier()
|
||||
.expect("saved partial coverage must have a valid frontier");
|
||||
assert!(cache.find(frontier).is_some(), "the persisted frontier must retain its child record");
|
||||
previous = total.objects;
|
||||
}
|
||||
panic!("large prefix must finish within repeated bounded scans");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user