fix(ecstore): fence decommission commit loss

This commit is contained in:
overtrue
2026-08-22 03:11:16 +08:00
parent d168b69a47
commit ffe4085d91
7 changed files with 764 additions and 34 deletions
+39 -18
View File
@@ -856,7 +856,6 @@ fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo)
fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool {
is_unversioned_data_movement_object(source)
&& is_unversioned_data_movement_object(target)
&& !target.delete_marker
&& source
.mod_time
.zip(target.mod_time)
@@ -3640,25 +3639,47 @@ mod tests {
}
#[test]
fn test_precondition_conflict_rejects_newer_delete_marker() {
let source = ObjectInfo {
size: 128,
etag: Some("etag-source".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
};
let target = ObjectInfo {
delete_marker: true,
etag: None,
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
..source.clone()
};
fn test_precondition_conflict_accepts_only_newer_null_delete_marker() {
for version_id in [None, Some(Uuid::nil())] {
let source = ObjectInfo {
version_id,
size: 128,
etag: Some("etag-source".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
};
let target = ObjectInfo {
delete_marker: true,
etag: None,
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
..source.clone()
};
let should_resume =
resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(target)), &source, 0, 1)
.expect("delete marker conflict should be evaluated");
assert!(
resolve_data_movement_overwrite_resume_result(
&Error::PreconditionFailed,
Ok(Some(target.clone())),
&source,
0,
1,
)
.expect("newer null delete marker should be evaluated")
);
assert!(!should_resume);
let mut same_time = target.clone();
same_time.mod_time = source.mod_time;
assert!(
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(same_time)), &source, 0, 1,)
.expect("same-generation null delete marker should be rejected")
);
let mut versioned = target;
versioned.version_id = Some(Uuid::new_v4());
assert!(
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(versioned)), &source, 0, 1,)
.expect("a UUID delete marker must not erase a null source version")
);
}
}
#[test]
+20 -10
View File
@@ -24,7 +24,7 @@ use crate::storage_api_contracts::{
pub struct NamespaceLockFence {
signals: Arc<Vec<Arc<rustfs_lock::distributed_lock::LockLostSignal>>>,
#[cfg(test)]
forced_lost: Arc<std::sync::atomic::AtomicBool>,
forced_lost: Arc<Vec<Arc<std::sync::atomic::AtomicBool>>>,
}
impl Debug for NamespaceLockFence {
@@ -40,13 +40,17 @@ impl NamespaceLockFence {
Self {
signals: Arc::default(),
#[cfg(test)]
forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)),
forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]),
}
}
pub(crate) fn is_lock_lost(&self) -> bool {
#[cfg(test)]
if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
if self
.forced_lost
.iter()
.any(|lost| lost.load(std::sync::atomic::Ordering::Acquire))
{
return true;
}
self.signals.iter().any(|signal| signal.is_lost())
@@ -57,27 +61,26 @@ impl NamespaceLockFence {
}
fn extend(&mut self, other: &Self) {
if Arc::ptr_eq(&self.signals, &other.signals) {
return;
if !Arc::ptr_eq(&self.signals, &other.signals) {
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
}
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
#[cfg(test)]
if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
self.forced_lost.store(true, std::sync::atomic::Ordering::Release);
if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) {
Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned());
}
}
#[cfg(test)]
pub(crate) fn lost_for_test() -> Self {
let fence = Self::new();
fence.forced_lost.store(true, std::sync::atomic::Ordering::Release);
fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release);
fence
}
#[cfg(test)]
pub(crate) fn loss_handle_for_test() -> (Self, Arc<std::sync::atomic::AtomicBool>) {
let fence = Self::new();
(fence.clone(), Arc::clone(&fence.forced_lost))
(fence.clone(), Arc::clone(&fence.forced_lost[0]))
}
}
@@ -411,6 +414,13 @@ impl ObjectOptions {
self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new);
}
#[cfg(test)]
pub(crate) fn add_namespace_lock_fence_for_test(&mut self, fence: &NamespaceLockFence) {
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.extend(fence);
}
pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) {
self.lifecycle_delete_all_journal
.get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default())));
+4
View File
@@ -735,8 +735,12 @@ pub(crate) use core::io_primitives::disk_call_counters;
mod ctx;
mod metadata;
mod ops;
#[cfg(test)]
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
#[cfg(any(test, feature = "test-util"))]
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
#[cfg(test)]
pub(crate) use ops::object::DeleteObjectCommitBarrier;
#[cfg(feature = "test-util")]
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
pub(crate) use ops::object::body_cache_plaintext_len;
@@ -32,6 +32,8 @@ use crate::crash_inject::{self, CrashPoint};
use crate::multipart_listing::paginate_multipart_listing;
use futures::{StreamExt, stream};
use std::future::Future;
#[cfg(test)]
use std::sync::atomic::AtomicBool;
#[cfg(any(test, feature = "test-util"))]
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
@@ -65,6 +67,7 @@ impl StaleMultipartCleanupGuard {
#[cfg(any(test, feature = "test-util"))]
#[derive(Clone, Copy, PartialEq, Eq)]
pub enum MultipartCommitPause {
NewUploadBeforeLockLost,
PutPartBeforeLockAcquire,
PutPartBeforeLockLost,
PutPartAfterRename,
@@ -156,6 +159,72 @@ impl Drop for MultipartCommitBarrier {
}
}
#[cfg(test)]
struct NewMultipartUploadCommitObservationState {
bucket: String,
object: String,
committed: AtomicBool,
}
#[cfg(test)]
pub(crate) struct NewMultipartUploadCommitObservation {
state: Arc<NewMultipartUploadCommitObservationState>,
}
#[cfg(test)]
static NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<NewMultipartUploadCommitObservationState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl NewMultipartUploadCommitObservation {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(NewMultipartUploadCommitObservationState {
bucket: bucket.to_string(),
object: object.to_string(),
committed: AtomicBool::new(false),
});
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison");
assert!(slot.is_none(), "new multipart upload commit observation must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn committed(&self) -> bool {
self.state.committed.load(Ordering::Acquire)
}
}
#[cfg(test)]
impl Drop for NewMultipartUploadCommitObservation {
fn drop(&mut self) {
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn observe_new_multipart_upload_commit(bucket: &str, object: &str) {
let state = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.committed.store(true, Ordering::Release);
}
}
#[cfg(any(test, feature = "test-util"))]
async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
let barrier = {
@@ -1615,6 +1684,30 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement);
#[cfg(any(test, feature = "test-util"))]
pause_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost).await;
if _object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "new_multipart_upload_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "new_multipart_upload_outer_lock",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
Self::write_unique_file_info(
&shuffle_disks,
@@ -1626,6 +1719,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
)
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
#[cfg(test)]
observe_new_multipart_upload_commit(bucket, object);
// evalDisks
+4 -4
View File
@@ -4773,7 +4773,7 @@ struct DeleteObjectCommitBarrierState {
}
#[cfg(test)]
struct DeleteObjectCommitBarrier {
pub(crate) struct DeleteObjectCommitBarrier {
state: Arc<DeleteObjectCommitBarrierState>,
}
@@ -4783,7 +4783,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option
#[cfg(test)]
impl DeleteObjectCommitBarrier {
fn install(bucket: &str, object: &str) -> Self {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(DeleteObjectCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
@@ -4799,13 +4799,13 @@ impl DeleteObjectCommitBarrier {
Self { state }
}
async fn wait_until_paused(&self) {
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("delete object should reach the deterministic commit barrier");
}
fn release(&self) {
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
+493
View File
@@ -1299,6 +1299,139 @@ mod tests {
(source_version, expected_source_versions)
}
async fn mark_test_pool_decommissioning(store: &Arc<crate::store::ECStore>, pool_idx: usize) {
let mut pool_meta = store.pool_meta.write().await;
pool_meta.pools[pool_idx].decommission = Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::now_utc()),
..Default::default()
});
}
async fn write_decommission_test_multipart_source(
store: &Arc<crate::store::ECStore>,
pool_idx: usize,
bucket: &str,
object: &str,
) {
let pool = &store.pools[pool_idx];
let upload = pool
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("create decommission multipart source upload");
let first_part = vec![b'm'; 5 * 1024 * 1024];
let second_part = b"decommission multipart tail".to_vec();
let mut completed_parts = Vec::with_capacity(2);
for (part_number, body) in [(1, first_part), (2, second_part)] {
let mut reader = PutObjReader::from_vec(body);
let part = pool
.put_object_part(bucket, object, &upload.upload_id, part_number, &mut reader, &ObjectOptions::default())
.await
.expect("write decommission multipart source part");
completed_parts.push(crate::storage_api_contracts::multipart::CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
});
}
pool.clone()
.complete_multipart_upload(bucket, object, &upload.upload_id, completed_parts, &ObjectOptions::default())
.await
.expect("complete decommission multipart source object");
}
async fn assert_pool_object_present(pool: &Arc<crate::core::sets::Sets>, bucket: &str, object: &str) {
pool.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("expected object generation must remain present");
}
async fn assert_pool_object_absent(pool: &Arc<crate::core::sets::Sets>, bucket: &str, object: &str) {
let err = pool
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect_err("fenced decommission target must remain absent");
assert!(
matches!(err, StorageError::ObjectNotFound(_, _) | StorageError::VersionNotFound(_, _, _)),
"unexpected fenced target result: {err:?}"
);
}
async fn write_suspended_decommission_source(store: &Arc<crate::store::ECStore>, bucket: &str, object: &str) {
let mut reader = PutObjReader::from_vec(b"suspended source generation".to_vec());
let source = store.pools[0]
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
version_suspended: true,
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
},
)
.await
.expect("write suspended null source version");
assert!(
source.version_id.is_none_or(|version_id| version_id.is_nil()),
"suspended source must use the null version identity"
);
}
async fn assert_suspended_null_source_present(store: &Arc<crate::store::ECStore>, bucket: &str, object: &str) {
let versions = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("suspended source versions should be readable")
.expect("suspended source must exist before worker convergence");
assert!(
versions
.versions
.iter()
.any(|version| !version.deleted && version.version_id.is_none_or(|version_id| version_id.is_nil())),
"the source pool must retain its null data version while DELETE owns the fixed fence"
);
}
async fn assert_suspended_decommission_converged(store: &Arc<crate::store::ECStore>, bucket: &str, object: &str) {
let source_versions = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("source versions should remain readable after suspended convergence");
assert!(
source_versions.is_none_or(|versions| versions.versions.is_empty()),
"worker convergence must remove only the decommissioned source null version"
);
let target_versions = store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("active target versions should be readable")
.expect("active target must retain the suspended DELETE marker");
assert!(
matches!(target_versions.versions.as_slice(), [marker] if marker.deleted && marker.version_id.is_none_or(|version_id| version_id.is_nil())),
"active target must contain only its null delete marker: {target_versions:?}"
);
let err = store
.get_object_info(
bucket,
object,
&ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
.expect_err("the active null delete marker must hide the migrated source generation");
assert!(
matches!(err, StorageError::ObjectNotFound(_, _)),
"unexpected suspended latest-object result: {err:?}"
);
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn tag_updates_skip_active_rebalance_source_pool() {
@@ -2969,6 +3102,206 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn decommission_outer_fence_loss_blocks_target_put_commit() {
let temp_dir = tempfile::tempdir().expect("create decommission PUT fence-loss store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-put-fence-loss", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-put-fence-loss-{}", uuid::Uuid::new_v4());
let object = "ordinary.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission PUT fence-loss 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 decommission PUT source");
mark_test_pool_decommissioning(&store, 0).await;
let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install(
&bucket,
object,
crate::store::object::DecommissionMutationFenceTestPhase::Migration,
);
let barrier = crate::set_disk::PutObjectCommitBarrier::install(
&bucket,
object,
crate::set_disk::PutObjectCommitPause::BeforeQuotaRename,
);
let source_set = store.pools[0].get_disks_by_key(object);
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
worker_bucket,
source_set,
)
.await
});
barrier.wait_until_paused().await;
loss_hook.mark_lost();
barrier.release();
drop(barrier);
worker
.await
.expect("decommission PUT fence-loss worker should join")
.expect("a fenced migration failure should remain retryable at entry scope");
assert_pool_object_absent(&store.pools[1], &bucket, object).await;
assert_pool_object_present(&store.pools[0], &bucket, object).await;
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn decommission_outer_fence_loss_blocks_multipart_commits() {
let temp_dir = tempfile::tempdir().expect("create decommission multipart fence-loss store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-multipart-fence-loss", &[4, 4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-multipart-fence-loss-{}", uuid::Uuid::new_v4());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission multipart fence-loss bucket");
for object in ["new-upload.bin", "complete.bin"] {
write_decommission_test_multipart_source(&store, 0, &bucket, object).await;
}
mark_test_pool_decommissioning(&store, 0).await;
for (object, pause) in [
("new-upload.bin", crate::set_disk::MultipartCommitPause::NewUploadBeforeLockLost),
("complete.bin", crate::set_disk::MultipartCommitPause::BeforeLockLost),
] {
let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install(
&bucket,
object,
crate::store::object::DecommissionMutationFenceTestPhase::Migration,
);
let commit_observation = (pause == crate::set_disk::MultipartCommitPause::NewUploadBeforeLockLost)
.then(|| crate::set_disk::NewMultipartUploadCommitObservation::install(&bucket, object));
let barrier = crate::set_disk::MultipartCommitBarrier::install(&bucket, object, pause);
let source_set = store.pools[0].get_disks_by_key(object);
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
worker_bucket,
source_set,
)
.await
});
barrier.wait_until_paused().await;
loss_hook.mark_lost();
barrier.release();
drop(barrier);
worker
.await
.expect("decommission multipart fence-loss worker should join")
.expect("a fenced multipart migration failure should remain retryable at entry scope");
if let Some(commit_observation) = commit_observation {
assert!(
!commit_observation.committed(),
"new multipart upload metadata must not commit after the outer fence is lost"
);
}
assert_pool_object_absent(&store.pools[1], &bucket, object).await;
assert_pool_object_present(&store.pools[0], &bucket, object).await;
let uploads = store.pools[1]
.list_multipart_uploads(&bucket, object, None, None, None, 100)
.await
.expect("list target multipart uploads after fenced migration");
assert!(uploads.uploads.is_empty(), "fenced multipart migration must not retain target staging");
}
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn decommission_outer_fence_loss_blocks_source_cleanup_delete_commit() {
let temp_dir = tempfile::tempdir().expect("create decommission cleanup fence-loss store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-cleanup-fence-loss", &[4, 4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-cleanup-fence-loss-{}", uuid::Uuid::new_v4());
let object = "cleanup.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission cleanup fence-loss 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 decommission cleanup source");
mark_test_pool_decommissioning(&store, 0).await;
let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install(
&bucket,
object,
crate::store::object::DecommissionMutationFenceTestPhase::SourceCleanup,
);
let barrier = crate::set_disk::DeleteObjectCommitBarrier::install(&bucket, object);
let source_set = store.pools[0].get_disks_by_key(object);
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
worker_bucket,
source_set,
)
.await
});
barrier.wait_until_paused().await;
loss_hook.mark_lost();
barrier.release();
drop(barrier);
let err = worker
.await
.expect("decommission cleanup fence-loss worker should join")
.expect_err("source cleanup must fail after its outer fence is lost");
assert!(
err.to_string().contains("delete_object_commit"),
"cleanup failure must come from the delete commit fence: {err:?}"
);
assert_pool_object_present(&store.pools[0], &bucket, object).await;
assert_pool_object_present(&store.pools[1], &bucket, object).await;
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn reverse_decommission_reuses_fixed_target_fence_for_put_and_multipart() {
@@ -3609,6 +3942,166 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn suspended_delete_marker_then_decommission_worker_converges_null_source() {
let temp_dir = tempfile::tempdir().expect("create suspended decommission DELETE store dir");
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
temp_dir.path(),
"suspended-decommission-delete-convergence",
&[4, 4],
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("suspended-decommission-delete-convergence-{}", uuid::Uuid::new_v4());
let object = "single.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create suspended decommission DELETE bucket");
write_suspended_decommission_source(&store, &bucket, object).await;
mark_test_pool_decommissioning(&store, 0).await;
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_object(
&delete_bucket,
object,
ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
});
delete_barrier.wait_until_paused().await;
assert_suspended_null_source_present(&store, &bucket, object).await;
let source_set = store.pools[0].get_disks_by_key(object);
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
worker_bucket,
source_set,
)
.await
});
delete_barrier.release();
let marker = delete
.await
.expect("suspended DELETE task should join")
.expect("suspended DELETE should commit its active-pool marker");
drop(delete_barrier);
assert!(marker.delete_marker, "suspended DELETE must create a marker");
assert!(
marker.version_id.is_none_or(|version_id| version_id.is_nil()),
"suspended DELETE marker must keep the null version identity"
);
worker
.await
.expect("suspended decommission worker should join")
.expect("worker must treat the newer active null marker as a completed migration");
assert_suspended_decommission_converged(&store, &bucket, object).await;
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn suspended_batch_delete_marker_then_decommission_worker_converges_null_source() {
let temp_dir = tempfile::tempdir().expect("create suspended batch decommission DELETE store dir");
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
temp_dir.path(),
"suspended-batch-decommission-delete-convergence",
&[4, 4],
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("suspended-batch-decommission-delete-convergence-{}", uuid::Uuid::new_v4());
let object = "batch.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create suspended batch decommission DELETE bucket");
write_suspended_decommission_source(&store, &bucket, object).await;
mark_test_pool_decommissioning(&store, 0).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::SUSPENDED)),
..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;
assert_suspended_null_source_present(&store, &bucket, object).await;
let source_set = store.pools[0].get_disks_by_key(object);
let worker_store = Arc::clone(&store);
let worker_bucket = bucket.clone();
let worker = tokio::spawn(async move {
worker_store
.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
worker_bucket,
source_set,
)
.await
});
delete_barrier.release();
let (deleted, errors) = delete.await.expect("suspended batch DELETE task should join");
drop(delete_barrier);
assert!(errors.iter().all(Option::is_none), "suspended batch DELETE should succeed: {errors:?}");
assert!(
matches!(deleted.as_slice(), [marker] if marker.delete_marker && marker.delete_marker_version_id.is_none_or(|version_id| version_id.is_nil())),
"suspended batch DELETE must create one null marker: {deleted:?}"
);
worker
.await
.expect("suspended batch decommission worker should join")
.expect("worker must treat the newer batch null marker as a completed migration");
assert_suspended_decommission_converged(&store, &bucket, object).await;
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
+109 -2
View File
@@ -353,6 +353,8 @@ impl fmt::Display for ObjectLockDiagMode {
pub(crate) struct ObjectLockDiagGuard {
guard: rustfs_lock::NamespaceLockGuard,
#[cfg(test)]
test_namespace_lock_fence: Option<NamespaceLockFence>,
enabled: bool,
op: &'static str,
bucket: Option<String>,
@@ -374,6 +376,8 @@ impl ObjectLockDiagGuard {
) -> Self {
Self {
guard,
#[cfg(test)]
test_namespace_lock_fence: None,
enabled,
op,
bucket,
@@ -400,9 +404,92 @@ impl ObjectLockDiagGuard {
if let Some(signal) = self.lock_lost_signal() {
opts.add_namespace_lock_lost_signal(signal);
}
#[cfg(test)]
if let Some(fence) = self.test_namespace_lock_fence.as_ref() {
opts.add_namespace_lock_fence_for_test(fence);
}
}
}
#[cfg(test)]
#[derive(Clone, Copy, PartialEq, Eq)]
pub(crate) enum DecommissionMutationFenceTestPhase {
Migration,
SourceCleanup,
}
#[cfg(test)]
struct DecommissionMutationFenceLossState {
bucket: String,
object: String,
phase: DecommissionMutationFenceTestPhase,
fence: NamespaceLockFence,
loss_handle: Arc<std::sync::atomic::AtomicBool>,
}
#[cfg(test)]
pub(crate) struct DecommissionMutationFenceLossHook {
state: Arc<DecommissionMutationFenceLossState>,
}
#[cfg(test)]
static DECOMMISSION_MUTATION_FENCE_LOSS_HOOK: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DecommissionMutationFenceLossState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl DecommissionMutationFenceLossHook {
pub(crate) fn install(bucket: &str, object: &str, phase: DecommissionMutationFenceTestPhase) -> Self {
let (fence, loss_handle) = NamespaceLockFence::loss_handle_for_test();
let state = Arc::new(DecommissionMutationFenceLossState {
bucket: bucket.to_string(),
object: object.to_string(),
phase,
fence,
loss_handle,
});
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison");
assert!(slot.is_none(), "decommission mutation fence loss hook must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn mark_lost(&self) {
self.state.loss_handle.store(true, Ordering::Release);
}
}
#[cfg(test)]
impl Drop for DecommissionMutationFenceLossHook {
fn drop(&mut self) {
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison");
if slot.as_ref().is_some_and(|hook| Arc::ptr_eq(hook, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn decommission_mutation_fence_for_test(
bucket: &str,
object: &str,
phase: DecommissionMutationFenceTestPhase,
) -> Option<NamespaceLockFence> {
DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison")
.as_ref()
.filter(|hook| hook.bucket == bucket && hook.object == object && hook.phase == phase)
.map(|hook| hook.fence.clone())
}
pub(crate) struct SourceCleanupMutationFence {
guard: ObjectLockDiagGuard,
source_lock_covered: bool,
@@ -1827,11 +1914,22 @@ impl ECStore {
return Err(Error::other("decommission object migration requires namespace locking"));
}
#[cfg(test)]
let test_namespace_lock_fence =
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::Migration);
let object = encode_dir_object(object);
let mut opts = ObjectOptions::default();
self.acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts)
let guard = 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"))
.ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))?;
#[cfg(test)]
let guard = {
let mut guard = guard;
guard.test_namespace_lock_fence = test_namespace_lock_fence;
guard
};
Ok(guard)
}
pub(super) fn apply_decommission_target_mutation_fence(
@@ -1865,6 +1963,9 @@ impl ECStore {
#[cfg(test)]
crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object);
#[cfg(test)]
let test_namespace_lock_fence =
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup);
let object = encode_dir_object(object);
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
// SetDisks namespaces include pool/set identity, so only the canonical
@@ -1875,6 +1976,12 @@ impl ECStore {
let guard = self
.acquire_object_write_lock("decommission_source_cleanup", bucket, &object)
.await?;
#[cfg(test)]
let guard = {
let mut guard = guard;
guard.test_namespace_lock_fence = test_namespace_lock_fence;
guard
};
Ok(SourceCleanupMutationFence {
guard,