mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-09 05:36:24 +00:00
feat(scanner): replay segment proof from root snapshots
Carry segment invalidation proof metadata into root snapshot set states so a complete published baseline can replay the durable producer evidence recorded by each set cache. Keep the field additive for older readers and leave incomplete or LKG set states unproven. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -593,6 +593,18 @@ pub struct DataUsageSnapshotIdentity {
|
||||
pub scanner_epoch: Option<u64>,
|
||||
}
|
||||
|
||||
#[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<DataUsageSegmentInvalidationProof>,
|
||||
}
|
||||
|
||||
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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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::<Option<Vec<_>>>()?;
|
||||
@@ -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,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
],
|
||||
|
||||
@@ -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()
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user