From 91a02914cce953920351d62ad358e2299888e59d Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 2 Jul 2026 22:35:56 +0800 Subject: [PATCH] fix(ecstore): retry multipart abort cleanup (#4194) --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 16 +++++ crates/ecstore/src/data_movement/mod.rs | 63 ++++++++++++++++++- 2 files changed, 78 insertions(+), 1 deletion(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index f340577e7..d0d69b75c 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -1817,6 +1817,22 @@ pub async fn run_stale_multipart_upload_cleanup_once(api: Arc) -> usize cleanup_stale_multipart_uploads_once_at(api, OffsetDateTime::now_utc(), stale_uploads_expiry()).await } +pub fn schedule_stale_multipart_upload_cleanup_once(api: Arc) { + tokio::spawn(async move { + let deleted = run_stale_multipart_upload_cleanup_once(api).await; + if deleted > 0 { + debug!( + event = EVENT_LIFECYCLE_STALE_MULTIPART_CLEANUP, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_LIFECYCLE, + deleted, + trigger = "on_demand", + "Completed stale multipart cleanup pass" + ); + } + }); +} + pub fn init_background_stale_multipart_upload_cleanup(api: Arc) { let cleanup_interval = stale_uploads_cleanup_interval(); let default_expiry = stale_uploads_expiry(); diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 099f73e62..dee12ec1f 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -17,7 +17,9 @@ pub(crate) mod backpressure; -use crate::error::{Error, Result, is_err_data_movement_overwrite, is_err_object_not_found, is_err_version_not_found}; +use crate::error::{ + Error, Result, is_err_data_movement_overwrite, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found, +}; use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}; use crate::set_disk::{SetDisks, get_lock_acquire_timeout}; use crate::storage_api_contracts::{ @@ -38,11 +40,14 @@ use std::sync::{ atomic::{AtomicBool, Ordering}, }; use std::task::{Context, Poll}; +use std::time::Duration as StdDuration; use time::format_description::well_known::Rfc3339; use tokio::io::{AsyncRead, BufReader, ReadBuf}; use tracing::{error, info}; type SharedDataMovementStream = Arc>>; +const DATA_MOVEMENT_MULTIPART_ABORT_RETRY_ATTEMPTS: usize = 3; +const DATA_MOVEMENT_MULTIPART_ABORT_RETRY_DELAY_SECS: u64 = 60; pub struct IndexedDataMovementReader { inner: R, @@ -283,6 +288,54 @@ fn data_movement_stage_error(op_label: &str, stage: &str, bucket: &str, object: Error::other(format!("{op_label}: {stage} failed for {bucket}/{object}: {err}")) } +fn schedule_data_movement_multipart_abort_cleanup( + store: Arc, + target_pool_idx: usize, + bucket: String, + object: String, + upload_id: String, + op_label: &str, +) { + let op_label = op_label.to_string(); + tokio::spawn(async move { + for attempt in 1..=DATA_MOVEMENT_MULTIPART_ABORT_RETRY_ATTEMPTS { + tokio::time::sleep(StdDuration::from_secs(DATA_MOVEMENT_MULTIPART_ABORT_RETRY_DELAY_SECS)).await; + + let Some(pool) = store.pools.get(target_pool_idx).cloned() else { + error!( + "{op_label}: background abort_multipart_upload cleanup skipped for {bucket}/{object} upload {upload_id}: target pool {target_pool_idx} is out of range" + ); + return; + }; + + match pool + .abort_multipart_upload(&bucket, &object, &upload_id, &ObjectOptions::default()) + .await + { + Ok(()) => { + info!( + "{op_label}: background abort_multipart_upload cleanup succeeded for {bucket}/{object} upload {upload_id} on attempt {attempt}" + ); + return; + } + Err(err) if is_err_invalid_upload_id(&err) => { + info!( + "{op_label}: background abort_multipart_upload cleanup found {bucket}/{object} upload {upload_id} already removed" + ); + return; + } + Err(err) => { + error!( + "{op_label}: background abort_multipart_upload cleanup attempt {attempt} failed for {bucket}/{object} upload {upload_id}: {err:?}" + ); + } + } + } + + crate::bucket::lifecycle::bucket_lifecycle_ops::schedule_stale_multipart_upload_cleanup_once(store); + }); +} + fn should_check_data_movement_overwrite_resume(err: &Error) -> bool { is_err_data_movement_overwrite(err) } @@ -807,6 +860,14 @@ pub(crate) async fn migrate_object( Ok(()) => Err(primary_err), Err(abort_err) => { error!("{op_label}: abort_multipart_upload err {:?}", &abort_err); + schedule_data_movement_multipart_abort_cleanup( + store.clone(), + target_pool_idx, + bucket.clone(), + object_info.name.clone(), + res.upload_id.clone(), + op_label, + ); Err(resolve_data_movement_abort_result( op_label, bucket.as_str(),