diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index bec66d9fe..e86d7c583 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -2640,8 +2640,12 @@ fn determine_decommission_final_state(items_failed: usize, was_cancelled: bool) } } -fn decommission_remaining_version_count(total_versions: usize, expired: usize) -> usize { - total_versions.saturating_sub(expired) +fn decommission_remaining_version_count(versions: &[rustfs_filemeta::FileInfo], expired: usize) -> usize { + versions + .iter() + .filter(|version| !version.tier_free_version()) + .count() + .saturating_sub(expired) } fn should_skip_decommission_delete_marker( @@ -3954,7 +3958,7 @@ impl ECStore { continue; } - let remaining_versions = decommission_remaining_version_count(fivs.versions.len(), expired); + let remaining_versions = decommission_remaining_version_count(&fivs.versions, expired); if should_skip_decommission_delete_marker(version, remaining_versions, replication_config.is_some()) { // decommissioned += 1; @@ -4419,6 +4423,39 @@ impl ECStore { .await } + #[cfg(test)] + pub(crate) async fn decommission_entry_for_test_with_bucket_incarnation( + self: &Arc, + idx: usize, + entry: MetaCacheEntry, + bucket: String, + set: Arc, + ) -> Result<()> { + let expected_bucket_incarnation_id = if is_meta_bucketname(&bucket) { + None + } else { + Some(self.bucket_incarnation_id_from_disk(&bucket).await?) + }; + self.decommission_entry( + CancellationToken::new(), + idx, + OffsetDateTime::now_utc(), + entry, + bucket, + set, + None, + None, + None, + expected_bucket_incarnation_id, + ) + .await + } + + #[cfg(test)] + pub(crate) async fn check_after_decommission_for_test(self: &Arc, idx: usize) -> Result<()> { + self.check_after_decommission(idx).await + } + #[tracing::instrument(skip(self, rx))] async fn decommission_pool( self: &Arc, @@ -5459,13 +5496,6 @@ mod tests { assert_eq!(determine_decommission_final_state(0, true), DecommissionFinalState::Failed); } - #[test] - fn decommission_remaining_version_count_excludes_only_expired_versions() { - assert_eq!(decommission_remaining_version_count(1, 0), 1); - assert_eq!(decommission_remaining_version_count(2, 1), 1); - assert_eq!(decommission_remaining_version_count(1, 1), 0); - } - #[test] fn lifecycle_action_removes_data_movement_version_rejects_delete_marker_action() { assert!(!lifecycle_action_removes_data_movement_version(IlmAction::DeleteAction)); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 4f73627e7..f5bb563a4 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -227,6 +227,44 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap Result<()> { + let raw = match disk.read_xl(bucket, object, false).await { + Ok(raw) => raw, + Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(()), + Err(err) => return Err(err.into()), + }; + let meta = FileMeta::load(&raw.buf)?; + let source_version_id = source.version_id.filter(|version_id| !version_id.is_nil()); + let mut matching_count = 0; + let mut all_matching_versions_equivalent = true; + for existing in meta + .versions + .iter() + .filter(|version| version.header.version_id.filter(|version_id| !version_id.is_nil()) == source_version_id) + { + matching_count += 1; + let existing = existing.into_fileinfo(bucket, object, true)?; + if !existing.tier_free_version() || !crate::store::object::tiered_data_movement_source_matches(source, &existing)? { + all_matching_versions_equivalent = false; + } + } + if matching_count == 0 || (matching_count == 1 && all_matching_versions_equivalent) { + return Ok(()); + } + + Err(StorageError::DataMovementOverwriteErr( + bucket.to_owned(), + object.to_owned(), + source_version_id.map(|version_id| version_id.to_string()).unwrap_or_default(), + ) + .into()) +} + impl SetDisks { pub(super) async fn require_current_restore_operation_id( &self, @@ -4698,6 +4736,9 @@ impl SetDisks { }); } + self.validate_decommission_tier_free_version_target(bucket, object, fi) + .await?; + let disks = self.disks.read().await.clone(); let write_quorum = self.default_write_quorum(); let futures = disks.into_iter().map(|disk| { @@ -4722,6 +4763,26 @@ impl SetDisks { resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object) } + pub(crate) async fn validate_decommission_tier_free_version_target( + &self, + bucket: &str, + object: &str, + fi: &FileInfo, + ) -> Result<()> { + // The caller holds the source and target object locks. Inspect every + // target disk before an idempotent return or metadata fan-out so a + // sub-quorum conflict cannot be hidden by a successful quorum. + let disks = self.disks.read().await.clone(); + let preflight = disks + .iter() + .flatten() + .map(|disk| check_decommission_tier_free_version_target(disk, bucket, object, fi)); + for result in join_all(preflight).await { + result?; + } + Ok(()) + } + #[tracing::instrument(skip(self, fi, opts))] pub(crate) async fn decommission_tiered_object( &self, diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 451e42242..98eab2c52 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -595,7 +595,7 @@ mod tests { use crate::{ bucket::replication::{ReplicationState, ReplicationStatusType, replication_statuses_map}, core::pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolStatus}, - disk::endpoint::Endpoint, + disk::{DiskAPI, endpoint::Endpoint}, error::{Error, Result, StorageError}, io_support::rio::{WritePlan, compression_metadata_value}, layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}, @@ -1307,6 +1307,77 @@ mod tests { }); } + #[cfg(feature = "test-util")] + async fn seed_transitioned_free_version( + ctx: &Arc, + store: &Arc, + bucket: &str, + object: &str, + ) -> (uuid::Uuid, uuid::Uuid) { + let tier_name = format!("DECOMFREE{}", uuid::Uuid::new_v4().simple()); + register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await; + + let mut reader = PutObjReader::from_vec(b"transitioned source bytes".to_vec()); + let source = store.pools[0] + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("write transitioned decommission source"); + let source_version = source.version_id.expect("transitioned source must be versioned"); + store.pools[0] + .transition_object( + bucket, + object, + &ObjectOptions { + versioned: true, + version_id: Some(source_version.to_string()), + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name, + etag: source.etag.clone().expect("transitioned source must have an ETag"), + ..Default::default() + }, + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect("transition source before decommission"); + store.pools[0] + .delete_object( + bucket, + object, + ObjectOptions { + versioned: true, + version_id: Some(source_version.to_string()), + ..Default::default() + }, + ) + .await + .expect("delete transitioned source version"); + + let versions = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("source versions should decode after transition delete") + .expect("source free version should remain after transition delete"); + let free_version = versions + .versions + .iter() + .find(|version| version.tier_free_version()) + .and_then(|version| version.version_id) + .expect("transition delete should create a free version"); + (source_version, free_version) + } + async fn write_decommission_test_multipart_source( store: &Arc, pool_idx: usize, @@ -3846,6 +3917,268 @@ mod tests { shutdown.cancel(); } + #[cfg(feature = "test-util")] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[serial_test::serial(storage_class_env)] + async fn decommission_entry_skips_cleanup_only_marker_when_free_version_is_present() { + let temp_dir = tempfile::tempdir().expect("create free-version decommission store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-marker", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let bucket = format!("decom-free-marker-{}", uuid::Uuid::new_v4()); + let object = "free-marker-object"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create free-version decommission bucket"); + let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await; + let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec()); + let target_w = store.pools[1] + .put_object( + &bucket, + object, + &mut target_reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("write unrelated target version"); + let target_w_version = target_w.version_id.expect("target version should have an id"); + let marker = store.pools[0] + .delete_object( + &bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("write cleanup-only delete marker"); + assert!(marker.delete_marker); + + mark_test_pool_decommissioning(&store, 0).await; + let source_set = store.pools[0].get_disks_by_key(object); + store + .decommission_entry_for_test_with_bucket_incarnation( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + bucket.clone(), + source_set.clone(), + ) + .await + .expect("real decommission entry should migrate the free version"); + + let target_versions = store.pools[1] + .get_disks_by_key(object) + .load_file_info_versions_exact(&bucket, object) + .await + .expect("target versions should decode") + .expect("target free version should be present"); + assert!( + target_versions + .versions + .iter() + .any(|version| { version.version_id == Some(free_version) && version.tier_free_version() }) + ); + let retained_w = target_versions + .versions + .iter() + .find(|version| version.version_id == Some(target_w_version)) + .expect("unrelated target version should remain"); + assert!(!retained_w.deleted && !retained_w.tier_free_version()); + assert_eq!(retained_w.size, target_w.size); + assert_eq!(retained_w.get_etag(), target_w.etag); + assert!( + target_versions + .versions + .iter() + .all(|version| { version.tier_free_version() || !version.deleted }) + ); + assert!( + source_set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("source versions should be readable after cleanup") + .is_none(), + "successful free-version migration should permit source cleanup" + ); + shutdown.cancel(); + } + + #[cfg(feature = "test-util")] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[serial_test::serial(storage_class_env)] + async fn decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source() { + let temp_dir = tempfile::tempdir().expect("create sub-quorum free-version store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-conflict", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let bucket = format!("decom-free-conflict-{}", uuid::Uuid::new_v4()); + let object = "free-conflict-object"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create sub-quorum conflict 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 should decode before crash replay setup") + .and_then(|versions| { + versions + .versions + .into_iter() + .find(|version| version.version_id == Some(free_version)) + }) + .expect("source free version should be available for crash replay setup"); + + let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec()); + let target = store.pools[1] + .put_object( + &bucket, + object, + &mut target_reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("seed ordinary target version"); + let target_version = target.version_id.expect("target version must have an ID"); + let target_disks = store.pools[1].get_disks_by_key(object).disks.read().await.clone(); + for disk in target_disks.iter().skip(1) { + disk.as_ref() + .expect("target crash replay quorum disk should be online") + .write_metadata("", &bucket, object, source_free.clone()) + .await + .expect("seed an equivalent free version on the target quorum"); + } + let conflict_path = temp_dir + .path() + .join(format!("pool1/set0/disk0/{bucket}/{object}/{STORAGE_FORMAT_FILE}")); + let encoded = tokio::fs::read(&conflict_path) + .await + .expect("target metadata should be readable"); + let mut metadata = FileMeta::load(&encoded).expect("target metadata should decode"); + let target_index = metadata + .versions + .iter() + .position(|version| version.header.version_id == Some(target_version)) + .expect("target version should be present on the conflict disk"); + let mut target_meta = metadata.versions[target_index] + .parse_version_meta() + .expect("target version metadata should decode"); + target_meta + .object + .as_mut() + .expect("target conflict must remain an ordinary object") + .version_id = Some(free_version); + metadata.versions[target_index] = target_meta.try_into().expect("conflict metadata should encode"); + let expected_conflict_meta = metadata.versions[target_index].meta.clone(); + let expected_conflict = metadata.versions[target_index] + .into_fileinfo(&bucket, object, true) + .expect("conflict metadata should decode as an ordinary object"); + let duplicate_free: rustfs_filemeta::FileMetaShallowVersion = rustfs_filemeta::FileMetaVersion::from(source_free.clone()) + .try_into() + .expect("duplicate free metadata should encode"); + metadata.versions.insert(target_index, duplicate_free); + tokio::fs::write(&conflict_path, metadata.marshal_msg().expect("conflict metadata should encode")) + .await + .expect("write sub-quorum conflict metadata"); + + mark_test_pool_decommissioning(&store, 0).await; + let source_set = store.pools[0].get_disks_by_key(object); + store + .decommission_entry_for_test_with_bucket_incarnation( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + bucket.clone(), + source_set.clone(), + ) + .await + .expect("conflicted decommission entry should retain the source and retry later"); + + let source_versions = source_set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("retained source versions should decode") + .expect("source free version should be retained after conflict"); + assert!( + source_versions + .versions + .iter() + .any(|version| { version.version_id == Some(free_version) && version.tier_free_version() }) + ); + let post_encoded = tokio::fs::read(&conflict_path) + .await + .expect("conflict metadata should remain readable"); + let post_metadata = FileMeta::load(&post_encoded).expect("post-conflict metadata should decode"); + let same_id = post_metadata + .versions + .iter() + .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); + let post_conflict = same_id + .into_iter() + .find(|version| !version.free_version()) + .expect("ordinary conflict version must remain addressable by the source ID"); + let post_conflict_info = post_conflict + .into_fileinfo(&bucket, object, true) + .expect("post-conflict ordinary metadata should decode"); + assert!(!post_conflict_info.deleted && !post_conflict_info.tier_free_version()); + assert_eq!(post_conflict.meta, expected_conflict_meta); + assert_eq!(post_conflict_info.size, expected_conflict.size); + assert_eq!(post_conflict_info.data_dir, expected_conflict.data_dir); + assert_eq!(post_conflict_info.metadata, expected_conflict.metadata); + assert_eq!(post_conflict_info.get_etag(), expected_conflict.get_etag()); + + for disk_index in 0..4 { + let target_path = temp_dir + .path() + .join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}")); + let target_encoded = tokio::fs::read(&target_path) + .await + .expect("target metadata should remain readable"); + let target_meta = FileMeta::load(&target_encoded).expect("target metadata should decode"); + let same_id = target_meta + .versions + .iter() + .filter(|version| version.header.version_id == Some(free_version)) + .collect::>(); + if disk_index == 0 { + assert_eq!(same_id.len(), 2); + assert!(same_id[0].free_version()); + assert!(!same_id[1].free_version()); + } else { + assert_eq!(same_id.len(), 1); + assert!(same_id[0].free_version()); + } + } + let sweep_err = store + .check_after_decommission_for_test(0) + .await + .expect_err("final sweep must report the retained free version"); + assert!( + sweep_err.to_string().contains("version(s) were found"), + "unexpected final sweep error: {sweep_err}" + ); + shutdown.cancel(); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial_test::serial(storage_class_env)] async fn versioned_batch_delete_marker_skips_decommission_source() { diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 566576034..8f30269d2 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1494,13 +1494,15 @@ fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo, && source_actual_size == target_actual_size } -fn tiered_data_movement_source_matches( +pub(crate) fn tiered_data_movement_source_matches( expected: &rustfs_filemeta::FileInfo, current: &rustfs_filemeta::FileInfo, ) -> Result { let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?; let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(¤t.metadata)?; Ok(expected.version_id == current.version_id + && expected.deleted == current.deleted + && expected.tier_free_version() == current.tier_free_version() && expected.data_dir == current.data_dir && expected.mod_time == current.mod_time && expected.size == current.size @@ -2363,6 +2365,10 @@ 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) .await? } else { @@ -2377,6 +2383,10 @@ 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) .await? diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index 9225a1edb..7e7f61f6d 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -462,6 +462,10 @@ impl FileMeta { }; if let Some(fidx) = existing_idx { + let existing = self.versions[fidx].parse_version_meta()?; + if existing.free_version() != version.free_version() { + return Err(Error::other("cannot replace a free version with a non-free version")); + } return self.set_idx(fidx, version); }