From 4255e0ca9aee7e893dd1ddcbe66a6e9df97ecb2b Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 2 Sep 2026 17:32:32 +0800 Subject: [PATCH] fix(ecstore): fail closed on multipart cleanup gaps --- crates/ecstore/src/core/pools_test.rs | 69 ++++++++-- crates/ecstore/src/data_movement/mod.rs | 19 ++- crates/ecstore/src/set_disk/ops/multipart.rs | 128 ++++++++++++++++++- crates/ecstore/src/store/init.rs | 27 +++- crates/ecstore/src/store/multipart.rs | 112 +++++++++++----- 5 files changed, 288 insertions(+), 67 deletions(-) diff --git a/crates/ecstore/src/core/pools_test.rs b/crates/ecstore/src/core/pools_test.rs index 59427eb49..aa6c5e92a 100644 --- a/crates/ecstore/src/core/pools_test.rs +++ b/crates/ecstore/src/core/pools_test.rs @@ -232,7 +232,13 @@ mod decommission_lock_order_tests { C: FnOnce() -> F + Send + 'static, F: Future + 'static, { - const STACK_SIZE: usize = 32 * 1024 * 1024; + const STACK_SIZE: usize = if cfg!(debug_assertions) { + 8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else if cfg!(target_os = "macos") { + 2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else { + rustfs_config::DEFAULT_THREAD_STACK_SIZE + }; std::thread::Builder::new() .name(name.to_string()) .stack_size(STACK_SIZE) @@ -255,7 +261,13 @@ mod decommission_lock_order_tests { C: FnOnce() -> F + Send + 'static, F: Future + 'static, { - const STACK_SIZE: usize = 32 * 1024 * 1024; + const STACK_SIZE: usize = if cfg!(debug_assertions) { + 8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else if cfg!(target_os = "macos") { + 2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else { + rustfs_config::DEFAULT_THREAD_STACK_SIZE + }; std::thread::Builder::new() .name(name.to_string()) .stack_size(STACK_SIZE) @@ -1496,7 +1508,7 @@ mod decommission_lock_order_tests { #[tokio::test] #[serial_test::serial] - async fn multipart_cleanup_proves_target_absence_after_acquiring_capacity_gate() { + async fn multipart_cleanup_rediscovers_uploads_and_proves_target_under_capacity_gate() { let (_temp_dirs, store, _other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await; let bucket = test_bucket("multipart-cleanup-target-race"); let object = "published-between-cleanup-proof-and-finalize.bin"; @@ -1549,6 +1561,7 @@ mod decommission_lock_order_tests { upload_identity.clone(), ); owner.apply_to(&mut upload_opts); + let concurrent_upload_opts = upload_opts.clone(); let mut cleanup_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, @@ -1627,6 +1640,9 @@ mod decommission_lock_order_tests { result = &mut cleanup => panic!("cleanup finished before its target-gate proof point: {result:?}"), } + let concurrent_upload = new_multipart_upload(&store, target_pool_index, &bucket, object, concurrent_upload_opts) + .await + .expect("stage a same-identity upload after cleanup reaches its target-gate boundary"); let mut target_data = PutObjReader::from_vec(body); let published = store.pools[target_pool_index] .put_object( @@ -1655,6 +1671,17 @@ mod decommission_lock_order_tests { .expect("cleanup should preserve the now-published target intent"); drop(barrier); + let remaining_uploads = store.pools[target_pool_index] + .get_disks_by_key(object) + .data_movement_multipart_upload_ids(&bucket, object, Some(incarnation), &upload_identity) + .await + .expect("check same-identity uploads after gated cleanup"); + assert!( + remaining_uploads.is_empty(), + "gated cleanup must rediscover and remove the concurrent same-identity upload {}", + concurrent_upload.upload_id + ); + let pool_meta = store.pool_meta.read().await; let reservation = pool_meta.pools[0] .decommission @@ -1672,7 +1699,7 @@ mod decommission_lock_order_tests { #[test] #[serial_test::serial] fn data_movement_multipart_restart_reconciles_published_capacity_before_new_upload() { - run_large_stack_async_test( + run_large_stack_current_thread_async_test( "multipart-published-capacity-restart", data_movement_multipart_restart_reconciles_published_capacity_before_new_upload_case, ); @@ -2009,7 +2036,7 @@ mod decommission_lock_order_tests { #[test] #[serial_test::serial] fn data_movement_multipart_restart_cleans_partial_upload_before_exact_fit_retry() { - run_large_stack_async_test( + run_large_stack_current_thread_async_test( "multipart-partial-upload-restart", data_movement_multipart_restart_cleans_partial_upload_before_exact_fit_retry_case, ); @@ -2311,18 +2338,22 @@ mod decommission_lock_order_tests { .await .expect("activate the abort fence reservation"); let owner = decommission_capacity_owner(&*store.pool_meta.read().await).with_mutation_id(uuid::Uuid::new_v4()); + let version_id = uuid::Uuid::new_v4().to_string(); + let mod_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(83); let mut upload_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, versioned: true, - version_id: Some(uuid::Uuid::new_v4().to_string()), + version_id: Some(version_id.clone()), + mod_time: Some(mod_time), expected_bucket_incarnation_id: Some(incarnation), ..Default::default() }; + let upload_identity = data_movement::data_movement_upload_identity_from_options(&upload_opts); rustfs_utils::http::insert_str( &mut upload_opts.user_defined, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, - "v1:abort-fence:1".to_string(), + upload_identity, ); let upload = new_multipart_upload(&store, 2, &bucket, object, upload_opts) .await @@ -2333,6 +2364,9 @@ mod decommission_lock_order_tests { let mut abort_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, + versioned: true, + version_id: Some(version_id), + mod_time: Some(mod_time), expected_bucket_incarnation_id: Some(incarnation), ..Default::default() }; @@ -2387,9 +2421,16 @@ mod decommission_lock_order_tests { assert!(upload_info.parts.is_empty()); } - #[tokio::test] + #[test] #[serial_test::serial] - async fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss() { + fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss() { + run_large_stack_current_thread_async_test( + "multipart-abort-restart", + data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss_case, + ); + } + + async fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss_case() { let (_temp_dirs, store, other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await; let bucket = test_bucket("multipart-abort-restart"); let object = "deleted-upload-before-release-save.bin"; @@ -2418,15 +2459,18 @@ mod decommission_lock_order_tests { .await .expect("activate the abort restart reservation"); let owner = decommission_capacity_owner(&*store.pool_meta.read().await).with_mutation_id(uuid::Uuid::new_v4()); - let upload_identity = format!("v1:{}:{}", uuid::Uuid::new_v4(), 91); + let version_id = uuid::Uuid::new_v4().to_string(); + let mod_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(91); let mut new_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, versioned: true, - version_id: Some(uuid::Uuid::new_v4().to_string()), + version_id: Some(version_id.clone()), + mod_time: Some(mod_time), expected_bucket_incarnation_id: Some(incarnation), ..Default::default() }; + let upload_identity = data_movement::data_movement_upload_identity_from_options(&new_opts); rustfs_utils::http::insert_str( &mut new_opts.user_defined, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, @@ -2496,6 +2540,9 @@ mod decommission_lock_order_tests { let mut abort_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, + versioned: true, + version_id: Some(version_id), + mod_time: Some(mod_time), expected_bucket_incarnation_id: Some(incarnation), ..Default::default() }; diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index b422a70aa..0b517da2e 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -269,16 +269,21 @@ fn insert_data_movement_checksum(user_defined: &mut HashMap, obj } } -fn data_movement_upload_identity(object_info: &ObjectInfo) -> String { - let version_id = object_info - .version_id - .map_or_else(|| "none".to_string(), |version_id| version_id.to_string()); - let mod_time = object_info - .mod_time - .map_or_else(|| "none".to_string(), |mod_time| mod_time.unix_timestamp_nanos().to_string()); +fn data_movement_upload_identity_parts(version_id: Option<&str>, mod_time: Option) -> String { + let version_id = version_id.unwrap_or("none"); + let mod_time = mod_time.map_or_else(|| "none".to_string(), |mod_time| mod_time.unix_timestamp_nanos().to_string()); format!("v1:{version_id}:{mod_time}") } +fn data_movement_upload_identity(object_info: &ObjectInfo) -> String { + let version_id = object_info.version_id.map(|version_id| version_id.to_string()); + data_movement_upload_identity_parts(version_id.as_deref(), object_info.mod_time) +} + +pub(crate) fn data_movement_upload_identity_from_options(opts: &ObjectOptions) -> String { + data_movement_upload_identity_parts(opts.version_id.as_deref(), opts.mod_time) +} + fn data_movement_new_multipart_opts(object_info: &ObjectInfo, src_pool_idx: usize) -> ObjectOptions { let mut user_defined = data_movement_user_defined(object_info); let upload_identity = data_movement_upload_identity(object_info); diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 96883d31d..aeb98af2f 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -837,7 +837,7 @@ impl SetDisks { orig_bucket: &str, error_path: &str, root_prefix: &str, - ) -> Result<(Vec>, Vec, usize)> { + ) -> Result<(Vec>, Vec, usize, bool)> { let disks = self.disks.read().await.clone(); if disks.is_empty() { return Err(Error::ErasureReadQuorum); @@ -882,16 +882,24 @@ impl SetDisks { return Err(to_object_err(err.into(), vec![orig_bucket, error_path])); } + let mut has_minority_candidate = false; let mut candidate_paths = candidate_counts .into_iter() - .filter_map(|(path, count)| (count >= discovery_quorum).then_some(path)) + .filter_map(|(path, count)| { + if count >= discovery_quorum { + Some(path) + } else { + has_minority_candidate = true; + None + } + }) .collect::>(); candidate_paths.sort_unstable(); - Ok((disks, candidate_paths, discovery_quorum)) + Ok((disks, candidate_paths, discovery_quorum, has_minority_candidate)) } pub(crate) async fn first_multipart_upload_path_for_decommission(&self, bucket: &str) -> Result> { - let (_, paths, _) = self + let (_, paths, _, _) = self .discover_multipart_upload_paths(bucket, RUSTFS_META_MULTIPART_BUCKET, "") .await?; Ok(paths.into_iter().next()) @@ -905,7 +913,14 @@ impl SetDisks { upload_identity: &str, ) -> Result> { let expected_parent = format!("{DATA_MOVEMENT_MULTIPART_PREFIX}/{}", Self::get_multipart_sha_dir(bucket, object)); - let (_, candidate_paths, _) = self.discover_multipart_upload_paths(bucket, object, &expected_parent).await?; + let (_, candidate_paths, _, has_minority_candidate) = + self.discover_multipart_upload_paths(bucket, object, &expected_parent).await?; + if has_minority_candidate { + return Err(Error::DecommissionCapacityBlocked { + message: "data movement multipart cleanup found an upload path on fewer than the discovery quorum of disks" + .to_string(), + }); + } let mut upload_ids = Vec::new(); for upload_path in candidate_paths { let Some((parent, raw_upload_id)) = upload_path.rsplit_once('/') else { @@ -921,7 +936,11 @@ impl SetDisks { { Ok((file_info, _)) => file_info, Err(err) if crate::error::is_err_invalid_upload_id(&err) || crate::error::is_err_object_not_found(&err) => { - continue; + return Err(Error::DecommissionCapacityBlocked { + message: format!( + "data movement multipart cleanup found quorum-visible upload path {upload_path} without verifiable metadata: {err}" + ), + }); } Err(err) => return Err(err), }; @@ -1190,7 +1209,7 @@ impl SetDisks { max_uploads: usize, expected_incarnation_id: Option, ) -> Result { - let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix, "").await?; + let (disks, candidate_paths, discovery_quorum, _) = self.discover_multipart_upload_paths(bucket, prefix, "").await?; let listed_uploads = stream::iter(candidate_paths) .map(|upload_path| { let disks = &disks; @@ -7551,6 +7570,101 @@ mod tests { assert_eq!(bucket_wide.uploads[0].object, "blobs/data/layer.bin"); } + #[tokio::test] + async fn data_movement_cleanup_discovery_rejects_quorum_visible_upload_without_metadata() { + let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "data-movement-unverifiable-upload"; + let object = "staged/object.bin"; + for disk in &disk_stores { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let upload_identity = format!("v1:{}:{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp_nanos()); + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, upload_identity.clone()); + let opts = ObjectOptions { + data_movement: true, + user_defined: metadata, + ..Default::default() + }; + let upload = set_disks + .new_multipart_upload(bucket, object, &opts) + .await + .expect("data movement upload should be created"); + let upload_path = SetDisks::get_multipart_upload_dir(bucket, object, &upload.upload_id, true); + for temp_dir in &temp_dirs { + tokio::fs::remove_file( + temp_dir + .path() + .join(RUSTFS_META_MULTIPART_BUCKET) + .join(&upload_path) + .join("xl.meta"), + ) + .await + .expect("upload metadata should be removable while preserving its quorum-visible directory"); + } + + let err = set_disks + .data_movement_multipart_upload_ids(bucket, object, None, &upload_identity) + .await + .expect_err("cleanup discovery must fail closed on an unverifiable upload path"); + assert!(matches!(err, Error::DecommissionCapacityBlocked { .. })); + } + + #[tokio::test] + async fn data_movement_cleanup_discovery_rejects_minority_upload_until_disks_recover() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "data-movement-minority-upload"; + let object = "staged/object.bin"; + for disk in &disk_stores { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + { + let mut disks = set_disks.disks.write().await; + disks[3] = None; + } + let upload_identity = format!("v1:{}:{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp_nanos()); + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, upload_identity.clone()); + let upload = set_disks + .new_multipart_upload( + bucket, + object, + &ObjectOptions { + data_movement: true, + user_defined: metadata, + ..Default::default() + }, + ) + .await + .expect("write quorum should create the upload while one disk is offline"); + + { + let mut disks = set_disks.disks.write().await; + disks[1] = None; + disks[2] = None; + disks[3] = Some(disk_stores[3].clone()); + } + let err = set_disks + .data_movement_multipart_upload_ids(bucket, object, None, &upload_identity) + .await + .expect_err("a minority-observed upload must block a destructive absence proof"); + assert!(matches!(err, Error::DecommissionCapacityBlocked { .. })); + + { + let mut disks = set_disks.disks.write().await; + disks[1] = Some(disk_stores[1].clone()); + disks[2] = Some(disk_stores[2].clone()); + } + let recovered = set_disks + .data_movement_multipart_upload_ids(bucket, object, None, &upload_identity) + .await + .expect("recovered quorum should make the staged upload verifiable"); + assert_eq!(recovered.len(), 1); + assert_eq!(upload_uuid_suffix(&recovered[0]), upload_uuid_suffix(&upload.upload_id)); + } + /// Regression (issue #5716): a single upload directory whose `xl.meta` was /// destroyed (crash mid-write, torn disk state) must degrade to that upload /// alone. Failing the whole ListMultipartUploads turns one piece of stale diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 2606bcf71..5d093e39b 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -923,7 +923,13 @@ mod tests { C: FnOnce() -> F + Send + 'static, F: Future + 'static, { - const STACK_SIZE: usize = 32 * 1024 * 1024; + const STACK_SIZE: usize = if cfg!(debug_assertions) { + 8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else if cfg!(target_os = "macos") { + 2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE + } else { + rustfs_config::DEFAULT_THREAD_STACK_SIZE + }; std::thread::Builder::new() .name(name.to_string()) .stack_size(STACK_SIZE) @@ -2418,19 +2424,23 @@ mod tests { let capacity_owner = test_decommission_capacity_owner(store.as_ref(), 0) .await .with_mutation_id(Uuid::new_v4()); - let mut upload_metadata = source.user_defined.as_ref().clone(); - rustfs_utils::http::insert_str( - &mut upload_metadata, - rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, - "mpu-staging-test".to_string(), - ); + let upload_metadata = source.user_defined.as_ref().clone(); let mut staging_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, + versioned: source.version_id.is_some(), + version_id: source.version_id.map(|version_id| version_id.to_string()), + mod_time: source.mod_time, user_defined: upload_metadata, expected_bucket_incarnation_id: Some(bucket_incarnation_id), ..Default::default() }; + let upload_identity = crate::data_movement::data_movement_upload_identity_from_options(&staging_opts); + rustfs_utils::http::insert_str( + &mut staging_opts.user_defined, + rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, + upload_identity, + ); capacity_owner.apply_to(&mut staging_opts); let (upload, target_pool_idx, staged_incarnation_id) = store .handle_new_multipart_upload_with_pool_idx(&bucket, object, &staging_opts, None) @@ -2500,6 +2510,9 @@ mod tests { let mut abort_opts = ObjectOptions { data_movement: true, src_pool_idx: 0, + versioned: source.version_id.is_some(), + version_id: source.version_id.map(|version_id| version_id.to_string()), + mod_time: source.mod_time, expected_bucket_incarnation_id: Some(bucket_incarnation_id), ..Default::default() }; diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index afe68f1c1..908572998 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -787,8 +787,16 @@ impl ECStore { upload_id: &str, opts: &ObjectOptions, ) -> Result<()> { - self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &[upload_id.to_owned()], None, opts) - .await + let upload_identity = crate::data_movement::data_movement_upload_identity_from_options(opts); + self.abort_multipart_uploads_for_data_movement( + target_pool_idx, + bucket, + object, + &[upload_id.to_owned()], + &upload_identity, + opts, + ) + .await } pub(crate) async fn reconcile_multipart_uploads_for_data_movement( @@ -799,8 +807,7 @@ impl ECStore { upload_identity: &str, opts: &ObjectOptions, ) -> Result<()> { - let pool = self - .pools + self.pools .get(target_pool_idx) .ok_or_else(|| Error::other(format!("data movement target pool {target_pool_idx} is out of range")))?; let owner = DecommissionCapacityOwner::from_options(opts) @@ -830,11 +837,7 @@ impl ECStore { .entry(self.id) .or_default() += 1; } - let set = pool.get_disks_by_key(object); - let upload_ids = set - .data_movement_multipart_upload_ids(bucket, object, opts.expected_bucket_incarnation_id, upload_identity) - .await?; - self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &upload_ids, Some(upload_identity), opts) + self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &[], upload_identity, opts) .await } @@ -844,7 +847,7 @@ impl ECStore { bucket: &str, object: &str, upload_ids: &[String], - expected_upload_identity: Option<&str>, + expected_upload_identity: &str, opts: &ObjectOptions, ) -> Result<()> { check_new_multipart_args(bucket, object)?; @@ -861,32 +864,78 @@ impl ECStore { .pools .get(target_pool_idx) .ok_or_else(|| Error::other(format!("data movement target pool {target_pool_idx} is out of range")))?; - let set = pool.get_disks_by_key(object); - let mut guards = Vec::with_capacity(upload_ids.len()); - for upload_id in upload_ids { - if let Some(guard) = set - .lock_data_movement_multipart_abort(bucket, object, upload_id, expected_upload_identity, &opts) - .await? - { - guard.add_namespace_lock_fence(&mut opts); - guards.push(guard); - } - } opts.no_lock = true; - // Keep every upload namespace guard alive through the final capacity progress save. - let cleanup_decision_error = self + let set = pool.get_disks_by_key(object); + // Discover and lock uploads only after the target capacity gate is held. + // Return the guards so they remain alive through the final capacity save. + let (cleanup_decision_error, guards) = self .run_decommission_capacity_temporary_release_with_capacity_lease(target_pool_idx, capacity_owner, |capacity_lease| { let mut delete_opts = opts.clone(); - let guards = &guards; let set = &set; let pool = &pool; async move { if let Some(capacity_lease) = capacity_lease.as_ref() { delete_opts.add_namespace_lock_lost_signal(Arc::clone(capacity_lease)); } - // The target capacity gate held by the surrounding - // transaction makes this proof and its pending-ledger - // decision one critical section with cleanup finalize. + let mut candidate_upload_ids = upload_ids.to_vec(); + candidate_upload_ids.extend( + set.data_movement_multipart_upload_ids( + bucket, + object, + delete_opts.expected_bucket_incarnation_id, + expected_upload_identity, + ) + .await?, + ); + candidate_upload_ids.sort_unstable(); + candidate_upload_ids.dedup(); + + let mut guards = Vec::with_capacity(candidate_upload_ids.len()); + for upload_id in &candidate_upload_ids { + match set + .lock_data_movement_multipart_abort( + bucket, + object, + upload_id, + Some(expected_upload_identity), + &delete_opts, + ) + .await + { + Ok(Some(guard)) => { + guard.add_namespace_lock_fence(&mut delete_opts); + guards.push(guard); + } + Ok(None) => {} + Err(err) => return Err(err), + } + } + for guard in &guards { + match guard.delete(set, bucket, object, &delete_opts).await { + Ok(()) => {} + Err(err) if is_err_invalid_upload_id(&err) => {} + Err(err) => return Err(err), + } + } + + if !set + .data_movement_multipart_upload_ids( + bucket, + object, + delete_opts.expected_bucket_incarnation_id, + expected_upload_identity, + ) + .await? + .is_empty() + { + return Err(Error::DecommissionCapacityBlocked { + message: "multipart cleanup could not prove the exact staged uploads are absent".to_string(), + }); + } + + // The target capacity gate makes the exact target proof, + // upload absence proof, and pending-ledger decision one + // critical section with cleanup finalize. let (clear_pending, cleanup_decision_error) = if capacity_owner.is_some() { let mut lookup_opts = ObjectOptions { versioned: delete_opts.versioned, @@ -911,14 +960,7 @@ impl ECStore { } else { (true, None) }; - for guard in guards { - match guard.delete(set, bucket, object, &delete_opts).await { - Ok(()) => {} - Err(err) if is_err_invalid_upload_id(&err) => {} - Err(err) => return Err(err), - } - } - Ok((cleanup_decision_error, clear_pending)) + Ok(((cleanup_decision_error, guards), clear_pending)) } }) .await?;