From 0cb9952aa0c8c1a034fb032bc96cba7b0b49950f Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 9 Aug 2026 12:26:09 +0800 Subject: [PATCH] fix(rpc): make authenticated file writes atomic (#5879) --- rustfs/src/storage/rpc/http_service.rs | 541 ++++++++++++++++++++++++- rustfs/src/storage/storage_api.rs | 15 +- 2 files changed, 536 insertions(+), 20 deletions(-) diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index e998ee8a6..3842cf3d2 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -16,9 +16,10 @@ use crate::server::RPC_PREFIX; use crate::storage::request_context::spawn_traced; use crate::storage::storage_api::DiskError; use crate::storage::storage_api::rpc_consumer::http_service::{ - DEFAULT_READ_BUFFER_SIZE, NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, PUT_FILE_AUTH_TRAILER_LEN, - PUT_FILE_AUTH_V1, StorageDiskRpcExt as _, WALK_DIR_STREAM_COMPLETION_V1, WalkDirOptions, check_and_record_signed_rpc_nonce, - find_local_disk_by_ref, sign_ns_scanner_capability, verify_put_file_auth_trailer, verify_rpc_signature, + DEFAULT_READ_BUFFER_SIZE, DeleteOptions, DiskStore, NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, + PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_V1, StorageDiskRpcExt as _, WALK_DIR_STREAM_COMPLETION_V1, WalkDirOptions, + check_and_record_signed_rpc_nonce, find_local_disk_by_ref, sign_ns_scanner_capability, verify_put_file_auth_trailer, + verify_rpc_signature, }; #[cfg(test)] use crate::storage::storage_api::rpc_consumer::http_service::{ @@ -44,13 +45,14 @@ use s3s::dto::StreamingBlob; use serde::de::DeserializeOwned; use serde_urlencoded::from_bytes; use sha2::{Digest, Sha256}; +use std::collections::HashMap; use std::future::Future; use std::pin::Pin; -use std::sync::LazyLock; +use std::sync::{Arc, LazyLock, Weak}; use std::task::{Context, Poll}; use std::time::{Duration, Instant}; use tokio::io::{self, AsyncWriteExt}; -use tokio::sync::oneshot; +use tokio::sync::{Mutex, oneshot}; use tokio_util::{io::ReaderStream, sync::CancellationToken}; use tower::Service; use tracing::{error, warn}; @@ -79,6 +81,8 @@ static PUT_FILE_AUTH_STRICT: LazyLock = LazyLock::new(|| { rustfs_config::DEFAULT_INTERNODE_RPC_BODY_DIGEST_STRICT, ) }); +static PUT_FILE_TARGET_LOCKS: LazyLock>>>> = + LazyLock::new(|| parking_lot::Mutex::new(HashMap::new())); macro_rules! log_internode_rpc_response_failure { ($status:expr, $rpc_path:expr, $method:expr, $operation:expr, $reason:expr, $result:expr, Some(($context_key:expr, $context_value:expr)), Some($error_text:expr)) => {{ @@ -1142,6 +1146,105 @@ where Body::from(StreamingBlob::wrap(stream.chain(completion))) } +fn put_file_target_lock(disk: &DiskStore, query: &PutFileQuery) -> Arc> { + let key = format!("{:p}\0{}\0{}", Arc::as_ptr(disk), query.volume, query.path); + let mut locks = PUT_FILE_TARGET_LOCKS.lock(); + locks.retain(|_, lock| lock.strong_count() > 0); + if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) { + return lock; + } + + let lock = Arc::new(Mutex::new(())); + locks.insert(key, Arc::downgrade(&lock)); + lock +} + +async fn remove_put_file_staging(disk: &DiskStore, volume: &str, path: &str) -> Result<(), BoxError> { + disk.delete( + volume, + path, + DeleteOptions { + immediate: true, + ..Default::default() + }, + ) + .await + .map_err(Into::into) +} + +async fn write_authenticated_put_file( + disk: &DiskStore, + body: S, + query: &PutFileQuery, + nonce: uuid::Uuid, + url: &str, +) -> Result +where + S: futures::TryStream + Unpin, + E: Into, +{ + let target_lock = put_file_target_lock(disk, query); + let _target_guard = target_lock.lock_owned().await; + let staging_name = format!(".rustfs-put-{}", uuid::Uuid::new_v4()); + let staging_path = match query.path.rsplit_once('/') { + Some((parent, _)) => format!("{parent}/{staging_name}"), + None => staging_name, + }; + + let result = async { + let mut file = disk + .create_file("", &query.volume, &staging_path, query.size) + .await + .map_err(|err| ("create_staging", Box::new(err) as BoxError))?; + + if query.append { + match disk.read_file(&query.volume, &query.path).await { + Ok(mut source) => { + tokio::io::copy(&mut source, &mut file) + .await + .map_err(|err| ("copy_existing", Box::new(err) as BoxError))?; + } + Err(DiskError::FileNotFound) => {} + Err(err) => return Err(("read_existing", Box::new(err) as BoxError)), + } + } + + let copied = write_put_file_body_chunks_to_writer(body, &mut file, query, Some(nonce), url) + .await + .map_err(|err| ("write_body", Box::new(err) as BoxError))?; + if put_body_size_mismatch(query, copied) { + return Err(( + "verify_size", + Box::new(io::Error::new( + io::ErrorKind::UnexpectedEof, + format!("body size mismatch: expected {} bytes, received {copied}", query.size), + )) as BoxError, + )); + } + file.shutdown().await.map_err(|err| ("shutdown", Box::new(err) as BoxError))?; + drop(file); + + disk.rename_file(&query.volume, &staging_path, &query.volume, &query.path) + .await + .map_err(|err| ("publish", Box::new(err) as BoxError))?; + Ok(copied) + } + .await; + + match result { + Ok(copied) => Ok(copied), + Err((stage, primary)) => { + if let Err(cleanup) = remove_put_file_staging(disk, &query.volume, &staging_path).await { + return Err(( + stage, + Box::new(io::Error::other(format!("{primary}; staging cleanup failed: {cleanup}"))) as BoxError, + )); + } + Err((stage, primary)) + } + } +} + async fn handle_put_file(req: Request) -> Response { let method = req.method().clone(); let path = req.uri().path().to_string(); @@ -1202,6 +1305,22 @@ async fn handle_put_file(req: Request) -> Response { return response_with_status(StatusCode::BAD_REQUEST, "disk not found"); }; + if let Some(nonce) = auth_nonce { + let copied = match write_authenticated_put_file(&disk, req.into_body().into_data_stream(), &query, nonce, &url).await { + Ok(copied) => copied, + Err((stage, e)) => { + let message = put_file_stage_error_message(stage, &query, e.as_ref()); + log_internode_put_file_stage_failure!(stage, query, e); + return response_with_status(StatusCode::INTERNAL_SERVER_ERROR, message); + } + }; + record_put_file_metrics(copied); + return empty_ok(); + } + + let target_lock = put_file_target_lock(&disk, &query); + let _target_guard = target_lock.lock_owned().await; + let mut file = if query.append { match disk.append_file(&query.volume, &query.path).await { Ok(file) => file, @@ -1233,16 +1352,7 @@ async fn handle_put_file(req: Request) -> Response { } }; - let metrics = runtime_sources::current_internode_metrics(); - metrics.record_incoming_request_for_operation_and_backend( - INTERNODE_OPERATION_PUT_FILE_STREAM, - INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, - ); - metrics.record_recv_bytes_for_operation_and_backend( - INTERNODE_OPERATION_PUT_FILE_STREAM, - INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, - usize::try_from(copied).unwrap_or(usize::MAX), - ); + record_put_file_metrics(copied); if put_body_size_mismatch(&query, copied) { let err = std::io::Error::new( @@ -1263,6 +1373,19 @@ async fn handle_put_file(req: Request) -> Response { empty_ok() } +fn record_put_file_metrics(copied: u64) { + let metrics = runtime_sources::current_internode_metrics(); + metrics.record_incoming_request_for_operation_and_backend( + INTERNODE_OPERATION_PUT_FILE_STREAM, + INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, + ); + metrics.record_recv_bytes_for_operation_and_backend( + INTERNODE_OPERATION_PUT_FILE_STREAM, + INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, + usize::try_from(copied).unwrap_or(usize::MAX), + ); +} + async fn write_body_chunks_to_writer(body: S, writer: &mut W) -> io::Result where S: futures::TryStream + Unpin, @@ -1483,11 +1606,13 @@ mod tests { READ_FILE_STREAM_PATH, WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_PATH, WalkDirQuery, append_walk_dir_completion, internode_http_operation, internode_rpc_subsystem, is_internode_rpc_path, ns_scanner_response_body, ns_scanner_server_epoch_matches, put_body_size_mismatch, put_file_auth_nonce, put_file_stage_error_message, - read_file_body_stream, remote_scanner_claim_rejection, response_with_disk_error, supports_walk_dir_stream_completion, - validate_walk_dir_completion_request, verify_internode_rpc_signature, verify_ns_scanner_body_digest, - verify_walk_dir_body_digest, walk_dir_response_body, write_body_chunks_to_writer, write_put_file_body_chunks_to_writer, + put_file_target_lock, read_file_body_stream, remote_scanner_claim_rejection, response_with_disk_error, + supports_walk_dir_stream_completion, validate_walk_dir_completion_request, verify_internode_rpc_signature, + verify_ns_scanner_body_digest, verify_walk_dir_body_digest, walk_dir_response_body, write_authenticated_put_file, + write_body_chunks_to_writer, write_put_file_body_chunks_to_writer, }; use crate::storage::storage_api::ecstore_rpc::build_put_file_auth_trailer; + use crate::storage::storage_api::rpc_consumer::http_service::{DiskAPI as _, DiskOption, DiskStore, Endpoint, new_disk}; use bytes::Bytes; use http::{HeaderMap, HeaderValue, Method, StatusCode, Uri}; use http_body_util::BodyExt; @@ -1499,11 +1624,25 @@ mod tests { }; use sha2::Digest as _; use std::collections::HashMap; + use std::future::Future as _; + use std::pin::Pin; + use std::task::{Context, Poll}; use tokio::io; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio_stream::StreamExt; use tokio_stream::iter; + async fn new_put_file_test_disk() -> (DiskStore, tempfile::TempDir) { + let dir = tempfile::tempdir().expect("temp directory should be created"); + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("temp path should be utf8")).expect("disk endpoint should parse"); + let disk = new_disk(&endpoint, &DiskOption::default()) + .await + .expect("local disk should be created"); + disk.make_volume("bucket").await.expect("test volume should be created"); + (disk, dir) + } + struct DropNotifier(Option>); impl Drop for DropNotifier { @@ -1514,6 +1653,30 @@ mod tests { } } + struct GatedPutBody { + started: Option>, + release: tokio::sync::oneshot::Receiver<()>, + payload: Option, + } + + impl futures_util::Stream for GatedPutBody { + type Item = io::Result; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + if let Some(started) = self.started.take() { + let _ = started.send(()); + } + if self.payload.is_none() { + return Poll::Ready(None); + } + match Pin::new(&mut self.release).poll(cx) { + Poll::Ready(Ok(())) => Poll::Ready(self.payload.take().map(Ok)), + Poll::Ready(Err(_)) => Poll::Ready(Some(Err(io::Error::other("gated body release dropped")))), + Poll::Pending => Poll::Pending, + } + } + } + #[test] fn internode_rpc_path_matches_rpc_prefix() { assert!(is_internode_rpc_path("/rustfs/rpc/read_file_stream")); @@ -1917,6 +2080,348 @@ mod tests { assert_eq!(err.to_string(), "put_file auth trailer is incomplete"); } + #[tokio::test] + async fn authenticated_create_tamper_preserves_existing_local_file() { + let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string()); + let (disk, _dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + let nonce = uuid::Uuid::parse_str("55555555-6666-4777-8888-999999999999").expect("nonce"); + let url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=false&size=11&put_file_auth=digest-trailer-v1&put_file_nonce=55555555-6666-4777-8888-999999999999" + ); + let signed_digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"hello world"), hex_simd::AsciiCase::Lower); + let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &signed_digest).expect("trailer should build"); + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: false, + size: 11, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(nonce), + }; + let mut payload = b"hello worle".to_vec(); + payload.extend_from_slice(&trailer); + + let err = + write_authenticated_put_file(&disk, iter(vec![Ok::(Bytes::from(payload))]), &query, nonce, url) + .await + .expect_err("tampered body must not publish"); + + assert_eq!(err.0, "write_body"); + assert_eq!( + disk.read_all("bucket", "object/part.1") + .await + .expect("existing data should remain"), + Bytes::from_static(b"old-data") + ); + } + + #[tokio::test] + async fn authenticated_create_truncation_preserves_existing_local_file() { + let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string()); + let (disk, _dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + let nonce = uuid::Uuid::parse_str("66666666-7777-4888-8999-aaaaaaaaaaaa").expect("nonce"); + let url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=false&size=11&put_file_auth=digest-trailer-v1&put_file_nonce=66666666-7777-4888-8999-aaaaaaaaaaaa" + ); + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: false, + size: 11, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(nonce), + }; + + let err = write_authenticated_put_file( + &disk, + iter(vec![Ok::(Bytes::from_static(b"short"))]), + &query, + nonce, + url, + ) + .await + .expect_err("truncated body must not publish"); + + assert_eq!(err.0, "write_body"); + assert_eq!( + disk.read_all("bucket", "object/part.1") + .await + .expect("existing data should remain"), + Bytes::from_static(b"old-data") + ); + } + + #[tokio::test] + async fn authenticated_append_missing_trailer_preserves_existing_local_file() { + let (disk, _dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + let nonce = uuid::Uuid::parse_str("77777777-8888-4999-8aaa-bbbbbbbbbbbb").expect("nonce"); + let url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=77777777-8888-4999-8aaa-bbbbbbbbbbbb" + ); + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: true, + size: 0, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(nonce), + }; + + let err = write_authenticated_put_file( + &disk, + iter(vec![Ok::(Bytes::from_static(b"append-data"))]), + &query, + nonce, + url, + ) + .await + .expect_err("missing trailer must not publish"); + + assert_eq!(err.0, "write_body"); + assert_eq!( + disk.read_all("bucket", "object/part.1") + .await + .expect("existing data should remain"), + Bytes::from_static(b"old-data") + ); + } + + #[tokio::test] + async fn authenticated_append_publishes_existing_and_new_bytes() { + let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string()); + let (disk, _dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + let nonce = uuid::Uuid::parse_str("88888888-9999-4aaa-8bbb-cccccccccccc").expect("nonce"); + let url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=88888888-9999-4aaa-8bbb-cccccccccccc" + ); + let digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"append-data"), hex_simd::AsciiCase::Lower); + let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &digest).expect("trailer should build"); + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: true, + size: 0, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(nonce), + }; + let mut payload = b"append-data".to_vec(); + payload.extend_from_slice(&trailer); + + let copied = + write_authenticated_put_file(&disk, iter(vec![Ok::(Bytes::from(payload))]), &query, nonce, url) + .await + .expect("authenticated append should publish"); + + assert_eq!(copied, 11); + assert_eq!( + disk.read_all("bucket", "object/part.1") + .await + .expect("published data should be readable"), + Bytes::from_static(b"old-dataappend-data") + ); + } + + #[tokio::test] + async fn concurrent_authenticated_appends_preserve_both_payloads_without_staging() { + let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string()); + let (disk, dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + + let first_nonce = uuid::Uuid::parse_str("99999999-aaaa-4bbb-8ccc-dddddddddddd").expect("first nonce"); + let first_url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=99999999-aaaa-4bbb-8ccc-dddddddddddd" + ); + let first_digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"first"), hex_simd::AsciiCase::Lower); + let first_trailer = + build_put_file_auth_trailer(first_url, &Method::PUT, first_nonce, &first_digest).expect("first trailer should build"); + let mut first_payload = b"first".to_vec(); + first_payload.extend_from_slice(&first_trailer); + + let second_nonce = uuid::Uuid::parse_str("aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee").expect("second nonce"); + let second_url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee" + ); + let second_digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"second"), hex_simd::AsciiCase::Lower); + let second_trailer = build_put_file_auth_trailer(second_url, &Method::PUT, second_nonce, &second_digest) + .expect("second trailer should build"); + let mut second_payload = b"second".to_vec(); + second_payload.extend_from_slice(&second_trailer); + + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: true, + size: 0, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(first_nonce), + }; + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let (release_tx, release_rx) = tokio::sync::oneshot::channel(); + let first_disk = disk.clone(); + let first_query = query.clone(); + let first = tokio::spawn(async move { + write_authenticated_put_file( + &first_disk, + GatedPutBody { + started: Some(started_tx), + release: release_rx, + payload: Some(Bytes::from(first_payload)), + }, + &first_query, + first_nonce, + first_url, + ) + .await + }); + started_rx.await.expect("first append should start reading its body"); + + let second_disk = disk.clone(); + let mut second_query = query; + second_query.put_file_nonce = Some(second_nonce); + let second = tokio::spawn(async move { + write_authenticated_put_file( + &second_disk, + iter(vec![Ok::(Bytes::from(second_payload))]), + &second_query, + second_nonce, + second_url, + ) + .await + }); + tokio::task::yield_now().await; + assert!(!second.is_finished(), "second append must wait while the first holds the target lock"); + + release_tx.send(()).expect("first append should still be waiting"); + assert_eq!( + first + .await + .expect("first append task should join") + .expect("first append should succeed"), + 5 + ); + assert_eq!( + second + .await + .expect("second append task should join") + .expect("second append should succeed"), + 6 + ); + let final_data = disk + .read_all("bucket", "object/part.1") + .await + .expect("final data should be readable"); + assert!( + final_data == Bytes::from_static(b"old-datafirstsecond") || final_data == Bytes::from_static(b"old-datasecondfirst"), + "both complete append payloads must be preserved: {final_data:?}" + ); + + let object_dir = dir.path().join("bucket/object"); + let entries = std::fs::read_dir(object_dir).expect("object directory should be readable"); + assert!( + entries + .map(|entry| entry.expect("directory entry should be readable").file_name()) + .all(|name| !name.to_string_lossy().starts_with(".rustfs-put-")), + "successful concurrent appends must not leave staging files" + ); + } + + #[tokio::test] + async fn legacy_and_authenticated_appends_share_the_target_lock() { + let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string()); + let (disk, _dir) = new_put_file_test_disk().await; + disk.write_all("bucket", "object/part.1", Bytes::from_static(b"old-data")) + .await + .expect("existing data should be written"); + let nonce = uuid::Uuid::parse_str("bbbbbbbb-cccc-4ddd-8eee-ffffffffffff").expect("nonce"); + let url = concat!( + "/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1", + "&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=bbbbbbbb-cccc-4ddd-8eee-ffffffffffff" + ); + let digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"authenticated"), hex_simd::AsciiCase::Lower); + let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &digest).expect("trailer should build"); + let mut payload = b"authenticated".to_vec(); + payload.extend_from_slice(&trailer); + let query = PutFileQuery { + disk: "disk-a".to_string(), + volume: "bucket".to_string(), + path: "object/part.1".to_string(), + append: true, + size: 0, + put_file_auth: Some("digest-trailer-v1".to_string()), + put_file_nonce: Some(nonce), + }; + + let legacy_lock = put_file_target_lock(&disk, &query); + let (legacy_started_tx, legacy_started_rx) = tokio::sync::oneshot::channel(); + let (legacy_release_tx, legacy_release_rx) = tokio::sync::oneshot::channel(); + let legacy_disk = disk.clone(); + let legacy = tokio::spawn(async move { + let _guard = legacy_lock.lock_owned().await; + legacy_started_tx.send(()).expect("test should observe legacy lock"); + legacy_release_rx.await.expect("legacy append should be released"); + let mut file = legacy_disk + .append_file("bucket", "object/part.1") + .await + .expect("legacy append should open"); + file.write_all(b"legacy").await.expect("legacy append should write"); + file.shutdown().await.expect("legacy append should finish"); + }); + legacy_started_rx.await.expect("legacy append should hold the target lock"); + + let authenticated_disk = disk.clone(); + let authenticated_query = query.clone(); + let authenticated = tokio::spawn(async move { + write_authenticated_put_file( + &authenticated_disk, + iter(vec![Ok::(Bytes::from(payload))]), + &authenticated_query, + nonce, + url, + ) + .await + }); + tokio::task::yield_now().await; + assert!(!authenticated.is_finished(), "authenticated append must wait for legacy append"); + + legacy_release_tx.send(()).expect("legacy append should still be waiting"); + legacy.await.expect("legacy append task should join"); + authenticated + .await + .expect("authenticated append task should join") + .expect("authenticated append should succeed"); + assert_eq!( + disk.read_all("bucket", "object/part.1") + .await + .expect("final data should be readable"), + Bytes::from_static(b"old-datalegacyauthenticated") + ); + } + #[tokio::test] async fn walk_dir_body_surfaces_background_failure_after_data() { let body = walk_dir_response_body(true, |mut writer| async move { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index d7fd898ea..ef7dac0ad 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -208,6 +208,12 @@ pub(crate) mod rpc_consumer { pub(crate) use super::super::rpc::InternodeRpcService; pub(crate) mod http_service { + #[cfg(test)] + pub(crate) use super::super::Endpoint; + #[cfg(test)] + pub(crate) use super::super::ecstore_disk::DiskAPI; + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::disk::{DiskOption, new_disk}; pub(crate) const DEFAULT_READ_BUFFER_SIZE: usize = super::super::DEFAULT_READ_BUFFER_SIZE; #[cfg(test)] pub(crate) use super::super::storage_contracts::{ @@ -220,8 +226,8 @@ pub(crate) mod rpc_consumer { WALK_DIR_STREAM_COMPLETION_V1, }; pub(crate) use super::super::{ - StorageDiskRpcExt, WalkDirOptions, check_and_record_signed_rpc_nonce, find_local_disk_by_ref, - sign_ns_scanner_capability, verify_put_file_auth_trailer, verify_rpc_signature, + DeleteOptions, DiskStore, StorageDiskRpcExt, WalkDirOptions, check_and_record_signed_rpc_nonce, + find_local_disk_by_ref, sign_ns_scanner_capability, verify_put_file_auth_trailer, verify_rpc_signature, }; } @@ -1118,6 +1124,7 @@ pub(crate) trait StorageDiskRpcExt { dst_path: &str, ) -> DiskResult; async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> DiskResult>; + async fn read_file(&self, volume: &str, path: &str) -> DiskResult; async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> DiskResult; async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> DiskResult<()>; async fn rename_part( @@ -1266,6 +1273,10 @@ where ecstore_disk::DiskAPI::list_dir(self, origvolume, volume, dir_path, count).await } + async fn read_file(&self, volume: &str, path: &str) -> DiskResult { + ecstore_disk::DiskAPI::read_file(self, volume, path).await + } + async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> DiskResult { ecstore_disk::DiskAPI::read_file_stream(self, volume, path, offset, length).await }