diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 080a86441..79a1d43a6 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -197,8 +197,8 @@ impl ECStore { registry: self.bucket_fence_registry.clone(), inner, }; - let memoized = pieces.enter(bucket); - let current = match memoized { + let registration = pieces.enter(bucket); + let current = match registration.memoized { Some(current) => current, None => match metadata_sys::get_bucket_incarnation_id_in(&self.ctx, bucket).await { Ok(current) => { @@ -210,16 +210,16 @@ impl ECStore { current } Err(err) => { - pieces.abandon(bucket); + pieces.abandon(bucket, registration.token); return Err(err); } }, }; if current != expected { - pieces.abandon(bucket); + pieces.abandon(bucket, registration.token); return Err(StorageError::BucketNotFound(bucket.to_string())); } - Ok(pieces.into_guard(bucket)) + Ok(pieces.into_guard(bucket, registration.token)) } pub(crate) async fn acquire_bucket_lifecycle_write_lock(&self, bucket: &str) -> Result { diff --git a/crates/ecstore/src/store/bucket_fence.rs b/crates/ecstore/src/store/bucket_fence.rs index b7d13a234..a705c4d5c 100644 --- a/crates/ecstore/src/store/bucket_fence.rs +++ b/crates/ecstore/src/store/bucket_fence.rs @@ -44,14 +44,51 @@ use std::collections::HashMap; use std::sync::{Arc, Mutex}; use rustfs_lock::NamespaceLockGuard; +use rustfs_lock::distributed_lock::LockLostSignal; use uuid::Uuid; #[derive(Default)] struct FenceEntry { - guards: usize, + next_token: u64, + guards: Vec, validated: Option, } +struct RegisteredGuard { + token: u64, + loss_probe: LockLossProbe, +} + +enum LockLossProbe { + Distributed(Arc), + Local, + #[cfg(test)] + Test(Arc), +} + +impl LockLossProbe { + fn from_guard(guard: &NamespaceLockGuard) -> Self { + match guard.lock_lost_signal() { + Some(signal) => Self::Distributed(signal), + None => Self::Local, + } + } + + fn is_lost(&self) -> bool { + match self { + Self::Distributed(signal) => signal.is_lost(), + Self::Local => false, + #[cfg(test)] + Self::Test(lost) => lost.load(std::sync::atomic::Ordering::SeqCst), + } + } +} + +pub(super) struct FenceRegistration { + pub(super) token: u64, + pub(super) memoized: Option, +} + /// Per-store registry tracking, per bucket, how many lifecycle read guards are /// live on this node and the incarnation id validated under that coverage. #[derive(Default)] @@ -62,11 +99,20 @@ pub(crate) struct BucketFenceRegistry { impl BucketFenceRegistry { /// Register a new live guard for `bucket` and return the memoized /// incarnation id if one is valid for the current coverage window. - fn enter(&self, bucket: &str) -> Option { + fn enter(&self, bucket: &str, loss_probe: LockLossProbe) -> FenceRegistration { let mut entries = self.entries.lock().expect("bucket fence registry poisoned"); let entry = entries.entry(bucket.to_string()).or_default(); - entry.guards += 1; - entry.validated + let token = entry.next_token; + entry.next_token = entry.next_token.wrapping_add(1); + entry.guards.push(RegisteredGuard { token, loss_probe }); + let has_lost_guard = entry.guards.iter().any(|guard| guard.loss_probe.is_lost()); + if has_lost_guard { + entry.validated = None; + } + FenceRegistration { + token, + memoized: if has_lost_guard { None } else { entry.validated }, + } } /// Memoize `incarnation` for `bucket`. Only meaningful while the caller @@ -74,23 +120,27 @@ impl BucketFenceRegistry { fn memoize(&self, bucket: &str, incarnation: Uuid) { let mut entries = self.entries.lock().expect("bucket fence registry poisoned"); if let Some(entry) = entries.get_mut(bucket) - && entry.guards > 0 + && !entry.guards.is_empty() { - entry.validated = Some(incarnation); + if entry.guards.iter().any(|guard| guard.loss_probe.is_lost()) { + entry.validated = None; + } else { + entry.validated = Some(incarnation); + } } } /// Deregister a guard. Clears the memo when the last guard leaves or when /// the leaving guard lost its lock (lost coverage means a lifecycle write /// lock may have been granted, so the memo can no longer be trusted). - fn exit(&self, bucket: &str, lock_lost: bool) { + fn exit(&self, bucket: &str, token: u64, lock_lost: bool) { let mut entries = self.entries.lock().expect("bucket fence registry poisoned"); if let Some(entry) = entries.get_mut(bucket) { - entry.guards = entry.guards.saturating_sub(1); + entry.guards.retain(|guard| guard.token != token); if lock_lost { entry.validated = None; } - if entry.guards == 0 { + if entry.guards.is_empty() { entries.remove(bucket); } } @@ -104,6 +154,7 @@ pub(crate) struct BucketIncarnationFenceGuard { inner: Option, registry: Arc, bucket: String, + token: u64, } impl BucketIncarnationFenceGuard { @@ -115,7 +166,7 @@ impl BucketIncarnationFenceGuard { impl Drop for BucketIncarnationFenceGuard { fn drop(&mut self) { let lost = self.is_lock_lost(); - self.registry.exit(&self.bucket, lost); + self.registry.exit(&self.bucket, self.token, lost); self.inner.take(); } } @@ -128,8 +179,8 @@ pub(super) struct FencePieces { impl FencePieces { /// Register the freshly acquired read lock and return the memoized /// incarnation for the coverage window, if any. - pub(super) fn enter(&self, bucket: &str) -> Option { - self.registry.enter(bucket) + pub(super) fn enter(&self, bucket: &str) -> FenceRegistration { + self.registry.enter(bucket, LockLossProbe::from_guard(&self.inner)) } pub(super) fn memoize(&self, bucket: &str, incarnation: Uuid) { @@ -140,18 +191,19 @@ impl FencePieces { self.inner.is_lock_lost() } - pub(super) fn into_guard(self, bucket: &str) -> BucketIncarnationFenceGuard { + pub(super) fn into_guard(self, bucket: &str, token: u64) -> BucketIncarnationFenceGuard { BucketIncarnationFenceGuard { inner: Some(self.inner), registry: self.registry, bucket: bucket.to_string(), + token, } } /// Abandon the acquisition (validation failed): deregister and release. - pub(super) fn abandon(self, bucket: &str) { + pub(super) fn abandon(self, bucket: &str, token: u64) { let lost = self.lock_lost(); - self.registry.exit(bucket, lost); + self.registry.exit(bucket, token, lost); drop(self.inner); } } @@ -159,57 +211,155 @@ impl FencePieces { #[cfg(test)] mod tests { use super::*; + use std::time::Duration; + + use rustfs_lock::{LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey}; fn uuid(n: u128) -> Uuid { Uuid::from_u128(n) } + fn live_probe() -> LockLossProbe { + LockLossProbe::Test(Arc::new(std::sync::atomic::AtomicBool::new(false))) + } + + fn controllable_probe() -> (LockLossProbe, Arc) { + let lost = Arc::new(std::sync::atomic::AtomicBool::new(false)); + (LockLossProbe::Test(lost.clone()), lost) + } + + fn lock_request(owner: &str) -> LockRequest { + LockRequest::new(ObjectKey::new("b", "lifecycle"), LockType::Shared, owner) + .with_acquire_timeout(Duration::from_millis(100)) + .with_ttl(Duration::from_millis(20)) + .with_refresh_interval(Duration::from_millis(50)) + } + #[test] fn memo_valid_only_while_guards_overlap() { let reg = BucketFenceRegistry::default(); - assert_eq!(reg.enter("b"), None, "first guard sees no memo"); + let first = reg.enter("b", live_probe()); + assert_eq!(first.memoized, None, "first guard sees no memo"); reg.memoize("b", uuid(1)); - assert_eq!(reg.enter("b"), Some(uuid(1)), "overlapping guard reuses memo"); - reg.exit("b", false); - reg.exit("b", false); + let second = reg.enter("b", live_probe()); + assert_eq!(second.memoized, Some(uuid(1)), "overlapping guard reuses memo"); + reg.exit("b", first.token, false); + reg.exit("b", second.token, false); // Coverage gap: all guards gone, memo must be dropped. - assert_eq!(reg.enter("b"), None, "post-gap guard must revalidate"); - reg.exit("b", false); + let third = reg.enter("b", live_probe()); + assert_eq!(third.memoized, None, "post-gap guard must revalidate"); + reg.exit("b", third.token, false); } #[test] fn lost_lock_clears_memo_but_keeps_other_guards_registered() { let reg = BucketFenceRegistry::default(); - assert_eq!(reg.enter("b"), None); + let first = reg.enter("b", live_probe()); + assert_eq!(first.memoized, None); reg.memoize("b", uuid(7)); - assert_eq!(reg.enter("b"), Some(uuid(7))); + let second = reg.enter("b", live_probe()); + assert_eq!(second.memoized, Some(uuid(7))); // First guard exits reporting a lost lock: memo cleared even though // a second guard is still live. - reg.exit("b", true); - assert_eq!(reg.enter("b"), None, "memo not trusted after a lost lock"); - reg.exit("b", false); - reg.exit("b", false); + reg.exit("b", first.token, true); + let third = reg.enter("b", live_probe()); + assert_eq!(third.memoized, None, "memo not trusted after a lost lock"); + reg.exit("b", second.token, false); + reg.exit("b", third.token, false); + } + + #[test] + fn live_lost_guard_blocks_memo_reuse_before_drop() { + let reg = BucketFenceRegistry::default(); + + let (first_probe, first_lost) = controllable_probe(); + let first = reg.enter("b", first_probe); + assert_eq!(first.memoized, None); + reg.memoize("b", uuid(7)); + + first_lost.store(true, std::sync::atomic::Ordering::SeqCst); + + let second = reg.enter("b", live_probe()); + assert_eq!(second.memoized, None, "live lost guard must force disk revalidation"); + reg.memoize("b", uuid(8)); + + let third = reg.enter("b", live_probe()); + assert_eq!(third.memoized, None, "memo remains blocked while the lost guard is live"); + + reg.exit("b", first.token, true); + reg.memoize("b", uuid(8)); + let fourth = reg.enter("b", live_probe()); + assert_eq!(fourth.memoized, Some(uuid(8)), "memo resumes after lost coverage leaves"); + + reg.exit("b", second.token, false); + reg.exit("b", third.token, false); + reg.exit("b", fourth.token, false); + } + + #[tokio::test] + async fn fence_pieces_forwards_distributed_lock_loss_to_registry() { + let registry = Arc::new(BucketFenceRegistry::default()); + let lock = NamespaceLock::new("bucket-fence-test".to_string(), Arc::new(LocalClient::new())); + let first_guard = lock + .acquire_guard(&lock_request("first")) + .await + .expect("distributed lock acquisition should not fail") + .expect("distributed lock quorum should be reached"); + let first_pieces = FencePieces { + registry: registry.clone(), + inner: first_guard, + }; + + let first = first_pieces.enter("b"); + assert_eq!(first.memoized, None); + first_pieces.memoize("b", uuid(7)); + + tokio::time::timeout(Duration::from_secs(2), first_pieces.inner.lock_lost_notified()) + .await + .expect("non-renewed distributed guard should lose its lease"); + assert!(first_pieces.lock_lost(), "test guard should observe lost refresh quorum"); + + let second_guard = lock + .acquire_guard(&lock_request("second")) + .await + .expect("second distributed lock acquisition should not fail") + .expect("second distributed lock quorum should be reached"); + let second_pieces = FencePieces { + registry: registry.clone(), + inner: second_guard, + }; + let second = second_pieces.enter("b"); + assert_eq!( + second.memoized, None, + "a live distributed guard whose signal is lost must block memo reuse" + ); + + second_pieces.abandon("b", second.token); + first_pieces.abandon("b", first.token); } #[test] fn buckets_are_isolated() { let reg = BucketFenceRegistry::default(); - assert_eq!(reg.enter("a"), None); + let first = reg.enter("a", live_probe()); + assert_eq!(first.memoized, None); reg.memoize("a", uuid(1)); - assert_eq!(reg.enter("b"), None, "memo does not leak across buckets"); - reg.exit("b", false); - reg.exit("a", false); + let second = reg.enter("b", live_probe()); + assert_eq!(second.memoized, None, "memo does not leak across buckets"); + reg.exit("b", second.token, false); + reg.exit("a", first.token, false); } #[test] fn memoize_without_live_guard_is_ignored() { let reg = BucketFenceRegistry::default(); reg.memoize("b", uuid(9)); - assert_eq!(reg.enter("b"), None); - reg.exit("b", false); + let first = reg.enter("b", live_probe()); + assert_eq!(first.memoized, None); + reg.exit("b", first.token, false); } }