diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 9d3c89fe8..eeb44d5ac 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -14,6 +14,7 @@ use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration}; use serde::{Deserialize, Serialize, ser::SerializeMap}; +use sha2::{Digest, Sha256}; use std::{ collections::{HashMap, HashSet}, future::Future, @@ -409,23 +410,27 @@ pub struct DataUsageScanIdentity { pub set_layout: DataUsageScanPlanDigest, pub publication_epoch: u64, pub tier_registry_generation: u64, + pub scan_mode: HealScanMode, } impl Serialize for DataUsageScanIdentity { fn serialize(&self, serializer: S) -> Result { - let mut map = serializer.serialize_map(Some(5))?; + let mut map = serializer.serialize_map(Some(6))?; map.serialize_entry("version", &self.version)?; map.serialize_entry("bucket_incarnation", &self.bucket_incarnation)?; map.serialize_entry("set_layout", &self.set_layout)?; map.serialize_entry("publication_epoch", &self.publication_epoch)?; map.serialize_entry("tier_registry_generation", &self.tier_registry_generation)?; + map.serialize_entry("scan_mode", &self.scan_mode)?; map.end() } } impl DataUsageScanIdentity { pub(crate) fn is_valid(&self) -> bool { - self.version == 1 && !self.bucket_incarnation.is_nil() + self.version == 1 + && !self.bucket_incarnation.is_nil() + && matches!(self.scan_mode, HealScanMode::Normal | HealScanMode::Deep) } } @@ -446,6 +451,35 @@ impl Serialize for DataUsageScanProgress { } } +#[derive(Clone, Debug, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct DataUsageScanCoverageReceipt { + pub through: String, + pub digest: [u8; 32], +} + +impl Serialize for DataUsageScanCoverageReceipt { + fn serialize(&self, serializer: S) -> Result { + let mut map = serializer.serialize_map(Some(2))?; + map.serialize_entry("through", &self.through)?; + map.serialize_entry("digest", &self.digest)?; + map.end() + } +} + +struct CheckpointDigestWriter(Sha256); + +impl std::io::Write for CheckpointDigestWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.0.update(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageEntryInfo { pub name: String, @@ -520,6 +554,8 @@ pub struct DataUsageCacheInfo { #[serde(default)] pub scan_progress: Option, #[serde(default)] + pub scan_coverage_receipt: Option, + #[serde(default)] pub pending_heals: Vec, #[serde(default)] pub object_lock: Option>, @@ -565,6 +601,7 @@ impl Serialize for DataUsageCacheInfo { let field_count = 16 + usize::from(self.scan_identity.is_some()) + usize::from(self.scan_progress.is_some()) + + usize::from(self.scan_coverage_receipt.is_some()) + usize::from(self.tier_registry_generation.is_some()) + usize::from(!self.size_reconciliation.is_empty()) + usize::from(self.lkg_snapshot_complete) @@ -589,6 +626,9 @@ impl Serialize for DataUsageCacheInfo { if let Some(progress) = self.scan_progress { state.serialize_entry("scan_progress", &progress)?; } + if let Some(receipt) = &self.scan_coverage_receipt { + state.serialize_entry("scan_coverage_receipt", receipt)?; + } state.serialize_entry("pending_heals", &self.pending_heals)?; state.serialize_entry("object_lock", &self.object_lock)?; state.serialize_entry("source", &self.source)?; @@ -777,7 +817,14 @@ impl DataUsageCache { && self.info.scan_identity == Some(identity) && self.info.tier_registry_generation == Some(identity.tier_registry_generation) && (self.cache.is_empty() || self.checked_flatten_complete_scope(name).is_some()); - if reusable && self.info.scan_progress.is_none() && self.info.scan_plan_digest == Some(scan_plan_digest) { + if reusable + && self.info.snapshot_complete + && self.info.scan_progress.is_none() + && self.info.scan_checkpoint.is_none() + && self.info.scan_resume_after.is_none() + && self.info.scan_coverage_receipt.is_none() + && self.info.scan_plan_digest == Some(scan_plan_digest) + { return self.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, true); } if !reusable { @@ -798,16 +845,10 @@ impl DataUsageCache { self.info.pending_heals = pending_heals; self.info.size_reconciliation = size_reconciliation; } - let cursor_is_valid = match (&self.info.scan_checkpoint, &self.info.scan_resume_after) { - (None, None) => true, - (Some(checkpoint), Some(resume)) => { - checkpoint.version == DATA_USAGE_SCAN_CHECKPOINT_VERSION - && checkpoint.resume_after == *resume - && resume.strip_prefix(name).is_some_and(|suffix| suffix.starts_with('/')) - && self.find(resume).is_some() - } - _ => false, - }; + let cursor_is_valid = (self.info.scan_checkpoint.is_none() + && self.info.scan_resume_after.is_none() + && self.info.scan_coverage_receipt.is_none()) + || self.validated_scan_frontier().is_some(); if !cursor_is_valid { self.info.scan_progress = None; } @@ -828,6 +869,7 @@ impl DataUsageCache { }); self.info.scan_resume_after = None; self.info.scan_checkpoint = None; + self.info.scan_coverage_receipt = None; } // Old readers do not understand coverage sweeps. An absent plan makes // their existing prepare path rebuild instead of promoting mixed data. @@ -839,6 +881,86 @@ impl DataUsageCache { } } + fn coverage_prefix_digest(&self, through: &str) -> Result<[u8; 32], serde_json::Error> { + let mut writer = CheckpointDigestWriter(Sha256::new()); + serde_json::to_writer( + &mut writer, + &( + &self.info.name, + self.info.scan_identity, + self.info.source, + self.info.leader_epoch, + self.info.cache_key_format, + self.info.scan_progress.map(|progress| progress.started_plan), + through, + ), + )?; + let mut prefix = self + .cache + .iter() + .filter(|(key, _)| { + let ancestor = through + .strip_prefix(key.as_str()) + .is_some_and(|suffix| suffix.starts_with('/')); + let descendant = key.strip_prefix(through).is_some_and(|suffix| suffix.starts_with('/')); + (key.as_str() <= through && !ancestor) || descendant + }) + .collect::>(); + prefix.sort_unstable_by(|(left, _), (right, _)| left.cmp(right)); + for (key, entry) in prefix { + let mut value = serde_json::to_value(entry)?; + value.sort_all_objects(); + if let Some(children) = value.get_mut("children").and_then(serde_json::Value::as_array_mut) { + children.sort_unstable_by(|left, right| left.as_str().cmp(&right.as_str())); + } + serde_json::to_writer(&mut writer, &(key, value))?; + } + Ok(writer.0.finalize().into()) + } + + pub(crate) fn validated_scan_frontier(&self) -> Option<&str> { + let receipt = self.info.scan_coverage_receipt.as_ref()?; + let checkpoint = self.info.scan_checkpoint.as_ref()?; + (self.info.scan_progress.is_some() + && self.info.scan_identity.is_some_and(|identity| identity.is_valid()) + && self.info.source.is_some() + && receipt.through.len() <= 16 * 1024 + && checkpoint.version == DATA_USAGE_SCAN_CHECKPOINT_VERSION + && checkpoint.resume_after == receipt.through + && self.info.scan_resume_after.as_deref() == Some(receipt.through.as_str()) + && receipt + .through + .strip_prefix(&self.info.name) + .is_some_and(|suffix| suffix.starts_with('/')) + && self.find(&receipt.through).is_some() + && self.coverage_prefix_digest(&receipt.through).ok() == Some(receipt.digest)) + .then_some(receipt.through.as_str()) + } + + /// Seal only the frontier supplied by completed traversal, never a restored cursor. + pub(crate) fn seal_scan_frontier(&mut self, frontier: Option<&str>) -> Result<(), serde_json::Error> { + if self.info.scan_progress.is_none() { + self.info.scan_coverage_receipt = None; + return Ok(()); + } + let frontier = frontier.filter(|path| path.len() <= 16 * 1024 && self.find(path).is_some()); + self.info.scan_coverage_receipt = match frontier { + Some(through) => Some(DataUsageScanCoverageReceipt { + through: through.to_owned(), + digest: self.coverage_prefix_digest(through)?, + }), + None => None, + }; + self.info.scan_resume_after = frontier.map(str::to_owned); + let reason = self + .info + .scan_checkpoint + .as_ref() + .map_or(DataUsageScanCheckpointReason::Unknown, |checkpoint| checkpoint.reason); + self.info.scan_checkpoint = frontier.map(|through| DataUsageScanCheckpoint::new(through.to_owned(), reason)); + Ok(()) + } + fn ensure_cache_save_metrics_registered() { CACHE_SAVE_METRICS_ONCE.call_once(|| { describe_counter!( diff --git a/crates/scanner/src/remote_scanner/stream.rs b/crates/scanner/src/remote_scanner/stream.rs index e7528e6e3..567490244 100644 --- a/crates/scanner/src/remote_scanner/stream.rs +++ b/crates/scanner/src/remote_scanner/stream.rs @@ -733,6 +733,7 @@ async fn scan_and_persist_local_bucket( &bucket, expected_publication_epoch, tier_registry_generation, + scan_mode, ) .await .ok(), @@ -835,6 +836,7 @@ async fn scan_and_persist_local_bucket( &bucket, expected_publication_epoch, tier_registry_generation, + scan_mode, ) .await .ok() diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 94ad28312..410cf6b48 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -707,6 +707,9 @@ pub struct FolderScanner { /// next scan and cannot mix generations in one aggregate. tier_registry: TierRegistrySnapshot, pending_heals_changed: bool, + coverage_frontier: Option, + resume_frontier: Option, + coverage_gap: bool, pending_size_reconciliation_keys: HashSet, pending_size_reconciliation_scopes: HashSet, pending_size_reconciliation_truncated: bool, @@ -906,6 +909,16 @@ impl FolderScanner { self.update_cache.info.scan_checkpoint = Some(checkpoint); } + fn record_completed_child(&mut self, folder: &str, healthy: bool) { + if self.old_cache.info.scan_progress.is_some() { + self.coverage_gap |= !healthy; + if !self.coverage_gap { + self.coverage_frontier = Some(folder.to_owned()); + } + } + self.record_scan_resume_hint(folder); + } + fn record_scan_resume_hint_if_not_ancestor(&mut self, folder: &str) { let keep_existing = self .new_cache @@ -1471,6 +1484,7 @@ impl FolderScanner { // (e.g. in the get_size error branch below). This branch only accounts // for subsequent skips of already-failed paths. if self.should_skip_failed(&item.path) { + self.coverage_gap |= self.old_cache.info.scan_progress.is_some(); continue; } @@ -1484,6 +1498,7 @@ impl FolderScanner { let failure_action = classify_get_size_failure(&item, &e); if failure_action != GetSizeFailureAction::Skip { + self.coverage_gap |= self.old_cache.info.scan_progress.is_some(); // Track failed objects to prevent infinite retry loops into.failed_objects += 1; self.record_failed(&item.path); @@ -1576,6 +1591,9 @@ impl FolderScanner { abandoned_children.remove(&path_join_buf(&[&item.bucket, &item.object_path()])); apply_scanner_size_summary(into, &sz); + if !sz.size_reconciliation.is_empty() { + self.coverage_gap |= self.old_cache.info.scan_progress.is_some(); + } self.apply_size_reconciliation(&sz); into.objects += 1; object_count += 1; @@ -1620,6 +1638,7 @@ impl FolderScanner { } if self.is_erasure_mode && found_erasure_data_directory && !found_object_metadata { + self.coverage_gap |= self.old_cache.info.scan_progress.is_some(); found_object_metadata = true; let metadata_path = path_join_buf(&[&dir_path, STORAGE_FORMAT_FILE]); @@ -1742,7 +1761,7 @@ impl FolderScanner { })); let has_queued_folders = !queued_folders.is_empty(); let forward_sweep = self.old_cache.info.scan_progress.is_some(); - let forward_resume_after = forward_sweep.then(|| scan_resume_after.map(str::to_owned)).flatten(); + let forward_resume_after = self.resume_frontier.clone(); let resume_order = if forward_sweep { order_queued_folders_for_resume(&mut queued_folders, None) } else { @@ -1828,7 +1847,7 @@ impl FolderScanner { // In compacted mode child totals are accumulated directly into the parent entry. let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), into)); fut.await.map_err(|e| ScannerError::Other(e.to_string()))?; - self.record_scan_resume_hint(&folder_item.name); + self.record_completed_child(&folder_item.name, into.failed_objects == 0); self.send_update_for_entry(&this_hash, &folder.parent, into).await; tokio::task::yield_now().await; } else { @@ -1853,12 +1872,13 @@ impl FolderScanner { error = %e, "Scanner child folder scan failed" ); + self.coverage_gap |= forward_sweep; continue; } tokio::task::yield_now().await; into.add_child(&h); - self.record_scan_resume_hint(&folder_item.name); + self.record_completed_child(&folder_item.name, dst.failed_objects == 0); // We scanned a folder, optionally send update. self.update_cache.delete_recursive(&h); self.update_cache.copy_with_children(&self.new_cache, &h, &folder_item.parent); @@ -2380,6 +2400,7 @@ pub async fn scan_data_folder( cache.fold_retired_tiers(&tier_registry.names); cache.info.tier_registry_generation = Some(tier_registry.generation); + let resume_frontier = cache.validated_scan_frontier().map(str::to_owned); // Create folder scanner let mut scanner = FolderScanner { root: base_path, @@ -2409,6 +2430,9 @@ pub async fn scan_data_folder( local_disk, tier_registry, pending_heals_changed: false, + coverage_frontier: resume_frontier.clone(), + resume_frontier, + coverage_gap: false, pending_size_reconciliation_keys: HashSet::new(), pending_size_reconciliation_scopes: HashSet::new(), pending_size_reconciliation_truncated: false, @@ -2439,11 +2463,13 @@ pub async fn scan_data_folder( match scanner.scan_folder(ctx.clone(), folder, &mut root).await { Ok(()) => { // Get the new cache and finalize it + let coverage_gap = scanner.coverage_gap; let new_cache = scanner.as_mut_new_cache(); new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); new_cache.info.last_update = Some(SystemTime::now()); new_cache.info.next_cycle = cache.info.next_cycle; - let unresolved_objects = root.failed_objects > 0 + let unresolved_objects = coverage_gap + || root.failed_objects > 0 || !new_cache.info.failed_objects.is_empty() || !new_cache.info.size_reconciliation.is_empty(); let mixed_coverage = new_cache @@ -2465,6 +2491,7 @@ pub async fn scan_data_folder( let had_scan_checkpoint = cache.info.scan_checkpoint.is_some() || new_cache.info.scan_checkpoint.is_some(); new_cache.info.scan_resume_after = None; new_cache.info.scan_checkpoint = None; + new_cache.info.scan_coverage_receipt = None; if had_scan_checkpoint { global_metrics().record_scanner_checkpoint_cleared(); } @@ -2484,6 +2511,7 @@ pub async fn scan_data_folder( if root_has_progress { scanner.carry_forward_old_children(&root_hash, &mut root); } + let coverage_frontier = scanner.coverage_frontier.clone(); let new_cache = scanner.as_mut_new_cache(); if root_has_progress { new_cache.replace_hashed(&root_hash, &None, &root); @@ -2498,6 +2526,18 @@ pub async fn scan_data_folder( if root_has_progress { set_scan_checkpoint(new_cache, checkpoint_reason_from_budget(budget.reason())); } + new_cache.seal_scan_frontier(coverage_frontier.as_deref())?; + if new_cache.info.scan_progress.is_some() { + if let Some(checkpoint) = &new_cache.info.scan_checkpoint { + global_metrics().record_scanner_checkpoint_set( + checkpoint.version, + checkpoint.resume_after.clone(), + checkpoint.reason.as_str(), + ); + } else { + global_metrics().record_scanner_checkpoint_cleared(); + } + } close_disk_guard.close().await; return Err(ScannerError::PartialCache(Box::new(new_cache.clone()))); } diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 650ccaab3..8e26429cc 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -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, diff --git a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs index bc234d2fd..f233bc19d 100644 --- a/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs +++ b/crates/scanner/src/scanner_folder/tests/checkpoint_fixture.rs @@ -241,6 +241,7 @@ fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) { 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); @@ -258,6 +259,9 @@ fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) { "bucket/static".into(), DataUsageScanCheckpointReason::Objects, )); + cache + .seal_scan_frontier(Some("bucket/static")) + .expect("completed fixture prefix receipt"); (cache, identity) } @@ -284,6 +288,7 @@ fn checkpoint_fixture_roundtrip_retains_verified_scope_but_old_reader_rebuilds() 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), @@ -300,6 +305,7 @@ fn checkpoint_fixture_unchanged_complete_plan_keeps_existing_rescan_policy() { 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), @@ -338,6 +344,10 @@ fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() { tier_registry_generation: 8, ..identity }, + crate::DataUsageScanIdentity { + scan_mode: HealScanMode::Deep, + ..identity + }, ] { let mut next = cache.clone(); assert_eq!( @@ -409,6 +419,332 @@ fn checkpoint_fixture_corrupt_cursor_restarts_validation_without_claiming_comple } } +#[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() { @@ -449,6 +785,7 @@ async fn run_checkpoint_fixture(change_digest: bool) { 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; diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 4e64e0250..1268bae15 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -737,6 +737,7 @@ pub(crate) async fn scanner_bucket_checkpoint_identity( bucket: &str, publication_epoch: u64, tier_registry_generation: u64, + scan_mode: HealScanMode, ) -> Result { let bucket_incarnation = set.bucket_incarnation_id_from_disk(bucket).await?; let disks = set @@ -760,6 +761,7 @@ pub(crate) async fn scanner_bucket_checkpoint_identity( set_layout: crate::DataUsageScanPlanDigest(digest.finalize().into()), publication_epoch, tier_registry_generation, + scan_mode, }) } diff --git a/crates/scanner/src/scanner_io/cache.rs b/crates/scanner/src/scanner_io/cache.rs index fc15b9359..b271c6551 100644 --- a/crates/scanner/src/scanner_io/cache.rs +++ b/crates/scanner/src/scanner_io/cache.rs @@ -113,6 +113,9 @@ pub(crate) fn current_cache_root_entry_with_generation( && cache.info.source == Some(source) && cache.info.snapshot_complete && cache.info.scan_progress.is_none() + && cache.info.scan_checkpoint.is_none() + && cache.info.scan_resume_after.is_none() + && cache.info.scan_coverage_receipt.is_none() && cache.info.scan_plan_digest == Some(scan_plan_digest) && cache.info.last_update.is_some() && cache.info.next_cycle == next_cycle diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index 7f6f29924..6085558bd 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -888,6 +888,7 @@ impl ScannerIOCache for SetDisks { &bucket.name, expected_publication_epoch_clone, tier_registry_generation, + scan_mode, ) .await .ok(); @@ -1057,6 +1058,7 @@ impl ScannerIOCache for SetDisks { &bucket.name, expected_publication_epoch_clone, tier_registry_generation, + scan_mode, ) .await .ok() diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 5c46dada0..2bbbc918f 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -131,26 +131,28 @@ async fn checkpoint_fixture_bucket_identity_uses_its_set_instance_owner() { .make_bucket("checkpoint-identity", &MakeBucketOptions::default()) .await .expect("first instance bucket"); - let first_identity = scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7) - .await - .expect("first durable identity"); + let first_identity = + scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal) + .await + .expect("first durable identity"); let (_second_dir, second) = setup_two_pool_scanner_store().await; second .make_bucket("checkpoint-identity", &MakeBucketOptions::default()) .await .expect("second instance bucket"); - let second_identity = scanner_bucket_checkpoint_identity(&second.pools[0].disk_set[0], "checkpoint-identity", 0, 7) - .await - .expect("second durable identity"); + let second_identity = + scanner_bucket_checkpoint_identity(&second.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal) + .await + .expect("second durable identity"); assert_ne!(first_identity.bucket_incarnation, second_identity.bucket_incarnation); assert_eq!( - scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7) + scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal) .await .expect("first owner remains bound"), first_identity ); assert!( - scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "missing-checkpoint-bucket", 0, 7) + scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "missing-checkpoint-bucket", 0, 7, HealScanMode::Normal) .await .is_err() ); diff --git a/docs/testing/scanner-checkpoint-fixture.md b/docs/testing/scanner-checkpoint-fixture.md index 9ffc7dcd3..1c413a4b3 100644 --- a/docs/testing/scanner-checkpoint-fixture.md +++ b/docs/testing/scanner-checkpoint-fixture.md @@ -13,7 +13,11 @@ Both the unchanged-plan and hot-plan cases require durable static coverage to in After three interrupted rounds, the fixture overwrites a previously visited object with two versions, deletes another visited object, and creates one more hot object. It then keeps the same four-object budget until the stable namespace is certified. Finishing a sweep that spans different mutation plans must first return partial; a subsequent verification sweep must produce exactly 25 objects, 2 versioned entries, and 34 logical bytes. There is no final unbudgeted sweep. -The new bucket checkpoint binds the persisted bucket incarnation, set layout, publication epoch and tier generation, with the existing source/leader/key-format checks. Its forward sweep records the starting and requested mutation plans separately. Partial sweeps omit the legacy `scan_plan_digest`, so older readers rebuild instead of treating mixed observations as a current complete snapshot. Completed sweeps restore that digest only after covering one mutation plan. Unsupported or missing identities retain the legacy rebuild path. Round-trip, stale identity, invalid cursor and instance-owner tests cover these boundaries. The new metadata remains map-encoded with optional top-level fields. +The new bucket checkpoint binds the persisted bucket incarnation, set layout, publication epoch, tier generation and scan mode, with the existing source/leader/key-format checks. Its forward sweep records the starting and requested mutation plans separately. Partial sweeps omit the legacy `scan_plan_digest`, so older readers rebuild instead of treating mixed observations as a current complete snapshot. Completed sweeps restore that digest only after covering one mutation plan. Unsupported or missing identities retain the legacy rebuild path. The stable-plan fast path requires a complete snapshot without unfinished checkpoint state; its interrupted result must enter forward validation on reload. + +A coverage receipt binds the completed traversal frontier to its scope, starting plan and canonical covered-prefix digest. It excludes ancestor aggregates and the unvisited suffix, so unrelated suffix changes cannot invalidate completed work. Cancellation seals only the completed frontier; failed child traversal and known failed-metadata skips block further frontier advancement. A saved cursor pointing at an existing but unvisited old subtree is rejected unless it agrees with that receipt. The receipt is a consistency check for storage owned by the scanner, not authentication against a party able to forge the entire cache and recompute its digest. New metadata remains map-encoded with optional top-level fields. + +The Normal-to-Deep regression holds the mutation plan fixed, changes metadata in a previously visited prefix, and verifies that the real Deep disk-scan entry point reads that prefix again. It also checks that a complete Normal cache cannot satisfy the Deep `Current` path. The fixture disables heal side effects; it proves traversal re-entry, not actual bitrot detection or repair. Additional tests cover map round-trips, stale identities, existing-but-uncovered cursors, coverage gaps, and per-instance metadata ownership. This fixture bounds object processing after directory enumeration. It does not prove fixed-budget enumeration of arbitrarily wide directories or real process-restart convergence. Those gates require a storage-owned resumable enumeration capability, including its initial construction cost; a readdir offset, an in-memory iterator or a last-name filter is not that capability.