diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 5112e066a..4f5d402d9 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -3911,6 +3911,7 @@ impl ECStore { if migrated { decommissioned += 1; + cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version)); free_version_disposition.record_migrated(); } else { free_version_disposition.record_retained(); @@ -6530,15 +6531,13 @@ mod pools_tests { use super::resolve_decommission_listing_error; use super::{ DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP, - DECOMMISSION_FREE_VERSION_MIGRATED_REASON, DECOMMISSION_FREE_VERSION_RETAINED_REASON, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler, - DecommissionEntryEnqueueResult, DecommissionFreeVersionDisposition, DecommissionStartPoolState, - DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, - QueuedDecommissionEntry, apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers, - bind_missing_decommission_cancelers, cancel_decommission_canceler, clamp_decommission_entry_concurrency, - classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, - decommission_entry_queue_capacity, decommission_item_size, decommission_meta_bucket_options, - decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, + DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, + PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info, + await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, + cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state, + count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size, + decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 4fdda942e..112e7fed0 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -1950,6 +1950,20 @@ mod tests { assert!(source_cleanup_versions_match_with_allowed_missing(&expected, ¤t, &allowed_missing)); } + #[test] + fn test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source() { + let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated"); + let mut free_version = cleanup_test_file_info("object.txt", Uuid::from_u128(2), "tier-cleanup"); + free_version.deleted = true; + free_version.set_tier_free_version(); + let mut expected = cleanup_test_versions(vec![migrated.clone()]); + expected.free_versions = vec![free_version.clone()]; + let current = cleanup_test_versions(vec![migrated]); + let allowed_missing = vec![source_cleanup_version_identity(&free_version)]; + + assert!(source_cleanup_versions_match_with_allowed_missing(&expected, ¤t, &allowed_missing)); + } + #[test] fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() { let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated"); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 2544c98a2..daf3f8175 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -227,15 +227,15 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap Result<()> { +) -> Result { let raw = match disk.read_xl(bucket, object, false).await { Ok(raw) => raw, - Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(()), + Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(false), Err(err) => return Err(err.into()), }; let meta = FileMeta::load(&raw.buf)?; @@ -253,8 +253,11 @@ async fn check_decommission_tier_free_version_target( all_matching_versions_equivalent = false; } } - if matching_count == 0 || (matching_count == 1 && all_matching_versions_equivalent) { - return Ok(()); + if matching_count == 0 { + return Ok(false); + } + if matching_count == 1 && all_matching_versions_equivalent { + return Ok(true); } Err(StorageError::DataMovementOverwriteErr( @@ -265,6 +268,28 @@ async fn check_decommission_tier_free_version_target( .into()) } +fn ensure_decommission_tier_free_version_commit_fence(bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { + if opts + .namespace_lock_fence + .as_ref() + .is_some_and(NamespaceLockFence::is_lock_lost) + || opts + .bucket_lifecycle_lock_fence + .as_ref() + .is_some_and(NamespaceLockFence::is_lock_lost) + { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "decommission_tier_free_version_commit", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } + + Ok(()) +} + impl SetDisks { pub(super) async fn require_current_restore_operation_id( &self, @@ -4718,26 +4743,11 @@ impl SetDisks { if !fi.deleted || !fi.tier_free_version() { return Err(Error::other("decommission tier free-version write requires a free version record")); } - if opts - .namespace_lock_fence - .as_ref() - .is_some_and(NamespaceLockFence::is_lock_lost) - || opts - .bucket_lifecycle_lock_fence - .as_ref() - .is_some_and(NamespaceLockFence::is_lock_lost) - { - return Err(StorageError::NamespaceLockQuorumUnavailable { - mode: "decommission_tier_free_version_commit", - bucket: bucket.to_string(), - object: object.to_string(), - required: 1, - achieved: 0, - }); - } + ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?; self.validate_decommission_tier_free_version_target(bucket, object, fi) .await?; + ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?; let disks = self.disks.read().await.clone(); let write_quorum = self.default_write_quorum(); @@ -4760,6 +4770,8 @@ impl SetDisks { } } + ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?; + resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object) } @@ -4776,13 +4788,38 @@ impl SetDisks { let preflight = disks .iter() .flatten() - .map(|disk| check_decommission_tier_free_version_target(disk, bucket, object, fi)); + .map(|disk| inspect_decommission_tier_free_version_target(disk, bucket, object, fi)); for result in join_all(preflight).await { result?; } Ok(()) } + pub(crate) async fn has_decommission_tier_free_version_write_quorum( + &self, + bucket: &str, + object: &str, + fi: &FileInfo, + opts: &ObjectOptions, + ) -> Result { + ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?; + let disks = self.disks.read().await.clone(); + let preflight = disks.iter().map(|disk| async { + match disk { + Some(disk) => inspect_decommission_tier_free_version_target(disk, bucket, object, fi).await, + None => Ok(false), + } + }); + let mut equivalent = 0; + for result in join_all(preflight).await { + if result? { + equivalent += 1; + } + } + ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?; + Ok(equivalent >= self.default_write_quorum()) + } + #[tracing::instrument(skip(self, fi, opts))] pub(crate) async fn decommission_tiered_object( &self, @@ -10056,6 +10093,72 @@ mod tests { assert_eq!(migrated.transitioned_objname, "remote/object"); } + #[tokio::test] + async fn decommission_tier_free_version_resume_requires_write_quorum() { + let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await; + let bucket = "free-version-decommission-resume"; + let object = "object.txt"; + let mut free_version = FileInfo { + name: object.to_string(), + volume: bucket.to_string(), + version_id: Some(Uuid::new_v4()), + mod_time: Some(time::OffsetDateTime::now_utc()), + deleted: true, + transition_tier: "WARM-TIER".to_string(), + transitioned_objname: "remote/object".to_string(), + ..Default::default() + }; + free_version.set_tier_free_version(); + let opts = ObjectOptions::default(); + + let disks = set_disks.get_disks_internal().await; + for disk in disks.iter().take(2).flatten() { + disk.write_metadata("", bucket, object, free_version.clone()) + .await + .expect("partial first attempt should leave equivalent metadata"); + } + assert!( + !set_disks + .has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts) + .await + .expect("partial target metadata should remain valid"), + "write-quorum-minus-one must not be accepted as an idempotent migration" + ); + + disks[2] + .as_ref() + .expect("third target disk should be online") + .write_metadata("", bucket, object, free_version.clone()) + .await + .expect("third equivalent target write should complete quorum"); + assert!( + set_disks + .has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts) + .await + .expect("write-quorum target metadata should remain valid") + ); + } + + #[test] + fn decommission_tier_free_version_commit_rejects_lost_fence() { + let opts = ObjectOptions { + namespace_lock_fence: Some(NamespaceLockFence::lost_for_test()), + ..Default::default() + }; + + let err = ensure_decommission_tier_free_version_commit_fence("bucket", "object", &opts) + .expect_err("lost target lock must fail the free-version commit"); + assert!(matches!( + err, + Error::NamespaceLockQuorumUnavailable { + mode: "decommission_tier_free_version_commit", + required: 1, + achieved: 0, + .. + } + )); + } + #[test] fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() { let errs = vec![None, None, Some(DiskError::DiskNotFound), None]; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 6fd983dfb..b889001d6 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -3934,6 +3934,18 @@ mod tests { .await .expect("create free-version decommission bucket"); let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await; + let source_free = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(&bucket, object) + .await + .expect("source free-version metadata should decode") + .and_then(|versions| { + versions + .versions + .into_iter() + .find(|version| version.version_id == Some(free_version) && version.tier_free_version()) + }) + .expect("source free-version identity should be present before decommission"); let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec()); let target_w = store.pools[1] .put_object( @@ -3988,6 +4000,15 @@ mod tests { .iter() .any(|version| { version.version_id == Some(free_version) && version.tier_free_version() }) ); + let migrated_free = target_versions + .versions + .iter() + .find(|version| version.version_id == Some(free_version) && version.tier_free_version()) + .expect("migrated free-version identity should remain readable from target disks"); + assert!( + crate::store::tiered_data_movement_source_matches(&source_free, migrated_free) + .expect("migrated free-version identity should decode") + ); let retained_w = target_versions .versions .iter() @@ -4010,6 +4031,22 @@ mod tests { .is_none(), "successful free-version migration should permit source cleanup" ); + let (heal_versions, _, _) = store + .heal_walk_versions_page(1, 0, &bucket, "", None, 2, 16, true) + .await + .expect("heal walk should decode the migrated free version"); + let free_version_string = free_version.to_string(); + let healed_free = heal_versions + .iter() + .find(|version| version.version_id.as_deref() == Some(free_version_string.as_str())) + .expect("heal walk should surface the migrated free version"); + let healed_info = healed_free + .lifecycle_object_info + .as_ref() + .expect("heal walk should retain lifecycle identity for the migrated free version"); + assert!(healed_info.transitioned_object.free_version); + assert_eq!(healed_info.transitioned_object.tier, source_free.transition_tier); + assert_eq!(healed_info.transitioned_object.name, source_free.transitioned_objname); shutdown.cancel(); } @@ -4132,11 +4169,11 @@ mod tests { .filter(|version| version.header.version_id == Some(free_version)) .collect::>(); assert_eq!(same_id.len(), 2, "conflict metadata should retain both same-ID records"); - assert_eq!(same_id.iter().filter(|version| version.free_version()).count(), 1); - assert_eq!(same_id.iter().filter(|version| !version.free_version()).count(), 1); + assert_eq!(same_id.iter().filter(|version| version.header.free_version()).count(), 1); + assert_eq!(same_id.iter().filter(|version| !version.header.free_version()).count(), 1); let post_conflict = same_id .into_iter() - .find(|version| !version.free_version()) + .find(|version| !version.header.free_version()) .expect("ordinary conflict version must remain addressable by the source ID"); let post_conflict_info = post_conflict .into_fileinfo(&bucket, object, true) @@ -4163,11 +4200,11 @@ mod tests { .collect::>(); if disk_index == 0 { assert_eq!(same_id.len(), 2); - assert!(same_id[0].free_version()); - assert!(!same_id[1].free_version()); + assert!(same_id[0].header.free_version()); + assert!(!same_id[1].header.free_version()); } else { assert_eq!(same_id.len(), 1); - assert!(same_id[0].free_version()); + assert!(same_id[0].header.free_version()); } } let sweep_err = store diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 90775a20d..dafcb56a0 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -2237,36 +2237,16 @@ impl ECStore { bucket: &str, object: &str, source: &rustfs_filemeta::FileInfo, + opts: &ObjectOptions, target_pool_idx: usize, ) -> Result { let pool = self .pools .get(target_pool_idx) .ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?; - let logical_object = decode_dir_object(object); - let Some(versions) = pool - .get_disks_by_key(object) - .load_file_info_versions_exact(bucket, &logical_object) - .await? - else { - return Ok(false); - }; - - let Some(target) = versions - .versions - .iter() - .find(|version| version.version_id == source.version_id) - else { - return Ok(false); - }; - if !target.tier_free_version() { - return Err(decommission_free_version_overwrite_error(bucket, object, source.version_id)); - } - if tiered_data_movement_source_matches(source, target)? { - Ok(true) - } else { - Err(decommission_free_version_overwrite_error(bucket, object, source.version_id)) - } + pool.get_disks_by_key(object) + .has_decommission_tier_free_version_write_quorum(bucket, object, source, opts) + .await } fn resolve_decommission_target_pool_idx_result(result: Result, bucket: &str, object: &str) -> Result { @@ -2365,11 +2345,7 @@ impl ECStore { return Err(Error::DiskFull); } let equivalent = if is_free_version { - self.pools[target_pool_idx] - .get_disks_by_key(&object) - .validate_decommission_tier_free_version_target(bucket, &object, &fi) - .await?; - self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, target_pool_idx) + self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, target_pool_idx) .await? } else { self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx) @@ -2383,12 +2359,8 @@ impl ECStore { } let result = if is_free_version { - self.pools[idx] - .get_disks_by_key(&object) - .validate_decommission_tier_free_version_target(bucket, &object, &fi) - .await?; if self - .has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, idx) + .has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, idx) .await? { return Ok(()); diff --git a/docs/architecture/decommission-compatibility.md b/docs/architecture/decommission-compatibility.md index f7d22acaa..35cb0caeb 100644 --- a/docs/architecture/decommission-compatibility.md +++ b/docs/architecture/decommission-compatibility.md @@ -235,6 +235,9 @@ report the enclosing object migration result. Regression guard: - `decommission_tier_free_version_preserves_remote_identity` +- `decommission_tier_free_version_resume_requires_write_quorum` +- `decommission_tier_free_version_commit_rejects_lost_fence` +- `test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source` - `decommission_entry_skips_cleanup_only_marker_when_free_version_is_present` - `decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source`