diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 4ca460c20..97edf02e5 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -1120,6 +1120,48 @@ fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: *pool_meta = previous_pool_meta; } +#[derive(Debug)] +struct DecommissionCancelCommit { + previous_start_time: Option, + previous_queued: bool, + previous_last_update: OffsetDateTime, + canceled_pool: PoolStatus, +} + +fn commit_decommission_cancel(pool_meta: &mut PoolMeta, idx: usize, commit: DecommissionCancelCommit) -> Result<()> { + let pool_count = pool_meta.pools.len(); + let Some(current) = pool_meta.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + // A peer reload can install the saved cancel while runtime-only fields are reconstructed. + let cancel_already_published = commit + .canceled_pool + .decommission + .as_ref() + .is_some_and(|info| info.canceled && !info.complete && !info.failed && info.start_time.is_none()) + && PersistedPoolStatus::from(current) == PersistedPoolStatus::from(&commit.canceled_pool); + if cancel_already_published { + return Ok(()); + } + + let matches_generation = current.id == commit.canceled_pool.id + && current.cmd_line == commit.canceled_pool.cmd_line + && current.decommission.as_ref().is_some_and(|info| { + info.start_time == commit.previous_start_time + && info.queued == commit.previous_queued + && is_decommission_active(info.complete, info.failed, info.canceled) + && (commit.previous_start_time.is_some() || current.last_update == commit.previous_last_update) + }); + if !matches_generation { + return Err(Error::other(format!( + "failed to publish decommission cancel for pool {idx}: operation generation changed" + ))); + } + + pool_meta.pools[idx] = commit.canceled_pool; + Ok(()) +} + fn rollback_start_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) { rollback_decommission_pool_meta(pool_meta, previous_pool_meta); } @@ -1570,7 +1612,7 @@ struct PersistedPoolMeta { pub pools: Vec, } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct PersistedPoolStatus { #[serde(rename = "id")] @@ -1583,7 +1625,7 @@ struct PersistedPoolStatus { pub decommission: Option, } -#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct PersistedPoolDecommissionInfo { #[serde(rename = "startTime", with = "time::serde::rfc3339::option")] @@ -3004,14 +3046,156 @@ impl ECStore { } #[tracing::instrument(skip(self))] - pub async fn decommission_cancel(&self, idx: usize) -> Result<()> { + pub async fn decommission_cancel(self: &Arc, idx: usize) -> Result<()> { self.decommission_cancel_with_owner(idx, None).await } - async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { + async fn decommission_cancel_for_operation(self: &Arc, idx: usize, owner: &DecommissionCanceler) -> Result<()> { self.decommission_cancel_with_owner(idx, Some(owner)).await } + async fn decommission_cancel_with_owner_and_save( + self: &Arc, + idx: usize, + owner: Option<&DecommissionCanceler>, + save_pool_meta: Save, + ) -> Result<()> + where + Save: FnOnce(PoolMeta) -> SaveFuture + Send + 'static, + SaveFuture: Future> + Send + 'static, + { + let store = self.clone(); + let owner = owner.cloned(); + // Dropping the RPC waiter detaches this task; the transaction retains + // the store and exact owner until persistence is resolved. + tokio::spawn(async move { store.decommission_cancel_transaction(idx, owner, save_pool_meta).await }) + .await + .map_err(|err| Error::other(format!("decommission cancel transaction task join error: {err}")))? + } + + async fn decommission_cancel_transaction( + &self, + idx: usize, + owner: Option, + save_pool_meta: Save, + ) -> Result<()> + where + Save: FnOnce(PoolMeta) -> SaveFuture, + SaveFuture: Future>, + { + let owner = owner.as_ref(); + ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; + let _start_guard = self.start_gate.lock().await; + let operation_gate = self.ctx.decommission_operation_gate(); + let operation_guard = operation_gate.write().await; + let save_guard = self.pool_meta_save_gate.lock().await; + + // Lock order: start gate, operation gate, save gate, + // decommission_cancelers, then pool_meta. The state guards stay held + // across persistence so publication and owner termination are + // synchronous once the save reports success. + let mut cancelers = self.decommission_cancelers.write().await; + let mut pool_meta = self.pool_meta.write().await; + let (pending, should_reload_pool_meta, already_canceled, terminal_canceler) = { + let mut already_canceled = false; + let (pool_present, decommission_present, terminal) = if let Some(pool) = pool_meta.pools.get(idx) { + if let Some(info) = pool.decommission.as_ref() { + already_canceled = info.canceled; + ( + true, + info.has_decommission_state(), + should_reject_decommission_cancel_as_terminal(info.complete, info.failed), + ) + } else { + (true, false, false) + } + } else { + (false, false, false) + }; + + ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?; + let previous_pool = pool_meta + .pools + .get(idx) + .cloned() + .ok_or_else(|| invalid_decommission_pool_index_error(pool_meta.pools.len(), idx))?; + let previous_decommission = previous_pool + .decommission + .as_ref() + .ok_or_else(|| decommission_metadata_not_initialized_error("cancel decommission"))?; + let mut snapshot = pool_meta.clone(); + let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut snapshot, idx, owner, |pool_meta| { + pool_meta.decommission_cancel(idx) + }) else { + return Ok(()); + }; + let pending = if changed { + let canceled_pool = snapshot + .pools + .get(idx) + .cloned() + .ok_or_else(|| invalid_decommission_pool_index_error(pool_meta.pools.len(), idx))?; + Some(( + snapshot, + DecommissionCancelCommit { + previous_start_time: previous_decommission.start_time, + previous_queued: previous_decommission.queued, + previous_last_update: previous_pool.last_update, + canceled_pool, + }, + )) + } else { + None + }; + let terminal_canceler = if let Some(owner) = owner { + Some(owner.clone()) + } else { + cancelers.get(idx).and_then(Option::as_ref).cloned() + }; + ( + pending, + should_retry_decommission_cancel_reload(changed, already_canceled), + already_canceled, + terminal_canceler, + ) + }; + let active_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active); + if !active_worker && !already_canceled { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "cancel_skipped", + reason = "no_active_canceler", + "Decommission cancel skipped" + ); + } + + let commit_result = if let Some((snapshot, commit)) = pending { + save_pool_meta(snapshot).await?; + commit_decommission_cancel(&mut pool_meta, idx, commit) + } else { + Ok(()) + }; + + if let Some(canceler) = terminal_canceler.as_ref() { + take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, canceler); + } + commit_result?; + drop(pool_meta); + drop(cancelers); + drop(save_guard); + drop(operation_guard); + + 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())?; + } + + Ok(()) + } + 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); @@ -3039,7 +3223,11 @@ impl ECStore { retryable } - async fn retry_decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) { + async fn retry_decommission_cancel_for_operation(self: &Arc, idx: usize, owner: &DecommissionCanceler) { + if !self.decommission_terminal_retryable_for_operation(idx, owner).await { + return; + } + let mut attempt = 0usize; loop { let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else { @@ -3089,87 +3277,10 @@ impl ECStore { } } - async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { - ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; - let _start_guard = self.start_gate.lock().await; - - // Lock order: decommission_cancelers before pool_meta. Holding both makes - // owner validation and the terminal transition one atomic operation. - let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta, terminal_canceler) = { - let cancelers = self.decommission_cancelers.read().await; - let mut lock = self.pool_meta.write().await; - let mut already_canceled = false; - let (pool_present, decommission_present, terminal) = if let Some(pool) = lock.pools.get(idx) { - if let Some(info) = pool.decommission.as_ref() { - already_canceled = info.canceled; - ( - true, - info.has_decommission_state(), - should_reject_decommission_cancel_as_terminal(info.complete, info.failed), - ) - } else { - (true, false, false) - } - } else { - (false, false, false) - }; - - ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?; - let previous_pool_meta = lock.clone(); - let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut lock, idx, owner, |pool_meta| { - pool_meta.decommission_cancel(idx) - }) else { - return Ok(()); - }; - let terminal_canceler = if let Some(owner) = owner { - Some(owner.clone()) - } else { - cancelers.get(idx).and_then(Option::as_ref).cloned() - }; - if let Some(canceler) = terminal_canceler.as_ref() { - canceler.cancel(); - } - ( - changed, - should_retry_decommission_cancel_reload(changed, already_canceled), - already_canceled, - changed.then_some(previous_pool_meta), - terminal_canceler, - ) - }; - let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active); - if !canceled_worker && !already_canceled { - warn!( - event = EVENT_DECOMMISSION_STATE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - state = "cancel_skipped", - reason = "no_active_canceler", - "Decommission cancel skipped" - ); - } - - self.wait_for_decommission_side_effects().await; - - if should_save_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 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())?; - } - - Ok(()) + async fn decommission_cancel_with_owner(self: &Arc, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { + let pools = self.pools.clone(); + self.decommission_cancel_with_owner_and_save(idx, owner, move |snapshot| async move { snapshot.save(pools).await }) + .await } #[tracing::instrument(skip(self))] @@ -9645,6 +9756,525 @@ mod pools_tests { assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none)); } + #[tokio::test] + async fn test_decommission_cancel_save_failure_preserves_generation_and_token_until_retry() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + + let err = store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), |_| async { Err(Error::Timeout) }) + .await + .expect_err("injected pool metadata timeout should fail cancel"); + assert!(matches!(err, Error::Timeout)); + + { + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("failed cancel must retain decommission metadata"); + assert_eq!(info.start_time, Some(generation)); + assert!(!info.canceled); + assert!(!info.complete); + assert!(!info.failed); + assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok()); + } + { + let cancelers = store.decommission_cancelers.read().await; + let current = cancelers[0].as_ref().expect("failed cancel must retain the worker owner"); + assert!(current.owns_same_operation(&canceler)); + assert!(current.is_active()); + } + assert!(!canceler.is_cancelled()); + assert!(store.decommission_terminal_retryable_for_operation(0, &canceler).await); + + let persisted = Arc::new(std::sync::Mutex::new(None)); + let persisted_for_save = persisted.clone(); + store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + let data = snapshot.encode_config_data()?; + *persisted_for_save + .lock() + .expect("persisted snapshot lock should not be poisoned") = Some(data); + Ok(()) + }) + .await + .expect("cancel retry should commit"); + + { + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("committed cancel metadata should remain present"); + assert!(info.canceled); + assert!(!info.complete); + assert!(!info.failed); + assert!(info.start_time.is_none()); + } + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + assert!(canceler.is_cancelled()); + + let repeated_save_called = Arc::new(AtomicBool::new(false)); + let repeated_save_called_by_closure = repeated_save_called.clone(); + store + .decommission_cancel_with_owner_and_save(0, None, move |_| async move { + repeated_save_called_by_closure.store(true, Ordering::SeqCst); + Ok(()) + }) + .await + .expect("repeated cancel should be idempotent"); + assert!(!repeated_save_called.load(Ordering::SeqCst)); + + let persisted = persisted + .lock() + .expect("persisted snapshot lock should not be poisoned") + .take() + .expect("successful retry should capture persisted bytes"); + let mut restarted = PoolMeta::default(); + restarted + .load_from_config_data(persisted) + .expect("a restarted process should decode the committed cancel"); + assert!( + restarted.pools[0] + .decommission + .as_ref() + .is_some_and(|info| info.canceled && info.start_time.is_none()) + ); + assert!(first_resumable_decommission_queue_indices(&restarted).is_empty()); + } + + #[tokio::test] + async fn test_decommission_cancel_serializes_reload_until_local_commit() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + stage: "migrate_object".to_string(), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + let (persisted_tx, persisted_rx) = tokio::sync::oneshot::channel(); + let save_release = Arc::new(tokio::sync::Notify::new()); + + let cancel = tokio::spawn({ + let store = store.clone(); + let canceler = canceler.clone(); + let save_release = save_release.clone(); + async move { + store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + persisted_tx + .send(snapshot.encode_config_data()?) + .map_err(|_| Error::other("failed to expose saved cancel snapshot"))?; + save_release.notified().await; + Ok(()) + }) + .await + } + }); + let persisted = persisted_rx.await.expect("save should expose the canceled snapshot"); + let reload_started = Arc::new(tokio::sync::Notify::new()); + let reload = tokio::spawn({ + let store = store.clone(); + let reload_started = reload_started.clone(); + async move { + let mut reloaded = PoolMeta::default(); + reloaded.load_from_config_data(persisted)?; + reload_started.notify_one(); + *store.pool_meta.write().await = reloaded; + Ok::<(), Error>(()) + } + }); + reload_started.notified().await; + assert!(!reload.is_finished(), "peer reload must wait for cancel publication"); + + save_release.notify_one(); + cancel + .await + .expect("cancel task should not panic") + .expect("cancel should commit before releasing the reload"); + reload + .await + .expect("reload task should not panic") + .expect("reload should install the saved cancel after publication"); + + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("reloaded cancel metadata should remain present"); + assert!(info.canceled); + assert!(!info.complete); + assert!(!info.failed); + assert!(info.start_time.is_none()); + assert!(info.stage.is_empty(), "the persisted snapshot should have been decoded before commit"); + drop(pool_meta); + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + assert!(canceler.is_cancelled()); + } + + #[tokio::test] + async fn test_decommission_cancel_transaction_survives_caller_abort_after_durable_save() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + let persisted = Arc::new(std::sync::Mutex::new(None)); + let (durable_tx, durable_rx) = tokio::sync::oneshot::channel(); + let save_release = Arc::new(tokio::sync::Notify::new()); + let cancel = tokio::spawn({ + let store = store.clone(); + let canceler = canceler.clone(); + let persisted = persisted.clone(); + let save_release = save_release.clone(); + async move { + store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + assert!( + snapshot.pools[0] + .decommission + .as_ref() + .is_some_and(|info| info.canceled && info.start_time.is_none()) + ); + *persisted.lock().expect("persisted cancel lock should not be poisoned") = + Some(snapshot.encode_config_data()?); + durable_tx + .send(()) + .map_err(|_| Error::other("failed to report durable cancel snapshot"))?; + save_release.notified().await; + Ok(()) + }) + .await + } + }); + durable_rx.await.expect("save hook should report the durable cancel snapshot"); + cancel.abort(); + let join_err = cancel.await.expect_err("caller cancel future should be aborted"); + assert!(join_err.is_cancelled()); + assert!( + store.pool_meta.try_read().is_err(), + "the detached transaction must retain its state guard" + ); + assert!( + store.decommission_cancelers.try_read().is_err(), + "the detached transaction must retain its owner guard" + ); + + save_release.notify_one(); + tokio::time::timeout(StdDuration::from_secs(1), canceler.token().cancelled()) + .await + .expect("detached cancel transaction should terminate the old token"); + + let persisted = persisted + .lock() + .expect("persisted cancel lock should not be poisoned") + .take() + .expect("save should capture the durable cancel snapshot"); + let mut durable = PoolMeta::default(); + durable + .load_from_config_data(persisted) + .expect("durable cancel snapshot should decode"); + assert!( + durable.pools[0] + .decommission + .as_ref() + .is_some_and(|info| info.canceled && info.start_time.is_none()) + ); + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + assert!(canceler.is_cancelled()); + + let side_effect_ran = Arc::new(AtomicBool::new(false)); + let operation_gate = store.ctx.decommission_operation_gate(); + let result = run_decommission_side_effect(canceler.token(), &operation_gate, { + let side_effect_ran = side_effect_ran.clone(); + move || async move { + side_effect_ran.store(true, Ordering::SeqCst); + Ok(()) + } + }) + .await; + assert!(matches!(result, Err(Error::OperationCanceled))); + assert!(!side_effect_ran.load(Ordering::SeqCst)); + assert!( + store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .is_some_and(|info| info.canceled && !info.failed) + ); + } + + #[tokio::test] + async fn test_decommission_cancel_quiesces_and_publishes_only_after_save() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + let operation_gate = store.ctx.decommission_operation_gate(); + let side_effect = operation_gate.read().await; + let save_started = Arc::new(tokio::sync::Notify::new()); + let save_release = Arc::new(tokio::sync::Notify::new()); + let save_entered = Arc::new(AtomicBool::new(false)); + + let mut cancel = tokio::spawn({ + let store = store.clone(); + let canceler = canceler.clone(); + let save_started = save_started.clone(); + let save_release = save_release.clone(); + let save_entered = save_entered.clone(); + async move { + store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_| async move { + save_entered.store(true, Ordering::SeqCst); + save_started.notify_one(); + save_release.notified().await; + Ok(()) + }) + .await + } + }); + + let mut cancel_holds_start_gate = false; + for _ in 0..100 { + if store.start_gate.try_lock().is_err() { + cancel_holds_start_gate = true; + break; + } + tokio::task::yield_now().await; + } + assert!(cancel_holds_start_gate, "cancel should reach the operation-gate wait"); + assert!(!cancel.is_finished(), "cancel must wait for the in-flight side effect"); + assert!(!save_entered.load(Ordering::SeqCst), "cancel must quiesce side effects before saving"); + { + let pool_meta = store.pool_meta.read().await; + assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok()); + } + assert!(!canceler.is_cancelled()); + + drop(side_effect); + tokio::time::timeout(StdDuration::from_secs(1), save_started.notified()) + .await + .expect("cancel should reach the injected save after side effects quiesce"); + assert!( + store.pool_meta_save_gate.try_lock().is_err(), + "the cancel save must exclude stale full-document saves until publication" + ); + assert!( + store.pool_meta.try_read().is_err(), + "the active generation must stay write-locked through persistence" + ); + assert!(!canceler.is_cancelled(), "the token must remain live until persistence commits"); + + let mut complete = tokio::spawn({ + let store = store.clone(); + async move { store.complete_decommission(0).await } + }); + let mut fail = tokio::spawn({ + let store = store.clone(); + async move { store.decommission_failed(0).await } + }); + tokio::task::yield_now().await; + assert!(!complete.is_finished(), "complete must serialize behind the pending cancel"); + assert!(!fail.is_finished(), "fail must serialize behind the pending cancel"); + + save_release.notify_one(); + tokio::time::timeout(StdDuration::from_secs(1), &mut cancel) + .await + .expect("cancel should finish after persistence commits") + .expect("cancel task should not panic") + .expect("cancel should commit"); + tokio::time::timeout(StdDuration::from_secs(1), &mut complete) + .await + .expect("complete should finish after cancel commits") + .expect("complete task should not panic") + .expect("stale complete should be a no-op"); + tokio::time::timeout(StdDuration::from_secs(1), &mut fail) + .await + .expect("fail should finish after cancel commits") + .expect("fail task should not panic") + .expect("stale fail should be a no-op"); + + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("cancel metadata should remain present"); + assert!(info.canceled); + assert!(!info.complete); + assert!(!info.failed); + assert!(info.start_time.is_none()); + drop(pool_meta); + assert!(canceler.is_cancelled()); + assert!(!canceler.is_active()); + assert!(store.decommission_cancelers.read().await[0].is_none()); + } + + #[test] + fn test_decommission_cancel_commit_rejects_a_replaced_generation() { + let old_generation = OffsetDateTime::UNIX_EPOCH; + let new_generation = old_generation + Duration::seconds(1); + let mut canceled = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(old_generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let previous_last_update = canceled.pools[0].last_update; + assert!(canceled.decommission_cancel(0)); + let canceled_pool = canceled.pools.remove(0); + let commit = super::DecommissionCancelCommit { + previous_start_time: Some(old_generation), + previous_queued: false, + previous_last_update, + canceled_pool, + }; + let mut current = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(new_generation), + ..Default::default() + }), + )], + ..Default::default() + }; + + let err = super::commit_decommission_cancel(&mut current, 0, commit) + .expect_err("an old cancel must not publish over a replacement generation"); + assert!(err.to_string().contains("operation generation changed")); + let info = current.pools[0] + .decommission + .as_ref() + .expect("replacement generation should remain present"); + assert_eq!(info.start_time, Some(new_generation)); + assert!(!info.canceled); + } + + #[tokio::test] + async fn test_decommission_cancel_rejects_stale_retry_after_queued_replacement() { + let old_generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(old_generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + let queued_replacement = Arc::new(std::sync::Mutex::new(None)); + let queued_replacement_for_save = queued_replacement.clone(); + + store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + let saved_cancel = snapshot.pools[0] + .decommission + .as_ref() + .expect("saved snapshot should contain decommission metadata"); + assert!(saved_cancel.canceled); + assert!(saved_cancel.start_time.is_none()); + let mut queued_replacement = snapshot.clone(); + let replacement = queued_replacement + .pools + .get_mut(0) + .expect("cancel snapshot should contain the pool"); + replacement.last_update += Duration::seconds(1); + let info = replacement + .decommission + .as_mut() + .expect("cancel snapshot should contain decommission metadata"); + info.canceled = false; + info.queued = true; + *queued_replacement_for_save + .lock() + .expect("queued replacement lock should not be poisoned") = Some(queued_replacement); + Ok(()) + }) + .await + .expect("cancel should commit before a queued replacement is installed"); + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + assert!(canceler.is_cancelled()); + + let queued_replacement = queued_replacement + .lock() + .expect("queued replacement lock should not be poisoned") + .take() + .expect("save should prepare the queued replacement"); + *store.pool_meta.write().await = queued_replacement; + + let replacement_revision = { + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("queued replacement should remain present"); + assert!(info.queued); + assert!(!info.canceled); + assert!(info.start_time.is_none()); + pool_meta.pools[0].last_update + }; + + store.retry_decommission_cancel_for_operation(0, &canceler).await; + + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("stale retry must preserve the queued replacement"); + assert_eq!(pool_meta.pools[0].last_update, replacement_revision); + assert!(info.queued); + assert!(!info.canceled); + assert!(info.start_time.is_none()); + } + #[tokio::test] async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() { let canceler = DecommissionCanceler::new(CancellationToken::new());