From 9815694301f447a20ab783813f41487b6e6b2638 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 22 Aug 2026 22:33:24 +0800 Subject: [PATCH] fix(ecstore): supervise decommission worker cleanup (#6372) --- crates/ecstore/src/core/pools.rs | 1092 +++++++++++++++++++++++------- crates/ecstore/src/store/mod.rs | 4 +- 2 files changed, 842 insertions(+), 254 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 04a94cc05..66fe2bbe9 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -58,7 +58,6 @@ use http::HeaderMap; #[cfg(test)] use rmp_serde::Deserializer; use rmp_serde::Serializer; -use rustfs_common::defer; use rustfs_common::heal_channel::HealOpts; use rustfs_filemeta::{FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams}; use rustfs_utils::path::{encode_dir_object, path_join, path_to_bucket_object, path_to_bucket_object_with_base_path}; @@ -66,13 +65,14 @@ use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, Replicatio use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet}; use std::fmt::Display; +use std::future::Future; #[cfg(test)] use std::io::Cursor; use std::io::Write; use std::path::PathBuf; use std::sync::{ Arc, - atomic::{AtomicUsize, Ordering}, + atomic::{AtomicBool, AtomicUsize, Ordering}, }; use time::{Duration, OffsetDateTime}; use tokio::sync::{OwnedSemaphorePermit, Semaphore}; @@ -96,6 +96,7 @@ const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4; 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); +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. const DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60); @@ -104,6 +105,74 @@ pub const POOL_META_NAME: &str = "pool.bin"; pub const POOL_META_FORMAT: u16 = 1; pub const POOL_META_VERSION: u16 = 1; +#[derive(Clone, Debug)] +pub struct DecommissionCanceler { + operation: Arc, +} + +#[derive(Debug)] +struct DecommissionOperation { + token: CancellationToken, + active: AtomicBool, +} + +impl DecommissionCanceler { + fn new(token: CancellationToken) -> Self { + Self { + operation: Arc::new(DecommissionOperation { + token, + active: AtomicBool::new(true), + }), + } + } + + fn token(&self) -> &CancellationToken { + &self.operation.token + } + + fn is_active(&self) -> bool { + self.operation.active.load(Ordering::Acquire) + } + + #[cfg(test)] + fn is_cancelled(&self) -> bool { + self.token().is_cancelled() + } + + fn cancel(&self) { + self.token().cancel(); + } + + fn release(&self) { + self.cancel(); + self.operation.active.store(false, Ordering::Release); + } + + fn owns_same_operation(&self, other: &Self) -> bool { + Arc::ptr_eq(&self.operation, &other.operation) + } +} + +struct DecommissionCancelerGuard { + canceler: DecommissionCanceler, +} + +impl DecommissionCancelerGuard { + fn new(canceler: DecommissionCanceler) -> Self { + Self { canceler } + } + + fn canceler(&self) -> &DecommissionCanceler { + &self.canceler + } +} + +impl Drop for DecommissionCancelerGuard { + fn drop(&mut self) { + self.canceler.release(); + } +} + fn dedup_indices(indices: &[usize]) -> Vec { let mut seen = HashSet::with_capacity(indices.len()); let mut output = Vec::with_capacity(indices.len()); @@ -119,18 +188,18 @@ fn dedup_indices(indices: &[usize]) -> Vec { fn bind_decommission_cancelers( indices: &[usize], parent: &CancellationToken, - cancelers: &mut [Option], -) -> Vec<(usize, CancellationToken)> { + cancelers: &mut [Option], +) -> Vec<(usize, DecommissionCanceler)> { let mut bound = Vec::with_capacity(indices.len()); for idx in indices { if let Some(slot) = cancelers.get_mut(*idx) { if let Some(existing) = slot.take() { - existing.cancel(); + existing.release(); } - let token = parent.child_token(); - *slot = Some(token.clone()); - bound.push((*idx, token)); + let canceler = DecommissionCanceler::new(parent.child_token()); + *slot = Some(canceler.clone()); + bound.push((*idx, canceler)); } } @@ -140,47 +209,104 @@ fn bind_decommission_cancelers( fn bind_missing_decommission_cancelers( indices: &[usize], parent: &CancellationToken, - cancelers: &mut [Option], -) -> Vec<(usize, CancellationToken)> { + cancelers: &mut [Option], +) -> Vec<(usize, DecommissionCanceler)> { let mut bound = Vec::with_capacity(indices.len()); for idx in indices { let Some(slot) = cancelers.get_mut(*idx) else { continue; }; - if slot.is_some() { + if slot.as_ref().is_some_and(DecommissionCanceler::is_active) { break; } - let token = parent.child_token(); - *slot = Some(token.clone()); - bound.push((*idx, token)); + if let Some(stale) = slot.take() { + stale.release(); + } + let canceler = DecommissionCanceler::new(parent.child_token()); + *slot = Some(canceler.clone()); + bound.push((*idx, canceler)); } bound } -fn take_decommission_canceler(cancelers: &mut [Option], idx: usize) -> Option { +fn take_decommission_canceler(cancelers: &mut [Option], idx: usize) -> Option { cancelers.get_mut(idx).and_then(Option::take) } -fn has_active_decommission_canceler(cancelers: &[Option]) -> bool { - cancelers.iter().any(Option::is_some) +fn take_decommission_canceler_for_operation( + cancelers: &mut [Option], + idx: usize, + owner: &DecommissionCanceler, +) -> Option { + let slot = cancelers.get_mut(idx)?; + if slot.as_ref().is_some_and(|canceler| canceler.owns_same_operation(owner)) { + slot.take() + } else { + None + } } -fn cancel_decommission_canceler(canceler: Option) -> bool { +fn decommission_canceler_is_owned_by( + cancelers: &[Option], + idx: usize, + owner: &DecommissionCanceler, +) -> bool { + cancelers + .get(idx) + .and_then(Option::as_ref) + .is_some_and(|canceler| canceler.owns_same_operation(owner)) +} + +fn update_decommission_for_operation( + cancelers: &[Option], + pool_meta: &mut PoolMeta, + idx: usize, + owner: Option<&DecommissionCanceler>, + update: impl FnOnce(&mut PoolMeta) -> T, +) -> Option { + if let Some(owner) = owner + && !decommission_canceler_is_owned_by(cancelers, idx, owner) + { + owner.release(); + return None; + } + + Some(update(pool_meta)) +} + +fn has_active_decommission_canceler(cancelers: &[Option]) -> bool { + cancelers.iter().flatten().any(DecommissionCanceler::is_active) +} + +fn cancel_decommission_canceler(canceler: Option) -> bool { if let Some(canceler) = canceler { - canceler.cancel(); + canceler.release(); true } else { false } } -fn take_and_cancel_decommission_canceler(cancelers: &mut [Option], idx: usize) -> bool { +fn take_and_cancel_decommission_canceler(cancelers: &mut [Option], idx: usize) -> bool { let canceler = take_decommission_canceler(cancelers, idx); cancel_decommission_canceler(canceler) } +fn take_and_cancel_decommission_canceler_for_operation( + cancelers: &mut [Option], + idx: usize, + owner: &DecommissionCanceler, +) -> bool { + let canceler = take_decommission_canceler_for_operation(cancelers, idx, owner); + if canceler.is_none() { + owner.release(); + return false; + } + cancel_decommission_canceler(canceler) +} + fn ensure_decommission_routines_scheduled(bound_count: usize, expected_count: usize) -> Result<()> { if bound_count == 0 || bound_count != expected_count { return Err(Error::other(format!( @@ -191,6 +317,36 @@ fn ensure_decommission_routines_scheduled(bound_count: usize, expected_count: us Ok(()) } +fn guard_decommission_cancelers(index_cancelers: Vec<(usize, DecommissionCanceler)>) -> Vec<(usize, DecommissionCancelerGuard)> { + index_cancelers + .into_iter() + .map(|(idx, canceler)| (idx, DecommissionCancelerGuard::new(canceler))) + .collect() +} + +async fn await_decommission_worker(idx: usize, worker: tokio::task::JoinHandle>) -> Result<()> { + worker + .await + .map_err(|err| Error::other(format!("decommission worker {idx} task join error: {err}")))? +} + +fn reserve_decommission_start_cancelers( + pool_meta: &PoolMeta, + indices: &[usize], + local_indices: &[usize], + parent: &CancellationToken, + cancelers: &mut [Option], +) -> Result> { + ensure_decommission_start_pool_states(pool_meta, indices)?; + if local_indices.is_empty() { + return Ok(Vec::new()); + } + let bound = bind_decommission_cancelers(local_indices, parent, cancelers); + let guards = guard_decommission_cancelers(bound); + ensure_decommission_routines_scheduled(guards.len(), local_indices.len())?; + Ok(guards) +} + fn default_decommission_bucket_concurrency(cpu_count: usize) -> usize { cpu_count.clamp(1, DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP) } @@ -309,11 +465,15 @@ fn first_resumable_decommission_queue_indices(meta: &PoolMeta) -> Vec { indices } -fn missing_decommission_worker_prefix(indices: &[usize], cancelers: &[Option]) -> Vec { +fn missing_decommission_worker_prefix(indices: &[usize], cancelers: &[Option]) -> Vec { let mut missing = Vec::with_capacity(indices.len()); for idx in indices { - if cancelers.get(*idx).and_then(Option::as_ref).is_some() { + if cancelers + .get(*idx) + .and_then(Option::as_ref) + .is_some_and(DecommissionCanceler::is_active) + { break; } missing.push(*idx); @@ -360,29 +520,25 @@ fn build_decommission_start_state( fn spawn_decommission_index_cancelers( store: Arc, rx: CancellationToken, - index_cancelers: Vec<(usize, CancellationToken)>, -) { + index_cancelers: Vec<(usize, DecommissionCancelerGuard)>, +) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { let mut stop_queue = false; - for (idx, canceler) in index_cancelers { + for (idx, canceler_guard) in index_cancelers { + let canceler = canceler_guard.canceler().clone(); if stop_queue || rx.is_cancelled() { canceler.cancel(); - if let Err(err) = store.decommission_cancel(idx).await { - warn!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "queued_cancel_failed", - error = %err, - "Failed to cancel queued decommission" - ); - } + store.retry_decommission_cancel_for_operation(idx, &canceler).await; continue; } - if let Err(err) = store.do_decommission_in_routine(canceler, idx).await { + let worker = tokio::spawn({ + let store = store.clone(); + let canceler = canceler.clone(); + async move { store.do_decommission_in_routine(canceler, idx).await } + }); + if let Err(err) = await_decommission_worker(idx, worker).await { error!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -392,6 +548,7 @@ fn spawn_decommission_index_cancelers( error = %err, "Decommission routine failed" ); + store.retry_decommission_failed_for_operation(idx, &canceler).await; stop_queue = true; continue; } @@ -401,7 +558,7 @@ fn spawn_decommission_index_cancelers( !should_continue_decommission_queue(&pool_meta, idx) }; } - }); + }) } fn decommission_meta_bucket_options() -> MakeBucketOptions { @@ -700,16 +857,6 @@ fn observe_decommission_terminal_reload_result(result: Result<()>, stage: &str) .map(|err| Error::other(format!("decommission terminal pool meta reload failed during {stage}: {err}"))) } -fn resolve_decommission_spawn_failure_result(spawn_err: Error, rollback_err: Option) -> Error { - if let Some(rollback_err) = rollback_err { - Error::other(format!( - "decommission spawn routines failed: {spawn_err}; rollback failed: {rollback_err}" - )) - } else { - spawn_err - } -} - fn decommission_item_size(size: T) -> usize where usize: TryFrom, @@ -2747,7 +2894,10 @@ impl ECStore { let active_workers = { let cancelers = self.decommission_cancelers.read().await; - cancelers.iter().map(Option::is_some).collect::>() + cancelers + .iter() + .map(|canceler| canceler.as_ref().is_some_and(DecommissionCanceler::is_active)) + .collect::>() }; let mut pool_meta = self.pool_meta.write().await; @@ -2783,9 +2933,98 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn decommission_cancel(&self, idx: usize) -> Result<()> { - ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; + self.decommission_cancel_with_owner(idx, None).await + } - let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta) = { + async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + self.decommission_cancel_with_owner(idx, Some(owner)).await + } + + async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) { + let mut cancelers = self.decommission_cancelers.write().await; + take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner); + } + + async fn decommission_terminal_retryable_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> bool { + let _start_guard = self.start_gate.lock().await; + let mut cancelers = self.decommission_cancelers.write().await; + if !decommission_canceler_is_owned_by(cancelers.as_slice(), idx, owner) { + owner.release(); + return false; + } + + let retryable = { + let pool_meta = self.pool_meta.read().await; + pool_meta + .pools + .get(idx) + .and_then(|pool| pool.decommission.as_ref()) + .is_some_and(|info| info.has_decommission_state() && !info.complete && !info.failed && !info.canceled) + }; + if !retryable { + take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner); + } + retryable + } + + async fn retry_decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) { + let mut attempt = 0usize; + loop { + let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else { + return; + }; + if !self.decommission_terminal_retryable_for_operation(idx, owner).await { + return; + } + attempt += 1; + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_save_retry", + terminal = "canceled", + attempt, + error = %err, + "Decommission terminal save will be retried" + ); + tokio::time::sleep(DECOMMISSION_TERMINAL_RETRY_DELAY).await; + } + } + + async fn retry_decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) { + let mut attempt = 0usize; + loop { + let Err(err) = self.decommission_failed_for_operation(idx, owner).await else { + return; + }; + if !self.decommission_terminal_retryable_for_operation(idx, owner).await { + return; + } + attempt += 1; + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_save_retry", + terminal = "failed", + attempt, + error = %err, + "Decommission terminal save will be retried" + ); + tokio::time::sleep(DECOMMISSION_TERMINAL_RETRY_DELAY).await; + } + } + + async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { + ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; + let _start_guard = self.start_gate.lock().await; + + // Lock order: decommission_cancelers before pool_meta. Holding both makes + // owner validation and the terminal transition one atomic operation. + let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta, terminal_canceler) = { + let cancelers = self.decommission_cancelers.read().await; let mut lock = self.pool_meta.write().await; let mut already_canceled = false; let (pool_present, decommission_present, terminal) = if let Some(pool) = lock.pools.get(idx) { @@ -2805,19 +3044,28 @@ impl ECStore { ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?; let previous_pool_meta = lock.clone(); - let changed = lock.decommission_cancel(idx); + let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut lock, idx, owner, |pool_meta| { + pool_meta.decommission_cancel(idx) + }) else { + return Ok(()); + }; + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) + } else { + cancelers.get(idx).and_then(Option::as_ref).cloned() + }; + if let Some(canceler) = terminal_canceler.as_ref() { + canceler.cancel(); + } ( changed, should_retry_decommission_cancel_reload(changed, already_canceled), already_canceled, changed.then_some(previous_pool_meta), + terminal_canceler, ) }; - - let canceled_worker = { - let mut cancelers = self.decommission_cancelers.write().await; - take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx) - }; + let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active); if !canceled_worker && !already_canceled { warn!( event = EVENT_DECOMMISSION_STATE, @@ -2838,6 +3086,10 @@ impl ECStore { return Err(err); } + if let Some(canceler) = terminal_canceler.as_ref() { + self.release_decommission_canceler_slot(idx, canceler).await; + } + if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() { let stage = format!("decommission_cancel for pool {idx}"); resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?; @@ -2849,6 +3101,7 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn clear_decommission(&self, idx: usize) -> Result<()> { ensure_decommission_terminal_operation_supported(self.single_pool(), "clear decommission")?; + let _start_guard = self.start_gate.lock().await; let (should_reload_pool_meta, previous_pool_meta) = { let mut pool_meta = self.pool_meta.write().await; @@ -2857,11 +3110,6 @@ impl ECStore { (changed, changed.then_some(previous_pool_meta)) }; - { - let mut cancelers = self.decommission_cancelers.write().await; - take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx); - } - if should_reload_pool_meta && let Err(err) = self.save_current_pool_meta().await { if let Some(previous_pool_meta) = previous_pool_meta { let mut pool_meta = self.pool_meta.write().await; @@ -2870,6 +3118,11 @@ impl ECStore { return Err(err); } + { + let mut cancelers = self.decommission_cancelers.write().await; + take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx); + } + if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() { let stage = format!("clear_decommission for pool {idx}"); resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?; @@ -2936,26 +3189,53 @@ impl ECStore { is_decommission_cancel_requested(rx.is_cancelled(), pool_meta.pools.get(idx)) } + async fn reserve_decommission_routines( + &self, + rx: &CancellationToken, + indices: &[usize], + ) -> Result> { + let indices = dedup_indices(indices); + if indices.is_empty() { + return Ok(Vec::new()); + } + + let _start_guard = self.start_gate.lock().await; + let indices = { + let pool_meta = self.pool_meta.read().await; + first_resumable_decommission_queue_indices(&pool_meta) + .into_iter() + .filter(|idx| indices.contains(idx)) + .collect::>() + }; + if indices.is_empty() { + return Ok(Vec::new()); + } + + let index_cancelers = { + let mut cancelers = self.decommission_cancelers.write().await; + let missing = missing_decommission_worker_prefix(indices.as_slice(), cancelers.as_slice()); + if missing.is_empty() { + return Ok(Vec::new()); + } + let bound = bind_missing_decommission_cancelers(missing.as_slice(), rx, cancelers.as_mut_slice()); + let guards = guard_decommission_cancelers(bound); + ensure_decommission_routines_scheduled(guards.len(), missing.len())?; + guards + }; + Ok(index_cancelers) + } + pub(crate) async fn spawn_decommission_routines( &self, store: Arc, rx: CancellationToken, indices: Vec, ) -> Result<()> { - let indices = dedup_indices(&indices); - if indices.is_empty() { - return Ok(()); + let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?; + if !index_cancelers.is_empty() { + std::mem::drop(spawn_decommission_index_cancelers(store, rx, index_cancelers)); } - let index_cancelers = { - let mut cancelers = self.decommission_cancelers.write().await; - bind_decommission_cancelers(indices.as_slice(), &rx, cancelers.as_mut_slice()) - }; - - ensure_decommission_routines_scheduled(index_cancelers.len(), indices.len())?; - - spawn_decommission_index_cancelers(store, rx, index_cancelers); - Ok(()) } @@ -2970,17 +3250,12 @@ impl ECStore { } let rx = CancellationToken::new(); - let index_cancelers = { - let mut cancelers = self.decommission_cancelers.write().await; - let missing = missing_decommission_worker_prefix(indices.as_slice(), cancelers.as_slice()); - bind_missing_decommission_cancelers(missing.as_slice(), &rx, cancelers.as_mut_slice()) - }; - + let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?; if index_cancelers.is_empty() { return Ok(()); } - spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers); + std::mem::drop(spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers)); Ok(()) } @@ -3002,28 +3277,10 @@ impl ECStore { let store = require_decommission_store(runtime_sources::object_store_handle(), "start decommission")?; let local_indices = local_decommission_queue_prefix(&self.endpoints(), &indices)?; - - self.start_decommission(indices.clone()).await?; - if let Err(err) = self.spawn_decommission_routines(store, rx, local_indices).await { - let mut rollback_err: Option = None; - for idx in indices { - if let Err(cancel_err) = self.decommission_cancel(idx).await { - error!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "rollback_failed", - error = ?cancel_err, - "Decommission rollback failed after spawn error" - ); - if rollback_err.is_none() { - rollback_err = Some(Error::other(format!("decommission rollback failed for idx {idx}: {cancel_err}"))); - } - } - } - return Err(resolve_decommission_spawn_failure_result(err, rollback_err)); - } + let index_cancelers = self + .start_decommission_with_routines(indices, &rx, local_indices.as_slice()) + .await?; + std::mem::drop(spawn_decommission_index_cancelers(store, rx, index_cancelers)); Ok(()) } @@ -3766,24 +4023,24 @@ impl ECStore { Ok(()) } - #[tracing::instrument(skip(self, rx))] - pub async fn do_decommission_in_routine(self: &Arc, rx: CancellationToken, idx: usize) -> Result<()> { - defer!(|| async { - let mut cancelers = self.decommission_cancelers.write().await; - if take_decommission_canceler(cancelers.as_mut_slice(), idx).is_none() { - warn!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "canceler_already_cleared", - "Decommission canceler already cleared" - ); - } - }); + #[tracing::instrument(skip(self, canceler))] + pub async fn do_decommission_in_routine(self: &Arc, canceler: DecommissionCanceler, idx: usize) -> Result<()> { + let rx = canceler.token().clone(); + self.run_decommission_in_routine(rx, idx, &canceler).await + } + async fn run_decommission_in_routine( + self: &Arc, + rx: CancellationToken, + idx: usize, + canceler: &DecommissionCanceler, + ) -> Result<()> { if let Err(err) = self.promote_queued_decommission(idx).await { - resolve_decommission_terminal_mark_after_error_result(self.decommission_failed(idx).await, idx, &err)?; + resolve_decommission_terminal_mark_after_error_result( + self.decommission_failed_for_operation(idx, canceler).await, + idx, + &err, + )?; return Err(err); } if rx.is_cancelled() { @@ -3802,8 +4059,12 @@ impl ECStore { ); return Ok(()); } - if let Err(err) = self.decommission_cancel(idx).await { - resolve_decommission_terminal_mark_after_error_result(self.decommission_failed(idx).await, idx, &err)?; + if let Err(err) = self.decommission_cancel_for_operation(idx, canceler).await { + resolve_decommission_terminal_mark_after_error_result( + self.decommission_failed_for_operation(idx, canceler).await, + idx, + &err, + )?; return Err(err); } return Ok(()); @@ -3862,7 +4123,11 @@ impl ECStore { return Ok(()); } - resolve_decommission_terminal_mark_after_error_result(self.decommission_failed(idx).await, idx, &err)?; + resolve_decommission_terminal_mark_after_error_result( + self.decommission_failed_for_operation(idx, canceler).await, + idx, + &err, + )?; warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -3909,7 +4174,11 @@ impl ECStore { "Decommission completion verification started" ); if let Err(err) = self.check_after_decommission(idx).await { - resolve_decommission_terminal_mark_result(self.decommission_failed(idx).await, "failed", &cmd_line)?; + resolve_decommission_terminal_mark_result( + self.decommission_failed_for_operation(idx, canceler).await, + "failed", + &cmd_line, + )?; return Err(Error::other(format!( "failed to finalize decommission for pool {cmd_line}: post-check failed: {err}" ))); @@ -3924,7 +4193,11 @@ impl ECStore { state = "marking_completed", "Decommission marking completed state" ); - resolve_decommission_terminal_mark_result(self.complete_decommission(idx).await, "completed", &cmd_line)?; + resolve_decommission_terminal_mark_result( + self.complete_decommission_for_operation(idx, canceler).await, + "completed", + &cmd_line, + )?; } DecommissionFinalState::Failed => { warn!( @@ -3936,7 +4209,11 @@ impl ECStore { state = "marking_failed", "Decommission marking failed state" ); - resolve_decommission_terminal_mark_result(self.decommission_failed(idx).await, "failed", &cmd_line)?; + resolve_decommission_terminal_mark_result( + self.decommission_failed_for_operation(idx, canceler).await, + "failed", + &cmd_line, + )?; } } @@ -3954,63 +4231,97 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn decommission_failed(&self, idx: usize) -> Result<()> { - ensure_decommission_terminal_operation_supported(self.single_pool(), "mark decommission failed")?; + self.decommission_failed_with_owner(idx, None).await + } - let (should_reload_pool_meta, previous_pool_meta) = { + async fn decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + self.decommission_failed_with_owner(idx, Some(owner)).await + } + + async fn decommission_failed_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { + self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta()) + .await + } + + async fn decommission_failed_with_owner_and_save( + &self, + idx: usize, + owner: Option<&DecommissionCanceler>, + save_pool_meta: SaveFuture, + ) -> Result<()> + where + SaveFuture: Future>, + { + ensure_decommission_terminal_operation_supported(self.single_pool(), "mark decommission failed")?; + let _start_guard = self.start_gate.lock().await; + + // Lock order: decommission_cancelers before pool_meta. Holding both makes + // owner validation and the terminal transition one atomic operation. + let (should_reload_pool_meta, previous_pool_meta, terminal_canceler) = { + let cancelers = self.decommission_cancelers.read().await; let mut pool_meta = self.pool_meta.write().await; let previous_pool_meta = pool_meta.clone(); - let changed = pool_meta.decommission_failed(idx); - (changed, changed.then_some(previous_pool_meta)) + let Some(changed) = + update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| { + pool_meta.decommission_failed(idx) + }) + else { + return Ok(()); + }; + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) + } else { + cancelers.get(idx).and_then(Option::as_ref).cloned() + }; + (changed, changed.then_some(previous_pool_meta), terminal_canceler) }; - { - let mut cancelers = self.decommission_cancelers.write().await; - take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx); - } - - if should_reload_pool_meta { - if let Err(err) = self.save_current_pool_meta().await { - if let Some(previous_pool_meta) = previous_pool_meta { - let mut pool_meta = self.pool_meta.write().await; - rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); - } - return Err(err); + if should_reload_pool_meta && let Err(err) = save_pool_meta.await { + if let Some(previous_pool_meta) = previous_pool_meta { + let mut pool_meta = self.pool_meta.write().await; + rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); } + return Err(err); + } + if should_reload_pool_meta { { let mut pool_meta = self.pool_meta.write().await; pool_meta.mark_decommission_progress_saved(); } - if let Some(notification_sys) = runtime_sources::notification_sys() { - let stage = format!("decommission_failed for pool {idx}"); - if let Some(err) = observe_decommission_terminal_reload_result( - resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), - stage.as_str(), - ) { - if let Err(record_err) = self - .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) - .await - { - warn!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "terminal_reload_record_failed", - error = %record_err, - original_error = %err, - "Decommission terminal reload failure record failed" - ); - } + } + if let Some(canceler) = terminal_canceler.as_ref() { + self.release_decommission_canceler_slot(idx, canceler).await; + } + if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() { + let stage = format!("decommission_failed for pool {idx}"); + if let Some(err) = observe_decommission_terminal_reload_result( + resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), + stage.as_str(), + ) { + if let Err(record_err) = self + .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) + .await + { warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_POOLS, pool_index = idx, - state = "terminal_reload_failed", - error = %err, - "Decommission terminal state saved but pool meta reload failed" + state = "terminal_reload_record_failed", + error = %record_err, + original_error = %err, + "Decommission terminal reload failure record failed" ); } + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_reload_failed", + error = %err, + "Decommission terminal state saved but pool meta reload failed" + ); } } @@ -4019,63 +4330,84 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn complete_decommission(&self, idx: usize) -> Result<()> { - ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?; + self.complete_decommission_with_owner(idx, None).await + } - let (should_reload_pool_meta, previous_pool_meta) = { + async fn complete_decommission_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + self.complete_decommission_with_owner(idx, Some(owner)).await + } + + async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { + ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?; + let _start_guard = self.start_gate.lock().await; + + // Lock order: decommission_cancelers before pool_meta. Holding both makes + // owner validation and the terminal transition one atomic operation. + let (should_reload_pool_meta, previous_pool_meta, terminal_canceler) = { + let cancelers = self.decommission_cancelers.read().await; let mut pool_meta = self.pool_meta.write().await; let previous_pool_meta = pool_meta.clone(); - let changed = pool_meta.decommission_complete(idx); - (changed, changed.then_some(previous_pool_meta)) + let Some(changed) = + update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| { + pool_meta.decommission_complete(idx) + }) + else { + return Ok(()); + }; + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) + } else { + cancelers.get(idx).and_then(Option::as_ref).cloned() + }; + (changed, changed.then_some(previous_pool_meta), terminal_canceler) }; - { - let mut cancelers = self.decommission_cancelers.write().await; - take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx); - } - - if should_reload_pool_meta { - if let Err(err) = self.save_current_pool_meta().await { - if let Some(previous_pool_meta) = previous_pool_meta { - let mut pool_meta = self.pool_meta.write().await; - rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); - } - return Err(err); + if should_reload_pool_meta && let Err(err) = self.save_current_pool_meta().await { + if let Some(previous_pool_meta) = previous_pool_meta { + let mut pool_meta = self.pool_meta.write().await; + rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); } + return Err(err); + } + if should_reload_pool_meta { { let mut pool_meta = self.pool_meta.write().await; pool_meta.mark_decommission_progress_saved(); } - if let Some(notification_sys) = runtime_sources::notification_sys() { - let stage = format!("complete_decommission for pool {idx}"); - if let Some(err) = observe_decommission_terminal_reload_result( - resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), - stage.as_str(), - ) { - if let Err(record_err) = self - .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) - .await - { - warn!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "terminal_reload_record_failed", - error = %record_err, - original_error = %err, - "Decommission terminal reload failure record failed" - ); - } + } + if let Some(canceler) = terminal_canceler.as_ref() { + self.release_decommission_canceler_slot(idx, canceler).await; + } + if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() { + let stage = format!("complete_decommission for pool {idx}"); + if let Some(err) = observe_decommission_terminal_reload_result( + resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), + stage.as_str(), + ) { + if let Err(record_err) = self + .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) + .await + { warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_POOLS, pool_index = idx, - state = "terminal_reload_failed", - error = %err, - "Decommission terminal state saved but pool meta reload failed" + state = "terminal_reload_record_failed", + error = %record_err, + original_error = %err, + "Decommission terminal reload failure record failed" ); } + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_reload_failed", + error = %err, + "Decommission terminal state saved but pool meta reload failed" + ); } } @@ -4189,6 +4521,23 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn start_decommission(&self, indices: Vec) -> Result<()> { + self.start_decommission_inner(indices, None).await.map(|_| ()) + } + + async fn start_decommission_with_routines( + &self, + indices: Vec, + rx: &CancellationToken, + local_indices: &[usize], + ) -> Result> { + self.start_decommission_inner(indices, Some((rx, local_indices))).await + } + + async fn start_decommission_inner( + &self, + indices: Vec, + reservation: Option<(&CancellationToken, &[usize])>, + ) -> Result> { let indices = dedup_indices(&indices); validate_start_decommission_request(&indices, self.single_pool())?; @@ -4231,11 +4580,19 @@ impl ECStore { self.ensure_decommission_rebalance_idle_after_refresh().await?; let all_space_infos = self.get_decommission_all_pool_space_infos().await?; - { + let index_cancelers = if let Some((rx, local_indices)) = reservation { + // Lock order matches terminal transitions: decommission_cancelers + // before pool_meta while start_gate excludes another start. + let mut cancelers = self.decommission_cancelers.write().await; + let pool_meta = self.pool_meta.read().await; + ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?; + reserve_decommission_start_cancelers(&pool_meta, &indices, local_indices, rx, cancelers.as_mut_slice())? + } else { let pool_meta = self.pool_meta.read().await; ensure_decommission_start_pool_states(&pool_meta, &indices)?; ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?; - } + Vec::new() + }; let mut space_infos = Vec::with_capacity(indices.len()); for (idx, pi) in all_space_infos.iter().copied() { @@ -4311,7 +4668,7 @@ impl ECStore { return Err(Error::other(format!("{err}; decommission start rollback succeeded"))); } - Ok(()) + Ok(index_cancelers) } async fn get_buckets_to_decommission(&self) -> Result> { @@ -5457,10 +5814,13 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { + use super::DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF; + use super::record_decommission_entry_error; + use super::resolve_decommission_listing_error; use super::{ - DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, - DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, - PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers, + DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler, + DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, + PoolStatus, apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, @@ -5470,35 +5830,36 @@ mod pools_tests { 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, get_by_index, - has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, + guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, - pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store, + pool_meta_has_active_decommission, require_decommission_store, reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, - resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error, + resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result, - resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result, - resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result, - resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta, - run_decommission_buckets_bounded, run_decommission_listing_with_retry, should_cleanup_decommission_source_entry, - should_continue_decommission_queue, should_count_decommission_version_complete, + resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result, + resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result, + rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry, + 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, - split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler, - track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request, - wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context, + 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, }; use crate::data_movement; use crate::disk::endpoint::Endpoint; use crate::error::{Error, StorageError}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; + use crate::runtime::instance::InstanceContext; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; - use rustfs_filemeta::{ - FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo, - }; + use crate::store::ECStore; + use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; + use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams}; use rustfs_rio::Index; use std::sync::{ Arc, @@ -5513,6 +5874,24 @@ mod pools_tests { Arc::new(|_| Box::pin(async {})) } + fn decommission_worker_test_store(pool_meta: PoolMeta, cancelers: Vec>) -> Arc { + let ctx = Arc::new(InstanceContext::new()); + let endpoint_pools = EndpointServerPools::default(); + Arc::new(ECStore { + id: uuid::Uuid::new_v4(), + disk_map: std::collections::HashMap::new(), + pools: Vec::new(), + peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, ctx.clone()), + pool_meta: tokio::sync::RwLock::new(pool_meta), + rebalance_meta: tokio::sync::RwLock::new(None), + decommission_cancelers: tokio::sync::RwLock::new(cancelers), + start_gate: tokio::sync::Mutex::new(()), + pool_meta_save_gate: tokio::sync::Mutex::new(()), + ctx, + bucket_fence_registry: Arc::default(), + }) + } + fn decommission_test_pool_endpoint(idx: usize, is_local: bool) -> PoolEndpoints { let port = 9000usize + idx; let mut endpoint = @@ -6383,20 +6762,6 @@ mod pools_tests { assert!(message.contains(Error::SlowDown.to_string().as_str())); } - #[test] - fn test_resolve_decommission_spawn_failure_result_keeps_primary_without_rollback_error() { - let err = resolve_decommission_spawn_failure_result(Error::SlowDown, None); - assert!(matches!(err, Error::SlowDown)); - } - - #[test] - fn test_resolve_decommission_spawn_failure_result_wraps_rollback_error() { - let err = resolve_decommission_spawn_failure_result(Error::SlowDown, Some(Error::OperationCanceled)); - let message = err.to_string(); - assert!(message.contains("decommission spawn routines failed")); - assert!(message.contains("rollback failed")); - } - #[test] fn test_decommission_item_size_converts_positive_values() { assert_eq!(decommission_item_size(42_i64), 42); @@ -8127,7 +8492,7 @@ mod pools_tests { #[test] fn test_bind_decommission_cancelers_replaces_existing_slot() { let parent = CancellationToken::new(); - let existing = CancellationToken::new(); + let existing = DecommissionCanceler::new(CancellationToken::new()); let mut cancelers = vec![Some(existing.clone())]; let bound = bind_decommission_cancelers(&[0], &parent, cancelers.as_mut_slice()); @@ -8144,7 +8509,7 @@ mod pools_tests { #[test] fn test_bind_missing_decommission_cancelers_stops_at_existing_slot() { let parent = CancellationToken::new(); - let existing = CancellationToken::new(); + let existing = DecommissionCanceler::new(CancellationToken::new()); let mut cancelers = vec![None, Some(existing.clone()), None]; let bound = bind_missing_decommission_cancelers(&[0, 1, 2], &parent, cancelers.as_mut_slice()); @@ -8157,6 +8522,38 @@ mod pools_tests { assert!(!existing.is_cancelled()); } + #[test] + fn test_serialized_decommission_double_start_preserves_first_operation() { + let mut pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status(0, None), decommission_test_pool_status(1, None)], + ..Default::default() + }; + let first_parent = CancellationToken::new(); + let second_parent = CancellationToken::new(); + let mut cancelers = vec![None, None]; + + let first = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &first_parent, cancelers.as_mut_slice()) + .expect("first start should reserve its worker"); + pool_meta + .decommission( + 0, + PoolSpaceInfo { + total: 100, + free: 40, + used: 60, + }, + ) + .expect("first start should install active metadata"); + + let second = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &second_parent, cancelers.as_mut_slice()); + + assert!(matches!(second, Err(Error::DecommissionAlreadyRunning))); + let current = cancelers[0].as_ref().expect("first operation should retain the slot"); + assert!(current.owns_same_operation(first[0].1.canceler())); + assert!(current.is_active()); + assert!(!first_parent.is_cancelled()); + } + #[test] fn test_local_decommission_queue_prefix_stops_at_remote_leader() { let endpoints = EndpointServerPools::from(vec![ @@ -8207,7 +8604,7 @@ mod pools_tests { #[test] fn test_missing_decommission_worker_prefix_stops_at_active_worker() { - let cancelers = vec![None, Some(CancellationToken::new()), None]; + let cancelers = vec![None, Some(DecommissionCanceler::new(CancellationToken::new())), None]; let missing = missing_decommission_worker_prefix(&[0, 1, 2], cancelers.as_slice()); @@ -8316,8 +8713,8 @@ mod pools_tests { #[test] fn test_take_decommission_canceler_takes_and_clears_slot() { - let token = CancellationToken::new(); - let mut cancelers = vec![Some(token)]; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let mut cancelers = vec![Some(canceler)]; let taken = take_decommission_canceler(cancelers.as_mut_slice(), 0); assert!(taken.is_some()); @@ -8326,13 +8723,13 @@ mod pools_tests { #[test] fn test_take_decommission_canceler_returns_none_for_missing_slot() { - let mut cancelers: Vec> = Vec::new(); + let mut cancelers: Vec> = Vec::new(); assert!(take_decommission_canceler(cancelers.as_mut_slice(), 0).is_none()); } #[test] fn test_has_active_decommission_canceler_true_when_any_slot_present() { - let cancelers = vec![None, Some(CancellationToken::new())]; + let cancelers = vec![None, Some(DecommissionCanceler::new(CancellationToken::new()))]; assert!(has_active_decommission_canceler(cancelers.as_slice())); } @@ -8344,11 +8741,12 @@ mod pools_tests { #[test] fn test_cancel_decommission_canceler_cancels_when_present() { - let token = CancellationToken::new(); - let canceled = cancel_decommission_canceler(Some(token.clone())); + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let canceled = cancel_decommission_canceler(Some(canceler.clone())); assert!(canceled); - assert!(token.is_cancelled()); + assert!(canceler.is_cancelled()); + assert!(!canceler.is_active()); } #[test] @@ -8358,12 +8756,13 @@ mod pools_tests { #[test] fn test_take_and_cancel_decommission_canceler_clears_slot() { - let token = CancellationToken::new(); - let mut cancelers = vec![Some(token.clone())]; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let mut cancelers = vec![Some(canceler.clone())]; assert!(take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), 0)); assert!(cancelers[0].is_none()); - assert!(token.is_cancelled()); + assert!(canceler.is_cancelled()); + assert!(!canceler.is_active()); } #[test] @@ -8374,6 +8773,195 @@ mod pools_tests { assert!(cancelers[0].is_none()); } + #[test] + fn test_guarded_decommission_future_releases_without_first_poll() { + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let cancelers = vec![Some(canceler.clone())]; + let guards = guard_decommission_cancelers(vec![(0, canceler.clone())]); + let unpolled = async move { + let _guards = guards; + std::future::pending::<()>().await; + }; + + drop(unpolled); + + assert!(canceler.is_cancelled()); + assert!(!has_active_decommission_canceler(cancelers.as_slice())); + } + + #[test] + fn test_partial_decommission_spawn_reservation_releases_bound_slot() { + let parent = CancellationToken::new(); + let mut cancelers = vec![None]; + let bound = bind_decommission_cancelers(&[0, 1], &parent, cancelers.as_mut_slice()); + let guards = guard_decommission_cancelers(bound); + + let result = super::ensure_decommission_routines_scheduled(guards.len(), 2); + drop(guards); + + assert!(result.is_err()); + assert!(!has_active_decommission_canceler(cancelers.as_slice())); + } + + #[tokio::test] + async fn test_decommission_supervisor_observes_worker_abort() { + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let worker = tokio::spawn(async move { + started_tx.send(()).expect("worker start should be observed"); + std::future::pending::<()>().await; + #[allow(unreachable_code)] + Ok(()) + }); + + started_rx.await.expect("worker start should be observed"); + worker.abort(); + let err = await_decommission_worker(3, worker) + .await + .expect_err("supervisor should observe aborted worker"); + + assert!(err.to_string().contains("decommission worker 3 task join error")); + } + + #[tokio::test] + async fn test_decommission_supervisor_observes_worker_panic() { + let worker = tokio::spawn(async move { + panic!("injected decommission worker panic"); + #[allow(unreachable_code)] + Ok(()) + }); + + let err = await_decommission_worker(4, worker) + .await + .expect_err("supervisor should observe panicked worker"); + + assert!(err.to_string().contains("decommission worker 4 task join error")); + } + + #[tokio::test] + async fn test_decommission_worker_metadata_missing_releases_owned_slot() { + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let store = decommission_worker_test_store(PoolMeta::default(), vec![Some(canceler.clone())]); + canceler.cancel(); + + let err = store + .do_decommission_in_routine(canceler.clone(), 0) + .await + .expect_err("missing worker metadata should fail the routine"); + + assert!(err.to_string().contains("target pool was not found")); + assert!(!canceler.is_active()); + assert!(store.decommission_cancelers.read().await[0].is_none()); + } + + #[tokio::test] + async fn test_decommission_supervisor_failure_cancels_queued_successor() { + let first = DecommissionCanceler::new(CancellationToken::new()); + let queued = DecommissionCanceler::new(CancellationToken::new()); + let store = decommission_worker_test_store(PoolMeta::default(), vec![Some(first.clone()), Some(queued.clone())]); + let guards = guard_decommission_cancelers(vec![(0, first.clone()), (1, queued.clone())]); + + spawn_decommission_index_cancelers(store.clone(), CancellationToken::new(), guards) + .await + .expect("decommission supervisor should finish after queued cleanup"); + + assert!(!first.is_active()); + assert!(!queued.is_active()); + assert!(queued.is_cancelled()); + assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none)); + } + + #[tokio::test] + async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() { + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::UNIX_EPOCH), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + + store + .decommission_failed_with_owner_and_save(0, Some(&canceler), async { Err(Error::SlowDown) }) + .await + .expect_err("injected terminal save failure should be returned"); + + { + let cancelers = store.decommission_cancelers.read().await; + let current = cancelers[0].as_ref().expect("failed save must retain the exact owner slot"); + assert!(current.owns_same_operation(&canceler)); + assert!(current.is_active()); + } + { + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("rollback must retain active decommission metadata"); + assert!(info.has_decommission_state()); + assert!(!info.failed); + assert!(!info.complete); + assert!(!info.canceled); + } + assert!(store.decommission_terminal_retryable_for_operation(0, &canceler).await); + + store + .decommission_failed_with_owner_and_save(0, Some(&canceler), async { Ok(()) }) + .await + .expect("terminal retry should commit"); + + let pool_meta = store.pool_meta.read().await; + assert!( + pool_meta.pools[0] + .decommission + .as_ref() + .expect("terminal metadata should remain") + .failed + ); + drop(pool_meta); + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + assert!(canceler.is_cancelled()); + } + + #[test] + fn test_stale_decommission_operation_cannot_cancel_replacement() { + let stale = DecommissionCanceler::new(CancellationToken::new()); + let replacement = DecommissionCanceler::new(CancellationToken::new()); + let cancelers = vec![Some(replacement.clone())]; + let mut pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::UNIX_EPOCH), + ..Default::default() + }), + )], + ..Default::default() + }; + + let changed = update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, 0, Some(&stale), |pool_meta| { + pool_meta.decommission_cancel(0) + }); + + assert!(changed.is_none()); + assert!( + !pool_meta.pools[0] + .decommission + .as_ref() + .expect("replacement metadata should remain") + .canceled + ); + assert!(replacement.is_active()); + assert!(!replacement.is_cancelled()); + assert!(!stale.is_active()); + assert!(stale.is_cancelled()); + } + #[test] fn test_ensure_decommission_routines_scheduled_accepts_positive_bound_count() { assert!(super::ensure_decommission_routines_scheduled(2, 2).is_ok()); diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 4b8a9cb71..757be0c53 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -33,7 +33,7 @@ use crate::bucket::utils::check_put_object_part_args; use crate::bucket::utils::{check_valid_bucket_name, check_valid_bucket_name_strict, is_meta_bucketname}; use crate::cluster::rpc::{RemoteClient, S3PeerSys}; use crate::config::storageclass; -use crate::core::pools::PoolMeta; +use crate::core::pools::{DecommissionCanceler, PoolMeta}; use crate::disk::endpoint::{Endpoint, EndpointType}; use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions}; use crate::error::{Error, Result}; @@ -176,7 +176,7 @@ pub struct ECStore { // pub local_disks: Vec, pub pool_meta: RwLock, pub rebalance_meta: RwLock>, - pub decommission_cancelers: RwLock>>, + pub decommission_cancelers: RwLock>>, /// Serializes rebalance/decommission start transitions. /// /// Lock order: acquire `start_gate` before `pool_meta`, `rebalance_meta`,