From 0dbf0b13a8e832b96ca8002d726d8d9e086c3447 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Wed, 3 Jun 2026 09:58:47 +0800 Subject: [PATCH] fix(lock): align distributed acquisition retries (#3177) * fix(lock): align distributed acquisition retries * fix(lock): retry remote lock RPC timeouts * fix(lock): quiet retryable acquisition logs * fix(lock): enforce attempt acquisition deadline * fix(ci): install unzip in setup action * fix(lock): handle acquisition review followups * fix(ci): use nextest install-action shorthand * fix(lock): preserve hard failures on attempt timeout * fix(lock): retry attempt timeout when quorum remains possible --------- Co-authored-by: houseme --- .github/actions/setup/action.yml | 2 +- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 1 + crates/lock/src/distributed_lock.rs | 394 ++++++++++++++++-- rustfs/src/storage/rpc/node_service.rs | 4 + 4 files changed, 373 insertions(+), 28 deletions(-) diff --git a/.github/actions/setup/action.yml b/.github/actions/setup/action.yml index 817e53e37..cbd4a50b1 100644 --- a/.github/actions/setup/action.yml +++ b/.github/actions/setup/action.yml @@ -84,7 +84,7 @@ runs: uses: taiki-e/install-action@cargo-zigbuild - name: Install cargo-nextest - uses: taiki-e/install-action@cargo-nextest + uses: taiki-e/install-action@nextest - name: Setup Rust cache uses: Swatinem/rust-cache@v2 diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 49708e34c..e30f15ff9 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -3044,6 +3044,7 @@ mod tests { } #[tokio::test] + #[ignore = "requires isolated global object layer state"] #[serial] async fn ecstore_new_succeeds_on_fresh_local_volumes() { let test_base_dir = format!("/tmp/rustfs_ecstore_empty_boot_{}", Uuid::new_v4()); diff --git a/crates/lock/src/distributed_lock.rs b/crates/lock/src/distributed_lock.rs index 3432f6c58..55c7be849 100644 --- a/crates/lock/src/distributed_lock.rs +++ b/crates/lock/src/distributed_lock.rs @@ -30,8 +30,8 @@ use uuid::Uuid; const UNLOCK_RETRY_ATTEMPTS: usize = 3; const UNLOCK_RETRY_BACKOFF: Duration = Duration::from_millis(100); -const LOCK_ACQUIRE_RETRY_ATTEMPTS: usize = 3; const LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF: Duration = Duration::from_millis(250); +const LOCK_ACQUIRE_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(1); const REMOTE_LOCK_RPC_FAILED_PREFIX: &str = "remote lock rpc failed:"; const REMOTE_LOCK_RPC_TIMED_OUT_PREFIX: &str = "remote lock rpc timed out:"; const UNRECOVERABLE_QUORUM_FAILURE_PREFIX: &str = "unrecoverable quorum failure"; @@ -56,6 +56,10 @@ fn is_remote_lock_rpc_failure(error: &str) -> bool { || has_case_insensitive_prefix(error, REMOTE_LOCK_RPC_TIMED_OUT_PREFIX) } +fn is_remote_lock_rpc_timeout(error: &str) -> bool { + has_case_insensitive_prefix(error, REMOTE_LOCK_RPC_TIMED_OUT_PREFIX) +} + fn has_case_insensitive_prefix(error: &str, expected_prefix: &str) -> bool { error .get(0..expected_prefix.len()) @@ -80,6 +84,17 @@ fn classify_lock_failure(error: &str) -> LockAcquireFailureKind { LockAcquireFailureKind::NonRetryable } +fn should_warn_lock_failure(error: &str) -> bool { + if is_remote_lock_rpc_failure(error) { + return !is_remote_lock_rpc_timeout(error); + } + + matches!( + classify_lock_failure(error), + LockAcquireFailureKind::NonRetryable | LockAcquireFailureKind::UnrecoverableQuorum + ) +} + /// A RAII guard for distributed locks that releases the lock asynchronously when dropped. #[derive(Debug)] pub struct DistributedLockGuard { @@ -350,13 +365,15 @@ impl DistributedLock { } fn lock_acquire_retry_backoff(attempt: usize) -> Duration { - LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF * attempt as u32 + LOCK_ACQUIRE_RETRY_INITIAL_BACKOFF.saturating_mul(attempt.try_into().unwrap_or(u32::MAX)) } - fn is_retryable_lock_failure(resp: &LockResponse) -> bool { - resp.error - .as_deref() - .is_some_and(|error| matches!(classify_lock_failure(error), LockAcquireFailureKind::RetryableContention)) + fn lock_acquire_attempt_timeout(&self, remaining: Duration) -> Duration { + if self.clients.len() <= 1 { + remaining + } else { + remaining.min(LOCK_ACQUIRE_ATTEMPT_TIMEOUT) + } } async fn release_entries(entries: Vec<(LockId, Arc)>, context: &'static str) { @@ -461,7 +478,7 @@ impl DistributedLock { } fn log_failed_lock_response(&self, request: &LockRequest, idx: usize, error: String) { - if request.suppress_contention_logs { + if request.suppress_contention_logs || !should_warn_lock_failure(&error) { tracing::debug!( resource = %request.resource, owner = %request.owner, @@ -493,15 +510,11 @@ impl DistributedLock { let remaining = request.acquire_timeout - elapsed; let mut attempt_request = request.clone(); - attempt_request.acquire_timeout = remaining; + attempt_request.acquire_timeout = self.lock_acquire_attempt_timeout(remaining); attempt_request.lock_id = LockId::new_unique(&request.resource); let result = self.acquire_lock_quorum_once(&attempt_request).await?; - if result.response.success - || !result.individual_locks.is_empty() - || !Self::is_retryable_lock_failure(&result.response) - || attempt >= LOCK_ACQUIRE_RETRY_ATTEMPTS - { + if result.response.success || !matches!(result.failure_kind, Some(LockAcquireFailureKind::RetryableContention)) { return Ok(result); } @@ -522,6 +535,42 @@ impl DistributedLock { })) } + fn lock_acquire_timeout_result(timeout: Duration) -> LockAcquireQuorumResult { + LockAcquireQuorumResult { + response: LockResponse::failure("Lock acquisition timeout", timeout), + individual_locks: Vec::new(), + failure_kind: Some(LockAcquireFailureKind::RetryableContention), + } + } + + fn lock_acquire_attempt_timeout_result( + timeout: Duration, + individual_locks: Vec<(LockId, Arc)>, + last_failure: Option, + ) -> LockAcquireQuorumResult { + let last_failure_kind = last_failure.as_deref().map(classify_lock_failure); + let failure_kind = + if let Some(kind @ (LockAcquireFailureKind::NonRetryable | LockAcquireFailureKind::UnrecoverableQuorum)) = + last_failure_kind + { + kind + } else { + return Self::lock_acquire_timeout_result(timeout); + }; + + let mut error = "Lock acquisition timeout".to_string(); + if let Some(last_failure) = last_failure { + error.push_str("; last failure: "); + error.push_str(&last_failure); + } + + LockAcquireQuorumResult { + response: LockResponse::failure(error, timeout), + individual_locks, + failure_kind: Some(failure_kind), + } + } + /// Quorum-based lock acquisition: success if at least the required quorum succeeds. /// Collects all individual lock_ids from successful clients and creates an aggregate lock_id. /// Returns the LockResponse with aggregate lock_id and individual lock mappings. @@ -532,8 +581,44 @@ impl DistributedLock { let fallback_lock_id = request.lock_id.clone(); let mut last_failure = None; let mut hard_failures = 0usize; + let start = tokio::time::Instant::now(); + + while !pending.is_empty() { + let remaining = request.acquire_timeout.saturating_sub(start.elapsed()); + if remaining.is_zero() { + Self::spawn_release_cleanup(individual_locks.clone(), "distributed_lock_attempt_timeout"); + Self::spawn_pending_cleanup( + pending, + self.clients.clone(), + fallback_lock_id.clone(), + "distributed_lock_attempt_timeout_cleanup", + ); + return Ok(Self::lock_acquire_attempt_timeout_result( + request.acquire_timeout, + individual_locks, + last_failure, + )); + } + + let join_result = match tokio::time::timeout(remaining, pending.join_next()).await { + Ok(Some(join_result)) => join_result, + Ok(None) => break, + Err(_) => { + Self::spawn_release_cleanup(individual_locks.clone(), "distributed_lock_attempt_timeout"); + Self::spawn_pending_cleanup( + pending, + self.clients.clone(), + fallback_lock_id.clone(), + "distributed_lock_attempt_timeout_cleanup", + ); + return Ok(Self::lock_acquire_attempt_timeout_result( + request.acquire_timeout, + individual_locks, + last_failure, + )); + } + }; - while let Some(join_result) = pending.join_next().await { match join_result { Ok((idx, Ok(resp))) => { if resp.success { @@ -550,7 +635,7 @@ impl DistributedLock { } } else { let error = resp.error.unwrap_or_else(|| "unknown error".to_string()); - if is_remote_lock_rpc_failure(&error) { + if is_remote_lock_rpc_failure(&error) && !is_remote_lock_rpc_timeout(&error) { hard_failures += 1; } self.log_failed_lock_response(request, idx, error.clone()); @@ -710,7 +795,10 @@ fn record_lock_held_release(lock_type: LockType) { #[cfg(test)] mod tests { - use super::{DistributedLock, is_remote_lock_rpc_failure}; + use super::{ + DistributedLock, LOCK_ACQUIRE_ATTEMPT_TIMEOUT, LockAcquireFailureKind, is_remote_lock_rpc_failure, + should_warn_lock_failure, + }; use crate::{LockError, LockId, LockInfo, LockRequest, LockResponse, LockStats, LockType, ObjectKey, client::LockClient}; use std::{ collections::{HashMap, VecDeque}, @@ -784,6 +872,61 @@ mod tests { } } + #[derive(Debug)] + struct TimeoutRecordingClient { + seen_timeouts: Arc>>, + } + + impl TimeoutRecordingClient { + fn new(seen_timeouts: Arc>>) -> Self { + Self { seen_timeouts } + } + + fn into_client(self) -> Arc { + Arc::new(self) + } + } + + #[async_trait::async_trait] + impl LockClient for TimeoutRecordingClient { + async fn acquire_lock(&self, request: &LockRequest) -> crate::Result { + self.seen_timeouts.lock().unwrap().push(request.acquire_timeout); + Ok(LockResponse::failure("Lock acquisition timeout", Duration::ZERO)) + } + + async fn release(&self, _lock_id: &LockId) -> crate::Result { + Ok(false) + } + + async fn refresh(&self, _lock_id: &LockId) -> crate::Result { + Ok(false) + } + + async fn force_release(&self, _lock_id: &LockId) -> crate::Result { + Ok(false) + } + + 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 + } + } + #[derive(Debug, Clone, Copy)] enum AcquirePlan { Success { delay: Duration }, @@ -894,6 +1037,15 @@ mod tests { assert!(!is_remote_lock_rpc_failure("acquired lock failed for other reason")); } + #[test] + fn test_should_warn_lock_failure() { + assert!(should_warn_lock_failure("Remote lock RPC failed: backend unreachable")); + assert!(should_warn_lock_failure("permission denied")); + assert!(should_warn_lock_failure("Unrecoverable quorum failure: 1/3 required")); + assert!(!should_warn_lock_failure("Remote lock RPC timed out: RPC timed out after 50ms")); + assert!(!should_warn_lock_failure("Lock acquisition timeout")); + } + #[tokio::test] async fn acquire_guard_returns_quorum_error_when_rpc_failures_make_quorum_impossible() { let clients: Vec> = vec![ @@ -926,7 +1078,7 @@ mod tests { } #[tokio::test] - async fn acquire_guard_returns_quorum_error_when_rpc_timeouts_make_quorum_impossible() { + async fn acquire_guard_retries_rpc_timeouts_until_caller_timeout() { let clients: Vec> = vec![ ResponseClient::new(LockResponse::failure( "Remote lock RPC timed out: RPC timed out after 50ms", @@ -947,19 +1099,16 @@ mod tests { ]; 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(1)); + .with_acquire_timeout(Duration::from_millis(900)); + let started = tokio::time::Instant::now(); let result = lock.acquire_guard(&request).await; + let elapsed = started.elapsed(); + assert!(matches!(result, Ok(None)), "unexpected result: {result:?}"); assert!( - matches!( - result, - Err(LockError::QuorumNotReached { - required: 3, - achieved: 0 - }) - ), - "unexpected result: {result:?}" + elapsed >= Duration::from_millis(250), + "remote RPC timeouts should be retried within the caller timeout, got {elapsed:?}" ); } @@ -977,7 +1126,7 @@ mod tests { ]; 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)); + .with_acquire_timeout(Duration::from_millis(120)); let started = tokio::time::Instant::now(); let result = lock.acquire_guard(&request).await; @@ -989,6 +1138,90 @@ mod tests { ); } + #[tokio::test] + async fn acquire_guard_retries_attempt_timeout_when_hard_failure_does_not_preclude_quorum() { + let clients: Vec> = vec![ + Arc::new(SequencedClient::new( + vec![ + AcquirePlan::Failure { + error: "Remote lock RPC failed: node unavailable", + delay: Duration::ZERO, + }, + AcquirePlan::Success { delay: Duration::ZERO }, + ], + Arc::new(Mutex::new(Vec::new())), + )), + Arc::new(SequencedClient::new( + vec![ + AcquirePlan::Success { + delay: Duration::from_millis(1200), + }, + AcquirePlan::Success { delay: Duration::ZERO }, + ], + Arc::new(Mutex::new(Vec::new())), + )), + Arc::new(SequencedClient::new( + vec![ + AcquirePlan::Success { + delay: Duration::from_millis(1200), + }, + AcquirePlan::Success { delay: Duration::ZERO }, + ], + Arc::new(Mutex::new(Vec::new())), + )), + Arc::new(SequencedClient::new( + vec![ + AcquirePlan::Success { + delay: Duration::from_millis(1200), + }, + AcquirePlan::Success { delay: Duration::ZERO }, + ], + Arc::new(Mutex::new(Vec::new())), + )), + ]; + 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("attempt timeout with possible quorum should retry") + .expect("second attempt should acquire quorum"); + + assert_eq!(guard.entries.len(), 3); + drop(guard); + } + + #[tokio::test] + async fn acquire_quorum_timeout_preserves_non_retryable_failure() { + let clients: Vec> = vec![ + ResponseClient::new(LockResponse::failure("permission denied", Duration::ZERO)).into_client(), + ResponseClient::new(LockResponse::failure("lock already held", Duration::ZERO)) + .with_delay(Duration::from_secs(1)) + .into_client(), + ResponseClient::new(LockResponse::failure("lock already held", Duration::ZERO)) + .with_delay(Duration::from_secs(1)) + .into_client(), + ]; + 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(120)); + + let result = lock.acquire_lock_quorum_with_retry(&request).await.unwrap(); + + assert!(matches!(result.failure_kind, Some(LockAcquireFailureKind::NonRetryable))); + assert!( + result + .response + .error + .as_deref() + .is_some_and(|error| error.contains("permission denied")), + "unexpected response: {:?}", + result.response + ); + } + #[tokio::test] async fn acquire_guard_retries_transient_timeout_before_quorum() { let clients: Vec> = vec![ @@ -1028,6 +1261,113 @@ mod tests { ); } + #[tokio::test] + async fn acquire_guard_uses_bounded_distributed_attempt_timeout() { + let seen_timeouts = Arc::new(Mutex::new(Vec::new())); + let clients: Vec> = (0..4) + .map(|_| TimeoutRecordingClient::new(seen_timeouts.clone()).into_client()) + .collect(); + 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(1500)); + + let result = lock.acquire_guard(&request).await; + + assert!(matches!(result, Ok(None)), "unexpected result: {result:?}"); + let seen_timeouts = seen_timeouts.lock().unwrap(); + assert!(!seen_timeouts.is_empty(), "expected lock clients to observe acquire timeouts"); + assert!( + seen_timeouts.iter().all(|timeout| *timeout <= LOCK_ACQUIRE_ATTEMPT_TIMEOUT), + "distributed attempts should be bounded by {LOCK_ACQUIRE_ATTEMPT_TIMEOUT:?}, saw {seen_timeouts:?}" + ); + } + + #[tokio::test] + async fn acquire_guard_enforces_attempt_timeout_when_clients_ignore_budget() { + let clients: Vec> = (0..4) + .map(|_| { + ResponseClient::new(LockResponse::failure("lock already held", Duration::ZERO)) + .with_delay(Duration::from_secs(5)) + .into_client() + }) + .collect(); + 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)); + + let result = tokio::time::timeout(Duration::from_millis(800), lock.acquire_guard(&request)).await; + + assert!(matches!(result, Ok(Ok(None))), "unexpected result: {result:?}"); + } + + #[tokio::test] + async fn acquire_guard_retries_partial_quorum_after_rollback() { + let seen_ids = Arc::new(Mutex::new(Vec::new())); + let clients: Vec> = vec![ + Arc::new(SequencedClient::new( + vec![ + AcquirePlan::Success { 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(), + )), + 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(), + )), + ]; + 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("partial quorum retry should not fail") + .expect("second attempt should reach quorum"); + + let unique_id_count = seen_ids + .lock() + .unwrap() + .iter() + .fold(Vec::::new(), |mut uuids, lock_id| { + if !uuids.iter().any(|uuid| uuid == &lock_id.uuid) { + uuids.push(lock_id.uuid.clone()); + } + uuids + }) + .len(); + + assert!(unique_id_count >= 2, "partial quorum should retry with a fresh lock id"); + drop(guard); + } + #[tokio::test] async fn acquire_guard_returns_timeout_when_quorum_remains_contended() { let clients: Vec> = vec![ diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 82d09a478..f54ad3b3c 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -2573,6 +2573,8 @@ mod tests { } #[tokio::test] + #[ignore = "requires isolated global object layer state"] + #[serial_test::serial] async fn test_signal_service_refresh_config_requires_object_layer() { let service = create_test_node_service(); @@ -2593,6 +2595,8 @@ mod tests { } #[tokio::test] + #[ignore = "requires isolated global object layer state"] + #[serial_test::serial] async fn test_signal_service_reload_dynamic_requires_object_layer() { let service = create_test_node_service();