From 081910e82544f14d7bee78491f62672f62b3819e Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 9 Sep 2026 05:17:03 +0800 Subject: [PATCH] feat(scanner): persist segment invalidation proof metadata (#7539) --- crates/scanner/src/data_usage_define.rs | 21 +++++++++++++++++++ crates/scanner/src/data_usage_define/tests.rs | 16 ++++++++++++++ crates/scanner/src/scanner_io.rs | 1 + crates/scanner/src/scanner_io/dirty_usage.rs | 17 +++++++++++++++ crates/scanner/src/scanner_io/io_cache.rs | 6 ++++++ crates/scanner/src/scanner_io/io_cycle.rs | 2 ++ crates/scanner/src/scanner_io/tests.rs | 10 +++++++++ 7 files changed, 73 insertions(+) diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 52193e834..95283b1cd 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -564,6 +564,18 @@ impl DataUsageCacheSource { #[serde(transparent)] pub struct DataUsageScanPlanDigest(pub [u8; 32]); +#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)] +pub struct DataUsageSegmentInvalidationProof { + #[serde(default)] + pub process_epoch: String, + #[serde(default)] + pub generation_start: u64, + #[serde(default)] + pub generation_end: u64, + #[serde(default)] + pub producer_identity_coverage_complete: bool, +} + #[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum PendingScannerHealKind { @@ -657,6 +669,11 @@ pub struct DataUsageCacheInfo { /// structural plan remains reusable across ordinary bucket writes. #[serde(default)] pub scan_execution_digest: Option, + /// Process-epoch and generation window that produced a complete set cache + /// with all known segment invalidation producers wired. This proof is + /// additive compatibility metadata; absence keeps segment reuse disabled. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub segment_invalidation_proof: Option, /// Durable bucket incarnations captured for a complete set aggregate. /// Missing or nil entries are legacy/unproven and cannot authorize /// skipping an unselected bucket in a later scoped set scan. @@ -686,6 +703,7 @@ impl Serialize for DataUsageCacheInfo { + usize::from(self.lkg_leader_epoch.is_some()) + usize::from(self.lkg_scan_plan_digest.is_some()) + usize::from(self.scan_execution_digest.is_some()) + + usize::from(self.segment_invalidation_proof.is_some()) + usize::from(!self.scan_bucket_incarnations.is_empty()); let mut state = serializer.serialize_map(Some(field_count))?; state.serialize_entry("name", &self.name)?; @@ -746,6 +764,9 @@ impl Serialize for DataUsageCacheInfo { if let Some(scan_execution_digest) = self.scan_execution_digest { state.serialize_entry("scan_execution_digest", &scan_execution_digest)?; } + if let Some(proof) = &self.segment_invalidation_proof { + state.serialize_entry("segment_invalidation_proof", proof)?; + } if !self.scan_bucket_incarnations.is_empty() { state.serialize_entry("scan_bucket_incarnations", &self.scan_bucket_incarnations)?; } diff --git a/crates/scanner/src/data_usage_define/tests.rs b/crates/scanner/src/data_usage_define/tests.rs index e1986af06..9b4713454 100644 --- a/crates/scanner/src/data_usage_define/tests.rs +++ b/crates/scanner/src/data_usage_define/tests.rs @@ -1095,6 +1095,7 @@ fn test_data_usage_cache_info_deserialize_defaults_scan_resume_after() { assert!(!decoded.snapshot_complete); assert!(decoded.scan_plan_digest.is_none()); assert!(decoded.scan_execution_digest.is_none()); + assert!(decoded.segment_invalidation_proof.is_none()); assert_eq!(decoded.cache_key_format, 0); } @@ -1183,6 +1184,12 @@ fn test_new_data_usage_cache_msgpack_round_trips_and_supports_old_reader() { snapshot_complete: true, scan_plan_digest: Some(TEST_PLAN_DIGEST), scan_execution_digest: Some(DataUsageScanPlanDigest([42; 32])), + segment_invalidation_proof: Some(DataUsageSegmentInvalidationProof { + process_epoch: "scanner-process".to_string(), + generation_start: 7, + generation_end: 9, + producer_identity_coverage_complete: true, + }), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, @@ -1212,6 +1219,15 @@ fn test_new_data_usage_cache_msgpack_round_trips_and_supports_old_reader() { assert!(current.info.snapshot_complete); assert_eq!(current.info.scan_plan_digest, Some(TEST_PLAN_DIGEST)); assert_eq!(current.info.scan_execution_digest, Some(DataUsageScanPlanDigest([42; 32]))); + assert_eq!( + current.info.segment_invalidation_proof, + Some(DataUsageSegmentInvalidationProof { + process_epoch: "scanner-process".to_string(), + generation_start: 7, + generation_end: 9, + producer_identity_coverage_complete: true, + }) + ); assert_eq!(current.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); assert_eq!(current.find("bucket").map(|entry| entry.objects), Some(3)); diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index b73365df8..c197e3e78 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -603,6 +603,7 @@ pub struct ScannerBucketScanPlan { pending_maintenance_work: Arc, cache_cycle_floor: Arc, cold_zero_walk_reuse_observed: Arc, + segment_invalidation_proof: Option, } #[derive(Clone, Default)] diff --git a/crates/scanner/src/scanner_io/dirty_usage.rs b/crates/scanner/src/scanner_io/dirty_usage.rs index 1c2af05e9..c8586c751 100644 --- a/crates/scanner/src/scanner_io/dirty_usage.rs +++ b/crates/scanner/src/scanner_io/dirty_usage.rs @@ -73,6 +73,21 @@ pub(super) struct DirtyUsageProducerEvidence { pub(super) durable_producer_identity: bool, pub(super) restart_gap_absent: bool, pub(super) generation_window_bound: bool, + pub(super) generation_start: u64, + pub(super) generation_end: u64, +} + +impl DirtyUsageProducerEvidence { + pub(super) fn segment_invalidation_proof(self) -> Option { + (self.generation_window_bound && self.producer_identity_coverage_complete).then(|| { + crate::DataUsageSegmentInvalidationProof { + process_epoch: scanner_activity_epoch().to_string(), + generation_start: self.generation_start, + generation_end: self.generation_end, + producer_identity_coverage_complete: true, + } + }) + } } /// A point-in-time view of the local dirty bucket generations. @@ -727,6 +742,8 @@ pub(super) fn dirty_usage_producer_evidence(snapshot: &DirtyUsageSnapshot) -> Di durable_producer_identity: false, restart_gap_absent: false, generation_window_bound, + generation_start: snapshot.generation, + generation_end: snapshot.generation, } } diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index c0371f85c..48fdcb433 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -173,6 +173,7 @@ impl ScannerIOCache for SetDisks { pending_maintenance_work, cache_cycle_floor, cold_zero_walk_reuse_observed, + segment_invalidation_proof, } = scan_plan; 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); @@ -243,6 +244,7 @@ impl ScannerIOCache for SetDisks { scan_plan_digest: Some(scan_plan_digest), scan_coverage_digest: Some(bucket_coverage_digest), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + segment_invalidation_proof: segment_invalidation_proof.clone(), scan_bucket_incarnations: current_bucket_incarnations.clone().unwrap_or_default(), ..Default::default() }, @@ -258,6 +260,7 @@ impl ScannerIOCache for SetDisks { cache.info.last_update = Some(now); cache.info.snapshot_complete = true; cache.info.scan_execution_digest = Some(execution_digest); + cache.info.segment_invalidation_proof = segment_invalidation_proof.clone(); cache.info.lkg_snapshot_complete = false; cache.info.lkg_next_cycle = None; cache.info.lkg_last_update = None; @@ -539,6 +542,7 @@ impl ScannerIOCache for SetDisks { lkg_last_update: old_cache.info.lkg_last_update, lkg_leader_epoch: old_cache.info.lkg_leader_epoch, lkg_scan_plan_digest: old_cache.info.lkg_scan_plan_digest, + segment_invalidation_proof: None, scan_bucket_incarnations: current_bucket_incarnations.clone().unwrap_or_default(), ..Default::default() }, @@ -1465,6 +1469,7 @@ impl ScannerIOCache for SetDisks { cache.info.last_update.get_or_insert_with(SystemTime::now); cache.info.snapshot_complete = true; cache.info.scan_execution_digest = Some(execution_digest); + cache.info.segment_invalidation_proof = segment_invalidation_proof.clone(); cache.info.lkg_snapshot_complete = false; cache.info.lkg_next_cycle = None; cache.info.lkg_last_update = None; @@ -1493,6 +1498,7 @@ impl ScannerIOCache for SetDisks { incomplete_scope.info.tier_registry_generation = Some(tier_registry_generation); incomplete_scope.info.source = Some(source); incomplete_scope.info.snapshot_complete = false; + incomplete_scope.info.segment_invalidation_proof = None; incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest); incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT; if let Err(e) = updates.send(incomplete_scope).await { diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 6bb6af130..1793bfb50 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -411,6 +411,7 @@ where let remote_dirty_usage_acknowledgements = scope_resolution.remote_dirty_usage_acknowledgements; let distributed_segment_invalidation_evidence = scope_resolution.distributed_segment_invalidation_evidence; let scan_scope = scope_resolution.scope; + let segment_invalidation_proof = dirty_usage_producer_evidence(&dirty_usage_snapshot).segment_invalidation_proof(); #[cfg(test)] if let Some(observer) = resolved_scope_observer { let _ = observer.send(scan_scope.clone()); @@ -600,6 +601,7 @@ where pending_maintenance_work: pending_maintenance_work.clone(), cache_cycle_floor: cache_cycle_floor.clone(), cold_zero_walk_reuse_observed: cold_zero_walk_reuse_observed.clone(), + segment_invalidation_proof: segment_invalidation_proof.clone(), }; // Spawn task to run the scanner let scanner_fut = tokio::spawn(async move { diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index fcbe1abd7..7bd1ae71b 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -270,6 +270,8 @@ fn complete_process_local_producer_evidence() -> DirtyUsageProducerEvidence { durable_producer_identity: false, restart_gap_absent: false, generation_window_bound: true, + generation_start: 7, + generation_end: 7, } } @@ -1557,6 +1559,12 @@ async fn set_snapshot_reuse_requires_execution_identity_and_fences_stale_writers let ctx = CancellationToken::new(); let empty_execution = DataUsageScanPlanDigest([5; 32]); + let segment_invalidation_proof = crate::DataUsageSegmentInvalidationProof { + process_epoch: scanner_activity_epoch().to_string(), + generation_start: 8, + generation_end: 8, + producer_identity_coverage_complete: true, + }; set.nsscanner_cache( ctx.clone(), ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()), @@ -1577,6 +1585,7 @@ async fn set_snapshot_reuse_requires_execution_identity_and_fences_stale_writers pending_maintenance_work: Arc::new(AtomicBool::new(false)), cache_cycle_floor: Arc::new(AtomicU64::new(8)), cold_zero_walk_reuse_observed: Arc::new(AtomicBool::new(false)), + segment_invalidation_proof: Some(segment_invalidation_proof.clone()), }, tx, 8, @@ -1586,6 +1595,7 @@ async fn set_snapshot_reuse_requires_execution_identity_and_fences_stale_writers .expect("empty set scope should replace its prior nonempty cache"); let empty = rx.try_recv().expect("empty set snapshot should be published"); assert_eq!(empty.info.scan_execution_digest, Some(empty_execution)); + assert_eq!(empty.info.segment_invalidation_proof, Some(segment_invalidation_proof)); assert!(empty.info.snapshot_complete); let root = empty.checked_flatten(DATA_USAGE_ROOT).expect("complete empty root"); assert_eq!((root.size, root.objects), (0, 0));