mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-10 14:16:01 +00:00
Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 56d14d9453 | |||
| 26affde51a | |||
| 51928eedcf | |||
| 84464f567c | |||
| 22fb79b6f2 |
@@ -64,6 +64,14 @@ pub const MAX_HEAL_REQUEST_SIZE: usize = 1024 * 1024; // 1 MB
|
||||
/// memory exhaustion from malicious or misconfigured remote services.
|
||||
pub const MAX_S3_CLIENT_RESPONSE_SIZE: usize = 10 * 1024 * 1024; // 10 MB
|
||||
|
||||
/// Maximum body size accepted by a single `PutObject` or `UploadPart` request (5 GiB).
|
||||
/// Used for: the s3s streaming-body limit and the request-header admission check.
|
||||
/// Rationale: matches the AWS S3 single-PUT / single-part ceiling. Larger objects
|
||||
/// must use multipart upload. The header check rejects an oversize
|
||||
/// `Content-Length` before any body byte is read so the client gets
|
||||
/// `EntityTooLarge` immediately instead of streaming 5 GiB into a mid-stream failure.
|
||||
pub const MAX_SINGLE_PUT_OBJECT_SIZE: u64 = 5 * 1024 * 1024 * 1024; // 5 GiB
|
||||
|
||||
/// Maximum size for OIDC provider response bodies (1 MB)
|
||||
/// Used for: discovery documents, JWKS documents and token endpoint responses
|
||||
/// Rationale: a hostile or compromised identity provider must not be able to exhaust
|
||||
|
||||
@@ -69,7 +69,7 @@ use super::storage_api::multipart_usecase::{
|
||||
};
|
||||
use crate::app::object::{
|
||||
ConcurrencyManager, ForegroundWriteAdmission, get_concurrency_manager, guard_put_object_body_read_timeout,
|
||||
put_object_body_read_timeout,
|
||||
put_object_body_read_timeout, reject_oversize_single_upload,
|
||||
};
|
||||
use crate::app::object_data_cache::{
|
||||
ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success,
|
||||
@@ -1169,6 +1169,9 @@ impl DefaultMultipartUsecase {
|
||||
validate_table_catalog_object_mutation(&bucket, &key).await?;
|
||||
|
||||
let mut size = resolve_upload_part_size(&req.headers, content_length)?;
|
||||
if let Some(size) = size {
|
||||
reject_oversize_single_upload(size)?;
|
||||
}
|
||||
let mut body_stream = body.ok_or_else(|| s3_error!(IncompleteBody))?;
|
||||
let Some(store) = self.object_store() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
@@ -3209,6 +3212,36 @@ mod tests {
|
||||
assert_eq!(err.code(), &S3ErrorCode::IncompleteBody);
|
||||
}
|
||||
|
||||
/// issue #7596: a part whose declared length exceeds the 5 GiB
|
||||
/// single-request ceiling is rejected before the body is polled or the
|
||||
/// store is consulted. Exact-cap and zero-length parts pass admission.
|
||||
#[tokio::test]
|
||||
async fn execute_upload_part_rejects_oversize_declared_part_before_reading_the_body() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
|
||||
for (declared, expect_too_large) in [(ceiling + 1, true), (ceiling, false), (0, false)] {
|
||||
let (body, polls) = crate::app::object::PollCountingBody::streaming_blob();
|
||||
let input = UploadPartInput::builder()
|
||||
.bucket("bucket".to_string())
|
||||
.key("object".to_string())
|
||||
.upload_id("upload-id".to_string())
|
||||
.part_number(1)
|
||||
.body(Some(body))
|
||||
.content_length(Some(declared))
|
||||
.build()
|
||||
.unwrap();
|
||||
let req = build_request(input, Method::PUT);
|
||||
|
||||
let err = make_usecase().execute_upload_part(req).await.unwrap_err();
|
||||
if expect_too_large {
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge, "declared {declared}");
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
} else {
|
||||
assert_ne!(err.code(), &S3ErrorCode::EntityTooLarge, "declared {declared} must pass admission");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_upload_part_rejects_invalid_part_number_before_body_lookup() {
|
||||
for part_number in [-1, 0, 10001] {
|
||||
|
||||
@@ -223,8 +223,10 @@ pub(crate) use self::extract::*;
|
||||
pub(crate) use self::get::*;
|
||||
pub(crate) use self::internal_put::*;
|
||||
pub(crate) use self::on_demand_migration_put::*;
|
||||
#[cfg(test)]
|
||||
pub(crate) use self::put::PollCountingBody;
|
||||
use self::put::*;
|
||||
pub(crate) use self::put::{guard_put_object_body_read_timeout, put_object_body_read_timeout};
|
||||
pub(crate) use self::put::{guard_put_object_body_read_timeout, put_object_body_read_timeout, reject_oversize_single_upload};
|
||||
#[cfg(test)]
|
||||
pub(crate) use self::restore::RestoreStatusCommitBarrier;
|
||||
pub(crate) use self::shared::*;
|
||||
|
||||
@@ -109,6 +109,21 @@ fn resolve_put_object_authoritative_size(headers: &HeaderMap, content_length: Op
|
||||
Ok(size)
|
||||
}
|
||||
|
||||
/// Reject a declared upload length above the single-request ceiling
|
||||
/// ([`rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE`]) with `EntityTooLarge`.
|
||||
///
|
||||
/// Applies to `PutObject` and `UploadPart`. A negative or unknown length is
|
||||
/// left to the caller's existing validation.
|
||||
pub(crate) fn reject_oversize_single_upload(size: i64) -> S3Result<()> {
|
||||
if u64::try_from(size).is_ok_and(|size| size > rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE) {
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::EntityTooLarge,
|
||||
ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Resolve the S3 request-body inter-chunk read timeout from the environment.
|
||||
///
|
||||
/// Returns `Duration::ZERO` when disabled (`RUSTFS_HTTP_REQUEST_BODY_READ_TIMEOUT=0`),
|
||||
@@ -1287,6 +1302,12 @@ impl DefaultObjectUsecase {
|
||||
// Resolve the authoritative decoded/plain object length (rejecting negative/unknown) before anything else consumes it.
|
||||
let size = resolve_put_object_authoritative_size(&req.headers, content_length)?;
|
||||
|
||||
// The streaming-body limit (s3s `put_object_max_size`) only fires once the
|
||||
// client has already streamed 5 GiB. The declared length is authoritative,
|
||||
// so reject an oversize single PUT here, before any body byte is read
|
||||
// (issue #7596).
|
||||
reject_oversize_single_upload(size)?;
|
||||
|
||||
if let Some(limit) = max_content_length
|
||||
&& u64::try_from(size).is_ok_and(|size| size > limit)
|
||||
{
|
||||
@@ -3318,6 +3339,77 @@ mod tests {
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass);
|
||||
}
|
||||
|
||||
/// issue #7596: a single PUT whose declared length exceeds the 5 GiB
|
||||
/// ceiling must be rejected from the headers, before any body byte is
|
||||
/// requested.
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_rejects_oversize_content_length_before_reading_the_body() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
let (body, polls) = PollCountingBody::streaming_blob();
|
||||
let input = PutObjectInput::builder()
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("huge.bin".to_string())
|
||||
.body(Some(body))
|
||||
.content_length(Some(ceiling + 1))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let req = build_request(input, Method::PUT);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
let fs = FS::new();
|
||||
|
||||
let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
}
|
||||
|
||||
/// Admission uses the logical object size, not the wire length: a signed
|
||||
/// aws-chunked request whose framed `Content-Length` exceeds the cap but
|
||||
/// whose decoded length is within it must not be rejected as oversize,
|
||||
/// while a decoded length above the cap must be.
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_oversize_admission_uses_decoded_length_for_aws_chunked() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
let framing_overhead = 1_000_000;
|
||||
|
||||
for (decoded, expect_too_large) in [(ceiling, false), (ceiling + 1, true)] {
|
||||
let (body, polls) = PollCountingBody::streaming_blob();
|
||||
let input = PutObjectInput::builder()
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("huge.bin".to_string())
|
||||
.body(Some(body))
|
||||
.content_length(Some(decoded + framing_overhead))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let mut req = build_request(input, Method::PUT);
|
||||
req.headers
|
||||
.insert(http::header::CONTENT_ENCODING, HeaderValue::from_static("aws-chunked"));
|
||||
req.headers.insert(
|
||||
HeaderName::from_static("x-amz-content-sha256"),
|
||||
HeaderValue::from_static("STREAMING-AWS4-HMAC-SHA256-PAYLOAD"),
|
||||
);
|
||||
req.headers.insert(
|
||||
HeaderName::from_static("x-amz-decoded-content-length"),
|
||||
HeaderValue::from_str(&decoded.to_string()).unwrap(),
|
||||
);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
let fs = FS::new();
|
||||
|
||||
let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err();
|
||||
if expect_too_large {
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge, "decoded {decoded}");
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
} else {
|
||||
assert_ne!(
|
||||
err.code(),
|
||||
&S3ErrorCode::EntityTooLarge,
|
||||
"framed wire length above the cap must not reject a decoded length at the cap"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_rejects_post_object_sse_kms_from_headers() {
|
||||
let input = PutObjectInput::builder()
|
||||
@@ -4184,3 +4276,55 @@ mod tests {
|
||||
assert!(is_err_object_not_found(&lookup_err), "{lookup_err}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Test-only request body that records how often it is polled, so admission
|
||||
/// tests can prove a rejection happened before any body byte was requested.
|
||||
#[cfg(test)]
|
||||
pub(crate) struct PollCountingBody {
|
||||
pub(crate) polls: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl PollCountingBody {
|
||||
pub(crate) fn streaming_blob() -> (StreamingBlob, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
|
||||
let polls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let body = StreamingBlob::new(Self {
|
||||
polls: std::sync::Arc::clone(&polls),
|
||||
});
|
||||
(body, polls)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Stream for PollCountingBody {
|
||||
type Item = Result<Bytes, StdError>;
|
||||
|
||||
fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
self.polls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
Poll::Ready(Some(Ok(Bytes::from_static(b"x"))))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl ByteStream for PollCountingBody {}
|
||||
|
||||
#[cfg(test)]
|
||||
mod oversize_single_upload_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn reject_oversize_single_upload_enforces_the_single_request_ceiling() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
|
||||
assert!(reject_oversize_single_upload(0).is_ok());
|
||||
assert!(reject_oversize_single_upload(ceiling).is_ok(), "exact ceiling is allowed");
|
||||
assert!(reject_oversize_single_upload(-1).is_ok(), "unknown length is left to later validation");
|
||||
|
||||
let err = reject_oversize_single_upload(ceiling + 1).expect_err("one byte over must be rejected");
|
||||
assert_eq!(*err.code(), S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(
|
||||
err.message(),
|
||||
Some(ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge).as_str())
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -478,6 +478,57 @@ fn error_chain_s3s_body_stream_error(err: &(dyn std::error::Error + 'static)) ->
|
||||
None
|
||||
}
|
||||
|
||||
/// Walk an error chain (including `io::Error` custom payloads) and return
|
||||
/// whether any link satisfies `pred`.
|
||||
fn error_chain_any(err: &(dyn std::error::Error + 'static), pred: &dyn Fn(&(dyn std::error::Error + 'static)) -> bool) -> bool {
|
||||
if pred(err) {
|
||||
return true;
|
||||
}
|
||||
if let Some(io_err) = err.downcast_ref::<std::io::Error>()
|
||||
&& let Some(inner) = io_err.get_ref()
|
||||
&& error_chain_any(inner, pred)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
let mut current = err.source();
|
||||
while let Some(err) = current {
|
||||
if error_chain_any(err, pred) {
|
||||
return true;
|
||||
}
|
||||
current = err.source();
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
/// s3s raises `BodySizeLimitExceeded` when the streaming-body budget
|
||||
/// (`put_object_max_size`) runs out mid-stream. The type lives in s3s's
|
||||
/// private `http` module, so it is recognised by its `Display` form
|
||||
/// (`body size {size} exceeds limit {limit}`), like the other s3s body-stream
|
||||
/// errors above. Switch to a typed downcast once s3s re-exports the type.
|
||||
fn is_body_size_limit_exceeded_display(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
let text = err.to_string();
|
||||
text.starts_with("body size ") && text.contains(" exceeds limit ")
|
||||
}
|
||||
|
||||
fn error_chain_has_body_size_limit_exceeded(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
error_chain_any(err, &is_body_size_limit_exceeded_display)
|
||||
}
|
||||
|
||||
/// hyper reports a request body whose connection hit EOF before
|
||||
/// `Content-Length` bytes arrived as a `Kind::Body` error carrying an
|
||||
/// `UnexpectedEof` `io::Error` (its `IncompleteBody` marker is private).
|
||||
/// That is a client-side short body, not a server fault.
|
||||
fn is_hyper_body_eof(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
err.downcast_ref::<hyper::Error>()
|
||||
.and_then(|hyper_err| std::error::Error::source(hyper_err))
|
||||
.and_then(|cause| cause.downcast_ref::<std::io::Error>())
|
||||
.is_some_and(|io_err| io_err.kind() == std::io::ErrorKind::UnexpectedEof)
|
||||
}
|
||||
|
||||
fn error_chain_has_hyper_body_eof(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
error_chain_any(err, &is_hyper_body_eof)
|
||||
}
|
||||
|
||||
impl From<ApiError> for S3Error {
|
||||
fn from(err: ApiError) -> Self {
|
||||
let status = custom_error_status(&err.code);
|
||||
@@ -535,6 +586,22 @@ impl From<StorageError> for ApiError {
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_body_size_limit_exceeded(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::EntityTooLarge,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_hyper_body_eof(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if matches!(s3s_body_stream_error, Some(S3sBodyStreamError::IncompleteBody)) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
@@ -678,6 +745,22 @@ impl From<std::io::Error> for ApiError {
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
if error_chain_has_body_size_limit_exceeded(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::EntityTooLarge,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_hyper_body_eof(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if matches!(s3s_body_stream_error, Some(S3sBodyStreamError::IncompleteBody)) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
@@ -950,6 +1033,117 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_size_limit_exceeded_maps_to_entity_too_large_across_io_boundaries() {
|
||||
// Shape observed in production (issue #7596):
|
||||
// Custom { UnexpectedEof, Custom { Other, BodySizeLimitExceeded { size, limit } } }
|
||||
let nested = || {
|
||||
IoError::new(
|
||||
ErrorKind::UnexpectedEof,
|
||||
IoError::other(MockS3sBodyStreamError("body size 16384 exceeds limit 6389")),
|
||||
)
|
||||
};
|
||||
|
||||
let direct: ApiError = nested().into();
|
||||
assert_eq!(direct.code, S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(direct.message, ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge));
|
||||
|
||||
let storage: ApiError = StorageError::Io(nested()).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::EntityTooLarge);
|
||||
assert!(storage.source.is_some());
|
||||
|
||||
// An unrelated message that merely mentions a limit stays internal.
|
||||
let other: ApiError = IoError::other(MockS3sBodyStreamError("limit exceeded for something else")).into();
|
||||
assert_eq!(other.code, S3ErrorCode::InternalError);
|
||||
}
|
||||
|
||||
/// Trip s3s's real streaming-body budget with a tiny limit so the
|
||||
/// display-based matcher is checked against the pinned dependency's
|
||||
/// actual error, not only the mocked string.
|
||||
#[tokio::test]
|
||||
async fn real_s3s_body_size_limit_error_maps_to_entity_too_large() {
|
||||
use futures::StreamExt;
|
||||
|
||||
let real_error = || async {
|
||||
let mut body = s3s::Body::from(bytes::Bytes::from_static(b"hello"));
|
||||
body.set_limit(Some(4));
|
||||
body.next()
|
||||
.await
|
||||
.expect("one frame")
|
||||
.expect_err("five bytes must exceed a four-byte budget")
|
||||
};
|
||||
|
||||
let err = real_error().await;
|
||||
assert!(is_body_size_limit_exceeded_display(err.as_ref()), "unexpected display: {err}");
|
||||
|
||||
let err = real_error().await;
|
||||
let storage: ApiError = StorageError::Io(IoError::new(ErrorKind::UnexpectedEof, IoError::other(err))).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(storage.message, ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge));
|
||||
|
||||
let err = real_error().await;
|
||||
let direct: ApiError = IoError::other(err).into();
|
||||
assert_eq!(direct.code, S3ErrorCode::EntityTooLarge);
|
||||
}
|
||||
|
||||
/// Drive a real hyper HTTP/1 server so the test sees hyper's own body EOF
|
||||
/// error (`hyper::Error(Body, UnexpectedEof, IncompleteBody)`), which has no
|
||||
/// public constructor.
|
||||
async fn capture_hyper_body_eof_error() -> hyper::Error {
|
||||
use http_body_util::BodyExt;
|
||||
use hyper::service::service_fn;
|
||||
use hyper_util::rt::TokioIo;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("bind");
|
||||
let addr = listener.local_addr().expect("local addr");
|
||||
let captured: Arc<Mutex<Option<hyper::Error>>> = Arc::new(Mutex::new(None));
|
||||
let server_slot = Arc::clone(&captured);
|
||||
let server = tokio::spawn(async move {
|
||||
let (stream, _) = listener.accept().await.expect("accept");
|
||||
let slot = server_slot;
|
||||
let service = service_fn(move |req: hyper::Request<hyper::body::Incoming>| {
|
||||
let slot = Arc::clone(&slot);
|
||||
async move {
|
||||
let err = req.into_body().collect().await.expect_err("short body must fail");
|
||||
*slot.lock().expect("slot") = Some(err);
|
||||
Ok::<_, std::convert::Infallible>(hyper::Response::new(String::new()))
|
||||
}
|
||||
});
|
||||
let _ = hyper::server::conn::http1::Builder::new()
|
||||
.serve_connection(TokioIo::new(stream), service)
|
||||
.await;
|
||||
});
|
||||
|
||||
let mut client = tokio::net::TcpStream::connect(addr).await.expect("connect");
|
||||
client
|
||||
.write_all(b"PUT /bucket/key HTTP/1.1\r\nHost: localhost\r\nContent-Length: 100\r\n\r\nabc")
|
||||
.await
|
||||
.expect("write partial body");
|
||||
client.shutdown().await.expect("shutdown write side");
|
||||
let _ = tokio::time::timeout(std::time::Duration::from_secs(10), server).await;
|
||||
let captured = captured.lock().expect("slot").take();
|
||||
captured.expect("hyper body error captured")
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn hyper_body_eof_maps_to_incomplete_body_across_io_boundaries() {
|
||||
let hyper_err = capture_hyper_body_eof_error().await;
|
||||
assert!(is_hyper_body_eof(&hyper_err), "unexpected hyper error shape: {hyper_err:?}");
|
||||
|
||||
// Shape observed in production (issue #7596):
|
||||
// Custom { UnexpectedEof, Custom { Other, hyper::Error(Body, UnexpectedEof, IncompleteBody) } }
|
||||
let nested = IoError::new(ErrorKind::UnexpectedEof, IoError::other(hyper_err));
|
||||
let storage: ApiError = StorageError::Io(nested).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::IncompleteBody);
|
||||
assert_eq!(storage.message, ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody));
|
||||
|
||||
let hyper_err = capture_hyper_body_eof_error().await;
|
||||
let direct: ApiError = IoError::other(hyper_err).into();
|
||||
assert_eq!(direct.code, S3ErrorCode::IncompleteBody);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn server_side_source_read_error_maps_to_service_unavailable_before_incomplete_body() {
|
||||
let short_source = IoError::new(ErrorKind::UnexpectedEof, rustfs_rio::IncompleteBody { remaining: 17 });
|
||||
|
||||
@@ -158,13 +158,11 @@ static HTTP_STATUS_CLASS_METRICS: std::sync::LazyLock<[HttpStatusClassMetrics; 6
|
||||
static HTTP_TRANSPORT_FAILURES_COUNTER: std::sync::LazyLock<metrics::Counter> =
|
||||
std::sync::LazyLock::new(|| counter!(METRIC_HTTP_SERVER_FAILURES_TOTAL, LABEL_HTTP_STATUS_CLASS => "transport"));
|
||||
|
||||
const RUSTFS_S3_PUT_OBJECT_MAX_SIZE: u64 = 5 * 1024 * 1024 * 1024;
|
||||
|
||||
fn rustfs_s3_config() -> S3Config {
|
||||
let mut s3_config = S3Config::default();
|
||||
s3_config.normalize_forward_slash_path = true;
|
||||
s3_config.enable_sig_v2 = true;
|
||||
s3_config.put_object_max_size = Some(RUSTFS_S3_PUT_OBJECT_MAX_SIZE);
|
||||
s3_config.put_object_max_size = Some(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE);
|
||||
s3_config.sig_v4_allowed_services.push("s3tables".to_string());
|
||||
s3_config
|
||||
}
|
||||
@@ -3051,7 +3049,7 @@ mod tests {
|
||||
assert!(s3_config.normalize_forward_slash_path);
|
||||
assert!(s3_config.normalize_content_length);
|
||||
assert!(s3_config.enable_sig_v2);
|
||||
assert_eq!(s3_config.put_object_max_size, Some(RUSTFS_S3_PUT_OBJECT_MAX_SIZE));
|
||||
assert_eq!(s3_config.put_object_max_size, Some(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "s3"));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "sts"));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "s3tables"));
|
||||
|
||||
Reference in New Issue
Block a user