Compare commits

...

5 Commits

Author SHA1 Message Date
唐小鸭 56d14d9453 Merge branch 'main' into fix/put-object-size-admission 2026-09-10 21:50:08 +08:00
唐小鸭 26affde51a Merge remote-tracking branch 'origin/fix/put-object-size-admission' into fix/put-object-size-admission 2026-09-10 20:18:48 +08:00
唐小鸭 51928eedcf test(s3): cover UploadPart admission, aws-chunked length, real s3s limit
- Poll-counting test body proves PutObject and UploadPart reject a
  declared size above the ceiling with zero body polls; exact-cap and
  zero-length parts pass admission.
- A STREAMING-* aws-chunked PUT whose framed Content-Length exceeds the
  cap is admitted when the decoded length is within it and rejected when
  the decoded length is over it.
- The display-based BodySizeLimitExceeded matcher is checked against the
  real error produced by the pinned s3s Body budget.
2026-09-10 20:18:14 +08:00
唐小鸭 84464f567c Merge branch 'main' into fix/put-object-size-admission 2026-09-10 19:46:21 +08:00
唐小鸭 22fb79b6f2 fix(s3): reject oversize single PUT early and map body errors to 4xx
A single PutObject above the 5 GiB single-request ceiling was only
rejected after the client had streamed 5 GiB into s3s's read-time body
budget, and the resulting BodySizeLimitExceeded surfaced from the erasure
writer as 500 InternalError. A body whose connection hit EOF before
Content-Length bytes arrived (hyper's IncompleteBody) was also a 500.
SDKs retry 500s, so one oversize upload was resent from offset 0 five
times.

- PutObject and UploadPart reject a declared length above
  MAX_SINGLE_PUT_OBJECT_SIZE with 400 EntityTooLarge before reading the
  body; the constant moves to rustfs_config so the s3s limit and the
  admission check share one value.
- ApiError maps BodySizeLimitExceeded to EntityTooLarge and a hyper body
  EOF to IncompleteBody across both io::Error conversions.

Fixes #7596.
2026-09-10 16:48:02 +08:00
6 changed files with 385 additions and 6 deletions
@@ -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
+34 -1
View File
@@ -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] {
+3 -1
View File
@@ -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::*;
+144
View File
@@ -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())
);
}
}
+194
View File
@@ -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 });
+2 -4
View File
@@ -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"));