mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 20:36:38 +00:00
fix(ecstore): preserve delete markers during source cleanup
This commit is contained in:
@@ -3275,6 +3275,9 @@ impl ECStore {
|
|||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
|
let source_cleanup_mutation_fence = self
|
||||||
|
.acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref())
|
||||||
|
.await?;
|
||||||
let cleanup_result = data_movement::cleanup_source_entry_if_unchanged(
|
let cleanup_result = data_movement::cleanup_source_entry_if_unchanged(
|
||||||
set.clone(),
|
set.clone(),
|
||||||
bucket.as_str(),
|
bucket.as_str(),
|
||||||
@@ -3286,6 +3289,7 @@ impl ECStore {
|
|||||||
lifecycle_guard: bucket_incarnation_fence
|
lifecycle_guard: bucket_incarnation_fence
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.and_then(|guard| guard.namespace_lock_guard()),
|
.and_then(|guard| guard.namespace_lock_guard()),
|
||||||
|
object_mutation_fence: Some(&source_cleanup_mutation_fence),
|
||||||
},
|
},
|
||||||
"decommission",
|
"decommission",
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -26,8 +26,7 @@ use crate::storage_api_contracts::{
|
|||||||
namespace::NamespaceLocking as _,
|
namespace::NamespaceLocking as _,
|
||||||
object::{HTTPPreconditions, ObjectOperations as _},
|
object::{HTTPPreconditions, ObjectOperations as _},
|
||||||
};
|
};
|
||||||
use crate::store::ECStore;
|
use crate::store::{ECStore, ObjectLockDiagGuard, SourceCleanupMutationFence};
|
||||||
use crate::store::ObjectLockDiagGuard;
|
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo};
|
use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo};
|
||||||
use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex};
|
use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex};
|
||||||
@@ -1029,6 +1028,7 @@ pub(crate) enum SourceCleanupError {
|
|||||||
pub(crate) struct SourceCleanupBucketFence<'a> {
|
pub(crate) struct SourceCleanupBucketFence<'a> {
|
||||||
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
|
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
|
||||||
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
|
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
|
||||||
|
pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>,
|
||||||
}
|
}
|
||||||
|
|
||||||
fn ensure_source_cleanup_versions_match(
|
fn ensure_source_cleanup_versions_match(
|
||||||
@@ -1066,7 +1066,9 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
|
|||||||
struct SourceCleanupDeleteBarrierState {
|
struct SourceCleanupDeleteBarrierState {
|
||||||
bucket: String,
|
bucket: String,
|
||||||
object: String,
|
object: String,
|
||||||
|
fence_pending: tokio::sync::Notify,
|
||||||
arrived: tokio::sync::Notify,
|
arrived: tokio::sync::Notify,
|
||||||
|
is_paused: AtomicBool,
|
||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1080,7 +1082,7 @@ pub(crate) struct SourceCleanupDeleteBarrier {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
|
static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<SourceCleanupDeleteBarrierState>>>> =
|
||||||
std::sync::OnceLock::new();
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -1093,15 +1095,22 @@ impl SourceCleanupDeleteBarrier {
|
|||||||
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
object: object.to_string(),
|
object: object.to_string(),
|
||||||
|
fence_pending: tokio::sync::Notify::new(),
|
||||||
arrived: tokio::sync::Notify::new(),
|
arrived: tokio::sync::Notify::new(),
|
||||||
|
is_paused: AtomicBool::new(false),
|
||||||
release: tokio::sync::Notify::new(),
|
release: tokio::sync::Notify::new(),
|
||||||
});
|
});
|
||||||
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison");
|
.expect("source cleanup delete barrier mutex should not poison");
|
||||||
assert!(slot.is_none(), "source cleanup delete barrier must be unique");
|
assert!(
|
||||||
*slot = Some(Arc::clone(&state));
|
!barriers
|
||||||
|
.iter()
|
||||||
|
.any(|barrier| barrier.bucket == bucket && barrier.object == object),
|
||||||
|
"source cleanup delete barrier must be unique per object"
|
||||||
|
);
|
||||||
|
barriers.push(Arc::clone(&state));
|
||||||
Self { state }
|
Self { state }
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1111,35 +1120,58 @@ impl SourceCleanupDeleteBarrier {
|
|||||||
.expect("source cleanup should reach the pre-delete barrier");
|
.expect("source cleanup should reach the pre-delete barrier");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn wait_until_fence_pending(&self) {
|
||||||
|
tokio::time::timeout(StdDuration::from_secs(30), self.state.fence_pending.notified())
|
||||||
|
.await
|
||||||
|
.expect("source cleanup should attempt the fixed mutation fence");
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn is_paused(&self) -> bool {
|
||||||
|
self.state.is_paused.load(Ordering::Acquire)
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn release(&self) {
|
pub(crate) fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object: &str) {
|
||||||
|
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||||
|
.lock()
|
||||||
|
.expect("source cleanup delete barrier mutex should not poison")
|
||||||
|
.iter()
|
||||||
|
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
|
||||||
|
.cloned();
|
||||||
|
if let Some(barrier) = barrier {
|
||||||
|
barrier.fence_pending.notify_one();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
impl Drop for SourceCleanupDeleteBarrier {
|
impl Drop for SourceCleanupDeleteBarrier {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison");
|
.expect("source cleanup delete barrier mutex should not poison");
|
||||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
barriers.retain(|state| !Arc::ptr_eq(state, &self.state));
|
||||||
*slot = None;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
|
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
|
||||||
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
|
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison")
|
.expect("source cleanup delete barrier mutex should not poison")
|
||||||
.as_ref()
|
.iter()
|
||||||
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
|
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
|
||||||
.cloned();
|
.cloned();
|
||||||
if let Some(barrier) = barrier {
|
if let Some(barrier) = barrier {
|
||||||
|
barrier.is_paused.store(true, Ordering::Release);
|
||||||
barrier.arrived.notify_one();
|
barrier.arrived.notify_one();
|
||||||
barrier.release.notified().await;
|
barrier.release.notified().await;
|
||||||
}
|
}
|
||||||
@@ -1155,11 +1187,20 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
op_label: &str,
|
op_label: &str,
|
||||||
) -> std::result::Result<ObjectInfo, SourceCleanupError> {
|
) -> std::result::Result<ObjectInfo, SourceCleanupError> {
|
||||||
let cleanup_key = encode_dir_object(object);
|
let cleanup_key = encode_dir_object(object);
|
||||||
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
|
let source_guard = if bucket_fence
|
||||||
let _guard = ns_lock
|
.object_mutation_fence
|
||||||
.get_write_lock(get_lock_acquire_timeout())
|
.is_some_and(SourceCleanupMutationFence::source_lock_covered)
|
||||||
.await
|
{
|
||||||
.map_err(Error::from)?;
|
None
|
||||||
|
} else {
|
||||||
|
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
|
||||||
|
Some(
|
||||||
|
ns_lock
|
||||||
|
.get_write_lock(get_lock_acquire_timeout())
|
||||||
|
.await
|
||||||
|
.map_err(Error::from)?,
|
||||||
|
)
|
||||||
|
};
|
||||||
|
|
||||||
if bucket_fence
|
if bucket_fence
|
||||||
.lifecycle_guard
|
.lifecycle_guard
|
||||||
@@ -1169,6 +1210,14 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
"{op_label}: bucket incarnation fence was lost before source cleanup"
|
"{op_label}: bucket incarnation fence was lost before source cleanup"
|
||||||
))));
|
))));
|
||||||
}
|
}
|
||||||
|
if bucket_fence
|
||||||
|
.object_mutation_fence
|
||||||
|
.is_some_and(SourceCleanupMutationFence::is_lock_lost)
|
||||||
|
{
|
||||||
|
return Err(SourceCleanupError::Storage(Error::other(format!(
|
||||||
|
"{op_label}: object mutation fence was lost before source cleanup"
|
||||||
|
))));
|
||||||
|
}
|
||||||
|
|
||||||
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
|
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
|
||||||
|
|
||||||
@@ -1183,7 +1232,12 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id,
|
expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
opts.add_namespace_lock_guard(&_guard);
|
if let Some(source_guard) = source_guard.as_ref() {
|
||||||
|
opts.add_namespace_lock_guard(source_guard);
|
||||||
|
}
|
||||||
|
if let Some(object_mutation_fence) = bucket_fence.object_mutation_fence {
|
||||||
|
object_mutation_fence.add_namespace_lock_fence(&mut opts);
|
||||||
|
}
|
||||||
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
|
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
|
||||||
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
|
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -334,6 +334,7 @@ impl ECStore {
|
|||||||
lifecycle_guard: bucket_incarnation_fence
|
lifecycle_guard: bucket_incarnation_fence
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.and_then(|guard| guard.namespace_lock_guard()),
|
.and_then(|guard| guard.namespace_lock_guard()),
|
||||||
|
..Default::default()
|
||||||
},
|
},
|
||||||
"rebalance",
|
"rebalance",
|
||||||
),
|
),
|
||||||
|
|||||||
@@ -11636,6 +11636,7 @@ mod transition_upload_integrity_tests {
|
|||||||
crate::data_movement::SourceCleanupBucketFence {
|
crate::data_movement::SourceCleanupBucketFence {
|
||||||
expected_incarnation_id: None,
|
expected_incarnation_id: None,
|
||||||
lifecycle_guard: Some(&bucket_guard),
|
lifecycle_guard: Some(&bucket_guard),
|
||||||
|
..Default::default()
|
||||||
},
|
},
|
||||||
"test_data_movement",
|
"test_data_movement",
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -2859,7 +2859,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
async fn versioned_delete_waits_for_decommission_commit_then_publishes_marker() {
|
async fn versioned_delete_marker_survives_decommission_source_cleanup() {
|
||||||
let temp_dir = tempfile::tempdir().expect("create versioned decommission delete-fence store dir");
|
let temp_dir = tempfile::tempdir().expect("create versioned decommission delete-fence store dir");
|
||||||
let (_ctx, store, shutdown) =
|
let (_ctx, store, shutdown) =
|
||||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "versioned-decommission-delete-fence", &[4, 4]))
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "versioned-decommission-delete-fence", &[4, 4]))
|
||||||
@@ -2886,6 +2886,12 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
.expect("write source version to the pool being decommissioned");
|
.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 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;
|
let mut pool_meta = store.pool_meta.write().await;
|
||||||
pool_meta.pools[0].decommission = Some(PoolDecommissionInfo {
|
pool_meta.pools[0].decommission = Some(PoolDecommissionInfo {
|
||||||
@@ -2929,8 +2935,13 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
barrier.wait_until_paused().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");
|
||||||
|
|
||||||
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
|
let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object);
|
||||||
let delete_store = Arc::clone(&store);
|
let delete_store = Arc::clone(&store);
|
||||||
let delete_bucket = bucket.clone();
|
let delete_bucket = bucket.clone();
|
||||||
let delete = tokio::spawn(async move {
|
let delete = tokio::spawn(async move {
|
||||||
@@ -2946,17 +2957,46 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
delete_barrier.wait_until_paused().await;
|
delete_barrier.wait_until_paused().await;
|
||||||
delete_barrier.release_and_wait_until_namespace_pending().await;
|
let cleanup_set = store.pools[0].get_disks_by_key(object);
|
||||||
|
crate::data_movement::ensure_source_cleanup_versions_unchanged(
|
||||||
|
Arc::clone(&cleanup_set),
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&expected_source_versions,
|
||||||
|
&[],
|
||||||
|
"test_versioned_decommission_delete_fence",
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("the committed delete marker must not be published to the suspended source pool");
|
||||||
|
|
||||||
|
let cleanup_delete_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, object);
|
||||||
|
let cleanup_store = Arc::clone(&store);
|
||||||
|
let cleanup_bucket = bucket.clone();
|
||||||
|
let cleanup = tokio::spawn(async move {
|
||||||
|
let mutation_fence = cleanup_store
|
||||||
|
.acquire_decommission_source_cleanup_fence(&cleanup_bucket, object, cleanup_set.as_ref())
|
||||||
|
.await?;
|
||||||
|
crate::data_movement::cleanup_source_entry_if_unchanged(
|
||||||
|
cleanup_set,
|
||||||
|
&cleanup_bucket,
|
||||||
|
object,
|
||||||
|
&expected_source_versions,
|
||||||
|
&[],
|
||||||
|
crate::data_movement::SourceCleanupBucketFence {
|
||||||
|
object_mutation_fence: Some(&mutation_fence),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
"test_versioned_decommission_delete_fence",
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
cleanup_delete_barrier.wait_until_fence_pending().await;
|
||||||
assert!(
|
assert!(
|
||||||
!delete.is_finished(),
|
!cleanup_delete_barrier.is_paused(),
|
||||||
"versioned DELETE must wait while the source version is being committed"
|
"source cleanup must wait for the versioned DELETE mutation fence"
|
||||||
);
|
);
|
||||||
|
|
||||||
barrier.release();
|
delete_barrier.release();
|
||||||
migration
|
|
||||||
.await
|
|
||||||
.expect("versioned decommission migration task should join")
|
|
||||||
.expect("versioned decommission migration should commit before DELETE");
|
|
||||||
let marker = delete
|
let marker = delete
|
||||||
.await
|
.await
|
||||||
.expect("versioned DELETE task should join")
|
.expect("versioned DELETE task should join")
|
||||||
@@ -2967,6 +3007,13 @@ mod tests {
|
|||||||
"the delete marker must have a non-nil version ID"
|
"the delete marker must have a non-nil version ID"
|
||||||
);
|
);
|
||||||
|
|
||||||
|
cleanup_delete_barrier.wait_until_paused().await;
|
||||||
|
cleanup_delete_barrier.release();
|
||||||
|
cleanup
|
||||||
|
.await
|
||||||
|
.expect("source cleanup task should join")
|
||||||
|
.expect("source cleanup should preserve the active-pool delete marker");
|
||||||
|
|
||||||
let err = store
|
let err = store
|
||||||
.get_object_info(
|
.get_object_info(
|
||||||
&bucket,
|
&bucket,
|
||||||
@@ -2994,6 +3041,18 @@ mod tests {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("the migrated source version must remain addressable below the delete marker");
|
.expect("the migrated source version must remain addressable below the delete marker");
|
||||||
|
store.pools[0]
|
||||||
|
.get_object_info(
|
||||||
|
&bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
versioned: true,
|
||||||
|
version_id: Some(source_version.to_string()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect_err("source cleanup must remove the decommissioned source versions");
|
||||||
|
|
||||||
shutdown.cancel();
|
shutdown.cancel();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
|
|||||||
pub(crate) mod list_objects;
|
pub(crate) mod list_objects;
|
||||||
mod multipart;
|
mod multipart;
|
||||||
mod object;
|
mod object;
|
||||||
pub(crate) use object::ObjectLockDiagGuard;
|
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
|
||||||
pub use object::{
|
pub use object::{
|
||||||
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
|
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
|
||||||
SnapshotConsistencyError,
|
SnapshotConsistencyError,
|
||||||
|
|||||||
@@ -36,7 +36,7 @@ use crate::bucket::replication::ReplicationObjectBridge;
|
|||||||
use crate::disk::OldCurrentSize;
|
use crate::disk::OldCurrentSize;
|
||||||
use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot};
|
use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot};
|
||||||
use crate::set_disk::{
|
use crate::set_disk::{
|
||||||
get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
|
SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
|
||||||
is_lock_optimization_enabled, is_object_lock_diag_enabled,
|
is_lock_optimization_enabled, is_object_lock_diag_enabled,
|
||||||
};
|
};
|
||||||
use crate::storage_api_contracts::{
|
use crate::storage_api_contracts::{
|
||||||
@@ -402,6 +402,25 @@ impl ObjectLockDiagGuard {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) struct SourceCleanupMutationFence {
|
||||||
|
guard: ObjectLockDiagGuard,
|
||||||
|
source_lock_covered: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SourceCleanupMutationFence {
|
||||||
|
pub(crate) fn source_lock_covered(&self) -> bool {
|
||||||
|
self.source_lock_covered
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn is_lock_lost(&self) -> bool {
|
||||||
|
self.guard.is_lock_lost()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
||||||
|
self.guard.add_namespace_lock_fence(opts);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Opaque write-lock guard for the RestoreObject accept path; see
|
/// Opaque write-lock guard for the RestoreObject accept path; see
|
||||||
/// [`ECStore::acquire_restore_accept_guard`]. Deliberately not a general lock
|
/// [`ECStore::acquire_restore_accept_guard`]. Deliberately not a general lock
|
||||||
/// API — it only exists so the accept path's restore-status compare-and-set
|
/// API — it only exists so the accept path's restore-status compare-and-set
|
||||||
@@ -894,6 +913,82 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
struct VersionedDeleteMarkerCommitBarrierState {
|
||||||
|
bucket: String,
|
||||||
|
object: String,
|
||||||
|
arrived: tokio::sync::Notify,
|
||||||
|
release: tokio::sync::Notify,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) struct VersionedDeleteMarkerCommitBarrier {
|
||||||
|
state: Arc<VersionedDeleteMarkerCommitBarrierState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static VERSIONED_DELETE_MARKER_COMMIT_BARRIER: std::sync::OnceLock<
|
||||||
|
std::sync::Mutex<Option<Arc<VersionedDeleteMarkerCommitBarrierState>>>,
|
||||||
|
> = std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl VersionedDeleteMarkerCommitBarrier {
|
||||||
|
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||||
|
let state = Arc::new(VersionedDeleteMarkerCommitBarrierState {
|
||||||
|
bucket: bucket.to_string(),
|
||||||
|
object: object.to_string(),
|
||||||
|
arrived: tokio::sync::Notify::new(),
|
||||||
|
release: tokio::sync::Notify::new(),
|
||||||
|
});
|
||||||
|
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("versioned delete-marker commit barrier mutex should not poison");
|
||||||
|
assert!(slot.is_none(), "versioned delete-marker commit barrier must be unique");
|
||||||
|
*slot = Some(Arc::clone(&state));
|
||||||
|
Self { state }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn wait_until_paused(&self) {
|
||||||
|
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
||||||
|
.await
|
||||||
|
.expect("versioned DELETE should reach the post-marker-commit barrier");
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn release(&self) {
|
||||||
|
self.state.release.notify_one();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl Drop for VersionedDeleteMarkerCommitBarrier {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.state.release.notify_one();
|
||||||
|
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("versioned delete-marker commit barrier mutex should not poison");
|
||||||
|
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
||||||
|
*slot = None;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
async fn pause_versioned_delete_marker_after_commit(bucket: &str, object: &str) {
|
||||||
|
let state = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("versioned delete-marker commit barrier mutex should not poison")
|
||||||
|
.as_ref()
|
||||||
|
.filter(|state| state.bucket == bucket && state.object == object)
|
||||||
|
.cloned();
|
||||||
|
if let Some(state) = state {
|
||||||
|
state.arrived.notify_one();
|
||||||
|
state.release.notified().await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Whether a delete-time lookup miss on a directory key should trigger an orphan
|
/// Whether a delete-time lookup miss on a directory key should trigger an orphan
|
||||||
/// empty-directory tree purge (issue #4189).
|
/// empty-directory tree purge (issue #4189).
|
||||||
///
|
///
|
||||||
@@ -1574,6 +1669,34 @@ impl ECStore {
|
|||||||
.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"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn acquire_decommission_source_cleanup_fence(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
source_set: &SetDisks,
|
||||||
|
) -> Result<SourceCleanupMutationFence> {
|
||||||
|
if self.ctx.lock_manager().is_disabled() {
|
||||||
|
return Err(Error::other("decommission source cleanup requires namespace locking"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object);
|
||||||
|
let object = encode_dir_object(object);
|
||||||
|
let distributed = self.ctx.is_dist_erasure().await;
|
||||||
|
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
|
||||||
|
let source_lock_covered = !distributed || same_distributed_lock_domain(&fixed_set.lockers, &source_set.lockers);
|
||||||
|
// Lock order: fixed store mutation domain first; source cleanup takes its
|
||||||
|
// hashed source-domain lock second only when this guard does not cover it.
|
||||||
|
let guard = self
|
||||||
|
.acquire_object_write_lock("decommission_source_cleanup", bucket, &object)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(SourceCleanupMutationFence {
|
||||||
|
guard,
|
||||||
|
source_lock_covered,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn acquire_all_object_read_locks(
|
pub(crate) async fn acquire_all_object_read_locks(
|
||||||
&self,
|
&self,
|
||||||
op: &'static str,
|
op: &'static str,
|
||||||
@@ -2681,12 +2804,14 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_pool_rebalancing(pool.pool_idx).await {
|
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
match pool.delete_object(bucket, object, opts.clone()).await {
|
match pool.delete_object(bucket, object, opts.clone()).await {
|
||||||
Ok(res) => {
|
Ok(res) => {
|
||||||
|
#[cfg(test)]
|
||||||
|
pause_versioned_delete_marker_after_commit(bucket, object).await;
|
||||||
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
||||||
commit_prepared_tier_delete_journal_entry(api, je).await;
|
commit_prepared_tier_delete_journal_entry(api, je).await;
|
||||||
}
|
}
|
||||||
@@ -4481,6 +4606,16 @@ mod tests {
|
|||||||
assert!(lookup_opts.no_lock);
|
assert!(lookup_opts.no_lock);
|
||||||
assert!(!lookup_opts.skip_decommissioned);
|
assert!(!lookup_opts.skip_decommissioned);
|
||||||
assert!(lookup_opts.skip_rebalancing);
|
assert!(lookup_opts.skip_rebalancing);
|
||||||
|
|
||||||
|
let explicit_version = delete_pool_lookup_opts(
|
||||||
|
&ObjectOptions {
|
||||||
|
versioned: true,
|
||||||
|
version_id: Some(uuid::Uuid::new_v4().to_string()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
true,
|
||||||
|
);
|
||||||
|
assert!(!explicit_version.skip_decommissioned);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
Reference in New Issue
Block a user