fix(ecstore): fence deletes against decommission commits

This commit is contained in:
overtrue
2026-08-21 23:16:21 +08:00
parent 98c4675617
commit 01a0a8b66b
4 changed files with 227 additions and 8 deletions
+105
View File
@@ -2752,6 +2752,111 @@ mod tests {
.expect_err("suspended delete must remove the requested UUID version");
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn delete_waits_for_decommission_commit_then_removes_every_copy() {
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;
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";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission delete-fence bucket");
let mut source = PutObjReader::from_vec(b"source generation".to_vec());
store.pools[0]
.put_object(&bucket, object, &mut source, &ObjectOptions::default())
.await
.expect("write source object to the pool being decommissioned");
{
let mut pool_meta = store.pool_meta.write().await;
pool_meta.pools[0].decommission = Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::now_utc()),
..Default::default()
});
}
assert!(store.is_suspended(0).await, "pool 0 must be a suspended decommission source");
let barrier = crate::set_disk::PutObjectCommitBarrier::install(
&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 {
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_decommission_delete_fence",
)
.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 mut delete = tokio::spawn(async move {
delete_store
.delete_object(&delete_bucket, object, ObjectOptions::default())
.await
});
delete_barrier.wait_until_paused().await;
delete_barrier.release();
assert!(
tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(),
"DELETE must wait while the decommission source generation is being committed"
);
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");
for pool in &store.pools {
let err = pool
.get_object_info(&bucket, object, &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())
.await
.expect_err("the migrated source generation must not become visible again");
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
+68 -7
View File
@@ -393,6 +393,14 @@ impl ObjectLockDiagGuard {
pub(crate) fn is_lock_lost(&self) -> bool {
self.guard.is_lock_lost()
}
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.no_lock = true;
opts.ensure_namespace_lock_fence();
if let Some(signal) = self.lock_lost_signal() {
opts.add_namespace_lock_lost_signal(signal);
}
}
}
/// Opaque write-lock guard for the RestoreObject accept path; see
@@ -410,10 +418,7 @@ impl RestoreAcceptGuard {
}
pub fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.ensure_namespace_lock_fence();
if let Some(signal) = self.0.lock_lost_signal() {
opts.add_namespace_lock_lost_signal(signal);
}
self.0.add_namespace_lock_fence(opts);
}
}
@@ -913,6 +918,16 @@ fn writer_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions
lookup_opts
}
fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
let mut lookup_opts = writer_pool_lookup_opts(opts, no_lock);
lookup_opts.skip_decommissioned = opts.data_movement;
lookup_opts
}
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 transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
let mut lookup_opts = opts.clone();
lookup_opts.skip_decommissioned = true;
@@ -1533,6 +1548,22 @@ impl ECStore {
)))
}
pub(crate) async fn acquire_decommission_object_mutation_fence(
&self,
bucket: &str,
object: &str,
) -> Result<ObjectLockDiagGuard> {
if self.ctx.lock_manager().is_disabled() {
return Err(Error::other("decommission object migration requires namespace locking"));
}
let object = encode_dir_object(object);
let mut opts = ObjectOptions::default();
self.acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts)
.await?
.ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))
}
pub(crate) async fn acquire_all_object_read_locks(
&self,
op: &'static str,
@@ -2479,7 +2510,7 @@ impl ECStore {
return Ok(ObjectInfo::default());
}
let gopts = writer_pool_lookup_opts(&opts, true);
let gopts = delete_pool_lookup_opts(&opts, true);
if opts.data_movement {
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
@@ -2622,7 +2653,7 @@ impl ECStore {
None
};
if !errs.is_empty() && !opts.versioned && !opts.version_suspended {
if should_delete_from_all_pools(&opts, errs.len()) {
let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await {
Ok(obj) => obj,
Err(err) => {
@@ -2640,7 +2671,7 @@ impl ECStore {
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
if self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
@@ -4433,6 +4464,36 @@ mod tests {
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
}
#[test]
fn ordinary_delete_lookup_includes_decommission_source_and_skips_rebalance_source() {
let lookup_opts = delete_pool_lookup_opts(&ObjectOptions::default(), true);
assert!(lookup_opts.no_lock);
assert!(!lookup_opts.skip_decommissioned);
assert!(lookup_opts.skip_rebalancing);
}
#[test]
fn delete_fans_out_for_unversioned_and_explicit_version_mutations() {
assert!(should_delete_from_all_pools(&ObjectOptions::default(), 1));
assert!(should_delete_from_all_pools(
&ObjectOptions {
versioned: true,
version_id: Some(uuid::Uuid::new_v4().to_string()),
..Default::default()
},
2,
));
assert!(!should_delete_from_all_pools(
&ObjectOptions {
versioned: true,
..Default::default()
},
1,
));
assert!(!should_delete_from_all_pools(&ObjectOptions::default(), 0));
}
#[test]
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
let lookup_opts = data_movement_pool_lookup_opts(