mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-08 21:25:59 +00:00
feat(scanner): admit durable segment proof evidence
Validate complete set snapshot segment invalidation proof metadata against the current dirty usage generation window and scanner process epoch before clearing durable producer and restart-gap activation blockers. Production segment reuse remains gated by the explicit activation flag. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -539,6 +539,47 @@ fn scanner_segment_reuse_activation_preflight_for_cycle(
|
||||
})
|
||||
}
|
||||
|
||||
fn scanner_durable_segment_invalidation_evidence(
|
||||
dirty_usage_snapshot: &DirtyUsageSnapshot,
|
||||
results: &[DataUsageCache],
|
||||
expected_sources: &HashSet<DataUsageCacheSource>,
|
||||
) -> DirtyUsageProducerEvidence {
|
||||
let mut evidence = dirty_usage_producer_evidence(dirty_usage_snapshot);
|
||||
if !evidence.generation_window_bound
|
||||
|| !evidence.producer_identity_coverage_complete
|
||||
|| !scanner_results_form_complete_snapshot(results, expected_sources)
|
||||
{
|
||||
return evidence;
|
||||
}
|
||||
|
||||
let mut covered_sources = HashSet::with_capacity(expected_sources.len());
|
||||
let all_sets_proved = results.iter().all(|result| {
|
||||
let Some(source) = result.info.source else {
|
||||
return false;
|
||||
};
|
||||
expected_sources.contains(&source)
|
||||
&& covered_sources.insert(source)
|
||||
&& scanner_segment_invalidation_proof_matches(result.info.segment_invalidation_proof.as_ref(), &evidence)
|
||||
});
|
||||
if all_sets_proved && covered_sources.len() == expected_sources.len() {
|
||||
evidence.durable_producer_identity = true;
|
||||
evidence.restart_gap_absent = true;
|
||||
}
|
||||
evidence
|
||||
}
|
||||
|
||||
fn scanner_segment_invalidation_proof_matches(
|
||||
proof: Option<&crate::DataUsageSegmentInvalidationProof>,
|
||||
evidence: &DirtyUsageProducerEvidence,
|
||||
) -> bool {
|
||||
proof.is_some_and(|proof| {
|
||||
proof.process_epoch == scanner_activity_epoch()
|
||||
&& proof.generation_start == evidence.generation_start
|
||||
&& proof.generation_end == evidence.generation_end
|
||||
&& proof.producer_identity_coverage_complete
|
||||
})
|
||||
}
|
||||
|
||||
fn scanner_segment_reuse_activated() -> bool {
|
||||
scanner_segment_reuse_activation_preflight().scanner_segment_reuse_activated
|
||||
}
|
||||
|
||||
@@ -715,7 +715,7 @@ where
|
||||
);
|
||||
let segment_reuse_activation_preflight = scanner_segment_reuse_activation_preflight_for_cycle(
|
||||
&dirty_usage_snapshot,
|
||||
dirty_usage_producer_evidence(&dirty_usage_snapshot),
|
||||
scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &results, &expected_sources),
|
||||
distributed,
|
||||
distributed_segment_invalidation_evidence,
|
||||
cold_zero_walk_oracle,
|
||||
|
||||
@@ -237,6 +237,49 @@ fn scanner_segment_reuse_activation_preflight_for_cycle_skips_distributed_blocke
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial]
|
||||
fn scanner_durable_segment_invalidation_evidence_requires_matching_complete_set_proofs() {
|
||||
use crate::segment_invalidation::SegmentInvalidationProducerIdentity;
|
||||
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
record_dirty_usage_bucket_from_producers("photos", SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION);
|
||||
let dirty_usage_snapshot = snapshot_dirty_usage_buckets(&[bucket_info("photos")], dirty_usage_generation());
|
||||
let process_proof = dirty_usage_producer_evidence(&dirty_usage_snapshot)
|
||||
.segment_invalidation_proof()
|
||||
.expect("complete process-local producer coverage should produce proof metadata");
|
||||
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(0, 1)]);
|
||||
let results = vec![
|
||||
complete_set_cache_with_segment_proof(DataUsageCacheSource::new(0, 0), process_proof.clone()),
|
||||
complete_set_cache_with_segment_proof(DataUsageCacheSource::new(0, 1), process_proof.clone()),
|
||||
];
|
||||
|
||||
let durable_evidence = scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &results, &expected_sources);
|
||||
|
||||
assert!(durable_evidence.producer_identity_coverage_complete);
|
||||
assert!(durable_evidence.durable_producer_identity);
|
||||
assert!(durable_evidence.restart_gap_absent);
|
||||
|
||||
let mut stale_epoch = results.clone();
|
||||
stale_epoch[0]
|
||||
.info
|
||||
.segment_invalidation_proof
|
||||
.as_mut()
|
||||
.expect("proof fixture should exist")
|
||||
.process_epoch = "stale-process".to_string();
|
||||
let stale_evidence = scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &stale_epoch, &expected_sources);
|
||||
assert!(stale_evidence.producer_identity_coverage_complete);
|
||||
assert!(!stale_evidence.durable_producer_identity);
|
||||
assert!(!stale_evidence.restart_gap_absent);
|
||||
|
||||
record_dirty_usage_bucket("videos");
|
||||
let changed_evidence = scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &results, &expected_sources);
|
||||
assert!(!changed_evidence.producer_identity_coverage_complete);
|
||||
assert!(!changed_evidence.durable_producer_identity);
|
||||
assert!(!changed_evidence.restart_gap_absent);
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_cycle_result_returns_segment_reuse_activation_preflight() {
|
||||
let proof = ScannerSegmentReuseActivationProof {
|
||||
@@ -275,6 +318,26 @@ fn complete_process_local_producer_evidence() -> DirtyUsageProducerEvidence {
|
||||
}
|
||||
}
|
||||
|
||||
fn complete_set_cache_with_segment_proof(
|
||||
source: DataUsageCacheSource,
|
||||
proof: crate::DataUsageSegmentInvalidationProof,
|
||||
) -> DataUsageCache {
|
||||
DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: DATA_USAGE_ROOT.to_string(),
|
||||
next_cycle: 7,
|
||||
last_update: Some(SystemTime::UNIX_EPOCH),
|
||||
leader_epoch: 11,
|
||||
source: Some(source),
|
||||
snapshot_complete: true,
|
||||
scan_plan_digest: Some(DataUsageScanPlanDigest([3; 32])),
|
||||
segment_invalidation_proof: Some(proof),
|
||||
..Default::default()
|
||||
},
|
||||
cache: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
async fn setup_two_pool_scanner_store() -> (tempfile::TempDir, Arc<ECStore>) {
|
||||
init_ecstore_config_for_scanner_tests();
|
||||
let temp_dir = tempfile::tempdir().expect("multi-pool scanner test directory should be created");
|
||||
|
||||
Reference in New Issue
Block a user