test(ecstore): exercise decommission delete fences

This commit is contained in:
overtrue
2026-08-22 02:01:55 +08:00
parent caa7eeac0b
commit 8a77034463
4 changed files with 365 additions and 54 deletions
+16
View File
@@ -3383,6 +3383,22 @@ impl ECStore {
Ok(())
}
#[cfg(test)]
pub(crate) async fn decommission_entry_for_test(
self: &Arc<Self>,
idx: usize,
entry: MetaCacheEntry,
bucket: String,
set: Arc<SetDisks>,
) -> Result<()> {
let worker_permit = Arc::new(Semaphore::new(1))
.acquire_owned()
.await
.map_err(|err| Error::other(format!("decommission test worker permit acquire failed: {err}")))?;
self.decommission_entry(CancellationToken::new(), idx, entry, bucket, set, worker_permit, None, None, None, None)
.await
}
#[tracing::instrument(skip(self, rx))]
async fn decommission_pool(
self: &Arc<Self>,
+23
View File
@@ -2483,6 +2483,7 @@ impl SetDisks {
})
.await?,
);
notify_put_object_commit_namespace_acquired(bucket, object);
}
#[cfg(not(any(test, feature = "test-util")))]
{
@@ -4630,6 +4631,7 @@ struct PutObjectCommitBarrierState {
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
namespace_pending: tokio::sync::Notify,
namespace_acquired: std::sync::atomic::AtomicBool,
}
#[cfg(any(test, feature = "test-util"))]
@@ -4651,6 +4653,7 @@ impl PutObjectCommitBarrier {
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
namespace_pending: tokio::sync::Notify::new(),
namespace_acquired: std::sync::atomic::AtomicBool::new(false),
});
let mut slot = PUT_OBJECT_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
@@ -4685,6 +4688,10 @@ impl PutObjectCommitBarrier {
.await
.expect("put object should wait for the namespace lock after leaving the commit barrier");
}
pub fn namespace_acquired(&self) -> bool {
self.state.namespace_acquired.load(std::sync::atomic::Ordering::Acquire)
}
}
#[cfg(any(test, feature = "test-util"))]
@@ -4741,6 +4748,22 @@ fn notify_put_object_commit_namespace_pending(bucket: &str, object: &str) {
}
}
#[cfg(any(test, feature = "test-util"))]
fn notify_put_object_commit_namespace_acquired(bucket: &str, object: &str) {
let barrier = PUT_OBJECT_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("put object commit barrier mutex should not poison")
.iter()
.find(|barrier| {
barrier.bucket == bucket && barrier.object == object && barrier.pause == PutObjectCommitPause::BeforeNamespace
})
.cloned();
if let Some(barrier) = barrier {
barrier.namespace_acquired.store(true, std::sync::atomic::Ordering::Release);
}
}
#[cfg(test)]
struct DeleteObjectCommitBarrierState {
bucket: String,
+205 -54
View File
@@ -612,7 +612,7 @@ mod tests {
use rustfs_config::server_config::KVS;
#[cfg(feature = "test-util")]
use rustfs_filemeta::{FileInfo, FileMeta};
use rustfs_filemeta::{FileInfoVersions, ObjectPartInfo};
use rustfs_filemeta::{FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
#[cfg(feature = "test-util")]
use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase};
use rustfs_rio::{Checksum, ChecksumType};
@@ -2827,21 +2827,30 @@ mod tests {
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn delete_waits_for_decommission_commit_then_removes_every_copy() {
async fn decommission_entry_carries_migration_and_cleanup_mutation_fences() {
let temp_dir = tempfile::tempdir().expect("create decommission delete-fence store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-delete-fence", &[4, 4])).await;
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store_with_layout(
temp_dir.path(),
"decommission-delete-fence",
&[(2, 4), (1, 4)],
CancellationToken::new(),
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-delete-fence-{}", uuid::Uuid::new_v4());
let object = "object.bin";
let object = (0..128)
.map(|index| format!("object-{index}.bin"))
.find(|candidate| store.pools[0].get_disks_by_key(candidate).set_index == 1)
.expect("the deterministic object search should select source set 1");
let source_body = b"source generation".to_vec();
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission delete-fence bucket");
let mut source = PutObjReader::from_vec(b"source generation".to_vec());
let mut source = PutObjReader::from_vec(source_body.clone());
store.pools[0]
.put_object(&bucket, object, &mut source, &ObjectOptions::default())
.put_object(&bucket, &object, &mut source, &ObjectOptions::default())
.await
.expect("write source object to the pool being decommissioned");
{
@@ -2855,77 +2864,219 @@ mod tests {
let barrier = crate::set_disk::PutObjectCommitBarrier::install(
&bucket,
object,
&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 {
no_lock: true,
data_movement: true,
raw_data_movement_read: true,
let cleanup_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, &object);
let source_set = store.pools[0].get_disks_by_key(&object);
assert_eq!(source_set.set_index, 1, "the source entry must exercise the non-fixed set cleanup lock");
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker_object = object.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: worker_object,
..Default::default()
},
worker_bucket,
source_set,
)
.await?;
crate::data_movement::migrate_decommission_object(
migration_store,
0,
migration_bucket,
source_reader,
None,
"test_decommission_delete_fence",
)
.await
.await
});
barrier.wait_until_paused().await;
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
let delete_store = Arc::clone(&store);
let delete_bucket = bucket.clone();
let delete_object = object.clone();
let delete = tokio::spawn(async move {
delete_store
.delete_object(&delete_bucket, object, ObjectOptions::default())
.delete_object(&delete_bucket, &delete_object, ObjectOptions::default())
.await
});
delete_barrier.wait_until_paused().await;
delete_barrier.release_and_wait_until_namespace_pending().await;
assert!(
!delete.is_finished(),
"DELETE must wait while the decommission source generation is being committed"
!delete_barrier.namespace_acquired() && !delete.is_finished(),
"DELETE must remain before namespace acquisition behind the decommission worker's target-commit mutation fence"
);
delete.abort();
assert!(
delete
.await
.expect_err("the blocked DELETE should be canceled")
.is_cancelled(),
"the competing DELETE must remain cancelable while blocked"
);
barrier.release();
migration
.await
.expect("decommission migration task should join")
.expect("decommission migration should commit before DELETE");
delete
.await
.expect("DELETE task should join")
.expect("DELETE should remove the committed migration generation");
cleanup_barrier.wait_until_paused().await;
drop(barrier);
for pool in &store.pools {
let err = pool
.get_object_info(&bucket, object, &ObjectOptions::default())
let fixed_set = Arc::clone(&store.pools[0].disk_set[0]);
let fixed_mutation_barrier = crate::set_disk::PutObjectCommitBarrier::install(
&bucket,
&object,
crate::set_disk::PutObjectCommitPause::BeforeNamespace,
);
let mutation_bucket = bucket.clone();
let mutation_object = object.clone();
let fixed_mutation = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(b"fixed-domain replacement".to_vec());
fixed_set
.put_object(&mutation_bucket, &mutation_object, &mut reader, &ObjectOptions::default())
.await
.expect_err("DELETE must remove the source and migrated target copies");
assert!(
matches!(err, StorageError::ObjectNotFound(_, _)),
"unexpected post-delete pool result: {err:?}"
);
}
store
.get_object_info(&bucket, object, &ObjectOptions::default())
});
fixed_mutation_barrier.wait_until_paused().await;
fixed_mutation_barrier.release_and_wait_until_namespace_pending().await;
assert!(
!fixed_mutation_barrier.namespace_acquired() && !fixed_mutation.is_finished(),
"the source cleanup must retain the fixed mutation fence before the set-0 mutation acquires its namespace"
);
fixed_mutation.abort();
assert!(
fixed_mutation
.await
.expect_err("the fixed-domain mutation should be canceled")
.is_cancelled(),
"the competing fixed-domain mutation must remain cancelable while blocked"
);
drop(fixed_mutation_barrier);
cleanup_barrier.release();
worker
.await
.expect_err("the migrated source generation must not become visible again");
.expect("decommission entry worker should join")
.expect("decommission entry should migrate and clean its source");
store.pools[0]
.get_object_info(&bucket, &object, &ObjectOptions::default())
.await
.expect_err("the real decommission entry must clean the source generation");
let mut target = store.pools[1]
.get_object_reader(&bucket, &object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("the real decommission entry must commit the target generation");
let mut actual = Vec::new();
target
.stream
.read_to_end(&mut actual)
.await
.expect("read the migrated target generation");
assert_eq!(actual, source_body, "the decommission worker must preserve the migrated object body");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn batch_delete_real_path_preserves_source_pool_errors_in_any_pool_order() {
let temp_dir = tempfile::tempdir().expect("create batch delete pool-error store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-delete-pool-errors", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
for source_pool_idx in [0, 1] {
{
let mut pool_meta = store.pool_meta.write().await;
for pool in &mut pool_meta.pools {
pool.decommission = None;
}
}
let bucket = format!("batch-delete-pool-errors-{source_pool_idx}-{}", uuid::Uuid::new_v4());
let object_names = vec![
format!("third-{source_pool_idx}.bin"),
format!("first-{source_pool_idx}.bin"),
format!("second-{source_pool_idx}.bin"),
];
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create batch delete pool-error bucket");
for pool in &store.pools {
for object_name in &object_names {
let mut reader = PutObjReader::from_vec(format!("pool {} {object_name}", pool.pool_idx).into_bytes());
pool.put_object(&bucket, object_name, &mut reader, &ObjectOptions::default())
.await
.expect("seed each object in both the source and active pools");
}
}
{
let mut pool_meta = store.pool_meta.write().await;
pool_meta.pools[source_pool_idx].decommission = Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::now_utc()),
..Default::default()
});
}
assert!(
store.is_suspended(source_pool_idx).await,
"the injected error pool must be the decommission source"
);
let expected_errors = vec![
StorageError::ErasureWriteQuorum,
StorageError::NamespaceLockQuorumUnavailable {
mode: "delete_objects_commit",
bucket: bucket.clone(),
object: object_names[1].clone(),
required: 3,
achieved: 2,
},
StorageError::ErasureWriteQuorum,
];
let injection = crate::store::object::BatchDeletePoolErrorInjection::install(
&bucket,
source_pool_idx,
object_names.iter().cloned().zip(expected_errors.iter().cloned()).collect(),
);
let requests = object_names
.iter()
.map(|object_name| ObjectToDelete {
object_name: object_name.clone(),
..Default::default()
})
.collect();
let (deleted, errors) = store.delete_objects(&bucket, requests, ObjectOptions::default()).await;
assert_eq!(
injection.observed(),
object_names.len(),
"the source pool must first complete every real delete"
);
assert_eq!(
errors,
expected_errors.iter().cloned().map(Some).collect::<Vec<_>>(),
"a successful pool must not clear a source pool failure at any request index"
);
assert_eq!(
deleted.iter().map(|object| object.object_name.as_str()).collect::<Vec<_>>(),
object_names.iter().map(String::as_str).collect::<Vec<_>>(),
"DeleteObjects must preserve request index mapping while aggregating pool failures"
);
assert!(
deleted.iter().all(|object| object.found),
"the injected source results must retain real delete success data"
);
for pool in &store.pools {
for object_name in &object_names {
let error = pool
.get_object_info(&bucket, object_name, &ObjectOptions::default())
.await
.expect_err("both the active and source pool delete calls must execute");
assert!(
matches!(error, StorageError::ObjectNotFound(_, _)),
"unexpected residual object: {error:?}"
);
}
}
drop(injection);
}
shutdown.cancel();
}
+121
View File
@@ -838,6 +838,7 @@ struct DeleteAfterObjectLockSnapshotBarrierState {
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
namespace_pending: tokio::sync::Notify,
namespace_acquired: AtomicBool,
}
#[cfg(test)]
@@ -858,6 +859,7 @@ impl DeleteAfterObjectLockSnapshotBarrier {
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
namespace_pending: tokio::sync::Notify::new(),
namespace_acquired: AtomicBool::new(false),
});
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
@@ -883,6 +885,10 @@ impl DeleteAfterObjectLockSnapshotBarrier {
.await
.expect("delete should proceed to its namespace lock after leaving the snapshot barrier");
}
pub(crate) fn namespace_acquired(&self) -> bool {
self.state.namespace_acquired.load(Ordering::Acquire)
}
}
#[cfg(test)]
@@ -914,6 +920,20 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
}
}
#[cfg(test)]
fn notify_delete_namespace_acquired(bucket: &str) {
let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("delete snapshot barrier mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket)
.cloned();
if let Some(state) = state {
state.namespace_acquired.store(true, Ordering::Release);
}
}
#[cfg(test)]
struct VersionedDeleteMarkerCommitBarrierState {
bucket: String,
@@ -1048,6 +1068,88 @@ fn batch_delete_targets_pool(creates_latest_marker: bool, marker_target_pool_idx
!creates_latest_marker || marker_target_pool_idx == Some(pool_idx)
}
#[cfg(test)]
struct BatchDeletePoolErrorInjectionState {
bucket: String,
pool_idx: usize,
errors: std::collections::HashMap<String, Error>,
observed: std::sync::atomic::AtomicUsize,
}
#[cfg(test)]
pub(crate) struct BatchDeletePoolErrorInjection {
state: Arc<BatchDeletePoolErrorInjectionState>,
}
#[cfg(test)]
static BATCH_DELETE_POOL_ERROR_INJECTION: std::sync::OnceLock<std::sync::Mutex<Option<Arc<BatchDeletePoolErrorInjectionState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl BatchDeletePoolErrorInjection {
pub(crate) fn install(bucket: &str, pool_idx: usize, errors: Vec<(String, Error)>) -> Self {
let state = Arc::new(BatchDeletePoolErrorInjectionState {
bucket: bucket.to_string(),
pool_idx,
errors: errors.into_iter().collect(),
observed: std::sync::atomic::AtomicUsize::new(0),
});
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison");
assert!(slot.is_none(), "batch delete pool error injection must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn observed(&self) -> usize {
self.state.observed.load(Ordering::Acquire)
}
}
#[cfg(test)]
impl Drop for BatchDeletePoolErrorInjection {
fn drop(&mut self) {
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn inject_batch_delete_pool_errors(
bucket: &str,
pool_idx: usize,
object_names: &[String],
result: &mut (Vec<DeletedObject>, Vec<Option<Error>>),
) {
let state = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.pool_idx == pool_idx)
.cloned();
let Some(state) = state else {
return;
};
for (idx, object_name) in object_names.iter().enumerate() {
let Some(error) = state.errors.get(object_name) else {
continue;
};
if result.1[idx].is_none() && result.0[idx].found {
result.1[idx] = Some(error.clone());
state.observed.fetch_add(1, Ordering::AcqRel);
}
}
}
fn resolve_batch_delete_pool_results<'a>(
initial_error: Option<Error>,
pool_results: impl IntoIterator<Item = (&'a DeletedObject, &'a Option<Error>)>,
@@ -2674,6 +2776,10 @@ impl ECStore {
} else {
None
};
#[cfg(test)]
if _object_lock_guard.is_some() {
notify_delete_namespace_acquired(bucket);
}
if let Some(trigger) = opts.lifecycle_delete_all.as_ref() {
let configs = delete_all_configs.as_ref().ok_or(StorageError::PreconditionFailed)?;
let expected_bucket_incarnation_id = opts.expected_bucket_incarnation_id.ok_or(StorageError::PreconditionFailed)?;
@@ -3006,6 +3112,10 @@ impl ECStore {
Ok(guards) => guards,
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
};
#[cfg(test)]
if !_object_lock_guards.is_empty() {
notify_delete_namespace_acquired(bucket);
}
let delete_config_snapshot = opts
.delete_replication_config_snapshot
@@ -3057,7 +3167,18 @@ impl ECStore {
let pool_opts = opts.clone();
futures.push(async move {
#[cfg(test)]
let pool_object_names = pool_objects
.iter()
.map(|object| object.object_name.clone())
.collect::<Vec<_>>();
let result = pool.delete_objects(bucket, pool_objects, pool_opts).await;
#[cfg(test)]
let result = {
let mut result = result;
inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result);
result
};
(object_indices, result)
});
}