mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 19:55:37 +00:00
fix(scanner): fence set snapshot reuse with the scan work proof
Prevent same-cycle set publication from replacing freshly scanned maintenance results with an older Normal aggregate. Recognize uniform completed maintenance baselines when planning later ordinary dirty-bucket work. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -183,6 +183,18 @@ fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<
|
||||
return None;
|
||||
}
|
||||
|
||||
// Completed maintenance also covers ordinary usage. Keep its exact stored
|
||||
// proof for cache reuse, and reject mixtures of different set work proofs.
|
||||
let baseline_plan_digest = DataUsageScanPlanDigest(baseline.usage_snapshot_set_states.first()?.scan_plan_digest?);
|
||||
if ![
|
||||
proof.scan_plan_digest,
|
||||
scanner_bucket_work_digest(proof.scan_plan_digest, HealScanMode::Normal, true),
|
||||
scanner_bucket_work_digest(proof.scan_plan_digest, HealScanMode::Deep, true),
|
||||
]
|
||||
.contains(&baseline_plan_digest)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
let mut states = HashSet::with_capacity(baseline.usage_snapshot_set_states.len());
|
||||
for state in &baseline.usage_snapshot_set_states {
|
||||
let source = DataUsageCacheSource::new(usize::try_from(state.pool_index).ok()?, usize::try_from(state.set_index).ok()?);
|
||||
@@ -192,13 +204,13 @@ fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<
|
||||
|| state.tombstone
|
||||
|| state.scanner_epoch != Some(proof.leader_epoch)
|
||||
|| state.scanner_cycle.is_none_or(|cycle| cycle > proof.want_cycle)
|
||||
|| state.scan_plan_digest != Some(proof.scan_plan_digest.0)
|
||||
|| state.scan_plan_digest != Some(baseline_plan_digest.0)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
}
|
||||
|
||||
(states == *proof.expected_sources).then_some(proof.scan_plan_digest)
|
||||
(states == *proof.expected_sources).then_some(baseline_plan_digest)
|
||||
}
|
||||
|
||||
fn scoped_scan_scope_from_dirty_buckets(
|
||||
|
||||
@@ -126,6 +126,7 @@ impl ScannerIOCache for SetDisks {
|
||||
cache_cycle_floor,
|
||||
} = scan_plan;
|
||||
let bucket_work_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan);
|
||||
let scan_plan_digest = scanner_bucket_work_digest(scan_plan_digest, scan_mode, requires_full_scan);
|
||||
let pool_label = self.pool_index.to_string();
|
||||
let set_label = self.set_index.to_string();
|
||||
|
||||
|
||||
@@ -275,10 +275,11 @@ where
|
||||
}
|
||||
bucket_plan_complete &= buckets_by_source.keys().copied().collect::<HashSet<_>>() == *expected_sources;
|
||||
bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source);
|
||||
let scan_plan_digest =
|
||||
let structural_scan_plan_digest =
|
||||
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
|
||||
let bucket_coverage_digest =
|
||||
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before));
|
||||
let scan_plan_digest = scanner_bucket_work_digest(structural_scan_plan_digest, scan_mode, requires_full_scan);
|
||||
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
|
||||
let scan_scope = resolve_scanner_bucket_scan_scope(
|
||||
store,
|
||||
@@ -290,7 +291,7 @@ where
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch,
|
||||
want_cycle,
|
||||
scan_plan_digest,
|
||||
scan_plan_digest: structural_scan_plan_digest,
|
||||
},
|
||||
activity_before: &activity_before,
|
||||
dirty_usage_snapshot: &dirty_usage_snapshot,
|
||||
@@ -431,7 +432,7 @@ where
|
||||
buckets: set_buckets,
|
||||
all_buckets: Arc::clone(&all_buckets),
|
||||
scope: scan_scope.clone(),
|
||||
digest: scan_plan_digest,
|
||||
digest: structural_scan_plan_digest,
|
||||
bucket_coverage_digest,
|
||||
requires_full_scan,
|
||||
leader_epoch,
|
||||
|
||||
@@ -1137,6 +1137,36 @@ fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
|
||||
assert_ne!(scope.baseline_scan_plan_digest, Some(baseline_scan_plan_digest));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_scan_baseline_work_proof_requires_uniform_known_set_identity() {
|
||||
let source = DataUsageCacheSource::new(1, 2);
|
||||
let second_source = DataUsageCacheSource::new(1, 3);
|
||||
let sources = HashSet::from([source, second_source]);
|
||||
let structural = DataUsageScanPlanDigest([9; 32]);
|
||||
let full = scanner_bucket_work_digest(structural, HealScanMode::Normal, true);
|
||||
let deep = scanner_bucket_work_digest(structural, HealScanMode::Deep, true);
|
||||
let encoded = complete_usage_baseline(source, full, 7, 11);
|
||||
let baseline: DataUsageInfo = serde_json::from_slice(&encoded).expect("baseline should decode");
|
||||
for (second_plan, expected) in [(full, Some(full)), (deep, None), (DataUsageScanPlanDigest([8; 32]), None)] {
|
||||
let mut candidate = baseline.clone();
|
||||
let mut second = candidate.usage_snapshot_set_states[0].clone();
|
||||
second.set_index = 3;
|
||||
second.scan_plan_digest = Some(second_plan.0);
|
||||
candidate.usage_snapshot_set_states.push(second);
|
||||
let data = Bytes::from(serde_json::to_vec(&candidate).expect("candidate should encode"));
|
||||
assert_eq!(
|
||||
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
|
||||
data: Some(&data),
|
||||
expected_sources: &sources,
|
||||
leader_epoch: 11,
|
||||
want_cycle: 8,
|
||||
scan_plan_digest: structural,
|
||||
}),
|
||||
expected
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn peer_dirty_usage_snapshot(
|
||||
instance_id: &str,
|
||||
generation: u64,
|
||||
|
||||
Reference in New Issue
Block a user