diff --git a/crates/config/src/constants/body_limits.rs b/crates/config/src/constants/body_limits.rs index 7e6adba74..a1a9236ac 100644 --- a/crates/config/src/constants/body_limits.rs +++ b/crates/config/src/constants/body_limits.rs @@ -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 diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index 0d442a48a..655eecd8f 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -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] { diff --git a/rustfs/src/app/object/mod.rs b/rustfs/src/app/object/mod.rs index 454fd5ff8..b9d584321 100644 --- a/rustfs/src/app/object/mod.rs +++ b/rustfs/src/app/object/mod.rs @@ -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::*; diff --git a/rustfs/src/app/object/put.rs b/rustfs/src/app/object/put.rs index 322a4309a..296814aa0 100644 --- a/rustfs/src/app/object/put.rs +++ b/rustfs/src/app/object/put.rs @@ -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, +} + +#[cfg(test)] +impl PollCountingBody { + pub(crate) fn streaming_blob() -> (StreamingBlob, std::sync::Arc) { + 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; + + fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + 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()) + ); + } +} diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index 91439e458..f8c69e30f 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -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::() + && 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::() + .and_then(|hyper_err| std::error::Error::source(hyper_err)) + .and_then(|cause| cause.downcast_ref::()) + .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 for S3Error { fn from(err: ApiError) -> Self { let status = custom_error_status(&err.code); @@ -535,6 +586,22 @@ impl From 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 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>> = 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| { + 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 }); diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 36583050a..8b0ef6988 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -158,13 +158,11 @@ static HTTP_STATUS_CLASS_METRICS: std::sync::LazyLock<[HttpStatusClassMetrics; 6 static HTTP_TRANSPORT_FAILURES_COUNTER: std::sync::LazyLock = 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"));