From 5583e8373ca69d1efcdda0e4f637893f6a59908e Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 13:32:23 +0800 Subject: [PATCH] fix(ecstore): supervise decommission worker exits --- crates/ecstore/src/core/pools.rs | 369 ++++++++++++++++++++++--------- 1 file changed, 267 insertions(+), 102 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 32b0ed2b5..8458b5497 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -95,6 +95,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); @@ -333,6 +334,12 @@ fn guard_decommission_cancelers( .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], @@ -524,7 +531,7 @@ fn spawn_decommission_index_cancelers( store: Arc, rx: CancellationToken, index_cancelers: Vec<(usize, DecommissionCancelerGuard)>, -) { +) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { let mut stop_queue = false; @@ -532,21 +539,16 @@ fn spawn_decommission_index_cancelers( let canceler = canceler_guard.canceler().clone(); if stop_queue || rx.is_cancelled() { canceler.cancel(); - if let Err(err) = store.decommission_cancel_for_operation(idx, &canceler).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, @@ -556,6 +558,7 @@ fn spawn_decommission_index_cancelers( error = %err, "Decommission routine failed" ); + store.retry_decommission_failed_for_operation(idx, &canceler).await; stop_queue = true; continue; } @@ -565,7 +568,7 @@ fn spawn_decommission_index_cancelers( !should_continue_decommission_queue(&pool_meta, idx) }; } - }); + }) } fn decommission_meta_bucket_options() -> MakeBucketOptions { @@ -2944,10 +2947,7 @@ impl ECStore { } 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 + self.decommission_cancel_with_owner(idx, Some(owner)).await } async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) { @@ -2955,6 +2955,78 @@ impl ECStore { 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, @@ -2965,8 +3037,8 @@ impl ECStore { // 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 (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) { @@ -2995,19 +3067,25 @@ impl ECStore { ) else { return Ok(()); }; - let canceler = if let Some(owner) = owner { - take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) } else { - take_decommission_canceler(cancelers.as_mut_slice(), idx) + 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), - cancel_decommission_canceler(canceler), + terminal_canceler, ) }; + let canceled_worker = terminal_canceler + .as_ref() + .is_some_and(DecommissionCanceler::is_active); if !canceled_worker && !already_canceled { warn!( event = EVENT_DECOMMISSION_STATE, @@ -3028,6 +3106,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())?; @@ -3039,6 +3121,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; @@ -3047,11 +3130,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; @@ -3060,6 +3138,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())?; @@ -3136,12 +3219,29 @@ impl ECStore { 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; - bind_decommission_cancelers(indices.as_slice(), rx, cancelers.as_mut_slice()) + 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 }; - let index_cancelers = guard_decommission_cancelers(index_cancelers); - ensure_decommission_routines_scheduled(index_cancelers.len(), indices.len())?; Ok(index_cancelers) } @@ -3152,7 +3252,9 @@ impl ECStore { indices: Vec, ) -> Result<()> { let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?; - spawn_decommission_index_cancelers(store, rx, index_cancelers); + if !index_cancelers.is_empty() { + let _ = spawn_decommission_index_cancelers(store, rx, index_cancelers); + } Ok(()) } @@ -3168,18 +3270,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(()); } - let index_cancelers = guard_decommission_cancelers(index_cancelers); - spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers); + let _ = spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers); Ok(()) } @@ -3204,7 +3300,7 @@ impl ECStore { let index_cancelers = self .start_decommission_with_routines(indices, &rx, local_indices.as_slice()) .await?; - spawn_decommission_index_cancelers(store, rx, index_cancelers); + let _ = spawn_decommission_index_cancelers(store, rx, index_cancelers); Ok(()) } @@ -3954,11 +4050,7 @@ impl ECStore { 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 + self.run_decommission_in_routine(rx, idx, &canceler).await } async fn run_decommission_in_routine( @@ -4167,10 +4259,7 @@ impl ECStore { } 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 + self.decommission_failed_with_owner(idx, Some(owner)).await } async fn decommission_failed_with_owner( @@ -4183,8 +4272,8 @@ impl ECStore { // 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 (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 Some(changed) = update_decommission_for_operation( @@ -4196,27 +4285,31 @@ impl ECStore { ) else { return Ok(()); }; - let canceler = if let Some(owner) = owner { - take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) } else { - take_decommission_canceler(cancelers.as_mut_slice(), idx) + cancelers.get(idx).and_then(Option::as_ref).cloned() }; - cancel_decommission_canceler(canceler); - (changed, changed.then_some(previous_pool_meta)) + (changed, changed.then_some(previous_pool_meta), terminal_canceler) }; - 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(canceler) = terminal_canceler.as_ref() { + self.release_decommission_canceler_slot(idx, canceler).await; + } + if should_reload_pool_meta { 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( @@ -4260,10 +4353,7 @@ impl ECStore { } 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 + self.complete_decommission_with_owner(idx, Some(owner)).await } async fn complete_decommission_with_owner( @@ -4276,8 +4366,8 @@ impl ECStore { // 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 (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 Some(changed) = update_decommission_for_operation( @@ -4289,27 +4379,31 @@ impl ECStore { ) else { return Ok(()); }; - let canceler = if let Some(owner) = owner { - take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner) + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) } else { - take_decommission_canceler(cancelers.as_mut_slice(), idx) + cancelers.get(idx).and_then(Option::as_ref).cloned() }; - cancel_decommission_canceler(canceler); - (changed, changed.then_some(previous_pool_meta)) + (changed, changed.then_some(previous_pool_meta), terminal_canceler) }; - 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(canceler) = terminal_canceler.as_ref() { + self.release_decommission_canceler_slot(idx, canceler).await; + } + if should_reload_pool_meta { 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( @@ -5754,12 +5848,11 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { use super::{ - DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, - DecommissionCanceler, DecommissionCancelerGuard, DecommissionStartPoolState, DecommissionTerminalState, + DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, + DecommissionCanceler, 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, + await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, + cancel_decommission_canceler, 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, @@ -5791,7 +5884,9 @@ DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, 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, + spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler, + take_decommission_canceler, + touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage, 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, @@ -5802,10 +5897,10 @@ track_decommission_current_object, track_decommission_current_object_stage, vali 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_rio::Index; use std::sync::{ Arc, @@ -5820,6 +5915,27 @@ track_decommission_current_object, track_decommission_current_object_stage, vali 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 = @@ -8748,24 +8864,73 @@ track_decommission_current_object, track_decommission_current_object_stage, vali } #[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())]; + async fn test_decommission_supervisor_observes_worker_abort() { 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; + #[allow(unreachable_code)] + Ok(()) }); - started_rx.await.expect("worker should install its guard"); + started_rx.await.expect("worker start should be observed"); worker.abort(); - let join_error = worker.await.expect_err("aborted worker should return a join error"); + let err = await_decommission_worker(3, worker) + .await + .expect_err("supervisor should observe aborted worker"); - assert!(join_error.is_cancelled()); - assert!(canceler.is_cancelled()); - assert!(!has_active_decommission_canceler(cancelers.as_slice())); + 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)); } #[test]