From 6f6d8a4d3e2e13c25bdab3fe323d0e0d55258abf Mon Sep 17 00:00:00 2001 From: cxymds Date: Fri, 24 Jul 2026 13:47:25 +0800 Subject: [PATCH] fix(s3): persist multipart upload storage class (#5170) Fixes rustfs/backlog#1464. --- .config/nextest.toml | 2 +- crates/e2e_test/src/lib.rs | 3 + .../src/multipart_storage_class_test.rs | 364 ++++++++++++++++++ docs/testing/e2e-suite-inventory.md | 3 +- rustfs/src/app/multipart_usecase.rs | 71 +++- 5 files changed, 430 insertions(+), 13 deletions(-) create mode 100644 crates/e2e_test/src/multipart_storage_class_test.rs diff --git a/.config/nextest.toml b/.config/nextest.toml index 84e6107ca..f6423de69 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -192,7 +192,7 @@ test-group = 'ecstore-serial-flaky' [profile.e2e-smoke] default-filter = """ package(e2e_test) & ( - test(/^(delete_marker_migration_semantics|version_id_regression|list_objects_v2_pagination|list_object_versions_regression|list_objects_duplicates|list_buckets_double_slash|leading_slash_key|special_chars|create_bucket_region|delete_objects_versioning|head_object_consistency|head_object_range|copy_object_metadata|copy_object_tagging|copy_source_invalid_date|content_encoding|anonymous_access|bucket_policy_check|presigned_negative|negative_sigv4|admin_auth|notification_webhook|tls_hot_reload|console_smoke|admin_iam_crud|admin_pools)_test::|^fake_s3_target::/) + test(/^(delete_marker_migration_semantics|version_id_regression|list_objects_v2_pagination|list_object_versions_regression|list_objects_duplicates|list_buckets_double_slash|leading_slash_key|special_chars|create_bucket_region|delete_objects_versioning|head_object_consistency|head_object_range|copy_object_metadata|copy_object_tagging|copy_source_invalid_date|content_encoding|multipart_storage_class|anonymous_access|bucket_policy_check|presigned_negative|negative_sigv4|admin_auth|notification_webhook|tls_hot_reload|console_smoke|admin_iam_crud|admin_pools)_test::|^fake_s3_target::/) | test(/^replication_extension_test::(test_replication_check_succeeds_with_remote_target|test_replication_check_rejects_target_without_object_lock|test_set_remote_target_rejects_unversioned_source_bucket|test_replication_check_rejects_unversioned_source_bucket|test_replication_check_rejects_missing_replication_config|test_replication_check_rejects_invalid_bucket|test_set_remote_target_rejects_same_bucket_on_same_deployment|test_set_remote_target_rejects_unversioned_target_bucket|test_set_remote_target_update_requires_arn|test_set_remote_target_update_rejects_missing_target|test_set_remote_target_rejects_invalid_target_url|test_set_remote_target_rejects_self_signed_https_target_without_skip_tls_verify|test_set_remote_target_rejects_private_ca_https_target_without_ca_cert_pem|test_list_remote_targets_rejects_empty_bucket|test_list_remote_targets_rejects_invalid_bucket|test_remove_remote_target_rejects_missing_target|test_remove_remote_target_rejects_missing_arn|test_remove_remote_target_rejects_invalid_bucket|test_remove_remote_target_rejects_target_used_by_replication|test_delete_bucket_replication_removes_remote_target)$/) | test(/^reliant::lifecycle::/) | test(/^reliant::tiering::/) diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 5e8c9043d..b727a8560 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -196,6 +196,9 @@ mod copy_object_version_restore_test; #[cfg(test)] mod copy_object_checksum_test; +#[cfg(test)] +mod multipart_storage_class_test; + // S3 dummy-compat bucket API tests #[cfg(test)] mod bucket_logging_test; diff --git a/crates/e2e_test/src/multipart_storage_class_test.rs b/crates/e2e_test/src/multipart_storage_class_test.rs new file mode 100644 index 000000000..21a5e52d0 --- /dev/null +++ b/crates/e2e_test/src/multipart_storage_class_test.rs @@ -0,0 +1,364 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! CreateMultipartUpload storage-class persistence regression tests. + +#[cfg(test)] +mod tests { + use crate::common::{RustFSTestEnvironment, init_logging}; + use aws_sdk_s3::Client; + use aws_sdk_s3::error::ProvideErrorMetadata; + use aws_sdk_s3::primitives::ByteStream; + use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart, StorageClass}; + + const PART_SIZE: usize = 5 * 1024 * 1024; + + async fn assert_completed_object( + client: &Client, + bucket: &str, + key: &str, + expected_storage_class: &str, + expected_body: &[u8], + ) -> Result<(), Box> { + let head = client.head_object().bucket(bucket).key(key).send().await?; + let expected_head_class = (expected_storage_class != "STANDARD").then_some(expected_storage_class); + assert_eq!( + head.storage_class().map(StorageClass::as_str), + expected_head_class, + "HeadObject should use S3's implicit STANDARD representation" + ); + + let listed = client.list_objects_v2().bucket(bucket).prefix(key).send().await?; + let object = listed + .contents() + .iter() + .find(|object| object.key() == Some(key)) + .ok_or("completed multipart object missing from ListObjectsV2")?; + assert_eq!( + object.storage_class().map(|storage_class| storage_class.as_str()), + Some(expected_storage_class), + "ListObjectsV2 should report the completed object's storage class" + ); + + let body = client + .get_object() + .bucket(bucket) + .key(key) + .send() + .await? + .body + .collect() + .await? + .into_bytes(); + assert_eq!(body.as_ref(), expected_body, "completed multipart body should be byte-exact"); + Ok(()) + } + + #[tokio::test] + async fn multipart_upload_preserves_standard_and_rrs_across_retry_and_resume() + -> Result<(), Box> { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(Vec::new()).await?; + let client = env.create_s3_client(); + let bucket = "multipart-storage-class-retry"; + env.create_test_bucket(bucket).await?; + + for storage_class in [StorageClass::Standard, StorageClass::ReducedRedundancy] { + let class_name = storage_class.as_str(); + let key = format!("retry-{class_name}.bin"); + let create = client + .create_multipart_upload() + .bucket(bucket) + .key(&key) + .storage_class(storage_class.clone()) + .content_type("application/octet-stream") + .metadata("content-type", "user-content-type") + .metadata("x-amz-storage-class", "user-storage-class") + .send() + .await?; + let upload_id = create.upload_id().ok_or("CreateMultipartUpload returned no upload ID")?; + + let original_part = vec![b'a'; PART_SIZE]; + let first_attempt = client + .upload_part() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(original_part)) + .send() + .await?; + + let resumed_client = env.create_s3_client(); + let before_retry = resumed_client + .list_parts() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .send() + .await?; + assert_eq!(before_retry.storage_class().map(StorageClass::as_str), Some(class_name)); + assert_eq!(before_retry.parts().len(), 1, "resume should find the previously uploaded part"); + + let retried_part = vec![b'b'; PART_SIZE]; + let retry = resumed_client + .upload_part() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(retried_part.clone())) + .send() + .await?; + assert_ne!( + first_attempt.e_tag(), + retry.e_tag(), + "retrying the same part number with different bytes should replace the part" + ); + + let tail = format!("-tail-{class_name}").into_bytes(); + let second = resumed_client + .upload_part() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .part_number(2) + .body(ByteStream::from(tail.clone())) + .send() + .await?; + let after_retry = resumed_client + .list_parts() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .send() + .await?; + assert_eq!(after_retry.storage_class().map(StorageClass::as_str), Some(class_name)); + assert_eq!(after_retry.parts().len(), 2); + assert_eq!(after_retry.parts()[0].e_tag(), retry.e_tag()); + + resumed_client + .complete_multipart_upload() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .multipart_upload( + CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .set_e_tag(retry.e_tag().map(str::to_owned)) + .build(), + ) + .parts( + CompletedPart::builder() + .part_number(2) + .set_e_tag(second.e_tag().map(str::to_owned)) + .build(), + ) + .build(), + ) + .send() + .await?; + + let mut expected_body = retried_part; + expected_body.extend_from_slice(&tail); + assert_completed_object(&resumed_client, bucket, &key, class_name, &expected_body).await?; + let metadata_head = resumed_client.head_object().bucket(bucket).key(&key).send().await?; + assert_eq!(metadata_head.content_type(), Some("application/octet-stream")); + assert_eq!( + metadata_head.metadata().and_then(|metadata| metadata.get("content-type")), + Some(&"user-content-type".to_string()) + ); + assert_eq!( + metadata_head + .metadata() + .and_then(|metadata| metadata.get("x-amz-storage-class")), + Some(&"user-storage-class".to_string()) + ); + } + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn multipart_copy_preserves_standard_and_rrs() -> Result<(), Box> { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(Vec::new()).await?; + let client = env.create_s3_client(); + let bucket = "multipart-storage-class-copy"; + let source_key = "source.bin"; + let source_body = vec![b'c'; 1024 * 1024]; + env.create_test_bucket(bucket).await?; + client + .put_object() + .bucket(bucket) + .key(source_key) + .body(ByteStream::from(source_body.clone())) + .send() + .await?; + + for storage_class in [StorageClass::Standard, StorageClass::ReducedRedundancy] { + let class_name = storage_class.as_str(); + let key = format!("copy-{class_name}.bin"); + let create = client + .create_multipart_upload() + .bucket(bucket) + .key(&key) + .storage_class(storage_class.clone()) + .send() + .await?; + let upload_id = create.upload_id().ok_or("CreateMultipartUpload returned no upload ID")?; + let copied = client + .upload_part_copy() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .part_number(1) + .copy_source(format!("{bucket}/{source_key}")) + .send() + .await?; + let e_tag = copied + .copy_part_result() + .and_then(|result| result.e_tag()) + .ok_or("UploadPartCopy returned no ETag")?; + + let parts = client + .list_parts() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .send() + .await?; + assert_eq!(parts.storage_class().map(StorageClass::as_str), Some(class_name)); + assert_eq!(parts.parts().len(), 1); + + client + .complete_multipart_upload() + .bucket(bucket) + .key(&key) + .upload_id(upload_id) + .multipart_upload( + CompletedMultipartUpload::builder() + .parts(CompletedPart::builder().part_number(1).e_tag(e_tag).build()) + .build(), + ) + .send() + .await?; + assert_completed_object(&client, bucket, &key, class_name, &source_body).await?; + } + + env.stop_server(); + Ok(()) + } + + #[tokio::test] + async fn invalid_and_aborted_uploads_leave_no_session_or_object() -> Result<(), Box> { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(Vec::new()).await?; + let client = env.create_s3_client(); + let bucket = "multipart-storage-class-errors"; + let invalid_key = "invalid.bin"; + let aborted_key = "aborted.bin"; + env.create_test_bucket(bucket).await?; + + let invalid = client + .create_multipart_upload() + .bucket(bucket) + .key(invalid_key) + .storage_class(StorageClass::from("INVALID")) + .send() + .await + .expect_err("invalid storage class should be rejected"); + assert_eq!( + invalid.as_service_error().and_then(ProvideErrorMetadata::code), + Some("InvalidStorageClass") + ); + let after_invalid = client + .list_multipart_uploads() + .bucket(bucket) + .prefix(invalid_key) + .send() + .await?; + assert!( + after_invalid.uploads().is_empty(), + "validation failure must not create a multipart session" + ); + + let create = client + .create_multipart_upload() + .bucket(bucket) + .key(aborted_key) + .storage_class(StorageClass::ReducedRedundancy) + .send() + .await?; + let upload_id = create.upload_id().ok_or("CreateMultipartUpload returned no upload ID")?; + client + .upload_part() + .bucket(bucket) + .key(aborted_key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from_static(b"aborted multipart part")) + .send() + .await?; + let before_abort = client + .list_parts() + .bucket(bucket) + .key(aborted_key) + .upload_id(upload_id) + .send() + .await?; + assert_eq!(before_abort.storage_class().map(StorageClass::as_str), Some("REDUCED_REDUNDANCY")); + + client + .abort_multipart_upload() + .bucket(bucket) + .key(aborted_key) + .upload_id(upload_id) + .send() + .await?; + let after_abort = client + .list_parts() + .bucket(bucket) + .key(aborted_key) + .upload_id(upload_id) + .send() + .await + .expect_err("aborted upload should not be resumable"); + assert_eq!(after_abort.as_service_error().and_then(ProvideErrorMetadata::code), Some("NoSuchUpload")); + let remaining_uploads = client + .list_multipart_uploads() + .bucket(bucket) + .prefix(aborted_key) + .send() + .await?; + assert!(remaining_uploads.uploads().is_empty(), "abort should remove the multipart session"); + let aborted_head = client + .head_object() + .bucket(bucket) + .key(aborted_key) + .send() + .await + .expect_err("aborted upload should not create an object"); + assert_eq!(aborted_head.raw_response().map(|response| response.status().as_u16()), Some(404)); + + env.stop_server(); + Ok(()) + } +} diff --git a/docs/testing/e2e-suite-inventory.md b/docs/testing/e2e-suite-inventory.md index 2d381b456..92e02fb14 100644 --- a/docs/testing/e2e-suite-inventory.md +++ b/docs/testing/e2e-suite-inventory.md @@ -61,6 +61,7 @@ | list_objects_v2_pagination_test | 12 | ✅ | | mc_mirror_small_bucket_test | 1 | | | multipart_auth_test | 109 | | +| multipart_storage_class_test | 3 | ✅ | | namespace_lock_quorum_test | 2 | | | negative_sigv4_test | 6 | ✅ | | notification_webhook_test | 2 | ✅ | @@ -84,4 +85,4 @@ `notification_webhook_test` also has 1 ignored store-and-forward regression tracked by rustfs#4852; ignored tests are excluded from the active counts above. -**Total listed: 472 tests across 64 modules · PR smoke subset: 119 tests / 29 modules** (27 full modules + 4 `reliant` tests + 20 of `replication_extension_test`) **· nightly `e2e-repl-nightly`: 27 tests** · generated 2026-07-24. +**Total listed: 475 tests across 65 modules · PR smoke subset: 122 tests / 30 modules** (28 full modules + 4 `reliant` tests + 20 of `replication_extension_test`) **· nightly `e2e-repl-nightly`: 27 tests** · generated 2026-07-24. diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index d740acf65..099a041e5 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -41,7 +41,7 @@ use super::storage_api::multipart_usecase::io::{HashReader, WriteEncryption, Wri use super::storage_api::multipart_usecase::object_utils::to_s3s_etag; use super::storage_api::multipart_usecase::options::{ copy_src_opts, extract_metadata_from_mime, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts, - parse_copy_source_range, put_opts, validate_archive_content_encoding, + namespace_reserved_user_metadata, parse_copy_source_range, put_opts, validate_archive_content_encoding, }; use super::storage_api::multipart_usecase::s3_api::multipart::{ ListMultipartUploadsParams, build_list_multipart_uploads_output, build_list_parts_output, @@ -74,7 +74,7 @@ use rustfs_targets::EventName; use rustfs_utils::CompressionAlgorithm; use rustfs_utils::http::{ SUFFIX_REPLICATION_STATUS, SUFFIX_REPLICATION_TIMESTAMP, get_source_scheme, - headers::{AMZ_DECODED_CONTENT_LENGTH, AMZ_OBJECT_TAGGING}, + headers::{AMZ_DECODED_CONTENT_LENGTH, AMZ_OBJECT_TAGGING, AMZ_STORAGE_CLASS}, insert_str, }; use s3s::dto::{ @@ -85,9 +85,7 @@ use s3s::dto::{ }; use s3s::header::{X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE}; use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; -#[cfg(test)] -use std::collections::HashMap; -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::str::FromStr; use std::sync::Arc; use tokio::sync::RwLock; @@ -178,6 +176,26 @@ fn complete_part_from_s3(value: CompletedPart) -> CompletePart { } } +fn create_multipart_upload_metadata( + input_metadata: Option>, + headers: &HeaderMap, + tagging: Option, + storage_class: Option<&s3s::dto::StorageClass>, +) -> HashMap { + let mut metadata = input_metadata.unwrap_or_default(); + namespace_reserved_user_metadata(&mut metadata); + extract_metadata_from_mime(headers, &mut metadata); + + if let Some(tags) = tagging { + metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags); + } + if let Some(storage_class) = storage_class { + metadata.insert(AMZ_STORAGE_CLASS.to_owned(), storage_class.as_str().to_owned()); + } + + metadata +} + async fn validate_table_catalog_object_mutation(bucket: &str, key: &str) -> S3Result<()> { table_catalog::validate_bucket_object_mutation(bucket, key) .await @@ -643,12 +661,7 @@ impl DefaultMultipartUsecase { req.headers.get("content-encoding").and_then(|value| value.to_str().ok()), )?; - let mut metadata = input_metadata.unwrap_or_default(); - extract_metadata_from_mime(&req.headers, &mut metadata); - - if let Some(tags) = tagging { - metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags); - } + let mut metadata = create_multipart_upload_metadata(input_metadata, &req.headers, tagging, storage_class.as_ref()); let has_explicit_object_lock_retention = object_lock_mode.is_some() || object_lock_retain_until_date.is_some(); if let Some(object_lock_metadata) = build_put_like_object_lock_metadata( @@ -1695,6 +1708,42 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass); } + #[test] + fn create_multipart_upload_metadata_persists_dto_storage_class_without_raw_header() { + let metadata = create_multipart_upload_metadata( + None, + &HeaderMap::new(), + None, + Some(&StorageClass::from_static("REDUCED_REDUNDANCY")), + ); + + assert_eq!(metadata.get(AMZ_STORAGE_CLASS), Some(&"REDUCED_REDUNDANCY".to_string())); + } + + #[test] + fn create_multipart_upload_metadata_keeps_user_and_system_namespaces_separate() { + let mut headers = HeaderMap::new(); + headers.insert("content-type", HeaderValue::from_static("application/octet-stream")); + headers.insert(AMZ_STORAGE_CLASS, HeaderValue::from_static("STANDARD")); + let input_metadata = HashMap::from([ + ("content-type".to_string(), "user-content-type".to_string()), + (AMZ_STORAGE_CLASS.to_string(), "user-storage-class".to_string()), + ]); + + let metadata = create_multipart_upload_metadata( + Some(input_metadata), + &headers, + Some("project=rustfs".to_string()), + Some(&StorageClass::from_static("REDUCED_REDUNDANCY")), + ); + + assert_eq!(metadata.get("content-type"), Some(&"application/octet-stream".to_string())); + assert_eq!(metadata.get("x-amz-meta-content-type"), Some(&"user-content-type".to_string())); + assert_eq!(metadata.get(AMZ_STORAGE_CLASS), Some(&"REDUCED_REDUNDANCY".to_string())); + assert_eq!(metadata.get("x-amz-meta-x-amz-storage-class"), Some(&"user-storage-class".to_string())); + assert_eq!(metadata.get(AMZ_OBJECT_TAGGING), Some(&"project=rustfs".to_string())); + } + #[tokio::test] async fn execute_complete_multipart_upload_rejects_missing_parts_payload() { let input = CompleteMultipartUploadInput::builder()