diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 86d40fd53..5ef7118aa 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1995,10 +1995,22 @@ where if let Some((notification_system, grants)) = remote_lease_probe.as_ref() && notification_system.validate_scanner_publication_leases(grants).await.is_err() { - // A remote restart or movement flip invalidates - // the token proof; usage_store interprets this - // as a publication barrier and performs no PUT. - return Some(ScannerCycleDeferReason::DataMovement); + let remote_lease_targets = grants + .iter() + .map(|grant| { + ( + grant.host.clone(), + grant.lease.session_id.clone(), + grant.lease.movement_generation, + ) + }) + .collect::>(); + let remote_leases_valid = grants.iter().all(|grant| grant.lease.is_valid()); + return Some(scanner_remote_publication_lease_failure_defer_reason( + &remote_lease_targets, + remote_leases_valid, + probe_scanner_activity(storeapi.as_ref(), true).await, + )); } scanner_local_publication_defer_reason(storeapi.as_ref()).await } @@ -3413,6 +3425,61 @@ fn scanner_post_lease_activity_defer_reason( } } +fn scanner_remote_publication_lease_failure_defer_reason( + remote_lease_targets: &[(String, String, u64)], + remote_leases_valid: bool, + activity_after_failure: Result, +) -> ScannerCycleDeferReason { + if !remote_leases_valid { + return ScannerCycleDeferReason::DataMovement; + } + let Ok(snapshot) = activity_after_failure else { + return ScannerCycleDeferReason::ActivityBaselineUnavailable; + }; + if !scanner_activity_allows_usage_publication(&snapshot) { + return ScannerCycleDeferReason::DataMovement; + } + if scanner_publication_lease_targets_match_activity(remote_lease_targets, &snapshot) { + ScannerCycleDeferReason::ActivityBaselineUnavailable + } else { + ScannerCycleDeferReason::DataMovement + } +} + +fn scanner_publication_lease_targets_match_activity( + remote_lease_targets: &[(String, String, u64)], + activity: &ScannerActivitySnapshot, +) -> bool { + let mut expected = BTreeMap::new(); + for (host, instance_id, movement_generation) in remote_lease_targets { + if host.is_empty() + || expected + .insert(host.as_str(), (instance_id.as_str(), *movement_generation)) + .is_some() + { + return false; + } + } + + let mut observed_remote_targets = 0usize; + for (host, node_activity) in activity { + if host == LOCAL_SCANNER_ACTIVITY_NODE { + continue; + } + observed_remote_targets = observed_remote_targets.saturating_add(1); + let Some((expected_instance_id, expected_movement_generation)) = expected.get(host.as_str()) else { + return false; + }; + if node_activity.instance_id != *expected_instance_id + || node_activity.movement_generation != *expected_movement_generation + { + return false; + } + } + + observed_remote_targets == remote_lease_targets.len() +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum ScannerCyclePreCommitOutcome { RecoverCacheCycle(u64), diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 21c22305a..35805d5d7 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -9226,6 +9226,64 @@ fn post_lease_activity_proof_rejects_a_put_tail_that_finished_before_lease_acqui ); } +#[test] +fn remote_lease_validation_failure_without_movement_debt_is_activity_baseline_unavailable() { + let before = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]); + let lease_targets = scanner_activity_publication_lease_targets(&before); + let mut after = before.clone(); + after + .get_mut("node-2") + .expect("writer should be present") + .namespace_generation += 1; + + assert_eq!( + before["node-2"].movement_generation, after["node-2"].movement_generation, + "ordinary namespace writes must not be reported as movement" + ); + assert!(scanner_activity_allows_usage_publication(&after)); + assert_eq!( + scanner_remote_publication_lease_failure_defer_reason(&lease_targets, true, Ok(after)), + ScannerCycleDeferReason::ActivityBaselineUnavailable + ); +} + +#[test] +fn remote_lease_validation_failure_preserves_movement_defer_for_remote_fence_loss() { + let before = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]); + let lease_targets = scanner_activity_publication_lease_targets(&before); + + let mut movement_changed = before.clone(); + movement_changed + .get_mut("node-2") + .expect("writer should be present") + .movement_generation += 1; + assert_eq!( + scanner_remote_publication_lease_failure_defer_reason(&lease_targets, true, Ok(movement_changed)), + ScannerCycleDeferReason::DataMovement + ); + + let mut restarted = before.clone(); + restarted.get_mut("node-2").expect("writer should be present").instance_id = "epoch-b".to_string(); + assert_eq!( + scanner_remote_publication_lease_failure_defer_reason(&lease_targets, true, Ok(restarted)), + ScannerCycleDeferReason::DataMovement + ); + + let mut blocked = before.clone(); + blocked + .get_mut("node-2") + .expect("writer should be present") + .publication_blocked = true; + assert_eq!( + scanner_remote_publication_lease_failure_defer_reason(&lease_targets, true, Ok(blocked)), + ScannerCycleDeferReason::DataMovement + ); + assert_eq!( + scanner_remote_publication_lease_failure_defer_reason(&lease_targets, false, Ok(before)), + ScannerCycleDeferReason::DataMovement + ); +} + #[test] fn post_lease_activity_proof_requires_a_complete_matching_baseline() { let before = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);