diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 04a94cc05..32b0ed2b5 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}; @@ -72,7 +71,7 @@ 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}; @@ -104,6 +103,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 +186,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 +207,113 @@ 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 +324,32 @@ 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() +} + +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 +468,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,15 +523,16 @@ 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::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 { + if let Err(err) = store.decommission_cancel_for_operation(idx, &canceler).await { warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -700,16 +864,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 +2901,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 +2940,33 @@ 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<()> { + let _canceler_guard = DecommissionCancelerGuard::new(owner.clone()); + let result = self.decommission_cancel_with_owner(idx, Some(owner)).await; + self.release_decommission_canceler_slot(idx, owner).await; + result + } + + 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_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, canceled_worker) = { + let mut cancelers = self.decommission_cancelers.write().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 +2986,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 canceler = if let Some(owner) = owner { + take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + } else { + take_decommission_canceler(cancelers.as_mut_slice(), idx) + }; ( changed, should_retry_decommission_cancel_reload(changed, already_canceled), already_canceled, changed.then_some(previous_pool_meta), + cancel_decommission_canceler(canceler), ) }; - - let canceled_worker = { - let mut cancelers = self.decommission_cancelers.write().await; - take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx) - }; if !canceled_worker && !already_canceled { warn!( event = EVENT_DECOMMISSION_STATE, @@ -2936,24 +3126,32 @@ 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 index_cancelers = { + let mut cancelers = self.decommission_cancelers.write().await; + bind_decommission_cancelers(indices.as_slice(), rx, cancelers.as_mut_slice()) + }; + let index_cancelers = guard_decommission_cancelers(index_cancelers); + ensure_decommission_routines_scheduled(index_cancelers.len(), indices.len())?; + 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 = { - 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())?; - + let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?; spawn_decommission_index_cancelers(store, rx, index_cancelers); Ok(()) @@ -2980,6 +3178,7 @@ impl ECStore { return Ok(()); } + let index_cancelers = guard_decommission_cancelers(index_cancelers); spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers); Ok(()) } @@ -3002,28 +3201,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?; + spawn_decommission_index_cancelers(store, rx, index_cancelers); Ok(()) } @@ -3766,24 +3947,32 @@ 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(); + let _canceler_guard = DecommissionCancelerGuard::new(canceler.clone()); + let result = self.run_decommission_in_routine(rx, idx, &canceler).await; + self.release_decommission_canceler_slot(idx, &canceler).await; + result + } + + 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 +3991,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 +4055,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 +4106,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 +4125,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 +4141,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,20 +4163,48 @@ 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 + } + async fn decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + let _canceler_guard = DecommissionCancelerGuard::new(owner.clone()); + let result = self.decommission_failed_with_owner(idx, Some(owner)).await; + self.release_decommission_canceler_slot(idx, owner).await; + result + } + + async fn decommission_failed_with_owner( + &self, + idx: usize, + owner: Option<&DecommissionCanceler>, + ) -> Result<()> { + 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) = { + let mut cancelers = self.decommission_cancelers.write().await; let mut pool_meta = self.pool_meta.write().await; let previous_pool_meta = pool_meta.clone(); - let changed = pool_meta.decommission_failed(idx); + 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 canceler = if let Some(owner) = owner { + take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + } else { + take_decommission_canceler(cancelers.as_mut_slice(), idx) + }; + cancel_decommission_canceler(canceler); (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 { if let Err(err) = self.save_current_pool_meta().await { if let Some(previous_pool_meta) = previous_pool_meta { @@ -4019,20 +4256,48 @@ 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 + } + async fn complete_decommission_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + let _canceler_guard = DecommissionCancelerGuard::new(owner.clone()); + let result = self.complete_decommission_with_owner(idx, Some(owner)).await; + self.release_decommission_canceler_slot(idx, owner).await; + result + } + + 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) = { + let mut cancelers = self.decommission_cancelers.write().await; let mut pool_meta = self.pool_meta.write().await; let previous_pool_meta = pool_meta.clone(); - let changed = pool_meta.decommission_complete(idx); + 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 canceler = if let Some(owner) = owner { + take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + } else { + take_decommission_canceler(cancelers.as_mut_slice(), idx) + }; + cancel_decommission_canceler(canceler); (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 { if let Err(err) = self.save_current_pool_meta().await { if let Some(previous_pool_meta) = previous_pool_meta { @@ -4189,6 +4454,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 +4513,25 @@ 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 +4607,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> { @@ -5458,11 +5754,17 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { use super::{ - DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, + DecommissionCanceler, DecommissionCancelerGuard, DecommissionStartPoolState, DecommissionTerminalState, + ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers, + bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, + classify_decommission_terminal_state, count_decommission_item, count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, + decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, +DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, +DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, @@ -5470,7 +5772,8 @@ 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, @@ -5480,16 +5783,20 @@ mod pools_tests { 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_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, + reserve_decommission_start_cancelers, 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, + update_decommission_for_operation, validate_start_decommission_request, wait_decommission_listing_retry, wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context, + wait_decommission_worker_drain, with_decommission_entry_context, +touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage, +track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request, }; use crate::data_movement; use crate::disk::endpoint::Endpoint; @@ -6383,20 +6690,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 +8420,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 +8437,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 +8450,50 @@ 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 +8544,11 @@ 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 +8657,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 +8667,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 +8685,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 +8700,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 +8717,95 @@ 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_canceler_guard_releases_operation_on_task_abort() { + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let cancelers = vec![Some(canceler.clone())]; + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let worker_canceler = canceler.clone(); + let worker = tokio::spawn(async move { + let _guard = DecommissionCancelerGuard::new(worker_canceler); + started_tx.send(()).expect("worker start should be observed"); + std::future::pending::<()>().await; + }); + + started_rx.await.expect("worker should install its guard"); + worker.abort(); + let join_error = worker.await.expect_err("aborted worker should return a join error"); + + assert!(join_error.is_cancelled()); + assert!(canceler.is_cancelled()); + assert!(!has_active_decommission_canceler(cancelers.as_slice())); + } + + #[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`,