fix(ecstore): fence multipart staging on rebalance lock loss

This commit is contained in:
overtrue
2026-08-22 04:37:52 +08:00
parent 4cce1b2e3b
commit 8d97f3570d
2 changed files with 144 additions and 58 deletions
+77 -58
View File
@@ -1323,71 +1323,90 @@ mod tests {
#[tokio::test]
#[serial_test::serial]
async fn real_rebalance_run_fence_loss_blocks_multipart_completion() {
async fn real_rebalance_run_fence_loss_blocks_multipart_publication() {
const REBALANCE_ID: &str = "rebalance-multipart-commit-fence";
let bucket = crate::disk::RUSTFS_META_BUCKET;
let object = "rebalance-multipart-commit-fence-object";
let (_temp_dirs, store, _unused_store) =
crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await;
let source_set = store.pools[0].get_disks_by_key(object);
let target_set = store.pools[1].get_disks_by_key(object);
let upload = source_set
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("source multipart upload should be created");
let mut reader = PutObjReader::from_vec(b"multipart source payload".repeat(1024));
let part = source_set
.put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default())
.await
.expect("source multipart part should be written");
source_set
.clone()
.complete_multipart_upload(
bucket,
object,
&upload.upload_id,
vec![CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.expect("source multipart object should commit");
let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await;
let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID);
let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::BeforeQuotaRename);
let task = spawn_real_rebalance_entry(
Arc::clone(&store),
Arc::clone(&source_set),
entry,
REBALANCE_ID,
Arc::new(RebalanceBucketConfigs::default()),
);
barrier.wait_until_paused().await;
run_signal_fence.mark_lost();
barrier.release();
drop(barrier);
assert_real_entry_rejected_after_run_fence_loss(task).await;
assert!(
target_set
.load_file_info_versions_exact(bucket, object)
for (object, pause, staged_commit) in [
(
"rebalance-multipart-new-upload-fence",
MultipartCommitPause::NewUploadBeforeLockLost,
Some("upload metadata"),
),
(
"rebalance-multipart-part-fence",
MultipartCommitPause::PutPartBeforeLockLost,
Some("part"),
),
("rebalance-multipart-completion-fence", MultipartCommitPause::BeforeQuotaRename, None),
] {
let source_set = store.pools[0].get_disks_by_key(object);
let target_set = store.pools[1].get_disks_by_key(object);
let upload = source_set
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("target metadata lookup should succeed")
.is_none(),
"lost run fence must not publish the target multipart object"
);
assert!(
.expect("source multipart upload should be created");
let mut reader = PutObjReader::from_vec(b"multipart source payload".repeat(1024));
let part = source_set
.put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default())
.await
.expect("source multipart part should be written");
source_set
.load_file_info_versions_exact(bucket, object)
.clone()
.complete_multipart_upload(
bucket,
object,
&upload.upload_id,
vec![CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.expect("source metadata lookup should succeed")
.is_some(),
"lost run fence must preserve the source multipart object"
);
.expect("source multipart object should commit");
let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await;
let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID);
let barrier = MultipartCommitBarrier::install(bucket, object, pause);
let task = spawn_real_rebalance_entry(
Arc::clone(&store),
Arc::clone(&source_set),
entry,
REBALANCE_ID,
Arc::new(RebalanceBucketConfigs::default()),
);
barrier.wait_until_paused().await;
run_signal_fence.mark_lost();
barrier.release();
assert_real_entry_rejected_after_run_fence_loss(task).await;
if let Some(staged_commit) = staged_commit {
assert!(
!barrier.commit_observed(),
"lost run fence must not publish target multipart {staged_commit}"
);
}
drop(barrier);
assert!(
target_set
.load_file_info_versions_exact(bucket, object)
.await
.expect("target metadata lookup should succeed")
.is_none(),
"lost run fence must not publish the target multipart object"
);
assert!(
source_set
.load_file_info_versions_exact(bucket, object)
.await
.expect("source metadata lookup should succeed")
.is_some(),
"lost run fence must preserve the source multipart object"
);
}
}
#[tokio::test]
@@ -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,
@@ -83,6 +86,8 @@ struct MultipartCommitBarrierState {
pause: MultipartCommitPause,
expected_arrivals: usize,
arrivals: AtomicUsize,
#[cfg(test)]
committed: AtomicBool,
arrived: tokio::sync::Notify,
release: tokio::sync::Semaphore,
}
@@ -110,6 +115,8 @@ impl MultipartCommitBarrier {
pause,
expected_arrivals,
arrivals: AtomicUsize::new(0),
#[cfg(test)]
committed: AtomicBool::new(false),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
});
@@ -140,6 +147,11 @@ impl MultipartCommitBarrier {
pub fn release(&self) {
self.state.release.add_permits(self.state.expected_arrivals);
}
#[cfg(test)]
pub(crate) fn commit_observed(&self) -> bool {
self.state.committed.load(Ordering::Acquire)
}
}
#[cfg(any(test, feature = "test-util"))]
@@ -194,6 +206,20 @@ async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartComm
}
}
#[cfg(test)]
fn observe_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
let slot = MULTIPART_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("multipart commit barrier mutex should not poison");
if let Some(barrier) = slot
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.pause == pause)
{
barrier.committed.store(true, Ordering::Release);
}
}
fn map_upload_id_metadata_error(bucket: &str, object: &str, upload_id: &str, err: DiskError) -> Error {
if err == DiskError::FileNotFound {
return StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned());
@@ -1251,6 +1277,19 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost).await;
fence_commit_on_lock_loss(_upload_commit_guard.as_ref(), "put_object_part_commit", &upload_id_path)?;
fence_commit_on_lock_loss(_part_commit_guard.as_ref(), "put_object_part_commit", &part_lock_path)?;
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "put_object_part_outer_lock",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
let _ = self
@@ -1271,6 +1310,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}),
)
.await?;
#[cfg(test)]
observe_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
#[cfg(test)]
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartAfterRename).await;
@@ -1615,6 +1656,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 +1691,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_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost);
// evalDisks