diff --git a/crates/e2e_test/src/namespace_lock_quorum_test.rs b/crates/e2e_test/src/namespace_lock_quorum_test.rs index fc3c6b0fb..9b9e1d53f 100644 --- a/crates/e2e_test/src/namespace_lock_quorum_test.rs +++ b/crates/e2e_test/src/namespace_lock_quorum_test.rs @@ -120,3 +120,115 @@ async fn test_concurrent_cluster_overwrites_do_not_fail_namespace_lock_quorum() clients[0].delete_object().bucket(BUCKET).key(KEY).send().await?; Ok(()) } + +/// Regression test: concurrent PUTs to the same key must return 503 (ServiceUnavailable) +/// on lock contention, never 500 (InternalError). +/// +/// Before the fix, `map_namespace_lock_error` wrapped lock timeout/conflict errors as +/// `StorageError::other(...)` → `StorageError::Io(...)`, which fell through to +/// `S3ErrorCode::InternalError` (500) in the error mapping. +#[tokio::test] +#[serial] +async fn test_concurrent_put_same_key_never_returns_500() -> TestResult { + crate::common::init_logging(); + info!("Starting concurrent PUT 500 regression test"); + + let mut cluster = RustFSTestClusterEnvironment::new(4).await?; + // Short lock timeout to trigger contention errors quickly + cluster.set_env("RUSTFS_OBJECT_LOCK_ACQUIRE_TIMEOUT", "3"); + cluster.start().await?; + cluster.create_test_bucket(BUCKET).await?; + + let clients = cluster.create_all_clients()?; + let writer_count = clients.len() * 4; // 16 writers for heavy contention + let barrier = Arc::new(Barrier::new(writer_count)); + + // Seed initial object + let first_payload = b"initial object for 500 regression".to_vec(); + put_object(clients[0].clone(), first_payload, 0).await?; + + let err_500_count = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let err_503_count = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let unexpected_err_count = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let ok_count = Arc::new(std::sync::atomic::AtomicU64::new(0)); + + let mut handles = Vec::with_capacity(writer_count); + for writer_id in 0..writer_count { + let client = clients[writer_id % clients.len()].clone(); + let barrier = barrier.clone(); + let payload = format!("payload from writer {writer_id:02}").into_bytes(); + let err_500 = err_500_count.clone(); + let err_503 = err_503_count.clone(); + let unexpected_err = unexpected_err_count.clone(); + let ok = ok_count.clone(); + + handles.push(tokio::spawn(async move { + barrier.wait().await; + match client + .put_object() + .bucket(BUCKET) + .key(KEY) + .body(Bytes::from(payload).into()) + .send() + .await + { + Ok(_) => { + ok.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } + Err(err) => match &err { + SdkError::ServiceError(service_err) => { + let code = service_err.err().meta().code().unwrap_or(""); + if code == "500" || code == "InternalError" { + err_500.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + warn!("writer {writer_id} returned 500: {code}"); + } else if code == "503" || code == "ServiceUnavailable" { + err_503.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } else { + unexpected_err.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + warn!("writer {writer_id} returned unexpected error: {code}"); + } + } + other => { + unexpected_err.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + warn!("writer {writer_id} returned SDK error: {other:?}"); + } + }, + } + })); + } + + for handle in handles { + if let Err(err) = handle.await { + unexpected_err_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + warn!("writer task join failed: {err}"); + } + } + + let ok = ok_count.load(std::sync::atomic::Ordering::Relaxed); + let err_503 = err_503_count.load(std::sync::atomic::Ordering::Relaxed); + let err_500 = err_500_count.load(std::sync::atomic::Ordering::Relaxed); + let unexpected_err = unexpected_err_count.load(std::sync::atomic::Ordering::Relaxed); + let total = ok + err_503 + err_500 + unexpected_err; + + info!("Concurrent PUT 500 regression: total={total}, ok={ok}, 503={err_503}, 500={err_500}, unexpected={unexpected_err}"); + + assert_eq!( + total, writer_count as u64, + "every concurrent PUT writer must be classified as success, 503, 500, or unexpected error" + ); + + assert_eq!( + unexpected_err, 0, + "Concurrent PUTs to the same key must only succeed or return 503. \ + Got {unexpected_err} unexpected errors out of {total} requests. 503 count: {err_503}, ok: {ok}" + ); + + assert_eq!( + err_500, 0, + "Concurrent PUTs to the same key must NEVER return 500 InternalError. \ + Got {err_500} out of {total} requests. 503 count: {err_503}, ok: {ok}" + ); + + clients[0].delete_object().bucket(BUCKET).key(KEY).send().await?; + Ok(()) +} diff --git a/crates/ecstore/src/set_disk/lock.rs b/crates/ecstore/src/set_disk/lock.rs index aa536c426..bd51d05f9 100644 --- a/crates/ecstore/src/set_disk/lock.rs +++ b/crates/ecstore/src/set_disk/lock.rs @@ -67,10 +67,7 @@ impl SetDisks { achieved, } } - other => StorageError::other(format!( - "Failed to acquire {mode} lock: {}", - self.format_lock_error_from_error(bucket, object, mode, &other) - )), + other => StorageError::Lock(other), } } diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 4828453dd..a20103ee7 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -81,7 +81,7 @@ impl ECStore { achieved, } } - other => StorageError::other(format!("make_bucket: failed to acquire write lock on {bucket}: {other}")), + other => StorageError::Lock(other), })?, ) } else { @@ -185,7 +185,7 @@ impl ECStore { achieved, } } - other => StorageError::other(format!("delete_bucket: failed to acquire write lock on {bucket}: {other}")), + other => StorageError::Lock(other), })?, ) } else { diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index acc9d1240..b04a4acb0 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -243,7 +243,7 @@ impl ECStore { required, achieved, }, - other => StorageError::other(format!("Failed to acquire {mode} lock on {bucket}/{object}: {other}")), + other => StorageError::Lock(other), } } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index de40a4dfc..e17c6f346 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -723,7 +723,7 @@ fn copy_namespace_lock_error(bucket: &str, object: &str, mode: &'static str, err required, achieved, }, - other => StorageError::other(format!("Failed to acquire {mode} lock on {bucket}/{object}: {other}")), + other => StorageError::Lock(other), } } diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index d5268e5c6..39ed78ae5 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -243,6 +243,7 @@ impl From for ApiError { StorageError::StorageFull => S3ErrorCode::ServiceUnavailable, StorageError::SlowDown => S3ErrorCode::SlowDown, StorageError::NamespaceLockQuorumUnavailable { .. } => S3ErrorCode::ServiceUnavailable, + StorageError::Lock(_) => S3ErrorCode::ServiceUnavailable, StorageError::DecommissionNotStarted => S3ErrorCode::InvalidRequest, StorageError::DecommissionAlreadyRunning => S3ErrorCode::InvalidRequest, StorageError::RebalanceAlreadyRunning => S3ErrorCode::InvalidRequest,