mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 13:16:28 +00:00
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 <housemecn@gmail.com> Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Generated
+1
@@ -9497,6 +9497,7 @@ dependencies = [
|
|||||||
"crossbeam-queue",
|
"crossbeam-queue",
|
||||||
"futures",
|
"futures",
|
||||||
"parking_lot",
|
"parking_lot",
|
||||||
|
"rand 0.10.2",
|
||||||
"rustfs-io-metrics",
|
"rustfs-io-metrics",
|
||||||
"rustfs-utils",
|
"rustfs-utils",
|
||||||
"serde",
|
"serde",
|
||||||
|
|||||||
@@ -41,9 +41,13 @@ tracing.workspace = true
|
|||||||
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
||||||
thiserror.workspace = true
|
thiserror.workspace = true
|
||||||
parking_lot.workspace = true
|
parking_lot.workspace = true
|
||||||
|
rand.workspace = true
|
||||||
smallvec = { workspace = true, features = ["serde"] }
|
smallvec = { workspace = true, features = ["serde"] }
|
||||||
smartstring.workspace = true
|
smartstring.workspace = true
|
||||||
crossbeam-queue = { workspace = true }
|
crossbeam-queue = { workspace = true }
|
||||||
|
|
||||||
|
[dev-dependencies]
|
||||||
|
tokio = { workspace = true, features = ["test-util"] }
|
||||||
|
|
||||||
[lib]
|
[lib]
|
||||||
doctest = false
|
doctest = false
|
||||||
|
|||||||
@@ -22,6 +22,7 @@ use futures::{
|
|||||||
future::join_all,
|
future::join_all,
|
||||||
stream::{FuturesUnordered, StreamExt},
|
stream::{FuturesUnordered, StreamExt},
|
||||||
};
|
};
|
||||||
|
use rand::RngExt as _;
|
||||||
use rustfs_io_metrics::{
|
use rustfs_io_metrics::{
|
||||||
record_lock_refresh_quorum_lost, record_read_lock_held_acquire, record_read_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,
|
record_write_lock_held_acquire, record_write_lock_held_release,
|
||||||
@@ -676,8 +677,15 @@ impl DistributedLock {
|
|||||||
pending
|
pending
|
||||||
}
|
}
|
||||||
|
|
||||||
fn lock_acquire_retry_backoff(attempt: usize) -> Duration {
|
fn lock_acquire_retry_backoff(attempt: usize, rng: &mut impl rand::Rng) -> Duration {
|
||||||
LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF.saturating_mul(attempt.try_into().unwrap_or(u32::MAX))
|
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 {
|
fn lock_acquire_attempt_timeout(&self, remaining: Duration) -> Duration {
|
||||||
@@ -831,7 +839,7 @@ impl DistributedLock {
|
|||||||
}
|
}
|
||||||
|
|
||||||
last_result = Some(result);
|
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 {
|
if start.elapsed().saturating_add(backoff) >= request.acquire_timeout {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -1154,11 +1162,13 @@ fn record_lock_held_release(lock_type: LockType) {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
DistributedLock, LOCK_ACQUIRE_ATTEMPT_TIMEOUT, LockAcquireFailureKind, LockLostSignal, is_remote_lock_rpc_failure,
|
DistributedLock, LOCK_ACQUIRE_ATTEMPT_TIMEOUT, LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF, LockAcquireFailureKind,
|
||||||
should_warn_lock_failure,
|
LockLostSignal, is_remote_lock_rpc_failure, should_warn_lock_failure,
|
||||||
};
|
};
|
||||||
use crate::{LockError, LockId, LockInfo, LockRequest, LockResponse, LockStats, LockType, ObjectKey, client::LockClient};
|
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::assert_matches;
|
||||||
|
use std::convert::Infallible;
|
||||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
use std::{
|
use std::{
|
||||||
collections::{HashMap, VecDeque},
|
collections::{HashMap, VecDeque},
|
||||||
@@ -1174,6 +1184,47 @@ mod tests {
|
|||||||
RpcError, // Err RPC jitter
|
RpcError, // Err RPC jitter
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct ConstantRng(u64);
|
||||||
|
|
||||||
|
impl TryRng for ConstantRng {
|
||||||
|
type Error = Infallible;
|
||||||
|
|
||||||
|
fn try_next_u32(&mut self) -> Result<u32, Self::Error> {
|
||||||
|
Ok(self.0 as u32)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn try_next_u64(&mut self) -> Result<u64, Self::Error> {
|
||||||
|
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::<u64>()) {
|
||||||
|
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::<Vec<_>>();
|
||||||
|
|
||||||
|
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
|
/// 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
|
/// the per-client id, so the guard heartbeat calls refresh with that id; refresh
|
||||||
/// increments a counter and returns per the configured outcome.
|
/// increments a counter and returns per the configured outcome.
|
||||||
@@ -2278,7 +2329,7 @@ mod tests {
|
|||||||
drop(guard);
|
drop(guard);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test(start_paused = true)]
|
||||||
async fn acquire_guard_returns_timeout_when_quorum_remains_contended() {
|
async fn acquire_guard_returns_timeout_when_quorum_remains_contended() {
|
||||||
let clients: Vec<Arc<dyn LockClient>> = vec![
|
let clients: Vec<Arc<dyn LockClient>> = vec![
|
||||||
ResponseClient::new(LockResponse::failure("lock already held", Duration::ZERO)).into_client(),
|
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")
|
let request = LockRequest::new(ObjectKey::new("bucket", "object"), LockType::Exclusive, "owner")
|
||||||
.with_acquire_timeout(Duration::from_millis(120));
|
.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]
|
#[tokio::test]
|
||||||
|
|||||||
Reference in New Issue
Block a user