From d7b6a8c10d219278319e9831b3a7a8d0692e58bb Mon Sep 17 00:00:00 2001 From: houseme Date: Thu, 10 Sep 2026 01:17:41 +0800 Subject: [PATCH] fix(scanner): avoid movement debt for stale remote lease (#7605) Classify a final remote scanner publication lease validation failure from a fresh activity snapshot instead of treating every validation error as data movement. If the peer session and movement generation still match the granted leases and publication is otherwise allowed, keep the publication rejected as an activity-baseline miss without creating pause backlog movement debt. Preserve DataMovement for expired leases, peer restarts, movement generation changes, and active publication blocks. Tests cover namespace-only validation invalidation and the remote fence-loss cases that must still defer as DataMovement. Co-authored-by: zhi22915 --- crates/scanner/src/scanner.rs | 75 +++++++++++++++++++++++++++-- crates/scanner/src/scanner/tests.rs | 58 ++++++++++++++++++++++ 2 files changed, 129 insertions(+), 4 deletions(-) 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))]);