diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index c401c8636..fc24ec023 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -1431,9 +1431,29 @@ mod tests { .await .expect("movement writer should proceed after lease expiry") .expect("expiry writer task should not panic"); + assert!( + store.validate_scanner_publication_lease(expiring_token, 0).await.is_err(), + "an expired lease must not validate after its read guard is released" + ); assert!(!store.release_scanner_publication_lease(expiring_token).await); } + #[tokio::test] + async fn scanner_publication_lease_rejects_a_new_movement_generation() { + let store = build_store_with_ctx(Arc::new(InstanceContext::new())); + let (token, generation) = store + .acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL) + .await + .expect("an idle store should grant a publication lease"); + + assert_eq!(store.ctx.advance_data_movement_generation(), Some(1)); + assert!( + store.validate_scanner_publication_lease(token, generation).await.is_err(), + "a lease from the prior movement generation must fail closed" + ); + assert!(store.release_scanner_publication_lease(token).await); + } + #[tokio::test] async fn scanner_publication_lease_rejects_stale_generation_before_install() { let store = build_store_with_ctx(Arc::new(InstanceContext::new())); diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 1aae284ee..49af15397 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -352,6 +352,7 @@ struct MemoryConfigStore { objects: Mutex>>, revisions: Mutex>, insert_after_gets: Mutex>>, + delayed_gets: Mutex>, non_regular_objects: Mutex>, fail_put_number: Mutex>, object_not_found_put_number: Mutex>, @@ -399,6 +400,9 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore { _opts: &ObjectOptions, ) -> EcstoreResult { let key = memory_config_key(bucket, object); + if let Some(delay) = self.delayed_gets.lock().await.remove(&key) { + tokio::time::sleep(delay).await; + } let inserted_data = self.insert_after_gets.lock().await.remove(&key); let data = { let mut objects = self.objects.lock().await; @@ -3532,6 +3536,47 @@ async fn coordinator_classifies_an_expired_publication_lease() { assert!(store.put_counts.lock().await.is_empty(), "expired lease must prevent a PUT"); } +#[tokio::test] +async fn backup_sync_checks_the_lease_deadline_after_a_slow_backup_read() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_path = DATA_USAGE_OBJ_NAME_PATH.as_str(); + let backup_path = format!("{primary_path}.bkp"); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, primary_path); + let backup_key = memory_config_key(RUSTFS_META_BUCKET, &backup_path); + let primary = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + store + .objects + .lock() + .await + .insert(primary_key, serde_json::to_vec(&primary).expect("primary usage snapshot should encode")); + store + .delayed_gets + .lock() + .await + .insert(backup_key.clone(), Duration::from_millis(20)); + + // The primary read is allowed to start, but the backup read consumes the + // remaining lease window. The second deadline check must prevent a stale + // backup PUT after that window has elapsed. + let deadline = std::time::Instant::now() + .checked_add(std::time::Duration::from_millis(5)) + .expect("test deadline should support a five-millisecond window"); + let result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence( + &CancellationToken::new(), + store.clone(), + None, + Some(deadline), + None, + ) + .await; + + assert!(scanner_publication_epoch_changed( + &result.expect_err("an expired backup lease must defer publication") + )); + assert!(!store.objects.lock().await.contains_key(&backup_key)); + assert_eq!(store.put_counts.lock().await.get(&backup_key), None); +} + #[tokio::test] #[serial] async fn test_deferred_usage_save_keeps_last_real_save_metric() { @@ -4807,6 +4852,29 @@ async fn data_usage_persist_wait_aborts_after_timeout() { assert!(task.is_finished()); } +#[tokio::test(start_paused = true)] +async fn data_usage_persist_timeout_drops_owned_task_without_a_late_commit() { + let ctx = CancellationToken::new(); + let commit_started = Arc::new(AtomicBool::new(false)); + let commit_started_by_task = commit_started.clone(); + let task_ready = Arc::new(tokio::sync::Notify::new()); + let task_ready_by_task = task_ready.clone(); + let mut task = AbortOnDropHandle::new(tokio::spawn(async move { + task_ready_by_task.notify_one(); + std::future::pending::<()>().await; + commit_started_by_task.store(true, Ordering::Release); + DataUsagePersistOutcome::Saved + })); + task_ready.notified().await; + + let result = wait_for_data_usage_persist_task(&ctx, &mut task, Duration::from_secs(1)).await; + + assert!(matches!(result, DataUsagePersistTaskResult::TimedOut)); + assert!(task.is_finished(), "the timed-out persistence task must be drained before return"); + tokio::task::yield_now().await; + assert!(!commit_started.load(Ordering::Acquire), "an owned task must not commit after its timeout"); +} + #[tokio::test(start_paused = true)] async fn maintenance_feature_inspection_preserves_base_cycle_after_timeout() { let ctx = CancellationToken::new();