diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 1268bae15..887cc3d79 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -273,6 +273,7 @@ pub struct ScannerBucketScanPlan { digest: DataUsageScanPlanDigest, /// Includes mutation generations even when the set planner uses a structural digest. bucket_coverage_digest: DataUsageScanPlanDigest, + requires_full_scan: bool, leader_epoch: u64, tier_registry_generation: u64, /// Epoch captured once for the whole scanner cycle. `None` is retained @@ -338,6 +339,29 @@ fn scanner_bucket_inventory_is_complete( covered.len() == inventory.len() } +// Bind known work requirements before both local and remote cache admission. +// Matching requirements remain reusable for the same intent; this is not a +// new deadline or a durable generation for newly due maintenance. +fn scanner_bucket_work_digest( + scan_plan_digest: DataUsageScanPlanDigest, + scan_mode: HealScanMode, + requires_full_scan: bool, +) -> DataUsageScanPlanDigest { + if scan_mode == HealScanMode::Normal && !requires_full_scan { + return scan_plan_digest; + } + let mut hasher = Sha256::new(); + hasher.update(b"scanner-bucket-work-v1"); + hasher.update(scan_plan_digest.0); + hasher.update([match scan_mode { + HealScanMode::Unknown => 0, + HealScanMode::Normal => 1, + HealScanMode::Deep => 2, + }]); + hasher.update([u8::from(requires_full_scan || scan_mode == HealScanMode::Deep)]); + DataUsageScanPlanDigest(hasher.finalize().into()) +} + fn scanner_bucket_cache_digest( scan_plan_digest: DataUsageScanPlanDigest, dirty_generation: Option, diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index 6085558bd..fb141d517 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -116,6 +116,7 @@ impl ScannerIOCache for SetDisks { scope, digest: scan_plan_digest, bucket_coverage_digest, + requires_full_scan, leader_epoch, tier_registry_generation, publication_epoch, @@ -124,6 +125,7 @@ impl ScannerIOCache for SetDisks { pending_maintenance_work, cache_cycle_floor, } = scan_plan; + let bucket_work_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan); let pool_label = self.pool_index.to_string(); let set_label = self.set_index.to_string(); @@ -635,7 +637,7 @@ impl ScannerIOCache for SetDisks { let cache_name = path_join_buf(&[&bucket.name, DATA_USAGE_CACHE_NAME]); let bucket_scan_plan_digest = - scanner_bucket_cache_digest(bucket_coverage_digest, dirty_usage_buckets_clone.get(&bucket.name).copied()); + scanner_bucket_cache_digest(bucket_work_digest, dirty_usage_buckets_clone.get(&bucket.name).copied()); if let Some(server_epoch) = remote_server_epoch { let request_sequence = remote_session_sequence; diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index cd9723d73..b06ccc8d4 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -433,6 +433,7 @@ where scope: scan_scope.clone(), digest: scan_plan_digest, bucket_coverage_digest, + requires_full_scan, leader_epoch, tier_registry_generation, publication_epoch, diff --git a/crates/scanner/src/scanner_io/publish_gate_tests.rs b/crates/scanner/src/scanner_io/publish_gate_tests.rs index f6aecdc15..254b561cc 100644 --- a/crates/scanner/src/scanner_io/publish_gate_tests.rs +++ b/crates/scanner/src/scanner_io/publish_gate_tests.rs @@ -1330,6 +1330,56 @@ fn dirty_bucket_cache_digest_changes_with_generation() { assert!(!cache_snapshot_is_current(&cache, "photos", source, 11, 0, second)); } +#[test] +fn scoped_scan_bucket_work_proof_fences_same_cycle_cache() { + let source = DataUsageCacheSource::new(0, 0); + let structural_plan = DataUsageScanPlanDigest([9; 32]); + let normal_plan = scanner_bucket_work_digest(structural_plan, HealScanMode::Normal, false); + assert_eq!(normal_plan, structural_plan, "ordinary work keeps the existing digest contract"); + for (scan_mode, full) in [(HealScanMode::Deep, false), (HealScanMode::Normal, true)] { + let requested_plan = scanner_bucket_work_digest(structural_plan, scan_mode, full); + assert_ne!(requested_plan, normal_plan); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "cold".to_string(), + next_cycle: 7, + leader_epoch: 11, + last_update: Some(SystemTime::UNIX_EPOCH), + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(normal_plan), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.replace("cold", "", DataUsageEntry::default()); + assert!(cache_snapshot_is_current(&cache, "cold", source, 7, 11, normal_plan)); + assert!(matches!( + current_cache_root_or_prepare(&mut cache, "cold", source, 7, 11, requested_plan, true), + DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + .. + } + )); + assert!(cache.cache.is_empty(), "different work requirements must enter a fresh walk"); + cache.replace("cold", "", DataUsageEntry::default()); + cache.info.snapshot_complete = true; + cache.info.last_update = Some(SystemTime::UNIX_EPOCH); + assert!( + matches!( + current_cache_root_or_prepare(&mut cache, "cold", source, 7, 11, requested_plan, true), + DataUsageCacheScanState::Current(_) + ), + "completed matching work may satisfy the same intent retry" + ); + } + assert_eq!( + scanner_bucket_work_digest(structural_plan, HealScanMode::Deep, false), + scanner_bucket_work_digest(structural_plan, HealScanMode::Deep, true) + ); +} + #[test] fn scanner_cache_lock_resource_is_scoped_to_cache_source() { let cache_name = "photos/.usage-cache.bin"; diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index b93a2e344..66a908bd2 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -415,6 +415,106 @@ async fn scoped_scan_production_entry_preserves_deep_and_full_maintenance_work() clear_dirty_usage_buckets_for_tests(); } +#[tokio::test] +#[serial] +async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure() { + for (scan_mode, requires_full_scan) in [(HealScanMode::Deep, false), (HealScanMode::Normal, true)] { + let (_temp_dir, store) = setup_two_pool_scanner_store().await; + clear_dirty_usage_buckets_for_tests(); + for bucket in ["hot-bucket", "cold-bucket"] { + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut reader = ScannerPutObjReader::from_vec(b"initial".to_vec()); + store.pools[0].disk_set[0] + .put_object(bucket, "initial", &mut reader, &ScannerObjectOptions::default()) + .await + .expect("initial object should persist"); + } + let ctx = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()); + let (updates, receiver) = mpsc::channel(1); + drop(receiver); + let failed = tokio::time::timeout( + Duration::from_secs(30), + nsscanner_with_storage_status_scoped( + store.as_ref(), + ScannerCycleRequest { + ctx, + budget, + updates, + want_cycle: 7, + leader_epoch: 11, + scan_mode: HealScanMode::Normal, + scan_scope: ScannerBucketScanScope::default(), + persisted_usage_baseline: None, + requires_full_scan: false, + resolved_scope_observer: None, + }, + ), + ) + .await + .expect("normal scan should finish") + .expect_err("root delivery must fail after bucket cache persistence"); + assert!(failed.to_string().contains("receiver closed"), "{failed}"); + let cache_name = path_join_buf(&["cold-bucket", DATA_USAGE_CACHE_NAME]); + let mut cached = DataUsageCache::default(); + cached + .load(store.pools[0].disk_set[0].clone(), &cache_name) + .await + .expect("normal bucket cache should have committed"); + assert!(cached.info.snapshot_complete); + assert_eq!(cached.info.next_cycle, 7); + assert_eq!( + cached + .checked_flatten("cold-bucket") + .expect("cached root should be valid") + .objects, + 1 + ); + + let mut reader = ScannerPutObjReader::from_vec(b"maintenance".to_vec()); + store.pools[0].disk_set[0] + .put_object("cold-bucket", "new", &mut reader, &ScannerObjectOptions::default()) + .await + .expect("new cold object should persist"); + record_dirty_usage_bucket("hot-bucket"); + let ctx = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()); + let (updates, mut receiver) = mpsc::channel(1); + let result = tokio::time::timeout( + Duration::from_secs(30), + nsscanner_with_storage_status_scoped( + store.as_ref(), + ScannerCycleRequest { + ctx, + budget, + updates, + want_cycle: 7, + leader_epoch: 11, + scan_mode, + scan_scope: ScannerBucketScanScope::default(), + persisted_usage_baseline: None, + requires_full_scan, + resolved_scope_observer: None, + }, + ), + ) + .await + .expect("maintenance scan should finish") + .expect("maintenance scan should succeed"); + assert_eq!(result.status, ScannerCycleStatus::Complete); + let snapshot = receiver.recv().await.expect("maintenance snapshot should be published"); + assert_eq!(snapshot.scanner_cycle, Some(7)); + assert_eq!( + snapshot.buckets_usage["cold-bucket"].objects_count, 2, + "{scan_mode:?}/full={requires_full_scan} must not replay the same-cycle Normal root" + ); + clear_dirty_usage_buckets_for_tests(); + } +} + #[tokio::test] async fn data_usage_publish_fails_when_receiver_is_closed() { let (updates, receiver) = mpsc::channel(1);