diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 77f4d4fea..d0087212c 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -74,7 +74,7 @@ use std::sync::{ atomic::{AtomicUsize, Ordering}, }; use time::{Duration, OffsetDateTime}; -use tokio::sync::{OwnedSemaphorePermit, Semaphore}; +use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc}; use tokio_util::sync::CancellationToken; use tracing::{debug, error, info, warn}; @@ -92,6 +92,10 @@ const DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD: usize = 1000; const DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF: Duration = Duration::seconds(1); const DECOMMISSION_BUCKET_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_BUCKET_CONCURRENCY"; const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4; +const DECOMMISSION_ENTRY_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_ENTRY_CONCURRENCY"; +const DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP: usize = 8; +const DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP: usize = 64; +const DECOMMISSION_ENTRY_WORKERS_PER_SET: usize = 2; const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30; const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3; const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); @@ -199,6 +203,19 @@ fn decommission_bucket_concurrency_limit() -> usize { rustfs_utils::get_env_usize(DECOMMISSION_BUCKET_CONCURRENCY_ENV, default_limit).max(1) } +fn default_decommission_entry_concurrency(cpu_count: usize) -> usize { + cpu_count.clamp(1, DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP) +} + +fn clamp_decommission_entry_concurrency(limit: usize) -> usize { + limit.clamp(1, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP) +} + +fn decommission_entry_concurrency_limit() -> usize { + let default_limit = default_decommission_entry_concurrency(num_cpus::get()); + clamp_decommission_entry_concurrency(rustfs_utils::get_env_usize(DECOMMISSION_ENTRY_CONCURRENCY_ENV, default_limit)) +} + fn is_decommission_meta_bucket(bucket: &DecomBucketInfo) -> bool { bucket.name == RUSTFS_META_BUCKET } @@ -360,6 +377,7 @@ fn spawn_decommission_index_cancelers( store: Arc, rx: CancellationToken, index_cancelers: Vec<(usize, CancellationToken)>, + entry_budget: Arc, ) { tokio::spawn(async move { let mut stop_queue = false; @@ -381,7 +399,7 @@ fn spawn_decommission_index_cancelers( continue; } - if let Err(err) = store.do_decommission_in_routine(canceler, idx).await { + if let Err(err) = store.do_decommission_in_routine(canceler, idx, entry_budget.clone()).await { error!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -612,6 +630,47 @@ fn count_decommission_item(meta: &mut PoolMeta, idx: usize, size: usize, failed: Ok(()) } +fn ensure_decommission_generation(meta: &PoolMeta, idx: usize, generation: OffsetDateTime) -> Result<()> { + let Some(pool) = meta.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(meta.pools.len(), idx)); + }; + let Some(info) = pool.decommission.as_ref() else { + return Err(decommission_metadata_not_initialized_error("check decommission generation")); + }; + + if info.start_time == Some(generation) && !info.queued && is_decommission_active(info.complete, info.failed, info.canceled) { + Ok(()) + } else { + Err(Error::OperationCanceled) + } +} + +async fn run_decommission_side_effect( + rx: &CancellationToken, + operation_gate: &Arc>, + operation: F, +) -> Result +where + F: FnOnce() -> Fut, + Fut: std::future::Future>, +{ + let _operation_guard = tokio::select! { + biased; + _ = rx.cancelled() => return Err(Error::OperationCanceled), + guard = operation_gate.read() => guard, + }; + + if rx.is_cancelled() { + return Err(Error::OperationCanceled); + } + + let result = operation().await; + if rx.is_cancelled() { + return Err(Error::OperationCanceled); + } + result +} + fn track_decommission_current_object_stage( meta: &mut PoolMeta, idx: usize, @@ -893,6 +952,7 @@ async fn wait_decommission_listing_retry(rx: &CancellationToken, delay: std::tim } } +#[cfg(test)] async fn run_decommission_listing_with_retry( rx: CancellationToken, bucket: String, @@ -900,11 +960,31 @@ async fn run_decommission_listing_with_retry( pool_idx: usize, set_idx: usize, max_attempts: usize, - mut list: List, + list: List, ) -> Result<()> where List: FnMut(ListCallback) -> ListFuture, ListFuture: std::future::Future>, +{ + run_decommission_listing_with_retry_and_drain(rx, bucket, cb, pool_idx, set_idx, max_attempts, list, || async { false }).await +} + +#[allow(clippy::too_many_arguments)] +async fn run_decommission_listing_with_retry_and_drain( + rx: CancellationToken, + bucket: String, + cb: ListCallback, + pool_idx: usize, + set_idx: usize, + max_attempts: usize, + mut list: List, + mut drain: Drain, +) -> Result<()> +where + List: FnMut(ListCallback) -> ListFuture, + ListFuture: std::future::Future>, + Drain: FnMut() -> DrainFuture, + DrainFuture: std::future::Future, { let max_attempts = max_attempts.max(1); @@ -936,7 +1016,12 @@ where "Decommission listing started" ); - match list(cb.clone()).await { + let list_result = list(cb.clone()).await; + if drain().await { + return Ok(()); + } + + match list_result { Ok(()) => { debug!( event = EVENT_DECOMMISSION_BUCKET, @@ -1168,6 +1253,7 @@ where Ok(()) } +#[cfg(test)] async fn wait_decommission_worker_drain(workers: &Semaphore, limit: usize) -> Result<()> { let permits = u32::try_from(limit) .map_err(|_| Error::other(format!("decommission worker limit {limit} exceeds semaphore drain capacity")))?; @@ -1841,7 +1927,7 @@ impl PoolMeta { pub fn decommission_failed(&mut self, idx: usize) -> bool { if let Some(stats) = self.pools.get_mut(idx) { if let Some(d) = &stats.decommission { - if !d.failed { + if is_decommission_active(d.complete, d.failed, d.canceled) { stats.last_update = OffsetDateTime::now_utc(); let mut pd = d.clone(); @@ -1889,7 +1975,7 @@ impl PoolMeta { pub fn decommission_complete(&mut self, idx: usize) -> bool { if let Some(stats) = self.pools.get_mut(idx) { if let Some(d) = &stats.decommission { - if !d.complete { + if is_decommission_active(d.complete, d.failed, d.canceled) { stats.last_update = OffsetDateTime::now_utc(); let mut pd = d.clone(); @@ -2542,12 +2628,13 @@ impl ECStore { snapshot.save(self.pools.clone()).await } - async fn save_decommission_progress_checkpoint(&self, idx: usize) -> Result { + async fn save_decommission_progress_checkpoint(&self, idx: usize, generation: OffsetDateTime) -> Result { // Lock order: save gate, then the short pool metadata read/write sections. Peer // reloads are intentionally performed by the caller after both locks are released. let _save_guard = self.pool_meta_save_gate.lock().await; let (snapshot, checkpoint) = { let pool_meta = self.pool_meta.read().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; let Some(checkpoint) = pool_meta.decommission_progress_checkpoint( idx, DECOMMISSION_PROGRESS_SAVE_INTERVAL, @@ -2715,6 +2802,9 @@ impl ECStore { pub async fn decommission_cancel(&self, idx: usize) -> Result<()> { ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; + let _start_guard = self.start_gate.lock().await; + let canceled_worker = self.cancel_decommission_routine_and_wait(idx).await; + let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta) = { let mut lock = self.pool_meta.write().await; let mut already_canceled = false; @@ -2744,10 +2834,6 @@ impl ECStore { ) }; - 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, @@ -2780,6 +2866,23 @@ impl ECStore { 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 pool_meta = self.pool_meta.read().await; + let pool_count = pool_meta.pools.len(); + ensure_valid_decommission_pool_index(pool_count, idx)?; + let Some(pool) = pool_meta.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let (decommission_present, complete, failed, canceled) = pool + .decommission + .as_ref() + .map(|info| (info.has_decommission_state(), info.complete, info.failed, info.canceled)) + .unwrap_or((false, false, false, false)); + ensure_decommission_clear_allowed(true, decommission_present, complete, failed, canceled)?; + } + self.cancel_decommission_routine_and_wait(idx).await; + let (should_reload_pool_meta, previous_pool_meta) = { let mut pool_meta = self.pool_meta.write().await; let previous_pool_meta = pool_meta.clone(); @@ -2787,11 +2890,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; @@ -2866,6 +2964,33 @@ impl ECStore { is_decommission_cancel_requested(rx.is_cancelled(), pool_meta.pools.get(idx)) } + async fn cancel_decommission_routine_and_wait(&self, idx: usize) -> bool { + let canceled_worker = { + let mut cancelers = self.decommission_cancelers.write().await; + take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx) + }; + if canceled_worker { + let operation_gate = self.ctx.decommission_operation_gate(); + let _operation_guard = operation_gate.write().await; + } + canceled_worker + } + + async fn cancel_decommission_routines_and_wait(&self, indices: &[usize]) { + let canceled_worker = { + let mut cancelers = self.decommission_cancelers.write().await; + let mut canceled_worker = false; + for idx in indices { + canceled_worker |= take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), *idx); + } + canceled_worker + }; + if canceled_worker { + let operation_gate = self.ctx.decommission_operation_gate(); + let _operation_guard = operation_gate.write().await; + } + } + pub(crate) async fn spawn_decommission_routines( &self, store: Arc, @@ -2884,7 +3009,12 @@ impl ECStore { ensure_decommission_routines_scheduled(index_cancelers.len(), indices.len())?; - spawn_decommission_index_cancelers(store, rx, index_cancelers); + spawn_decommission_index_cancelers( + store, + rx, + index_cancelers, + Arc::new(Semaphore::new(decommission_entry_concurrency_limit())), + ); Ok(()) } @@ -2910,7 +3040,12 @@ impl ECStore { return Ok(()); } - spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers); + spawn_decommission_index_cancelers( + self.clone(), + rx, + index_cancelers, + Arc::new(Semaphore::new(decommission_entry_concurrency_limit())), + ); Ok(()) } @@ -2958,15 +3093,329 @@ impl ECStore { Ok(()) } + async fn active_decommission_generation(&self, idx: usize) -> Result { + let pool_meta = self.pool_meta.read().await; + let Some(pool) = pool_meta.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(pool_meta.pools.len(), idx)); + }; + let Some(info) = pool.decommission.as_ref() else { + return Err(decommission_metadata_not_initialized_error("load decommission generation")); + }; + let Some(generation) = info.start_time else { + return Err(Error::OperationCanceled); + }; + ensure_decommission_generation(&pool_meta, idx, generation)?; + Ok(generation) + } + + async fn ensure_decommission_generation_current(&self, idx: usize, generation: OffsetDateTime) -> Result<()> { + let pool_meta = self.pool_meta.read().await; + ensure_decommission_generation(&pool_meta, idx, generation) + } + + #[allow(clippy::too_many_arguments)] + async fn decommission_entry_worker( + self: Arc, + rx: CancellationToken, + idx: usize, + set_idx: usize, + generation: OffsetDateTime, + bucket: String, + set: Arc, + lifecycle_config: Option, + object_lock_config: Option, + replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, + expected_bucket_incarnation_id: Option, + entry_budget: Arc, + queue: Arc>>, + entry_error: Arc>>, + ) { + loop { + let queued = tokio::select! { + biased; + _ = rx.cancelled() => return, + item = async { + let mut queue = queue.lock().await; + queue.recv().await + } => item, + }; + let Some(QueuedDecommissionEntry { entry, queue_permit }) = queued else { + return; + }; + let object_name = entry.name.clone(); + + if entry_error.lock().await.is_some() { + drop(queue_permit); + continue; + } + + if let Err(err) = self.ensure_decommission_generation_current(idx, generation).await { + if matches!(err, Error::OperationCanceled) { + rx.cancel(); + } else { + record_decommission_entry_error(&entry_error, &rx, err).await; + } + return; + } + + if let Err(err) = backpressure::wait_for_data_movement_admission(DataMovementOperation::Decommission, idx, &rx).await + { + if matches!(err, Error::OperationCanceled) { + return; + } + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + set_index = set_idx, + bucket = %bucket, + object = %object_name, + state = "entry_admission_failed", + error = %err, + "Decommission entry admission failed" + ); + record_decommission_entry_error(&entry_error, &rx, err).await; + return; + } + + let entry_budget_permit = match tokio::select! { + biased; + _ = rx.cancelled() => return, + permit = entry_budget.clone().acquire_owned() => permit, + } { + Ok(permit) => permit, + Err(err) => { + let err = Error::other(format!("decommission entry budget permit acquire failed: {err}")); + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + set_index = set_idx, + bucket = %bucket, + object = %object_name, + state = "entry_budget_acquire_failed", + error = %err, + "Decommission entry budget permit acquire failed" + ); + record_decommission_entry_error(&entry_error, &rx, err).await; + return; + } + }; + + let result = self + .decommission_entry( + rx.clone(), + idx, + generation, + entry, + bucket.clone(), + set.clone(), + lifecycle_config.clone(), + object_lock_config.clone(), + replication_config.clone(), + expected_bucket_incarnation_id, + ) + .await; + drop(entry_budget_permit); + drop(queue_permit); + + if let Err(err) = result { + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + set_index = set_idx, + bucket = %bucket, + object = %object_name, + state = "entry_failed", + error = %err, + "Decommission entry failed" + ); + record_decommission_entry_error(&entry_error, &rx, err).await; + return; + } + } + } + + #[allow(clippy::too_many_arguments)] + async fn decommission_set( + self: Arc, + rx: CancellationToken, + idx: usize, + set_idx: usize, + generation: OffsetDateTime, + set: Arc, + bi: DecomBucketInfo, + lifecycle_config: Option, + object_lock_config: Option, + replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, + expected_bucket_incarnation_id: Option, + entry_budget: Arc, + entry_error: Arc>>, + ) -> Result<()> { + let worker_count = DECOMMISSION_ENTRY_WORKERS_PER_SET; + let queue_capacity = decommission_entry_queue_capacity(worker_count); + let outstanding_capacity = queue_capacity.saturating_add(worker_count); + let outstanding = Arc::new(Semaphore::new(outstanding_capacity)); + let (tx, rx_queue) = mpsc::channel(queue_capacity); + let queue = Arc::new(tokio::sync::Mutex::new(rx_queue)); + + let mut entry_workers = tokio::task::JoinSet::new(); + for _ in 0..worker_count { + let this = self.clone(); + let rx = rx.clone(); + let bucket = bi.name.clone(); + let set = set.clone(); + let lifecycle_config = lifecycle_config.clone(); + let object_lock_config = object_lock_config.clone(); + let replication_config = replication_config.clone(); + let queue = queue.clone(); + let entry_budget = entry_budget.clone(); + let entry_error = entry_error.clone(); + entry_workers.spawn(async move { + this.decommission_entry_worker( + rx, + idx, + set_idx, + generation, + bucket, + set, + lifecycle_config, + object_lock_config, + replication_config, + expected_bucket_incarnation_id, + entry_budget, + queue, + entry_error, + ) + .await; + }); + } + + let callback: ListCallback = Arc::new({ + let tx = tx.clone(); + let outstanding = outstanding.clone(); + let callback_rx = rx.clone(); + let entry_error = entry_error.clone(); + let bucket = bi.name.clone(); + move |entry: MetaCacheEntry| { + let tx = tx.clone(); + let outstanding = outstanding.clone(); + let callback_rx = callback_rx.clone(); + let entry_error = entry_error.clone(); + let bucket = bucket.clone(); + Box::pin(async move { + if callback_rx.is_cancelled() || entry_error.lock().await.is_some() { + return; + } + + if matches!( + enqueue_decommission_entry(&callback_rx, &outstanding, &tx, entry).await, + DecommissionEntryEnqueueResult::Closed + ) { + let err = Error::other("decommission entry queue closed"); + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + set_index = set_idx, + bucket = %bucket, + state = "entry_queue_closed", + error = %err, + "Decommission entry queue closed" + ); + record_decommission_entry_error(&entry_error, &callback_rx, err).await; + } + }) + } + }); + + let list_set = set.clone(); + let list_rx = rx.clone(); + let list_rx_for_list = list_rx.clone(); + let list_rx_for_drain = list_rx.clone(); + let list_bi = bi.clone(); + let list_outstanding = outstanding.clone(); + let mut listing = tokio::spawn(async move { + run_decommission_listing_with_retry_and_drain( + list_rx.clone(), + list_bi.name.clone(), + callback, + idx, + set_idx, + DECOMMISSION_LISTING_MAX_ATTEMPTS, + move |callback| { + let set = list_set.clone(); + let rx = list_rx_for_list.clone(); + let bucket = list_bi.clone(); + async move { set.list_objects_to_decommission(rx, bucket, callback).await } + }, + move || { + let rx = list_rx_for_drain.clone(); + let outstanding = list_outstanding.clone(); + async move { drain_decommission_entry_queue(&rx, &outstanding, outstanding_capacity).await } + }, + ) + .await + }); + + let mut listing_result = None; + let mut workers_left = worker_count; + let mut sender = Some(tx); + while listing_result.is_none() || workers_left > 0 { + tokio::select! { + biased; + result = &mut listing, if listing_result.is_none() => { + let result = resolve_decommission_listing_worker_result(set_idx, result); + if result.is_err() { + rx.cancel(); + } + listing_result = Some(result); + drop(sender.take()); + } + worker_result = entry_workers.join_next(), if workers_left > 0 => { + workers_left -= 1; + if let Some(Err(err)) = worker_result { + let err = Error::other(format!("decommission entry worker {set_idx} task join error: {err}")); + error!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + set_index = set_idx, + bucket = %bi.name, + state = "entry_worker_join_failed", + error = %err, + "Decommission entry worker task failed" + ); + record_decommission_entry_error(&entry_error, &rx, err).await; + } + } + } + } + + let listing_result = listing_result.unwrap_or_else(|| Err(Error::other("decommission listing task did not complete"))); + if let Some(err) = entry_error.lock().await.clone() { + return Err(err); + } + listing_result + } + async fn track_decommission_entry_progress_stage( &self, idx: usize, + generation: OffsetDateTime, bucket: &str, object: &str, stage: &'static str, ) -> Result<()> { { let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; track_decommission_current_object_stage(&mut pool_meta, idx, bucket, object, stage) .map_err(|err| with_decommission_entry_context(stage, bucket, object, err))?; } @@ -2975,15 +3424,15 @@ impl ECStore { } #[allow(unused_assignments, clippy::too_many_arguments)] - #[tracing::instrument(skip(self, set, _worker_permit, lifecycle_config, object_lock_config, replication_config))] + #[tracing::instrument(skip(self, set, lifecycle_config, object_lock_config, replication_config))] async fn decommission_entry( self: &Arc, rx: CancellationToken, idx: usize, + generation: OffsetDateTime, entry: MetaCacheEntry, bucket: String, set: Arc, - _worker_permit: OwnedSemaphorePermit, lifecycle_config: Option, object_lock_config: Option, replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, @@ -3016,6 +3465,8 @@ impl ECStore { rx.cancel(); } decommission_cancel_signal_result(rx.is_cancelled())?; + self.ensure_decommission_generation_current(idx, generation).await?; + let operation_gate = self.ctx.decommission_operation_gate(); let bucket_incarnation_fence = match expected_bucket_incarnation_id { Some(expected) => Some(self.acquire_bucket_incarnation_fence(&bucket, expected).await?), @@ -3037,15 +3488,18 @@ impl ECStore { } decommission_cancel_signal_result(rx.is_cancelled())?; - if should_skip_lifecycle_for_data_movement( - self.clone(), - &bucket, - version, - lifecycle_config.as_ref(), - object_lock_config.as_ref(), - true, - &LcEventSrc::Decom, - ) + if run_decommission_side_effect(&rx, &operation_gate, || async { + should_skip_lifecycle_for_data_movement( + self.clone(), + &bucket, + version, + lifecycle_config.as_ref(), + object_lock_config.as_ref(), + true, + &LcEventSrc::Decom, + ) + .await + }) .await .map_err(|err| with_decommission_entry_context("lifecycle_expiry", bucket.as_str(), version.name.as_str(), err))? { @@ -3078,13 +3532,15 @@ impl ECStore { let mut failure = false; let mut error = None; if version.deleted { - if let Err(err) = self - .delete_object( + if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { + self.delete_object( bucket.as_str(), &version.name, decommission_delete_marker_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), ) .await + }) + .await { if is_decommission_copy_cleanup_safe_error(&err) { warn!( @@ -3136,6 +3592,7 @@ impl ECStore { { let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; if let Err(err) = count_decommission_item(&mut pool_meta, idx, 0, failure) { return Err(with_decommission_entry_context( "count_decommission_item", @@ -3167,14 +3624,16 @@ impl ECStore { for _i in 0..3 { if version.is_remote() { - if let Err(err) = self - .decommission_tiered_object( + if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { + self.decommission_tiered_object( bucket.as_str(), &version.name, version, &decommission_remote_tiered_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), ) .await + }) + .await { if is_decommission_copy_cleanup_safe_error(&err) { ignore = true; @@ -3238,16 +3697,19 @@ impl ECStore { self.track_decommission_entry_progress_stage( idx, + generation, bucket_name.as_str(), object_name.as_str(), DECOMMISSION_STAGE_MIGRATE_OBJECT, ) .await?; - if let Err(err) = self - .clone() - .decommission_object(idx, bucket, rd, expected_bucket_incarnation_id) - .await + if let Err(err) = run_decommission_side_effect(&rx, &operation_gate, || async { + self.clone() + .decommission_object(idx, bucket, rd, expected_bucket_incarnation_id) + .await + }) + .await { if is_decommission_copy_cleanup_safe_error(&err) { ignore = true; @@ -3305,6 +3767,7 @@ impl ECStore { { let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; if let Err(err) = count_decommission_item(&mut pool_meta, idx, decommission_item_size(version.size), failure) { return Err(with_decommission_entry_context( "count_decommission_item", @@ -3329,9 +3792,11 @@ impl ECStore { return Err(Error::other("decommission bucket incarnation fence was lost before source cleanup")); } decommission_cancel_signal_result(rx.is_cancelled())?; + self.ensure_decommission_generation_current(idx, generation).await?; self.track_decommission_entry_progress_stage( idx, + generation, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_CLEANUP_PREFLIGHT, @@ -3340,34 +3805,38 @@ impl ECStore { self.track_decommission_entry_progress_stage( idx, + generation, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_SOURCE_CLEANUP, ) .await?; - let cleanup_result = data_movement::cleanup_source_entry_if_unchanged( - set.clone(), - bucket.as_str(), - entry.name.as_str(), - &fivs, - &cleanup_preflight_allowed_missing, - data_movement::SourceCleanupBucketFence { - expected_incarnation_id: expected_bucket_incarnation_id, - lifecycle_guard: bucket_incarnation_fence - .as_ref() - .and_then(|guard| guard.namespace_lock_guard()), - }, - "decommission", - ) - .await - .map_err(|err| match err { - data_movement::SourceCleanupError::SourceChanged => Error::other(format!( - "decommission: source cleanup preflight failed for {}/{}: source versions changed after migration started", - bucket, entry.name - )), - data_movement::SourceCleanupError::Storage(err) => err, - }); + let cleanup_result = run_decommission_side_effect(&rx, &operation_gate, || async { + data_movement::cleanup_source_entry_if_unchanged( + set.clone(), + bucket.as_str(), + entry.name.as_str(), + &fivs, + &cleanup_preflight_allowed_missing, + data_movement::SourceCleanupBucketFence { + expected_incarnation_id: expected_bucket_incarnation_id, + lifecycle_guard: bucket_incarnation_fence + .as_ref() + .and_then(|guard| guard.namespace_lock_guard()), + }, + "decommission", + ) + .await + .map_err(|err| match err { + data_movement::SourceCleanupError::SourceChanged => Error::other(format!( + "decommission: source cleanup preflight failed for {}/{}: source versions changed after migration started", + bucket, entry.name + )), + data_movement::SourceCleanupError::Storage(err) => err, + }) + }) + .await; resolve_decommission_entry_cleanup_delete_result(cleanup_result, bucket.as_str(), entry.name.as_str())? } else if decommissioned != fivs.versions.len() || expired > 0 { warn!( @@ -3387,6 +3856,7 @@ impl ECStore { let should_save_progress = { let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; if let Err(err) = track_decommission_current_object(&mut pool_meta, idx, bucket.as_str(), entry.name.as_str()) { return Err(with_decommission_entry_context( @@ -3407,6 +3877,7 @@ impl ECStore { self.track_decommission_entry_progress_stage( idx, + generation, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_ENTRY_FINISHED, @@ -3414,7 +3885,7 @@ impl ECStore { .await?; if should_save_progress { - match self.save_decommission_progress_checkpoint(idx).await { + match self.save_decommission_progress_checkpoint(idx, generation).await { Ok(true) => { if let Some(notification_sys) = runtime_sources::notification_sys() && let Err(err) = resolve_decommission_entry_reload_result( @@ -3465,13 +3936,10 @@ impl ECStore { idx: usize, pool: Arc, bi: DecomBucketInfo, + entry_budget: Arc, ) -> Result<()> { - let worker_limit = pool.disk_set.len() * 2; - if worker_limit == 0 { - return Err(Error::other("decommission worker limit must be greater than zero")); - } - let workers = Arc::new(Semaphore::new(worker_limit)); let entry_error = Arc::new(tokio::sync::Mutex::new(None::)); + let generation = self.active_decommission_generation(idx).await?; let mut listing_workers = Vec::with_capacity(pool.disk_set.len()); let mut lifecycle_config = None; @@ -3500,12 +3968,6 @@ impl ECStore { } for (set_idx, set) in pool.disk_set.iter().enumerate() { - let listing_permit = workers - .clone() - .acquire_owned() - .await - .map_err(|err| Error::other(format!("decommission listing worker permit acquire failed: {err}")))?; - debug!( event = EVENT_DECOMMISSION_BUCKET, component = LOG_COMPONENT_ECSTORE, @@ -3517,125 +3979,34 @@ impl ECStore { "Decommission listing worker started" ); - let decommission_entry: ListCallback = Arc::new({ - let this = Arc::clone(self); - let bucket = bi.name.clone(); - let workers = workers.clone(); - let set = set.clone(); - let lifecycle_config = lifecycle_config.clone(); - let object_lock_config = object_lock_config.clone(); - let replication_config = replication_config.clone(); - let entry_error = entry_error.clone(); - let callback_rx = rx.clone(); - move |entry: MetaCacheEntry| { - let this = this.clone(); - let bucket = bucket.clone(); - let workers = workers.clone(); - let set = set.clone(); - let lifecycle_config = lifecycle_config.clone(); - let object_lock_config = object_lock_config.clone(); - let replication_config = replication_config.clone(); - let expected_bucket_incarnation_id = expected_bucket_incarnation_id; - let entry_error = entry_error.clone(); - let callback_rx = callback_rx.clone(); - - Box::pin(async move { - if callback_rx.is_cancelled() { - return; - } - if entry_error.lock().await.is_some() { - return; - } - - if let Err(err) = - backpressure::wait_for_data_movement_admission(DataMovementOperation::Decommission, idx, &callback_rx) - .await - { - if matches!(err, Error::OperationCanceled) { - return; - } - error!("decommission_pool: data movement admission failed: {err}"); - let mut first_err = entry_error.lock().await; - if first_err.is_none() { - *first_err = Some(err); - callback_rx.cancel(); - } - return; - } - - if entry_error.lock().await.is_some() { - return; - } - - let worker_permit = match tokio::select! { - _ = callback_rx.cancelled() => return, - permit = workers.clone().acquire_owned() => permit, - } { - Ok(permit) => permit, - Err(err) => { - let err = Error::other(format!("decommission entry worker permit acquire failed: {err}")); - error!("decommission_pool: decommission_entry failed: {err}"); - let mut first_err = entry_error.lock().await; - if first_err.is_none() { - *first_err = Some(err); - callback_rx.cancel(); - } - return; - } - }; - if entry_error.lock().await.is_some() { - return; - } - let entry_rx = callback_rx.clone(); - if let Err(err) = this - .decommission_entry( - entry_rx, - idx, - entry, - bucket, - set, - worker_permit, - lifecycle_config, - object_lock_config, - replication_config, - expected_bucket_incarnation_id, - ) - .await - { - error!("decommission_pool: decommission_entry failed: {err}"); - let mut first_err = entry_error.lock().await; - if first_err.is_none() { - *first_err = Some(err); - callback_rx.cancel(); - } - } - }) - } - }); - let set = set.clone(); + let store = Arc::clone(self); let rx_clone = rx.clone(); - let bi = bi.clone(); - let set_id = set_idx; + let bi_clone = bi.clone(); + let lifecycle_config = lifecycle_config.clone(); + let object_lock_config = object_lock_config.clone(); + let replication_config = replication_config.clone(); + let entry_budget = entry_budget.clone(); + let entry_error = entry_error.clone(); let worker = tokio::spawn(async move { - let _listing_permit = listing_permit; - run_decommission_listing_with_retry( - rx_clone.clone(), - bi.name.clone(), - decommission_entry.clone(), - idx, - set_id, - DECOMMISSION_LISTING_MAX_ATTEMPTS, - |callback| { - let set = set.clone(); - let rx = rx_clone.clone(); - let bucket = bi.clone(); - async move { set.list_objects_to_decommission(rx, bucket, callback).await } - }, - ) - .await + store + .decommission_set( + rx_clone, + idx, + set_idx, + generation, + set, + bi_clone, + lifecycle_config, + object_lock_config, + replication_config, + expected_bucket_incarnation_id, + entry_budget, + entry_error, + ) + .await }); - listing_workers.push((set_id, worker)); + listing_workers.push((set_idx, worker)); } debug!( @@ -3658,8 +4029,6 @@ impl ECStore { } } - wait_decommission_worker_drain(&workers, worker_limit).await?; - if let Some(err) = listing_worker_error { return Err(err); } @@ -3696,7 +4065,12 @@ impl ECStore { } #[tracing::instrument(skip(self, rx))] - pub async fn do_decommission_in_routine(self: &Arc, rx: CancellationToken, idx: usize) -> Result<()> { + pub async fn do_decommission_in_routine( + self: &Arc, + rx: CancellationToken, + idx: usize, + entry_budget: Arc, + ) -> Result<()> { defer!(|| async { let mut cancelers = self.decommission_cancelers.write().await; if take_decommission_canceler(cancelers.as_mut_slice(), idx).is_none() { @@ -3737,7 +4111,8 @@ impl ECStore { } return Ok(()); } - let result = self.decommission_in_background(rx.clone(), idx).await; + let generation = self.active_decommission_generation(idx).await?; + let result = self.decommission_in_background(rx.clone(), idx, entry_budget).await; let (final_state, canceled, cmd_line) = { let pool_meta = self.pool_meta.read().await; @@ -3838,12 +4213,22 @@ impl ECStore { "Decommission completion verification started" ); if let Err(err) = self.check_after_decommission(idx).await { + if self.decommission_cancel_requested(idx, &rx).await { + rx.cancel(); + return Ok(()); + } resolve_decommission_terminal_mark_result(self.decommission_failed(idx).await, "failed", &cmd_line)?; return Err(Error::other(format!( "failed to finalize decommission for pool {cmd_line}: post-check failed: {err}" ))); } + if self.decommission_cancel_requested(idx, &rx).await { + rx.cancel(); + } + decommission_cancel_signal_result(rx.is_cancelled())?; + self.ensure_decommission_generation_current(idx, generation).await?; + info!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -3950,6 +4335,8 @@ impl ECStore { pub async fn complete_decommission(&self, idx: usize) -> Result<()> { ensure_decommission_terminal_operation_supported(self.single_pool(), "complete 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; let previous_pool_meta = pool_meta.clone(); @@ -4017,6 +4404,7 @@ impl ECStore { idx: usize, pool: Arc, bucket: DecomBucketInfo, + entry_budget: Arc, ) -> Result<()> { let is_decommissioned = { let pool_meta = self.pool_meta.read().await; @@ -4042,7 +4430,10 @@ impl ECStore { warn!("decommission: currently on bucket {}", &bucket.name); - if let Err(err) = self.decommission_pool(rx.clone(), idx, pool, bucket.clone()).await { + if let Err(err) = self + .decommission_pool(rx.clone(), idx, pool, bucket.clone(), entry_budget) + .await + { error!("decommission: decommission_pool err {:?}", &err); return Err(err); } else { @@ -4076,40 +4467,53 @@ impl ECStore { pool: Arc, buckets: Vec, limit: usize, + entry_budget: Arc, ) -> Result<()> { let store = Arc::clone(self); run_decommission_buckets_bounded(rx, buckets, limit, move |bucket, rx| { let store = Arc::clone(&store); let pool = pool.clone(); - Box::pin(async move { store.decommission_pending_bucket(rx, idx, pool, bucket).await }) + let entry_budget = entry_budget.clone(); + Box::pin(async move { store.decommission_pending_bucket(rx, idx, pool, bucket, entry_budget).await }) }) .await } #[tracing::instrument(skip(self, rx))] - async fn decommission_in_background(self: &Arc, rx: CancellationToken, idx: usize) -> Result<()> { + async fn decommission_in_background( + self: &Arc, + rx: CancellationToken, + idx: usize, + entry_budget: Arc, + ) -> Result<()> { let pool = get_by_index(self.pools.as_slice(), idx, "load decommission background pool")?.clone(); let pending = { let pool_meta = self.pool_meta.read().await; pool_meta.pending_buckets(idx) }; - let bucket_concurrency = decommission_bucket_concurrency_limit(); if bucket_concurrency <= 1 { for bucket in pending { - self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket) + self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone()) .await?; } return Ok(()); } let (regular_buckets, meta_buckets) = split_decommission_buckets(pending); - self.decommission_buckets_concurrently(rx.clone(), idx, pool.clone(), regular_buckets, bucket_concurrency) - .await?; + self.decommission_buckets_concurrently( + rx.clone(), + idx, + pool.clone(), + regular_buckets, + bucket_concurrency, + entry_budget.clone(), + ) + .await?; for bucket in meta_buckets { - self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket) + self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone()) .await?; } @@ -4166,6 +4570,8 @@ impl ECStore { ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?; } + self.cancel_decommission_routines_and_wait(&indices).await; + let mut space_infos = Vec::with_capacity(indices.len()); for (idx, pi) in all_space_infos.iter().copied() { if indices.contains(&idx) { @@ -4716,6 +5122,14 @@ mod tests { let mut pool_meta = build_pool_meta(); assert!(pool_meta.decommission_cancel(0)); assert_eq!(pool_meta.pools[0].decommission.as_ref().and_then(|info| info.start_time), None); + + let mut pool_meta = build_pool_meta(); + assert!(pool_meta.decommission_cancel(0)); + assert!(!pool_meta.decommission_complete(0)); + + let mut pool_meta = build_pool_meta(); + assert!(pool_meta.decommission_failed(0)); + assert!(!pool_meta.decommission_complete(0)); } #[test] @@ -5099,6 +5513,75 @@ mod tests { pub type ListCallback = Arc BoxFuture<'static, ()> + Send + Sync + 'static>; +const DECOMMISSION_ENTRY_QUEUE_HARD_CAP: usize = 256; + +struct QueuedDecommissionEntry { + entry: MetaCacheEntry, + queue_permit: OwnedSemaphorePermit, +} + +enum DecommissionEntryEnqueueResult { + Enqueued, + Canceled, + Closed, +} + +fn decommission_entry_queue_capacity(worker_limit: usize) -> usize { + worker_limit.saturating_mul(2).clamp(1, DECOMMISSION_ENTRY_QUEUE_HARD_CAP) +} + +async fn enqueue_decommission_entry( + rx: &CancellationToken, + outstanding: &Arc, + tx: &mpsc::Sender, + entry: MetaCacheEntry, +) -> DecommissionEntryEnqueueResult { + let queue_permit = match tokio::select! { + biased; + _ = rx.cancelled() => return DecommissionEntryEnqueueResult::Canceled, + permit = outstanding.clone().acquire_owned() => permit, + } { + Ok(permit) => permit, + Err(_) => return DecommissionEntryEnqueueResult::Closed, + }; + + let queued = QueuedDecommissionEntry { entry, queue_permit }; + tokio::select! { + biased; + _ = rx.cancelled() => DecommissionEntryEnqueueResult::Canceled, + result = tx.send(queued) => { + if result.is_ok() { + DecommissionEntryEnqueueResult::Enqueued + } else { + DecommissionEntryEnqueueResult::Closed + } + } + } +} + +async fn drain_decommission_entry_queue(rx: &CancellationToken, outstanding: &Arc, capacity: usize) -> bool { + let Ok(permits) = u32::try_from(capacity) else { + return true; + }; + + tokio::select! { + _ = rx.cancelled() => true, + result = outstanding.acquire_many(permits) => result.is_err(), + } +} + +async fn record_decommission_entry_error( + entry_error: &Arc>>, + rx: &CancellationToken, + err: Error, +) { + let mut first_err = entry_error.lock().await; + if first_err.is_none() { + *first_err = Some(err); + rx.cancel(); + } +} + impl SetDisks { #[tracing::instrument(skip(self, rx, cb_func))] async fn list_objects_to_decommission( @@ -5343,19 +5826,22 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { use super::{ + DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, - DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, - PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, 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, - 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, - ensure_decommission_start_local_leader, ensure_decommission_start_pool_states, - ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity, - ensure_decommission_terminal_operation_supported, ensure_local_decommission_pool_leaders, - ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, get_by_index, - has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, + DecomBucketInfo, DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, + PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info, + bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, + clamp_decommission_entry_concurrency, classify_decommission_terminal_state, count_decommission_item, + decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size, + decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, + default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry, + ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation, + ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, + ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, + ensure_decommission_start_pool_states, ensure_decommission_start_rebalance_meta_allowed, + ensure_decommission_start_target_capacity, ensure_decommission_terminal_operation_supported, + ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, + get_by_index, 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, require_decommission_store, resolve_decommission_bucket_done_save_result, @@ -5367,13 +5853,14 @@ mod pools_tests { 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, - 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, + run_decommission_buckets_bounded, run_decommission_listing_with_retry, run_decommission_listing_with_retry_and_drain, + run_decommission_side_effect, should_cleanup_decommission_source_entry, should_continue_decommission_queue, + should_count_decommission_version_complete, should_preserve_decommission_canceled_state, + should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload, + should_retry_decommission_listing, should_skip_canceled_decommission_routine, 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, }; use crate::data_movement; use crate::disk::endpoint::Endpoint; @@ -5608,6 +6095,25 @@ mod pools_tests { assert_eq!(default_decommission_bucket_concurrency(8), 4); } + #[test] + fn test_default_decommission_entry_concurrency_is_conservative() { + assert_eq!(default_decommission_entry_concurrency(0), 1); + assert_eq!(default_decommission_entry_concurrency(1), 1); + assert_eq!(default_decommission_entry_concurrency(4), 4); + assert_eq!(default_decommission_entry_concurrency(16), DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP); + } + + #[test] + fn test_decommission_entry_concurrency_clamps_operator_configuration() { + assert_eq!(clamp_decommission_entry_concurrency(0), 1); + assert_eq!(clamp_decommission_entry_concurrency(1), 1); + assert_eq!( + clamp_decommission_entry_concurrency(DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP), + DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP + ); + assert_eq!(clamp_decommission_entry_concurrency(usize::MAX), DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP); + } + #[test] fn test_split_decommission_buckets_keeps_meta_buckets_last() { let (regular, meta) = split_decommission_buckets(vec![ @@ -5779,6 +6285,167 @@ mod pools_tests { assert!(result.is_ok()); } + #[test] + fn test_decommission_entry_queue_capacity_is_bounded() { + assert_eq!(decommission_entry_queue_capacity(0), 1); + assert_eq!(decommission_entry_queue_capacity(1), 2); + assert_eq!( + decommission_entry_queue_capacity(DECOMMISSION_ENTRY_QUEUE_HARD_CAP), + DECOMMISSION_ENTRY_QUEUE_HARD_CAP + ); + assert_eq!(decommission_entry_queue_capacity(usize::MAX), DECOMMISSION_ENTRY_QUEUE_HARD_CAP); + } + + #[tokio::test] + async fn test_drain_decommission_entry_queue_waits_for_all_outstanding_entries() { + let outstanding = Arc::new(Semaphore::new(1)); + let held = outstanding + .clone() + .acquire_owned() + .await + .expect("test outstanding permit should acquire"); + let rx = CancellationToken::new(); + let drain = tokio::spawn({ + let outstanding = outstanding.clone(); + let rx = rx.clone(); + async move { drain_decommission_entry_queue(&rx, &outstanding, 1).await } + }); + + tokio::task::yield_now().await; + assert!(!drain.is_finished(), "queue drain must wait for active entry work"); + drop(held); + + let drained = tokio::time::timeout(StdDuration::from_secs(1), drain) + .await + .expect("queue drain should finish after entry completion") + .expect("queue drain task should not panic"); + assert!(!drained); + } + + #[tokio::test] + async fn test_enqueue_decommission_entry_observes_cancellation_when_queue_is_full() { + let outstanding = Arc::new(Semaphore::new(2)); + let (tx, mut queue) = tokio::sync::mpsc::channel(1); + let held = outstanding + .clone() + .acquire_owned() + .await + .expect("first queue permit should acquire"); + tx.send(QueuedDecommissionEntry { + entry: MetaCacheEntry::default(), + queue_permit: held, + }) + .await + .expect("first entry should fill the queue"); + + let rx = CancellationToken::new(); + let enqueue = tokio::spawn({ + let rx = rx.clone(); + let outstanding = outstanding.clone(); + let tx = tx.clone(); + async move { enqueue_decommission_entry(&rx, &outstanding, &tx, MetaCacheEntry::default()).await } + }); + + tokio::task::yield_now().await; + rx.cancel(); + let result = tokio::time::timeout(StdDuration::from_secs(1), enqueue) + .await + .expect("full queue enqueue should observe cancellation") + .expect("enqueue task should not panic"); + assert!(matches!(result, DecommissionEntryEnqueueResult::Canceled)); + drop(queue.recv().await); + } + + #[tokio::test] + async fn test_decommission_side_effect_gate_quiesces_before_transition() { + let operation_gate = Arc::new(tokio::sync::RwLock::new(())); + let rx = CancellationToken::new(); + let started = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let operation = tokio::spawn({ + let operation_gate = operation_gate.clone(); + let rx = rx.clone(); + let started = started.clone(); + let release = release.clone(); + async move { + run_decommission_side_effect(&rx, &operation_gate, || async { + started.notify_one(); + release.notified().await; + Ok::<_, Error>(()) + }) + .await + } + }); + + started.notified().await; + rx.cancel(); + let transition = tokio::spawn({ + let operation_gate = operation_gate.clone(); + async move { + let _guard = operation_gate.write().await; + } + }); + tokio::task::yield_now().await; + assert!(!transition.is_finished(), "transition must wait for the in-flight side effect"); + + release.notify_one(); + let operation_result = operation.await.expect("operation task should not panic"); + assert!(matches!(operation_result, Err(Error::OperationCanceled))); + transition.await.expect("transition task should not panic"); + + let called = Arc::new(AtomicBool::new(false)); + let result = run_decommission_side_effect(&rx, &operation_gate, { + let called = called.clone(); + move || async move { + called.store(true, Ordering::SeqCst); + Ok::<_, Error>(()) + } + }) + .await; + assert!(matches!(result, Err(Error::OperationCanceled))); + assert!(!called.load(Ordering::SeqCst)); + } + + #[tokio::test(start_paused = true)] + async fn test_run_decommission_listing_with_retry_drains_before_each_retry() { + let attempts = Arc::new(AtomicUsize::new(0)); + let drains = Arc::new(AtomicUsize::new(0)); + let err = run_decommission_listing_with_retry_and_drain( + CancellationToken::new(), + "bucket-a".to_string(), + noop_decommission_list_callback(), + 1, + 2, + 2, + { + let attempts = attempts.clone(); + move |_| { + let attempts = attempts.clone(); + async move { + attempts.fetch_add(1, Ordering::SeqCst); + Err(Error::SlowDown) + } + } + }, + { + let drains = drains.clone(); + move || { + let drains = drains.clone(); + async move { + drains.fetch_add(1, Ordering::SeqCst); + false + } + } + }, + ) + .await + .expect_err("permanent listing failure must be returned"); + + assert!(err.to_string().contains("attempt 2/2")); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + assert_eq!(drains.load(Ordering::SeqCst), 2); + } + #[test] fn test_get_by_index_returns_value_when_in_range() { let values = vec!["a", "b", "c"]; @@ -6961,6 +7628,33 @@ mod pools_tests { assert!(!is_decommission_active(false, false, true)); } + #[test] + fn test_ensure_decommission_generation_rejects_stale_or_queued_workers() { + let generation = OffsetDateTime::UNIX_EPOCH; + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: generation, + decommission: Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + }], + ..Default::default() + }; + + assert!(ensure_decommission_generation(&meta, 0, generation).is_ok()); + assert!(ensure_decommission_generation(&meta, 0, generation + Duration::seconds(1)).is_err()); + + meta.pools[0] + .decommission + .as_mut() + .expect("decommission metadata should exist") + .queued = true; + assert!(ensure_decommission_generation(&meta, 0, generation).is_err()); + } + #[test] fn test_pool_meta_has_active_decommission_counts_running_and_queued_states() { let active_meta = PoolMeta { diff --git a/crates/ecstore/src/runtime/instance.rs b/crates/ecstore/src/runtime/instance.rs index 71a1898cf..9453ed02b 100644 --- a/crates/ecstore/src/runtime/instance.rs +++ b/crates/ecstore/src/runtime/instance.rs @@ -160,6 +160,10 @@ pub struct InstanceContext { /// workers (scanner/heal/tier/lifecycle) without touching another instance. /// Replaces the process-global cancel-token static. background_cancel_token: OnceLock, + /// Serializes decommission data-movement operations with cancellation and + /// a subsequent restart. Readers are held across one object side effect; + /// the transition path takes the writer after cancelling the routine. + decommission_operation_gate: Arc>, /// Resolves object-encryption material at the application boundary. object_encryption_resolver: OnceLock>, tier_delete_journal_recovery_stores: std::sync::Mutex>, @@ -200,6 +204,7 @@ impl InstanceContext { local_disk_set_drives: Arc::new(RwLock::new(Vec::new())), bucket_metadata_sys: std::sync::Mutex::new(None), background_cancel_token: OnceLock::new(), + decommission_operation_gate: Arc::new(RwLock::new(())), object_encryption_resolver: OnceLock::new(), tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()), transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()), @@ -218,6 +223,10 @@ impl InstanceContext { self.lock_manager.clone() } + pub(crate) fn decommission_operation_gate(&self) -> Arc> { + Arc::clone(&self.decommission_operation_gate) + } + /// Install the application-owned object-encryption resolver once. pub fn set_object_encryption_resolver( &self,