diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 887cc3d79..d1cc7a073 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -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( diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index fb141d517..4289d1efd 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -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(); diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index b06ccc8d4..1d3ce57b5 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -275,10 +275,11 @@ where } bucket_plan_complete &= buckets_by_source.keys().copied().collect::>() == *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, diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 66a908bd2..0391c4928 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -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,