mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-26 05:56:50 +00:00
fix(ecstore): fence data movement source cleanup (#5794)
This commit is contained in:
@@ -546,6 +546,81 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
struct SourceCleanupDeleteBarrierState {
|
||||||
|
bucket: String,
|
||||||
|
object: String,
|
||||||
|
arrived: tokio::sync::Notify,
|
||||||
|
release: tokio::sync::Notify,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) struct SourceCleanupDeleteBarrier {
|
||||||
|
state: Arc<SourceCleanupDeleteBarrierState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
|
||||||
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl SourceCleanupDeleteBarrier {
|
||||||
|
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||||
|
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
||||||
|
bucket: bucket.to_string(),
|
||||||
|
object: object.to_string(),
|
||||||
|
arrived: tokio::sync::Notify::new(),
|
||||||
|
release: tokio::sync::Notify::new(),
|
||||||
|
});
|
||||||
|
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("source cleanup delete barrier mutex should not poison");
|
||||||
|
assert!(slot.is_none(), "source cleanup delete barrier must be unique");
|
||||||
|
*slot = Some(Arc::clone(&state));
|
||||||
|
Self { state }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn wait_until_paused(&self) {
|
||||||
|
tokio::time::timeout(StdDuration::from_secs(30), self.state.arrived.notified())
|
||||||
|
.await
|
||||||
|
.expect("source cleanup should reach the pre-delete barrier");
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn release(&self) {
|
||||||
|
self.state.release.notify_one();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl Drop for SourceCleanupDeleteBarrier {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.state.release.notify_one();
|
||||||
|
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("source cleanup delete 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_source_cleanup_before_delete(bucket: &str, object: &str) {
|
||||||
|
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("source cleanup delete barrier mutex should not poison")
|
||||||
|
.as_ref()
|
||||||
|
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
|
||||||
|
.cloned();
|
||||||
|
if let Some(barrier) = barrier {
|
||||||
|
barrier.arrived.notify_one();
|
||||||
|
barrier.release.notified().await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn cleanup_source_entry_if_unchanged(
|
pub(crate) async fn cleanup_source_entry_if_unchanged(
|
||||||
set: Arc<SetDisks>,
|
set: Arc<SetDisks>,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -560,19 +635,18 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
|
|
||||||
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?;
|
||||||
|
|
||||||
let result = set
|
#[cfg(test)]
|
||||||
.delete_object(
|
pause_source_cleanup_before_delete(bucket, object).await;
|
||||||
bucket,
|
|
||||||
cleanup_key.as_str(),
|
let mut opts = ObjectOptions {
|
||||||
ObjectOptions {
|
delete_prefix: true,
|
||||||
delete_prefix: true,
|
delete_prefix_object: true,
|
||||||
delete_prefix_object: true,
|
data_movement: true,
|
||||||
data_movement: true,
|
no_lock: true,
|
||||||
no_lock: true,
|
..Default::default()
|
||||||
..Default::default()
|
};
|
||||||
},
|
opts.add_namespace_lock_guard(&_guard);
|
||||||
)
|
let result = set.delete_object(bucket, cleanup_key.as_str(), opts).await;
|
||||||
.await;
|
|
||||||
if result.is_ok() {
|
if result.is_ok() {
|
||||||
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
|
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7578,6 +7578,55 @@ mod transition_upload_integrity_tests {
|
|||||||
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "current_thread", start_paused = true)]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn data_movement_cleanup_aborts_after_outer_lock_loss() {
|
||||||
|
let refresh_calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let lockers: Vec<Arc<dyn LockClient>> = (0..4)
|
||||||
|
.map(|_| Arc::new(LockLostRefreshClient::new(Arc::clone(&refresh_calls))) as Arc<dyn LockClient>)
|
||||||
|
.collect();
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
|
||||||
|
let bucket = "data-movement-cleanup-lock-lost";
|
||||||
|
let object = "object.bin";
|
||||||
|
let payload = b"lost data movement cleanup lock must preserve the source".repeat(1024);
|
||||||
|
write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
|
||||||
|
let expected = set_disks
|
||||||
|
.load_file_info_versions_exact(bucket, object)
|
||||||
|
.await
|
||||||
|
.expect("source versions should be readable")
|
||||||
|
.expect("source versions should exist");
|
||||||
|
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
|
||||||
|
let barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(bucket, object);
|
||||||
|
|
||||||
|
let cleanup_set = Arc::clone(&set_disks);
|
||||||
|
let cleanup = tokio::spawn(async move {
|
||||||
|
crate::data_movement::cleanup_source_entry_if_unchanged(
|
||||||
|
cleanup_set,
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
&expected,
|
||||||
|
&[],
|
||||||
|
"test_data_movement",
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
barrier.wait_until_paused().await;
|
||||||
|
tokio::time::advance(Duration::from_secs(11)).await;
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
assert!(
|
||||||
|
refresh_calls.load(Ordering::SeqCst) > 0,
|
||||||
|
"test must drive the real distributed-lock heartbeat before cleanup commit"
|
||||||
|
);
|
||||||
|
barrier.release();
|
||||||
|
|
||||||
|
let error = cleanup
|
||||||
|
.await
|
||||||
|
.expect("cleanup task should not panic")
|
||||||
|
.expect_err("cleanup must fail after its outer namespace lock loses refresh quorum");
|
||||||
|
assert!(matches!(error, StorageError::NamespaceLockQuorumUnavailable { .. }));
|
||||||
|
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn partial_remote_acceptance_cleans_exact_candidate_and_preserves_source() {
|
async fn partial_remote_acceptance_cleans_exact_candidate_and_preserves_source() {
|
||||||
|
|||||||
Reference in New Issue
Block a user