mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 20:59:05 +00:00
acce8b2253
* fix(lock): let waiters hear releases and let acquisition succeed past registered waiters Same-key write contention scaled superlinearly with writer count: 8 concurrent conditional PUTs on one key cost ~340-460 ms, 16 cost ~700 ms, 32 cost ~5 s, against ~4 ms per uncontended write and ~10 ms actual lock holds (measured via RUSTFS_OBJECT_LOCK_DIAG at 1 ms thresholds). Outcomes were always correct; the cost was pure waiting. Two coupled defects in fast_lock caused it: 1. The slow path's early retries slept without subscribing to anything. notify_writer()/notify_readers() are gated on the waiter counters, which a sleeper never increments, so a release during the backoff woke nobody. The lock sat free while every loser slept out its full backoff, and the ladder compounded: successive acquires landed at the cumulative ladder offsets (10+20+40+80+100... ms). 2. try_acquire_exclusive demanded the entire packed state word be zero, including the readers_waiting/writers_waiting counter bits. A lock with registered waiters could be acquired by no one - including the waiters themselves, each blocked by the others' registration - so contended acquisition only succeeded in windows where every waiter happened to be unregistered. This is also why (1) could not be fixed by simply registering the sleepers: registration alone deadlocks acquisition until the acquire deadline. try_acquire_shared already masks correctly and preserves the counter bits in its CAS; the exclusive path now mirrors it. The fix: mask the acquisition CAS to ownership bits only (writer flag, active readers), and turn the early-retry sleep into a notification wait bounded by the same backoff, so a release wakes a waiter immediately while the bound still protects against lost or stolen wakeups exactly as NOTIFY_WAIT_CAP does for the post-retry wait. With both changes, 8 concurrent same-key CAS writers resolve in 17-29 ms (was 340-460 ms) and 32 resolve in 20-53 ms (was ~5 s), with per-racer cost now decreasing in N. Outcomes remain exactly one winner, N-1 precondition failures, zero errors at every width. cargo test -p rustfs-lock passes 113/113 at pristine-parity runtime, including test_concurrent_write_lock_contention, which previously only passed because sleepers were invisible to it. * test(lock): pin both halves of the waiter-starvation fix The fix commit touched only production files, so reverting either half left the suite green: test_concurrent_write_lock_contention only waits for five writers to finish and never asserts that acquisition happens before the backoff ladder runs out. Three tests, one per revert: * exclusive_acquisition_ignores_registered_waiters (state.rs) - a free lock with registered waiters must be acquirable, and the CAS must preserve the counters. Fails against the all-zero `expected`. * early_retry_registers_as_waiter (shard.rs) - a waiter in the early-retry backoff must appear in the writer waiter count within the ~750ms early-retry phase, since notify_writer/notify_readers are gated on those counters. Fails against a bare `sleep`, which registers nowhere. * contended_writers_drain_promptly_after_release (tests.rs) - 16 same-key writers, all registered behind one holder, must drain within 1s of the release rather than sit out their 5s acquire deadlines. Fails against the all-zero `expected` end to end. Wakeup latency is deliberately not asserted anywhere. NOTIFY_POOL is a process-global of 128 Notify slots shared by every lock, so a waiter in a concurrently-running test can consume another's notify_one and push it to the end of its rung: a 24-key latency probe measured ~150us in isolation and ~92ms - a full unexpired rung - alongside the existing 64-key missed-wakeup test. That is the stolen wakeup NOTIFY_WAIT_CAP already exists to bound, and it makes any in-suite latency budget flaky. cargo test -p rustfs-lock: 116/116. Signed-off-by: Miguel Amador <miguel@amador.one> --------- Signed-off-by: Miguel Amador <miguel@amador.one>
681 lines
26 KiB
Rust
681 lines
26 KiB
Rust
// Copyright 2024 RustFS Team
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
#[cfg(test)]
|
|
mod fast_lock_tests {
|
|
use crate::LockError;
|
|
use crate::fast_lock::types::{LockConfig, LockMode, LockPriority, LockResult, ObjectKey, ObjectLockRequest};
|
|
use crate::fast_lock::{DEFAULT_SHARD_COUNT, FastObjectLockManager};
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, Instant};
|
|
use tokio::time::sleep;
|
|
|
|
/// Helper function to create a test lock manager
|
|
fn create_test_manager() -> FastObjectLockManager {
|
|
let config = LockConfig {
|
|
shard_count: 4, // Use smaller shard count for tests
|
|
default_lock_timeout: Duration::from_secs(30),
|
|
default_acquire_timeout: Duration::from_secs(5),
|
|
..LockConfig::default()
|
|
};
|
|
FastObjectLockManager::with_config(config)
|
|
}
|
|
|
|
#[test]
|
|
fn try_with_config_returns_error_for_invalid_shard_count() {
|
|
let config = LockConfig {
|
|
shard_count: 3,
|
|
..LockConfig::default()
|
|
};
|
|
|
|
let err = FastObjectLockManager::try_with_config(config)
|
|
.expect_err("non-power-of-two shard counts should return an explicit configuration error");
|
|
|
|
assert!(matches!(err, LockError::Configuration { .. }));
|
|
assert!(err.to_string().contains("shard count must be a non-zero power of 2"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn with_config_falls_back_to_default_for_invalid_shard_count() {
|
|
let config = LockConfig {
|
|
shard_count: 3,
|
|
..LockConfig::default()
|
|
};
|
|
|
|
let manager = FastObjectLockManager::with_config(config);
|
|
|
|
assert_eq!(manager.shards.len(), DEFAULT_SHARD_COUNT);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_basic_write_lock_acquire_release() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner: Arc<str> = Arc::from("test-owner");
|
|
|
|
// Acquire write lock
|
|
let mut guard = manager
|
|
.acquire_write_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire write lock");
|
|
|
|
// Verify guard properties
|
|
assert_eq!(guard.key(), &key);
|
|
assert_eq!(guard.mode(), LockMode::Exclusive);
|
|
assert_eq!(guard.owner(), &owner);
|
|
assert!(!guard.is_released());
|
|
|
|
// Manually release lock
|
|
assert!(guard.release(), "Should release lock successfully");
|
|
assert!(guard.is_released(), "Guard should be marked as released");
|
|
|
|
// Try to acquire again - should succeed
|
|
let guard2 = manager
|
|
.acquire_write_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire write lock again after release");
|
|
drop(guard2);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_basic_read_lock_acquire_release() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner: Arc<str> = Arc::from("test-owner");
|
|
|
|
// Acquire read lock
|
|
let mut guard = manager
|
|
.acquire_read_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire read lock");
|
|
|
|
// Verify guard properties
|
|
assert_eq!(guard.key(), &key);
|
|
assert_eq!(guard.mode(), LockMode::Shared);
|
|
assert_eq!(guard.owner(), &owner);
|
|
assert!(!guard.is_released());
|
|
|
|
// Manually release lock
|
|
assert!(guard.release(), "Should release lock successfully");
|
|
assert!(guard.is_released(), "Guard should be marked as released");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_lock_auto_release_on_drop() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner1: Arc<str> = Arc::from("owner1");
|
|
let owner2: Arc<str> = Arc::from("owner2");
|
|
|
|
// Acquire lock and drop guard
|
|
{
|
|
let guard = manager
|
|
.acquire_write_lock(key.clone(), owner1.clone())
|
|
.await
|
|
.expect("Should acquire write lock");
|
|
assert!(!guard.is_released());
|
|
// Guard is dropped here, lock should be automatically released
|
|
}
|
|
|
|
// Wait a bit to ensure cleanup
|
|
sleep(Duration::from_millis(10)).await;
|
|
|
|
// Another owner should be able to acquire the lock
|
|
let guard2 = manager
|
|
.acquire_write_lock(key.clone(), owner2.clone())
|
|
.await
|
|
.expect("Should acquire write lock after previous guard dropped");
|
|
drop(guard2);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_multiple_read_locks() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner1: Arc<str> = Arc::from("owner1");
|
|
let owner2: Arc<str> = Arc::from("owner2");
|
|
let owner3: Arc<str> = Arc::from("owner3");
|
|
|
|
// Multiple read locks should be allowed
|
|
let mut guard1 = manager
|
|
.acquire_read_lock(key.clone(), owner1.clone())
|
|
.await
|
|
.expect("Should acquire first read lock");
|
|
|
|
let mut guard2 = manager
|
|
.acquire_read_lock(key.clone(), owner2.clone())
|
|
.await
|
|
.expect("Should acquire second read lock");
|
|
|
|
let mut guard3 = manager
|
|
.acquire_read_lock(key.clone(), owner3.clone())
|
|
.await
|
|
.expect("Should acquire third read lock");
|
|
|
|
// All guards should be valid
|
|
assert_eq!(guard1.mode(), LockMode::Shared);
|
|
assert_eq!(guard2.mode(), LockMode::Shared);
|
|
assert_eq!(guard3.mode(), LockMode::Shared);
|
|
|
|
// Release all
|
|
assert!(guard1.release());
|
|
assert!(guard2.release());
|
|
assert!(guard3.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_write_lock_excludes_read_lock() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let writer: Arc<str> = Arc::from("writer");
|
|
let reader: Arc<str> = Arc::from("reader");
|
|
|
|
// Acquire write lock
|
|
let mut write_guard = manager
|
|
.acquire_write_lock(key.clone(), writer.clone())
|
|
.await
|
|
.expect("Should acquire write lock");
|
|
|
|
// Try to acquire read lock - should timeout
|
|
let read_request =
|
|
ObjectLockRequest::new_read(key.clone(), reader.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(read_request).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"Read lock should timeout when write lock is held"
|
|
);
|
|
|
|
// Release write lock
|
|
assert!(write_guard.release());
|
|
|
|
// Now read lock should succeed
|
|
let mut read_guard = manager
|
|
.acquire_read_lock(key.clone(), reader.clone())
|
|
.await
|
|
.expect("Should acquire read lock after write lock released");
|
|
assert!(read_guard.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_read_lock_excludes_write_lock() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let reader: Arc<str> = Arc::from("reader");
|
|
let writer: Arc<str> = Arc::from("writer");
|
|
|
|
// Acquire read lock
|
|
let mut read_guard = manager
|
|
.acquire_read_lock(key.clone(), reader.clone())
|
|
.await
|
|
.expect("Should acquire read lock");
|
|
|
|
// Try to acquire write lock - should timeout
|
|
let write_request =
|
|
ObjectLockRequest::new_write(key.clone(), writer.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(write_request).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"Write lock should timeout when read lock is held"
|
|
);
|
|
|
|
// Release read lock
|
|
assert!(read_guard.release());
|
|
|
|
// Now write lock should succeed
|
|
let mut write_guard = manager
|
|
.acquire_write_lock(key.clone(), writer.clone())
|
|
.await
|
|
.expect("Should acquire write lock after read lock released");
|
|
assert!(write_guard.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_write_lock_excludes_write_lock() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner1: Arc<str> = Arc::from("owner1");
|
|
let owner2: Arc<str> = Arc::from("owner2");
|
|
|
|
// Acquire first write lock
|
|
let mut guard1 = manager
|
|
.acquire_write_lock(key.clone(), owner1.clone())
|
|
.await
|
|
.expect("Should acquire first write lock");
|
|
|
|
// Try to acquire second write lock - should timeout
|
|
let request2 = ObjectLockRequest::new_write(key.clone(), owner2.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(request2).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"Second write lock should timeout when first write lock is held"
|
|
);
|
|
|
|
// Release first lock
|
|
assert!(guard1.release());
|
|
|
|
// Now second write lock should succeed
|
|
let mut guard2 = manager
|
|
.acquire_write_lock(key.clone(), owner2.clone())
|
|
.await
|
|
.expect("Should acquire second write lock after first released");
|
|
assert!(guard2.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_same_owner_reentrant_write_lock() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner: Arc<str> = Arc::from("owner");
|
|
|
|
// Acquire first write lock
|
|
let mut guard1 = manager
|
|
.acquire_write_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire first write lock");
|
|
|
|
// Same owner trying to acquire again - should timeout (not reentrant)
|
|
let request2 = ObjectLockRequest::new_write(key.clone(), owner.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(request2).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"Same owner should not be able to acquire lock again (not reentrant)"
|
|
);
|
|
|
|
assert!(guard1.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_different_keys_no_conflict() {
|
|
let manager = create_test_manager();
|
|
let key1 = ObjectKey::new("bucket1", "object1");
|
|
let key2 = ObjectKey::new("bucket2", "object2");
|
|
let owner: Arc<str> = Arc::from("owner");
|
|
|
|
// Acquire locks on different keys simultaneously
|
|
let mut guard1 = manager
|
|
.acquire_write_lock(key1.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock on key1");
|
|
|
|
let mut guard2 = manager
|
|
.acquire_write_lock(key2.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock on key2");
|
|
|
|
// Both should be valid
|
|
assert_eq!(guard1.key(), &key1);
|
|
assert_eq!(guard2.key(), &key2);
|
|
|
|
assert!(guard1.release());
|
|
assert!(guard2.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_versioned_keys() {
|
|
let manager = create_test_manager();
|
|
let base_key = ObjectKey::new("bucket", "object");
|
|
let versioned_key = ObjectKey::with_version("bucket", "object", "v1");
|
|
let owner: Arc<str> = Arc::from("owner");
|
|
|
|
// Acquire lock on base key
|
|
let mut guard1 = manager
|
|
.acquire_write_lock(base_key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock on base key");
|
|
|
|
// Should be able to acquire lock on versioned key (different keys)
|
|
let mut guard2 = manager
|
|
.acquire_write_lock(versioned_key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock on versioned key");
|
|
|
|
assert_eq!(guard1.key(), &base_key);
|
|
assert_eq!(guard2.key(), &versioned_key);
|
|
|
|
assert!(guard1.release());
|
|
assert!(guard2.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_concurrent_read_locks() {
|
|
let manager = Arc::new(create_test_manager());
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let num_readers = 10;
|
|
|
|
let mut handles = Vec::new();
|
|
|
|
// Spawn multiple readers
|
|
for i in 0..num_readers {
|
|
let manager = manager.clone();
|
|
let key = key.clone();
|
|
let owner: Arc<str> = Arc::from(format!("reader-{}", i));
|
|
|
|
let handle = tokio::spawn(async move {
|
|
let mut guard = manager.acquire_read_lock(key, owner).await.expect("Should acquire read lock");
|
|
// Hold lock for a bit
|
|
sleep(Duration::from_millis(10)).await;
|
|
assert!(guard.release());
|
|
});
|
|
|
|
handles.push(handle);
|
|
}
|
|
|
|
// Wait for all readers
|
|
for handle in handles {
|
|
handle.await.expect("Reader task should complete");
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_concurrent_write_lock_contention() {
|
|
let manager = Arc::new(create_test_manager());
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let num_writers = 5;
|
|
|
|
let mut handles = Vec::new();
|
|
|
|
// Spawn multiple writers - they should serialize
|
|
for i in 0..num_writers {
|
|
let manager = manager.clone();
|
|
let key = key.clone();
|
|
let owner: Arc<str> = Arc::from(format!("writer-{}", i));
|
|
|
|
let handle = tokio::spawn(async move {
|
|
let mut guard = manager
|
|
.acquire_write_lock(key, owner)
|
|
.await
|
|
.expect("Should acquire write lock");
|
|
// Hold lock for a bit
|
|
sleep(Duration::from_millis(10)).await;
|
|
assert!(guard.release());
|
|
});
|
|
|
|
handles.push(handle);
|
|
}
|
|
|
|
// Wait for all writers - they should complete sequentially
|
|
for handle in handles {
|
|
handle.await.expect("Writer task should complete");
|
|
}
|
|
}
|
|
|
|
// Regression for the fast-lock lost-wakeup (backlog#853 follow-up).
|
|
//
|
|
// A waiter that enters the notification wait can miss its wakeup: the release
|
|
// path only calls `notify_one` when `writer_waiters > 0`, so if the holder
|
|
// releases in the narrow gap after the waiter's `try_acquire` fails but before
|
|
// it registers as a waiter, no notification (and no stored permit) is produced.
|
|
// The waiter then blocks until the acquire deadline even though the lock is free
|
|
// and stays free — surfacing as a spurious `LockResult::Timeout`.
|
|
//
|
|
// Each key here has one long holder plus one waiter that starts slightly later,
|
|
// so the waiter is pushed past the early backoff phase into the notification
|
|
// wait and there is no re-contention after the single release. With the fix
|
|
// (bounded notification wait + re-poll) the waiter acquires within ~50ms of the
|
|
// release; without it, the missed wakeup strands the waiter until timeout.
|
|
// Many independent keys make hitting the narrow race reliable.
|
|
#[tokio::test(flavor = "multi_thread")]
|
|
async fn write_lock_waiter_is_not_stranded_by_missed_wakeup() {
|
|
let manager = Arc::new(create_test_manager());
|
|
const KEYS: usize = 64;
|
|
// Long enough to push the waiter past the ~850ms backoff phase into the
|
|
// notification wait, where the missed-wakeup bug lives.
|
|
const HOLDER_HOLD: Duration = Duration::from_millis(950);
|
|
|
|
let mut handles = Vec::new();
|
|
for k in 0..KEYS {
|
|
let key = ObjectKey::new("bucket", format!("object-{k}"));
|
|
|
|
// Holder: grabs the lock immediately and holds it across the waiter's
|
|
// backoff-to-notification transition, then releases exactly once.
|
|
let holder_mgr = manager.clone();
|
|
let holder_key = key.clone();
|
|
handles.push(tokio::spawn(async move {
|
|
let mut guard = holder_mgr
|
|
.acquire_write_lock(holder_key, "holder")
|
|
.await
|
|
.expect("holder should acquire immediately");
|
|
sleep(HOLDER_HOLD).await;
|
|
assert!(guard.release());
|
|
}));
|
|
|
|
// Waiter: starts a touch later so the holder wins the lock first, then
|
|
// must survive the whole hold and acquire promptly after the release.
|
|
let waiter_mgr = manager.clone();
|
|
let waiter_key = key.clone();
|
|
handles.push(tokio::spawn(async move {
|
|
sleep(Duration::from_millis(10)).await;
|
|
let request = ObjectLockRequest::new_write(waiter_key, "waiter").with_acquire_timeout(Duration::from_secs(5));
|
|
let mut guard = waiter_mgr
|
|
.acquire_lock(request)
|
|
.await
|
|
.expect("waiter must acquire after the holder releases, not time out on a missed wakeup");
|
|
assert!(guard.release());
|
|
}));
|
|
}
|
|
|
|
for handle in handles {
|
|
handle.await.expect("lock task should not panic");
|
|
}
|
|
}
|
|
|
|
// Regression for the waiter-preserving exclusive CAS, end to end
|
|
// (rustfs#5657 same-key write contention).
|
|
//
|
|
// `try_acquire_exclusive` used to demand a fully-zero state word, which
|
|
// includes the waiting counters. Once the slow path's retries register as
|
|
// waiters — as they must, to hear a release — a lock with waiters becomes
|
|
// acquirable by no one, each waiter blocked by the others' registration, so
|
|
// every waiter here sits out its full acquire deadline instead of draining.
|
|
//
|
|
// The holder is held long enough for all waiters to be registered before
|
|
// the single release. After it they only serialize on each other, holding
|
|
// nothing, so they should drain in tens of milliseconds.
|
|
#[tokio::test(flavor = "multi_thread")]
|
|
async fn contended_writers_drain_promptly_after_release() {
|
|
let manager = Arc::new(create_test_manager());
|
|
let key = ObjectKey::new("bucket", "hot-object");
|
|
const WAITERS: usize = 16;
|
|
const HOLD: Duration = Duration::from_millis(300);
|
|
// Generous next to the ~20ms the fixed path needs for these waiters, and
|
|
// far below the 5s acquire deadline a zero-state CAS makes them all
|
|
// sit out.
|
|
const DRAIN_BUDGET: Duration = Duration::from_millis(1000);
|
|
|
|
let mut holder = manager
|
|
.acquire_write_lock(key.clone(), "holder")
|
|
.await
|
|
.expect("holder should acquire immediately");
|
|
|
|
let mut handles = Vec::new();
|
|
for i in 0..WAITERS {
|
|
let manager = manager.clone();
|
|
let key = key.clone();
|
|
handles.push(tokio::spawn(async move {
|
|
let request =
|
|
ObjectLockRequest::new_write(key, format!("waiter-{i}")).with_acquire_timeout(Duration::from_secs(5));
|
|
let mut guard = manager
|
|
.acquire_lock(request)
|
|
.await
|
|
.expect("every waiter must acquire once the holder releases");
|
|
assert!(guard.release());
|
|
}));
|
|
}
|
|
|
|
// Let every waiter fail its fast path and register in the retry ladder.
|
|
sleep(HOLD).await;
|
|
|
|
let released_at = Instant::now();
|
|
assert!(holder.release());
|
|
|
|
for handle in handles {
|
|
handle.await.expect("waiter task should not panic");
|
|
}
|
|
|
|
let drain = released_at.elapsed();
|
|
assert!(
|
|
drain < DRAIN_BUDGET,
|
|
"{WAITERS} waiters took {drain:?} to drain after the release (budget {DRAIN_BUDGET:?}) - \
|
|
registered waiters are blocking acquisition"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_lock_timeout() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner1: Arc<str> = Arc::from("owner1");
|
|
let owner2: Arc<str> = Arc::from("owner2");
|
|
|
|
// Acquire first lock
|
|
let mut guard1 = manager
|
|
.acquire_write_lock(key.clone(), owner1.clone())
|
|
.await
|
|
.expect("Should acquire first lock");
|
|
|
|
// Try to acquire with short timeout - should timeout
|
|
let request = ObjectLockRequest::new_write(key.clone(), owner2.clone()).with_acquire_timeout(Duration::from_millis(50));
|
|
let result = manager.acquire_lock(request).await;
|
|
assert!(matches!(result, Err(LockResult::Timeout)), "Should timeout when lock is held");
|
|
|
|
assert!(guard1.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_lock_priority() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let normal_owner: Arc<str> = Arc::from("normal");
|
|
let high_owner: Arc<str> = Arc::from("high");
|
|
|
|
// Acquire normal priority lock
|
|
let normal_request = ObjectLockRequest::new_write(key.clone(), normal_owner.clone())
|
|
.with_priority(LockPriority::Normal)
|
|
.with_acquire_timeout(Duration::from_secs(1));
|
|
let mut normal_guard = manager
|
|
.acquire_lock(normal_request)
|
|
.await
|
|
.expect("Should acquire normal priority lock");
|
|
|
|
// Try high priority lock - should still timeout (write locks are exclusive)
|
|
let high_request = ObjectLockRequest::new_write(key.clone(), high_owner.clone())
|
|
.with_priority(LockPriority::High)
|
|
.with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(high_request).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"High priority write lock should still timeout when normal write lock is held"
|
|
);
|
|
|
|
assert!(normal_guard.release());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_double_release() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner: Arc<str> = Arc::from("owner");
|
|
|
|
let mut guard = manager
|
|
.acquire_write_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock");
|
|
|
|
// First release should succeed
|
|
assert!(guard.release(), "First release should succeed");
|
|
assert!(guard.is_released(), "Guard should be marked as released");
|
|
|
|
// Second release should fail
|
|
assert!(!guard.release(), "Second release should fail");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_lock_info() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let owner: Arc<str> = Arc::from("owner");
|
|
|
|
let mut guard = manager
|
|
.acquire_write_lock(key.clone(), owner.clone())
|
|
.await
|
|
.expect("Should acquire lock");
|
|
|
|
// Get lock info
|
|
let lock_info = guard.lock_info();
|
|
assert!(lock_info.is_some(), "Should have lock info");
|
|
if let Some(info) = lock_info {
|
|
assert_eq!(info.key, key);
|
|
assert_eq!(info.mode, LockMode::Exclusive);
|
|
assert_eq!(info.owner, owner);
|
|
}
|
|
|
|
// Release lock
|
|
assert!(guard.release());
|
|
|
|
// Lock info should be None after release
|
|
let lock_info_after = guard.lock_info();
|
|
assert!(lock_info_after.is_none(), "Lock info should be None after release");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_read_write_mixed_scenario() {
|
|
let manager = create_test_manager();
|
|
let key = ObjectKey::new("test-bucket", "test-object");
|
|
let reader1: Arc<str> = Arc::from("reader1");
|
|
let reader2: Arc<str> = Arc::from("reader2");
|
|
let writer: Arc<str> = Arc::from("writer");
|
|
|
|
// Acquire two read locks
|
|
let mut read_guard1 = manager
|
|
.acquire_read_lock(key.clone(), reader1.clone())
|
|
.await
|
|
.expect("Should acquire first read lock");
|
|
let mut read_guard2 = manager
|
|
.acquire_read_lock(key.clone(), reader2.clone())
|
|
.await
|
|
.expect("Should acquire second read lock");
|
|
|
|
// Writer should timeout
|
|
let write_request =
|
|
ObjectLockRequest::new_write(key.clone(), writer.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result = manager.acquire_lock(write_request).await;
|
|
assert!(
|
|
matches!(result, Err(LockResult::Timeout)),
|
|
"Write lock should timeout when read locks are held"
|
|
);
|
|
|
|
// Release one read lock
|
|
assert!(read_guard1.release());
|
|
|
|
// Writer should still timeout (other read lock still held)
|
|
let write_request2 =
|
|
ObjectLockRequest::new_write(key.clone(), writer.clone()).with_acquire_timeout(Duration::from_millis(100));
|
|
let result2 = manager.acquire_lock(write_request2).await;
|
|
assert!(
|
|
matches!(result2, Err(LockResult::Timeout)),
|
|
"Write lock should still timeout when read lock is held"
|
|
);
|
|
|
|
// Release second read lock
|
|
assert!(read_guard2.release());
|
|
|
|
// Now writer should succeed
|
|
let mut write_guard = manager
|
|
.acquire_write_lock(key.clone(), writer.clone())
|
|
.await
|
|
.expect("Should acquire write lock after all read locks released");
|
|
assert!(write_guard.release());
|
|
}
|
|
}
|