diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index c197e3e78..7811c9fb6 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -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, +) -> 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 } diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 1793bfb50..3b476d8d1 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -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, diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 7bd1ae71b..9c73180e0 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -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) { init_ecstore_config_for_scanner_tests(); let temp_dir = tempfile::tempdir().expect("multi-pool scanner test directory should be created");