From 6cecf7d26f5c8551104d5e815dbbfd75ee670dcd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Thu, 13 Aug 2026 08:12:49 +0800 Subject: [PATCH] fix(rebalance): isolate internal multipart uploads --- crates/ecstore/src/data_movement/mod.rs | 34 ++++--- crates/ecstore/src/set_disk/ops/multipart.rs | 102 +++++++++++++++++-- 2 files changed, 113 insertions(+), 23 deletions(-) diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 28a9b331a..7ea6cfdba 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -34,7 +34,7 @@ use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIn use rustfs_utils::http::{ AMZ_OBJECT_TAGGING, SUFFIX_ACTUAL_SIZE, SUFFIX_COMPRESSION_SIZE, SUFFIX_CRC, SUFFIX_DATA_MOVED, SUFFIX_DATA_MOVEMENT_UPLOAD, SUFFIX_PART_CHECKSUMS, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITIONED_OBJECTNAME, - SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, strip_internal_prefix_preserving_case, + SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, has_internal_suffix, }; use rustfs_utils::path::encode_dir_object; use std::collections::{BTreeMap, HashMap}; @@ -198,8 +198,7 @@ fn data_movement_user_defined(object_info: &ObjectInfo) -> HashMap>(); @@ -227,7 +226,7 @@ fn data_movement_user_defined(object_info: &ObjectInfo) -> HashMap { @@ -571,10 +578,6 @@ fn are_equivalent_data_movement_parts_for(source: &[ObjectPartInfo], target: &[O }) } -fn is_data_movement_internal_metadata(key: &str, suffix: &str) -> bool { - strip_internal_prefix_preserving_case(key).is_some_and(|candidate| candidate.eq_ignore_ascii_case(suffix)) -} - fn is_canonical_data_movement_internal_metadata(key: &str, suffix: &str) -> bool { key.strip_prefix(rustfs_utils::http::RUSTFS_INTERNAL_PREFIX) == Some(suffix) || key.strip_prefix(rustfs_utils::http::MINIO_INTERNAL_PREFIX) == Some(suffix) @@ -598,7 +601,7 @@ fn data_movement_layout_marker_presence(object_info: &ObjectInfo) -> Option Option { let mut present = false; for (key, value) in object_info.user_defined.iter() { - if !is_data_movement_internal_metadata(key, suffix) { + if !has_internal_suffix(key, suffix) { continue; } present = true; @@ -612,7 +615,7 @@ fn data_movement_size_marker_presence(object_info: &ObjectInfo, suffix: &str, ex fn data_movement_checksum_marker_presence(object_info: &ObjectInfo) -> Option { let mut present = false; for (key, value) in object_info.user_defined.iter() { - if !is_data_movement_internal_metadata(key, SUFFIX_CRC) { + if !has_internal_suffix(key, SUFFIX_CRC) { continue; } let Some(checksum) = object_info.checksum.as_deref().filter(|checksum| !checksum.is_empty()) else { @@ -663,7 +666,7 @@ fn is_data_movement_rewritten_transition_metadata(object_info: &ObjectInfo, key: canonical && !preserves_unusable_version || expected.is_some_and(|expected| { !expected.is_empty() - && is_data_movement_internal_metadata(key, suffix) + && has_internal_suffix(key, suffix) && rustfs_utils::http::get_consistent_str(&object_info.user_defined, suffix) == Some(expected) }) }) @@ -679,11 +682,11 @@ fn is_data_movement_rewritten_metadata(object_info: &ObjectInfo, key: &str, norm crate::object_api::ENCRYPTED_PART_LAYOUT_QUORUM_SUFFIX, ] .iter() - .any(|suffix| is_data_movement_internal_metadata(key, suffix)) + .any(|suffix| has_internal_suffix(key, suffix)) || is_data_movement_rewritten_transition_metadata(object_info, key) || key == rustfs_rio::RUSTFS_MULTIPART_CHECKSUM || key == rustfs_rio::RUSTFS_MULTIPART_CHECKSUM_TYPE - || normalize_compression_size && is_data_movement_internal_metadata(key, SUFFIX_COMPRESSION_SIZE) + || normalize_compression_size && has_internal_suffix(key, SUFFIX_COMPRESSION_SIZE) } pub(crate) fn is_equivalent_data_movement_metadata( @@ -771,6 +774,7 @@ fn is_equivalent_data_movement_object_identity( && are_equivalent_data_movement_parts_for(&source.parts, &target.parts, compare_part_checksums) } +#[cfg(test)] fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool { is_equivalent_data_movement_object_identity(source, target, true, true) } @@ -3273,7 +3277,7 @@ mod tests { upload_identity, ); assert!( - !resolve_data_movement_overwrite_resume_result(&err, Ok(Some(legacy_target)), &legacy_source, 0, 1) + !resolve_data_movement_overwrite_resume_result_for(&err, Ok(Some(legacy_target)), &legacy_source, 0, 1, true) .expect("an old-node takeover must not discard legacy part checksums") ); } diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index caf8e1e73..8b01b7797 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -207,6 +207,20 @@ fn validate_multipart_bucket_incarnation( Err(StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned())) } +fn ensure_data_movement_upload_access( + fi: &FileInfo, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, +) -> Result<()> { + if rustfs_utils::http::contains_key_str(&fi.metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD) && !opts.data_movement + { + return Err(StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned())); + } + Ok(()) +} + async fn ensure_multipart_bucket_incarnation( ctx: &crate::runtime::instance::InstanceContext, fi: &FileInfo, @@ -688,6 +702,12 @@ impl SetDisks { { return Ok(None); } + if rustfs_utils::http::contains_key_str( + &file_info.metadata, + rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, + ) { + return Ok(None); + } let object = match ( file_info.metadata.get(RUSTFS_MULTIPART_BUCKET_KEY), @@ -839,6 +859,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, true).await?; + ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id) .await?; @@ -1091,6 +1112,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { ) }; let (commit_fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?; + ensure_data_movement_upload_access(&commit_fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation( &self.ctx, &commit_fi, @@ -1173,6 +1195,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .acquire_multipart_upload_read_lock("list_object_parts", bucket, object, upload_id, opts) .await?; let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?; + ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id) .await?; @@ -1507,6 +1530,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .check_upload_id_exists(bucket, object, upload_id, false) .await .map_err(|e| to_object_err(e, vec![bucket, object, upload_id]))?; + ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id) .await?; ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?; @@ -1529,6 +1553,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts) .await?; let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, true).await?; + ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id) .await?; ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?; @@ -1583,11 +1608,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let expected_restore_operation_id = restore_commit_operation_id_from_metadata(&opts.user_defined)?; let (mut fi, files_metas) = self.check_upload_id_exists(bucket, object, upload_id, true).await?; - if rustfs_utils::http::contains_key_str(&fi.metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD) - && !opts.data_movement - { - return Err(StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned())); - } + ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?; ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id) .await?; let has_layout_candidate = range_seek_rollout_enabled @@ -2843,7 +2864,7 @@ mod tests { #[tokio::test] #[serial] - async fn data_movement_upload_rejects_external_completion() { + async fn data_movement_upload_is_hidden_from_external_multipart_operations() { let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; let bucket = "data-movement-upload-complete-bucket"; make_bucket_on_all(&disk_stores, bucket).await; @@ -2904,8 +2925,52 @@ mod tests { user_defined: metadata, ..Default::default() }; - let (upload_id, parts) = - stage_upload_with_create_opts(&set_disks, bucket, object, b"data movement multipart body", &create_opts).await; + let upload = set_disks + .new_multipart_upload(bucket, object, &create_opts) + .await + .expect("data movement upload should be created"); + let upload_id = upload.upload_id; + let mut part_reader = PutObjReader::from_vec(b"data movement multipart body".to_vec()); + let uploaded_part = set_disks + .put_object_part(bucket, object, &upload_id, 1, &mut part_reader, &create_opts) + .await + .expect("data movement part upload should retain ownership of its upload"); + let parts = vec![CompletePart { + part_num: uploaded_part.part_num, + etag: uploaded_part.etag, + ..Default::default() + }]; + + let listed = set_disks + .list_multipart_uploads_for_incarnation(bucket, "", None, None, None, 1000, None) + .await + .expect("external multipart listing should succeed"); + assert!(!listed.uploads.iter().any(|upload| upload.upload_id == upload_id)); + + let get_err = set_disks + .get_multipart_info(bucket, object, &upload_id, &ObjectOptions::default()) + .await + .expect_err("external multipart metadata reads must not expose a data movement upload"); + assert!(matches!(get_err, StorageError::InvalidUploadID(..))); + + let list_parts_err = set_disks + .list_object_parts(bucket, object, &upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default()) + .await + .expect_err("external part listings must not expose a data movement upload"); + assert!(matches!(list_parts_err, StorageError::InvalidUploadID(..))); + + let mut external_part = PutObjReader::from_vec(b"external overwrite".to_vec()); + let put_err = set_disks + .put_object_part(bucket, object, &upload_id, 1, &mut external_part, &ObjectOptions::default()) + .await + .expect_err("external part uploads must not modify a data movement upload"); + assert!(matches!(put_err, StorageError::InvalidUploadID(..))); + + let abort_err = set_disks + .abort_multipart_upload(bucket, object, &upload_id, &ObjectOptions::default()) + .await + .expect_err("external aborts must not remove a data movement upload"); + assert!(matches!(abort_err, StorageError::InvalidUploadID(..))); let external_err = set_disks .clone() @@ -2914,6 +2979,12 @@ mod tests { .expect_err("external completion must not finalize a data movement upload"); assert!(matches!(external_err, StorageError::InvalidUploadID(..))); + let internal_parts = set_disks + .list_object_parts(bucket, object, &upload_id, None, MAX_PARTS_COUNT, &create_opts) + .await + .expect("data movement part listing should retain ownership of its upload"); + assert_eq!(internal_parts.parts.len(), 1); + set_disks .clone() .complete_multipart_upload(bucket, object, &upload_id, parts, &create_opts) @@ -2957,6 +3028,21 @@ mod tests { .map(String::as_str), Some("AAAAAA==") ); + + let abort_object = "owned-abort-object"; + let abort_upload = set_disks + .new_multipart_upload(bucket, abort_object, &create_opts) + .await + .expect("data movement abort upload should be created"); + let abort_err = set_disks + .abort_multipart_upload(bucket, abort_object, &abort_upload.upload_id, &ObjectOptions::default()) + .await + .expect_err("external abort must not remove the second data movement upload"); + assert!(matches!(abort_err, StorageError::InvalidUploadID(..))); + set_disks + .abort_multipart_upload(bucket, abort_object, &abort_upload.upload_id, &create_opts) + .await + .expect("data movement abort should retain ownership of its upload"); } #[tokio::test]