From 09c6412593f50aa5ce46d67041182f308ea78cc2 Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 8 Sep 2026 00:55:23 +0800 Subject: [PATCH] fix(ecstore): bind scanner leases with drain-safe fixtures --- crates/ecstore/src/runtime/instance.rs | 62 +++++---- crates/ecstore/src/store/mod.rs | 111 ++++++++++++++- crates/ecstore/src/store/peer.rs | 181 +++++++++++++++++++++++++ 3 files changed, 324 insertions(+), 30 deletions(-) diff --git a/crates/ecstore/src/runtime/instance.rs b/crates/ecstore/src/runtime/instance.rs index fd4007698..16a974af7 100644 --- a/crates/ecstore/src/runtime/instance.rs +++ b/crates/ecstore/src/runtime/instance.rs @@ -75,6 +75,7 @@ pub(crate) const SCANNER_PUBLICATION_LEASE_TTL: std::time::Duration = std::time: pub(crate) struct ScannerPublicationLeaseEntry { pub(crate) expires_at: Instant, pub(crate) movement_generation: u64, + pub(crate) namespace_generation: u64, pub(crate) _operation_guard: OwnedRwLockReadGuard<()>, } @@ -305,6 +306,7 @@ impl InstanceContext { token: Uuid, expires_at: Instant, movement_generation: u64, + namespace_generation: u64, operation_guard: OwnedRwLockReadGuard<()>, ) -> bool { let mut leases = self.scanner_publication_leases.lock().await; @@ -316,6 +318,7 @@ impl InstanceContext { ScannerPublicationLeaseEntry { expires_at, movement_generation, + namespace_generation, _operation_guard: operation_guard, }, ); @@ -326,39 +329,21 @@ impl InstanceContext { self.scanner_publication_leases.lock().await.remove(&token).is_some() } - /// Check a lease token while the caller holds the movement read guard. - /// - /// The token table is deliberately process-owned and non-persistent: a - /// restarted instance has no entries from the previous process, so an old - /// coordinator proof cannot become valid again merely because the - /// movement generation counter restarted at zero. - pub(crate) async fn scanner_publication_lease_is_active(&self, token: Uuid) -> bool { + /// Return both generations from the same live lease while the caller holds + /// the movement read guard. Namespace commits do not take that guard, so + /// the caller must compare the saved namespace generation after this await. + /// The process-owned table rejects tokens from a prior instance or expiry. + pub(crate) async fn scanner_publication_lease_generations(&self, token: Uuid) -> Option<(u64, u64)> { let mut leases = self.scanner_publication_leases.lock().await; let now = Instant::now(); - let Some(expires_at) = leases.get(&token).map(|entry| entry.expires_at) else { - return false; - }; - if expires_at <= now { - leases.remove(&token); - return false; - } - true - } - - /// Return the generation bound to a live lease. The lease entry owns the - /// movement read guard, so a successful lookup remains valid for the - /// caller's guard-protected operation; expiry is still fail-closed. - pub(crate) async fn scanner_publication_lease_generation(&self, token: Uuid) -> Option { - let mut leases = self.scanner_publication_leases.lock().await; - let now = Instant::now(); - let (expires_at, movement_generation) = leases + let (expires_at, movement_generation, namespace_generation) = leases .get(&token) - .map(|entry| (entry.expires_at, entry.movement_generation))?; + .map(|entry| (entry.expires_at, entry.movement_generation, entry.namespace_generation))?; if expires_at <= now { leases.remove(&token); return None; } - Some(movement_generation) + Some((movement_generation, namespace_generation)) } pub(crate) async fn expire_scanner_publication_lease(&self, token: Uuid, expires_at: Instant) { @@ -516,6 +501,11 @@ impl InstanceContext { .store(SCANNER_PUBLICATION_STATE_UNKNOWN, Ordering::Release); } + #[cfg(test)] + pub(crate) fn set_namespace_commit_generation_for_test(&self, generation: u64) { + self.namespace_commit_generation.store(generation, Ordering::Release); + } + #[cfg(test)] pub(crate) fn set_data_movement_generation_for_test(&self, generation: u64) { self.data_movement_generation.store(generation, Ordering::Release); @@ -862,6 +852,26 @@ mod tests { } } + #[tokio::test(start_paused = true)] + async fn scanner_lease_generations_remain_bound_until_expiry() { + let ctx = Arc::new(InstanceContext::new()); + let token = Uuid::new_v4(); + let gate = ctx.data_movement_operation_gate(); + let permit = gate.clone().read_owned().await; + assert!( + ctx.install_scanner_publication_lease(token, Instant::now() + SCANNER_PUBLICATION_LEASE_TTL, 7, 11, permit) + .await + ); + drop(ctx.begin_namespace_commit()); + assert_eq!(ctx.namespace_commit_generation(), 2); + assert_eq!(ctx.scanner_publication_lease_generations(token).await, Some((7, 11))); + assert!(gate.clone().try_write_owned().is_err(), "lookup must retain the stored permit"); + tokio::time::advance(SCANNER_PUBLICATION_LEASE_TTL).await; + assert_eq!(ctx.scanner_publication_lease_generations(token).await, None); + assert!(!ctx.remove_scanner_publication_lease(token).await); + assert!(gate.try_write_owned().is_ok(), "expiry releases the stored permit"); + } + // The SetupType inputs must derive the exact (is_erasure, // is_dist_erasure, is_erasure_sd) triples that the original three // process-global erasure bools produced via update_erasure_type(). diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 202002903..69cde8b7f 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -1089,11 +1089,15 @@ impl ECStore { return Err(Error::other("scanner publication lease TTL is not supported")); } + // Bind the original activity generation across the asynchronous checks; + // a completed namespace commit must never refresh an existing proof. + let namespace_generation = self.scanner_namespace_mutation_generation(); let operation_gate = self.ctx.data_movement_operation_gate(); let operation_guard = operation_gate.read_owned().await; if self.ctx.data_movement_generation_exhausted() || self.ctx.data_movement_operation_epoch_exhausted() || self.ctx.data_movement_generation() != expected_generation + || namespace_generation == u64::MAX { return Err(Error::other("scanner publication lease generation is stale")); } @@ -1101,11 +1105,15 @@ impl ECStore { return Err(Error::other("scanner publication lease is blocked by data movement")); } + if self.scanner_namespace_mutation_generation() != namespace_generation { + return Err(Error::other("scanner publication lease generation is stale")); + } + let token = Uuid::new_v4(); let expires_at = tokio::time::Instant::now() + ttl; if !self .ctx - .install_scanner_publication_lease(token, expires_at, expected_generation, operation_guard) + .install_scanner_publication_lease(token, expires_at, expected_generation, namespace_generation, operation_guard) .await { return Err(Error::other("scanner publication lease capacity is exhausted")); @@ -1139,8 +1147,16 @@ impl ECStore { if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() { return Err(Error::other("scanner publication lease is blocked by data movement")); } - if !self.ctx.scanner_publication_lease_is_active(token).await { + let Some((lease_generation, lease_namespace_generation)) = self.ctx.scanner_publication_lease_generations(token).await + else { return Err(Error::other("scanner publication lease is unknown or expired")); + }; + let namespace_generation = self.scanner_namespace_mutation_generation(); + if lease_generation != self.ctx.data_movement_generation() + || lease_namespace_generation != namespace_generation + || namespace_generation == u64::MAX + { + return Err(Error::other("scanner publication lease generation is stale")); } Ok(()) } @@ -1159,10 +1175,15 @@ impl ECStore { if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() { return Err(Error::other("scanner publication lease is blocked by data movement")); } - let Some(lease_generation) = self.ctx.scanner_publication_lease_generation(token).await else { + let Some((lease_generation, lease_namespace_generation)) = self.ctx.scanner_publication_lease_generations(token).await + else { return Err(Error::other("scanner publication lease is unknown or expired")); }; - if lease_generation != self.ctx.data_movement_generation() { + let namespace_generation = self.scanner_namespace_mutation_generation(); + if lease_generation != self.ctx.data_movement_generation() + || lease_namespace_generation != namespace_generation + || namespace_generation == u64::MAX + { return Err(Error::other("scanner publication lease generation is stale")); } Ok(operation_guard) @@ -2556,6 +2577,88 @@ mod tests { assert!(error.to_string().contains("generation is stale")); } + #[tokio::test] + async fn scanner_publication_lease_rejects_namespace_change_during_admission() { + let ctx = Arc::new(InstanceContext::new()); + let store = build_store_with_ctx(ctx.clone()); + let snapshot_blocker = store.rebalance_meta.write().await; + let mut acquire = + Box::pin(store.acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)); + assert!(futures::poll!(&mut acquire).is_pending(), "acquisition reaches the blocked snapshot"); + assert!(ctx.data_movement_operation_gate().try_write_owned().is_err()); + drop(ctx.begin_namespace_commit()); + assert_eq!(ctx.namespace_commit_generation(), 2); + assert!(!ctx.namespace_commits_pending()); + drop(snapshot_blocker); + let error = tokio::time::timeout(Duration::from_secs(1), acquire) + .await + .expect("snapshot admission must finish") + .expect_err("acquisition must preserve its original namespace generation"); + assert_eq!(error.to_string(), "Io error: scanner publication lease generation is stale"); + assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok(), "no lease was installed"); + } + + #[tokio::test] + async fn scanner_publication_lease_rechecks_namespace_after_validate_snapshot() { + let ctx = Arc::new(InstanceContext::new()); + let store = build_store_with_ctx(ctx.clone()); + let (token, generation) = store + .acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL) + .await + .expect("current lease"); + let snapshot_blocker = store.rebalance_meta.write().await; + let mut validate = Box::pin(store.validate_scanner_publication_lease(token, generation)); + assert!(futures::poll!(&mut validate).is_pending(), "target guard queues its first snapshot read"); + let mut next_writer = Box::pin(store.rebalance_meta.write()); + assert!(futures::poll!(&mut next_writer).is_pending()); + drop(snapshot_blocker); + // Fair lock order admits the first read, then this queued writer, then + // validate's second snapshot. The first target check has already passed. + assert!(futures::poll!(&mut validate).is_pending(), "validate reaches its second snapshot"); + let next_writer = next_writer.await; + drop(ctx.begin_namespace_commit()); + assert_eq!(ctx.namespace_commit_generation(), 2); + assert!(!ctx.namespace_commits_pending()); + drop(next_writer); + let error = tokio::time::timeout(Duration::from_secs(1), validate) + .await + .expect("validation must finish without nesting movement read locks") + .expect_err("the final snapshot must reject a completed namespace commit"); + assert_eq!(error.to_string(), "Io error: scanner publication lease generation is stale"); + assert!( + ctx.data_movement_operation_gate().try_write_owned().is_err(), + "stale lookup retains the lease permit" + ); + assert!(store.release_scanner_publication_lease(token).await); + assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok()); + } + + #[tokio::test] + async fn scanner_publication_lease_rejects_namespace_generation_exhaustion() { + let ctx = Arc::new(InstanceContext::new()); + let store = build_store_with_ctx(ctx.clone()); + let (token, generation) = store + .acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL) + .await + .expect("current lease"); + ctx.set_namespace_commit_generation_for_test(u64::MAX); + assert_eq!(store.scanner_namespace_mutation_generation(), u64::MAX); + let acquire = store + .acquire_scanner_publication_lease(generation, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL) + .await + .expect_err("exhausted namespace cannot become a new lease baseline"); + assert_eq!(acquire.to_string(), "Io error: scanner publication lease generation is stale"); + let validate = store.validate_scanner_publication_lease(token, generation).await; + let target = store.acquire_scanner_publication_lease_guard(token).await; + assert!(validate.is_err() && target.is_err(), "exhaustion rejects both old-token entrances"); + assert!( + ctx.data_movement_operation_gate().try_write_owned().is_err(), + "rejection must retain the old permit" + ); + assert!(store.release_scanner_publication_lease(token).await); + assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok()); + } + #[tokio::test(start_paused = true)] async fn scanner_publication_commit_scope_owns_permit_until_terminal_drain() { let store = build_store_with_ctx(Arc::new(InstanceContext::new())); diff --git a/crates/ecstore/src/store/peer.rs b/crates/ecstore/src/store/peer.rs index 06b33d2d5..8cbadeda8 100644 --- a/crates/ecstore/src/store/peer.rs +++ b/crates/ecstore/src/store/peer.rs @@ -1275,6 +1275,187 @@ mod tests { ); } + #[tokio::test] + async fn renewed_namespace_lease_allows_real_metadata_publication() { + use futures::FutureExt; + use std::time::Duration; + use tokio::time::timeout; + + let ctx = Arc::new(InstanceContext::new()); + let store = super::super::tests::build_store_with_ctx(ctx.clone()); + let root = tempfile::tempdir().expect("target root"); + let ttl = crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL; + let user_volume = "target-bucket"; + let metadata_volume = ".rustfs.sys/tmp"; + let setup = std::panic::AssertUnwindSafe(timeout(ttl / 2, async { + let disk = target_disk(&ctx, root.path(), Uuid::new_v4()).await; + let user_new = target_file_info("destination", Uuid::new_v4(), b"completed-user-write"); + let mut metadata_old = target_file_info("destination", Uuid::new_v4(), b"old-metadata"); + let mut metadata_new = target_file_info("destination", Uuid::new_v4(), b"metadata-from-renewed-scan"); + metadata_old.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixed old time")); + metadata_new.mod_time = metadata_old.mod_time.map(|old| old + time::Duration::seconds(1)); + seed_target(&disk, user_volume, "staged", user_new.clone()).await; + let metadata_before = seed_target(&disk, metadata_volume, "destination", metadata_old).await; + seed_target(&disk, metadata_volume, "staged", metadata_new.clone()).await; + (disk, user_new, metadata_new, metadata_before) + })) + .catch_unwind() + .await; + let (disk, user_new, metadata_new, metadata_before) = match setup { + Ok(Ok(setup)) => setup, + Ok(Err(error)) => { + let retained = root.keep(); + panic!("fixture initialization must finish before lease acquisition: {error}; retained={retained:?}"); + } + Err(panic) => { + let retained = root.keep(); + eprintln!("fixture setup panicked; retained={retained:?}"); + std::panic::resume_unwind(panic); + } + }; + let disk_ref = disk.endpoint().to_string(); + let movement_generation = ctx.data_movement_generation(); + let namespace_generation = store.scanner_namespace_mutation_generation(); + let read_options = crate::disk::ReadOptions { + read_data: true, + ..Default::default() + }; + let mut tokens = Vec::new(); + let result = std::panic::AssertUnwindSafe(timeout(ttl / 2, async { + let (old_token, _) = store + .acquire_scanner_publication_lease(movement_generation, ttl) + .await + .expect("original lease"); + tokens.push(old_token); + store + .rename_local_data(&disk_ref, (user_volume, "staged"), &user_new, (user_volume, "destination"), None) + .await + .expect("ordinary namespace write"); + while ctx.namespace_commits_pending() { + tokio::task::yield_now().await; + } + assert_eq!(ctx.namespace_commit_generation(), 2); + assert_eq!(ctx.data_movement_generation(), movement_generation); + let user_latest = disk + .read_version(user_volume, user_volume, "destination", "", &read_options) + .await + .expect("latest ordinary committed object"); + assert_eq!(user_latest.version_id, user_new.version_id); + assert_eq!(user_latest.data, user_new.data); + assert!( + store + .validate_scanner_publication_lease(old_token, movement_generation) + .await + .is_err() + ); + assert_eq!( + ctx.scanner_publication_lease_generations(old_token).await, + Some((movement_generation, namespace_generation)), + "rejection must not refresh or discard the old lease" + ); + let (fresh_token, fresh_generation) = store + .acquire_scanner_publication_lease(movement_generation, ttl) + .await + .expect("a new scan can acquire a current lease after the completed write"); + tokens.push(fresh_token); + store + .validate_scanner_publication_lease(fresh_token, fresh_generation) + .await + .expect("fresh lease validates"); + store + .rename_local_data( + &disk_ref, + (metadata_volume, "staged"), + &metadata_new, + (metadata_volume, "destination"), + Some(fresh_token), + ) + .await + .expect("fresh lease authorizes real internal metadata publication"); + let latest = disk + .read_version(metadata_volume, metadata_volume, "destination", "", &read_options) + .await + .expect("latest published metadata"); + let raw = tokio::fs::read(root.path().join(metadata_volume).join("destination/xl.meta")) + .await + .expect("raw published metadata"); + assert_eq!(latest.version_id, metadata_new.version_id); + assert_eq!(latest.data, metadata_new.data); + assert_ne!(raw, metadata_before); + assert!(!ctx.namespace_commits_pending()); + assert_eq!( + ctx.namespace_commit_generation(), + 2, + "internal metadata does not mutate the user namespace" + ); + assert_eq!(ctx.data_movement_generation(), movement_generation); + })) + .catch_unwind() + .await; + // A table release does not drain an independently owned metadata call. + // Collect every release outcome before checking either ownership chain. + let cleanup = std::panic::AssertUnwindSafe(async { + let mut releases = Vec::new(); + for token in tokens { + releases.push( + std::panic::AssertUnwindSafe(timeout(Duration::from_secs(5), store.release_scanner_publication_lease(token))) + .catch_unwind() + .await, + ); + } + let movement_guard = timeout(Duration::from_secs(5), ctx.data_movement_operation_gate().write_owned()).await; + let namespace_drained = timeout(Duration::from_secs(5), async { + while ctx.namespace_commits_pending() { + tokio::task::yield_now().await; + } + }) + .await; + let drained = releases.iter().all(|release| matches!(release, Ok(Ok(_)))) + && movement_guard.is_ok() + && namespace_drained.is_ok(); + #[cfg(not(windows))] + let drained = { + use crate::disk::os::prepared_publication_test_hooks as hooks; + + let mut keys_drained = true; + for volume in [user_volume, metadata_volume] { + // The physical lease uses the actual IO object directory, + // including descriptor-rooted aliases, not its xl.meta file. + let key_drained = match disk.get_object_path_for_io_if_local(volume, "destination") { + Some(Ok(path)) => timeout(Duration::from_secs(5), hooks::drain_namespace_key(&path)) + .await + .is_ok(), + _ => false, + }; + keys_drained &= key_drained; + } + drained && keys_drained + }; + drop(movement_guard); + if let Some(panic) = releases.into_iter().find_map(|release| release.err()) { + std::panic::resume_unwind(panic); + } + drained + }) + .catch_unwind() + .await; + // Generic cancelled IO cannot be proved drained by a namespace key. + // Retain on any failed observation, even if best-effort cleanup succeeds. + if !matches!(&result, Ok(Ok(()))) || !matches!(&cleanup, Ok(true)) { + let retained = root.keep(); + eprintln!("fixture observations or physical cleanup incomplete; retained={retained:?}"); + } + match result { + Ok(result) => result.expect("publication must finish before the original lease can expire"), + Err(panic) => std::panic::resume_unwind(panic), + } + match cleanup { + Ok(drained) => assert!(drained, "fixture physical cleanup must finish before deleting its root"), + Err(panic) => std::panic::resume_unwind(panic), + } + assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok()); + } + #[cfg(not(windows))] #[tokio::test] #[serial_test::serial]