mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-20 19:42:17 +00:00
fix(lock): isolate retry attempt lock ids (#3102)
* fix(lock): isolate retry attempt lock ids Generate a fresh lock id for each distributed lock retry attempt so late pending cleanup from an earlier attempt cannot release a lock acquired by a later retry. Also add a regression test covering late-success cleanup versus retried acquisition. * style(lock): normalize retry race regression test Normalize the distributed lock retry race regression test formatting without changing behavior. * test(lock): remove quorum-order assumption Make the retry race regression test assert the safety property without assuming which delayed client enters the retry quorum first.
This commit is contained in:
@@ -494,6 +494,7 @@ impl DistributedLock {
|
|||||||
let remaining = request.acquire_timeout - elapsed;
|
let remaining = request.acquire_timeout - elapsed;
|
||||||
let mut attempt_request = request.clone();
|
let mut attempt_request = request.clone();
|
||||||
attempt_request.acquire_timeout = remaining;
|
attempt_request.acquire_timeout = remaining;
|
||||||
|
attempt_request.lock_id = LockId::new_unique(&request.resource);
|
||||||
|
|
||||||
let result = self.acquire_lock_quorum_once(&attempt_request).await?;
|
let result = self.acquire_lock_quorum_once(&attempt_request).await?;
|
||||||
if result.response.success
|
if result.response.success
|
||||||
@@ -711,7 +712,11 @@ fn record_lock_held_release(lock_type: LockType) {
|
|||||||
mod tests {
|
mod tests {
|
||||||
use super::{DistributedLock, is_remote_lock_rpc_failure};
|
use super::{DistributedLock, is_remote_lock_rpc_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 std::{sync::Arc, time::Duration};
|
use std::{
|
||||||
|
collections::{HashMap, VecDeque},
|
||||||
|
sync::{Arc, Mutex},
|
||||||
|
time::Duration,
|
||||||
|
};
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
struct ResponseClient {
|
struct ResponseClient {
|
||||||
@@ -779,6 +784,107 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy)]
|
||||||
|
enum AcquirePlan {
|
||||||
|
Success { delay: Duration },
|
||||||
|
Failure { error: &'static str, delay: Duration },
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct SequencedClient {
|
||||||
|
plans: Mutex<VecDeque<AcquirePlan>>,
|
||||||
|
active: tokio::sync::Mutex<HashMap<LockId, LockInfo>>,
|
||||||
|
seen_ids: Arc<Mutex<Vec<LockId>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SequencedClient {
|
||||||
|
fn new(plans: Vec<AcquirePlan>, seen_ids: Arc<Mutex<Vec<LockId>>>) -> Self {
|
||||||
|
Self {
|
||||||
|
plans: Mutex::new(plans.into()),
|
||||||
|
active: tokio::sync::Mutex::new(HashMap::new()),
|
||||||
|
seen_ids,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl LockClient for SequencedClient {
|
||||||
|
async fn acquire_lock(&self, request: &LockRequest) -> crate::Result<LockResponse> {
|
||||||
|
self.seen_ids.lock().unwrap().push(request.lock_id.clone());
|
||||||
|
|
||||||
|
let plan = self
|
||||||
|
.plans
|
||||||
|
.lock()
|
||||||
|
.unwrap()
|
||||||
|
.pop_front()
|
||||||
|
.unwrap_or(AcquirePlan::Success { delay: Duration::ZERO });
|
||||||
|
|
||||||
|
match plan {
|
||||||
|
AcquirePlan::Success { delay } => {
|
||||||
|
if !delay.is_zero() {
|
||||||
|
tokio::time::sleep(delay).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
let lock_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,
|
||||||
|
};
|
||||||
|
|
||||||
|
self.active.lock().await.insert(request.lock_id.clone(), lock_info.clone());
|
||||||
|
|
||||||
|
Ok(LockResponse::success(lock_info, Duration::ZERO))
|
||||||
|
}
|
||||||
|
AcquirePlan::Failure { error, delay } => {
|
||||||
|
if !delay.is_zero() {
|
||||||
|
tokio::time::sleep(delay).await;
|
||||||
|
}
|
||||||
|
Ok(LockResponse::failure(error, Duration::ZERO))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn release(&self, lock_id: &LockId) -> crate::Result<bool> {
|
||||||
|
Ok(self.active.lock().await.remove(lock_id).is_some())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn refresh(&self, _lock_id: &LockId) -> crate::Result<bool> {
|
||||||
|
Ok(false)
|
||||||
|
}
|
||||||
|
|
||||||
|
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(self.active.lock().await.get(lock_id).cloned())
|
||||||
|
}
|
||||||
|
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_is_remote_lock_rpc_failure() {
|
fn test_is_remote_lock_rpc_failure() {
|
||||||
assert!(is_remote_lock_rpc_failure("Remote lock RPC failed: backend unreachable"));
|
assert!(is_remote_lock_rpc_failure("Remote lock RPC failed: backend unreachable"));
|
||||||
@@ -938,4 +1044,105 @@ mod tests {
|
|||||||
|
|
||||||
assert!(matches!(result, Ok(None)), "unexpected result: {result:?}");
|
assert!(matches!(result, Ok(None)), "unexpected result: {result:?}");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn acquire_guard_uses_fresh_lock_ids_across_retry_attempts() {
|
||||||
|
let seen_ids = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let delayed_success_a = Arc::new(SequencedClient::new(
|
||||||
|
vec![
|
||||||
|
AcquirePlan::Success {
|
||||||
|
delay: Duration::from_millis(400),
|
||||||
|
},
|
||||||
|
AcquirePlan::Success { delay: Duration::ZERO },
|
||||||
|
],
|
||||||
|
seen_ids.clone(),
|
||||||
|
));
|
||||||
|
let delayed_success_b = Arc::new(SequencedClient::new(
|
||||||
|
vec![
|
||||||
|
AcquirePlan::Success {
|
||||||
|
delay: Duration::from_millis(400),
|
||||||
|
},
|
||||||
|
AcquirePlan::Success { delay: Duration::ZERO },
|
||||||
|
],
|
||||||
|
seen_ids.clone(),
|
||||||
|
));
|
||||||
|
|
||||||
|
let clients: Vec<Arc<dyn LockClient>> = vec![
|
||||||
|
Arc::new(SequencedClient::new(
|
||||||
|
vec![
|
||||||
|
AcquirePlan::Failure {
|
||||||
|
error: "Lock acquisition timeout",
|
||||||
|
delay: Duration::ZERO,
|
||||||
|
},
|
||||||
|
AcquirePlan::Success { delay: Duration::ZERO },
|
||||||
|
],
|
||||||
|
seen_ids.clone(),
|
||||||
|
)),
|
||||||
|
Arc::new(SequencedClient::new(
|
||||||
|
vec![
|
||||||
|
AcquirePlan::Failure {
|
||||||
|
error: "Lock acquisition timeout",
|
||||||
|
delay: Duration::ZERO,
|
||||||
|
},
|
||||||
|
AcquirePlan::Success { delay: Duration::ZERO },
|
||||||
|
],
|
||||||
|
seen_ids.clone(),
|
||||||
|
)),
|
||||||
|
delayed_success_a.clone(),
|
||||||
|
delayed_success_b.clone(),
|
||||||
|
];
|
||||||
|
|
||||||
|
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_secs(2));
|
||||||
|
|
||||||
|
let guard = lock
|
||||||
|
.acquire_guard(&request)
|
||||||
|
.await
|
||||||
|
.expect("retry path should not fail")
|
||||||
|
.expect("second attempt should reach quorum");
|
||||||
|
|
||||||
|
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||||
|
|
||||||
|
let mut unique_ids = Vec::new();
|
||||||
|
for lock_id in seen_ids.lock().unwrap().iter() {
|
||||||
|
if unique_ids.iter().all(|id: &LockId| id.uuid != lock_id.uuid) {
|
||||||
|
unique_ids.push(lock_id.clone());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
unique_ids.len() >= 2,
|
||||||
|
"expected retry attempts to use distinct lock ids, saw {unique_ids:?}"
|
||||||
|
);
|
||||||
|
|
||||||
|
let retry_lock_ids = unique_ids.iter().skip(1).cloned().collect::<Vec<_>>();
|
||||||
|
assert!(
|
||||||
|
!retry_lock_ids.is_empty(),
|
||||||
|
"expected at least one retry-attempt lock id, saw {unique_ids:?}"
|
||||||
|
);
|
||||||
|
|
||||||
|
let delayed_a_active = delayed_success_a.active.lock().await;
|
||||||
|
let delayed_b_active = delayed_success_b.active.lock().await;
|
||||||
|
let remaining_delayed_lock_ids = delayed_a_active
|
||||||
|
.keys()
|
||||||
|
.chain(delayed_b_active.keys())
|
||||||
|
.cloned()
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
remaining_delayed_lock_ids.len() == 1,
|
||||||
|
"exactly one delayed client lock should remain held by the retry guard after late cleanup"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
retry_lock_ids
|
||||||
|
.iter()
|
||||||
|
.any(|retry_lock_id| retry_lock_id.uuid == remaining_delayed_lock_ids[0].uuid),
|
||||||
|
"late cleanup must not leave only a first-attempt delayed lock active"
|
||||||
|
);
|
||||||
|
drop(delayed_b_active);
|
||||||
|
drop(delayed_a_active);
|
||||||
|
|
||||||
|
drop(guard);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user