diff --git a/crates/data-usage/src/data_usage.rs b/crates/data-usage/src/data_usage.rs index 68a4d9822..93f325a1c 100644 --- a/crates/data-usage/src/data_usage.rs +++ b/crates/data-usage/src/data_usage.rs @@ -593,6 +593,18 @@ pub struct DataUsageSnapshotIdentity { pub scanner_epoch: Option, } +#[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, Debug, Default, Serialize, Deserialize, PartialEq, Eq)] pub struct DataUsageSnapshotSetState { pub pool_index: u64, @@ -607,6 +619,8 @@ pub struct DataUsageSnapshotSetState { pub complete: bool, #[serde(default)] pub tombstone: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub segment_invalidation_proof: Option, } impl DataUsageInfo { @@ -3073,6 +3087,7 @@ mod tests { scan_plan_digest: Some([1; 32]), complete: false, tombstone: false, + segment_invalidation_proof: None, }]; assert!(observed_data_usage_is_newer(&partial, &authoritative)); } @@ -3095,6 +3110,7 @@ mod tests { scan_plan_digest: Some([1; 32]), complete: true, tombstone: false, + segment_invalidation_proof: None, }, DataUsageSnapshotSetState { pool_index: 1, @@ -3104,6 +3120,7 @@ mod tests { scan_plan_digest: Some([2; 32]), complete: false, tombstone: false, + segment_invalidation_proof: None, }, ], ..Default::default() @@ -3113,6 +3130,41 @@ mod tests { assert!(partial.is_valid_partial_snapshot()); } + #[test] + fn set_state_segment_invalidation_proof_is_additive() { + #[derive(Deserialize)] + struct LegacySetState { + pool_index: u64, + set_index: u64, + complete: bool, + } + + let proof = DataUsageSegmentInvalidationProof { + process_epoch: "scanner-process".to_string(), + generation_start: 3, + generation_end: 5, + producer_identity_coverage_complete: true, + }; + let state = DataUsageSnapshotSetState { + pool_index: 1, + set_index: 2, + scanner_cycle: Some(9), + scanner_epoch: Some(4), + scan_plan_digest: Some([7; 32]), + complete: true, + tombstone: false, + segment_invalidation_proof: Some(proof.clone()), + }; + let encoded = rmp_serde::to_vec_named(&state).expect("set state should encode with additive proof"); + let legacy: LegacySetState = rmp_serde::from_slice(&encoded).expect("legacy readers should ignore proof metadata"); + assert_eq!(legacy.pool_index, 1); + assert_eq!(legacy.set_index, 2); + assert!(legacy.complete); + + let decoded: DataUsageSnapshotSetState = rmp_serde::from_slice(&encoded).expect("new readers should restore proof"); + assert_eq!(decoded.segment_invalidation_proof, Some(proof)); + } + #[test] fn completeness_marker_requires_a_snapshot_timestamp() { let untimestamped = DataUsageInfo { diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 95283b1cd..3634dfd35 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -28,10 +28,10 @@ use metrics::{counter, describe_counter, describe_histogram, histogram}; use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS; pub use rustfs_data_usage::{ AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME, - DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME, - PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeReconciliationEntry, - SizeReconciliationScope, SizeSummary, TierAccountingProof, TierStats, UNKNOWN_TIER, UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP, - UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP, UnknownTierStats, hash_path, prefix_usage_in_cache, + DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSegmentInvalidationProof, DataUsageSnapshotSetState, + LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, + SizeReconciliationEntry, SizeReconciliationScope, SizeSummary, TierAccountingProof, TierStats, UNKNOWN_TIER, + UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP, UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP, UnknownTierStats, hash_path, prefix_usage_in_cache, }; use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; @@ -564,18 +564,6 @@ 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 { diff --git a/crates/scanner/src/scanner_io/cache.rs b/crates/scanner/src/scanner_io/cache.rs index 328b91ce9..2eb88089f 100644 --- a/crates/scanner/src/scanner_io/cache.rs +++ b/crates/scanner/src/scanner_io/cache.rs @@ -464,6 +464,7 @@ pub(super) fn completed_usage_candidate( scan_plan_digest: Some(result.info.scan_plan_digest?.0), complete: true, tombstone: false, + segment_invalidation_proof: result.info.segment_invalidation_proof.clone(), }) }) .collect::>>()?; @@ -642,13 +643,14 @@ pub(super) fn observational_data_usage_info( let current_snapshot = current.is_some(); let selected = current.or(lkg); if let Some(selected) = selected { - let (cycle, epoch, digest, last_update, complete) = if current_snapshot { + let (cycle, epoch, digest, last_update, complete, segment_invalidation_proof) = if current_snapshot { ( Some(selected.info.next_cycle), Some(selected.info.leader_epoch), selected.info.scan_plan_digest.map(|digest| digest.0), selected.info.last_update, true, + selected.info.segment_invalidation_proof.clone(), ) } else { ( @@ -657,6 +659,7 @@ pub(super) fn observational_data_usage_info( selected.info.lkg_scan_plan_digest.map(|digest| digest.0), selected.info.lkg_last_update, false, + None, ) }; set_states.push(DataUsageSnapshotSetState { @@ -667,6 +670,7 @@ pub(super) fn observational_data_usage_info( scan_plan_digest: digest, complete, tombstone: false, + segment_invalidation_proof, }); usable.push((selected, last_update)); } else { @@ -678,6 +682,7 @@ pub(super) fn observational_data_usage_info( scan_plan_digest: Some(expected_plan_digest.0), complete: false, tombstone: false, + segment_invalidation_proof: None, }); } } diff --git a/crates/scanner/src/scanner_io/publish_gate_tests.rs b/crates/scanner/src/scanner_io/publish_gate_tests.rs index 6274f27fa..8a662fcdf 100644 --- a/crates/scanner/src/scanner_io/publish_gate_tests.rs +++ b/crates/scanner/src/scanner_io/publish_gate_tests.rs @@ -13,7 +13,7 @@ // limitations under the License. use super::*; -use crate::data_usage_define::{UNKNOWN_TIER, UnknownTierStats, hash_path}; +use crate::data_usage_define::{DataUsageSegmentInvalidationProof, UNKNOWN_TIER, UnknownTierStats, hash_path}; use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccountingProof}; const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]); @@ -176,6 +176,26 @@ fn completed_data_usage_info_rejects_duplicate_bucket_inventory() { assert!(completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none()); } +#[test] +fn completed_data_usage_info_carries_segment_invalidation_proof_to_set_state() { + let source = DataUsageCacheSource::new(0, 0); + let proof = DataUsageSegmentInvalidationProof { + process_epoch: "scanner-process".to_string(), + generation_start: 5, + generation_end: 8, + producer_identity_coverage_complete: true, + }; + let mut set = completed_root_cache("bucket", 2, 10, source); + set.info.segment_invalidation_proof = Some(proof.clone()); + + let (usage, _) = + completed_usage_for_scope(&[set], &HashSet::from([source]), &["bucket".to_string()], &[], true, false, false) + .expect("complete set should publish root usage"); + + assert_eq!(usage.usage_snapshot_set_states.len(), 1); + assert_eq!(usage.usage_snapshot_set_states[0].segment_invalidation_proof, Some(proof)); +} + #[test] fn completed_data_usage_info_rejects_extra_or_detached_bucket_data() { let buckets = vec!["bucket".to_string()]; @@ -350,6 +370,7 @@ fn set_membership_add_remove_uses_generation_and_tombstone() { scan_plan_digest: Some(TEST_PLAN_DIGEST.0), complete: false, tombstone: true, + segment_invalidation_proof: None, }; let encoded = serde_json::to_vec(&state).expect("set state should serialize"); let decoded: DataUsageSnapshotSetState = serde_json::from_slice(&encoded).expect("set state should deserialize"); @@ -371,6 +392,7 @@ fn set_membership_add_remove_uses_generation_and_tombstone() { scan_plan_digest: Some(TEST_PLAN_DIGEST.0), complete: true, tombstone: false, + segment_invalidation_proof: None, }, state, ], diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 9c73180e0..ee3d1cc63 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -1686,6 +1686,7 @@ fn complete_usage_baseline( scan_plan_digest: Some(scan_plan_digest.0), complete: true, tombstone: false, + segment_invalidation_proof: None, }], ..Default::default() };