From e3fd1ae19aa702c5b46c2ebc74dfcdcb7ce899c1 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 5 Sep 2026 20:09:08 +0800 Subject: [PATCH] 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 Co-Authored-By: zhi22915 --- crates/scanner/src/data_usage_define.rs | 9 +++ crates/scanner/src/data_usage_define/tests.rs | 25 +++++++++ crates/scanner/src/scanner.rs | 1 - crates/scanner/src/scanner/activity.rs | 1 - crates/scanner/src/scanner/tests.rs | 56 +++++++++++++++++++ crates/scanner/src/scanner_io/cache.rs | 29 ++++++---- crates/scanner/src/scanner_io/io_cache.rs | 7 ++- crates/scanner/src/scanner_io/io_cycle.rs | 3 +- .../src/scanner_io/publish_gate_tests.rs | 42 ++++++++++++++ 9 files changed, 156 insertions(+), 17 deletions(-) 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";