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).
This commit is contained in:
Zhengchao An
2026-07-08 06:01:45 +08:00
committed by GitHub
parent 45435d83ab
commit ddf197ba57
10 changed files with 714 additions and 14 deletions
@@ -63,6 +63,7 @@ impl GrpcLockClient {
priority: LockPriority::Normal,
deadlock_detection: false,
suppress_contention_logs: false,
refresh_interval: None,
}
}
@@ -73,6 +73,7 @@ impl RemoteClient {
priority: LockPriority::Normal,
deadlock_detection: false,
suppress_contention_logs: false,
refresh_interval: None,
}
}
+3 -2
View File
@@ -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};
@@ -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<u64>,
@@ -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()
+23 -2
View File
@@ -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 {
+461 -4
View File
@@ -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<Duration>) -> Option<Duration> {
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<JoinHandle<()>>,
/// Lock-loss signal, shared with the heartbeat task.
lock_lost: Arc<LockLostSignal>,
}
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<dyn LockClient>)>, 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<dyn LockClient>)>,
lock_type: LockType,
refresh_interval: Option<Duration>,
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<dyn LockClient>)>,
interval: Duration,
refresh_quorum: usize,
lock_lost: Arc<LockLostSignal>,
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<LockLostSignal> {
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<AtomicUsize>,
outcome: RefreshOutcome,
active: tokio::sync::Mutex<std::collections::HashSet<LockId>>,
}
impl RefreshCountingClient {
fn new(outcome: RefreshOutcome, counter: Arc<AtomicUsize>) -> Arc<Self> {
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<LockResponse> {
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<bool> {
Ok(self.active.lock().await.remove(lock_id))
}
async fn refresh(&self, _lock_id: &LockId) -> crate::Result<bool> {
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<bool> {
self.release(lock_id).await
}
async fn check_status(&self, _lock_id: &LockId) -> crate::Result<Option<LockInfo>> {
Ok(None)
}
async fn get_stats(&self) -> crate::Result<LockStats> {
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<Arc<dyn LockClient>>, Vec<Arc<AtomicUsize>>) {
let mut clients: Vec<Arc<dyn LockClient>> = 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<dyn LockClient>);
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<usize> = 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<usize> = 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,
+7
View File
@@ -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);
+91
View File
@@ -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::<Vec<_>>();
let clients = managers
.iter()
.map(|m| Arc::new(LocalClient::with_manager(m.clone())) as Arc<dyn LockClient>)
.collect::<Vec<_>>();
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::<Vec<_>>();
let node_clients = managers
.iter()
.map(|m| Arc::new(LocalClient::with_manager(m.clone())))
.collect::<Vec<_>>();
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<dyn LockClient>)
.collect::<Vec<_>>();
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);
}
+14
View File
@@ -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<Duration>,
}
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
+84 -6
View File
@@ -33,6 +33,41 @@ fn lock_result_from_error(error: impl Into<String>) -> 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<dyn rustfs_lock::client::LockClient>,
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<GenerallyLockRequest>,
) -> Result<Response<GenerallyLockResponse>, 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<dyn LockClient> = 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<dyn LockClient> = 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)"
);
}
}