From 693800db8ee255b773e602220a14cd6548c78e2f Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 23 Aug 2026 06:46:28 +0800 Subject: [PATCH] fix(ecstore): retry decommission entries safely on source changes A single object's SourceChanged during decommission cleanup no longer cancels the shared worker token and fails the whole pool operation. Cleanup preflight and source-cleanup outcomes now retry per entry with bounded attempts and cancellation-aware backoff, applied uniformly to ordinary versions, delete markers, and tiered copies (removing the try-once-only branches); every retry re-lists the entry and redoes version multiset validation before touching the source. Only quorum loss, unrecoverable system errors, or exceeding a pool-level SourceChanged exhaustion threshold still fails the decommission, and exhausted entries never delete their source versions. Retry attempts, backoff, and deferred/exhausted reasons are logged per entry for observability. Heavy regression tests spawn on dedicated 32MiB stacks following the existing store-test pattern. Fixes rustfs/backlog#1913 --- crates/ecstore/src/core/pools.rs | 737 ++++++++++++++++++++++++----- crates/ecstore/src/set_disk/mod.rs | 10 +- crates/ecstore/src/store/init.rs | 612 +++++++++++++++++++++++- 3 files changed, 1246 insertions(+), 113 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 4ca460c20..c6b70d004 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -100,6 +100,11 @@ const DECOMMISSION_ENTRY_WORKERS_PER_SET: usize = 2; const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30; const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3; const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); +pub(crate) const DECOMMISSION_ENTRY_MAX_ATTEMPTS: usize = 3; +const DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(100); +pub(crate) const DECOMMISSION_VERSION_COPY_ATTEMPTS: usize = 3; +const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50); +const DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT: usize = 100; const DECOMMISSION_TERMINAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1); /// Background decommission walks must tolerate slow object migrations; the /// stall timeout is the drive-health bound, not the total listing duration. @@ -804,28 +809,29 @@ fn ensure_decommission_generation(meta: &PoolMeta, idx: usize, generation: Offse } } -async fn run_decommission_side_effect( +async fn run_decommission_side_effect( rx: &CancellationToken, operation_gate: &Arc>, operation: F, -) -> Result +) -> std::result::Result where F: FnOnce() -> Fut, - Fut: std::future::Future>, + Fut: std::future::Future>, + E: From, { let _operation_guard = tokio::select! { biased; - _ = rx.cancelled() => return Err(Error::OperationCanceled), + _ = rx.cancelled() => return Err(Error::OperationCanceled.into()), guard = operation_gate.read() => guard, }; if rx.is_cancelled() { - return Err(Error::OperationCanceled); + return Err(Error::OperationCanceled.into()); } let result = operation().await; if rx.is_cancelled() { - return Err(Error::OperationCanceled); + return Err(Error::OperationCanceled.into()); } result } @@ -1147,13 +1153,17 @@ fn should_retry_decommission_listing(err: &Error, attempt: usize, max_attempts: !is_err_bucket_not_found(err) && attempt + 1 < max_attempts } -async fn wait_decommission_listing_retry(rx: &CancellationToken, delay: std::time::Duration) -> bool { +async fn wait_decommission_retry_backoff(rx: &CancellationToken, delay: std::time::Duration) -> bool { tokio::select! { _ = rx.cancelled() => true, _ = tokio::time::sleep(delay) => false, } } +fn decommission_retry_backoff_delay(base: std::time::Duration, attempt: usize) -> std::time::Duration { + base.saturating_mul(u32::try_from(attempt).unwrap_or(u32::MAX)) +} + #[cfg(test)] async fn run_decommission_listing_with_retry( rx: CancellationToken, @@ -1269,7 +1279,7 @@ where error = ?err, "Decommission listing failed; retrying" ); - if wait_decommission_listing_retry(&rx, DECOMMISSION_LISTING_RETRY_DELAY).await { + if wait_decommission_retry_backoff(&rx, DECOMMISSION_LISTING_RETRY_DELAY).await { debug!( event = EVENT_DECOMMISSION_BUCKET, component = LOG_COMPONENT_ECSTORE, @@ -1347,6 +1357,123 @@ fn should_cleanup_decommission_source_entry(decommissioned: usize, total_version decommissioned.saturating_add(expired) == total_versions } +fn should_fail_decommission_pool_after_exhausted_source_changed(exhausted_entries: usize) -> bool { + exhausted_entries > DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum DecommissionEntryAttemptOutcome { + Complete, + SourceChanged, +} + +#[cfg(test)] +pub(crate) type DecommissionTestFaultDecision = Arc bool + Send + Sync>; + +#[cfg(test)] +static DECOMMISSION_TEST_FAULT_HOOK: std::sync::OnceLock>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +pub(crate) struct DecommissionTestFaultGuard(DecommissionTestFaultDecision); + +#[cfg(test)] +impl DecommissionTestFaultGuard { + pub(crate) fn install(decision: DecommissionTestFaultDecision) -> Self { + let mut slot = DECOMMISSION_TEST_FAULT_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission test fault hook mutex should not poison"); + assert!(slot.is_none(), "decommission test fault hook must be unique"); + let stored = Arc::clone(&decision); + *slot = Some(decision); + Self(stored) + } +} + +#[cfg(test)] +impl Drop for DecommissionTestFaultGuard { + fn drop(&mut self) { + let mut slot = DECOMMISSION_TEST_FAULT_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission test fault hook mutex should not poison"); + if slot.as_ref().is_some_and(|decision| Arc::ptr_eq(decision, &self.0)) { + *slot = None; + } + } +} + +#[cfg(test)] +fn decommission_test_wrap_result( + stage: &'static str, + bucket: &str, + object: &str, + attempt: usize, + result: Result, +) -> Result { + let decision = DECOMMISSION_TEST_FAULT_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission test fault hook mutex should not poison") + .clone(); + if result.is_ok() && decision.is_some_and(|decision| decision(stage, bucket, object, attempt)) { + return Err(Error::other(format!( + "injected decommission test fault at {stage} attempt {attempt} for {bucket}/{object}" + ))); + } + result +} + +#[cfg(test)] +pub(crate) type DecommissionCleanupMutationHook = Arc BoxFuture<'static, ()> + Send + Sync>; + +#[cfg(test)] +static DECOMMISSION_CLEANUP_MUTATION_HOOK: std::sync::OnceLock>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +pub(crate) struct DecommissionCleanupMutationGuard(DecommissionCleanupMutationHook); + +#[cfg(test)] +impl DecommissionCleanupMutationGuard { + pub(crate) fn install(hook: DecommissionCleanupMutationHook) -> Self { + let mut slot = DECOMMISSION_CLEANUP_MUTATION_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission cleanup mutation hook mutex should not poison"); + assert!(slot.is_none(), "decommission cleanup mutation hook must be unique"); + let stored = Arc::clone(&hook); + *slot = Some(hook); + Self(stored) + } +} + +#[cfg(test)] +impl Drop for DecommissionCleanupMutationGuard { + fn drop(&mut self) { + let mut slot = DECOMMISSION_CLEANUP_MUTATION_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission cleanup mutation hook mutex should not poison"); + if slot.as_ref().is_some_and(|hook| Arc::ptr_eq(hook, &self.0)) { + *slot = None; + } + } +} + +#[cfg(test)] +async fn run_decommission_cleanup_mutation_hook(bucket: &str, object: &str, attempt: usize) { + let hook = DECOMMISSION_CLEANUP_MUTATION_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission cleanup mutation hook mutex should not poison") + .clone(); + if let Some(hook) = hook { + hook(bucket, object, attempt).await; + } +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] #[allow( dead_code, @@ -3463,6 +3590,7 @@ impl ECStore { object_lock_config: Option, replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, expected_bucket_incarnation_id: Option, + source_changed_exhaustions: Arc, entry_budget: Arc, queue: Arc>>, entry_error: Arc>>, @@ -3553,6 +3681,7 @@ impl ECStore { object_lock_config.clone(), replication_config.clone(), expected_bucket_incarnation_id, + Arc::clone(&source_changed_exhaustions), ) .await; drop(entry_budget_permit); @@ -3590,6 +3719,7 @@ impl ECStore { object_lock_config: Option, replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, expected_bucket_incarnation_id: Option, + source_changed_exhaustions: Arc, entry_budget: Arc, entry_error: Arc>>, ) -> Result<()> { @@ -3609,6 +3739,7 @@ impl ECStore { let lifecycle_config = lifecycle_config.clone(); let object_lock_config = object_lock_config.clone(); let replication_config = replication_config.clone(); + let source_changed_exhaustions = Arc::clone(&source_changed_exhaustions); let queue = queue.clone(); let entry_budget = entry_budget.clone(); let entry_error = entry_error.clone(); @@ -3624,6 +3755,7 @@ impl ECStore { object_lock_config, replication_config, expected_bucket_incarnation_id, + source_changed_exhaustions, entry_budget, queue, entry_error, @@ -3765,8 +3897,15 @@ impl ECStore { Ok(()) } - #[allow(unused_assignments, clippy::too_many_arguments)] - #[tracing::instrument(skip(self, set, lifecycle_config, object_lock_config, replication_config))] + #[allow(clippy::too_many_arguments)] + #[tracing::instrument(skip( + self, + set, + lifecycle_config, + object_lock_config, + replication_config, + source_changed_exhaustions + ))] async fn decommission_entry( self: &Arc, rx: CancellationToken, @@ -3779,15 +3918,85 @@ impl ECStore { object_lock_config: Option, replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, expected_bucket_incarnation_id: Option, + source_changed_exhaustions: Arc, ) -> Result<()> { + let mut counted_versions = HashSet::new(); + + for entry_attempt in 1..=DECOMMISSION_ENTRY_MAX_ATTEMPTS { + match self + .decommission_entry_attempt( + rx.clone(), + idx, + generation, + entry.clone(), + bucket.clone(), + Arc::clone(&set), + lifecycle_config.clone(), + object_lock_config.clone(), + replication_config.clone(), + expected_bucket_incarnation_id, + entry_attempt, + source_changed_exhaustions.as_ref(), + &mut counted_versions, + ) + .await? + { + DecommissionEntryAttemptOutcome::Complete => return Ok(()), + DecommissionEntryAttemptOutcome::SourceChanged => { + let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, entry_attempt); + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "source_changed_retry", + pool_index = idx, + bucket = %bucket, + object = %entry.name, + attempt = entry_attempt, + max_attempts = DECOMMISSION_ENTRY_MAX_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + "Decommission source changed during cleanup preflight; retrying entry" + ); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } + } + } + } + + Err(Error::other(format!( + "decommission entry retry loop ended without a terminal result for {bucket}/{}", + entry.name + ))) + } + + #[allow(unused_assignments, clippy::too_many_arguments)] + async fn decommission_entry_attempt( + self: &Arc, + rx: CancellationToken, + idx: usize, + generation: OffsetDateTime, + entry: MetaCacheEntry, + bucket: String, + set: Arc, + lifecycle_config: Option, + object_lock_config: Option, + replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, + expected_bucket_incarnation_id: Option, + entry_attempt: usize, + source_changed_exhaustions: &AtomicUsize, + counted_versions: &mut HashSet<(Option, bool)>, + ) -> Result { debug!( event = EVENT_DECOMMISSION_ENTRY, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_POOLS, + state = "started", pool_index = idx, bucket = %bucket, object = %entry.name, - state = "started", + attempt = entry_attempt, + max_attempts = DECOMMISSION_ENTRY_MAX_ATTEMPTS, "Decommission entry started" ); if entry.is_dir() { @@ -3801,7 +4010,7 @@ impl ECStore { state = "skipped_directory", "Decommission entry skipped directory" ); - return Ok(()); + return Ok(DecommissionEntryAttemptOutcome::Complete); } if self.decommission_cancel_requested(idx, &rx).await { rx.cancel(); @@ -3823,6 +4032,7 @@ impl ECStore { let mut decommissioned: usize = 0; let mut expired: usize = 0; let mut cleanup_preflight_allowed_missing = Vec::new(); + let mut entry_blocked = false; for version in fivs.versions.iter() { if self.decommission_cancel_requested(idx, &rx).await { @@ -3874,33 +4084,49 @@ impl ECStore { let mut failure = false; let mut error = None; if version.deleted { - if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { - self.delete_object( + for version_attempt in 1..=DECOMMISSION_VERSION_COPY_ATTEMPTS { + let result = run_decommission_side_effect(&rx, &operation_gate, || async { + self.delete_object( + bucket.as_str(), + &version.name, + decommission_delete_marker_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), + ) + .await + }) + .await; + #[cfg(test)] + let result = decommission_test_wrap_result( + "delete_marker_copy", bucket.as_str(), - &version.name, - decommission_delete_marker_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), - ) - .await - }) - .await - { - if is_decommission_copy_cleanup_safe_error(&err) { - warn!( - event = EVENT_DECOMMISSION_ENTRY, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - bucket = %bucket, - object = %version.name, - version_id = ?version_id, - state = "ignored_delete_marker_copy", - error = ?err, - "Decommission delete marker copy ignored" - ); - ignore = true; - cleanup_ignored = true; - } else { - if is_decommission_target_capacity_error(&err) { + version.name.as_str(), + version_attempt, + result, + ); + + match result { + Ok(_) => { + failure = false; + error = None; + break; + } + Err(err) if is_decommission_copy_cleanup_safe_error(&err) => { + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "ignored_delete_marker_copy", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + error = ?err, + "Decommission delete marker copy ignored" + ); + ignore = true; + cleanup_ignored = true; + break; + } + Err(err) if is_decommission_target_capacity_error(&err) => { return Err(with_decommission_entry_context( "delete_marker_copy", bucket.as_str(), @@ -3908,10 +4134,47 @@ impl ECStore { err, )); } + Err(err) => { + failure = true; + if version_attempt == DECOMMISSION_VERSION_COPY_ATTEMPTS { + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "delete_marker_copy_failed", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + error = ?err, + "Decommission delete marker copy failed" + ); + error = Some(err); + break; + } - failure = true; - - error = Some(err) + let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_COPY_RETRY_DELAY, version_attempt); + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "delete_marker_copy_retry", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempt = version_attempt, + max_attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + error = ?err, + "Decommission delete marker copy failed; retrying" + ); + error = Some(err); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } + } } } @@ -3932,7 +4195,7 @@ impl ECStore { continue; } - { + if counted_versions.insert((version.version_id, version.deleted)) { let mut pool_meta = self.pool_meta.write().await; ensure_decommission_generation(&pool_meta, idx, generation)?; if let Err(err) = count_decommission_item(&mut pool_meta, idx, 0, failure) { @@ -3949,24 +4212,26 @@ impl ECStore { decommissioned += 1; } - debug!( - event = EVENT_DECOMMISSION_ENTRY, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - bucket = %bucket, - object = %version.name, - version_id = ?version_id, - result = ?error, - state = "delete_marker_copied", - "Decommission delete marker copied" - ); + if !failure { + debug!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "delete_marker_copied", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + result = ?error, + "Decommission delete marker copied" + ); + } continue; } - for _i in 0..3 { + for version_attempt in 1..=DECOMMISSION_VERSION_COPY_ATTEMPTS { if version.is_remote() { - if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { + let result = run_decommission_side_effect(&rx, &operation_gate, || async { self.decommission_tiered_object( bucket.as_str(), &version.name, @@ -3975,15 +4240,26 @@ impl ECStore { ) .await }) - .await - { - if is_decommission_copy_cleanup_safe_error(&err) { + .await; + #[cfg(test)] + let result = decommission_test_wrap_result( + "decommission_tiered_object", + bucket.as_str(), + version.name.as_str(), + version_attempt, + result, + ); + + match result { + Ok(_) => { + failure = false; + error = None; + } + Err(err) if is_decommission_copy_cleanup_safe_error(&err) => { ignore = true; cleanup_ignored = true; - break; } - - if is_decommission_target_capacity_error(&err) { + Err(err) if is_decommission_target_capacity_error(&err) => { return Err(with_decommission_entry_context( "decommission_tiered_object", bucket.as_str(), @@ -3991,10 +4267,48 @@ impl ECStore { err, )); } - - failure = true; - error!("decommission_pool: decommission_tiered_object err {:?}", &err); - error = Some(err); + Err(err) => { + failure = true; + if version_attempt == DECOMMISSION_VERSION_COPY_ATTEMPTS { + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "tiered_copy_failed", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + error = ?err, + "Decommission tiered version copy failed" + ); + error = Some(err); + } else { + let retry_delay = + decommission_retry_backoff_delay(DECOMMISSION_COPY_RETRY_DELAY, version_attempt); + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "tiered_copy_retry", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempt = version_attempt, + max_attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + error = ?err, + "Decommission tiered version copy failed; retrying" + ); + error = Some(err); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } + continue; + } + } } break; } @@ -4029,7 +4343,44 @@ impl ECStore { } failure = true; - error!("decommission_pool: get_object_reader err {:?}", &err); + if version_attempt == DECOMMISSION_VERSION_COPY_ATTEMPTS { + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "object_read_failed", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + error = ?err, + "Decommission source object read failed" + ); + error = Some(err); + continue; + } + + let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_COPY_RETRY_DELAY, version_attempt); + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "object_read_retry", + pool_index = idx, + bucket = %bucket, + object = %version.name, + version_id = ?version_id, + attempt = version_attempt, + max_attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + error = ?err, + "Decommission source object read failed; retrying" + ); + error = Some(err); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } continue; } }; @@ -4046,13 +4397,22 @@ impl ECStore { ) .await?; - if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { + let migrate_result = run_decommission_side_effect(&rx, &operation_gate, || async { self.clone() .decommission_object(idx, bucket, rd, expected_bucket_incarnation_id) .await }) - .await - { + .await; + #[cfg(test)] + let migrate_result = decommission_test_wrap_result( + DECOMMISSION_STAGE_MIGRATE_OBJECT, + bucket_name.as_str(), + object_name.as_str(), + version_attempt, + migrate_result, + ); + + if let Err(err) = migrate_result { if is_decommission_copy_cleanup_safe_error(&err) { ignore = true; cleanup_ignored = true; @@ -4069,8 +4429,45 @@ impl ECStore { } failure = true; + if version_attempt == DECOMMISSION_VERSION_COPY_ATTEMPTS { + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "object_migration_failed", + pool_index = idx, + bucket = %bucket_name, + object = %object_name, + version = %version.name, + attempt = version_attempt, + max_attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + error = ?err, + "Decommission object migration failed" + ); + error = Some(err); + continue; + } - error!("decommission_pool: decommission_object err {:?}", &err); + let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_COPY_RETRY_DELAY, version_attempt); + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "object_migration_retry", + pool_index = idx, + bucket = %bucket_name, + object = %object_name, + version = %version.name, + attempt = version_attempt, + max_attempts = DECOMMISSION_VERSION_COPY_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + error = ?err, + "Decommission object migration failed; retrying" + ); + error = Some(err); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } continue; } @@ -4107,7 +4504,7 @@ impl ECStore { continue; } - { + if counted_versions.insert((version.version_id, version.deleted)) { let mut pool_meta = self.pool_meta.write().await; ensure_decommission_generation(&pool_meta, idx, generation)?; if let Err(err) = count_decommission_item(&mut pool_meta, idx, decommission_item_size(version.size), failure) { @@ -4154,6 +4551,9 @@ impl ECStore { ) .await?; + #[cfg(test)] + run_decommission_cleanup_mutation_hook(bucket.as_str(), entry.name.as_str(), entry_attempt).await; + let source_cleanup_mutation_fence = self .acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref()) .await?; @@ -4174,16 +4574,58 @@ impl ECStore { "decommission", ) .await - .map_err(|err| match err { - data_movement::SourceCleanupError::SourceChanged => Error::other(format!( - "decommission: source cleanup preflight failed for {}/{}: source versions changed after migration started", - bucket, entry.name - )), - data_movement::SourceCleanupError::Storage(err) => err, - }) }) .await; - resolve_decommission_entry_cleanup_delete_result(cleanup_result, bucket.as_str(), entry.name.as_str())? + match cleanup_result { + Ok(_) => {} + Err(data_movement::SourceCleanupError::Storage(err)) => { + resolve_decommission_entry_cleanup_delete_result( + Err::<(), Error>(err), + bucket.as_str(), + entry.name.as_str(), + )?; + } + Err(data_movement::SourceCleanupError::SourceChanged) if entry_attempt < DECOMMISSION_ENTRY_MAX_ATTEMPTS => { + return Ok(DecommissionEntryAttemptOutcome::SourceChanged); + } + Err(data_movement::SourceCleanupError::SourceChanged) => { + let exhausted_entries = source_changed_exhaustions.fetch_add(1, Ordering::Relaxed) + 1; + if should_fail_decommission_pool_after_exhausted_source_changed(exhausted_entries) { + return Err(Error::other(format!( + "decommission source cleanup retries exhausted for {}/{} on all {} attempts; exhausted entries exceed pool limit {}", + bucket, entry.name, DECOMMISSION_ENTRY_MAX_ATTEMPTS, DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT + ))); + } + + { + let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; + count_decommission_item(&mut pool_meta, idx, 0, true).map_err(|err| { + with_decommission_entry_context( + "count_source_changed_exhaustion", + bucket.as_str(), + entry.name.as_str(), + err, + ) + })?; + } + + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "source_cleanup_exhausted", + pool_index = idx, + bucket = %bucket, + object = %entry.name, + attempts = DECOMMISSION_ENTRY_MAX_ATTEMPTS, + exhausted_entries, + exhaustion_limit = DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT, + "Decommission source cleanup retries exhausted; source retained and entry marked failed" + ); + entry_blocked = true; + } + } } else if decommissioned != fivs.versions.len() || expired > 0 { warn!( event = EVENT_DECOMMISSION_ENTRY, @@ -4266,13 +4708,13 @@ impl ECStore { event = EVENT_DECOMMISSION_ENTRY, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_POOLS, + state = if entry_blocked { "blocked" } else { "completed" }, pool_index = idx, bucket = %bucket, object = %entry.name, - state = "completed", - "Decommission entry completed" + "Decommission entry finished" ); - Ok(()) + Ok(DecommissionEntryAttemptOutcome::Complete) } #[cfg(test)] @@ -4283,17 +4725,43 @@ impl ECStore { bucket: String, set: Arc, ) -> Result<()> { - self.decommission_entry( + self.decommission_entry_with_retry_state_for_test( CancellationToken::new(), idx, - OffsetDateTime::now_utc(), + entry, + bucket, + set, + None, + Arc::new(AtomicUsize::new(0)), + ) + .await + } + + #[cfg(test)] + #[allow(clippy::too_many_arguments)] + pub(crate) async fn decommission_entry_with_retry_state_for_test( + self: &Arc, + rx: CancellationToken, + idx: usize, + entry: MetaCacheEntry, + bucket: String, + set: Arc, + expected_bucket_incarnation_id: Option, + source_changed_exhaustions: Arc, + ) -> Result<()> { + let generation = self.active_decommission_generation(idx).await?; + self.decommission_entry( + rx, + idx, + generation, entry, bucket, set, None, None, None, - None, + expected_bucket_incarnation_id, + source_changed_exhaustions, ) .await } @@ -4306,6 +4774,7 @@ impl ECStore { pool: Arc, bi: DecomBucketInfo, entry_budget: Arc, + source_changed_exhaustions: Arc, ) -> Result<()> { let entry_error = Arc::new(tokio::sync::Mutex::new(None::)); let generation = self.active_decommission_generation(idx).await?; @@ -4356,6 +4825,7 @@ impl ECStore { let object_lock_config = object_lock_config.clone(); let replication_config = replication_config.clone(); let entry_budget = entry_budget.clone(); + let source_changed_exhaustions = Arc::clone(&source_changed_exhaustions); let entry_error = entry_error.clone(); let worker = tokio::spawn(async move { store @@ -4370,6 +4840,7 @@ impl ECStore { object_lock_config, replication_config, expected_bucket_incarnation_id, + source_changed_exhaustions, entry_budget, entry_error, ) @@ -4847,6 +5318,7 @@ impl ECStore { pool: Arc, bucket: DecomBucketInfo, entry_budget: Arc, + source_changed_exhaustions: Arc, ) -> Result<()> { let is_decommissioned = { let pool_meta = self.pool_meta.read().await; @@ -4873,7 +5345,7 @@ impl ECStore { warn!("decommission: currently on bucket {}", &bucket.name); if let Err(err) = self - .decommission_pool(rx.clone(), idx, pool, bucket.clone(), entry_budget) + .decommission_pool(rx.clone(), idx, pool, bucket.clone(), entry_budget, source_changed_exhaustions) .await { error!("decommission: decommission_pool err {:?}", &err); @@ -4910,13 +5382,19 @@ impl ECStore { buckets: Vec, limit: usize, entry_budget: Arc, + source_changed_exhaustions: Arc, ) -> Result<()> { let store = Arc::clone(self); run_decommission_buckets_bounded(rx, buckets, limit, move |bucket, rx| { let store = Arc::clone(&store); let pool = pool.clone(); let entry_budget = entry_budget.clone(); - Box::pin(async move { store.decommission_pending_bucket(rx, idx, pool, bucket, entry_budget).await }) + let source_changed_exhaustions = Arc::clone(&source_changed_exhaustions); + Box::pin(async move { + store + .decommission_pending_bucket(rx, idx, pool, bucket, entry_budget, source_changed_exhaustions) + .await + }) }) .await } @@ -4934,11 +5412,19 @@ impl ECStore { let pool_meta = self.pool_meta.read().await; pool_meta.pending_buckets(idx) }; + let source_changed_exhaustions = Arc::new(AtomicUsize::new(0)); let bucket_concurrency = decommission_bucket_concurrency_limit(); if bucket_concurrency <= 1 { for bucket in pending { - self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone()) - .await?; + self.decommission_pending_bucket( + rx.clone(), + idx, + pool.clone(), + bucket, + entry_budget.clone(), + Arc::clone(&source_changed_exhaustions), + ) + .await?; } return Ok(()); } @@ -4951,12 +5437,20 @@ impl ECStore { regular_buckets, bucket_concurrency, entry_budget.clone(), + Arc::clone(&source_changed_exhaustions), ) .await?; for bucket in meta_buckets { - self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone()) - .await?; + self.decommission_pending_bucket( + rx.clone(), + idx, + pool.clone(), + bucket, + entry_budget.clone(), + Arc::clone(&source_changed_exhaustions), + ) + .await?; } Ok(()) @@ -6350,17 +6844,18 @@ mod pools_tests { use super::resolve_decommission_listing_error; use super::{ DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP, - DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler, - DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, - PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info, - await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, - cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state, - count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size, - decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, - default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry, - ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation, - ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, - ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, + DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, + DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT, DecomBucketInfo, DecommissionCanceler, DecommissionEntryEnqueueResult, + DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, + PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info, await_decommission_worker, + bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, + clamp_decommission_entry_concurrency, classify_decommission_terminal_state, count_decommission_item, + decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size, + decommission_meta_bucket_options, decommission_retry_backoff_delay, decommission_start_pool_state, dedup_indices, + default_decommission_bucket_concurrency, default_decommission_entry_concurrency, drain_decommission_entry_queue, + enqueue_decommission_entry, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, + ensure_decommission_generation, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, + ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, ensure_decommission_start_pool_states, ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity, ensure_decommission_terminal_operation_supported, ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, @@ -6379,12 +6874,13 @@ mod pools_tests { rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry, run_decommission_listing_with_retry_and_drain, run_decommission_side_effect, should_cleanup_decommission_source_entry, should_continue_decommission_queue, should_count_decommission_version_complete, - should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal, - should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine, - spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler, - take_decommission_canceler, track_decommission_current_object, track_decommission_current_object_stage, - update_decommission_for_operation, validate_start_decommission_request, wait_decommission_listing_retry, - wait_decommission_worker_drain, with_decommission_entry_context, + should_fail_decommission_pool_after_exhausted_source_changed, should_preserve_decommission_canceled_state, + should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload, + should_retry_decommission_listing, should_skip_canceled_decommission_routine, spawn_decommission_index_cancelers, + split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler, + track_decommission_current_object, track_decommission_current_object_stage, update_decommission_for_operation, + validate_start_decommission_request, wait_decommission_retry_backoff, wait_decommission_worker_drain, + with_decommission_entry_context, }; use crate::data_movement; use crate::disk::endpoint::Endpoint; @@ -7753,11 +8249,30 @@ mod pools_tests { } #[tokio::test] - async fn test_wait_decommission_listing_retry_reports_canceled_without_sleeping() { + async fn test_wait_decommission_retry_backoff_reports_canceled_without_sleeping() { let token = CancellationToken::new(); token.cancel(); - assert!(wait_decommission_listing_retry(&token, StdDuration::from_secs(30)).await); + assert!(wait_decommission_retry_backoff(&token, StdDuration::from_secs(30)).await); + } + + #[test] + fn test_decommission_retry_backoff_delay_grows_linearly() { + let base = StdDuration::from_millis(100); + + assert_eq!(decommission_retry_backoff_delay(base, 1), base); + assert_eq!(decommission_retry_backoff_delay(base, 3), StdDuration::from_millis(300)); + assert_eq!(decommission_retry_backoff_delay(base, usize::MAX), base.saturating_mul(u32::MAX)); + } + + #[test] + fn test_source_changed_exhaustion_fails_only_after_pool_limit() { + assert!(!should_fail_decommission_pool_after_exhausted_source_changed( + DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT + )); + assert!(should_fail_decommission_pool_after_exhausted_source_changed( + DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT + 1 + )); } #[tokio::test(start_paused = true)] diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 0422ba94e..67d28dfe0 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -4736,7 +4736,15 @@ impl SetDisks { achieved: 0, }); } - let parts_metadata = vec![fi.clone(); disks.len()]; + // Rebuilt tiered metadata starts with index zero, but shuffling validates + // each source slot before assigning the shuffled index below. + let parts_metadata: Vec = (0..disks.len()) + .map(|disk_index| { + let mut part = fi.clone(); + part.erasure.index = fi.erasure.distribution[disk_index]; + part + }) + .collect(); let (shuffle_disks, parts_metadata) = Self::shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi); let mut errs = Vec::with_capacity(shuffle_disks.len()); diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 451e42242..bca8647db 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -626,7 +626,7 @@ mod tests { io::Cursor, sync::{ Arc, - atomic::{AtomicBool, Ordering}, + atomic::{AtomicBool, AtomicUsize, Ordering}, }, time::Duration, }; @@ -1307,6 +1307,85 @@ mod tests { }); } + const DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER: &str = "delete_marker_copy"; + const DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT: &str = "migrate_object"; + #[cfg(feature = "test-util")] + const DECOMMISSION_TEST_FAULT_STAGE_TIERED: &str = "decommission_tiered_object"; + + async fn seed_decommission_source( + store: &Arc, + bucket: &str, + object: &str, + body: Vec, + opts: &ObjectOptions, + ) { + let mut reader = PutObjReader::from_vec(body); + store.pools[0] + .put_object(bucket, object, &mut reader, opts) + .await + .expect("seed decommission source object"); + } + + async fn run_decommission_entry_retry_test( + store: &Arc, + rx: CancellationToken, + bucket: &str, + object: &str, + expected_bucket_incarnation_id: Option, + source_changed_exhaustions: Arc, + ) -> crate::error::Result<()> { + let source_set = store.pools[0].get_disks_by_key(object); + store + .decommission_entry_with_retry_state_for_test( + rx, + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + bucket.to_string(), + source_set, + expected_bucket_incarnation_id, + source_changed_exhaustions, + ) + .await + } + + async fn read_decommission_target_body( + store: &Arc, + bucket: &str, + object: &str, + opts: &ObjectOptions, + ) -> Vec { + let mut reader = store.pools[1] + .get_object_reader(bucket, object, None, HeaderMap::new(), opts) + .await + .expect("read decommission target object"); + let mut body = Vec::new(); + reader + .stream + .read_to_end(&mut body) + .await + .expect("drain decommission target body"); + body + } + + async fn assert_decommission_source_absent( + store: &Arc, + bucket: &str, + object: &str, + opts: &ObjectOptions, + ) { + let err = store.pools[0] + .get_object_info(bucket, object, opts) + .await + .expect_err("decommission source must be retained until target commit, then removed"); + assert!( + matches!(err, StorageError::ObjectNotFound(_, _) | StorageError::VersionNotFound(_, _, _)), + "unexpected decommission source result: {err:?}" + ); + } + async fn write_decommission_test_multipart_source( store: &Arc, pool_idx: usize, @@ -3102,6 +3181,537 @@ mod tests { shutdown.cancel(); } + #[test] + #[serial_test::serial(storage_class_env)] + fn decommission_entry_retries_source_changed_without_canceling_other_bucket() { + let handle = std::thread::Builder::new() + .name("decommission_entry_retries_source_changed_without_canceling_other_bucket".to_string()) + .stack_size(32 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(2) + .build() + .expect("test runtime should build"); + runtime.block_on(async { + let temp_dir = tempfile::tempdir().expect("create decommission retry store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "decommission-entry-retry", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let changed_bucket = format!("decom-retry-a-{}", uuid::Uuid::new_v4()); + let other_bucket = format!("decom-retry-b-{}", uuid::Uuid::new_v4()); + for bucket in [&changed_bucket, &other_bucket] { + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission retry bucket"); + } + + let changed_object = "changed.bin"; + let first_version = uuid::Uuid::new_v4(); + let second_version = uuid::Uuid::new_v4(); + let base_time = OffsetDateTime::now_utc(); + seed_decommission_source( + &store, + &changed_bucket, + changed_object, + b"first generation".to_vec(), + &ObjectOptions { + versioned: true, + version_id: Some(first_version.to_string()), + mod_time: Some(base_time), + ..Default::default() + }, + ) + .await; + let other_object = "other.bin"; + seed_decommission_source( + &store, + &other_bucket, + other_object, + b"other bucket generation".to_vec(), + &ObjectOptions::default(), + ) + .await; + mark_test_pool_decommissioning(&store, 0).await; + + let mutation_calls = Arc::new(AtomicUsize::new(0)); + let mutation_calls_for_hook = Arc::clone(&mutation_calls); + let mutation_store = Arc::clone(&store); + let mutation_bucket = changed_bucket.clone(); + let _mutation_guard = crate::core::pools::DecommissionCleanupMutationGuard::install(Arc::new( + move |bucket, object, attempt| { + let is_target = bucket == mutation_bucket.as_str() && object == changed_object; + let calls = Arc::clone(&mutation_calls_for_hook); + let store = Arc::clone(&mutation_store); + let bucket = mutation_bucket.clone(); + Box::pin(async move { + if !is_target { + return; + } + calls.fetch_add(1, Ordering::SeqCst); + if attempt == 1 { + seed_decommission_source( + &store, + &bucket, + changed_object, + b"second generation".to_vec(), + &ObjectOptions { + versioned: true, + version_id: Some(second_version.to_string()), + mod_time: Some(base_time + time::Duration::seconds(1)), + ..Default::default() + }, + ) + .await; + } + }) + }, + )); + + let ordinary_faults = Arc::new(AtomicUsize::new(0)); + let ordinary_faults_for_hook = Arc::clone(&ordinary_faults); + let fault_bucket = other_bucket.clone(); + let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new( + move |stage, bucket, object, attempt| { + let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT + && bucket == fault_bucket.as_str() + && object == other_object + && attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS; + if injected { + ordinary_faults_for_hook.fetch_add(1, Ordering::SeqCst); + } + injected + }, + )); + + let rx = CancellationToken::new(); + let source_changed_exhaustions = Arc::new(AtomicUsize::new(0)); + let changed_incarnation = Some( + store + .bucket_incarnation_id(&changed_bucket) + .await + .expect("changed bucket incarnation"), + ); + let other_incarnation = Some( + store + .bucket_incarnation_id(&other_bucket) + .await + .expect("other bucket incarnation"), + ); + let (changed_result, other_result) = tokio::join!( + run_decommission_entry_retry_test( + &store, + rx.clone(), + &changed_bucket, + changed_object, + changed_incarnation, + Arc::clone(&source_changed_exhaustions), + ), + run_decommission_entry_retry_test( + &store, + rx.clone(), + &other_bucket, + other_object, + other_incarnation, + Arc::clone(&source_changed_exhaustions), + ) + ); + changed_result.expect("SourceChanged entry retry must converge"); + other_result.expect("other bucket entry must continue through ordinary copy retries"); + + assert!(!rx.is_cancelled(), "entry-level SourceChanged must not cancel the shared worker token"); + assert_eq!(mutation_calls.load(Ordering::SeqCst), 2, "entry must be re-listed after SourceChanged"); + assert_eq!(ordinary_faults.load(Ordering::SeqCst), 2, "ordinary copy must consume the retry budget"); + assert_eq!(source_changed_exhaustions.load(Ordering::SeqCst), 0); + + for (version_id, expected_body) in [ + (first_version, b"first generation".as_slice()), + (second_version, b"second generation".as_slice()), + ] { + let opts = ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }; + assert_decommission_source_absent(&store, &changed_bucket, changed_object, &opts).await; + assert_eq!( + read_decommission_target_body(&store, &changed_bucket, changed_object, &opts).await, + expected_body + ); + } + assert_decommission_source_absent(&store, &other_bucket, other_object, &ObjectOptions::default()).await; + assert_eq!( + read_decommission_target_body(&store, &other_bucket, other_object, &ObjectOptions::default()).await, + b"other bucket generation" + ); + + shutdown.cancel(); + }); + }) + .expect("spawn decommission retry test thread"); + if let Err(payload) = handle.join() { + std::panic::resume_unwind(payload); + } + } + + #[test] + #[serial_test::serial(storage_class_env)] + fn decommission_entry_exhausted_source_changed_retains_source_and_records_failure() { + let handle = std::thread::Builder::new() + .name("decommission_entry_exhausted_source_changed_retains_source_and_records_failure".to_string()) + .stack_size(32 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(2) + .build() + .expect("test runtime should build"); + runtime.block_on(async { + let temp_dir = tempfile::tempdir().expect("create decommission exhaustion store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "decommission-entry-exhaustion", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decom-exhausted-{}", uuid::Uuid::new_v4()); + let object = "exhausted.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission exhaustion bucket"); + let original_version = uuid::Uuid::new_v4(); + let base_time = OffsetDateTime::now_utc(); + seed_decommission_source( + &store, + &bucket, + object, + b"original generation".to_vec(), + &ObjectOptions { + versioned: true, + version_id: Some(original_version.to_string()), + mod_time: Some(base_time), + ..Default::default() + }, + ) + .await; + mark_test_pool_decommissioning(&store, 0).await; + + let mutation_calls = Arc::new(AtomicUsize::new(0)); + let mutation_calls_for_hook = Arc::clone(&mutation_calls); + let mutation_store = Arc::clone(&store); + let mutation_bucket = bucket.clone(); + let _mutation_guard = crate::core::pools::DecommissionCleanupMutationGuard::install(Arc::new( + move |called_bucket, called_object, _attempt| { + let is_target = called_bucket == mutation_bucket.as_str() && called_object == object; + let calls = Arc::clone(&mutation_calls_for_hook); + let store = Arc::clone(&mutation_store); + let bucket = mutation_bucket.clone(); + Box::pin(async move { + if !is_target { + return; + } + let call = calls.fetch_add(1, Ordering::SeqCst) + 1; + let offset = i64::try_from(call).expect("entry retry count should fit i64"); + seed_decommission_source( + &store, + &bucket, + object, + format!("concurrent generation {call}").into_bytes(), + &ObjectOptions { + versioned: true, + version_id: Some(uuid::Uuid::new_v4().to_string()), + mod_time: Some(base_time + time::Duration::seconds(offset)), + ..Default::default() + }, + ) + .await; + }) + }, + )); + + let rx = CancellationToken::new(); + let source_changed_exhaustions = Arc::new(AtomicUsize::new(0)); + let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation")); + run_decommission_entry_retry_test( + &store, + rx.clone(), + &bucket, + object, + incarnation, + Arc::clone(&source_changed_exhaustions), + ) + .await + .expect("entry-level exhaustion must stay local below the pool threshold"); + + assert!(!rx.is_cancelled(), "one exhausted entry must not cancel other bucket workers"); + assert_eq!(mutation_calls.load(Ordering::SeqCst), crate::core::pools::DECOMMISSION_ENTRY_MAX_ATTEMPTS); + assert_eq!(source_changed_exhaustions.load(Ordering::SeqCst), 1); + store.pools[0] + .get_object_info( + &bucket, + object, + &ObjectOptions { + versioned: true, + version_id: Some(original_version.to_string()), + ..Default::default() + }, + ) + .await + .expect("retry exhaustion must retain the original source version"); + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("decommission progress must be initialized"); + assert_eq!(info.items_decommission_failed, 1, "exhausted entry must be visible as failed"); + drop(pool_meta); + + shutdown.cancel(); + }); + }) + .expect("spawn decommission retry test thread"); + if let Err(payload) = handle.join() { + std::panic::resume_unwind(payload); + } + } + + #[test] + #[serial_test::serial(storage_class_env)] + fn decommission_entry_delete_marker_copy_retries_real_path() { + let handle = std::thread::Builder::new() + .name("decommission_entry_delete_marker_copy_retries_real_path".to_string()) + .stack_size(32 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(2) + .build() + .expect("test runtime should build"); + runtime.block_on(async { + let temp_dir = tempfile::tempdir().expect("create delete marker retry store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "decommission-delete-marker-retry", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decom-marker-{}", uuid::Uuid::new_v4()); + let object = "marker.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create delete marker retry bucket"); + let data_version = uuid::Uuid::new_v4(); + let marker_version = uuid::Uuid::new_v4(); + let base_time = OffsetDateTime::now_utc(); + seed_decommission_source( + &store, + &bucket, + object, + b"delete marker data".to_vec(), + &ObjectOptions { + versioned: true, + version_id: Some(data_version.to_string()), + mod_time: Some(base_time), + ..Default::default() + }, + ) + .await; + store.pools[0] + .delete_object( + &bucket, + object, + ObjectOptions { + versioned: true, + version_id: Some(marker_version.to_string()), + delete_marker: true, + mod_time: Some(base_time + time::Duration::seconds(1)), + ..Default::default() + }, + ) + .await + .expect("seed source delete marker"); + mark_test_pool_decommissioning(&store, 0).await; + + let fault_calls = Arc::new(AtomicUsize::new(0)); + let fault_calls_for_hook = Arc::clone(&fault_calls); + let fault_bucket = bucket.clone(); + let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new( + move |stage, called_bucket, called_object, attempt| { + let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER + && called_bucket == fault_bucket.as_str() + && called_object == object + && attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS; + if injected { + fault_calls_for_hook.fetch_add(1, Ordering::SeqCst); + } + injected + }, + )); + + let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation")); + run_decommission_entry_retry_test( + &store, + CancellationToken::new(), + &bucket, + object, + incarnation, + Arc::new(AtomicUsize::new(0)), + ) + .await + .expect("delete marker copy retries must converge"); + + assert_eq!( + fault_calls.load(Ordering::SeqCst), + crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS - 1 + ); + let marker_opts = ObjectOptions { + versioned: true, + version_id: Some(marker_version.to_string()), + ..Default::default() + }; + let target_marker = store.pools[1] + .get_object_info(&bucket, object, &marker_opts) + .await + .expect("target delete marker must exist"); + assert!(target_marker.delete_marker); + assert_decommission_source_absent(&store, &bucket, object, &marker_opts).await; + + let data_opts = ObjectOptions { + versioned: true, + version_id: Some(data_version.to_string()), + ..Default::default() + }; + assert_eq!( + read_decommission_target_body(&store, &bucket, object, &data_opts).await, + b"delete marker data" + ); + assert_decommission_source_absent(&store, &bucket, object, &data_opts).await; + + shutdown.cancel(); + }); + }) + .expect("spawn decommission retry test thread"); + if let Err(payload) = handle.join() { + std::panic::resume_unwind(payload); + } + } + + #[cfg(feature = "test-util")] + #[test] + #[serial_test::serial(storage_class_env)] + fn decommission_entry_tiered_copy_retries_real_path() { + let handle = std::thread::Builder::new() + .name("decommission_entry_tiered_copy_retries_real_path".to_string()) + .stack_size(32 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(2) + .build() + .expect("test runtime should build"); + runtime.block_on(async { + let temp_dir = tempfile::tempdir().expect("create tiered retry store dir"); + let (ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "decommission-tiered-retry", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decom-tiered-{}", uuid::Uuid::new_v4()); + let object = "tiered.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create tiered retry bucket"); + let mut reader = PutObjReader::from_vec(b"tiered generation".to_vec()); + let original = store.pools[0] + .put_object(&bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("seed tiered source object"); + let tier_name = format!("DECOM{}", &uuid::Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await; + store.pools[0] + .transition_object( + &bucket, + object, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name, + etag: original.etag.clone().expect("tiered source ETag"), + ..Default::default() + }, + version_id: original.version_id.map(|version_id| version_id.to_string()), + mod_time: original.mod_time, + ..Default::default() + }, + ) + .await + .expect("transition source object to mock tier"); + mark_test_pool_decommissioning(&store, 0).await; + + let fault_calls = Arc::new(AtomicUsize::new(0)); + let fault_calls_for_hook = Arc::clone(&fault_calls); + let fault_bucket = bucket.clone(); + let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new( + move |stage, called_bucket, called_object, attempt| { + let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_TIERED + && called_bucket == fault_bucket.as_str() + && called_object == object + && attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS; + if injected { + fault_calls_for_hook.fetch_add(1, Ordering::SeqCst); + } + injected + }, + )); + + let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation")); + run_decommission_entry_retry_test( + &store, + CancellationToken::new(), + &bucket, + object, + incarnation, + Arc::new(AtomicUsize::new(0)), + ) + .await + .expect("tiered copy retries must converge"); + + assert_eq!( + fault_calls.load(Ordering::SeqCst), + crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS - 1 + ); + let target = store.pools[1] + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect("tiered target metadata must exist"); + assert_eq!(target.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE); + assert_decommission_source_absent(&store, &bucket, object, &ObjectOptions::default()).await; + + shutdown.cancel(); + }); + }) + .expect("spawn decommission retry test thread"); + if let Err(payload) = handle.join() { + std::panic::resume_unwind(payload); + } + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial_test::serial(storage_class_env)] async fn decommission_outer_fence_loss_blocks_target_put_commit() {