diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index ae6379cd0..d4dfddb1a 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1662,14 +1662,17 @@ async fn run_data_scanner_cycle_with_budget( .is_some_and(|(_, grants)| grants.iter().any(|grant| !grant.lease.is_valid())); if let Some((notification_system, grants)) = remote_publication_leases.take() { let release_result = notification_system.release_scanner_publication_leases(grants).await; - if lease_expired || release_result.is_err() { + let lease_release_failed = release_result.is_err(); + if lease_expired || lease_release_failed { // A lease that expired or could not be released is never treated // as a successful authoritative publication. The peer may have // admitted movement immediately after the lease ended. usage_persist_outcome = if usage_persist_outcome == DataUsagePersistOutcome::Failed { DataUsagePersistOutcome::Failed + } else if lease_release_failed { + DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseReleaseFailed) } else { - DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded) }; } } diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 0c02caec0..88932cab9 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -3393,6 +3393,44 @@ async fn coordinator_does_not_put_after_remote_generation_flip() { assert_eq!(store.put_counts.lock().await.get(&key), None); } +#[tokio::test] +async fn coordinator_classifies_an_expired_publication_lease() { + let store = Arc::new(MemoryConfigStore::default()); + let (sender, receiver) = mpsc::channel(1); + sender + .send(complete_usage_with_bucket_count( + Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)), + 1, + )) + .await + .expect("usage snapshot should enqueue"); + drop(sender); + + let expired = std::time::Instant::now() + .checked_sub(std::time::Duration::from_secs(1)) + .expect("test instant should support a one-second subtraction"); + let outcome = + store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe_for_publication_epoch_and_lease_fence( + CancellationToken::new(), + store.clone(), + receiver, + None, + Some(DataUsagePersistBaseline { + data: None, + revision: DataUsageCacheRevision::Missing, + }), + ScannerPublicationFence::new(None, Some(expired), None), + || async { false }, + ) + .await; + + assert_eq!( + outcome, + DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded) + ); + assert!(store.put_counts.lock().await.is_empty(), "expired lease must prevent a PUT"); +} + #[tokio::test] #[serial] async fn test_deferred_usage_save_keeps_last_real_save_metric() { @@ -4446,6 +4484,7 @@ fn scanner_cycle_cache_floor_stays_pending_during_deferred_usage_publication() { ScannerCycleDeferReason::ActivityBaselineUnavailable, ScannerCycleDeferReason::PublicationLeaseBudgetExceeded, ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded, + ScannerCycleDeferReason::PublicationLeaseReleaseFailed, ] { let deferred = DataUsagePersistOutcome::Deferred(reason); assert_eq!( @@ -4632,6 +4671,10 @@ fn scanner_publication_lease_budget_has_a_strict_ttl_boundary() { ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded.as_str(), "publication_lease_deadline_exceeded" ); + assert_eq!( + ScannerCycleDeferReason::PublicationLeaseReleaseFailed.as_str(), + "publication_lease_release_failed" + ); } #[tokio::test] diff --git a/crates/scanner/src/scanner/usage_store.rs b/crates/scanner/src/scanner/usage_store.rs index a32e0cecf..6edda880c 100644 --- a/crates/scanner/src/scanner/usage_store.rs +++ b/crates/scanner/src/scanner/usage_store.rs @@ -225,7 +225,7 @@ where data_usage_info.scanner_epoch = Some(leader_epoch); } if remote_lease_expired(remote_lease_deadline) { - outcome = DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable); + outcome = DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded); break 'updates; } if let Some(expected_epoch) = expected_publication_epoch @@ -497,7 +497,7 @@ where break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement); } if remote_lease_expired(remote_lease_deadline) { - break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable); + break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded); } let done_save = Metrics::time(Metric::SaveUsage); @@ -510,7 +510,7 @@ where }; if remote_lease_expired(remote_lease_deadline) { done_save(); - break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable); + break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded); } save_config_shared_with_preconditions_and_lease_fence( storeapi.clone(), diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 9036a9fee..e1abbf474 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -582,6 +582,9 @@ pub(crate) enum ScannerCycleDeferReason { /// operation. This can occur even when the configured budget fits the /// nominal TTL because lease acquisition consumed part of the window. PublicationLeaseDeadlineExceeded, + /// A remote lease could not be released after the persistence attempt. + /// Keep the cycle deferred because the peer may still admit movement. + PublicationLeaseReleaseFailed, } impl ScannerCycleDeferReason { @@ -591,6 +594,7 @@ impl ScannerCycleDeferReason { Self::DataMovement => "data_movement", Self::PublicationLeaseBudgetExceeded => "publication_lease_budget_exceeded", Self::PublicationLeaseDeadlineExceeded => "publication_lease_deadline_exceeded", + Self::PublicationLeaseReleaseFailed => "publication_lease_release_failed", } } }