From e73ed4f2a1a20eee32e411dbcaf3592849afb3a2 Mon Sep 17 00:00:00 2001 From: Hiroaki KAWAI Date: Sat, 25 Jul 2026 21:26:37 +0900 Subject: [PATCH] fix(lock): jitter distributed lock retries (#5113) * fix(lock): jitter distributed lock retries * fix(lock): keep retry jitter within bounds * fix(lock): qualify test size_of usage --------- Signed-off-by: houseme Co-authored-by: houseme --- Cargo.lock | 1 + crates/lock/Cargo.toml | 4 ++ crates/lock/src/distributed_lock.rs | 67 +++++++++++++++++++++++++---- 3 files changed, 64 insertions(+), 8 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a252e293a..7685dabaa 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9497,6 +9497,7 @@ dependencies = [ "crossbeam-queue", "futures", "parking_lot", + "rand 0.10.2", "rustfs-io-metrics", "rustfs-utils", "serde", diff --git a/crates/lock/Cargo.toml b/crates/lock/Cargo.toml index dd807364d..7b0993adb 100644 --- a/crates/lock/Cargo.toml +++ b/crates/lock/Cargo.toml @@ -41,9 +41,13 @@ tracing.workspace = true uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] } thiserror.workspace = true parking_lot.workspace = true +rand.workspace = true smallvec = { workspace = true, features = ["serde"] } smartstring.workspace = true crossbeam-queue = { workspace = true } +[dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } + [lib] doctest = false diff --git a/crates/lock/src/distributed_lock.rs b/crates/lock/src/distributed_lock.rs index d6fb896f6..97b8af06d 100644 --- a/crates/lock/src/distributed_lock.rs +++ b/crates/lock/src/distributed_lock.rs @@ -22,6 +22,7 @@ use futures::{ future::join_all, stream::{FuturesUnordered, StreamExt}, }; +use rand::RngExt as _; use rustfs_io_metrics::{ 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, @@ -676,8 +677,15 @@ impl DistributedLock { pending } - fn lock_acquire_retry_backoff(attempt: usize) -> Duration { - LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF.saturating_mul(attempt.try_into().unwrap_or(u32::MAX)) + fn lock_acquire_retry_backoff(attempt: usize, rng: &mut impl rand::Rng) -> Duration { + let base = LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF.saturating_mul(attempt.try_into().unwrap_or(u32::MAX)); + let jitter = base / 4; + let lower = base.saturating_sub(jitter); + let upper = base.saturating_add(jitter); + let lower_nanos = u64::try_from(lower.as_nanos()).unwrap_or(u64::MAX); + let upper_nanos = u64::try_from(upper.as_nanos()).unwrap_or(u64::MAX); + + Duration::from_nanos(rng.random_range(lower_nanos..=upper_nanos)) } fn lock_acquire_attempt_timeout(&self, remaining: Duration) -> Duration { @@ -831,7 +839,7 @@ impl DistributedLock { } last_result = Some(result); - let backoff = Self::lock_acquire_retry_backoff(attempt); + let backoff = Self::lock_acquire_retry_backoff(attempt, &mut rand::rng()); if start.elapsed().saturating_add(backoff) >= request.acquire_timeout { break; } @@ -1154,11 +1162,13 @@ fn record_lock_held_release(lock_type: LockType) { #[cfg(test)] mod tests { use super::{ - DistributedLock, LOCK_ACQUIRE_ATTEMPT_TIMEOUT, LockAcquireFailureKind, LockLostSignal, is_remote_lock_rpc_failure, - should_warn_lock_failure, + DistributedLock, LOCK_ACQUIRE_ATTEMPT_TIMEOUT, LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF, LockAcquireFailureKind, + LockLostSignal, is_remote_lock_rpc_failure, should_warn_lock_failure, }; use crate::{LockError, LockId, LockInfo, LockRequest, LockResponse, LockStats, LockType, ObjectKey, client::LockClient}; + use rand::{SeedableRng as _, TryRng, rngs::StdRng}; use std::assert_matches; + use std::convert::Infallible; use std::sync::atomic::{AtomicUsize, Ordering}; use std::{ collections::{HashMap, VecDeque}, @@ -1174,6 +1184,47 @@ mod tests { RpcError, // Err RPC jitter } + struct ConstantRng(u64); + + impl TryRng for ConstantRng { + type Error = Infallible; + + fn try_next_u32(&mut self) -> Result { + Ok(self.0 as u32) + } + + fn try_next_u64(&mut self) -> Result { + Ok(self.0) + } + + fn try_fill_bytes(&mut self, dst: &mut [u8]) -> Result<(), Self::Error> { + for chunk in dst.chunks_mut(std::mem::size_of::()) { + chunk.copy_from_slice(&self.0.to_ne_bytes()[..chunk.len()]); + } + Ok(()) + } + } + + #[test] + fn lock_acquire_retry_backoff_uses_bounded_jitter() { + let mut rng = StdRng::seed_from_u64(42); + let base = LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF * 3; + let jitter = base / 4; + let lower = base - jitter; + let upper = base + jitter; + let samples = (0..32) + .map(|_| DistributedLock::lock_acquire_retry_backoff(3, &mut rng)) + .collect::>(); + + assert!(samples.iter().all(|delay| *delay >= lower && *delay <= upper)); + assert!(samples.windows(2).any(|pair| pair[0] != pair[1])); + + let mut lower_rng = ConstantRng(0); + let mut upper_rng = ConstantRng(u64::MAX); + assert_eq!(DistributedLock::lock_acquire_retry_backoff(3, &mut lower_rng), lower); + assert_eq!(DistributedLock::lock_acquire_retry_backoff(3, &mut upper_rng), upper); + } + /// 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. @@ -2278,7 +2329,7 @@ mod tests { drop(guard); } - #[tokio::test] + #[tokio::test(start_paused = true)] async fn acquire_guard_returns_timeout_when_quorum_remains_contended() { let clients: Vec> = vec![ ResponseClient::new(LockResponse::failure("lock already held", Duration::ZERO)).into_client(), @@ -2290,9 +2341,9 @@ mod tests { let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner") .with_acquire_timeout(Duration::from_millis(120)); - let result = lock.acquire_guard(&request).await; + let result = tokio::time::timeout(Duration::from_secs(1), lock.acquire_guard(&request)).await; - assert!(matches!(result, Ok(None)), "unexpected result: {result:?}"); + assert_matches!(result, Ok(Ok(None))); } #[tokio::test]