diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index d1112f9c5..f6fb5e709 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=294350518743cac8d7c41880a2835216e4b697908d7b0b1bc92b62816d94c59d +sha256=dbebfbab9b9efd4eff31211e69dd32235dc00e207f2ab0dd919a1b2ac9e724c2 diff --git a/.config/nextest.toml b/.config/nextest.toml index 90e09740a..4f92bf95e 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -481,23 +481,12 @@ path = "junit.xml" # parallel-safe — the same property e2e-smoke relies on. The exceptions are the # 4-disk reliability / degraded-read fault-injection tests and the fixed-port # Vault tests, both serialized below. -# KNOWN-FAILURE EXCLUSIONS (characterization run 29381309848, 2026-07-15: -# 341 ran / 32 failed on the suites' first automated run ever). Deterministic -# product failures cannot be quarantined away with retries, so each family is -# excluded here with its tracking issue, under the same discipline as the -# ci-profile quarantine (docs/testing/README.md): every entry MUST cite one -# OPEN issue, and the fixing PR MUST delete the exclusion. The passing -# negative-path siblings of each family stay in as regression guards. -# * rustfs#4843 — over-limit archive entry paths hard-reject the whole -# archive even under ignore-errors semantics. [profile.e2e-full] default-filter = """ package(e2e_test) & !test(/^protocols::/) & !test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|heal_erasure_disk_rebuild_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/) & !test(/^replication_extension_test::/) - & !test(/^multipart_auth_test::test_signed_put_object_extract_skips_invalid_entry_when_ignore_errors_enabled$/) - & !test(/^snowball_auto_extract_test::tests::snowball_auto_extract_(ignores_invalid_entries_when_requested|supports_standard_headers_with_combined_extract_options)$/) """ fail-fast = false diff --git a/Cargo.lock b/Cargo.lock index 1375bf3b3..ab5e7893c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -627,8 +627,7 @@ dependencies = [ [[package]] name = "astral-tokio-tar" version = "0.7.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6f2e989b33246fe9240d39accf4dd9a01e0b6c1f3ce9dd095e0a47fa02505523" +source = "git+https://github.com/cxymds/tokio-tar.git?rev=603756478b7668436e464519c77ccac22a99ba96#603756478b7668436e464519c77ccac22a99ba96" dependencies = [ "futures-core", "libc", diff --git a/Cargo.toml b/Cargo.toml index 52032c0ae..2c7f42ce8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -232,7 +232,8 @@ tokio-postgres-rustls = "0.14.0" # Utilities and Tools anyhow = "1.0.104" arc-swap = "1.9.2" -astral-tokio-tar = "0.7.0" +# RUSTFS_COMPAT_TODO(tokio-tar-extension-limits): keep the fork pin until bounded extension parsing is released upstream. Remove after astral-sh/tokio-tar#118 is merged and a published tokio-tar release exposes the extension limits used here. +astral-tokio-tar = { git = "https://github.com/cxymds/tokio-tar.git", rev = "603756478b7668436e464519c77ccac22a99ba96" } atoi = "3.1.0" atomic_enum = "0.3.0" aws-config = { version = "1.11.0" } diff --git a/crates/e2e_test/src/multipart_auth_test.rs b/crates/e2e_test/src/multipart_auth_test.rs index 320b3ff3a..b15998512 100644 --- a/crates/e2e_test/src/multipart_auth_test.rs +++ b/crates/e2e_test/src/multipart_auth_test.rs @@ -3456,6 +3456,62 @@ async fn test_signed_put_object_extract_expands_tar_entries_with_prefix_headers( Ok(()) } +#[tokio::test] +async fn test_signed_put_object_extract_ignore_dirs_skips_unauthorized_directory() +-> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let bucket = "signed-extract-ignore-dirs-auth"; + let archive_key = "bundle.tar"; + let allowed_member = "allowed/member.txt"; + let denied_directory = "denied/"; + let username = "snowball-ignore-dirs"; + let secret_key = "snowball-ignore-dirs-secret"; + let expected_body = b"allowed-body"; + + let admin_client = env.create_s3_client(); + admin_client.create_bucket().bucket(bucket).send().await?; + create_restricted_user(&env, username, secret_key).await?; + + let policy = serde_json::json!({ + "Version": "2012-10-17", + "Statement": [{ + "Effect": "Allow", + "Principal": { "AWS": [username] }, + "Action": ["s3:PutObject"], + "Resource": [ + format!("arn:aws:s3:::{bucket}/{archive_key}"), + format!("arn:aws:s3:::{bucket}/{allowed_member}") + ] + }] + }) + .to_string(); + admin_client.put_bucket_policy().bucket(bucket).policy(policy).send().await?; + + let restricted_client = restricted_user_client(&env, username, secret_key); + let tar_bytes = make_tar(&[(allowed_member, expected_body)], &[denied_directory]).await; + restricted_client + .put_object() + .bucket(bucket) + .key(archive_key) + .body(ByteStream::from(tar_bytes)) + .customize() + .mutate_request(|req| { + req.headers_mut().insert("x-amz-meta-snowball-auto-extract", "true"); + req.headers_mut().insert("x-amz-meta-snowball-ignore-dirs", "true"); + }) + .send() + .await?; + + let stored = admin_client.get_object().bucket(bucket).key(allowed_member).send().await?; + assert_eq!(stored.body.collect().await?.into_bytes().as_ref(), expected_body); + + Ok(()) +} + #[tokio::test] async fn test_signed_put_object_extract_preserves_request_metadata_on_extracted_objects() -> Result<(), Box> { diff --git a/crates/e2e_test/src/negative_sigv4_test.rs b/crates/e2e_test/src/negative_sigv4_test.rs index 19c3b4b0c..ca4d1e63c 100644 --- a/crates/e2e_test/src/negative_sigv4_test.rs +++ b/crates/e2e_test/src/negative_sigv4_test.rs @@ -36,9 +36,10 @@ use crate::common::{RustFSTestEnvironment, init_logging, local_http_client}; use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; -use rustfs_signer::constants::UNSIGNED_PAYLOAD; +use rustfs_signer::constants::{UNSIGNED_PAYLOAD, UNSIGNED_PAYLOAD_TRAILER}; use rustfs_signer::request_signature_v4::{SIGN_V4_ALGORITHM, get_scope, get_signature, get_signing_key}; use std::fmt::Write as _; +use std::io::Cursor; use time::macros::format_description; use time::{Duration, OffsetDateTime}; use tracing::info; @@ -98,15 +99,37 @@ impl SigV4 { /// header AND folded into the canonical request — pass the hash of the /// body you *claim* to send, which may differ from what you actually send. fn sign(&self, method: &str, path: &str, canonical_query: &str, content_sha256: &str) -> SignedHeaders { - let amz_date = amz_datetime(self.time); - let signed_headers = "host;x-amz-content-sha256;x-amz-date"; + self.sign_with_extra_headers(method, path, canonical_query, content_sha256, &[]) + } - let canonical_headers = format!( - "host:{host}\nx-amz-content-sha256:{sha}\nx-amz-date:{date}\n", - host = self.host, - sha = content_sha256, - date = amz_date, - ); + /// Sign additional request headers while preserving SigV4's lowercase, + /// lexicographically sorted canonical-header representation. + fn sign_with_extra_headers( + &self, + method: &str, + path: &str, + canonical_query: &str, + content_sha256: &str, + extra_signed_headers: &[(&str, &str)], + ) -> SignedHeaders { + let amz_date = amz_datetime(self.time); + let mut canonical_header_values = vec![ + ("host", self.host.as_str()), + ("x-amz-content-sha256", content_sha256), + ("x-amz-date", amz_date.as_str()), + ]; + canonical_header_values.extend(extra_signed_headers.iter().copied()); + canonical_header_values.sort_unstable_by(|left, right| left.0.cmp(right.0)); + + let signed_headers = canonical_header_values + .iter() + .map(|(name, _)| *name) + .collect::>() + .join(";"); + let mut canonical_headers = String::new(); + for (name, value) in canonical_header_values { + let _ = writeln!(canonical_headers, "{name}:{value}"); + } let canonical_request = format!("{method}\n{path}\n{canonical_query}\n{canonical_headers}\n{signed_headers}\n{content_sha256}"); @@ -179,6 +202,34 @@ async fn setup(env: &mut RustFSTestEnvironment) -> Result<(), Box Result, Box> { + let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new())); + let mut header = tokio_tar::Header::new_gnu(); + header.set_size(member_body.len() as u64); + header.set_mode(0o644); + header.set_cksum(); + builder.append_data(&mut header, member_key, Cursor::new(member_body)).await?; + Ok(builder.into_inner().await?.into_inner()) +} + +fn sha256_base64(data: &[u8]) -> String { + use sha2::{Digest, Sha256}; + + base64_simd::STANDARD.encode_to_string(Sha256::digest(data)) +} + +fn encode_unsigned_aws_chunked_with_sha256_trailer(decoded: &[u8]) -> Vec { + let checksum = sha256_base64(decoded); + let mut encoded = format!("{:x}\r\n", decoded.len()).into_bytes(); + encoded.extend_from_slice(decoded); + encoded.extend_from_slice(b"\r\n0\r\n\r\n"); + encoded.extend_from_slice(format!("x-amz-checksum-sha256:{checksum}").as_bytes()); + encoded +} + /// Positive control: a correctly hand-signed request must succeed. Without /// this, every negative assertion below could pass for the wrong reason (a /// broken signer that never produces a valid signature). @@ -249,6 +300,128 @@ async fn tampered_signature_returns_signature_does_not_match() -> Result<(), Box Ok(()) } +/// `STREAMING-UNSIGNED-PAYLOAD-TRAILER` disables per-chunk signatures, not the +/// seed/header SigV4 signature. A forged request must be rejected before the +/// Snowball handler can publish any archive member. +#[tokio::test] +async fn snowball_streaming_unsigned_trailer_rejects_forged_signature() -> Result<(), Box> { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + setup(&mut env).await?; + + let archive_key = "forged-streaming-snowball.tar"; + let member_key = "must-not-be-published.txt"; + let archive = build_single_member_archive(member_key, b"forged request payload").await?; + let decoded_content_length = archive.len().to_string(); + let encoded_body = encode_unsigned_aws_chunked_with_sha256_trailer(&archive); + let path = format!("/{BUCKET}/{archive_key}"); + + let mut signer = SigV4::new(&env); + signer.secret_key = "wrong-secret-for-forged-streaming-request".to_string(); + let extra_signed_headers = [ + ("content-encoding", "aws-chunked"), + ("x-amz-decoded-content-length", decoded_content_length.as_str()), + ("x-amz-meta-snowball-auto-extract", "true"), + ("x-amz-trailer", "x-amz-checksum-sha256"), + ]; + let headers = signer.sign_with_extra_headers("PUT", &path, "", UNSIGNED_PAYLOAD_TRAILER, &extra_signed_headers); + + let response = local_http_client() + .put(format!("{}{}", env.url, path)) + .header("authorization", &headers.authorization) + .header("content-encoding", "aws-chunked") + .header("x-amz-content-sha256", &headers.content_sha256) + .header("x-amz-date", &headers.amz_date) + .header("x-amz-decoded-content-length", &decoded_content_length) + .header("x-amz-meta-snowball-auto-extract", "true") + .header("x-amz-trailer", "x-amz-checksum-sha256") + .body(encoded_body) + .send() + .await?; + let status = response.status(); + let body = response.text().await?; + assert_eq!(status.as_u16(), 403, "forged streaming signature must be 403, body:\n{body}"); + assert_error_code(&body, "SignatureDoesNotMatch"); + + let absent = env + .create_s3_client() + .get_object() + .bucket(BUCKET) + .key(member_key) + .send() + .await + .expect_err("a forged streaming request must not publish a Snowball member"); + assert_eq!(absent.raw_response().map(|response| response.status().as_u16()), Some(404)); + assert_eq!(absent.as_service_error().and_then(ProvideErrorMetadata::code), Some("NoSuchKey")); + + env.stop_server(); + Ok(()) +} + +/// Snowball must consume the complete aws-chunked body before reading the +/// trailing checksum exported by s3s into the PutObject response. +#[tokio::test] +async fn snowball_streaming_unsigned_trailer_returns_sha256_checksum() -> Result<(), Box> { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + setup(&mut env).await?; + + let archive_key = "valid-streaming-snowball.tar"; + let member_key = "streaming-checksum-member.txt"; + let member_body = b"valid streaming Snowball payload"; + let archive = build_single_member_archive(member_key, member_body).await?; + let expected_checksum = sha256_base64(&archive); + let decoded_content_length = archive.len().to_string(); + let encoded_body = encode_unsigned_aws_chunked_with_sha256_trailer(&archive); + let path = format!("/{BUCKET}/{archive_key}"); + + let signer = SigV4::new(&env); + let extra_signed_headers = [ + ("content-encoding", "aws-chunked"), + ("x-amz-decoded-content-length", decoded_content_length.as_str()), + ("x-amz-meta-snowball-auto-extract", "true"), + ("x-amz-sdk-checksum-algorithm", "SHA256"), + ("x-amz-trailer", "x-amz-checksum-sha256"), + ]; + let headers = signer.sign_with_extra_headers("PUT", &path, "", UNSIGNED_PAYLOAD_TRAILER, &extra_signed_headers); + + let response = local_http_client() + .put(format!("{}{}", env.url, path)) + .header("authorization", &headers.authorization) + .header("content-encoding", "aws-chunked") + .header("x-amz-content-sha256", &headers.content_sha256) + .header("x-amz-date", &headers.amz_date) + .header("x-amz-decoded-content-length", &decoded_content_length) + .header("x-amz-meta-snowball-auto-extract", "true") + .header("x-amz-sdk-checksum-algorithm", "SHA256") + .header("x-amz-trailer", "x-amz-checksum-sha256") + .body(encoded_body) + .send() + .await?; + let status = response.status(); + let response_checksum = response + .headers() + .get("x-amz-checksum-sha256") + .and_then(|value| value.to_str().ok()) + .map(str::to_owned); + let response_body = response.text().await?; + assert_eq!(status.as_u16(), 200, "valid streaming Snowball PUT failed, body:\n{response_body}"); + assert_eq!(response_checksum.as_deref(), Some(expected_checksum.as_str())); + + let member = env + .create_s3_client() + .get_object() + .bucket(BUCKET) + .key(member_key) + .send() + .await?; + let stored = member.body.collect().await?.into_bytes(); + assert_eq!(stored.as_ref(), member_body); + + env.stop_server(); + Ok(()) +} + /// (b) A valid AccessKeyId paired with the wrong secret key must be rejected /// with SignatureDoesNotMatch / 403. #[tokio::test] diff --git a/crates/e2e_test/src/snowball_auto_extract_test.rs b/crates/e2e_test/src/snowball_auto_extract_test.rs index 2077af321..5e705cc8b 100644 --- a/crates/e2e_test/src/snowball_auto_extract_test.rs +++ b/crates/e2e_test/src/snowball_auto_extract_test.rs @@ -17,8 +17,9 @@ mod tests { use crate::common::{RustFSTestEnvironment, init_logging}; use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; + use flate2::{Compression, write::GzEncoder}; use std::error::Error; - use std::io::Cursor; + use std::io::{Cursor, Write}; async fn build_test_archive() -> Result, Box> { let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new())); @@ -69,12 +70,50 @@ mod tests { Ok(builder.into_inner().await?.into_inner()) } - fn build_archive_with_parent_dir_entry(victim_bucket: &str) -> Vec { - let path = format!("../{victim_bucket}/evil-injected.txt"); - let data = b"injected-body"; + async fn build_archive_with_invalid_checksum() -> Result, Box> { + let mut archive = build_test_archive().await?; + archive[0] ^= 1; + Ok(archive) + } + + async fn build_archive_with_negative_gnu_mtime() -> Result, Box> { + let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new())); + let mut header = tokio_tar::Header::new_gnu(); + header.set_size(b"negative-mtime-body".len() as u64); + header.set_mode(0o644); + header.as_old_mut().mtime.fill(0xff); + builder + .append_data(&mut header, "negative-mtime.txt", Cursor::new(b"negative-mtime-body".as_slice())) + .await?; + Ok(builder.into_inner().await?.into_inner()) + } + + fn gzip_member(payload: &[u8]) -> Result, Box> { + let mut encoder = GzEncoder::new(Vec::new(), Compression::default()); + encoder.write_all(payload)?; + Ok(encoder.finish()?) + } + + async fn build_concatenated_gzip_archive() -> Result, Box> { + let archive = build_test_archive().await?; + let split_at = archive.len() / 2; + let mut encoded = gzip_member(&archive[..split_at])?; + encoded.extend(gzip_member(&archive[split_at..])?); + Ok(encoded) + } + + async fn build_gzip_archive_with_invalid_crc() -> Result, Box> { + let mut encoded = gzip_member(&build_test_archive().await?)?; + let crc_offset = encoded.len().checked_sub(8).expect("gzip fixture must contain a trailer"); + encoded[crc_offset] ^= 1; + Ok(encoded) + } + + fn append_raw_tar_entry_with_type(archive: &mut Vec, path: &[u8], data: &[u8], entry_type: u8) { + assert!(path.len() <= 100, "raw TAR fixture path must fit in the name field"); let mut header = [0u8; 512]; - header[..path.len()].copy_from_slice(path.as_bytes()); + header[..path.len()].copy_from_slice(path); header[100..108].copy_from_slice(b"0000644\0"); header[108..116].copy_from_slice(b"0000000\0"); header[116..124].copy_from_slice(b"0000000\0"); @@ -82,7 +121,7 @@ mod tests { header[124..136].copy_from_slice(size.as_bytes()); header[136..148].copy_from_slice(b"00000000000\0"); header[148..156].fill(b' '); - header[156] = b'0'; + header[156] = entry_type; header[257..263].copy_from_slice(b"ustar\0"); header[263..265].copy_from_slice(b"00"); @@ -90,11 +129,36 @@ mod tests { let checksum = format!("{:06o}\0 ", checksum); header[148..156].copy_from_slice(checksum.as_bytes()); - let mut archive = Vec::new(); archive.extend_from_slice(&header); archive.extend_from_slice(data); let padding = (512 - (data.len() % 512)) % 512; archive.extend(std::iter::repeat_n(0, padding)); + } + + fn append_raw_tar_entry(archive: &mut Vec, path: &[u8], data: &[u8]) { + append_raw_tar_entry_with_type(archive, path, data, b'0'); + } + + fn build_archive_with_parent_dir_entry(victim_bucket: &str) -> Vec { + let path = format!("../{victim_bucket}/evil-injected.txt"); + let mut archive = Vec::new(); + append_raw_tar_entry(&mut archive, path.as_bytes(), b"injected-body"); + archive.extend_from_slice(&[0u8; 1024]); + archive + } + + fn build_archive_with_invalid_utf8_entry() -> Vec { + let mut archive = Vec::new(); + append_raw_tar_entry(&mut archive, b"invalid-\xff.txt", b"ignored-body"); + append_raw_tar_entry(&mut archive, b"valid.txt", b"valid-body"); + archive.extend_from_slice(&[0u8; 1024]); + archive + } + + fn build_archive_with_invalid_utf8_symlink() -> Vec { + let mut archive = Vec::new(); + append_raw_tar_entry_with_type(&mut archive, b"invalid-\xff-link", b"", b'2'); + append_raw_tar_entry(&mut archive, b"valid.txt", b"valid-body"); archive.extend_from_slice(&[0u8; 1024]); archive } @@ -263,6 +327,113 @@ mod tests { Ok(()) } + #[tokio::test] + async fn snowball_auto_extract_accepts_negative_gnu_mtime() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-negative-mtime"; + client.create_bucket().bucket(bucket).send().await?; + client + .put_object() + .bucket(bucket) + .key("fixture.tar") + .metadata("Snowball-Auto-Extract", "true") + .body(ByteStream::from(build_archive_with_negative_gnu_mtime().await?)) + .send() + .await?; + + let object = client.get_object().bucket(bucket).key("negative-mtime.txt").send().await?; + assert_eq!(object.body.collect().await?.into_bytes().as_ref(), b"negative-mtime-body"); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_consumes_concatenated_gzip_members() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-concatenated-gzip"; + client.create_bucket().bucket(bucket).send().await?; + client + .put_object() + .bucket(bucket) + .key("fixture.tar.gz") + .metadata("Snowball-Auto-Extract", "true") + .body(ByteStream::from(build_concatenated_gzip_archive().await?)) + .send() + .await?; + + let object = client.get_object().bucket(bucket).key("root.txt").send().await?; + assert_eq!(object.body.collect().await?.into_bytes().as_ref(), b"root payload\n"); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_rejects_gzip_crc_error_when_ignore_errors_enabled() -> Result<(), Box> + { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-gzip-crc-ignore-errors"; + client.create_bucket().bucket(bucket).send().await?; + let err = client + .put_object() + .bucket(bucket) + .key("fixture.tar.gz") + .metadata("Snowball-Auto-Extract", "true") + .metadata("Minio-Snowball-Ignore-Errors", "true") + .body(ByteStream::from(build_gzip_archive_with_invalid_crc().await?)) + .send() + .await + .expect_err("gzip integrity failures must remain fatal under ignore-errors"); + + assert_eq!(err.into_service_error().code(), Some("InvalidArgument")); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_rejects_mismatched_content_md5() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-content-md5"; + client.create_bucket().bucket(bucket).send().await?; + let err = client + .put_object() + .bucket(bucket) + .key("fixture.tar") + .metadata("Snowball-Auto-Extract", "true") + .content_md5("AAAAAAAAAAAAAAAAAAAAAA==") + .body(ByteStream::from(build_test_archive().await?)) + .send() + .await + .expect_err("mismatched Content-MD5 must fail after the raw body reaches EOF"); + + assert_eq!(err.into_service_error().code(), Some("BadDigest")); + + env.stop_server(); + Ok(()) + } + #[tokio::test] async fn snowball_auto_extract_ignores_invalid_entries_when_requested() -> Result<(), Box> { init_logging(); @@ -299,7 +470,100 @@ mod tests { } #[tokio::test] - async fn snowball_auto_extract_rejects_parent_dir_entry_without_cross_bucket_write() + async fn snowball_auto_extract_skips_non_utf8_symlink_without_ignore_errors() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-invalid-utf8-link"; + client.create_bucket().bucket(bucket).send().await?; + + client + .put_object() + .bucket(bucket) + .key("fixture.tar") + .metadata("Snowball-Auto-Extract", "true") + .body(ByteStream::from(build_archive_with_invalid_utf8_symlink())) + .send() + .await?; + + let valid = client.get_object().bucket(bucket).key("valid.txt").send().await?; + assert_eq!(valid.body.collect().await?.into_bytes().as_ref(), b"valid-body"); + let listed = client.list_objects_v2().bucket(bucket).send().await?; + let keys: Vec<_> = listed.contents().iter().filter_map(|entry| entry.key()).collect(); + assert_eq!(keys, vec!["valid.txt"]); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_skips_non_utf8_member_without_lossy_key_collision() -> Result<(), Box> + { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-invalid-utf8"; + client.create_bucket().bucket(bucket).send().await?; + + client + .put_object() + .bucket(bucket) + .key("fixture.tar") + .metadata("Snowball-Auto-Extract", "true") + .metadata("Minio-Snowball-Ignore-Errors", "true") + .body(ByteStream::from(build_archive_with_invalid_utf8_entry())) + .send() + .await?; + + let valid = client.get_object().bucket(bucket).key("valid.txt").send().await?; + assert_eq!(valid.body.collect().await?.into_bytes().as_ref(), b"valid-body"); + let listed = client.list_objects_v2().bucket(bucket).send().await?; + let keys: Vec<_> = listed.contents().iter().filter_map(|entry| entry.key()).collect(); + assert_eq!(keys, vec!["valid.txt"]); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_rejects_corrupt_tar_when_ignore_errors_enabled() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + let bucket = "snowball-corrupt-ignore-errors"; + let archive = build_archive_with_invalid_checksum().await?; + client.create_bucket().bucket(bucket).send().await?; + + let err = client + .put_object() + .bucket(bucket) + .key("fixture.tar") + .metadata("Snowball-Auto-Extract", "true") + .metadata("Minio-Snowball-Ignore-Errors", "true") + .body(ByteStream::from(archive)) + .send() + .await + .expect_err("corrupt TAR structure must remain fatal under ignore-errors"); + assert_eq!(err.into_service_error().code(), Some("InvalidArgument")); + + let listed = client.list_objects_v2().bucket(bucket).send().await?; + assert!(listed.contents().is_empty(), "corrupt archive must not produce objects"); + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn snowball_auto_extract_rejects_parent_dir_entry_even_when_ignore_errors_enabled() -> Result<(), Box> { init_logging(); @@ -319,6 +583,7 @@ mod tests { .bucket(attacker_bucket) .key("fixture.tar") .metadata("Snowball-Auto-Extract", "true") + .metadata("Minio-Snowball-Ignore-Errors", "true") .body(ByteStream::from(archive)) .send() .await diff --git a/crates/zip/src/lib.rs b/crates/zip/src/lib.rs index 82aa8867a..1bcc78506 100644 --- a/crates/zip/src/lib.rs +++ b/crates/zip/src/lib.rs @@ -47,7 +47,12 @@ pub struct ArchiveLimits { pub max_entries: usize, pub max_entry_size: u64, pub max_total_unpacked_size: u64, + pub max_decoded_size: u64, pub max_path_length: usize, + pub max_pax_metadata_size: u64, + pub max_total_pax_metadata_size: u64, + pub max_pax_metadata_records: usize, + pub max_total_pax_metadata_records: usize, pub validate_entry_paths: bool, } @@ -57,7 +62,12 @@ impl Default for ArchiveLimits { max_entries: 100_000, max_entry_size: 1_073_741_824, max_total_unpacked_size: 10_737_418_240, + max_decoded_size: 11_811_160_064, max_path_length: 1024, + max_pax_metadata_size: 1_048_576, + max_total_pax_metadata_size: 67_108_864, + max_pax_metadata_records: 4_096, + max_total_pax_metadata_records: 100_000, validate_entry_paths: true, } } @@ -100,11 +110,31 @@ impl CompressionFormat { let reader = BufReader::new(input); let decoder: Box = match self { - CompressionFormat::Gzip => Box::new(GzipDecoder::new(reader)), - CompressionFormat::Bzip2 => Box::new(BzDecoder::new(reader)), - CompressionFormat::Zlib => Box::new(ZlibDecoder::new(reader)), - CompressionFormat::Xz => Box::new(XzDecoder::new(reader)), - CompressionFormat::Zstd => Box::new(ZstdDecoder::new(reader)), + CompressionFormat::Gzip => { + let mut decoder = GzipDecoder::new(reader); + decoder.multiple_members(true); + Box::new(decoder) + } + CompressionFormat::Bzip2 => { + let mut decoder = BzDecoder::new(reader); + decoder.multiple_members(true); + Box::new(decoder) + } + CompressionFormat::Zlib => { + let mut decoder = ZlibDecoder::new(reader); + decoder.multiple_members(true); + Box::new(decoder) + } + CompressionFormat::Xz => { + let mut decoder = XzDecoder::new(reader); + decoder.multiple_members(true); + Box::new(decoder) + } + CompressionFormat::Zstd => { + let mut decoder = ZstdDecoder::new(reader); + decoder.multiple_members(true); + Box::new(decoder) + } CompressionFormat::Tar => Box::new(reader), CompressionFormat::Zip => { return Err(ZipError::UnsupportedFormat { @@ -160,6 +190,30 @@ mod tests { assert_eq!(decoded, b"payload"); } + #[tokio::test] + async fn test_get_decoder_consumes_concatenated_gzip_members() { + async fn gzip_member(payload: &[u8]) -> Vec { + let mut encoder = GzipEncoder::new(Vec::new()); + encoder.write_all(payload).await.expect("gzip encode should succeed"); + encoder.shutdown().await.expect("gzip encoder shutdown should succeed"); + encoder.into_inner() + } + + let mut encoded = gzip_member(b"first-").await; + encoded.extend(gzip_member(b"second").await); + let mut decoder = CompressionFormat::Gzip + .get_decoder(std::io::Cursor::new(encoded)) + .expect("gzip decoder should be created"); + let mut decoded = Vec::new(); + + decoder + .read_to_end(&mut decoded) + .await + .expect("concatenated gzip members should decode"); + + assert_eq!(decoded, b"first-second"); + } + #[tokio::test] async fn test_get_decoder_rejects_zip_and_unknown_formats() { let zip_err = CompressionFormat::Zip diff --git a/deny.toml b/deny.toml index c6fb8facc..b8e825836 100644 --- a/deny.toml +++ b/deny.toml @@ -36,6 +36,10 @@ unknown-registry = "deny" unknown-git = "deny" allow-registry = ["https://github.com/rust-lang/crates.io-index"] allow-git = [ + # Temporary tokio-tar fork pinned to the reviewed bounded extension parser + # change while astral-sh/tokio-tar#118 awaits an upstream release. + # owner: cxymds review: 2026-10 + "https://github.com/cxymds/tokio-tar.git", # Official s3s repository. Temporarily pinned to the merged generic REST # SigV4 payload-checksum fix until it is available in a crates.io release. # owner: marshawcoco review: 2026-10 diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index 8b56a5939..297e7da3d 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -12,6 +12,7 @@ for later deletion. ## Open Items +- `tokio-tar-extension-limits` bounded archive extension parsing: Snowball extraction depends on explicit limits for GNU long-name, GNU long-link, and PAX extension payloads, while the released tokio-tar API does not expose those limits. Keep the reviewed fork pin so untrusted archives cannot allocate unbounded extension metadata. Remove the fork pin after astral-sh/tokio-tar#118 is merged and a published tokio-tar release exposes the extension limits used here. - `backlog-2102` rc.2/rc.3 empty scanner usage floor recovery: old DeleteBucket cleanup could synthesize an empty incomplete v2 usage primary/backup before leadership added an epoch, while newer scanners require a durable authoritative baseline identity. New scanners recognize only that exact serialized empty-fence shape, preserve its epoch through a CAS-protected recovery marker, and rebuild namespace coverage without treating zero usage as authoritative. Remove this recovery path and marker after rc.2 and rc.3 are no longer supported direct-upgrade sources. - `s3gate-metadata-xml` persisted bucket XML migration: mixed-version site-replication peers, retained `.metadata.bin` objects, and backup archives can all carry XML written by the s3s codec, so the gateway migration must keep the legacy codec available until every stored form has crossed a verified rewrite boundary. Remove the legacy s3s parser and serializer only after the minimum supported direct-upgrade release reads and writes every persisted XML configuration family through the gateway codec, the four-way D1-D5 gate has remained clean for one full support window, every supported mixed-version site-replication topology has completed its writer upgrade, and migration tooling has verified or rewritten every retained bucket metadata object and restorable backup archive. - `rustfs-6339` legacy bucket policy ID casing: earlier RustFS releases persisted the top-level policy identifier as "ID", while current writes use the S3-compatible "Id" spelling. Readers accept both spellings so retained bucket metadata remains usable after upgrade. Remove the legacy alias after migration tooling has rewritten every retained bucket policy using "ID". diff --git a/rustfs/src/app/object/extract.rs b/rustfs/src/app/object/extract.rs index 0ec1dd160..478489223 100644 --- a/rustfs/src/app/object/extract.rs +++ b/rustfs/src/app/object/extract.rs @@ -16,6 +16,14 @@ use super::*; +// One logical member can be preceded by local PAX, GNU long-name, and GNU +// long-link records. Count all four physical headers without rejecting that +// compatible extension combination. +const EXTRACT_ARCHIVE_PHYSICAL_ENTRY_MULTIPLIER: u64 = 4; +// Sparse maps are metadata, so bound them independently of object byte quotas. +const EXTRACT_ARCHIVE_MAX_SPARSE_ENTRIES: u64 = 4_096; +const EXTRACT_ARCHIVE_MAX_SPARSE_CONTINUATION_BLOCKS: u64 = 256; + fn ensure_legacy_archive_size_within_quota(result: &QuotaCheckResult, total_unpacked_size: u64) -> S3Result<()> { if result.uses_durable_reservations { return Ok(()); @@ -40,45 +48,232 @@ pin_project! { #[pin] inner: R, md5: Md5, + expected_length: u64, + bytes_read: u64, + pending_final_byte: Option, + validating_eof: bool, finished: bool, - etag: Arc>>, + state: Arc>, } } +#[derive(Debug, Default)] +struct ExtractArchiveUploadState { + etag: Option, + body_complete: bool, +} + impl ExtractArchiveEtagReader { - fn new(inner: R, etag: Arc>>) -> Self { + fn new(inner: R, expected_length: u64, state: Arc>) -> Self { Self { inner, md5: Md5::new(), + expected_length, + bytes_read: 0, + pending_final_byte: None, + validating_eof: expected_length == 0, finished: false, - etag, + state, } } } +fn extract_archive_incomplete_body(remaining: u64) -> std::io::Error { + let Ok(remaining) = i64::try_from(remaining) else { + return std::io::Error::new(std::io::ErrorKind::InvalidData, "archive remaining body length exceeds i64"); + }; + std::io::Error::new(std::io::ErrorKind::UnexpectedEof, rustfs_rio::IncompleteBody { remaining }) +} + impl AsyncRead for ExtractArchiveEtagReader { fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { - let this = self.project(); - let before = buf.filled().len(); - match this.inner.poll_read(cx, buf) { - Poll::Pending => Poll::Pending, - Poll::Ready(Ok(())) => { - let filled = &buf.filled()[before..]; - if !filled.is_empty() { - this.md5.update(filled); - } else if !*this.finished { - *this.finished = true; - if let Ok(mut etag) = this.etag.lock() { - *etag = Some(hex_simd::encode_to_string(this.md5.clone().finalize(), hex_simd::AsciiCase::Lower)); + let mut this = self.project(); + if buf.remaining() == 0 || *this.finished { + return Poll::Ready(Ok(())); + } + + loop { + if *this.validating_eof { + let mut probe = [0u8; 1]; + let mut probe_buf = ReadBuf::new(&mut probe); + match this.inner.as_mut().poll_read(cx, &mut probe_buf) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Err(err)) => return Poll::Ready(Err(err)), + Poll::Ready(Ok(())) if !probe_buf.filled().is_empty() => { + return Poll::Ready(Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "archive body exceeded expected Content-Length", + ))); + } + Poll::Ready(Ok(())) => { + if let Ok(mut state) = this.state.lock() + && !state.body_complete + { + state.etag = + Some(hex_simd::encode_to_string(this.md5.clone().finalize(), hex_simd::AsciiCase::Lower)); + state.body_complete = true; + } + *this.validating_eof = false; + *this.finished = true; + if let Some(final_byte) = this.pending_final_byte.take() { + buf.put_slice(&[final_byte]); + } + return Poll::Ready(Ok(())); } } - Poll::Ready(Ok(())) } - Poll::Ready(Err(err)) => Poll::Ready(Err(err)), + + let remaining = *this.expected_length - *this.bytes_read; + if remaining == 1 { + let mut final_byte = [0u8; 1]; + let mut final_buf = ReadBuf::new(&mut final_byte); + match this.inner.as_mut().poll_read(cx, &mut final_buf) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Err(err)) => return Poll::Ready(Err(err)), + Poll::Ready(Ok(())) if final_buf.filled().is_empty() => { + return Poll::Ready(Err(extract_archive_incomplete_body(*this.expected_length - *this.bytes_read))); + } + Poll::Ready(Ok(())) => { + this.md5.update(final_buf.filled()); + *this.bytes_read = match this.bytes_read.checked_add(1) { + Some(bytes_read) => bytes_read, + None => return Poll::Ready(Err(std::io::Error::other("archive read length overflow"))), + }; + *this.pending_final_byte = Some(final_buf.filled()[0]); + *this.validating_eof = true; + continue; + } + } + } + + let max_read = usize::try_from(remaining - 1).unwrap_or(usize::MAX).min(buf.remaining()); + let read_len = { + let target = buf.initialize_unfilled_to(max_read); + let mut limited_buf = ReadBuf::new(target); + match this.inner.as_mut().poll_read(cx, &mut limited_buf) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Err(err)) => return Poll::Ready(Err(err)), + Poll::Ready(Ok(())) if limited_buf.filled().is_empty() => { + return Poll::Ready(Err(extract_archive_incomplete_body(*this.expected_length - *this.bytes_read))); + } + Poll::Ready(Ok(())) => { + this.md5.update(limited_buf.filled()); + limited_buf.filled().len() + } + } + }; + let read = match u64::try_from(read_len) { + Ok(read) => read, + Err(_) => return Poll::Ready(Err(std::io::Error::other("archive read length exceeds u64"))), + }; + *this.bytes_read = match this.bytes_read.checked_add(read) { + Some(bytes_read) => bytes_read, + None => return Poll::Ready(Err(std::io::Error::other("archive read length overflow"))), + }; + buf.advance(read_len); + return Poll::Ready(Ok(())); } } } +pin_project! { + struct ExtractMemberReadTracker { + #[pin] + inner: HashReader, + failed: Arc, + } +} + +impl AsyncRead for ExtractMemberReadTracker { + fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + let this = self.project(); + match this.inner.poll_read(cx, buf) { + Poll::Ready(Err(err)) => { + this.failed.store(true, Ordering::Release); + Poll::Ready(Err(err)) + } + other => other, + } + } +} + +impl rustfs_rio::EtagResolvable for ExtractMemberReadTracker { + fn try_resolve_etag(&mut self) -> Option { + rustfs_rio::EtagResolvable::try_resolve_etag(&mut self.inner) + } +} + +impl rustfs_rio::HashReaderDetector for ExtractMemberReadTracker {} + +impl rustfs_rio::TryGetIndex for ExtractMemberReadTracker { + fn try_get_index(&self) -> Option<&rustfs_rio::Index> { + rustfs_rio::TryGetIndex::try_get_index(&self.inner) + } +} + +fn track_extract_member_read_errors(reader: HashReader) -> std::io::Result<(HashReader, Arc)> { + let size = reader.size(); + let actual_size = reader.actual_size(); + let failed = Arc::new(AtomicBool::new(false)); + let tracker = ExtractMemberReadTracker { + inner: reader, + failed: failed.clone(), + }; + let mut tracked = HashReader::from_reader(tracker, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)?; + tracked.update_params(size, actual_size, None); + Ok((tracked, failed)) +} + +fn should_ignore_extract_member_write_error(ignore_errors: bool, member_read_failed: &AtomicBool) -> bool { + ignore_errors && !member_read_failed.load(Ordering::Acquire) +} + +pin_project! { + struct ExtractDecodedLimitReader { + #[pin] + inner: R, + remaining: u64, + } +} + +impl ExtractDecodedLimitReader { + fn new(inner: R, limit: u64) -> Self { + Self { inner, remaining: limit } + } +} + +impl AsyncRead for ExtractDecodedLimitReader { + fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + if buf.remaining() == 0 { + return Poll::Ready(Ok(())); + } + + let mut this = self.project(); + let allowed = this.remaining.saturating_add(1); + let max_read = usize::try_from(allowed).unwrap_or(usize::MAX).min(buf.remaining()); + let read_len = { + let unfilled = buf.initialize_unfilled_to(max_read); + let mut limited = ReadBuf::new(unfilled); + match this.inner.as_mut().poll_read(cx, &mut limited) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Err(err)) => return Poll::Ready(Err(err)), + Poll::Ready(Ok(())) => limited.filled().len(), + } + }; + let read = u64::try_from(read_len).unwrap_or(u64::MAX); + if read > *this.remaining { + return Poll::Ready(Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "archive decoded size exceeds limit", + ))); + } + + *this.remaining -= read; + buf.advance(read_len); + Poll::Ready(Ok(())) + } +} + const AMZ_SNOWBALL_EXTRACT_COMPAT: &str = "X-Amz-Snowball-Auto-Extract"; #[cfg(test)] @@ -198,8 +393,186 @@ pub fn normalize_extract_entry_key(path: &str, prefix: Option<&str>, is_dir: boo rustfs_utils::path::normalize_extract_entry_key(path, prefix, is_dir).map_err(|msg| s3_error!(InvalidArgument, "{msg}")) } -fn map_extract_archive_error(err: impl std::fmt::Display) -> S3Error { - s3_error!(InvalidArgument, "Failed to process archive entry: {}", err) +fn map_extract_archive_error(err: std::io::Error) -> S3Error { + let message = err.to_string(); + let api_error = ApiError::from(err); + if matches!(api_error.code, S3ErrorCode::BadDigest | S3ErrorCode::IncompleteBody) { + return api_error.into(); + } + + let mut archive_error = s3_error!(InvalidArgument, "Failed to process archive entry: {}", message); + archive_error.set_source(Box::new(api_error)); + archive_error +} + +#[derive(Debug)] +enum ExtractEntryError { + Fatal(S3Error), + Recoverable(S3Error), +} + +impl ExtractEntryError { + fn into_s3_error(self) -> S3Error { + match self { + Self::Fatal(err) | Self::Recoverable(err) => err, + } + } + + fn ignore_or_return(self, ignore_errors: bool) -> S3Result<()> { + match self { + Self::Recoverable(_) if ignore_errors => Ok(()), + Self::Fatal(err) | Self::Recoverable(err) => Err(err), + } + } + + #[cfg(test)] + fn is_recoverable(&self) -> bool { + matches!(self, Self::Recoverable(_)) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ExtractEntryDisposition { + File, + Directory, + FormatSkip, +} + +fn classify_extract_entry_type(entry_type: tokio_tar::EntryType) -> ExtractEntryDisposition { + use tokio_tar::EntryType; + + match entry_type { + EntryType::Regular | EntryType::Char | EntryType::Block | EntryType::Fifo | EntryType::GNUSparse => { + ExtractEntryDisposition::File + } + EntryType::Directory => ExtractEntryDisposition::Directory, + EntryType::Link + | EntryType::Symlink + | EntryType::GNULongName + | EntryType::GNULongLink + | EntryType::Continuous + | EntryType::XGlobalHeader + | EntryType::XHeader + | EntryType::SolarisXHeader + | EntryType::Other(_) => ExtractEntryDisposition::FormatSkip, + _ => ExtractEntryDisposition::FormatSkip, + } +} + +fn extract_entry_quota_growth(disposition: ExtractEntryDisposition, entry_size: u64) -> u64 { + match disposition { + ExtractEntryDisposition::File => entry_size, + ExtractEntryDisposition::Directory | ExtractEntryDisposition::FormatSkip => 0, + } +} + +fn extract_archive_entry_mod_time(header: &tokio_tar::Header) -> S3Result> { + let modified_at_secs = header.mtime().map_err(map_extract_archive_error)?; + // GNU base-256 represents negative values with an all-sign-extended first + // byte. MinIO treats non-positive mtimes as unset, while the tar parser's + // unsigned API exposes `-1` as `u64::MAX`. + if modified_at_secs == 0 || header.as_old().mtime[0] == 0xff { + return Ok(None); + } + + let modified_at_secs = i64::try_from(modified_at_secs) + .map_err(|_| object_s3_error(S3ErrorCode::InvalidArgument, "Archive entry modification time is out of range"))?; + OffsetDateTime::from_unix_timestamp(modified_at_secs) + .map(Some) + .map_err(|_| object_s3_error(S3ErrorCode::InvalidArgument, "Archive entry modification time is out of range")) +} + +fn strict_extract_entry_path(path: &[u8]) -> Result<&str, ExtractEntryError> { + std::str::from_utf8(path) + .map_err(|_| ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Archive entry path must be valid UTF-8"))) +} + +fn is_empty_extract_entry_path(path: &str) -> bool { + path.is_empty() || path == "." || path == "./" +} + +fn validate_extract_member_key(path: &str, limits: ArchiveLimits) -> Result<(), ExtractEntryError> { + validate_put_object_extract_entry_path(path, limits).map_err(ExtractEntryError::Recoverable)?; + validate_object_key(path, "PUT").map_err(ExtractEntryError::Recoverable) +} + +fn record_extract_pax_metadata_bytes( + entry_size: &mut u64, + total_size: &mut u64, + key_size: usize, + value_size: usize, + limits: ArchiveLimits, +) -> Result<(), ExtractEntryError> { + let record_size = key_size + .checked_add(value_size) + .and_then(|size| u64::try_from(size).ok()) + .ok_or_else(|| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive PAX metadata size overflowed while processing entry", + )) + })?; + *entry_size = (*entry_size).checked_add(record_size).ok_or_else(|| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive PAX metadata size overflowed while processing entry", + )) + })?; + *total_size = (*total_size).checked_add(record_size).ok_or_else(|| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive total PAX metadata size overflowed", + )) + })?; + + if *entry_size > limits.max_pax_metadata_size { + return Err(ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive PAX metadata exceeds per-entry limit", + ))); + } + if *total_size > limits.max_total_pax_metadata_size { + return Err(ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive total PAX metadata exceeds limit", + ))); + } + + Ok(()) +} + +fn record_extract_pax_metadata_record( + entry_records: &mut usize, + total_records: &mut usize, + limits: ArchiveLimits, +) -> Result<(), ExtractEntryError> { + *entry_records = entry_records.checked_add(1).ok_or_else(|| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive PAX metadata record count overflowed", + )) + })?; + *total_records = total_records.checked_add(1).ok_or_else(|| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive total PAX metadata record count overflowed", + )) + })?; + + if *entry_records > limits.max_pax_metadata_records { + return Err(ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive PAX metadata record count exceeds per-entry limit", + ))); + } + if *total_records > limits.max_total_pax_metadata_records { + return Err(ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + "Archive total PAX metadata record count exceeds limit", + ))); + } + + Ok(()) } #[derive(Debug, Default)] @@ -210,6 +583,39 @@ struct ExtractEntryPaxAuthorization { object_lock_retain_until_date: Option, } +async fn count_extract_entry_pax_metadata( + entry: &mut tokio_tar::Entry>, + total_pax_metadata_size: &mut u64, + total_pax_metadata_records: &mut usize, + limits: ArchiveLimits, +) -> Result<(), ExtractEntryError> +where + R: AsyncRead + Send + Unpin + 'static, +{ + let Some(extensions) = entry + .pax_extensions() + .await + .map_err(|err| ExtractEntryError::Fatal(map_extract_archive_error(err)))? + else { + return Ok(()); + }; + + let mut entry_pax_metadata_size = 0u64; + let mut entry_pax_metadata_records = 0usize; + for ext in extensions { + let ext = ext.map_err(|err| ExtractEntryError::Fatal(map_extract_archive_error(err)))?; + record_extract_pax_metadata_record(&mut entry_pax_metadata_records, total_pax_metadata_records, limits)?; + record_extract_pax_metadata_bytes( + &mut entry_pax_metadata_size, + total_pax_metadata_size, + ext.key_bytes().len(), + ext.value_bytes().len(), + limits, + )?; + } + Ok(()) +} + async fn apply_extract_entry_pax_extensions( entry: &mut tokio_tar::Entry>, bucket: &str, @@ -217,27 +623,43 @@ async fn apply_extract_entry_pax_extensions( object_lock_config_state: &metadata_sys::ObjectLockConfigState, metadata: &mut HashMap, opts: &mut ObjectOptions, -) -> S3Result +) -> Result where R: AsyncRead + Send + Unpin + 'static, { - let Some(extensions) = entry.pax_extensions().await.map_err(map_extract_archive_error)? else { + let Some(extensions) = entry + .pax_extensions() + .await + .map_err(|err| ExtractEntryError::Fatal(map_extract_archive_error(err)))? + else { return Ok(ExtractEntryPaxAuthorization::default()); }; let mut pax_headers = HeaderMap::new(); let mut pax_version_id = None; for ext in extensions { - let ext = ext.map_err(map_extract_archive_error)?; - let key = ext.key().map_err(map_extract_archive_error)?; - let value = ext.value().map_err(map_extract_archive_error)?; + let ext = ext.map_err(|err| ExtractEntryError::Fatal(map_extract_archive_error(err)))?; + let key = ext.key().map_err(|err| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + format!("Failed to process archive PAX key: {}", err), + )) + })?; + let value = ext.value().map_err(|err| { + ExtractEntryError::Fatal(object_s3_error( + S3ErrorCode::InvalidArgument, + format!("Failed to process archive PAX value: {}", err), + )) + })?; if let Some(meta_key) = key.strip_prefix("minio.metadata.") { if !meta_key.is_empty() { - let name = http::HeaderName::from_bytes(meta_key.as_bytes()) - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball PAX metadata header"))?; - let header_value = HeaderValue::from_str(value) - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball PAX metadata value"))?; + let name = http::HeaderName::from_bytes(meta_key.as_bytes()).map_err(|_| { + ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball PAX metadata header")) + })?; + let header_value = HeaderValue::from_str(value).map_err(|_| { + ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball PAX metadata value")) + })?; preserve_unclassified_user_metadata(metadata, name.as_str(), value); pax_headers.insert(name, header_value); } @@ -246,7 +668,10 @@ where if key == "minio.versionId" && !value.is_empty() { if Uuid::parse_str(value).is_err() { - return Err(s3_error!(InvalidArgument, "Invalid Snowball PAX version ID")); + return Err(ExtractEntryError::Recoverable(s3_error!( + InvalidArgument, + "Invalid Snowball PAX version ID" + ))); } pax_version_id = Some(value.to_string()); } @@ -256,9 +681,12 @@ where if let Some(value) = pax_headers.get(AMZ_BUCKET_REPLICATION_STATUS) { let status = value .to_str() - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball replication status"))?; + .map_err(|_| ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball replication status")))?; if !status.eq_ignore_ascii_case(ReplicationStatusType::Replica.as_str()) { - return Err(s3_error!(InvalidArgument, "Invalid Snowball replication status")); + return Err(ExtractEntryError::Recoverable(s3_error!( + InvalidArgument, + "Invalid Snowball replication status" + ))); } pax_headers.insert(AMZ_BUCKET_REPLICATION_STATUS, HeaderValue::from_static("REPLICA")); } @@ -268,7 +696,7 @@ where if let Some(value) = pax_headers.remove("x-amz-tagging") { let value = value .to_str() - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball object tagging value"))?; + .map_err(|_| ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball object tagging value")))?; metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), value.to_owned()); } @@ -278,18 +706,18 @@ where value .to_str() .map(|value| ObjectLockMode::from(value.to_string())) - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball Object Lock mode")) + .map_err(|_| ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball Object Lock mode"))) }) .transpose()?; let object_lock_retain_until_date = pax_headers .remove(AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER) .map(|value| { - let value = value - .to_str() - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball Object Lock retain-until date"))?; - OffsetDateTime::parse(value, &Rfc3339) - .map(Timestamp::from) - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball Object Lock retain-until date")) + let value = value.to_str().map_err(|_| { + ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball Object Lock retain-until date")) + })?; + OffsetDateTime::parse(value, &Rfc3339).map(Timestamp::from).map_err(|_| { + ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball Object Lock retain-until date")) + }) }) .transpose()?; let object_lock_legal_hold_status = pax_headers @@ -298,7 +726,9 @@ where value .to_str() .map(|value| ObjectLockLegalHoldStatus::from(value.to_string())) - .map_err(|_| s3_error!(InvalidArgument, "Invalid Snowball Object Lock legal-hold status")) + .map_err(|_| { + ExtractEntryError::Recoverable(s3_error!(InvalidArgument, "Invalid Snowball Object Lock legal-hold status")) + }) }) .transpose()?; opts.version_id = pax_version_id; @@ -317,7 +747,9 @@ where object_lock_legal_hold_status.clone(), object_lock_mode.clone(), object_lock_retain_until_date.clone(), - )? { + ) + .map_err(ExtractEntryError::Recoverable)? + { metadata.extend(object_lock_metadata); } @@ -329,6 +761,30 @@ where }) } +#[cfg(test)] +async fn apply_extract_entry_pax_extensions_for_test( + entry: &mut tokio_tar::Entry>, + bucket: &str, + object_name: &str, + object_lock_config_state: &metadata_sys::ObjectLockConfigState, + metadata: &mut HashMap, + opts: &mut ObjectOptions, +) -> Result +where + R: AsyncRead + Send + Unpin + 'static, +{ + let mut total_pax_metadata_size = 0; + let mut total_pax_metadata_records = 0; + count_extract_entry_pax_metadata( + entry, + &mut total_pax_metadata_size, + &mut total_pax_metadata_records, + ArchiveLimits::default(), + ) + .await?; + apply_extract_entry_pax_extensions(entry, bucket, object_name, object_lock_config_state, metadata, opts).await +} + fn resolve_put_object_extract_options(headers: &HeaderMap) -> S3Result { let prefix = snowball_meta_value(headers, SNOWBALL_PREFIX_HEADER_KEYS, SNOWBALL_PREFIX_SUFFIX_LOWER) .map(|value| normalize_snowball_prefix(&value)) @@ -348,6 +804,23 @@ fn put_object_extract_limits() -> ArchiveLimits { ArchiveLimits::default() } +fn build_put_object_extract_archive(decoder: R, limits: ArchiveLimits) -> Archive +where + R: AsyncRead + Unpin, +{ + let max_physical_entries = u64::try_from(limits.max_entries) + .unwrap_or(u64::MAX) + .saturating_mul(EXTRACT_ARCHIVE_PHYSICAL_ENTRY_MULTIPLIER); + + tokio_tar::ArchiveBuilder::new(decoder) + .set_max_extension_entry_size(limits.max_pax_metadata_size) + .set_max_total_extension_size(limits.max_total_pax_metadata_size) + .set_max_physical_entries(max_physical_entries) + .set_max_sparse_entries(EXTRACT_ARCHIVE_MAX_SPARSE_ENTRIES) + .set_max_sparse_continuation_blocks(EXTRACT_ARCHIVE_MAX_SPARSE_CONTINUATION_BLOCKS) + .build() +} + fn validate_put_object_extract_entry_count(count: usize, limits: ArchiveLimits) -> S3Result<()> { if count > limits.max_entries { return Err(s3_error!( @@ -493,6 +966,14 @@ impl DefaultObjectUsecase { true, )?; let Some(body) = body else { return Err(s3_error!(IncompleteBody)) }; + let body = guard_put_object_body_read_timeout( + body, + &bucket, + &key, + &request_context.request_id, + content_length, + put_object_body_read_timeout(), + ); let size = match content_length { Some(c) => c, @@ -552,15 +1033,22 @@ impl DefaultObjectUsecase { return Err(ApiError::from(err).into()); } - let archive_etag = Arc::new(Mutex::new(None)); + let expected_archive_length = u64::try_from(size).map_err(|_| S3Error::new(S3ErrorCode::UnexpectedContent))?; + let archive_upload_state = Arc::new(Mutex::new(ExtractArchiveUploadState::default())); + let extract_limits = put_object_extract_limits(); let decoder = CompressionFormat::from_extension(&ext) - .get_decoder(ExtractArchiveEtagReader::new(archive_reader, archive_etag.clone())) + .get_decoder(ExtractArchiveEtagReader::new( + archive_reader, + expected_archive_length, + archive_upload_state.clone(), + )) .map_err(|e| { error!(error = ?e, "Archive decoder creation failed"); s3_error!(InvalidArgument, "get_decoder err") })?; + let decoder = ExtractDecodedLimitReader::new(decoder, extract_limits.max_decoded_size); - let mut ar = Archive::new(decoder); + let mut ar = build_put_object_extract_archive(decoder, extract_limits); let mut entries = ar.entries().map_err(|e| { error!(error = ?e, "Archive entry listing failed"); s3_error!(InvalidArgument, "get entries err") @@ -571,7 +1059,6 @@ impl DefaultObjectUsecase { }; let extract_options = resolve_put_object_extract_options(&req.headers)?; - let extract_limits = put_object_extract_limits(); let extract_quota_check = if let Some(metadata_sys) = self.bucket_metadata_sys() { let quota_checker = QuotaChecker::new(metadata_sys); let check_result = @@ -595,7 +1082,10 @@ impl DefaultObjectUsecase { let user_agent = get_request_user_agent(&req.headers); let mut wrote_any_entry = false; let mut extracted_entry_count = 0usize; - let mut total_unpacked_size = 0u64; + let mut resource_total_size = 0u64; + let mut legacy_quota_growth = 0u64; + let mut total_pax_metadata_size = 0u64; + let mut total_pax_metadata_records = 0usize; let object_lock_config_snapshot = store.object_lock_config_snapshot(&bucket).await.map_err(ApiError::from)?; let object_lock_config_state = object_lock_config_snapshot.state(); @@ -603,40 +1093,54 @@ impl DefaultObjectUsecase { let mut f = match entry { Ok(f) => f, Err(e) => { - if extract_options.ignore_errors { - warn!(error = %e, "Archive entry read skipped due to ignore-errors"); - continue; - } error!(error = %e, "Archive entry read failed"); return Err(s3_error!(InvalidArgument, "Failed to read archive entry: {:?}", e)); } }; extracted_entry_count = extracted_entry_count.saturating_add(1); validate_put_object_extract_entry_count(extracted_entry_count, extract_limits)?; + let entry_size = f.effective_size(); + validate_put_object_extract_entry_size("archive member", entry_size, extract_limits)?; + resource_total_size = resource_total_size + .checked_add(entry_size) + .ok_or_else(|| s3_error!(InvalidArgument, "Archive total unpacked size overflowed while processing entries"))?; + validate_put_object_extract_total_size(resource_total_size, extract_limits)?; + count_extract_entry_pax_metadata( + &mut f, + &mut total_pax_metadata_size, + &mut total_pax_metadata_records, + extract_limits, + ) + .await + .map_err(ExtractEntryError::into_s3_error)?; - let fpath = match f.path() { - Ok(path) => path, - Err(e) => { - if extract_options.ignore_errors { - warn!(error = %e, "Archive path decode skipped due to ignore-errors"); + let entry_type = classify_extract_entry_type(f.header().entry_type()); + if entry_type == ExtractEntryDisposition::FormatSkip { + continue; + } + let is_dir = entry_type == ExtractEntryDisposition::Directory; + if is_dir && extract_options.ignore_dirs { + continue; + } + let fpath = { + let path_bytes = f.path_bytes().map_err(map_extract_archive_error)?; + let path = match strict_extract_entry_path(path_bytes.as_ref()) { + Ok(path) => path, + Err(err) => { + err.ignore_or_return(extract_options.ignore_errors)?; continue; } - return Err(s3_error!(InvalidArgument, "Failed to decode archive entry path")); + }; + if is_empty_extract_entry_path(path) { + continue; } + normalize_extract_entry_key(path, extract_options.prefix.as_deref(), is_dir)? }; - let is_dir = f.header().entry_type().is_dir(); - let fpath = match normalize_extract_entry_key(&fpath.to_string_lossy(), extract_options.prefix.as_deref(), is_dir) { - Ok(fpath) => fpath, - Err(err) => { - if extract_options.ignore_errors { - warn!(error = %err, "Unsafe archive path skipped due to ignore-errors"); - continue; - } - return Err(err); - } - }; - validate_put_object_extract_entry_path(&fpath, extract_limits)?; + if let Err(err) = validate_extract_member_key(&fpath, extract_limits) { + err.ignore_or_return(extract_options.ignore_errors)?; + continue; + } validate_table_catalog_object_mutation(&bucket, &fpath).await?; let mut auth_req = S3Request { @@ -656,26 +1160,12 @@ impl DefaultObjectUsecase { req_info.object = Some(fpath.clone()); req_info.version_id = None; } - let entry_size = f.effective_size(); - validate_put_object_extract_entry_size(&fpath, entry_size, extract_limits)?; - total_unpacked_size = total_unpacked_size - .checked_add(entry_size) - .ok_or_else(|| s3_error!(InvalidArgument, "Archive total unpacked size overflowed while processing entries"))?; - validate_put_object_extract_total_size(total_unpacked_size, extract_limits)?; - if let Some(quota_check) = extract_quota_check.as_ref() { - ensure_legacy_archive_size_within_quota(quota_check, total_unpacked_size)?; - } let mut size = i64::try_from(entry_size).map_err(|_| s3_error!(InvalidArgument, "Archive entry size does not fit into i64"))?; - // mtime 0 means "unset" in tar headers, and xl.meta cannot represent an - // epoch mod_time anyway (0 nanos decodes as no-mod_time, making the version - // unreadable — rustfs#4842), so fall back to the upload time instead. - let archive_entry_mod_time = f - .header() - .mtime() - .ok() - .filter(|&modified_at_secs| modified_at_secs > 0) - .and_then(|modified_at_secs| OffsetDateTime::from_unix_timestamp(modified_at_secs as i64).ok()); + // mtime 0 or a negative GNU base-256 value means "unset". xl.meta + // also cannot represent the Unix epoch as an object mod_time, so + // those cases fall back to the upload time (rustfs#4842). + let archive_entry_mod_time = extract_archive_entry_mod_time(f.header())?; let mut metadata = HashMap::new(); let has_explicit_object_lock_retention = object_lock_mode.is_some() || object_lock_retain_until_date.is_some(); apply_put_request_metadata( @@ -713,9 +1203,31 @@ impl DefaultObjectUsecase { } opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id; opts.object_lock_config_snapshot = Some(Arc::clone(&object_lock_config_snapshot)); - let pax_authorization = - apply_extract_entry_pax_extensions(&mut f, &bucket, &fpath, object_lock_config_state, &mut metadata, &mut opts) - .await?; + let pax_authorization = match apply_extract_entry_pax_extensions( + &mut f, + &bucket, + &fpath, + object_lock_config_state, + &mut metadata, + &mut opts, + ) + .await + { + Ok(authorization) => authorization, + Err(err) => { + err.ignore_or_return(extract_options.ignore_errors)?; + continue; + } + }; + if let Some(quota_check) = extract_quota_check.as_ref() { + let next_legacy_quota_growth = legacy_quota_growth + .checked_add(extract_entry_quota_growth(entry_type, entry_size)) + .ok_or_else(|| { + object_s3_error(S3ErrorCode::InvalidArgument, "Archive quota growth overflowed while processing entries") + })?; + ensure_legacy_archive_size_within_quota(quota_check, next_legacy_quota_growth)?; + legacy_quota_growth = next_legacy_quota_growth; + } for (name, value) in &pax_authorization.headers { auth_req.headers.insert(name.clone(), value.clone()); } @@ -752,10 +1264,6 @@ impl DefaultObjectUsecase { debug!("Extracting file: {}, size: {} bytes", fpath, size); if is_dir { - if extract_options.ignore_dirs { - debug!("Skipping directory entry during archive extract: {}", fpath); - continue; - } size = 0; } @@ -811,6 +1319,7 @@ impl DefaultObjectUsecase { opts.user_defined.extend(encryption_metadata); } hrd = write_plan.apply(hrd, actual_size).map_err(ApiError::from)?; + let (hrd, member_read_failed) = track_extract_member_read_errors(hrd).map_err(ApiError::from)?; opts.user_defined.extend(metadata); // Each extracted member is an independent user write and joins @@ -847,27 +1356,32 @@ impl DefaultObjectUsecase { { Ok(result) => result, Err(e) => { - if extract_options.ignore_errors { + if should_ignore_extract_member_write_error(extract_options.ignore_errors, &member_read_failed) { warn!(error = %e, "Archive object write skipped due to ignore-errors"); continue; } return Err(ApiError::from(e).into()); } }; - let committed_size = quota_accounting_object_size(&obj_info, extract_quota_enabled)?; let extract_versioned = BucketVersioningSys::prefix_enabled(&bucket, &fpath).await; - match previous_current_size_from_backfill(backfilled_old_current_size) { - Some(previous_current_size) => { - if extract_versioned { - record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await; - } else { - record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await; + let post_commit_error = match quota_accounting_object_size(&obj_info, extract_quota_enabled) { + Ok(committed_size) => { + match previous_current_size_from_backfill(backfilled_old_current_size) { + Some(previous_current_size) => { + if extract_versioned { + record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await; + } else { + record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await; + } + } + None => { + record_bucket_object_write_unknown_previous_memory(&bucket, committed_size, extract_versioned).await; + } } + None } - None => { - record_bucket_object_write_unknown_previous_memory(&bucket, committed_size, extract_versioned).await; - } - } + Err(err) => Some(err), + }; let _ = invalidate_object_data_cache_after_put_success(&cache_adapter, &bucket, &fpath).await; // Reuse the per-entry pre-commit decision (see `dsc` above) so the @@ -907,6 +1421,10 @@ impl DefaultObjectUsecase { spawn_background_with_context(Some(request_context.clone()), async move { notify.notify(event_args).await; }); + + if let Some(err) = post_commit_error { + return Err(err); + } } let mut checksums = PutObjectChecksums { @@ -916,12 +1434,6 @@ impl DefaultObjectUsecase { sha256: input.checksum_sha256, crc64nvme: input.checksum_crc64nvme, }; - apply_trailing_checksums( - input.checksum_algorithm.as_ref().map(|a| a.as_str()), - &req.trailing_headers, - &mut checksums, - ); - warn!( "put object extract checksum_crc32={:?}, checksum_crc32c={:?}, checksum_sha1={:?}, checksum_sha256={:?}, checksum_crc64nvme={:?}", checksums.crc32, checksums.crc32c, checksums.sha1, checksums.sha256, checksums.crc64nvme, @@ -935,11 +1447,23 @@ impl DefaultObjectUsecase { tokio::io::copy(&mut decoder, &mut tokio::io::sink()) .await .map_err(map_extract_archive_error)?; - let archive_etag = archive_etag - .lock() - .ok() - .and_then(|etag| etag.clone()) - .map(|etag| to_s3s_etag(&etag)); + let archive_etag = { + let state = archive_upload_state + .lock() + .map_err(|_| object_s3_error(S3ErrorCode::InternalError, "Archive upload state lock was poisoned"))?; + if !state.body_complete { + return Err(object_s3_error( + S3ErrorCode::UnexpectedContent, + "Archive decoder did not consume the complete request body", + )); + } + state.etag.as_ref().map(|etag| to_s3s_etag(etag)) + }; + apply_trailing_checksums( + input.checksum_algorithm.as_ref().map(|a| a.as_str()), + &req.trailing_headers, + &mut checksums, + ); let output = PutObjectOutput { e_tag: archive_etag, @@ -961,6 +1485,7 @@ mod tests { use super::*; use http::{HeaderMap, HeaderName, HeaderValue}; use s3s::dto::{ObjectLockConfiguration, ObjectLockEnabled}; + use tokio::io::AsyncReadExt; use tokio_tar::{Builder, EntryType, Header}; fn pax_record(key: &str, value: &[u8]) -> Vec { @@ -981,6 +1506,134 @@ mod tests { record } + async fn entry_with_local_pax(record: &[u8], entry_type: EntryType) -> tokio_tar::Entry>>> { + let mut builder = Builder::new(std::io::Cursor::new(Vec::new())); + let mut extension = Header::new_ustar(); + extension.set_size(record.len() as u64); + extension.set_entry_type(EntryType::XHeader); + builder + .append_data(&mut extension, "pax", record) + .await + .expect("local PAX fixture should be appended"); + + let mut member = Header::new_ustar(); + member.set_size(0); + member.set_entry_type(entry_type); + if entry_type == EntryType::Symlink { + member.set_link_name("target").expect("symlink fixture should have a target"); + } + builder + .append_data(&mut member, "member", std::io::Cursor::new(Vec::new())) + .await + .expect("member fixture should be appended"); + + let bytes = builder + .into_inner() + .await + .expect("fixture builder should finish") + .into_inner(); + let mut archive = Archive::new(std::io::Cursor::new(bytes)); + let mut entries = archive.entries().expect("fixture archive should be iterable"); + entries + .next() + .await + .expect("fixture should contain a logical member") + .expect("fixture member should parse") + } + + #[tokio::test] + async fn archive_etag_reader_validates_raw_eof_before_returning_final_byte() { + let payload = b"decoder-consumed-exact-body".to_vec(); + let state = Arc::new(Mutex::new(ExtractArchiveUploadState::default())); + let mut reader = ExtractArchiveEtagReader::new( + std::io::Cursor::new(payload.clone()), + u64::try_from(payload.len()).expect("fixture length must fit u64"), + state.clone(), + ); + let mut output = Vec::new(); + + reader + .read_to_end(&mut output) + .await + .expect("exact body should validate through raw EOF"); + + assert_eq!(output, payload); + let expected_etag = hex_simd::encode_to_string(Md5::digest(&payload), hex_simd::AsciiCase::Lower); + let state = state.lock().expect("archive state lock must remain healthy"); + assert!(state.body_complete); + assert_eq!(state.etag.as_deref(), Some(expected_etag.as_str())); + } + + #[tokio::test] + async fn archive_etag_reader_rejects_short_and_overlong_bodies() { + let payload = b"body-length-fixture".to_vec(); + + let short_state = Arc::new(Mutex::new(ExtractArchiveUploadState::default())); + let mut short = ExtractArchiveEtagReader::new( + std::io::Cursor::new(payload.clone()), + u64::try_from(payload.len() + 1).expect("fixture length must fit u64"), + short_state.clone(), + ); + let short_err = short + .read_to_end(&mut Vec::new()) + .await + .expect_err("short body must be rejected"); + assert_eq!(short_err.kind(), std::io::ErrorKind::UnexpectedEof); + assert!( + !short_state + .lock() + .expect("archive state lock must remain healthy") + .body_complete + ); + + let overlong_state = Arc::new(Mutex::new(ExtractArchiveUploadState::default())); + let mut overlong = ExtractArchiveEtagReader::new( + std::io::Cursor::new(payload.clone()), + u64::try_from(payload.len() - 1).expect("fixture length must fit u64"), + overlong_state.clone(), + ); + let overlong_err = overlong + .read_to_end(&mut Vec::new()) + .await + .expect_err("overlong body must be rejected"); + assert_eq!(overlong_err.kind(), std::io::ErrorKind::InvalidData); + assert!( + !overlong_state + .lock() + .expect("archive state lock must remain healthy") + .body_complete + ); + } + + #[tokio::test] + async fn archive_etag_reader_validates_content_md5_before_completion() { + let payload = b"archive-with-wrong-content-md5".to_vec(); + let expected_length = i64::try_from(payload.len()).expect("fixture length must fit i64"); + let hash_reader = HashReader::from_stream( + std::io::Cursor::new(payload), + expected_length, + expected_length, + Some("00000000000000000000000000000000".to_string()), + None, + false, + ) + .expect("hash reader should be created"); + let state = Arc::new(Mutex::new(ExtractArchiveUploadState::default())); + let mut reader = ExtractArchiveEtagReader::new( + hash_reader, + u64::try_from(expected_length).expect("fixture length must fit u64"), + state.clone(), + ); + + let err = reader + .read_to_end(&mut Vec::new()) + .await + .expect_err("Content-MD5 must be checked before upload completion"); + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert!(!state.lock().expect("archive state lock must remain healthy").body_complete); + } + #[tokio::test] async fn snowball_pax_rejects_unpaired_object_lock_retention() { let record = pax_record("minio.metadata.x-amz-object-lock-mode", b"GOVERNANCE"); @@ -1009,9 +1662,10 @@ mod tests { updated_at: OffsetDateTime::now_utc(), }; - let err = apply_extract_entry_pax_extensions(&mut entry, "bucket", "object", &state, &mut metadata, &mut opts) + let err = apply_extract_entry_pax_extensions_for_test(&mut entry, "bucket", "object", &state, &mut metadata, &mut opts) .await - .unwrap_err(); + .unwrap_err() + .into_s3_error(); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert_eq!(metadata.get(AMZ_OBJECT_LOCK_MODE_LOWER).map(String::as_str), Some("COMPLIANCE")); @@ -1061,10 +1715,16 @@ mod tests { let mut entry = entries.next().await.unwrap().unwrap(); let mut opts = ObjectOptions::default(); - let authorization = - apply_extract_entry_pax_extensions(&mut entry, "bucket", "object", &state, &mut HashMap::new(), &mut opts) - .await - .unwrap(); + let authorization = apply_extract_entry_pax_extensions_for_test( + &mut entry, + "bucket", + "object", + &state, + &mut HashMap::new(), + &mut opts, + ) + .await + .unwrap(); assert_eq!( ( @@ -1129,8 +1789,7 @@ mod tests { let mut archive = Archive::new(std::io::Cursor::new(builder.into_inner().await.unwrap())); let mut entries = archive.entries().unwrap(); let mut entry = entries.next().await.unwrap().unwrap(); - - let err = apply_extract_entry_pax_extensions( + let err = apply_extract_entry_pax_extensions_for_test( &mut entry, "bucket", "object", @@ -1139,7 +1798,8 @@ mod tests { &mut ObjectOptions::default(), ) .await - .unwrap_err(); + .unwrap_err() + .into_s3_error(); assert!( err.code() == &S3ErrorCode::InvalidArgument || err.code() == &S3ErrorCode::MalformedXML, @@ -1185,7 +1845,7 @@ mod tests { }; let authorization = - apply_extract_entry_pax_extensions(&mut entry, "bucket", "object.txt", &state, &mut metadata, &mut opts) + apply_extract_entry_pax_extensions_for_test(&mut entry, "bucket", "object.txt", &state, &mut metadata, &mut opts) .await .unwrap(); @@ -1227,8 +1887,7 @@ mod tests { let mut entries = archive.entries().unwrap(); let mut entry = entries.next().await.unwrap().unwrap(); let mut metadata = HashMap::new(); - - let err = apply_extract_entry_pax_extensions( + let err = apply_extract_entry_pax_extensions_for_test( &mut entry, "bucket", "object", @@ -1237,7 +1896,8 @@ mod tests { &mut ObjectOptions::default(), ) .await - .unwrap_err(); + .unwrap_err() + .into_s3_error(); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert!(!metadata.contains_key(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER)); @@ -1266,8 +1926,7 @@ mod tests { }, updated_at: OffsetDateTime::now_utc(), }; - - let err = apply_extract_entry_pax_extensions( + let err = apply_extract_entry_pax_extensions_for_test( &mut entry, "bucket", "object", @@ -1276,7 +1935,8 @@ mod tests { &mut ObjectOptions::default(), ) .await - .unwrap_err(); + .unwrap_err() + .into_s3_error(); assert_eq!(err.code(), &S3ErrorCode::MalformedXML); assert!(!metadata.contains_key(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER)); @@ -1483,6 +2143,281 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } + #[test] + fn extract_entry_error_boundary_only_ignores_recoverable_members() { + let recoverable = ExtractEntryError::Recoverable(object_s3_error(S3ErrorCode::InvalidArgument, "invalid member")); + recoverable + .ignore_or_return(true) + .expect("recoverable member should be skipped"); + + let recoverable = ExtractEntryError::Recoverable(object_s3_error(S3ErrorCode::InvalidArgument, "invalid member")); + assert_eq!( + recoverable + .ignore_or_return(false) + .expect_err("recoverable member should fail without ignore-errors") + .code(), + &S3ErrorCode::InvalidArgument + ); + + let fatal = ExtractEntryError::Fatal(object_s3_error(S3ErrorCode::InvalidArgument, "invalid archive")); + assert_eq!( + fatal + .ignore_or_return(true) + .expect_err("fatal archive errors must ignore ignore-errors") + .code(), + &S3ErrorCode::InvalidArgument + ); + } + + #[test] + fn strict_extract_entry_path_rejects_non_utf8_without_lossy_replacement() { + let invalid = strict_extract_entry_path(b"same-\xff-key").expect_err("non-UTF-8 member path must be rejected"); + assert!(invalid.is_recoverable()); + invalid + .ignore_or_return(true) + .expect("ignore-errors should skip an invalid member key"); + } + + #[test] + fn classify_extract_entry_type_skips_links_extensions_and_continuous_entries() { + assert_eq!(classify_extract_entry_type(EntryType::Regular), ExtractEntryDisposition::File); + assert_eq!(classify_extract_entry_type(EntryType::Directory), ExtractEntryDisposition::Directory); + for entry_type in [ + EntryType::Link, + EntryType::Symlink, + EntryType::Continuous, + EntryType::XGlobalHeader, + EntryType::Other(b'V'), + ] { + assert_eq!( + classify_extract_entry_type(entry_type), + ExtractEntryDisposition::FormatSkip, + "{entry_type:?} must not be materialized as an object" + ); + } + } + + #[test] + fn extract_entry_quota_growth_counts_only_materialized_files() { + assert_eq!(extract_entry_quota_growth(ExtractEntryDisposition::File, 9), 9); + assert_eq!(extract_entry_quota_growth(ExtractEntryDisposition::Directory, 9), 0); + assert_eq!(extract_entry_quota_growth(ExtractEntryDisposition::FormatSkip, 9), 0); + } + + #[test] + fn archive_mod_time_treats_negative_gnu_base256_as_unset() { + let mut header = Header::new_gnu(); + header.as_old_mut().mtime.fill(0xff); + + assert_eq!( + extract_archive_entry_mod_time(&header).expect("negative GNU mtime should be accepted"), + None + ); + } + + #[test] + fn archive_mod_time_keeps_malformed_octal_fatal() { + let mut header = Header::new_ustar(); + header.as_old_mut().mtime.fill(b'9'); + + let err = extract_archive_entry_mod_time(&header).expect_err("malformed octal mtime must be rejected"); + assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); + } + + #[test] + fn pax_metadata_budget_is_fatal_even_with_ignore_errors() { + let limits = ArchiveLimits { + max_pax_metadata_size: 3, + ..ArchiveLimits::default() + }; + let mut entry_size = 0; + let mut total_size = 0; + let err = record_extract_pax_metadata_bytes(&mut entry_size, &mut total_size, 2, 2, limits) + .expect_err("PAX metadata over the resource budget must fail"); + assert!(!err.is_recoverable()); + assert!(err.ignore_or_return(true).is_err(), "ignore-errors must not bypass resource limits"); + } + + #[test] + fn pax_metadata_record_budget_is_fatal_even_with_ignore_errors() { + let limits = ArchiveLimits { + max_pax_metadata_records: 1, + ..ArchiveLimits::default() + }; + let mut entry_records = 0; + let mut total_records = 0; + + record_extract_pax_metadata_record(&mut entry_records, &mut total_records, limits) + .expect("first PAX record should fit the budget"); + let err = record_extract_pax_metadata_record(&mut entry_records, &mut total_records, limits) + .expect_err("second PAX record must exceed the per-entry budget"); + + assert!(!err.is_recoverable()); + assert!(err.ignore_or_return(true).is_err(), "ignore-errors must not bypass record limits"); + + let limits = ArchiveLimits { + max_pax_metadata_records: 2, + max_total_pax_metadata_records: 1, + ..ArchiveLimits::default() + }; + let mut first_entry_records = 0; + let mut second_entry_records = 0; + let mut total_records = 0; + record_extract_pax_metadata_record(&mut first_entry_records, &mut total_records, limits) + .expect("first archive PAX record should fit the total budget"); + let err = record_extract_pax_metadata_record(&mut second_entry_records, &mut total_records, limits) + .expect_err("second archive PAX record must exceed the total budget"); + assert!(!err.is_recoverable()); + } + + #[tokio::test] + async fn extract_archive_builder_rejects_oversized_extension_payload() { + let record = pax_record("comment", b"value"); + let mut builder = Builder::new(Vec::new()); + let mut extension = Header::new_ustar(); + extension.set_entry_type(EntryType::XHeader); + extension.set_size(u64::try_from(record.len()).expect("fixture size must fit u64")); + extension.set_cksum(); + builder + .append_data(&mut extension, "pax", record.as_slice()) + .await + .expect("PAX extension fixture should be appended"); + let bytes = builder.into_inner().await.expect("PAX extension fixture should finalize"); + let limits = ArchiveLimits { + max_pax_metadata_size: u64::try_from(record.len() - 1).expect("fixture size must fit u64"), + ..ArchiveLimits::default() + }; + + let mut archive = build_put_object_extract_archive(std::io::Cursor::new(bytes), limits); + let err = archive + .entries() + .expect("archive entry stream should be created") + .next() + .await + .expect("extension header should produce a result") + .expect_err("dependency must reject the extension before buffering its payload"); + + assert_eq!(err.to_string(), "archive extension entry size limit exceeded"); + } + + #[tokio::test] + async fn pax_metadata_budget_counts_format_skipped_members() { + let record = pax_record("comment", b"oversized-symlink-metadata"); + let mut entry = entry_with_local_pax(&record, EntryType::Symlink).await; + let limits = ArchiveLimits { + max_pax_metadata_size: 8, + ..ArchiveLimits::default() + }; + let mut total_size = 0; + let mut total_records = 0; + + let err = count_extract_entry_pax_metadata(&mut entry, &mut total_size, &mut total_records, limits) + .await + .expect_err("format-skipped members must still consume the PAX budget"); + + assert!(!err.is_recoverable()); + assert!(err.ignore_or_return(true).is_err()); + } + + #[tokio::test] + async fn pax_metadata_budget_precedes_recoverable_semantic_errors() { + let invalid_key = "minio.metadata.x-amz-meta-owner"; + let mut record = pax_record(invalid_key, b"\0"); + record.extend(pax_record("comment", b"oversized-tail")); + let mut entry = entry_with_local_pax(&record, EntryType::Regular).await; + + let semantic_err = apply_extract_entry_pax_extensions( + &mut entry, + "bucket", + "member", + &metadata_sys::ObjectLockConfigState::ConfirmedAbsent, + &mut HashMap::new(), + &mut ObjectOptions::default(), + ) + .await + .expect_err("NUL metadata value should be a recoverable member error"); + assert!(semantic_err.is_recoverable()); + + let limits = ArchiveLimits { + max_pax_metadata_size: u64::try_from(invalid_key.len() + 1).expect("fixture size must fit u64"), + ..ArchiveLimits::default() + }; + let mut total_size = 0; + let mut total_records = 0; + let budget_err = count_extract_entry_pax_metadata(&mut entry, &mut total_size, &mut total_records, limits) + .await + .expect_err("the oversized tail must be counted before ignore-errors can skip the member"); + + assert!(!budget_err.is_recoverable()); + assert!(budget_err.ignore_or_return(true).is_err()); + } + + #[tokio::test] + async fn extract_decoded_reader_enforces_exact_byte_limit() { + let mut exact = ExtractDecodedLimitReader::new(std::io::Cursor::new(b"1234"), 4); + let mut exact_bytes = Vec::new(); + exact + .read_to_end(&mut exact_bytes) + .await + .expect("decoded stream at the limit should succeed"); + assert_eq!(exact_bytes, b"1234"); + + let mut oversized = ExtractDecodedLimitReader::new(std::io::Cursor::new(b"12345"), 4); + let mut oversized_bytes = Vec::new(); + let err = oversized + .read_to_end(&mut oversized_bytes) + .await + .expect_err("decoded stream over the limit must fail"); + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + } + + #[tokio::test] + async fn extract_member_read_tracker_keeps_integrity_errors_fatal_under_ignore_errors() { + struct FailingReader; + + impl AsyncRead for FailingReader { + fn poll_read(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &mut ReadBuf<'_>) -> Poll> { + Poll::Ready(Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "member decoder failed"))) + } + } + + let reader = HashReader::from_stream(FailingReader, 1, 1, None, None, false).unwrap(); + let (mut tracked, failed) = track_extract_member_read_errors(reader).unwrap(); + let mut output = Vec::new(); + let err = tracked.read_to_end(&mut output).await.unwrap_err(); + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert!(failed.load(Ordering::Acquire)); + assert!(!should_ignore_extract_member_write_error(true, &failed)); + + let storage_only_failure = AtomicBool::new(false); + assert!(should_ignore_extract_member_write_error(true, &storage_only_failure)); + assert!(!should_ignore_extract_member_write_error(false, &storage_only_failure)); + } + + #[tokio::test] + async fn snowball_extract_body_guard_aborts_stalled_upload() { + let body = StreamingBlob::wrap(futures::stream::pending::>()); + let mut guarded = guard_put_object_body_read_timeout( + body, + "test-bucket", + "archive.tar", + "snowball-timeout", + Some(512), + Duration::from_millis(1), + ); + + let err = guarded + .next() + .await + .expect("stalled Snowball body should yield an error") + .expect_err("stalled Snowball body must not hang"); + let io_err = err + .downcast_ref::() + .expect("stall error should retain its I/O kind"); + assert_eq!(io_err.kind(), std::io::ErrorKind::TimedOut); + } + #[test] fn legacy_archive_quota_rejects_cumulative_size_and_overflow() { let legacy = QuotaCheckResult { diff --git a/rustfs/src/app/object/mod.rs b/rustfs/src/app/object/mod.rs index eaf083ee7..d28777a99 100644 --- a/rustfs/src/app/object/mod.rs +++ b/rustfs/src/app/object/mod.rs @@ -178,6 +178,10 @@ use s3s::header::{X_AMZ_RESTORE, X_AMZ_RESTORE_OUTPUT_PATH}; use s3s::stream::{ByteStream, DynByteStream, RemainingLength}; use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; +fn object_s3_error(code: S3ErrorCode, message: impl Into>) -> S3Error { + S3Error::with_message(code, message) +} + mod copy; mod delete; mod extract;