From e7853e37a986c8de5a39645f8cac59f03afe6b50 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Wed, 12 Aug 2026 17:33:37 +0800 Subject: [PATCH] fix(rebalance): harden multipart retry replacement --- crates/ecstore/src/core/pools.rs | 11 +- crates/ecstore/src/data_movement/mod.rs | 288 ++++++++++--- .../ecstore/src/services/rebalance/entry.rs | 14 +- .../src/services/rebalance/migration.rs | 34 -- .../rebalance/rebalance_unit_tests.rs | 22 +- .../src/set_disk/core/io_primitives.rs | 14 +- crates/ecstore/src/set_disk/ops/multipart.rs | 167 +++++++ crates/ecstore/src/set_disk/ops/object.rs | 9 +- crates/ecstore/src/store/init.rs | 408 +++++++++++++++++- crates/ecstore/src/store/multipart.rs | 81 ++-- crates/ecstore/src/store/object.rs | 4 +- crates/filemeta/src/filemeta.rs | 16 +- crates/filemeta/src/filemeta/version.rs | 165 ++++--- 13 files changed, 997 insertions(+), 236 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 684112e64..70aa73bde 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -3178,10 +3178,12 @@ impl ECStore { entry.name.as_str(), &fivs, &cleanup_preflight_allowed_missing, - expected_bucket_incarnation_id, - bucket_incarnation_fence - .as_ref() - .and_then(|guard| guard.namespace_lock_guard()), + data_movement::SourceCleanupBucketFence { + expected_incarnation_id: expected_bucket_incarnation_id, + lifecycle_guard: bucket_incarnation_fence + .as_ref() + .and_then(|guard| guard.namespace_lock_guard()), + }, "decommission", ) .await @@ -3341,7 +3343,6 @@ impl ECStore { let lifecycle_config = lifecycle_config.clone(); let object_lock_config = object_lock_config.clone(); let replication_config = replication_config.clone(); - let expected_bucket_incarnation_id = expected_bucket_incarnation_id; let entry_error = entry_error.clone(); let callback_rx = rx.clone(); move |entry: MetaCacheEntry| { diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index d6f3de58a..28a9b331a 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -294,12 +294,15 @@ pub(crate) fn prepare_tiered_data_movement_file_info(file_info: &mut rustfs_file } fn prepare_tiered_data_movement_file_info_for(file_info: &mut rustfs_filemeta::FileInfo, writer_enabled: bool) -> Result<()> { + if !writer_enabled { + rustfs_utils::http::remove_str(&mut file_info.metadata, SUFFIX_PART_CHECKSUMS); + for part in &mut file_info.parts { + part.checksums = None; + } + return Ok(()); + } + SetDisks::hydrate_selected_fileinfo_part_checksums(file_info).map_err(|_| Error::FileCorrupt)?; - let has_part_checksums = file_info - .parts - .iter() - .any(|part| part.checksums.as_ref().is_some_and(|checksums| !checksums.is_empty())); - ensure_data_movement_part_checksum_writer_allowed(has_part_checksums, writer_enabled)?; rustfs_utils::http::remove_str(&mut file_info.metadata, SUFFIX_PART_CHECKSUMS); if let Some(encoded) = data_movement_part_checksums(&file_info.parts)? { rustfs_utils::http::insert_str(&mut file_info.metadata, SUFFIX_PART_CHECKSUMS, encoded); @@ -324,15 +327,6 @@ fn data_movement_part_checksum_writer_enabled() -> bool { ) } -fn ensure_data_movement_part_checksum_writer_allowed(has_part_checksums: bool, enabled: bool) -> Result<()> { - if has_part_checksums && !enabled { - return Err(Error::other( - "data movement per-part checksums require the fleet-confirmed writer capability", - )); - } - Ok(()) -} - fn should_use_multipart_data_movement(object_info: &ObjectInfo, has_part_checksums: bool) -> bool { object_info.is_multipart() || has_part_checksums @@ -340,7 +334,11 @@ fn should_use_multipart_data_movement(object_info: &ObjectInfo, has_part_checksu || object_info.parts.first().is_some_and(|part| part.number != 1) } -fn data_movement_complete_multipart_opts(object_info: &ObjectInfo, src_pool_idx: usize) -> Result { +fn data_movement_complete_multipart_opts( + object_info: &ObjectInfo, + src_pool_idx: usize, + preserve_part_checksums: bool, +) -> Result { let mut user_defined = HashMap::new(); insert_data_movement_checksum(&mut user_defined, object_info); let actual_size = object_info @@ -350,7 +348,7 @@ fn data_movement_complete_multipart_opts(object_info: &ObjectInfo, src_pool_idx: return Err(Error::other("data movement source actual size is unknown")); } rustfs_utils::http::insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, actual_size.to_string()); - if let Some(encoded) = data_movement_part_checksums(&object_info.parts)? { + if preserve_part_checksums && let Some(encoded) = data_movement_part_checksums(&object_info.parts)? { rustfs_utils::http::insert_str(&mut user_defined, SUFFIX_PART_CHECKSUMS, encoded); } Ok(ObjectOptions { @@ -391,6 +389,40 @@ pub(crate) fn data_movement_target_precondition() -> HTTPPreconditions { } } +pub(crate) fn can_replace_stale_data_movement_target(target: &ObjectInfo, opts: &ObjectOptions) -> bool { + let Some(preconditions) = opts.http_preconditions.as_ref() else { + return false; + }; + if !opts.data_movement + || preconditions.if_none_match_value() != Some("*") + || preconditions.if_match_value().is_some() + || target.delete_marker + { + return false; + } + + let rustfs_marker = rustfs_utils::http::internal_key_rustfs(SUFFIX_DATA_MOVED); + let minio_marker = format!("{}{SUFFIX_DATA_MOVED}", rustfs_utils::http::MINIO_INTERNAL_PREFIX); + if rustfs_utils::http::get_consistent_str(&target.user_defined, SUFFIX_DATA_MOVED) != Some("true") + || target.user_defined.get(&rustfs_marker).map(String::as_str) != Some("true") + || target.user_defined.get(&minio_marker).map(String::as_str) != Some("true") + { + return false; + } + + let version_matches = match (opts.version_id.as_deref(), target.version_id) { + (None, None) => true, + (Some(expected), Some(actual)) => uuid::Uuid::parse_str(expected).ok() == Some(actual), + _ => false, + }; + + version_matches + && target + .mod_time + .zip(opts.mod_time) + .is_some_and(|(target_time, source_time)| target_time < source_time) +} + fn data_movement_put_object_reader( bucket: &str, object_info: &ObjectInfo, @@ -489,7 +521,7 @@ fn effective_part_actual_size(part: &ObjectPartInfo) -> Option { .or_else(|| i64::try_from(part.size).ok()) } -fn is_equivalent_data_movement_part(source: &ObjectPartInfo, target: &ObjectPartInfo) -> bool { +fn is_equivalent_data_movement_part(source: &ObjectPartInfo, target: &ObjectPartInfo, compare_checksums: bool) -> bool { // Multipart migration rewrites part timestamps. source.number == target.number && source.etag == target.etag @@ -500,8 +532,9 @@ fn is_equivalent_data_movement_part(source: &ObjectPartInfo, target: &ObjectPart ) // A missing target compression index selects the safe full-read fallback. && (target.index.is_none() || source.index == target.index) - && source.checksums.as_ref().filter(|checksums| !checksums.is_empty()) - == target.checksums.as_ref().filter(|checksums| !checksums.is_empty()) + && (!compare_checksums + || source.checksums.as_ref().filter(|checksums| !checksums.is_empty()) + == target.checksums.as_ref().filter(|checksums| !checksums.is_empty())) } fn data_movement_parts_by_number(parts: &[ObjectPartInfo]) -> Option> { @@ -516,6 +549,10 @@ fn data_movement_parts_by_number(parts: &[ObjectPartInfo]) -> Option bool { + are_equivalent_data_movement_parts_for(source, target, true) +} + +fn are_equivalent_data_movement_parts_for(source: &[ObjectPartInfo], target: &[ObjectPartInfo], compare_checksums: bool) -> bool { if source.len() != target.len() { return false; } @@ -530,7 +567,7 @@ pub(crate) fn are_equivalent_data_movement_parts(source: &[ObjectPartInfo], targ source_parts.iter().all(|(number, source_part)| { target_parts .get(number) - .is_some_and(|target_part| is_equivalent_data_movement_part(source_part, target_part)) + .is_some_and(|target_part| is_equivalent_data_movement_part(source_part, target_part, compare_checksums)) }) } @@ -699,7 +736,12 @@ pub(crate) fn is_equivalent_data_movement_metadata( .all(|(key, value)| source.user_defined.get(key) == Some(value)) } -fn is_equivalent_data_movement_object_identity(source: &ObjectInfo, target: &ObjectInfo, compare_mod_time: bool) -> bool { +fn is_equivalent_data_movement_object_identity( + source: &ObjectInfo, + target: &ObjectInfo, + compare_mod_time: bool, + compare_part_checksums: bool, +) -> bool { let (Some(source_actual_size), Some(target_actual_size)) = (effective_actual_size(source), effective_actual_size(target)) else { return false; @@ -726,11 +768,11 @@ fn is_equivalent_data_movement_object_identity(source: &ObjectInfo, target: &Obj && source.transitioned_object.free_version == target.transitioned_object.free_version && source.transitioned_object.status == target.transitioned_object.status && source.transition_version_state == target.transition_version_state - && are_equivalent_data_movement_parts(&source.parts, &target.parts) + && are_equivalent_data_movement_parts_for(&source.parts, &target.parts, compare_part_checksums) } fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool { - is_equivalent_data_movement_object_identity(source, target, true) + is_equivalent_data_movement_object_identity(source, target, true, true) } fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool { @@ -743,11 +785,11 @@ fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: .is_some_and(|(source_time, target_time)| target_time > source_time) } -fn is_data_movement_upload_takeover_target(source: &ObjectInfo, target: &ObjectInfo) -> bool { +fn is_data_movement_upload_takeover_target(source: &ObjectInfo, target: &ObjectInfo, compare_part_checksums: bool) -> bool { let identity = data_movement_upload_identity(source); source.mod_time.is_some() && rustfs_utils::http::get_consistent_str(&target.user_defined, SUFFIX_DATA_MOVEMENT_UPLOAD) == Some(identity.as_str()) - && is_equivalent_data_movement_object_identity(source, target, false) + && is_equivalent_data_movement_object_identity(source, target, false, compare_part_checksums) } #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] @@ -895,6 +937,12 @@ pub(crate) enum SourceCleanupError { Storage(#[from] Error), } +#[derive(Clone, Copy, Default)] +pub(crate) struct SourceCleanupBucketFence<'a> { + pub(crate) expected_incarnation_id: Option, + pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>, +} + fn ensure_source_cleanup_versions_match( expected: &FileInfoVersions, current: &FileInfoVersions, @@ -1007,8 +1055,7 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( object: &str, expected: &FileInfoVersions, allowed_missing: &[SourceCleanupVersionIdentity], - expected_bucket_incarnation_id: Option, - bucket_lifecycle_guard: Option<&rustfs_lock::NamespaceLockGuard>, + bucket_fence: SourceCleanupBucketFence<'_>, op_label: &str, ) -> std::result::Result { let cleanup_key = encode_dir_object(object); @@ -1018,7 +1065,10 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( .await .map_err(Error::from)?; - if bucket_lifecycle_guard.is_some_and(rustfs_lock::NamespaceLockGuard::is_lock_lost) { + if bucket_fence + .lifecycle_guard + .is_some_and(rustfs_lock::NamespaceLockGuard::is_lock_lost) + { return Err(SourceCleanupError::Storage(Error::other(format!( "{op_label}: bucket incarnation fence was lost before source cleanup" )))); @@ -1034,11 +1084,11 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( delete_prefix_object: true, data_movement: true, no_lock: true, - expected_bucket_incarnation_id, + expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id, ..Default::default() }; opts.add_namespace_lock_guard(&_guard); - if let Some(bucket_lifecycle_guard) = bucket_lifecycle_guard { + if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard { opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard); } let result = set.delete_object(bucket, cleanup_key.as_str(), opts).await; @@ -1086,6 +1136,24 @@ fn resolve_data_movement_overwrite_resume_result( source: &ObjectInfo, src_pool_idx: usize, target_pool_idx: usize, +) -> Result { + resolve_data_movement_overwrite_resume_result_for( + err, + target_result, + source, + src_pool_idx, + target_pool_idx, + data_movement_part_checksum_writer_enabled(), + ) +} + +fn resolve_data_movement_overwrite_resume_result_for( + err: &Error, + target_result: Result>, + source: &ObjectInfo, + src_pool_idx: usize, + target_pool_idx: usize, + compare_part_checksums: bool, ) -> Result { if !should_check_data_movement_overwrite_resume(err) || !should_check_data_movement_resume_target(src_pool_idx, target_pool_idx) @@ -1097,11 +1165,11 @@ fn resolve_data_movement_overwrite_resume_result( return Ok(false); }; - if is_equivalent_data_movement_object(source, &target) { + if is_equivalent_data_movement_object_identity(source, &target, true, compare_part_checksums) { return Ok(true); } - if is_data_movement_upload_takeover_target(source, &target) { + if is_data_movement_upload_takeover_target(source, &target, compare_part_checksums) { return Ok(true); } @@ -1115,17 +1183,19 @@ async fn should_treat_data_movement_overwrite_as_complete( bucket: &str, object_info: &ObjectInfo, err: &Error, + compare_part_checksums: bool, ) -> Result { if !should_check_data_movement_overwrite_resume(err) { return Ok(false); } - resolve_data_movement_overwrite_resume_result( + resolve_data_movement_overwrite_resume_result_for( err, find_data_movement_target_info(store, target_pool_idx, bucket, object_info).await, object_info, src_pool_idx, target_pool_idx, + compare_part_checksums, ) } @@ -1174,9 +1244,7 @@ pub(crate) async fn migrate_object( .iter() .any(|part| part.checksums.as_ref().is_some_and(|checksums| !checksums.is_empty())); - ensure_data_movement_part_checksum_writer_allowed(has_part_checksums, data_movement_part_checksum_writer_enabled()).map_err( - |err| data_movement_stage_error(op_label, "prepare_new_multipart", bucket.as_str(), object_info.name.as_str(), err), - )?; + let preserve_part_checksums = data_movement_part_checksum_writer_enabled(); if should_use_multipart_data_movement(&object_info, has_part_checksums) { let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx); @@ -1225,21 +1293,22 @@ pub(crate) async fn migrate_object( err, ) })?; + let part_opts = ObjectOptions { + part_number: Some(part.number), + preserve_etag: Some(part.etag.clone()), + data_movement: true, + src_pool_idx: pool_idx, + expected_bucket_incarnation_id, + ..Default::default() + }; let pi = match store .put_object_part_for_data_movement( target_pool_idx, &bucket, &object_info.name, &res.upload_id, - part.number, &mut data, - &ObjectOptions { - preserve_etag: Some(part.etag.clone()), - data_movement: true, - src_pool_idx: pool_idx, - expected_bucket_incarnation_id, - ..Default::default() - }, + &part_opts, ) .await { @@ -1265,9 +1334,16 @@ pub(crate) async fn migrate_object( }; } - let mut complete_multipart_opts = data_movement_complete_multipart_opts(&object_info, pool_idx).map_err(|err| { - data_movement_stage_error(op_label, "prepare_complete_multipart", bucket.as_str(), object_info.name.as_str(), err) - })?; + let mut complete_multipart_opts = + data_movement_complete_multipart_opts(&object_info, pool_idx, preserve_part_checksums).map_err(|err| { + data_movement_stage_error( + op_label, + "prepare_complete_multipart", + bucket.as_str(), + object_info.name.as_str(), + err, + ) + })?; complete_multipart_opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id; if let Err(err) = store .clone() @@ -1288,6 +1364,7 @@ pub(crate) async fn migrate_object( bucket.as_str(), &object_info, &err, + preserve_part_checksums, ) .await? { @@ -1339,6 +1416,7 @@ pub(crate) async fn migrate_object( bucket.as_str(), &object_info, &abort_err, + preserve_part_checksums, ) .await? { @@ -1442,6 +1520,7 @@ pub(crate) async fn migrate_object( bucket.as_str(), &object_info, &err, + preserve_part_checksums, ) .await? { @@ -1897,10 +1976,17 @@ mod tests { }; let new_opts = data_movement_new_multipart_opts(&object_info, 0); - let complete_opts = data_movement_complete_multipart_opts(&object_info, 0).expect("complete opts should be created"); + let compatible_opts = + data_movement_complete_multipart_opts(&object_info, 0, false).expect("compatible opts should be created"); + let complete_opts = + data_movement_complete_multipart_opts(&object_info, 0, true).expect("complete opts should be created"); assert!(!rustfs_utils::http::contains_key_str(&new_opts.user_defined, SUFFIX_PART_CHECKSUMS)); assert!(rustfs_utils::http::contains_key_str(&new_opts.user_defined, SUFFIX_DATA_MOVEMENT_UPLOAD)); + assert!(!rustfs_utils::http::contains_key_str( + &compatible_opts.user_defined, + SUFFIX_PART_CHECKSUMS + )); assert_eq!( rustfs_utils::http::get_consistent_str(&complete_opts.user_defined, SUFFIX_PART_CHECKSUMS), Some(r#"[[2,[["CRC32C","crc32c-value"]]]]"#) @@ -1921,11 +2007,14 @@ mod tests { assert!(!data_movement_part_checksum_writer_enabled_for(true, false)); assert!(!data_movement_part_checksum_writer_enabled_for(false, true)); assert!(data_movement_part_checksum_writer_enabled_for(true, true)); - let has_part_checksums = object_info.parts.iter().any(|part| part.checksums.is_some()); - assert!(has_part_checksums); - assert!(ensure_data_movement_part_checksum_writer_allowed(has_part_checksums, false).is_err()); - assert!(ensure_data_movement_part_checksum_writer_allowed(has_part_checksums, true).is_ok()); - assert!(ensure_data_movement_part_checksum_writer_allowed(false, false).is_ok()); + let mut compatible = FileInfo { + parts: object_info.parts.as_ref().clone(), + ..Default::default() + }; + prepare_tiered_data_movement_file_info_for(&mut compatible, false) + .expect("disabled sidecar writer should preserve data movement compatibility"); + assert!(compatible.parts.iter().all(|part| part.checksums.is_none())); + assert!(!rustfs_utils::http::contains_key_str(&compatible.metadata, SUFFIX_PART_CHECKSUMS)); let empty = ObjectInfo { parts: Arc::new(vec![ObjectPartInfo { @@ -1952,12 +2041,16 @@ mod tests { let mut valid = FileInfo { parts: vec![ObjectPartInfo { number: 1, - checksums: Some(valid_checksums.clone()), + checksums: Some(valid_checksums), ..Default::default() }], ..Default::default() }; - assert!(prepare_tiered_data_movement_file_info_for(&mut valid.clone(), false).is_err()); + let mut compatible = valid.clone(); + prepare_tiered_data_movement_file_info_for(&mut compatible, false) + .expect("disabled sidecar writer should omit optional part checksums"); + assert!(compatible.parts.iter().all(|part| part.checksums.is_none())); + assert!(!rustfs_utils::http::contains_key_str(&compatible.metadata, SUFFIX_PART_CHECKSUMS)); prepare_tiered_data_movement_file_info_for(&mut valid, true).expect("valid legacy part checksums should be encoded"); assert_eq!( rustfs_utils::http::get_consistent_str(&valid.metadata, SUFFIX_PART_CHECKSUMS), @@ -2308,7 +2401,7 @@ mod tests { ..Default::default() }; - let opts = data_movement_complete_multipart_opts(&object_info, 7).expect("complete opts should encode metadata"); + let opts = data_movement_complete_multipart_opts(&object_info, 7, false).expect("complete opts should encode metadata"); assert!(opts.versioned); assert!(opts.data_movement); @@ -2387,7 +2480,7 @@ mod tests { let put_opts = data_movement_put_object_opts(&object_info, 9); let complete_opts = - data_movement_complete_multipart_opts(&object_info, 9).expect("complete opts should encode metadata"); + data_movement_complete_multipart_opts(&object_info, 9, false).expect("complete opts should encode metadata"); assert_eq!( put_opts @@ -2406,6 +2499,58 @@ mod tests { } } + #[test] + fn test_stale_data_movement_target_replacement_requires_exact_owned_generation() { + let version_id = Uuid::from_u128(41); + let source_time = OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(2); + let opts = ObjectOptions { + data_movement: true, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(source_time), + http_preconditions: Some(data_movement_target_precondition()), + ..Default::default() + }; + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, SUFFIX_DATA_MOVED, "true".to_string()); + let target = ObjectInfo { + version_id: Some(version_id), + mod_time: Some(source_time - time::Duration::SECOND), + user_defined: Arc::new(metadata), + ..Default::default() + }; + + assert!(can_replace_stale_data_movement_target(&target, &opts)); + + let mut client_target = target.clone(); + client_target.user_defined = Arc::new(HashMap::new()); + assert!(!can_replace_stale_data_movement_target(&client_target, &opts)); + + let mut single_marker = target.clone(); + Arc::make_mut(&mut single_marker.user_defined) + .remove(&format!("{}{SUFFIX_DATA_MOVED}", rustfs_utils::http::MINIO_INTERNAL_PREFIX)); + assert!(!can_replace_stale_data_movement_target(&single_marker, &opts)); + + let mut conflicting_marker = target.clone(); + Arc::make_mut(&mut conflicting_marker.user_defined).insert( + format!("{}{SUFFIX_DATA_MOVED}", rustfs_utils::http::MINIO_INTERNAL_PREFIX), + "false".to_string(), + ); + assert!(!can_replace_stale_data_movement_target(&conflicting_marker, &opts)); + + let mut different_version = target.clone(); + different_version.version_id = Some(Uuid::from_u128(42)); + assert!(!can_replace_stale_data_movement_target(&different_version, &opts)); + + let mut same_generation = target.clone(); + same_generation.mod_time = Some(source_time); + assert!(!can_replace_stale_data_movement_target(&same_generation, &opts)); + + let mut delete_marker = target; + delete_marker.delete_marker = true; + assert!(!can_replace_stale_data_movement_target(&delete_marker, &opts)); + } + #[test] fn test_is_equivalent_data_movement_object_accepts_matching_metadata() { let version_id = Uuid::nil(); @@ -2528,8 +2673,12 @@ mod tests { } fn overwrite_resume_for_target(source: &ObjectInfo, target: ObjectInfo) -> bool { + overwrite_resume_for_target_with_checksums(source, target, data_movement_part_checksum_writer_enabled()) + } + + fn overwrite_resume_for_target_with_checksums(source: &ObjectInfo, target: ObjectInfo, compare_part_checksums: bool) -> bool { let err = Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string()); - resolve_data_movement_overwrite_resume_result(&err, Ok(Some(target)), source, 0, 1) + resolve_data_movement_overwrite_resume_result_for(&err, Ok(Some(target)), source, 0, 1, compare_part_checksums) .expect("overwrite target should be evaluated") } @@ -2890,7 +3039,7 @@ mod tests { parts[0].checksums = None; target.parts = Arc::new(parts); - assert!(!overwrite_resume_for_target(&source, target)); + assert!(!overwrite_resume_for_target_with_checksums(&source, target, true)); } #[test] @@ -3032,6 +3181,27 @@ mod tests { assert!(should_resume); } + #[test] + fn test_overwrite_resume_omits_part_checksums_only_in_compatible_mode() { + let source = overwrite_equivalence_source(); + let mut target = source.clone(); + let mut target_parts = target.parts.as_ref().clone(); + for part in &mut target_parts { + part.checksums = None; + } + target.parts = Arc::new(target_parts); + let err = Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string()); + + assert!( + resolve_data_movement_overwrite_resume_result_for(&err, Ok(Some(target.clone())), &source, 0, 1, false) + .expect("compatible migration should accept an omitted optional checksum sidecar") + ); + assert!( + !resolve_data_movement_overwrite_resume_result_for(&err, Ok(Some(target)), &source, 0, 1, true) + .expect("fleet-confirmed migration should compare persisted part checksums") + ); + } + #[test] fn test_invalid_upload_accepts_versioned_target_taken_over_by_old_node() { let mut source = overwrite_equivalence_source(); diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 9aea35e58..764a68500 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -16,7 +16,7 @@ use super::meta::{ clone_arc_by_index, ensure_valid_rebalance_pool_index, invalid_rebalance_pool_index_error, rebalance_metadata_not_initialized_error, should_ignore_rebalance_data_usage_cache, }; -use super::migration::{RebalanceMigrationBackend, migrate_entry_version_with_incarnation}; +use super::migration::{RebalanceMigrationBackend, migrate_entry_version}; use super::worker::{ RebalanceEntryCleanupResult, RebalanceEntryTask, load_rebalance_bucket_configs, rebalance_max_attempts, resolve_rebalance_bucket_error, resolve_rebalance_entry_cleanup_delete_result, resolve_rebalance_file_info_versions_result, @@ -223,7 +223,7 @@ impl ECStore { let store = self.clone(); async move { store.delete_object(&bucket, &object, opts).await } }; - let result = migrate_entry_version_with_incarnation( + let result = migrate_entry_version( &RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()), bucket.clone(), pool_index, @@ -329,10 +329,12 @@ impl ECStore { entry.name.as_str(), &fivs, &cleanup_preflight_allowed_missing, - bucket_configs.bucket_incarnation_id, - bucket_incarnation_fence - .as_ref() - .and_then(|guard| guard.namespace_lock_guard()), + data_movement::SourceCleanupBucketFence { + expected_incarnation_id: bucket_configs.bucket_incarnation_id, + lifecycle_guard: bucket_incarnation_fence + .as_ref() + .and_then(|guard| guard.namespace_lock_guard()), + }, "rebalance", ), ) diff --git a/crates/ecstore/src/services/rebalance/migration.rs b/crates/ecstore/src/services/rebalance/migration.rs index 4859c2faf..7e23c9ed1 100644 --- a/crates/ecstore/src/services/rebalance/migration.rs +++ b/crates/ecstore/src/services/rebalance/migration.rs @@ -136,40 +136,6 @@ impl MigrationBackend for RebalanceMigrationBackend<'_> { #[allow(clippy::too_many_arguments)] pub(crate) async fn migrate_entry_version( - set: &Backend, - bucket: String, - pool_index: usize, - version: &FileInfo, - version_id: Option, - max_attempts: usize, - ignore_data_usage_cache: bool, - transfer: F, - delete_marker: D, -) -> MigrationVersionResult -where - Backend: MigrationBackend + ?Sized, - F: FnMut(usize, String, GetObjectReader) -> Fut + Send, - Fut: Future> + Send, - D: FnMut(String, String, ObjectOptions) -> DFut + Send, - DFut: Future> + Send, -{ - migrate_entry_version_with_incarnation( - set, - bucket, - pool_index, - version, - version_id, - None, - max_attempts, - ignore_data_usage_cache, - transfer, - delete_marker, - ) - .await -} - -#[allow(clippy::too_many_arguments)] -pub(crate) async fn migrate_entry_version_with_incarnation( set: &Backend, bucket: String, pool_index: usize, diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index dcceb3b03..38d4d7460 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -29,8 +29,8 @@ use super::meta::{ validate_init_rebalance_state, validate_start_rebalance_state, }; use super::migration::{ - MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_incarnation, - migrate_entry_version_with_retry_wait, rebalance_delete_marker_opts, + MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait, + rebalance_delete_marker_opts, }; use super::runtime::{should_fail_repeated_rebalance_bucket_defer, source_cleanup_defer_attempt}; use super::worker::{ @@ -285,7 +285,7 @@ async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() { }; let incarnation = uuid::Uuid::new_v4(); - let result = migrate_entry_version_with_incarnation( + let result = migrate_entry_version( &backend, "bucket".to_string(), 0, @@ -336,6 +336,7 @@ async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() { 0, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -372,6 +373,7 @@ async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() { 0, &version, Some("vid-1".to_string()), + None, 3, false, &mut transfer, @@ -410,6 +412,7 @@ async fn test_migrate_entry_version_remote_failure_is_reported() { 0, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -452,6 +455,7 @@ async fn test_migrate_entry_version_deleted_version_routes_delete_through_store_ 1, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -491,6 +495,7 @@ async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -533,6 +538,7 @@ async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() { 1, &version, Some("vid-1".to_string()), + None, 3, false, &mut transfer, @@ -562,6 +568,7 @@ async fn test_migrate_entry_version_reader_not_found_is_ignored() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -689,6 +696,7 @@ async fn test_migrate_entry_version_reader_fails_after_retries() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -727,6 +735,7 @@ async fn test_migrate_entry_version_zero_max_attempts_still_attempts_once() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 0, false, &mut transfer, @@ -871,6 +880,7 @@ async fn test_migrate_entry_version_transfer_fails_after_retries() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 2, false, &mut transfer, @@ -909,6 +919,7 @@ async fn test_migrate_entry_version_transfer_not_found_is_ignored() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 3, false, &mut transfer, @@ -950,6 +961,7 @@ async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() { 1, &version, Some("vid-1".to_string()), + None, 3, false, &mut transfer, @@ -992,6 +1004,7 @@ async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 2, true, &mut transfer, @@ -1034,6 +1047,7 @@ async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled( 1, &version, version.version_id.map(|v| v.to_string()), + None, 2, false, &mut transfer, @@ -2075,6 +2089,7 @@ async fn test_migrate_entry_version_transfer_failure_reports_write_target_stage( 1, &version, version.version_id.map(|v| v.to_string()), + None, 1, false, &mut transfer, @@ -2099,6 +2114,7 @@ async fn test_migrate_entry_version_reader_failure_reports_read_source_stage() { 1, &version, version.version_id.map(|v| v.to_string()), + None, 1, false, &mut transfer, diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 7b6be13b3..0b4a60f0b 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -4524,16 +4524,16 @@ impl SetDisks { object: &str, opts: &ObjectOptions, ) -> Option { - let mut opts = opts.clone(); + let mut lookup_opts = opts.clone(); - let http_preconditions = opts.http_preconditions?; - opts.http_preconditions = None; + let http_preconditions = lookup_opts.http_preconditions?; + lookup_opts.http_preconditions = None; // Never claim a lock here, to avoid deadlock // - If no_lock is false, we must have obtained the lock out side of this function // - If no_lock is true, we should not obtain locks - opts.no_lock = true; - let oi = self.get_object_info(bucket, object, &opts).await; + lookup_opts.no_lock = true; + let oi = self.get_object_info(bucket, object, &lookup_opts).await; match oi { Ok(oi) => { @@ -4544,7 +4544,9 @@ impl SetDisks { } let if_none_match = http_preconditions.if_none_match_value().map(str::to_owned); let if_match = http_preconditions.if_match_value().map(str::to_owned); - if should_prevent_write(&oi, if_none_match, if_match) { + if should_prevent_write(&oi, if_none_match, if_match) + && !crate::data_movement::can_replace_stale_data_movement_target(&oi, opts) + { return Some(StorageError::PreconditionFailed); } } diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 8fd6ec781..caf8e1e73 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2067,6 +2067,19 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { achieved: 0, }); } + if opts + .namespace_lock_fence + .as_ref() + .is_some_and(NamespaceLockFence::is_lock_lost) + { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "complete_multipart_upload_outer_lock", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } if upload_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { return Err(StorageError::NamespaceLockQuorumUnavailable { mode: "complete_multipart_upload_commit", @@ -2087,6 +2100,74 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { ) .await?; + if opts.data_movement + && opts.http_preconditions.is_some() + && let Some(err) = self.check_write_precondition(bucket, object, opts).await + { + return Err(err); + } + if opts.data_movement && opts.http_preconditions.is_some() && !crate::bucket::utils::is_meta_bucketname(bucket) { + let current = self + .get_object_info( + bucket, + object, + &ObjectOptions { + version_id: opts.version_id.clone(), + no_lock: true, + metadata_cache_safe: false, + versioned: opts.versioned, + version_suspended: opts.version_suspended, + ..Default::default() + }, + ) + .await; + match current { + Ok(existing) if crate::data_movement::can_replace_stale_data_movement_target(&existing, opts) => { + let object_lock_config = opts.object_lock_config_snapshot.as_deref().ok_or_else(|| { + Error::other("data movement completion is missing its Object Lock configuration snapshot") + })?; + if check_object_lock_for_deletion_with_state(object_lock_config.state(), &existing, false)?.is_some() { + return Err(StorageError::PrefixAccessDenied(bucket.to_string(), object.to_string())); + } + } + Ok(_) => return Err(StorageError::PreconditionFailed), + Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {} + Err(err) => return Err(err), + } + } + if object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "complete_multipart_upload_commit", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } + if opts + .namespace_lock_fence + .as_ref() + .is_some_and(NamespaceLockFence::is_lock_lost) + { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "complete_multipart_upload_outer_lock", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } + if upload_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "complete_multipart_upload_commit", + bucket: RUSTFS_META_MULTIPART_BUCKET.to_string(), + object: upload_id_path.clone(), + required: 1, + achieved: 0, + }); + } + ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?; + let complete_tail_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); // Crash-consistency injection: hard power loss after the upload is fully @@ -2878,6 +2959,92 @@ mod tests { ); } + #[tokio::test] + async fn stale_data_movement_replacement_fails_before_commit_on_outer_fence_loss() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "data-movement-stale-fence-bucket"; + make_bucket_on_all(&disk_stores, bucket).await; + let old_time = OffsetDateTime::UNIX_EPOCH + time::Duration::SECOND; + let new_time = old_time + time::Duration::SECOND; + + for (object, namespace_lock_fence, bucket_lifecycle_lock_fence) in [ + ("metadata-fence", Some(NamespaceLockFence::lost_for_test()), None), + ("bucket-fence", None, Some(NamespaceLockFence::lost_for_test())), + ] { + let version_id = Uuid::new_v4(); + let old_body = format!("old-{object}").into_bytes(); + let mut old_reader = PutObjReader::from_vec(old_body.clone()); + set_disks + .put_object( + bucket, + object, + &mut old_reader, + &ObjectOptions { + data_movement: true, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(old_time), + ..Default::default() + }, + ) + .await + .expect("seed old data movement target"); + + let replacement_body = format!("new-{object}").into_bytes(); + let create_opts = ObjectOptions { + data_movement: true, + ..Default::default() + }; + let (upload_id, parts) = + stage_upload_with_create_opts(&set_disks, bucket, object, &replacement_body, &create_opts).await; + let complete_opts = ObjectOptions { + data_movement: true, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(new_time), + http_preconditions: Some(crate::data_movement::data_movement_target_precondition()), + namespace_lock_fence, + bucket_lifecycle_lock_fence, + object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new( + ObjectLockConfigState::ConfirmedAbsent, + ))), + ..Default::default() + }; + let err = set_disks + .clone() + .complete_multipart_upload(bucket, object, &upload_id, parts, &complete_opts) + .await + .expect_err("a lost outer fence must abort stale target replacement"); + assert!(matches!(err, StorageError::NamespaceLockQuorumUnavailable { .. })); + + let mut preserved = set_disks + .get_object_reader( + bucket, + object, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read target after rejected replacement"); + let mut preserved_body = Vec::new(); + preserved + .stream + .read_to_end(&mut preserved_body) + .await + .expect("drain target after rejected replacement"); + assert_eq!(preserved_body, old_body); + set_disks + .check_upload_id_exists(bucket, object, &upload_id, true) + .await + .expect("fence loss must leave replacement staging retryable"); + } + } + #[tokio::test] #[serial] async fn data_movement_complete_accepts_unknown_compressed_part_actual_size() { diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 282eb2841..87beb844c 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -7998,8 +7998,7 @@ mod transition_upload_integrity_tests { object, &expected, &[], - None, - None, + crate::data_movement::SourceCleanupBucketFence::default(), "test_data_movement", ) .await @@ -8061,8 +8060,10 @@ mod transition_upload_integrity_tests { object, &expected, &[], - None, - Some(&bucket_guard), + crate::data_movement::SourceCleanupBucketFence { + expected_incarnation_id: None, + lifecycle_guard: Some(&bucket_guard), + }, "test_data_movement", ) .await; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 4e959b15b..bf493a137 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -1706,6 +1706,106 @@ mod tests { .contains_key(rustfs_rio::RUSTFS_MULTIPART_CHECKSUM_TYPE) ); retry_reader.object_info = retry_source_info.clone(); + temp_env::async_with_vars( + [ + (rustfs_config::ENV_DATA_MOVEMENT_PART_CHECKSUMS_WRITE, None::<&str>), + (rustfs_config::ENV_DATA_MOVEMENT_PART_CHECKSUMS_FLEET_CONFIRMED, None::<&str>), + ], + crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + retry_reader, + None, + "test_data_movement_compatible_retry", + ), + ) + .await + .expect("default multipart migration should not require the checksum sidecar capability"); + + let compatible_target = store.pools[1] + .get_object_info( + &bucket, + retry_object, + &ObjectOptions { + include_part_checksums: true, + ..Default::default() + }, + ) + .await + .expect("read compatible multipart target metadata"); + assert_eq!(compatible_target.checksum, retry_source_info.checksum); + assert!(compatible_target.parts.iter().all(|part| part.checksums.is_none())); + assert!(!rustfs_utils::http::contains_key_str( + &compatible_target.user_defined, + rustfs_utils::http::SUFFIX_PART_CHECKSUMS + )); + let mut compatible_reader = store.pools[1] + .get_object_reader(&bucket, retry_object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("read compatible multipart target body"); + let mut compatible_body = Vec::new(); + compatible_reader + .stream + .read_to_end(&mut compatible_body) + .await + .expect("drain compatible multipart target body"); + assert_eq!(compatible_body, retry_source_body); + + let mut compatible_retry_reader = store.pools[0] + .get_object_reader( + &bucket, + retry_object, + None, + HeaderMap::new(), + &ObjectOptions { + raw_data_movement_read: true, + ..Default::default() + }, + ) + .await + .expect("read source multipart object for compatible retry"); + compatible_retry_reader.object_info = retry_source_info.clone(); + temp_env::async_with_vars( + [ + (rustfs_config::ENV_DATA_MOVEMENT_PART_CHECKSUMS_WRITE, None::<&str>), + (rustfs_config::ENV_DATA_MOVEMENT_PART_CHECKSUMS_FLEET_CONFIRMED, None::<&str>), + ], + crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + compatible_retry_reader, + None, + "test_data_movement_compatible_retry", + ), + ) + .await + .expect("compatible multipart migration retry should converge"); + let compatible_uploads = store.pools[1] + .list_multipart_uploads(&bucket, retry_object, None, None, None, 100) + .await + .expect("list compatible target multipart uploads"); + assert!(compatible_uploads.uploads.is_empty(), "compatible retry staging must be aborted"); + + store.pools[1] + .delete_object(&bucket, retry_object, ObjectOptions::default()) + .await + .expect("remove compatible target before fleet-confirmed migration"); + let mut retry_reader = store.pools[0] + .get_object_reader( + &bucket, + retry_object, + None, + HeaderMap::new(), + &ObjectOptions { + raw_data_movement_read: true, + ..Default::default() + }, + ) + .await + .expect("read source multipart object for fleet-confirmed migration"); + retry_reader.object_info = retry_source_info.clone(); temp_env::async_with_vars( [ (rustfs_config::ENV_DATA_MOVEMENT_PART_CHECKSUMS_WRITE, Some("true")), @@ -1921,6 +2021,301 @@ mod tests { .expect("join data movement test thread"); } + #[test] + #[serial_test::serial(storage_class_env)] + fn data_movement_multipart_replaces_only_unlocked_owned_generation() { + std::thread::Builder::new() + .stack_size(16 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(2) + .build() + .expect("build data movement replacement test runtime"); + runtime.block_on(async move { + let temp_dir = tempfile::tempdir().expect("create data movement replacement store dir"); + let (_ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "data-movement-stale-replacement", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("dm-stale-replacement-{}", uuid::Uuid::new_v4()); + store + .make_bucket( + &bucket, + &MakeBucketOptions { + lock_enabled: true, + ..Default::default() + }, + ) + .await + .expect("create data movement replacement bucket"); + *store.rebalance_meta.write().await = Some(active_rebalance_meta_for_pool(store.pools.len(), 0)); + + let version_id = uuid::Uuid::new_v4(); + let old_time = OffsetDateTime::UNIX_EPOCH + time::Duration::SECOND; + let new_time = old_time + time::Duration::SECOND; + let first_part_size = 5 * 1024 * 1024; + let mut old_body = vec![b'a'; first_part_size]; + old_body.push(b'b'); + let mut new_body = vec![b'c'; first_part_size]; + new_body.push(b'd'); + + let source_reader = |name: &str, body: Vec, mod_time, metadata: HashMap| { + let size = i64::try_from(body.len()).expect("source body size should fit i64"); + GetObjectReader { + stream: Box::new(Cursor::new(body)), + object_info: ObjectInfo { + bucket: bucket.clone(), + name: name.to_string(), + version_id: Some(version_id), + size, + actual_size: size, + etag: Some(format!("{name}-multipart-etag-2")), + mod_time: Some(mod_time), + user_defined: Arc::new(metadata), + parts: Arc::new(vec![ + ObjectPartInfo { + number: 1, + size: first_part_size, + actual_size: i64::try_from(first_part_size).expect("first part size should fit i64"), + etag: format!("{name}-part-1"), + ..Default::default() + }, + ObjectPartInfo { + number: 2, + size: 1, + actual_size: 1, + etag: format!("{name}-part-2"), + ..Default::default() + }, + ]), + ..Default::default() + }, + buffered_body: None, + body_source: Default::default(), + } + }; + + let replaceable = "replaceable.bin"; + crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + source_reader( + replaceable, + old_body.clone(), + old_time, + HashMap::from([("x-amz-meta-generation".to_string(), "old".to_string())]), + ), + None, + "test_stale_target_seed", + ) + .await + .expect("seed old migrated target"); + let seeded_target = store.pools[1] + .get_object_info( + &bucket, + replaceable, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read old migrated target"); + let replacement_opts = ObjectOptions { + data_movement: true, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(new_time), + http_preconditions: Some(crate::data_movement::data_movement_target_precondition()), + ..Default::default() + }; + assert!( + crate::data_movement::can_replace_stale_data_movement_target(&seeded_target, &replacement_opts), + "seeded target should be replaceable: {seeded_target:?}" + ); + crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + source_reader( + replaceable, + new_body.clone(), + new_time, + HashMap::from([("x-amz-meta-generation".to_string(), "new".to_string())]), + ), + None, + "test_stale_target_replace", + ) + .await + .expect("newer source generation should replace the old migrated target"); + + let replacement = store.pools[1] + .get_object_info( + &bucket, + replaceable, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read replaced target metadata"); + assert_eq!(replacement.mod_time, Some(new_time)); + assert_eq!(replacement.user_defined.get("x-amz-meta-generation").map(String::as_str), Some("new")); + let mut replacement_reader = store.pools[1] + .get_object_reader( + &bucket, + replaceable, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read replaced target body"); + let mut replacement_body = Vec::new(); + replacement_reader + .stream + .read_to_end(&mut replacement_body) + .await + .expect("drain replaced target body"); + assert_eq!(replacement_body, new_body); + + let client_target = "client-target.bin"; + let client_body = b"client-owned exact version".to_vec(); + let mut client_reader = PutObjReader::from_vec(client_body.clone()); + store.pools[1] + .put_object( + &bucket, + client_target, + &mut client_reader, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(old_time), + user_defined: HashMap::from([("x-amz-meta-owner".to_string(), "client".to_string())]), + ..Default::default() + }, + ) + .await + .expect("seed client-owned exact version"); + let err = crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + source_reader(client_target, new_body.clone(), new_time, HashMap::new()), + None, + "test_client_target_reject", + ) + .await + .expect_err("data movement must not replace a client-owned exact version"); + assert!(err.to_string().contains("complete_multipart_upload"), "unexpected migration error: {err}"); + let mut preserved_reader = store.pools[1] + .get_object_reader( + &bucket, + client_target, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read preserved client target"); + let mut preserved_body = Vec::new(); + preserved_reader + .stream + .read_to_end(&mut preserved_body) + .await + .expect("drain preserved client target"); + assert_eq!(preserved_body, client_body); + + let retain_until = (OffsetDateTime::now_utc() + time::Duration::days(1)) + .format(&time::format_description::well_known::Rfc3339) + .expect("retain-until date should format"); + for (object, mode) in [ + ("compliance-target.bin", s3s::dto::ObjectLockRetentionMode::COMPLIANCE), + ("governance-target.bin", s3s::dto::ObjectLockRetentionMode::GOVERNANCE), + ] { + let retained_metadata = HashMap::from([ + (s3s::header::X_AMZ_OBJECT_LOCK_MODE.as_str().to_string(), mode.to_string()), + ( + s3s::header::X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE.as_str().to_string(), + retain_until.clone(), + ), + ]); + crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + source_reader(object, old_body.clone(), old_time, retained_metadata.clone()), + None, + "test_retained_target_seed", + ) + .await + .expect("seed retained migrated target"); + let err = crate::data_movement::migrate_object( + store.clone(), + 0, + bucket.clone(), + source_reader(object, new_body.clone(), new_time, retained_metadata), + None, + "test_retained_target_replace", + ) + .await + .expect_err("active retention must block stale target replacement"); + assert!(err.to_string().contains("complete_multipart_upload"), "unexpected retention error: {err}"); + + let mut retained_reader = store.pools[1] + .get_object_reader( + &bucket, + object, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("read retained target"); + let mut retained_body = Vec::new(); + retained_reader + .stream + .read_to_end(&mut retained_body) + .await + .expect("drain retained target"); + assert_eq!(retained_body, old_body); + } + + for object in [replaceable, client_target, "compliance-target.bin", "governance-target.bin"] { + let uploads = store.pools[1] + .list_multipart_uploads(&bucket, object, None, None, None, 100) + .await + .expect("list target multipart staging"); + assert!(uploads.uploads.is_empty(), "data movement staging must be cleaned for {object}"); + } + }); + }) + .expect("spawn data movement replacement test thread") + .join() + .expect("join data movement replacement test thread"); + } + #[tokio::test] #[serial_test::serial(storage_class_env)] async fn data_movement_delete_marker_retries_converge_safely() { @@ -2275,12 +2670,7 @@ mod tests { let mut source_body = vec![b'a'; first_part_size]; source_body.push(b'b'); let source_size = i64::try_from(source_body.len()).expect("multipart source size should fit i64"); - let guard_barrier = crate::store::multipart::DataMovementMultipartGuardBarrier::install(&bucket, 2); - let part_barrier = crate::set_disk::MultipartCommitBarrier::install( - &bucket, - object, - crate::set_disk::MultipartCommitPause::PutPartAfterRename, - ); + let completion_barrier = crate::store::multipart::DataMovementMultipartCompletionBarrier::install(&bucket); let migration_store = store.clone(); let migration_bucket = bucket.clone(); let migration = tokio::spawn(async move { @@ -2325,7 +2715,7 @@ mod tests { ) .await }); - guard_barrier.wait_until_paused().await; + completion_barrier.wait_until_paused().await; let unrelated_set = store.pools[1].get_disks_by_key(object); let original_disks = { let mut disks = unrelated_set.disks.write().await; @@ -2335,9 +2725,7 @@ mod tests { } original }; - drop(guard_barrier); - part_barrier.wait_until_paused().await; - drop(part_barrier); + drop(completion_barrier); let err = migration .await .expect("multipart migration task should join") diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index 0893d4d4d..3d460d4d5 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -67,40 +67,35 @@ fn ensure_multipart_bucket_lifecycle_guard_held( } #[cfg(test)] -struct DataMovementMultipartGuardBarrierState { +struct DataMovementMultipartCompletionBarrierState { bucket: String, - pause_at_arrival: usize, - arrivals: std::sync::atomic::AtomicUsize, arrived: tokio::sync::Notify, release: tokio::sync::Notify, } #[cfg(test)] -pub(crate) struct DataMovementMultipartGuardBarrier { - state: Arc, +pub(crate) struct DataMovementMultipartCompletionBarrier { + state: Arc, } #[cfg(test)] -static DATA_MOVEMENT_MULTIPART_GUARD_BARRIER: std::sync::OnceLock< - std::sync::Mutex>>, +static DATA_MOVEMENT_MULTIPART_COMPLETION_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, > = std::sync::OnceLock::new(); #[cfg(test)] -impl DataMovementMultipartGuardBarrier { - pub(crate) fn install(bucket: &str, pause_at_arrival: usize) -> Self { - assert!(pause_at_arrival > 0, "data movement multipart guard barrier needs an arrival"); - let state = Arc::new(DataMovementMultipartGuardBarrierState { +impl DataMovementMultipartCompletionBarrier { + pub(crate) fn install(bucket: &str) -> Self { + let state = Arc::new(DataMovementMultipartCompletionBarrierState { bucket: bucket.to_string(), - pause_at_arrival, - arrivals: std::sync::atomic::AtomicUsize::new(0), arrived: tokio::sync::Notify::new(), release: tokio::sync::Notify::new(), }); - let mut slot = DATA_MOVEMENT_MULTIPART_GUARD_BARRIER + let mut slot = DATA_MOVEMENT_MULTIPART_COMPLETION_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) .lock() - .expect("data movement multipart guard barrier mutex should not poison"); - assert!(slot.is_none(), "data movement multipart guard barrier must be unique"); + .expect("data movement multipart completion barrier mutex should not poison"); + assert!(slot.is_none(), "data movement multipart completion barrier must be unique"); *slot = Some(Arc::clone(&state)); Self { state } } @@ -108,18 +103,18 @@ impl DataMovementMultipartGuardBarrier { pub(crate) async fn wait_until_paused(&self) { tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified()) .await - .expect("data movement multipart operation should pass the bucket guard"); + .expect("data movement multipart operation should reach selected completion"); } } #[cfg(test)] -impl Drop for DataMovementMultipartGuardBarrier { +impl Drop for DataMovementMultipartCompletionBarrier { fn drop(&mut self) { self.state.release.notify_one(); - let mut slot = DATA_MOVEMENT_MULTIPART_GUARD_BARRIER + let mut slot = DATA_MOVEMENT_MULTIPART_COMPLETION_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) .lock() - .expect("data movement multipart guard barrier mutex should not poison"); + .expect("data movement multipart completion barrier mutex should not poison"); if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { *slot = None; } @@ -127,22 +122,15 @@ impl Drop for DataMovementMultipartGuardBarrier { } #[cfg(test)] -async fn pause_data_movement_multipart_after_bucket_guard(bucket: &str, opts: &ObjectOptions) { - if !opts.data_movement { - return; - } - let barrier = DATA_MOVEMENT_MULTIPART_GUARD_BARRIER +async fn pause_data_movement_multipart_before_selected_completion(bucket: &str) { + let barrier = DATA_MOVEMENT_MULTIPART_COMPLETION_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) .lock() - .expect("data movement multipart guard barrier mutex should not poison") + .expect("data movement multipart completion barrier mutex should not poison") .as_ref() .filter(|barrier| barrier.bucket == bucket) .cloned(); if let Some(barrier) = barrier { - let arrival = barrier.arrivals.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1; - if arrival != barrier.pause_at_arrival { - return; - } barrier.arrived.notify_one(); barrier.release.notified().await; } @@ -259,8 +247,6 @@ impl ECStore { if opts.expected_bucket_incarnation_id != Some(current) { return Err(StorageError::BucketNotFound(bucket.to_string())); } - #[cfg(test)] - Box::pin(pause_data_movement_multipart_after_bucket_guard(bucket, &opts)).await; Ok((opts, guard)) } @@ -559,10 +545,12 @@ impl ECStore { bucket: &str, object: &str, upload_id: &str, - part_id: usize, data: &mut PutObjReader, opts: &ObjectOptions, ) -> Result { + let part_id = opts + .part_number + .ok_or_else(|| Error::other("targeted multipart upload requires a part number"))?; check_put_object_part_args(bucket, object, upload_id)?; if !opts.data_movement { return Err(Error::other("targeted multipart upload requires data_movement options")); @@ -727,7 +715,32 @@ impl ECStore { if !opts.data_movement { return Err(Error::other("targeted multipart completion requires data_movement options")); } - let (opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?; + let (mut opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?; + if opts.overwrites_existing_version() && !is_meta_bucketname(bucket) { + let expected_incarnation_id = opts + .expected_bucket_incarnation_id + .ok_or_else(|| Error::other("data movement completion is missing its bucket incarnation"))?; + let lifecycle_fence = opts + .bucket_lifecycle_lock_fence + .as_ref() + .ok_or_else(|| Error::other("data movement completion is missing its bucket lifecycle fence"))?; + let snapshot = match opts.object_lock_config_snapshot.as_ref() { + Some(snapshot) => Arc::clone(snapshot), + None => { + self.object_lock_config_snapshot_under_lifecycle_fence(bucket, lifecycle_fence) + .await? + } + }; + if !snapshot.is_valid_for_destructive_put(self.id, bucket, expected_incarnation_id) { + return Err(Error::other( + "data movement Object Lock snapshot does not match the target bucket generation", + )); + } + snapshot.add_lock_fences(&mut opts); + opts.object_lock_config_snapshot = Some(snapshot); + } + #[cfg(test)] + pause_data_movement_multipart_before_selected_completion(bucket).await; let pool = self .pools .get(target_pool_idx) diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 61395c654..0d5df89c2 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1216,7 +1216,7 @@ impl ECStore { ))) } - async fn object_lock_config_snapshot_under_lifecycle_fence( + pub(super) async fn object_lock_config_snapshot_under_lifecycle_fence( &self, bucket: &str, lifecycle_fence: &NamespaceLockFence, @@ -3569,7 +3569,7 @@ mod tests { ReplicationStatusType::Replica.to_string(), ); rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_REPLICA_TIMESTAMP, timestamp.clone()); - rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_REPLICATION_TIMESTAMP, timestamp.clone()); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_REPLICATION_TIMESTAMP, timestamp); rustfs_utils::http::insert_str( &mut metadata, rustfs_utils::http::SUFFIX_REPLICATION_STATUS, diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index a79193058..d2569447c 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -729,7 +729,7 @@ impl FileMeta { include_free_versions: bool, all_parts: bool, ) -> Result { - self.into_fileinfo_with_part_checksums( + self.to_fileinfo_with_part_checksums( volume, path, version_id, @@ -750,7 +750,7 @@ impl FileMeta { read_data: bool, include_free_versions: bool, ) -> Result { - self.into_fileinfo_with_part_checksums( + self.to_fileinfo_with_part_checksums( volume, path, version_id, @@ -763,7 +763,7 @@ impl FileMeta { ) } - fn into_fileinfo_with_part_checksums( + fn to_fileinfo_with_part_checksums( &self, volume: &str, path: &str, @@ -802,12 +802,8 @@ impl FileMeta { // Known side effect: if a disk holds only free versions and they are // corrupt, `into_fileinfo` falls through to `FileNotFound` (not // `FileCorrupt`), so that disk is not enqueued for heal. - match found_free_fi.into_fileinfo_with_part_checksums( - volume, - path, - opts.all_parts, - opts.include_part_checksums, - ) { + match found_free_fi.to_fileinfo_with_part_checksums(volume, path, opts.all_parts, opts.include_part_checksums) + { Ok(mut free_fi) => { free_fi.is_latest = true; found_free_version = Some(free_fi); @@ -835,7 +831,7 @@ impl FileMeta { found = true; - let mut fi = ver.into_fileinfo_with_part_checksums(volume, path, opts.all_parts, opts.include_part_checksums)?; + let mut fi = ver.to_fileinfo_with_part_checksums(volume, path, opts.all_parts, opts.include_part_checksums)?; fi.is_latest = is_latest; if let Some(_d) = succ_mod_time { diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index 462a84c9e..37c22b125 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -273,9 +273,7 @@ fn transition_version_state_from_bytes(value: Option<&[u8]>) -> Result, state: TransitionVersionState) -> Option { - let Some(value) = value else { - return None; - }; + let value = value?; if value.is_empty() { return None; } @@ -290,7 +288,7 @@ fn transitioned_version_from_bytes(value: Option<&[u8]>, state: TransitionVersio if value.is_empty() || value.len() > MAX_TRANSITION_VERSION_LEN || value.chars().any(char::is_control) - || Uuid::parse_str(&value).is_ok_and(|id| id.is_nil()) + || Uuid::parse_str(value).is_ok_and(|id| id.is_nil()) { None } else { @@ -320,34 +318,62 @@ struct DerivedInternalMetadata<'a> { impl<'a> DerivedInternalMetadata<'a> { fn from_meta_sys(meta_sys: &'a HashMap>) -> Result { - let mut metadata = Self::default(); + let mut canonical = Self::default(); + let mut legacy = Self::default(); for (key, value) in meta_sys { let Some(suffix) = rustfs_utils::http::strip_internal_prefix_preserving_case(key) else { continue; }; - let slot = if suffix.eq_ignore_ascii_case(SUFFIX_CRC) { - &mut metadata.checksum + let (canonical_slot, legacy_slot, expected_suffix) = if suffix.eq_ignore_ascii_case(SUFFIX_CRC) { + (&mut canonical.checksum, &mut legacy.checksum, SUFFIX_CRC) } else if suffix.eq_ignore_ascii_case(SUFFIX_PART_CHECKSUMS) { - &mut metadata.part_checksums + (&mut canonical.part_checksums, &mut legacy.part_checksums, SUFFIX_PART_CHECKSUMS) } else if suffix.eq_ignore_ascii_case(SUFFIX_TRANSITION_STATUS) { - &mut metadata.transition_status + (&mut canonical.transition_status, &mut legacy.transition_status, SUFFIX_TRANSITION_STATUS) } else if suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_OBJECTNAME) { - &mut metadata.transitioned_object + ( + &mut canonical.transitioned_object, + &mut legacy.transitioned_object, + SUFFIX_TRANSITIONED_OBJECTNAME, + ) } else if suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID) { - &mut metadata.transitioned_version + ( + &mut canonical.transitioned_version, + &mut legacy.transitioned_version, + SUFFIX_TRANSITIONED_VERSION_ID, + ) } else if suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_STATE) { - &mut metadata.transitioned_version_state + ( + &mut canonical.transitioned_version_state, + &mut legacy.transitioned_version_state, + SUFFIX_TRANSITIONED_VERSION_STATE, + ) } else if suffix.eq_ignore_ascii_case(SUFFIX_TRANSITION_TIER) { - &mut metadata.transition_tier + (&mut canonical.transition_tier, &mut legacy.transition_tier, SUFFIX_TRANSITION_TIER) } else { continue; }; + let slot = if suffix == expected_suffix + && (key.starts_with(RUSTFS_INTERNAL_PREFIX) || key.starts_with(rustfs_utils::http::MINIO_INTERNAL_PREFIX)) + { + canonical_slot + } else { + legacy_slot + }; if slot.is_some_and(|current| current != value.as_slice()) { return Err(Error::FileCorrupt); } *slot = Some(value.as_slice()); } - Ok(metadata) + Ok(Self { + checksum: canonical.checksum.or(legacy.checksum), + part_checksums: canonical.part_checksums.or(legacy.part_checksums), + transition_status: canonical.transition_status.or(legacy.transition_status), + transitioned_object: canonical.transitioned_object.or(legacy.transitioned_object), + transitioned_version: canonical.transitioned_version.or(legacy.transitioned_version), + transitioned_version_state: canonical.transitioned_version_state.or(legacy.transitioned_version_state), + transition_tier: canonical.transition_tier.or(legacy.transition_tier), + }) } } @@ -531,7 +557,7 @@ impl FileMetaShallowVersion { self.parse_version_meta()?.into_fileinfo(volume, path, all_parts) } - pub(super) fn into_fileinfo_with_part_checksums( + pub(super) fn to_fileinfo_with_part_checksums( &self, volume: &str, path: &str, @@ -539,7 +565,7 @@ impl FileMetaShallowVersion { include_part_checksums: bool, ) -> Result { self.parse_version_meta()? - .into_fileinfo_with_part_checksums(volume, path, all_parts, include_part_checksums) + .to_fileinfo_with_part_checksums(volume, path, all_parts, include_part_checksums) } } @@ -862,18 +888,16 @@ impl FileMetaVersion { } pub fn into_fileinfo(&self, volume: &str, path: &str, all_parts: bool) -> Result { - self.into_fileinfo_with_part_checksums(volume, path, all_parts, true) + self.to_fileinfo_with_part_checksums(volume, path, all_parts, true) } - pub(super) fn into_fileinfo_with_part_checksums( + pub(super) fn to_fileinfo_with_part_checksums( &self, volume: &str, path: &str, all_parts: bool, include_part_checksums: bool, ) -> Result { - // Only the Object arm carries part arrays and can fail the length guard; the - // Legacy and Delete arms have no part arrays and stay infallible. let mut fi = match self.version_type { VersionType::Invalid | VersionType::Legacy => { if let Some(ref legacy) = self.legacy_object { @@ -891,14 +915,14 @@ impl FileMetaVersion { self.object .as_ref() .unwrap_or(&default_object) - .into_fileinfo_with_part_checksums(volume, path, all_parts, include_part_checksums)? + .to_fileinfo_with_part_checksums(volume, path, all_parts, include_part_checksums)? } VersionType::Delete => { let default_marker = MetaDeleteMarker::default(); self.delete_marker .as_ref() .unwrap_or(&default_marker) - .into_fileinfo(volume, path, all_parts) + .into_fileinfo(volume, path, all_parts)? } }; fi.uses_legacy_checksum = self.uses_legacy_checksum; @@ -2493,10 +2517,10 @@ impl MetaObject { } pub fn into_fileinfo(&self, volume: &str, path: &str, all_parts: bool) -> Result { - self.into_fileinfo_with_part_checksums(volume, path, all_parts, true) + self.to_fileinfo_with_part_checksums(volume, path, all_parts, true) } - fn into_fileinfo_with_part_checksums( + fn to_fileinfo_with_part_checksums( &self, volume: &str, path: &str, @@ -2908,7 +2932,7 @@ impl MetaDeleteMarker { contains_key_bytes(&self.meta_sys, SUFFIX_FREE_VERSION) } - pub fn into_fileinfo(&self, volume: &str, path: &str, _all_parts: bool) -> FileInfo { + pub fn into_fileinfo(&self, volume: &str, path: &str, _all_parts: bool) -> Result { let metadata = self .meta_sys .clone() @@ -2930,33 +2954,27 @@ impl MetaDeleteMarker { if self.free_version() { fi.set_tier_free_version(); - if let Ok(derived_metadata) = DerivedInternalMetadata::from_meta_sys(&self.meta_sys) { - fi.transition_tier = derived_metadata - .transition_tier - .filter(|value| !value.is_empty()) - .map(|value| String::from_utf8_lossy(value).to_string()) - .unwrap_or_default(); - fi.transitioned_objname = derived_metadata - .transitioned_object - .filter(|value| !value.is_empty()) - .map(|value| String::from_utf8_lossy(value).to_string()) - .unwrap_or_default(); - fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state) - .unwrap_or(TransitionVersionState::Unknown); - fi.transition_version = - transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state); - fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); - if derived_metadata.transitioned_version_state.is_some() - && validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref()).is_err() - { - fi.transition_version = None; - fi.transition_version_id = None; - fi.transition_version_state = TransitionVersionState::Unknown; - } + let derived_metadata = DerivedInternalMetadata::from_meta_sys(&self.meta_sys)?; + fi.transition_tier = derived_metadata + .transition_tier + .filter(|value| !value.is_empty()) + .map(|value| String::from_utf8_lossy(value).to_string()) + .unwrap_or_default(); + fi.transitioned_objname = derived_metadata + .transitioned_object + .filter(|value| !value.is_empty()) + .map(|value| String::from_utf8_lossy(value).to_string()) + .unwrap_or_default(); + fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state)?; + fi.transition_version = + transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state); + fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); + if derived_metadata.transitioned_version_state.is_some() { + validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref())?; } } - fi + Ok(fi) } pub fn encode_to(&self, wr: &mut W) -> Result<()> { @@ -3744,7 +3762,9 @@ mod tests { .is_some_and(|state| state.target_delete_marker_version_ids_corrupt) ); - let roundtrip = MetaDeleteMarker::from(conflicting).into_fileinfo("bucket", "object", false); + let roundtrip = MetaDeleteMarker::from(conflicting) + .into_fileinfo("bucket", "object", false) + .expect("dynamic replication aliases should remain decodable"); assert!( roundtrip .replication_state_internal @@ -3904,6 +3924,22 @@ mod tests { assert_eq!(file_info.checksum.as_deref(), Some(checksum.as_slice())); } + #[test] + fn into_fileinfo_prefers_canonical_rewrite_over_stale_mixed_case_alias() { + let checksum = vec![0xff, 0x00, 0x80, 0x01]; + let mut object = MetaObject { + meta_sys: HashMap::from([("X-Minio-Internal-crc".to_string(), b"stale".to_vec())]), + ..Default::default() + }; + insert_bytes(&mut object.meta_sys, SUFFIX_CRC, checksum.clone()); + + let file_info = object + .into_fileinfo("bucket", "key", false) + .expect("canonical rewrites should supersede legacy mixed-case aliases"); + + assert_eq!(file_info.checksum.as_deref(), Some(checksum.as_slice())); + } + #[test] fn into_fileinfo_recovers_data_movement_part_checksums() { let mut object = object_with_parts(vec![1], vec![16], vec![16]); @@ -3923,7 +3959,7 @@ mod tests { ); let mut deferred = object - .into_fileinfo_with_part_checksums("bucket", "key", true, false) + .to_fileinfo_with_part_checksums("bucket", "key", true, false) .expect("quorum candidates should retain raw checksum metadata"); assert!(deferred.parts[0].checksums.is_none()); deferred @@ -4613,7 +4649,7 @@ mod tests { ); let persisted_version = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID); assert_eq!( - transitioned_version_from_bytes(persisted_version.as_deref(), TransitionVersionState::Unknown) + transitioned_version_from_bytes(persisted_version, TransitionVersionState::Unknown) .and_then(|value| Uuid::parse_str(&value).ok()), Some(id), "UUID exact writes must remain readable by the legacy UUID consumer" @@ -4737,7 +4773,8 @@ mod tests { mod_time: None, meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect("nil tier version should remain an absent remote version"); assert_eq!(fi.transition_version_id, None); } @@ -4752,7 +4789,8 @@ mod tests { mod_time: None, meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect("legacy binary UUID tier version should decode"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); } @@ -4769,7 +4807,8 @@ mod tests { mod_time: Some(sample_mod_time()), meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect("opaque tier version should remain readable"); assert_eq!(fi.transition_version_id, None); assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42")); @@ -4791,7 +4830,8 @@ mod tests { mod_time: Some(sample_mod_time()), meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect("mixed-case tier aliases should decode"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); @@ -4817,7 +4857,8 @@ mod tests { mod_time: Some(sample_mod_time()), meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect("mixed-case tier aliases should decode"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); @@ -4839,17 +4880,15 @@ mod tests { b"target-version".to_vec(), ); - let fi = MetaDeleteMarker { + let err = MetaDeleteMarker { version_id: Some(sample_version_id()), mod_time: Some(sample_mod_time()), meta_sys: sys, } - .into_fileinfo("b", "k", false); + .into_fileinfo("b", "k", false) + .expect_err("conflicting transition aliases must fail closed"); - assert_eq!(fi.transition_version, None); - assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); - assert!(fi.transition_tier.is_empty()); - assert!(fi.transitioned_objname.is_empty()); + assert_eq!(err, Error::FileCorrupt); } #[test]