From 7bd09a00b0de6c6741bb4ddf69f58733edec9cb2 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 12 Sep 2026 09:21:14 +0800 Subject: [PATCH] fix(get): reject pre-header read quorum failures (#7670) * fix(get): reject pre-header read quorum failures * fix(ci): unblock get pre-header quorum checks --- crates/ecstore/src/disk/error.rs | 2 +- rustfs/src/app/object/get.rs | 216 +++++++++++++++++++++++++++---- rustfs/src/app/object/shared.rs | 2 +- rustfs/src/error.rs | 31 ++++- 4 files changed, 216 insertions(+), 35 deletions(-) diff --git a/crates/ecstore/src/disk/error.rs b/crates/ecstore/src/disk/error.rs index fad6e561d..8c00ab93e 100644 --- a/crates/ecstore/src/disk/error.rs +++ b/crates/ecstore/src/disk/error.rs @@ -866,7 +866,7 @@ mod tests { "staging rejected", ))); assert!(marked.is_conditional_file_not_committed()); - assert!(marked.clone().is_conditional_file_not_committed()); + assert!(marked.is_conditional_file_not_committed()); assert!(!DiskError::Timeout.is_conditional_file_not_committed()); assert!( !DiskError::Io(io::Error::new(io::ErrorKind::PermissionDenied, "rename rejected")) diff --git a/rustfs/src/app/object/get.rs b/rustfs/src/app/object/get.rs index 34814c33c..f9a0ca36e 100644 --- a/rustfs/src/app/object/get.rs +++ b/rustfs/src/app/object/get.rs @@ -2133,7 +2133,7 @@ impl DefaultObjectUsecase { } #[allow(clippy::too_many_arguments)] - fn build_reader_blob( + async fn build_reader_blob( reader: R, response_content_length: i64, request_id: &str, @@ -2144,10 +2144,12 @@ impl DefaultObjectUsecase { key: &str, lifecycle: GetObjectBodyLifecycle, resume: Option>, - ) -> StreamingBlob + ) -> S3Result where R: AsyncRead + Send + Sync + Unpin + 'static, { + use tokio::io::AsyncReadExt as _; + let streaming_blob_start = rustfs_io_metrics::get_stage_metrics_enabled().then(std::time::Instant::now); let expected = usize::try_from(response_content_length.max(0)).unwrap_or(usize::MAX); let tuned_stream_buffer_size = @@ -2163,7 +2165,7 @@ impl DefaultObjectUsecase { ); } let handoff_start = get_stage_metrics_enabled.then(std::time::Instant::now); - let reader = GetObjectStreamingReader::new( + let mut reader = GetObjectStreamingReader::new( reader, bucket, key, @@ -2174,6 +2176,17 @@ impl DefaultObjectUsecase { lifecycle, resume, ); + let mut prefix = [0_u8; 1]; + let prefix_len = if expected == 0 { + 0 + } else { + reader + .read_exact(&mut prefix) + .await + .map_err(|error| map_get_object_reader_error(StorageError::from(error)))?; + 1 + }; + let reader = std::io::Cursor::new(prefix).take(prefix_len).chain(reader); let stream = GetObjectReaderStream::new(reader, stream_buffer_size, expected, stream_strategy.as_str(), buffer_source) .with_diagnostics(bucket, key, request_id); let blob = StreamingBlob::new(stream); @@ -2187,7 +2200,7 @@ impl DefaultObjectUsecase { ); } record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_BODY_STREAMING_BLOB, streaming_blob_start); - blob + Ok(blob) } fn init_get_object_bootstrap(&self, bucket: &str, key: &str, request_id: &str) -> S3Result { @@ -3161,7 +3174,7 @@ impl DefaultObjectUsecase { let (stream_buffer_size, stream_strategy) = Self::select_stream_buffer_strategy(response_content_length, optimal_buffer_size, enable_readahead, has_range); record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_BODY_STREAM_STRATEGY, stream_strategy_start); - return Ok(Self::build_reader_blob( + return Self::build_reader_blob( final_stream, response_content_length, request_id, @@ -3172,7 +3185,8 @@ impl DefaultObjectUsecase { key, lifecycle, resume(info), - )); + ) + .await; } if let Some(buffered_body) = buffered_body { @@ -3239,7 +3253,7 @@ impl DefaultObjectUsecase { let (stream_buffer_size, stream_strategy) = Self::select_stream_buffer_strategy(response_content_length, optimal_buffer_size, enable_readahead, has_range); record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_BODY_STREAM_STRATEGY, stream_strategy_start); - Ok(Self::build_reader_blob( + Self::build_reader_blob( final_stream, response_content_length, request_id, @@ -3250,7 +3264,8 @@ impl DefaultObjectUsecase { key, lifecycle, resume(info), - )) + ) + .await } #[allow(clippy::too_many_arguments)] @@ -5847,8 +5862,11 @@ mod tests { } impl AsyncRead for ReadProbeReader { - fn poll_read(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &mut ReadBuf<'_>) -> Poll> { + fn poll_read(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { self.reads.fetch_add(1, AtomicOrdering::Relaxed); + if buf.remaining() > 0 { + buf.put_slice(b"x"); + } Poll::Ready(Ok(())) } } @@ -6448,12 +6466,11 @@ mod tests { ) .await .expect("reservation bypass must construct the normal streaming fallback"); - let chunk = fallback_body - .next() - .await - .expect("fallback stream must yield a body chunk") - .expect("fallback stream must not fail"); - assert_eq!(chunk, Bytes::from_static(b"body")); + let mut received = Vec::new(); + while let Some(chunk) = fallback_body.next().await { + received.extend_from_slice(&chunk.expect("fallback stream must not fail")); + } + assert_eq!(received, b"body"); assert!(fallback_reads.load(AtomicOrdering::Relaxed) > 0); assert_eq!(readers.load(AtomicOrdering::Relaxed), 0, "cold-fill materialization must remain unopened"); assert_eq!(coordinator.active_session_count_for_test(), 0); @@ -7773,6 +7790,151 @@ mod tests { ) } + #[tokio::test] + async fn build_reader_blob_rejects_quorum_failure_before_handoff() { + let result = DefaultObjectUsecase::build_reader_blob( + FailAtEndReader::new( + b"", + Some(std::io::Error::other(StorageError::InsufficientReadQuorum( + "test-bucket".to_string(), + "unavailable-object".to_string(), + ))), + ), + 5, + "req-preheader-quorum", + None, + 64, + GetObjectStreamStrategy::Standard, + "test-bucket", + "unavailable-object", + GetObjectBodyLifecycle::disabled(), + None, + ) + .await; + + let error = result.expect_err("a read quorum failure before the first byte must reject response construction"); + assert_eq!(error.code(), &S3ErrorCode::Custom("SlowDownRead".into())); + assert_eq!(error.status_code(), Some(StatusCode::SERVICE_UNAVAILABLE)); + assert_eq!(error.message(), Some("Resource requested is unreadable, please reduce your request rate")); + } + + #[tokio::test] + async fn build_reader_blob_preserves_primed_byte_for_full_and_range() { + for (request_id, content_range) in [("req-preheader-full", None), ("req-preheader-range", Some("bytes 10-14/100"))] { + let reads = Arc::new(AtomicUsize::new(0)); + let mut body = DefaultObjectUsecase::build_reader_blob( + DataProbeReader { + reads: Arc::clone(&reads), + data: std::io::Cursor::new(b"hello".to_vec()), + }, + 5, + request_id, + content_range, + 64, + GetObjectStreamStrategy::Standard, + "test-bucket", + "test-object", + GetObjectBodyLifecycle::disabled(), + None, + ) + .await + .expect("the first byte should be available before response handoff"); + + assert_eq!(reads.load(AtomicOrdering::Relaxed), 1, "response construction must prime one byte"); + let mut received = Vec::new(); + while let Some(chunk) = body.next().await { + received.extend_from_slice(&chunk.expect("the primed body should remain readable")); + } + assert_eq!(received, b"hello", "the primed byte must be delivered exactly once"); + } + } + + #[tokio::test] + async fn build_reader_blob_does_not_poll_empty_object() { + let reads = Arc::new(AtomicUsize::new(0)); + let mut body = DefaultObjectUsecase::build_reader_blob( + ReadProbeReader { + reads: Arc::clone(&reads), + }, + 0, + "req-empty-object", + None, + 64, + GetObjectStreamStrategy::Standard, + "test-bucket", + "empty-object", + GetObjectBodyLifecycle::disabled(), + None, + ) + .await + .expect("an empty response should not require a storage read"); + + assert_eq!(reads.load(AtomicOrdering::Relaxed), 0); + assert!(body.next().await.is_none()); + assert_eq!(reads.load(AtomicOrdering::Relaxed), 0); + } + + #[tokio::test] + async fn build_reader_blob_leaves_later_failure_in_body_stream() { + let mut body = DefaultObjectUsecase::build_reader_blob( + FailAtEndReader::new(b"h", Some(std::io::Error::other("failure after handoff"))), + 5, + "req-postheader-failure", + None, + 64, + GetObjectStreamStrategy::Standard, + "test-bucket", + "later-failure-object", + GetObjectBodyLifecycle::disabled(), + None, + ) + .await + .expect("the available first byte should allow response handoff"); + + let first = body + .next() + .await + .expect("the primed byte must be present") + .expect("the primed byte must be successful"); + assert_eq!(first, Bytes::from_static(b"h")); + let error = body + .next() + .await + .expect("the later read must produce a body result") + .expect_err("a failure after the first byte must stay in the body stream"); + assert!(error.to_string().contains("failure after handoff")); + } + + #[tokio::test] + async fn build_reader_blob_resume_offset_includes_primed_byte() { + let reopen_count = Arc::new(AtomicUsize::new(0)); + let control = counting_resume_control(Arc::clone(&reopen_count), |emitted| { + assert_eq!(emitted, 1, "resume must start after the byte consumed before response handoff"); + Ok(FailAtEndReader::new(b"ello", None)) + }); + let mut body = DefaultObjectUsecase::build_reader_blob( + FailAtEndReader::new(b"h", Some(relocation_read_error())), + 5, + "req-preheader-resume-offset", + None, + 64, + GetObjectStreamStrategy::Standard, + "test-bucket", + "relocated-object", + GetObjectBodyLifecycle::disabled(), + Some(control), + ) + .await + .expect("the first byte should permit response handoff before relocation"); + + let mut received = Vec::new(); + while let Some(chunk) = body.next().await { + received.extend_from_slice(&chunk.expect("resume should complete the body")); + } + assert_eq!(received, b"hello"); + assert_eq!(reopen_count.load(Ordering::Relaxed), 1); + } + #[tokio::test] async fn get_object_streaming_reader_resumes_after_relocation_error() { use tokio::io::AsyncReadExt; @@ -9041,7 +9203,7 @@ mod tests { } #[tokio::test] - async fn build_get_object_body_keeps_large_objects_on_streaming_path_without_preread() { + async fn build_get_object_body_primes_large_stream_before_handoff() { let reads = Arc::new(AtomicUsize::new(0)); let reader = ReadProbeReader { reads: Arc::clone(&reads), @@ -9074,13 +9236,13 @@ mod tests { assert_eq!( reads.load(AtomicOrdering::Relaxed), - 0, - "large-object response construction should not pre-read object data" + 1, + "large-object response construction should prime exactly one byte" ); } #[tokio::test] - async fn build_get_object_body_keeps_large_encrypted_objects_on_streaming_path_without_preread() { + async fn build_get_object_body_primes_large_encrypted_stream_before_handoff() { let reads = Arc::new(AtomicUsize::new(0)); let reader = ReadProbeReader { reads: Arc::clone(&reads), @@ -9113,8 +9275,8 @@ mod tests { assert_eq!( reads.load(AtomicOrdering::Relaxed), - 0, - "large encrypted object response construction should not pre-read object data" + 1, + "large encrypted object response construction should prime exactly one byte" ); } @@ -9284,8 +9446,8 @@ mod tests { assert_eq!(fill, rustfs_object_data_cache::ObjectDataCacheFillResult::SkippedSizeMismatch); assert_eq!( reads.load(AtomicOrdering::Relaxed), - 0, - "size-mismatched rejected fill should construct the fallback stream without pre-reading" + 1, + "size-mismatched rejected fill should prime the fallback stream before handoff" ); assert!( matches!(lookup_after_mismatch, rustfs_object_data_cache::ObjectDataCacheLookup::Miss), @@ -9993,8 +10155,8 @@ mod tests { assert_eq!( reads.load(AtomicOrdering::Relaxed), - 0, - "too-large materialize-fill candidate must not pre-read the fallback reader" + 1, + "too-large materialize-fill candidate must prime the streaming fallback" ); } @@ -10032,8 +10194,8 @@ mod tests { assert_eq!( reads.load(AtomicOrdering::Relaxed), - 0, - "default GetObject response construction should not pre-read small plain object data" + 1, + "default GetObject response construction should prime exactly one byte" ); } diff --git a/rustfs/src/app/object/shared.rs b/rustfs/src/app/object/shared.rs index 063203e07..79edcfa32 100644 --- a/rustfs/src/app/object/shared.rs +++ b/rustfs/src/app/object/shared.rs @@ -465,7 +465,7 @@ mod bucket_default_sse_lookup_tests { let err = classify_bucket_default_sse_lookup("bucket", Err(StorageError::ErasureReadQuorum)) .expect_err("an unreadable metadata subsystem must never degrade to plaintext"); - assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable); + assert_eq!(err.code(), &S3ErrorCode::Custom("SlowDownRead".into())); } #[test] diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index 91439e458..9531f795e 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -20,6 +20,8 @@ use s3s::{S3Error, S3ErrorCode}; const MAX_VERSIONS_EXCEEDED_CODE: &str = "MaxVersionsExceeded"; const MAX_VERSIONS_EXCEEDED_MESSAGE: &str = "You've exceeded the limit on the number of versions you can create on this object"; +const SLOW_DOWN_READ_CODE: &str = "SlowDownRead"; +const SLOW_DOWN_READ_MESSAGE: &str = "Resource requested is unreadable, please reduce your request rate"; /// S3 error code for a request that names a KMS key the KMS does not hold. pub const KMS_KEY_NOT_FOUND_ERROR_CODE: &str = "KMS.NotFoundException"; @@ -105,6 +107,7 @@ fn custom_error_status(code: &S3ErrorCode) -> Option { S3ErrorCode::Custom(custom) if &**custom == KMS_KEY_NOT_FOUND_ERROR_CODE || &**custom == MAX_VERSIONS_EXCEEDED_CODE => { Some(StatusCode::BAD_REQUEST) } + S3ErrorCode::Custom(custom) if &**custom == SLOW_DOWN_READ_CODE => Some(StatusCode::SERVICE_UNAVAILABLE), _ => None, } } @@ -403,6 +406,7 @@ impl ApiError { S3ErrorCode::Custom(code) if &**code == MAX_VERSIONS_EXCEEDED_CODE => { MAX_VERSIONS_EXCEEDED_MESSAGE.to_string() } + S3ErrorCode::Custom(code) if &**code == SLOW_DOWN_READ_CODE => SLOW_DOWN_READ_MESSAGE.to_string(), _ => code.as_str().to_string(), } } @@ -577,10 +581,10 @@ impl From for ApiError { | StorageError::FaultyRemoteDisk | StorageError::DiskNotFound | StorageError::TooManyOpenFiles => S3ErrorCode::ServiceUnavailable, - StorageError::ErasureReadQuorum - | StorageError::InsufficientReadQuorum(_, _) - | StorageError::ErasureWriteQuorum - | StorageError::InsufficientWriteQuorum(_, _) => S3ErrorCode::ServiceUnavailable, + StorageError::ErasureReadQuorum | StorageError::InsufficientReadQuorum(_, _) => { + S3ErrorCode::Custom(SLOW_DOWN_READ_CODE.into()) + } + StorageError::ErasureWriteQuorum | StorageError::InsufficientWriteQuorum(_, _) => S3ErrorCode::ServiceUnavailable, StorageError::NamespaceLockQuorumUnavailable { .. } => S3ErrorCode::ServiceUnavailable, StorageError::QuotaExceeded { .. } => S3ErrorCode::InvalidRequest, StorageError::MaxVersionsExceeded => S3ErrorCode::Custom(MAX_VERSIONS_EXCEEDED_CODE.into()), @@ -1288,10 +1292,10 @@ mod tests { (StorageError::FaultyRemoteDisk, S3ErrorCode::ServiceUnavailable), (StorageError::DiskNotFound, S3ErrorCode::ServiceUnavailable), (StorageError::TooManyOpenFiles, S3ErrorCode::ServiceUnavailable), - (StorageError::ErasureReadQuorum, S3ErrorCode::ServiceUnavailable), + (StorageError::ErasureReadQuorum, S3ErrorCode::Custom(SLOW_DOWN_READ_CODE.into())), ( StorageError::InsufficientReadQuorum("test".into(), "test".into()), - S3ErrorCode::ServiceUnavailable, + S3ErrorCode::Custom(SLOW_DOWN_READ_CODE.into()), ), (StorageError::ErasureWriteQuorum, S3ErrorCode::ServiceUnavailable), ( @@ -1429,6 +1433,21 @@ mod tests { assert_eq!(s3_error.status_code(), Some(http::StatusCode::SERVICE_UNAVAILABLE)); } + #[test] + fn read_quorum_failure_matches_minio_slow_down_read_response() { + for error in [ + StorageError::ErasureReadQuorum, + StorageError::InsufficientReadQuorum("bucket".into(), "object".into()), + ] { + let api_error = ApiError::from(error); + assert_eq!(api_error.code, S3ErrorCode::Custom(SLOW_DOWN_READ_CODE.into())); + assert_eq!(api_error.message, SLOW_DOWN_READ_MESSAGE); + + let s3_error: S3Error = api_error.into(); + assert_eq!(s3_error.status_code(), Some(StatusCode::SERVICE_UNAVAILABLE)); + } + } + #[test] fn quota_exceeded_preserves_existing_s3_error_contract() { let api_error: ApiError = StorageError::QuotaExceeded { current: 5, limit: 10 }.into();