mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-09 21:56:03 +00:00
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 <qiuzgang@gmail.com>
This commit is contained in:
@@ -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::<Vec<_>>();
|
||||
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<ScannerActivitySnapshot, String>,
|
||||
) -> 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),
|
||||
|
||||
@@ -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))]);
|
||||
|
||||
Reference in New Issue
Block a user