diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index f7710d68a..2f44a941e 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -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; diff --git a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture/deep_compacted.rs b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture/deep_compacted.rs index 7a2d6ae9c..23d0332dd 100644 --- a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture/deep_compacted.rs +++ b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture/deep_compacted.rs @@ -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"); +}