fix(ecstore): route batch delete markers to active pools

This commit is contained in:
overtrue
2026-08-22 00:54:26 +08:00
parent d139c8a8c3
commit cd06111dcb
2 changed files with 299 additions and 76 deletions
+177 -68
View File
@@ -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<crate::store::ECStore>,
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)]