From 35ef4ea7dd8cb6227f093001644cc75ae053f143 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 9 Sep 2026 01:25:55 +0800 Subject: [PATCH] feat(scanner): persist segment invalidation proof metadata Add compatible set cache metadata for complete scanner segment invalidation producer evidence and carry it through completed set snapshot publication. Keep incomplete snapshots unproven so segment reuse activation remains fail-closed until baseline admission consumes the durable proof. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- 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));