From ddf197ba57f96416938b3668d663915b43fa20be Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 8 Jul 2026 06:01:45 +0800 Subject: [PATCH] fix(lock): refresh object write locks with a heartbeat (#4388) fix(lock): renew distributed locks via heartbeat; wire server refresh (backlog#899) A held distributed write lock had a fixed 30s TTL and was never renewed, so any operation exceeding it got its per-node lease reclaimed and stolen by a contender, causing split-brain writes. This implements Phase 0 + Phase 1 of the #899 design (Phase 2 abort-on-loss is deferred and tracked in code comments). Phase 0 (server refresh wiring, P0): - node_service handle_refresh was a no-op stub that parsed args then always returned success=true. Extract a testable `refresh_lock` free function that actually delegates to the node's lock backend. Not-found maps to success=false with no error_info, so RemoteClient::refresh keeps its Ok(resp.success) semantics and yields a real not_found signal to the coordinator heartbeat. Without this, client heartbeats were silently no-ops. Phase 1 (client heartbeat + safe interval + observability): - DistributedLockGuard spawns a heartbeat that refreshes every per-client lease on a derived interval; interval is derived without Duration::clamp (entries<=1 or interval>=ttl => no spawn), fixing the sub-second-ttl panic. - Add LockLostSignal: declare the lock lost when not_found exceeds entries.len() - refresh_quorum; RPC errors are not counted (absorbed by the ttl > interval margin). Expose is_lock_lost()/lock_lost() for observers. - disarm(), release(), and Drop now abort the heartbeat before releasing so no refresh races the unlock (refresh only extends, never creates, a lease). - Reclaim path stays behaviorally unchanged but now warns with owner/resource/ lease age and records a metric (#698 scavenger preserved). - Add DEFAULT_LOCK_REFRESH_INTERVAL, LockRequest.refresh_interval (serde default for RPC back-compat) + builder, and lock lifecycle metrics. Open questions adopt the design's documented defaults (marked TODO in code): refresh not-found -> Ok(false)/error_info=None; lost-quorum base entries.len(); DEFAULT_LOCK_REFRESH_INTERVAL=10s. Tests: heartbeat keepalive/quorum-loss/jitter/boundary/disarm (lock crate), server refresh delegation (rustfs), and end-to-end survives-past-ttl plus crashed-owner-reclaim regressions (namespace). --- .../e2e_test/src/reliant/grpc_lock_client.rs | 1 + .../ecstore/src/cluster/rpc/remote_locker.rs | 1 + crates/io-metrics/src/lib.rs | 5 +- crates/io-metrics/src/process_lock_metrics.rs | 29 ++ crates/lock/src/client/local.rs | 25 +- crates/lock/src/distributed_lock.rs | 465 +++++++++++++++++- crates/lock/src/fast_lock/mod.rs | 7 + crates/lock/src/namespace/tests.rs | 91 ++++ crates/lock/src/types.rs | 14 + rustfs/src/storage/rpc/node_service/lock.rs | 90 +++- 10 files changed, 714 insertions(+), 14 deletions(-) diff --git a/crates/e2e_test/src/reliant/grpc_lock_client.rs b/crates/e2e_test/src/reliant/grpc_lock_client.rs index bb1279fab..e1acd50d6 100644 --- a/crates/e2e_test/src/reliant/grpc_lock_client.rs +++ b/crates/e2e_test/src/reliant/grpc_lock_client.rs @@ -63,6 +63,7 @@ impl GrpcLockClient { priority: LockPriority::Normal, deadlock_detection: false, suppress_contention_logs: false, + refresh_interval: None, } } diff --git a/crates/ecstore/src/cluster/rpc/remote_locker.rs b/crates/ecstore/src/cluster/rpc/remote_locker.rs index ffdfa9834..05e946766 100644 --- a/crates/ecstore/src/cluster/rpc/remote_locker.rs +++ b/crates/ecstore/src/cluster/rpc/remote_locker.rs @@ -73,6 +73,7 @@ impl RemoteClient { priority: LockPriority::Normal, deadlock_detection: false, suppress_contention_logs: false, + refresh_interval: None, } } diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index 4148b1b7e..fa84c3d7b 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -157,8 +157,9 @@ pub use lock_metrics::{ }; pub use process_lock_metrics::{ - ProcessLockSnapshot, ProcessPlatformSnapshot, record_read_lock_held_acquire, record_read_lock_held_release, - record_write_lock_held_acquire, record_write_lock_held_release, snapshot_process_lock_counts, + ProcessLockEventSnapshot, ProcessLockSnapshot, ProcessPlatformSnapshot, record_lock_reclaimed, + record_lock_refresh_quorum_lost, record_read_lock_held_acquire, record_read_lock_held_release, + record_write_lock_held_acquire, record_write_lock_held_release, snapshot_process_lock_counts, snapshot_process_lock_events, snapshot_process_platform_stats, }; pub use s3_api_metrics::{init_s3_metrics, record_s3_op}; diff --git a/crates/io-metrics/src/process_lock_metrics.rs b/crates/io-metrics/src/process_lock_metrics.rs index d5b2b146a..e60386b36 100644 --- a/crates/io-metrics/src/process_lock_metrics.rs +++ b/crates/io-metrics/src/process_lock_metrics.rs @@ -28,6 +28,10 @@ use std::process::Command; static READ_LOCKS_HELD: AtomicU64 = AtomicU64::new(0); static WRITE_LOCKS_HELD: AtomicU64 = AtomicU64::new(0); +/// Monotonic count of expired lock guards forcibly reclaimed by a contender (#899 observability). +static LOCKS_RECLAIMED: AtomicU64 = AtomicU64::new(0); +/// Monotonic count of guard heartbeats that observed a refresh-quorum loss (#899 observability). +static LOCK_REFRESH_QUORUM_LOST: AtomicU64 = AtomicU64::new(0); #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct ProcessLockSnapshot { @@ -35,6 +39,13 @@ pub struct ProcessLockSnapshot { pub write_locks_held: u64, } +/// Monotonic lock lifecycle counters (distinct from the held gauges above). +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct ProcessLockEventSnapshot { + pub locks_reclaimed: u64, + pub lock_refresh_quorum_lost: u64, +} + #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct ProcessPlatformSnapshot { pub io_rchar_bytes: Option, @@ -66,6 +77,16 @@ pub fn record_write_lock_held_release() { let _ = WRITE_LOCKS_HELD.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |value| Some(value.saturating_sub(1))); } +#[inline(always)] +pub fn record_lock_reclaimed() { + LOCKS_RECLAIMED.fetch_add(1, Ordering::Relaxed); +} + +#[inline(always)] +pub fn record_lock_refresh_quorum_lost() { + LOCK_REFRESH_QUORUM_LOST.fetch_add(1, Ordering::Relaxed); +} + #[inline(always)] pub fn snapshot_process_lock_counts() -> ProcessLockSnapshot { ProcessLockSnapshot { @@ -74,6 +95,14 @@ pub fn snapshot_process_lock_counts() -> ProcessLockSnapshot { } } +#[inline(always)] +pub fn snapshot_process_lock_events() -> ProcessLockEventSnapshot { + ProcessLockEventSnapshot { + locks_reclaimed: LOCKS_RECLAIMED.load(Ordering::Relaxed), + lock_refresh_quorum_lost: LOCK_REFRESH_QUORUM_LOST.load(Ordering::Relaxed), + } +} + #[inline] pub fn snapshot_process_platform_stats() -> ProcessPlatformSnapshot { platform::snapshot() diff --git a/crates/lock/src/client/local.rs b/crates/lock/src/client/local.rs index ea6299ca5..2c65f2040 100644 --- a/crates/lock/src/client/local.rs +++ b/crates/lock/src/client/local.rs @@ -42,15 +42,18 @@ struct LocalGuardEntry { guard: FastLockGuard, expires_at: SystemTime, ttl: Duration, + /// Owner recorded at acquire time; used only for reclaim diagnostics (#899). + owner: String, } impl LocalGuardEntry { - fn new(guard: FastLockGuard, ttl: Duration) -> Self { + fn new(guard: FastLockGuard, ttl: Duration, owner: String) -> Self { let now = SystemTime::now(); Self { guard, expires_at: now + ttl, ttl, + owner, } } @@ -136,6 +139,24 @@ impl LocalClient { }; for mut entry in expired_entries { + // An expired entry whose owner never refreshed it (a dead coordinator, #698) is + // reclaimed so a live contender can re-form quorum. With guard heartbeats in place + // (#899) a live owner keeps its entry from expiring, so reaching here means the + // lease genuinely lapsed. Surface it for observability; the reclaim decision itself + // is unchanged. + let since_last_refresh = entry + .expires_at + .checked_sub(entry.ttl) + .and_then(|last_refresh| SystemTime::now().duration_since(last_refresh).ok()) + .unwrap_or(entry.ttl); + tracing::warn!( + owner = %entry.owner, + resource = %resource, + ttl_ms = entry.ttl.as_millis() as u64, + since_last_refresh_ms = since_last_refresh.as_millis() as u64, + "reclaiming expired lock guard whose lease was not refreshed" + ); + rustfs_io_metrics::record_lock_reclaimed(); let _ = entry.guard.release(); reclaimed = reclaimed.saturating_add(1); } @@ -175,7 +196,7 @@ impl LockClient for LocalClient { { let shard = self.get_shard(&lock_id); let mut guards = shard.write().await; - guards.insert(lock_id.clone(), LocalGuardEntry::new(guard, request.ttl)); + guards.insert(lock_id.clone(), LocalGuardEntry::new(guard, request.ttl, request.owner.clone())); } let lock_info = LockInfo { diff --git a/crates/lock/src/distributed_lock.rs b/crates/lock/src/distributed_lock.rs index a056b3b8a..c17a7bac6 100644 --- a/crates/lock/src/distributed_lock.rs +++ b/crates/lock/src/distributed_lock.rs @@ -20,11 +20,14 @@ use crate::{ }; use futures::future::join_all; use rustfs_io_metrics::{ - record_read_lock_held_acquire, record_read_lock_held_release, record_write_lock_held_acquire, record_write_lock_held_release, + record_lock_refresh_quorum_lost, record_read_lock_held_acquire, record_read_lock_held_release, + record_write_lock_held_acquire, record_write_lock_held_release, }; use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; -use tokio::task::JoinSet; +use tokio::sync::Notify; +use tokio::task::{JoinHandle, JoinSet}; use tracing::{debug, warn}; use uuid::Uuid; @@ -72,6 +75,24 @@ fn is_unrecoverable_quorum_error(error: &str) -> bool { .is_some_and(|prefix| prefix.eq_ignore_ascii_case(UNRECOVERABLE_QUORUM_FAILURE_PREFIX)) } +/// Derive a safe heartbeat interval for a freshly acquired guard. +/// +/// Deliberately avoids `Duration::clamp` (which asserts `min <= max` and would panic for +/// sub-second ttls where `ttl - 1s` underflows to zero). Returns `None` for degenerate cases so +/// no heartbeat is spawned: +/// - `entries_len <= 1`: single/degenerate path (matches a local, non-distributed lock), +/// - `interval.is_zero()` or `interval >= ttl`: too small a ttl / too large an interval to renew. +fn derive_refresh_interval(entries_len: usize, ttl: Duration, injected: Option) -> Option { + if entries_len <= 1 { + return None; + } + let interval = injected.unwrap_or(ttl / 3); + if interval.is_zero() || interval >= ttl { + return None; + } + Some(interval) +} + fn classify_lock_failure(error: &str) -> LockAcquireFailureKind { if is_unrecoverable_quorum_error(error) { return LockAcquireFailureKind::UnrecoverableQuorum; @@ -103,6 +124,35 @@ fn should_warn_lock_failure(error: &str) -> bool { ) } +/// Observability channel for lock-loss: set once when a guard's heartbeat detects it has +/// lost refresh quorum. Phase 1 only signals (warn + metric + this flag); Phase 2 (backlog#899) +/// will `select!` on `notified()` inside long write operations to abort and avoid split-brain. +#[derive(Debug, Default)] +pub struct LockLostSignal { + lost: AtomicBool, + notify: Notify, +} + +impl LockLostSignal { + fn mark_lost(&self) { + self.lost.store(true, Ordering::SeqCst); + self.notify.notify_waiters(); + } + + /// Whether refresh quorum has been lost for the associated guard. + pub fn is_lost(&self) -> bool { + self.lost.load(Ordering::SeqCst) + } + + /// Resolves once the lock is declared lost (immediately if already lost). + pub async fn notified(&self) { + if self.is_lost() { + return; + } + self.notify.notified().await; + } +} + /// A RAII guard for distributed locks that releases the lock asynchronously when dropped. #[derive(Debug)] pub struct DistributedLockGuard { @@ -115,6 +165,11 @@ pub struct DistributedLockGuard { lock_type: LockType, /// If true, Drop will not try to release (used if user manually released). disarmed: bool, + /// Background heartbeat task that periodically refreshes the per-client leases. + /// `None` for degenerate/single-client paths that do not need renewal. + refresh_task: Option>, + /// Lock-loss signal, shared with the heartbeat task. + lock_lost: Arc, } impl DistributedLockGuard { @@ -123,13 +178,95 @@ impl DistributedLockGuard { /// - `lock_id` is the id returned to the caller (`lock_id()`). /// - `entries` is the full list of underlying (LockId, client) pairs /// that should be released when this guard is dropped. - pub(crate) fn new(lock_id: LockId, entries: Vec<(LockId, Arc)>, lock_type: LockType) -> Self { + /// - `refresh_interval`: `Some` spawns a heartbeat that refreshes every entry on that + /// cadence; `None` (degenerate/single-client, or interval derived away) spawns nothing. + /// - `refresh_quorum`: minimum refreshes that must keep succeeding; if `not_found` exceeds + /// `entries.len() - refresh_quorum` the guard is declared lost. + /// - `owner`/`resource`: diagnostics only. + pub(crate) fn new( + lock_id: LockId, + entries: Vec<(LockId, Arc)>, + lock_type: LockType, + refresh_interval: Option, + refresh_quorum: usize, + owner: String, + resource: ObjectKey, + ) -> Self { record_lock_held_acquire(lock_type); + let lock_lost = Arc::new(LockLostSignal::default()); + let refresh_task = refresh_interval.and_then(|interval| { + // Only spawn when a tokio runtime is available; guard construction off-runtime + // (e.g. some tests) simply skips the heartbeat. + tokio::runtime::Handle::try_current().ok().map(|handle| { + let entries = entries.clone(); + let lock_lost = lock_lost.clone(); + handle.spawn(Self::run_heartbeat(entries, interval, refresh_quorum, lock_lost, owner, resource)) + }) + }); Self { lock_id, entries, lock_type, disarmed: false, + refresh_task, + lock_lost, + } + } + + /// Heartbeat loop: every `interval`, refresh all entries and classify the outcomes. + /// `Ok(true)` = refreshed, `Ok(false)` = not_found, `Err` = RPC jitter (ignored, absorbed by + /// the ttl > interval margin and retried next tick). Declares the lock lost when + /// `not_found > entries.len() - refresh_quorum`. Phase 1 does not release on loss: the guarded + /// operation is still running, and tearing the lock down here would only widen the window. + async fn run_heartbeat( + entries: Vec<(LockId, Arc)>, + interval: Duration, + refresh_quorum: usize, + lock_lost: Arc, + owner: String, + resource: ObjectKey, + ) { + let tolerable_not_found = entries.len().saturating_sub(refresh_quorum); + let mut ticker = tokio::time::interval(interval); + ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + // The first tick fires immediately; skip it so the first refresh lands one interval + // after acquisition rather than right away. + ticker.tick().await; + loop { + ticker.tick().await; + + let results = join_all(entries.iter().map(|(lock_id, client)| { + let lock_id = lock_id.clone(); + let client = client.clone(); + async move { client.refresh(&lock_id).await } + })) + .await; + + let mut refreshed = 0usize; + let mut not_found = 0usize; + for result in &results { + match result { + Ok(true) => refreshed += 1, + Ok(false) => not_found += 1, + // RPC jitter: count as neither; the ttl > interval margin covers a transient + // dip and the next tick retries. + Err(_) => {} + } + } + + if not_found > tolerable_not_found { + warn!( + resource = %resource, + owner = %owner, + refreshed, + not_found, + entries = entries.len(), + refresh_quorum, + "lock refresh lost quorum" + ); + record_lock_refresh_quorum_lost(); + lock_lost.mark_lost(); + } } } @@ -138,9 +275,32 @@ impl DistributedLockGuard { &self.lock_id } + /// Whether the guard's heartbeat has observed a refresh-quorum loss. + pub fn is_lock_lost(&self) -> bool { + self.lock_lost.is_lost() + } + + /// Shared lock-loss signal (used by callers to `select!` on `notified()`; Phase 2). + pub fn lock_lost(&self) -> Arc { + self.lock_lost.clone() + } + + /// Abort the heartbeat task, if running. Idempotent. + /// + /// Aborting may truncate an in-flight refresh, but refresh is idempotent and only extends + /// an existing lease (never creates one): a late refresh reaching a backend that already + /// dropped the entry returns `Ok(false)` with no side effect. + fn stop_heartbeat(&mut self) { + if let Some(task) = self.refresh_task.take() { + task.abort(); + } + } + /// Manually disarm the guard so dropping it won't release the lock. /// Call this if you explicitly released the lock elsewhere. pub fn disarm(&mut self) { + // Stop renewing a lock the caller says it released elsewhere. + self.stop_heartbeat(); if !self.disarmed { record_lock_held_release(self.lock_type); } @@ -157,6 +317,10 @@ impl DistributedLockGuard { /// to prevent double-release on drop. /// Returns true if release was scheduled or the guard was already disarmed. pub fn release(&mut self) -> bool { + // Stop the heartbeat before releasing so no refresh races the unlock. Abort is safe even + // if disarmed (take() makes it idempotent / a no-op when never spawned). + self.stop_heartbeat(); + if self.disarmed { // Lock was already released, return true to indicate lock is in released state return true; @@ -272,7 +436,18 @@ impl DistributedLock { .map(|info| info.id.clone()) .unwrap_or_else(|| LockId::new_unique(&request.resource)); - Ok(Some(DistributedLockGuard::new(aggregate_lock_id, individual_locks, request.lock_type))) + // Heartbeat operates on the per-client individual locks (their per-client ids), never + // the aggregate id, so refreshes round-trip to the exact backend entries. + let refresh_interval = derive_refresh_interval(individual_locks.len(), request.ttl, request.refresh_interval); + Ok(Some(DistributedLockGuard::new( + aggregate_lock_id, + individual_locks, + request.lock_type, + refresh_interval, + required_quorum, + request.owner.clone(), + request.resource.clone(), + ))) } else { // Check if it's a timeout or quorum failure if let Some(error_msg) = &resp.error { @@ -836,12 +1011,294 @@ mod tests { }; use crate::{LockError, LockId, LockInfo, LockRequest, LockResponse, LockStats, LockType, ObjectKey, client::LockClient}; use std::assert_matches; + use std::sync::atomic::{AtomicUsize, Ordering}; use std::{ collections::{HashMap, VecDeque}, sync::{Arc, Mutex}, time::Duration, }; + #[derive(Clone, Copy, Debug)] + enum RefreshOutcome { + Alive, // Ok(true) refresh succeeded + NotFound, // Ok(false) remote already lost the lock (reclaimed / never held) + RpcError, // Err RPC jitter + } + + /// Counting test client: acquires successfully and echoes back request.lock_id as + /// the per-client id, so the guard heartbeat calls refresh with that id; refresh + /// increments a counter and returns per the configured outcome. + #[derive(Debug)] + struct RefreshCountingClient { + refresh_calls: Arc, + outcome: RefreshOutcome, + active: tokio::sync::Mutex>, + } + + impl RefreshCountingClient { + fn new(outcome: RefreshOutcome, counter: Arc) -> Arc { + Arc::new(Self { + refresh_calls: counter, + outcome, + active: tokio::sync::Mutex::new(std::collections::HashSet::new()), + }) + } + } + + #[async_trait::async_trait] + impl LockClient for RefreshCountingClient { + async fn acquire_lock(&self, request: &LockRequest) -> crate::Result { + self.active.lock().await.insert(request.lock_id.clone()); + let info = LockInfo { + id: request.lock_id.clone(), + resource: request.resource.clone(), + lock_type: request.lock_type, + status: crate::types::LockStatus::Acquired, + owner: request.owner.clone(), + acquired_at: std::time::SystemTime::now(), + expires_at: std::time::SystemTime::now() + request.ttl, + last_refreshed: std::time::SystemTime::now(), + metadata: request.metadata.clone(), + priority: request.priority, + wait_start_time: None, + }; + Ok(LockResponse::success(info, Duration::ZERO)) + } + async fn release(&self, lock_id: &LockId) -> crate::Result { + Ok(self.active.lock().await.remove(lock_id)) + } + async fn refresh(&self, _lock_id: &LockId) -> crate::Result { + self.refresh_calls.fetch_add(1, Ordering::SeqCst); + match self.outcome { + RefreshOutcome::Alive => Ok(true), + RefreshOutcome::NotFound => Ok(false), + RefreshOutcome::RpcError => Err(LockError::internal("refresh rpc jitter")), + } + } + async fn force_release(&self, lock_id: &LockId) -> crate::Result { + self.release(lock_id).await + } + async fn check_status(&self, _lock_id: &LockId) -> crate::Result> { + Ok(None) + } + async fn get_stats(&self) -> crate::Result { + Ok(LockStats::default()) + } + async fn close(&self) -> crate::Result<()> { + Ok(()) + } + async fn is_online(&self) -> bool { + true + } + async fn is_local(&self) -> bool { + false + } + } + + fn counting_clients(outcomes: &[RefreshOutcome]) -> (Vec>, Vec>) { + let mut clients: Vec> = Vec::new(); + let mut counters = Vec::new(); + for outcome in outcomes { + let counter = Arc::new(AtomicUsize::new(0)); + clients.push(RefreshCountingClient::new(*outcome, counter.clone()) as Arc); + counters.push(counter); + } + (clients, counters) + } + + // A1 -- heartbeat broadcasts refresh every interval; a live lock stays refreshed (#899 keepalive). + #[tokio::test] + async fn heartbeat_keeps_distributed_lock_refreshed() { + let (clients, counters) = counting_clients(&[ + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + ]); + let lock = DistributedLock::new("test".to_string(), clients, 3); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_acquire_timeout(Duration::from_millis(200)) + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); // interval < ttl -> spawn + + let guard = lock + .acquire_guard(&request) + .await + .expect("acquire should not error") + .expect("quorum should be reached"); + + tokio::time::sleep(Duration::from_millis(90)).await; // spans >2 intervals + + // The guard tracks exactly the quorum of per-client locks (excess successes are cleaned + // up on acquire), so the heartbeat broadcasts to at least `quorum` clients. + let refreshed_clients = counters.iter().filter(|c| c.load(Ordering::SeqCst) >= 1).count(); + assert!( + refreshed_clients >= 3, + "heartbeat should refresh at least the quorum of tracked clients, got {refreshed_clients}" + ); + assert!(!guard.is_lock_lost(), "all-alive refresh must not signal lock loss"); + drop(guard); + } + + // A2 -- single client (degenerate path) spawns no heartbeat (matches localLockInstance). + #[tokio::test] + async fn single_client_guard_spawns_no_heartbeat() { + let (clients, counters) = counting_clients(&[RefreshOutcome::Alive]); + let lock = DistributedLock::new("test".to_string(), clients, 1); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); + + let guard = lock.acquire_guard(&request).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(90)).await; + + assert_eq!(counters[0].load(Ordering::SeqCst), 0, "single-client guard must not run a heartbeat"); + drop(guard); + } + + // A3 -- dropping the guard aborts the heartbeat before release; no refresh lands after drop. + #[tokio::test] + async fn dropping_guard_stops_heartbeat() { + let (clients, counters) = counting_clients(&[ + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + ]); + let lock = DistributedLock::new("test".to_string(), clients, 3); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); + + let guard = lock.acquire_guard(&request).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + drop(guard); // should abort the heartbeat first + let after_drop: Vec = counters.iter().map(|c| c.load(Ordering::SeqCst)).collect(); + + tokio::time::sleep(Duration::from_millis(80)).await; // wait >2 more intervals + for (idx, counter) in counters.iter().enumerate() { + assert_eq!( + counter.load(Ordering::SeqCst), + after_drop[idx], + "client {idx} must not be refreshed after guard drop" + ); + } + } + + // A4 -- losing refresh quorum (not_found > n - refresh_quorum) signals lock loss. + // 4 clients / write-quorum 3 -> n-quorum = 1; 2 NotFound (>1) -> lost. + #[tokio::test] + async fn heartbeat_quorum_loss_signals_lock_lost() { + let (clients, _counters) = counting_clients(&[ + RefreshOutcome::NotFound, + RefreshOutcome::NotFound, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + ]); + let lock = DistributedLock::new("test".to_string(), clients, 3); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); + + let guard = lock.acquire_guard(&request).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(60)).await; // at least one tick + + assert!( + guard.is_lock_lost(), + "losing refresh quorum (2 not_found of 4, quorum 3) must signal lock loss" + ); + let signal = guard.lock_lost(); + tokio::time::timeout(Duration::from_millis(20), signal.notified()) + .await + .expect("lock_lost().notified() must resolve once quorum lost"); + drop(guard); + } + + // A5 -- refresh RPC jitter (Err) is not counted as not_found; the lock is not declared lost. + #[tokio::test] + async fn heartbeat_rpc_error_not_counted_as_lock_lost() { + let (clients, _counters) = counting_clients(&[ + RefreshOutcome::RpcError, + RefreshOutcome::RpcError, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + ]); + let lock = DistributedLock::new("test".to_string(), clients, 3); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); + + let guard = lock.acquire_guard(&request).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(60)).await; + + assert!( + !guard.is_lock_lost(), + "RPC jitter (Err) must NOT be counted as not_found; lock must not be declared lost" + ); + drop(guard); + } + + // A6 -- sub-second ttl boundary: no panic, and skip heartbeat when interval>=ttl + // (fixes clamp panic; aligns with the 50ms-ttl distributed path in namespace tests). + #[tokio::test] + async fn subsecond_ttl_guard_does_not_panic_and_skips_heartbeat() { + // tiny ttl; the point is that acquire must not panic. + let (clients, _counters) = counting_clients(&[RefreshOutcome::Alive, RefreshOutcome::Alive, RefreshOutcome::Alive]); + let lock = DistributedLock::new("test".to_string(), clients, 2); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_acquire_timeout(Duration::from_millis(100)) + .with_ttl(Duration::from_millis(50)); // a clamp(1s, ttl-1s) would panic here + let guard = lock + .acquire_guard(&request) + .await + .expect("subsecond ttl acquire must not panic") + .expect("quorum should be reached"); + drop(guard); + + // inject interval >= ttl -> degenerate to no spawn + let (clients2, counters2) = counting_clients(&[RefreshOutcome::Alive, RefreshOutcome::Alive, RefreshOutcome::Alive]); + let lock2 = DistributedLock::new("test".to_string(), clients2, 2); + let request2 = LockRequest::new(ObjectKey::new("bucket", "object2"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(50)) + .with_refresh_interval(Duration::from_millis(50)); // interval >= ttl -> None + let guard2 = lock2.acquire_guard(&request2).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(60)).await; + for c in &counters2 { + assert_eq!(c.load(Ordering::SeqCst), 0, "interval>=ttl must not spawn heartbeat"); + } + drop(guard2); + } + + // A7 -- disarm() must also stop the heartbeat (public API gap). + #[tokio::test] + async fn disarm_stops_heartbeat() { + let (clients, counters) = counting_clients(&[ + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + RefreshOutcome::Alive, + ]); + let lock = DistributedLock::new("test".to_string(), clients, 3); + let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(Duration::from_millis(120)) + .with_refresh_interval(Duration::from_millis(20)); + + let mut guard = lock.acquire_guard(&request).await.unwrap().unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + guard.disarm(); // should abort the heartbeat + let after_disarm: Vec = counters.iter().map(|c| c.load(Ordering::SeqCst)).collect(); + + tokio::time::sleep(Duration::from_millis(80)).await; + for (idx, counter) in counters.iter().enumerate() { + assert_eq!( + counter.load(Ordering::SeqCst), + after_disarm[idx], + "disarm() must stop the heartbeat for client {idx}" + ); + } + drop(guard); + } + #[derive(Debug)] struct ResponseClient { response: LockResponse, diff --git a/crates/lock/src/fast_lock/mod.rs b/crates/lock/src/fast_lock/mod.rs index dcc355d43..a03201e6c 100644 --- a/crates/lock/src/fast_lock/mod.rs +++ b/crates/lock/src/fast_lock/mod.rs @@ -58,6 +58,13 @@ pub const DEFAULT_SHARD_COUNT: usize = 1024; /// Default lock timeout (lease TTL; lock is released if not refreshed within this duration) pub const DEFAULT_LOCK_TIMEOUT: Duration = Duration::from_secs(30); +/// Default heartbeat interval for refreshing a held distributed lock (lease renewal cadence). +/// +/// Kept at ttl/3 of `DEFAULT_LOCK_TIMEOUT` so a lock tolerates losing two consecutive refresh +/// cycles before its lease expires. TODO(backlog#899): maintainers to confirm this default vs +/// MinIO's 10s/20s interval-vs-validity and the #814 latency budget. +pub const DEFAULT_LOCK_REFRESH_INTERVAL: Duration = Duration::from_secs(10); + /// Default acquire timeout - common value for local/low-latency; use env to increase for slow storage pub const DEFAULT_ACQUIRE_TIMEOUT: Duration = Duration::from_secs(DEFAULT_RUSTFS_ACQUIRE_TIMEOUT); diff --git a/crates/lock/src/namespace/tests.rs b/crates/lock/src/namespace/tests.rs index c99d4fecc..b1b99df5e 100644 --- a/crates/lock/src/namespace/tests.rs +++ b/crates/lock/src/namespace/tests.rs @@ -1235,3 +1235,94 @@ async fn test_namespace_lock_distributed_failure_retries_and_cleans_up_late_succ drop(write_guard); } + +// C1 -- a long-held write lock survives past its ttl thanks to the heartbeat; a +// contender fails to steal it while it is held (#899 direct regression). +// Without a heartbeat: owner-a's per-node entries expire after ttl, owner-b triggers +// reclaim and steals it -> contended.is_some() -> assertion fails (red). Once fixed the +// entries stay refreshed -> is_none() (green). +#[tokio::test] +async fn distributed_write_lock_survives_past_ttl_with_heartbeat() { + let managers = (0..3).map(|_| Arc::new(GlobalLockManager::new())).collect::>(); + let clients = managers + .iter() + .map(|m| Arc::new(LocalClient::with_manager(m.clone())) as Arc) + .collect::>(); + let lock = Arc::new(NamespaceLock::with_clients("heartbeat-keepalive".to_string(), clients)); + let resource = create_test_object_key("bucket", "object-heartbeat-keepalive"); + + let req_a = LockRequest::new(resource.clone(), LockType::Exclusive, "owner-a") + .with_acquire_timeout(Duration::from_millis(200)) + .with_ttl(Duration::from_millis(150)) + .with_refresh_interval(Duration::from_millis(40)); // < ttl -> spawn heartbeat + + let guard_a = lock + .acquire_guard(&req_a) + .await + .expect("owner-a acquire should not error") + .expect("owner-a should hold the distributed write lock"); + + tokio::time::sleep(Duration::from_millis(400)).await; // hold well past a single ttl window + + let req_b = LockRequest::new(resource.clone(), LockType::Exclusive, "owner-b") + .with_acquire_timeout(Duration::from_millis(150)) + .with_ttl(Duration::from_millis(150)); + let contended = lock.acquire_guard(&req_b).await.expect("owner-b acquire should not error"); + assert!( + contended.is_none(), + "heartbeat must keep owner-a's lock alive; owner-b must not steal it past ttl" + ); + + drop(guard_a); + let mut acquired_after_release = false; + for _ in 0..20 { + let req_c = LockRequest::new(resource.clone(), LockType::Exclusive, "owner-b") + .with_acquire_timeout(Duration::from_millis(150)) + .with_ttl(Duration::from_millis(150)); + if let Some(g) = lock.acquire_guard(&req_c).await.expect("acquire should not error") { + drop(g); + acquired_after_release = true; + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + assert!(acquired_after_release, "after owner-a releases, owner-b must eventually acquire"); +} + +// C2 -- a crashed owner (leaving orphan entries that are never refreshed again) has its +// lock eventually reclaimed; no permanent deadlock (#698 scavenger). Must stay green. +#[tokio::test] +async fn crashed_owner_distributed_lock_is_reclaimed() { + let managers = (0..3).map(|_| Arc::new(GlobalLockManager::new())).collect::>(); + let node_clients = managers + .iter() + .map(|m| Arc::new(LocalClient::with_manager(m.clone()))) + .collect::>(); + let resource = create_test_object_key("bucket", "object-crashed-owner"); + + // Simulate a dead coordinator: plant short-ttl orphan entries on each node backend, + // then never refresh them. + for c in &node_clients { + let orphan = LockRequest::new(resource.clone(), LockType::Exclusive, "dead-owner").with_ttl(Duration::from_millis(60)); + let resp = c.acquire_lock(&orphan).await.expect("orphan acquire"); + assert!(resp.success, "orphan entry should be planted on each node"); + } + + tokio::time::sleep(Duration::from_millis(120)).await; // ttl elapses + + let clients = node_clients + .iter() + .map(|c| c.clone() as Arc) + .collect::>(); + let lock = NamespaceLock::with_clients("crashed-owner-recovery".to_string(), clients); + + let req_new = LockRequest::new(resource.clone(), LockType::Exclusive, "owner-live") + .with_acquire_timeout(Duration::from_millis(300)) + .with_ttl(Duration::from_millis(150)); + let recovered = lock + .acquire_guard(&req_new) + .await + .expect("acquire should not error") + .expect("stale orphan entries must be reclaimed; no permanent deadlock (#698)"); + drop(recovered); +} diff --git a/crates/lock/src/types.rs b/crates/lock/src/types.rs index b531082e5..eebe158ff 100644 --- a/crates/lock/src/types.rs +++ b/crates/lock/src/types.rs @@ -227,6 +227,13 @@ pub struct LockRequest { pub deadlock_detection: bool, /// Suppress warning logs for expected contention failures pub suppress_contention_logs: bool, + /// Optional override for the guard heartbeat/refresh interval. + /// + /// `None` means the holding guard derives an interval from `ttl` (ttl/3). Mainly injected + /// by tests to drive very small intervals; `#[serde(default)]` keeps older RPC payloads + /// (which never encoded this field) deserializable. + #[serde(default)] + pub refresh_interval: Option, } impl LockRequest { @@ -243,6 +250,7 @@ impl LockRequest { priority: LockPriority::default(), deadlock_detection: false, suppress_contention_logs: false, + refresh_interval: None, } } @@ -281,6 +289,12 @@ impl LockRequest { self.suppress_contention_logs = enabled; self } + + /// Override the guard heartbeat/refresh interval (otherwise derived from `ttl`). + pub fn with_refresh_interval(mut self, interval: Duration) -> Self { + self.refresh_interval = Some(interval); + self + } } /// Lock response structure diff --git a/rustfs/src/storage/rpc/node_service/lock.rs b/rustfs/src/storage/rpc/node_service/lock.rs index ff53e706e..cb2e423ee 100644 --- a/rustfs/src/storage/rpc/node_service/lock.rs +++ b/rustfs/src/storage/rpc/node_service/lock.rs @@ -33,6 +33,41 @@ fn lock_result_from_error(error: impl Into) -> GenerallyLockResult { } } +/// Delegate a refresh RPC to the node's lock backend (`LocalClient::refresh`). +/// +/// Extracted as a free function (decoupled from the `current_lock_client()` OnceLock +/// global) so it can be unit-tested against a real `LocalClient`, mirroring the existing +/// `lock_result_from_release` free-function test style. +/// +/// Contract (see backlog#899 design, open question on not-found signalling): a lock the +/// backend no longer holds returns `success=false` with `error_info=None`, so +/// `RemoteClient::refresh` keeps mapping it to `Ok(false)` (a real not_found signal for the +/// coordinator heartbeat) rather than `Err`. `error_info` is reserved for genuine RPC/backend +/// errors. TODO(backlog#899): maintainers to confirm whether not-found should instead carry a +/// distinguishable error_info prefix that `RemoteClient::refresh` decodes explicitly. +async fn refresh_lock( + lock_client: &std::sync::Arc, + args: &LockRequest, +) -> GenerallyLockResponse { + match lock_client.refresh(&args.lock_id).await { + Ok(true) => GenerallyLockResponse { + success: true, + error_info: None, + lock_info: None, + }, + Ok(false) => GenerallyLockResponse { + success: false, + error_info: None, + lock_info: None, + }, + Err(err) => GenerallyLockResponse { + success: false, + error_info: Some(format!("can not refresh, resource: {0}, err: {1}", args.resource, err)), + lock_info: None, + }, + } +} + fn lock_result_from_release(lock_id: &rustfs_lock::LockId, success: bool) -> GenerallyLockResult { if success { GenerallyLockResult { @@ -51,7 +86,7 @@ impl NodeService { request: Request, ) -> Result, Status> { let request = request.into_inner(); - let _args: LockRequest = match serde_json::from_str(&request.args) { + let args: LockRequest = match serde_json::from_str(&request.args) { Ok(args) => args, Err(err) => { return Ok(Response::new(GenerallyLockResponse { @@ -62,11 +97,8 @@ impl NodeService { } }; - Ok(Response::new(GenerallyLockResponse { - success: true, - error_info: None, - lock_info: None, - })) + let lock_client = self.get_lock_client()?; + Ok(Response::new(refresh_lock(&lock_client, &args).await)) } pub(super) async fn handle_force_un_lock( @@ -303,4 +335,50 @@ mod tests { assert_eq!(result.error_info.as_deref(), Some("lock conflict")); assert!(result.lock_info.is_none()); } + + // B1 -- refresh_lock delegates to LocalClient::refresh for a held lock, returns + // success=true, and actually extends expires_at (inverts the "parse then drop" stub). + #[tokio::test] + async fn refresh_lock_delegates_to_lock_client_and_extends_entry() { + use super::refresh_lock; + use rustfs_lock::client::LockClient; + use rustfs_lock::{LocalClient, LockRequest, LockStatus, LockType, ObjectKey}; + use std::sync::Arc; + + let client: Arc = Arc::new(LocalClient::new()); + let req = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") + .with_ttl(std::time::Duration::from_secs(30)); + let acquired = client.acquire_lock(&req).await.expect("acquire"); + assert!(acquired.success, "precondition: lock acquired"); + + let resp = refresh_lock(&client, &req).await; + assert!( + resp.success, + "handle_refresh must delegate to lock_client.refresh and succeed for a held lock" + ); + assert!(resp.error_info.is_none()); + + let status = client.check_status(&req.lock_id).await.unwrap().expect("entry present"); + assert_eq!(status.status, LockStatus::Acquired, "refresh must extend the entry, keeping it Acquired"); + } + + // B2 -- refresh_lock reports a non-existent lock_id as success=false + // (feeds broadcast_refresh a real not_found signal; the stub's unconditional + // true would permanently mask lock loss). + #[tokio::test] + async fn refresh_lock_reports_missing_lock_as_failure() { + use super::refresh_lock; + use rustfs_lock::client::LockClient; + use rustfs_lock::{LocalClient, LockRequest, LockType, ObjectKey}; + use std::sync::Arc; + + let client: Arc = Arc::new(LocalClient::new()); + let req = LockRequest::new(ObjectKey::new("bucket", "ghost"), LockType::Exclusive, "owner"); + + let resp = refresh_lock(&client, &req).await; + assert!( + !resp.success, + "refresh of a non-existent lock must return success=false (not the stub's unconditional true)" + ); + } }