From cd06111dcbec6c27ac8c71c572ff976f4ec21baf Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 00:54:26 +0800 Subject: [PATCH] fix(ecstore): route batch delete markers to active pools --- crates/ecstore/src/store/init.rs | 245 +++++++++++++++++++++-------- crates/ecstore/src/store/object.rs | 130 ++++++++++++++- 2 files changed, 299 insertions(+), 76 deletions(-) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 5c581a275..11877ee25 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -604,15 +604,15 @@ mod tests { storage_api_contracts::{ bucket::{BucketOperations as _, MakeBucketOptions}, multipart::MultipartOperations as _, - object::{ObjectIO, ObjectOperations as _}, + object::{ObjectIO, ObjectOperations as _, ObjectToDelete}, range::HTTPRangeSpec, }, }; use http::HeaderMap; use rustfs_config::server_config::KVS; - use rustfs_filemeta::ObjectPartInfo; #[cfg(feature = "test-util")] use rustfs_filemeta::{FileInfo, FileMeta}; + use rustfs_filemeta::{FileInfoVersions, ObjectPartInfo}; #[cfg(feature = "test-util")] use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase}; use rustfs_rio::{Checksum, ChecksumType}; @@ -1226,6 +1226,79 @@ mod tests { shutdown.cancel(); } + async fn migrate_versioned_decommission_test_object( + store: &Arc, + bucket: &str, + object: &str, + payload: &[u8], + op_label: &'static str, + ) -> (uuid::Uuid, FileInfoVersions) { + let mut source = PutObjReader::from_vec(payload.to_vec()); + let source_info = store.pools[0] + .put_object( + bucket, + object, + &mut source, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("write versioned source to the pool being decommissioned"); + let source_version = source_info.version_id.expect("versioned source must have a version ID"); + let expected_source_versions = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("source versions should be readable before migration") + .expect("source versions should exist before migration"); + { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + } + + let barrier = crate::set_disk::PutObjectCommitBarrier::install( + bucket, + object, + crate::set_disk::PutObjectCommitPause::BeforeNamespace, + ); + let migration_store = Arc::clone(store); + let migration_bucket = bucket.to_string(); + let migration_object = object.to_string(); + let migration = tokio::spawn(async move { + let source_reader = migration_store.pools[0] + .get_object_reader( + &migration_bucket, + &migration_object, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(source_version.to_string()), + no_lock: true, + data_movement: true, + raw_data_movement_read: true, + ..Default::default() + }, + ) + .await?; + crate::data_movement::migrate_decommission_object(migration_store, 0, migration_bucket, source_reader, None, op_label) + .await + }); + barrier.wait_until_paused().await; + barrier.release(); + migration + .await + .expect("versioned decommission migration task should join") + .expect("versioned decommission migration should commit"); + + (source_version, expected_source_versions) + } + #[tokio::test] #[serial_test::serial(storage_class_env)] async fn tag_updates_skip_active_rebalance_source_pool() { @@ -2872,74 +2945,14 @@ mod tests { .make_bucket(&bucket, &MakeBucketOptions::default()) .await .expect("create versioned decommission delete-fence bucket"); - let mut source = PutObjReader::from_vec(b"source generation".to_vec()); - let source_info = store.pools[0] - .put_object( - &bucket, - object, - &mut source, - &ObjectOptions { - versioned: true, - ..Default::default() - }, - ) - .await - .expect("write source version to the pool being decommissioned"); - let source_version = source_info.version_id.expect("versioned source must have a version ID"); - let expected_source_versions = store.pools[0] - .get_disks_by_key(object) - .load_file_info_versions_exact(&bucket, object) - .await - .expect("source versions should be readable before migration") - .expect("source versions should exist before migration"); - { - let mut pool_meta = store.pool_meta.write().await; - pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { - start_time: Some(OffsetDateTime::now_utc()), - ..Default::default() - }); - } - - let barrier = crate::set_disk::PutObjectCommitBarrier::install( + let (source_version, expected_source_versions) = migrate_versioned_decommission_test_object( + &store, &bucket, object, - crate::set_disk::PutObjectCommitPause::BeforeNamespace, - ); - let migration_store = Arc::clone(&store); - let migration_bucket = bucket.clone(); - let migration = tokio::spawn(async move { - let source_reader = migration_store.pools[0] - .get_object_reader( - &migration_bucket, - object, - None, - HeaderMap::new(), - &ObjectOptions { - versioned: true, - version_id: Some(source_version.to_string()), - no_lock: true, - data_movement: true, - raw_data_movement_read: true, - ..Default::default() - }, - ) - .await?; - crate::data_movement::migrate_decommission_object( - migration_store, - 0, - migration_bucket, - source_reader, - None, - "test_versioned_decommission_delete_fence", - ) - .await - }); - barrier.wait_until_paused().await; - barrier.release(); - migration - .await - .expect("versioned decommission migration task should join") - .expect("versioned decommission migration should commit before DELETE"); + b"source generation", + "test_versioned_decommission_delete_fence", + ) + .await; let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); let delete_store = Arc::clone(&store); @@ -3057,6 +3070,102 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn versioned_batch_delete_marker_skips_decommission_source() { + let temp_dir = tempfile::tempdir().expect("create versioned batch decommission store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "versioned-batch-decommission-delete-fence", + &[4, 4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("versioned-batch-decommission-delete-fence-{}", uuid::Uuid::new_v4()); + let object = "batch-object.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create versioned batch decommission bucket"); + let (_source_version, expected_source_versions) = migrate_versioned_decommission_test_object( + &store, + &bucket, + object, + b"batch source generation", + "test_versioned_batch_decommission_delete_fence", + ) + .await; + + let delete_config_snapshot = + Arc::new(crate::bucket::replication::DeleteReplicationConfigSnapshot::from_configs_for_test( + s3s::dto::VersioningConfiguration { + status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::ENABLED)), + ..Default::default() + }, + None, + )); + let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); + let delete_store = Arc::clone(&store); + let delete_bucket = bucket.clone(); + let delete = tokio::spawn(async move { + delete_store + .delete_objects( + &delete_bucket, + vec![ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + ObjectOptions { + delete_replication_config_snapshot: Some(delete_config_snapshot), + ..Default::default() + }, + ) + .await + }); + delete_barrier.wait_until_paused().await; + + let source_set = store.pools[0].get_disks_by_key(object); + crate::data_movement::ensure_source_cleanup_versions_unchanged( + source_set, + &bucket, + object, + &expected_source_versions, + &[], + "test_versioned_batch_decommission_delete_fence", + ) + .await + .expect("batch DELETE must not publish a marker to the suspended source"); + + delete_barrier.release(); + let (deleted, errors) = delete.await.expect("versioned batch DELETE task should join"); + assert!(errors.iter().all(Option::is_none), "versioned batch DELETE should succeed: {errors:?}"); + assert_eq!(deleted.len(), 1); + assert!(deleted[0].delete_marker, "versioned batch DELETE must return a marker"); + assert!( + deleted[0] + .delete_marker_version_id + .is_some_and(|version_id| !version_id.is_nil()), + "versioned batch DELETE marker must have a non-nil version ID" + ); + + let mut active_marker_count = 0; + for pool in store.pools.iter().skip(1) { + let Some(versions) = pool + .get_disks_by_key(object) + .load_file_info_versions_exact(&bucket, object) + .await + .expect("active-pool versions should be readable") + else { + continue; + }; + active_marker_count += versions.versions.iter().filter(|version| version.deleted).count(); + } + assert_eq!(active_marker_count, 1, "batch DELETE must publish exactly one active-pool marker"); + + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index e913cfe2d..7d4e797e1 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -32,7 +32,8 @@ use crate::bucket::metadata_sys::{ use crate::bucket::object_lock::objectlock_sys::{ check_object_lock_for_deletion_with_state, ensure_recursive_force_delete_allowed_for_state, }; -use crate::bucket::replication::ReplicationObjectBridge; +use crate::bucket::replication::{DeleteReplicationConfigSnapshot, ReplicationObjectBridge}; +use crate::bucket::versioning::VersioningApi; use crate::disk::OldCurrentSize; use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot}; use crate::set_disk::{ @@ -1033,6 +1034,20 @@ fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some()) } +fn batch_delete_creates_latest_marker(object: &ObjectToDelete, delete_config_snapshot: &DeleteReplicationConfigSnapshot) -> bool { + if object.version_id.is_some() { + return false; + } + + let object_name = decode_dir_object(&object.object_name); + let (versioned, version_suspended) = delete_config_snapshot.versioning_config().delete_state(&object_name); + versioned || version_suspended +} + +fn batch_delete_targets_pool(creates_latest_marker: bool, marker_target_pool_idx: Option, pool_idx: usize) -> bool { + !creates_latest_marker || marker_target_pool_idx == Some(pool_idx) +} + fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions { let mut lookup_opts = opts.clone(); lookup_opts.skip_decommissioned = true; @@ -2943,30 +2958,98 @@ impl ECStore { Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), }; - let mut futures = Vec::with_capacity(self.pools.len()); + let delete_config_snapshot = opts + .delete_replication_config_snapshot + .as_deref() + .expect("batch delete replication config snapshot should be loaded"); + let latest_marker_objects = objects + .iter() + .map(|object| batch_delete_creates_latest_marker(object, delete_config_snapshot)) + .collect::>(); + let marker_target_results = join_all(objects.iter().zip(&latest_marker_objects).map( + |(object, creates_marker)| async move { + if *creates_marker { + Some(self.get_pool_idx_no_lock(bucket, &object.object_name, 0).await) + } else { + None + } + }, + )) + .await; + let mut marker_target_pool_indices = Vec::with_capacity(objects.len()); + for (idx, target_result) in marker_target_results.into_iter().enumerate() { + match target_result { + Some(Ok(pool_idx)) => marker_target_pool_indices.push(Some(pool_idx)), + Some(Err(err)) => { + del_errs[idx] = Some(err); + marker_target_pool_indices.push(None); + } + None => marker_target_pool_indices.push(None), + } + } + let mut futures = Vec::with_capacity(self.pools.len()); for pool in self.pools.iter() { if self.is_pool_rebalancing(pool.pool_idx).await { continue; } - futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone())); + + let (object_indices, pool_objects): (Vec<_>, Vec<_>) = objects + .iter() + .enumerate() + .filter(|(idx, _)| { + batch_delete_targets_pool(latest_marker_objects[*idx], marker_target_pool_indices[*idx], pool.pool_idx) + }) + .map(|(idx, object)| (idx, object.clone())) + .unzip(); + if pool_objects.is_empty() { + continue; + } + + let pool_opts = opts.clone(); + futures.push(async move { + let result = pool.delete_objects(bucket, pool_objects, pool_opts).await; + (object_indices, result) + }); } let results = join_all(futures).await; + let mut attempted = vec![false; del_objects.len()]; for idx in 0..del_objects.len() { - for (dels, errs) in results.iter() { - if errs[idx].is_none() && dels[idx].found { + for (object_indices, (dels, errs)) in results.iter() { + let Ok(pool_object_idx) = object_indices.binary_search(&idx) else { + continue; + }; + attempted[idx] = true; + + if errs[pool_object_idx].is_none() && dels[pool_object_idx].found { del_errs[idx] = None; - del_objects[idx] = dels[idx].clone(); + del_objects[idx] = dels[pool_object_idx].clone(); break; } if del_errs[idx].is_none() { - del_errs[idx] = errs[idx].clone(); - del_objects[idx] = dels[idx].clone(); + del_errs[idx] = errs[pool_object_idx].clone(); + del_objects[idx] = dels[pool_object_idx].clone(); } } + + if !attempted[idx] && del_errs[idx].is_none() && latest_marker_objects[idx] { + del_objects[idx] = DeletedObject { + object_name: objects[idx].object_name.clone(), + version_id: objects[idx].version_id, + ..Default::default() + }; + del_errs[idx] = Some(StorageError::ObjectNotFound(bucket.to_owned(), objects[idx].object_name.clone())); + } + } + + #[cfg(test)] + for (idx, object) in objects.iter().enumerate() { + if del_errs[idx].is_none() && del_objects[idx].delete_marker { + pause_versioned_delete_marker_after_commit(bucket, &object.object_name).await; + } } del_objects.iter_mut().for_each(|v| { @@ -4639,6 +4722,37 @@ mod tests { assert!(!should_delete_from_all_pools(&ObjectOptions::default(), 0)); } + #[test] + fn batch_delete_identifies_only_latest_versioned_markers() { + let versioned = DeleteReplicationConfigSnapshot::from_configs_for_test( + s3s::dto::VersioningConfiguration { + status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::ENABLED)), + ..Default::default() + }, + None, + ); + let latest = ObjectToDelete { + object_name: "latest".to_string(), + ..Default::default() + }; + assert!(batch_delete_creates_latest_marker(&latest, &versioned)); + assert!(!batch_delete_targets_pool(true, Some(1), 0)); + assert!(batch_delete_targets_pool(true, Some(1), 1)); + assert!(!batch_delete_targets_pool(true, Some(1), 2)); + + let explicit = ObjectToDelete { + object_name: "explicit".to_string(), + version_id: Some(uuid::Uuid::new_v4()), + ..Default::default() + }; + assert!(!batch_delete_creates_latest_marker(&explicit, &versioned)); + assert!(batch_delete_targets_pool(false, Some(1), 0)); + + let unversioned = DeleteReplicationConfigSnapshot::default(); + assert!(!batch_delete_creates_latest_marker(&latest, &unversioned)); + assert!(batch_delete_targets_pool(false, None, 0)); + } + #[test] fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() { let lookup_opts = data_movement_pool_lookup_opts(