diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 2dd0d3f99..345f5b60c 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -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] diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 6af71db17..283fd8f88 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -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