diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index eeb44d5ac..44230f9c0 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -567,6 +567,11 @@ pub struct DataUsageCacheInfo { pub snapshot_complete: bool, #[serde(default)] pub scan_plan_digest: Option, + /// Full activity and inventory scope of a set scan; only a complete + /// snapshot proves coverage. Bucket caches bind this scope into their + /// opaque scan plan digest instead. + #[serde(default)] + pub scan_coverage_digest: Option, #[serde(default)] pub cache_key_format: u16, /// Registry generation used for the completed/partial scan. This is @@ -602,6 +607,7 @@ impl Serialize for DataUsageCacheInfo { + 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.scan_coverage_digest.is_some()) + usize::from(self.tier_registry_generation.is_some()) + usize::from(!self.size_reconciliation.is_empty()) + usize::from(self.lkg_snapshot_complete) @@ -634,6 +640,9 @@ impl Serialize for DataUsageCacheInfo { state.serialize_entry("source", &self.source)?; state.serialize_entry("snapshot_complete", &self.snapshot_complete)?; state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?; + if let Some(coverage) = self.scan_coverage_digest { + state.serialize_entry("scan_coverage_digest", &coverage)?; + } state.serialize_entry("cache_key_format", &self.cache_key_format)?; if let Some(generation) = self.tier_registry_generation { state.serialize_entry("tier_registry_generation", &generation)?; diff --git a/crates/scanner/src/data_usage_define/tests.rs b/crates/scanner/src/data_usage_define/tests.rs index 624ba13a3..dc38b9822 100644 --- a/crates/scanner/src/data_usage_define/tests.rs +++ b/crates/scanner/src/data_usage_define/tests.rs @@ -29,6 +29,31 @@ use tokio::sync::Mutex; const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([3; 32]); +#[test] +fn scoped_scan_coverage_metadata_preserves_map_compatibility() { + #[derive(serde::Deserialize)] + struct LegacyInfo { + name: String, + next_cycle: u64, + } + let mut info = DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: 7, + ..Default::default() + }; + let old = serde_json::to_value(&info).expect("legacy metadata should encode"); + assert!(old.get("scan_coverage_digest").is_none()); + let old: DataUsageCacheInfo = serde_json::from_value(old).expect("missing coverage must remain readable"); + assert!(old.scan_coverage_digest.is_none()); + info.scan_coverage_digest = Some(TEST_PLAN_DIGEST); + let encoded = rmp_serde::to_vec(&info).expect("coverage metadata should remain map encoded"); + let legacy: LegacyInfo = rmp_serde::from_slice(&encoded).expect("old map readers should ignore additive proof fields"); + assert_eq!(legacy.name, DATA_USAGE_ROOT); + assert_eq!(legacy.next_cycle, 7); + let decoded: DataUsageCacheInfo = rmp_serde::from_slice(&encoded).expect("new reader should restore the coverage proof"); + assert_eq!(decoded.scan_coverage_digest, Some(TEST_PLAN_DIGEST)); +} + #[derive(Debug, PartialEq, Eq)] struct CachePutRecord { object: String, diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 9268a92ce..18970ac37 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -3438,7 +3438,6 @@ use cycle_state::*; use leadership::*; use usage_store::*; -#[cfg(test)] pub(crate) use activity::scanner_activity_snapshot_digest; pub use activity::scanner_topology_digest; pub(crate) use activity::{ diff --git a/crates/scanner/src/scanner/activity.rs b/crates/scanner/src/scanner/activity.rs index ffcbc5313..58c79f3ed 100644 --- a/crates/scanner/src/scanner/activity.rs +++ b/crates/scanner/src/scanner/activity.rs @@ -902,7 +902,6 @@ where observation } -#[cfg(test)] pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] { let mut hasher = Sha256::new(); hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes()); diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 189fa56d1..97da5e853 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -8585,6 +8585,62 @@ fn scanner_node_activity(epoch: &str, namespace_generation: u64, maintenance_gen } } +#[test] +fn scoped_scan_remote_dirty_coverage_invalidates_local_bucket_current() { + let before = BTreeMap::from([("remote".to_string(), scanner_node_activity("epoch-a", 7, 3))]); + let mut after = before.clone(); + let remote = after.get_mut("remote").expect("remote activity should exist"); + remote.dirty_usage_generation += 1; + remote.dirty_usage_pending = true; + assert_eq!(scanner_activity_structural_digest(&before), scanner_activity_structural_digest(&after)); + let old_plan = crate::scanner_io::checkpoint_fixture_bucket_digest( + DataUsageScanPlanDigest(scanner_activity_snapshot_digest(&before)), + None, + ); + let new_plan = crate::scanner_io::checkpoint_fixture_bucket_digest( + DataUsageScanPlanDigest(scanner_activity_snapshot_digest(&after)), + None, + ); + assert_ne!( + old_plan, new_plan, + "remote dirty changes must fence Current even without a local bucket hint" + ); + let source = DataUsageCacheSource::new(0, 0); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + leader_epoch: 11, + source: Some(source), + last_update: Some(std::time::SystemTime::UNIX_EPOCH), + snapshot_complete: true, + scan_plan_digest: Some(old_plan), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.replace("bucket", "", DataUsageEntry::default()); + assert!(matches!( + crate::scanner_io::current_cache_root_or_prepare_with_generation( + &mut cache, + "bucket", + source, + 7, + 11, + new_plan, + crate::scanner_io::DataUsageCacheReuseOptions { + require_source: true, + tier_registry_generation: None + }, + ), + crate::scanner_io::DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + .. + } + )); +} + #[test] fn scanner_activity_snapshot_digest_fences_storage_topology() { let first = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]); diff --git a/crates/scanner/src/scanner_io/cache.rs b/crates/scanner/src/scanner_io/cache.rs index b271c6551..8b13c558a 100644 --- a/crates/scanner/src/scanner_io/cache.rs +++ b/crates/scanner/src/scanner_io/cache.rs @@ -235,6 +235,7 @@ pub(super) struct ScannerSnapshotIdentity { pub(super) cycle: u64, pub(super) leader_epoch: u64, pub(super) plan_digest: DataUsageScanPlanDigest, + pub(super) coverage_digest: DataUsageScanPlanDigest, pub(super) tier_registry_generation: Option, } @@ -286,6 +287,7 @@ impl<'a> ValidatedScannerSnapshot<'a> { if result.info.next_cycle != scope.identity.cycle || result.info.leader_epoch != scope.identity.leader_epoch || result.info.scan_plan_digest != Some(scope.identity.plan_digest) + || result.info.scan_coverage_digest != Some(scope.identity.coverage_digest) || result.info.tier_registry_generation != scope.identity.tier_registry_generation { return Err(ScannerSnapshotValidationError::GenerationMismatch); @@ -684,6 +686,7 @@ pub(super) async fn persist_and_publish_cache_snapshot( expected_publication_epoch: u64, ) -> Option { let source = cache_snapshot.info.source?; + let coverage_digest = cache_snapshot.info.scan_coverage_digest?; let guard = match acquire_scanner_cache_locks(store.as_ref(), DATA_USAGE_CACHE_NAME, source).await { Ok(guard) => guard, Err(err) => { @@ -748,18 +751,20 @@ pub(super) async fn persist_and_publish_cache_snapshot( ); return None; } - if matches!( - current_cache_root_entry_with_generation( - &persisted, - DATA_USAGE_ROOT, - source, - cache_snapshot.info.next_cycle, - cache_snapshot.info.leader_epoch, - scan_plan_digest, - cache_snapshot.info.tier_registry_generation, - ), - Ok(Some(_)) - ) { + if persisted.info.scan_coverage_digest == Some(coverage_digest) + && matches!( + current_cache_root_entry_with_generation( + &persisted, + DATA_USAGE_ROOT, + source, + cache_snapshot.info.next_cycle, + cache_snapshot.info.leader_epoch, + scan_plan_digest, + cache_snapshot.info.tier_registry_generation, + ), + Ok(Some(_)) + ) + { cache_snapshot = persisted; } else { if guard.is_lock_lost() { diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index 4289d1efd..bfd59ea95 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -125,8 +125,8 @@ impl ScannerIOCache for SetDisks { pending_maintenance_work, cache_cycle_floor, } = scan_plan; - let bucket_work_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan); let scan_plan_digest = scanner_bucket_work_digest(scan_plan_digest, scan_mode, requires_full_scan); + let bucket_work_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan); let pool_label = self.pool_index.to_string(); let set_label = self.set_index.to_string(); @@ -165,8 +165,9 @@ impl ScannerIOCache for SetDisks { scan_plan_digest, }, ); - let mut scoped_cache = scoped_scan.map(|prepared| { + let mut scoped_cache = scoped_scan.map(|mut prepared| { buckets = prepared.buckets; + prepared.cache.info.scan_coverage_digest = Some(bucket_coverage_digest); prepared.cache }); if buckets.is_empty() { @@ -182,6 +183,7 @@ impl ScannerIOCache for SetDisks { tier_registry_generation: Some(tier_registry_generation), source: Some(source), scan_plan_digest: Some(scan_plan_digest), + scan_coverage_digest: Some(bucket_coverage_digest), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, @@ -469,6 +471,7 @@ impl ScannerIOCache for SetDisks { source: Some(source), snapshot_complete: false, scan_plan_digest: Some(scan_plan_digest), + scan_coverage_digest: Some(bucket_coverage_digest), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, lkg_snapshot_complete: old_cache.info.lkg_snapshot_complete, lkg_next_cycle: old_cache.info.lkg_next_cycle, diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 1d3ce57b5..5cb452d6d 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -277,9 +277,9 @@ where bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source); let structural_scan_plan_digest = scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before)); + let scan_plan_digest = scanner_bucket_work_digest(structural_scan_plan_digest, scan_mode, requires_full_scan); let bucket_coverage_digest = scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before)); - let scan_plan_digest = scanner_bucket_work_digest(structural_scan_plan_digest, scan_mode, requires_full_scan); let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list)); let scan_scope = resolve_scanner_bucket_scan_scope( store, @@ -568,6 +568,7 @@ where cycle: want_cycle, leader_epoch, plan_digest: scan_plan_digest, + coverage_digest: bucket_coverage_digest, tier_registry_generation: Some(tier_registry_generation), }, }, diff --git a/crates/scanner/src/scanner_io/publish_gate_tests.rs b/crates/scanner/src/scanner_io/publish_gate_tests.rs index 254b561cc..bec3a40e5 100644 --- a/crates/scanner/src/scanner_io/publish_gate_tests.rs +++ b/crates/scanner/src/scanner_io/publish_gate_tests.rs @@ -17,6 +17,7 @@ use crate::data_usage_define::{UNKNOWN_TIER, UnknownTierStats, hash_path}; use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccountingProof}; const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]); +const TEST_COVERAGE_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([6; 32]); #[test] fn scanner_bucket_inventory_requires_exact_unique_set_union() { @@ -105,6 +106,7 @@ fn completed_root_cache(bucket: &str, objects: usize, update_secs: u64, source: source: Some(source), snapshot_complete: true, scan_plan_digest: Some(TEST_PLAN_DIGEST), + scan_coverage_digest: Some(TEST_COVERAGE_DIGEST), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, @@ -156,6 +158,7 @@ fn completed_usage_for_scope( cycle: first.info.next_cycle, leader_epoch: first.info.leader_epoch, plan_digest: TEST_PLAN_DIGEST, + coverage_digest: TEST_COVERAGE_DIGEST, tier_registry_generation: first.info.tier_registry_generation, }, }, @@ -227,6 +230,7 @@ fn completed_data_usage_info_binds_all_results_to_requested_identity() { cycle: 0, leader_epoch: 0, plan_digest: TEST_PLAN_DIGEST, + coverage_digest: TEST_COVERAGE_DIGEST, tier_registry_generation: None, }; let results = [set]; @@ -244,6 +248,10 @@ fn completed_data_usage_info_binds_all_results_to_requested_identity() { tier_registry_generation: Some(1), ..identity }, + ScannerSnapshotIdentity { + coverage_digest: DataUsageScanPlanDigest([4; 32]), + ..identity + }, ] { let scope = ScannerSnapshotScope { sources: &sources, @@ -1380,6 +1388,40 @@ fn scoped_scan_bucket_work_proof_fences_same_cycle_cache() { ); } +#[test] +fn scoped_scan_complete_root_requires_current_coverage_from_every_set() { + let sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]); + let buckets = vec!["bucket".to_string()]; + let coverage = DataUsageScanPlanDigest([4; 32]); + let scope = ScannerSnapshotScope { + sources: &sources, + buckets: &buckets, + identity: ScannerSnapshotIdentity { + cycle: 0, + leader_epoch: 0, + plan_digest: TEST_PLAN_DIGEST, + coverage_digest: coverage, + tier_registry_generation: None, + }, + }; + for (first_coverage, second_coverage, valid) in [ + (Some(coverage), Some(coverage), true), + (None, Some(coverage), false), + (Some(coverage), None, false), + (None, None, false), + (Some(coverage), Some(DataUsageScanPlanDigest([5; 32])), false), + ] { + let mut first = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0)); + let mut second = completed_root_cache("bucket", 3, 10, DataUsageCacheSource::new(1, 0)); + first.info.scan_coverage_digest = first_coverage; + second.info.scan_coverage_digest = second_coverage; + assert_eq!( + completed_data_usage_info(&[first, second], &scope, &[], true, false, false).is_some(), + valid + ); + } +} + #[test] fn scanner_cache_lock_resource_is_scoped_to_cache_source() { let cache_name = "photos/.usage-cache.bin";