mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-25 13:36:50 +00:00
fix(ecstore): preserve decommission target write locks
This commit is contained in:
@@ -2817,15 +2817,15 @@ mod tests {
|
|||||||
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
|
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
|
||||||
let delete_store = Arc::clone(&store);
|
let delete_store = Arc::clone(&store);
|
||||||
let delete_bucket = bucket.clone();
|
let delete_bucket = bucket.clone();
|
||||||
let mut delete = tokio::spawn(async move {
|
let delete = tokio::spawn(async move {
|
||||||
delete_store
|
delete_store
|
||||||
.delete_object(&delete_bucket, object, ObjectOptions::default())
|
.delete_object(&delete_bucket, object, ObjectOptions::default())
|
||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
delete_barrier.wait_until_paused().await;
|
delete_barrier.wait_until_paused().await;
|
||||||
delete_barrier.release();
|
delete_barrier.release_and_wait_until_namespace_pending().await;
|
||||||
assert!(
|
assert!(
|
||||||
tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(),
|
!delete.is_finished(),
|
||||||
"DELETE must wait while the decommission source generation is being committed"
|
"DELETE must wait while the decommission source generation is being committed"
|
||||||
);
|
);
|
||||||
|
|
||||||
@@ -2857,6 +2857,147 @@ mod tests {
|
|||||||
shutdown.cancel();
|
shutdown.cancel();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn versioned_delete_waits_for_decommission_commit_then_publishes_marker() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create versioned decommission delete-fence store dir");
|
||||||
|
let (_ctx, store, shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "versioned-decommission-delete-fence", &[4, 4]))
|
||||||
|
.await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||||
|
|
||||||
|
let bucket = format!("versioned-decommission-delete-fence-{}", uuid::Uuid::new_v4());
|
||||||
|
let object = "object.bin";
|
||||||
|
store
|
||||||
|
.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 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.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;
|
||||||
|
|
||||||
|
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
|
||||||
|
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 {
|
||||||
|
versioned: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
delete_barrier.wait_until_paused().await;
|
||||||
|
delete_barrier.release_and_wait_until_namespace_pending().await;
|
||||||
|
assert!(
|
||||||
|
!delete.is_finished(),
|
||||||
|
"versioned DELETE must wait while the source version is being committed"
|
||||||
|
);
|
||||||
|
|
||||||
|
barrier.release();
|
||||||
|
migration
|
||||||
|
.await
|
||||||
|
.expect("versioned decommission migration task should join")
|
||||||
|
.expect("versioned decommission migration should commit before DELETE");
|
||||||
|
let marker = delete
|
||||||
|
.await
|
||||||
|
.expect("versioned DELETE task should join")
|
||||||
|
.expect("versioned DELETE should publish a delete marker after migration");
|
||||||
|
assert!(marker.delete_marker, "versioned DELETE must publish a delete marker");
|
||||||
|
assert!(
|
||||||
|
marker.version_id.is_some_and(|version_id| !version_id.is_nil()),
|
||||||
|
"the delete marker must have a non-nil version ID"
|
||||||
|
);
|
||||||
|
|
||||||
|
let err = store
|
||||||
|
.get_object_info(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
versioned: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect_err("the post-migration delete marker must hide the migrated version");
|
||||||
|
assert!(
|
||||||
|
matches!(err, StorageError::ObjectNotFound(_, _)),
|
||||||
|
"unexpected latest-version result: {err:?}"
|
||||||
|
);
|
||||||
|
store
|
||||||
|
.get_object_info(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
versioned: true,
|
||||||
|
version_id: Some(source_version.to_string()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("the migrated source version must remain addressable below the delete marker");
|
||||||
|
|
||||||
|
shutdown.cancel();
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
|||||||
@@ -395,7 +395,6 @@ impl ObjectLockDiagGuard {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
||||||
opts.no_lock = true;
|
|
||||||
opts.ensure_namespace_lock_fence();
|
opts.ensure_namespace_lock_fence();
|
||||||
if let Some(signal) = self.lock_lost_signal() {
|
if let Some(signal) = self.lock_lost_signal() {
|
||||||
opts.add_namespace_lock_lost_signal(signal);
|
opts.add_namespace_lock_lost_signal(signal);
|
||||||
@@ -818,6 +817,7 @@ struct DeleteAfterObjectLockSnapshotBarrierState {
|
|||||||
bucket: String,
|
bucket: String,
|
||||||
arrived: tokio::sync::Notify,
|
arrived: tokio::sync::Notify,
|
||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
|
namespace_pending: tokio::sync::Notify,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -837,6 +837,7 @@ impl DeleteAfterObjectLockSnapshotBarrier {
|
|||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
arrived: tokio::sync::Notify::new(),
|
arrived: tokio::sync::Notify::new(),
|
||||||
release: tokio::sync::Notify::new(),
|
release: tokio::sync::Notify::new(),
|
||||||
|
namespace_pending: tokio::sync::Notify::new(),
|
||||||
});
|
});
|
||||||
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
|
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
@@ -854,6 +855,14 @@ impl DeleteAfterObjectLockSnapshotBarrier {
|
|||||||
pub(crate) fn release(&self) {
|
pub(crate) fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn release_and_wait_until_namespace_pending(&self) {
|
||||||
|
let namespace_pending = self.state.namespace_pending.notified();
|
||||||
|
self.release();
|
||||||
|
tokio::time::timeout(Duration::from_secs(5), namespace_pending)
|
||||||
|
.await
|
||||||
|
.expect("delete should proceed to its namespace lock after leaving the snapshot barrier");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -881,6 +890,7 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
|
|||||||
if let Some(state) = state {
|
if let Some(state) = state {
|
||||||
state.arrived.notify_one();
|
state.arrived.notify_one();
|
||||||
state.release.notified().await;
|
state.release.notified().await;
|
||||||
|
state.namespace_pending.notify_one();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user