From c1538cf1c3227761c5bbd749694eb00cac516c59 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 29 Jul 2026 09:49:19 +0800 Subject: [PATCH] fix(multipart): serialize complete and abort (#5356) * fix(multipart): serialize complete and abort * test(multipart): order abort-first finalization * fix(multipart): enforce quorum staging cleanup * fix(multipart): remove stale mutable binding --- crates/ecstore/src/set_disk/ops/list.rs | 23 ++ crates/ecstore/src/set_disk/ops/multipart.rs | 309 ++++++++++++++++++- 2 files changed, 323 insertions(+), 9 deletions(-) diff --git a/crates/ecstore/src/set_disk/ops/list.rs b/crates/ecstore/src/set_disk/ops/list.rs index 11384914c..85eaaa188 100644 --- a/crates/ecstore/src/set_disk/ops/list.rs +++ b/crates/ecstore/src/set_disk/ops/list.rs @@ -29,6 +29,12 @@ impl SetDisks { pub async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> { ListOperations::new(self.ctx()).delete_all(bucket, prefix).await } + + pub(crate) async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> { + ListOperations::new(self.ctx()) + .delete_all_with_quorum(bucket, prefix, write_quorum) + .await + } } /// List/prefix maintenance operations, borrowing the `SetDisks` core state @@ -48,6 +54,14 @@ impl<'a> ListOperations<'a> { } pub(crate) async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> { + self.delete_all_inner(bucket, prefix, None).await + } + + async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> { + self.delete_all_inner(bucket, prefix, Some(write_quorum)).await + } + + async fn delete_all_inner(&self, bucket: &str, prefix: &str, write_quorum: Option) -> Result<()> { let disks = self.ctx.disks().read().await; let disks = disks.clone(); @@ -79,6 +93,9 @@ impl<'a> ListOperations<'a> { Ok(_) => { errors.push(None); } + Err(DiskError::FileNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound) => { + errors.push(None); + } Err(e) => { errors.push(Some(e)); } @@ -97,6 +114,12 @@ impl<'a> ListOperations<'a> { ); } + if let Some(write_quorum) = write_quorum + && let Some(err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum) + { + return Err(err.into()); + } + Ok(()) } } diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 790e3a519..6eaaa3399 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -783,6 +783,9 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { mut max_parts: usize, opts: &ObjectOptions, ) -> Result { + let _upload_guard = self + .acquire_multipart_upload_read_lock("list_object_parts", bucket, object, upload_id, opts) + .await?; let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?; let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); @@ -1260,10 +1263,15 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let _upload_guard = self .acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts) .await?; - self.check_upload_id_exists(bucket, object, upload_id, false).await?; + let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, true).await?; let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); - self.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await + self.delete_all_with_quorum( + RUSTFS_META_MULTIPART_BUCKET, + &upload_id_path, + fi.write_quorum(self.default_write_quorum()), + ) + .await } // complete_multipart_upload finished #[tracing::instrument(skip(self))] @@ -1821,7 +1829,36 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { #[cfg(test)] pause_multipart_commit(bucket, object, MultipartCommitPause::AfterRename).await; - drop(upload_guard); + + let cleanup_store = self.clone(); + let cleanup_upload_id_path = upload_id_path.clone(); + let cleanup_bucket = bucket.to_owned(); + let cleanup_object = object.to_owned(); + let cleanup_upload_id = upload_id.to_owned(); + let cleanup_handle = tokio::spawn(async move { + let _upload_guard = upload_guard; + if let Err(err) = cleanup_store + .delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &cleanup_upload_id_path, write_quorum) + .await + { + warn!( + bucket = %cleanup_bucket, + object = %cleanup_object, + upload_id = %cleanup_upload_id, + error = ?err, + "completed multipart upload staging cleanup did not reach write quorum" + ); + } + }); + if let Err(err) = cleanup_handle.await { + warn!( + bucket = %bucket, + object = %object, + upload_id = %upload_id, + error = ?err, + "completed multipart upload staging cleanup task failed" + ); + } drop(object_lock_guard); // drop object lock guard to release the lock // backlog#1321: enqueue heal only when the committed replicas actually @@ -1866,12 +1903,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { }); } - let upload_id_path = upload_id_path.clone(); - let store = self.clone(); - let _cleanup_handle = tokio::spawn(async move { - let _ = store.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await; - }); - for (i, op_disk) in online_disks.iter().enumerate() { if let Some(disk) = op_disk && disk.is_online().await @@ -2183,6 +2214,122 @@ mod tests { ) } + async fn assert_complete_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager)))); + let lockers: Vec> = vec![signaling.clone()]; + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await; + make_bucket_on_all(&disk_stores, bucket).await; + let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x47; 4096], &create_opts).await; + let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id); + signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path)); + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await; + let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::AfterRename); + + let complete_store = set_disks.clone(); + let complete_upload_id = upload_id.clone(); + let complete = tokio::spawn(async move { + complete_store + .complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default()) + .await + }); + barrier.wait_until_paused().await; + + let abort_store = set_disks.clone(); + let abort_upload_id = upload_id.clone(); + let abort = tokio::spawn(async move { + abort_store + .abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(2).await; + assert!(!abort.is_finished(), "abort must wait for the completion upload lock"); + + barrier.release(); + complete + .await + .expect("completion task should not panic") + .expect("completion should win the upload finalization"); + let abort_err = abort + .await + .expect("abort task should not panic") + .expect_err("abort must observe the upload as finalized"); + assert!(matches!(abort_err, StorageError::InvalidUploadID(..))); + set_disks + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect("complete-first must leave the committed object readable"); + assert!(matches!( + set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await, + Err(StorageError::InvalidUploadID(..)) + )); + } + + async fn assert_abort_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager)))); + let lockers: Vec> = vec![signaling.clone()]; + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await; + make_bucket_on_all(&disk_stores, bucket).await; + let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x48; 4096], &create_opts).await; + let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id); + signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path.clone())); + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await; + let object_holder = set_disks + .new_ns_lock(bucket, object) + .await + .expect("object namespace lock should be created") + .get_write_lock(Duration::from_secs(5)) + .await + .expect("test should hold the object lock"); + let holder = set_disks + .new_ns_lock(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await + .expect("upload namespace lock should be created") + .get_write_lock(Duration::from_secs(5)) + .await + .expect("test should hold the upload lock"); + signaling.wait_for_attempts(1).await; + + let abort_store = set_disks.clone(); + let abort_upload_id = upload_id.clone(); + let abort = tokio::spawn(async move { + abort_store + .abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(2).await; + + let complete_store = set_disks.clone(); + let complete_upload_id = upload_id.clone(); + let complete = tokio::spawn(async move { + complete_store + .complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default()) + .await + }); + drop(holder); + + abort + .await + .expect("abort task should not panic") + .expect("abort should win the upload finalization"); + drop(object_holder); + let complete_err = complete + .await + .expect("completion task should not panic") + .expect_err("completion must observe the aborted upload"); + assert!(matches!(complete_err, StorageError::InvalidUploadID(..))); + let object_err = set_disks + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect_err("abort-first must not publish an object"); + assert!(matches!(object_err, StorageError::ObjectNotFound(..))); + assert!(matches!( + set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await, + Err(StorageError::InvalidUploadID(..)) + )); + } + async fn assert_quorum_minus_one_retry_preserves_completable_part( disk_count: usize, parity: usize, @@ -3045,6 +3192,81 @@ mod tests { .await; } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn abort_and_complete_linearize_for_plain_sse_and_legacy_layouts() { + assert_complete_first_linearizes("multipart-complete-first-plain", "object", ObjectOptions::default()).await; + assert_abort_first_linearizes("multipart-abort-first-plain", "object", ObjectOptions::default()).await; + + let encrypted_opts = ObjectOptions { + user_defined: HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]), + ..Default::default() + }; + temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async { + assert_complete_first_linearizes("multipart-complete-first-sse", "object", encrypted_opts.clone()).await; + assert_abort_first_linearizes("multipart-abort-first-sse", "object", encrypted_opts.clone()).await; + }) + .await; + temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("false"))], async { + assert_complete_first_linearizes("multipart-complete-first-legacy", "object", encrypted_opts.clone()).await; + assert_abort_first_linearizes("multipart-abort-first-legacy", "object", encrypted_opts).await; + }) + .await; + } + + #[tokio::test] + async fn abort_enforces_delete_write_quorum_boundary() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "multipart-abort-delete-quorum"; + let object = "object"; + make_bucket_on_all(&disk_stores, bucket).await; + let quorum_upload = set_disks + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("multipart upload should be created"); + + let saved_disks = { + let mut disks = set_disks.disks.write().await; + let saved = disks.clone(); + disks[3] = None; + saved + }; + set_disks + .abort_multipart_upload(bucket, object, &quorum_upload.upload_id, &ObjectOptions::default()) + .await + .expect("abort should succeed at the exact delete write quorum"); + *set_disks.disks.write().await = saved_disks; + assert!(matches!( + set_disks + .check_upload_id_exists(bucket, object, &quorum_upload.upload_id, false) + .await, + Err(StorageError::InvalidUploadID(..)) + )); + + let below_quorum_upload = set_disks + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("second multipart upload should be created"); + let saved_disks = { + let mut disks = set_disks.disks.write().await; + let saved = disks.clone(); + disks[2] = None; + disks[3] = None; + saved + }; + let err = set_disks + .abort_multipart_upload(bucket, object, &below_quorum_upload.upload_id, &ObjectOptions::default()) + .await + .expect_err("abort must report a delete below write quorum"); + assert!(matches!(err, StorageError::ErasureWriteQuorum)); + + *set_disks.disks.write().await = saved_disks; + set_disks + .check_upload_id_exists(bucket, object, &below_quorum_upload.upload_id, false) + .await + .expect("failed abort must leave quorum-visible staging on the restored disks"); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn complete_revalidates_layout_candidate_after_upload_lock() { @@ -3176,6 +3398,17 @@ mod tests { tokio::task::yield_now().await; assert!(!abort.is_finished(), "abort must wait until completion releases the upload lock"); + let list_store = set_disks.clone(); + let list_upload_id = upload_id.clone(); + let list = tokio::spawn(async move { + list_store + .list_object_parts(bucket, object, &list_upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(3).await; + tokio::task::yield_now().await; + assert!(!list.is_finished(), "ListParts must wait until completion releases the upload lock"); + barrier.release(); complete .await @@ -3186,10 +3419,68 @@ mod tests { .expect("abort task should not panic") .expect_err("the committed upload should no longer exist when abort acquires the lock"); assert!(matches!(abort_err, StorageError::InvalidUploadID(..))); + let list_err = list + .await + .expect("ListParts task should not panic") + .expect_err("the committed upload should no longer exist when ListParts acquires the lock"); + assert!(matches!(list_err, StorageError::InvalidUploadID(..))); }) .await; } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn complete_validates_parts_after_an_inflight_upload_part_commit() { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager)))); + let lockers: Vec> = vec![signaling.clone()]; + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await; + let bucket = "multipart-complete-put-part-race-bucket"; + let object = "object"; + make_bucket_on_all(&disk_stores, bucket).await; + let (upload_id, original_parts) = + stage_upload_with_create_opts(&set_disks, bucket, object, &[0x49; 4096], &ObjectOptions::default()).await; + let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id); + signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path)); + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await; + let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::PutPartBeforeLockLost); + + let put_store = set_disks.clone(); + let put_upload_id = upload_id.clone(); + let put = tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(vec![0x4a; 4096]); + put_store + .put_object_part(bucket, object, &put_upload_id, 1, &mut reader, &ObjectOptions::default()) + .await + }); + barrier.wait_until_paused().await; + + let complete_store = set_disks.clone(); + let complete_upload_id = upload_id.clone(); + let complete = tokio::spawn(async move { + complete_store + .complete_multipart_upload(bucket, object, &complete_upload_id, original_parts, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(2).await; + tokio::task::yield_now().await; + assert!(!complete.is_finished(), "completion must wait for the UploadPart commit lock"); + + barrier.release(); + put.await + .expect("UploadPart task should not panic") + .expect("UploadPart replacement should commit"); + let err = complete + .await + .expect("completion task should not panic") + .expect_err("completion must reject the stale ETag after UploadPart wins"); + assert!(matches!(err, StorageError::InvalidPart(..))); + set_disks + .list_object_parts(bucket, object, &upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default()) + .await + .expect("failed completion must leave the upload retryable"); + } + #[tokio::test(start_paused = true)] #[serial] async fn complete_fences_upload_lock_loss_before_commit() {