mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 21:26:28 +00:00
fix(ecstore): persist cancel before signaling (#6423)
* fix(ecstore): persist decommission cancel before signaling * test(ecstore): align cancel regressions with movement gate --------- Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
+726
-110
@@ -1546,6 +1546,48 @@ fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta:
|
|||||||
*pool_meta = previous_pool_meta;
|
*pool_meta = previous_pool_meta;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct DecommissionCancelCommit {
|
||||||
|
previous_start_time: Option<OffsetDateTime>,
|
||||||
|
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) {
|
fn rollback_start_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) {
|
||||||
rollback_decommission_pool_meta(pool_meta, previous_pool_meta);
|
rollback_decommission_pool_meta(pool_meta, previous_pool_meta);
|
||||||
}
|
}
|
||||||
@@ -2117,7 +2159,7 @@ struct PersistedPoolMeta {
|
|||||||
pub pools: Vec<PersistedPoolStatus>,
|
pub pools: Vec<PersistedPoolStatus>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||||
#[serde(deny_unknown_fields)]
|
#[serde(deny_unknown_fields)]
|
||||||
struct PersistedPoolStatus {
|
struct PersistedPoolStatus {
|
||||||
#[serde(rename = "id")]
|
#[serde(rename = "id")]
|
||||||
@@ -2130,7 +2172,7 @@ struct PersistedPoolStatus {
|
|||||||
pub decommission: Option<PersistedPoolDecommissionInfo>,
|
pub decommission: Option<PersistedPoolDecommissionInfo>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
|
||||||
#[serde(deny_unknown_fields)]
|
#[serde(deny_unknown_fields)]
|
||||||
struct PersistedPoolDecommissionInfo {
|
struct PersistedPoolDecommissionInfo {
|
||||||
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
|
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
|
||||||
@@ -3555,14 +3597,174 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[tracing::instrument(skip(self))]
|
#[tracing::instrument(skip(self))]
|
||||||
pub async fn decommission_cancel(&self, idx: usize) -> Result<()> {
|
pub async fn decommission_cancel(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||||
self.decommission_cancel_with_owner(idx, None).await
|
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<Self>, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||||
self.decommission_cancel_with_owner(idx, Some(owner)).await
|
self.decommission_cancel_with_owner(idx, Some(owner)).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn decommission_cancel_with_owner_and_save<Save, SaveFuture>(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
idx: usize,
|
||||||
|
owner: Option<&DecommissionCanceler>,
|
||||||
|
save_pool_meta: Save,
|
||||||
|
) -> Result<()>
|
||||||
|
where
|
||||||
|
Save: FnOnce(PoolMeta) -> SaveFuture + Send + 'static,
|
||||||
|
SaveFuture: Future<Output = Result<()>> + 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<Save, SaveFuture>(
|
||||||
|
&self,
|
||||||
|
idx: usize,
|
||||||
|
owner: Option<DecommissionCanceler>,
|
||||||
|
save_pool_meta: Save,
|
||||||
|
) -> Result<()>
|
||||||
|
where
|
||||||
|
Save: FnOnce(PoolMeta) -> SaveFuture,
|
||||||
|
SaveFuture: Future<Output = Result<()>>,
|
||||||
|
{
|
||||||
|
let owner = owner.as_ref();
|
||||||
|
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
|
||||||
|
let _start_guard = self.start_gate.lock().await;
|
||||||
|
let save_guard = self.pool_meta_save_gate.lock().await;
|
||||||
|
|
||||||
|
// Lock order: start gate, save gate, decommission_cancelers, then
|
||||||
|
// pool_meta. The state guards stay held across persistence so the
|
||||||
|
// active generation cannot change before the cancel is published.
|
||||||
|
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 changed = pending.is_some();
|
||||||
|
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);
|
||||||
|
|
||||||
|
if changed {
|
||||||
|
// Persistence must commit before the cancellation signal. Wait for
|
||||||
|
// in-flight movement only after signaling so readers can release
|
||||||
|
// the shared gate without deadlocking the terminal transition.
|
||||||
|
let movement_gate = self.ctx.data_movement_operation_gate();
|
||||||
|
let _movement_guard = movement_gate.write().await;
|
||||||
|
self.ctx.advance_data_movement_operation_epoch();
|
||||||
|
}
|
||||||
|
|
||||||
|
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||||
|
let stage = format!("decommission_cancel for pool {idx}");
|
||||||
|
if let Err(err) =
|
||||||
|
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())
|
||||||
|
{
|
||||||
|
warn!(
|
||||||
|
event = EVENT_DECOMMISSION_STATE,
|
||||||
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
|
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||||
|
pool_index = idx,
|
||||||
|
state = "terminal_reload_failed",
|
||||||
|
error = %err,
|
||||||
|
"Decommission cancel saved locally but pool meta reload failed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) {
|
async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||||
let mut cancelers = self.decommission_cancelers.write().await;
|
let mut cancelers = self.decommission_cancelers.write().await;
|
||||||
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner);
|
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner);
|
||||||
@@ -3590,7 +3792,11 @@ impl ECStore {
|
|||||||
retryable
|
retryable
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn retry_decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) {
|
async fn retry_decommission_cancel_for_operation(self: &Arc<Self>, idx: usize, owner: &DecommissionCanceler) {
|
||||||
|
if !self.decommission_terminal_retryable_for_operation(idx, owner).await {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
let mut attempt = 0usize;
|
let mut attempt = 0usize;
|
||||||
loop {
|
loop {
|
||||||
let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else {
|
let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else {
|
||||||
@@ -3640,111 +3846,10 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
async fn decommission_cancel_with_owner(self: &Arc<Self>, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
||||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
|
let pools = self.pools.clone();
|
||||||
let _start_guard = self.start_gate.lock().await;
|
self.decommission_cancel_with_owner_and_save(idx, owner, move |snapshot| async move { snapshot.save(pools).await })
|
||||||
// Signal cancellation before waiting for the movement writer. A worker
|
.await
|
||||||
// may still hold a publication read guard while it observes this
|
|
||||||
// signal; waiting for the writer first would deadlock that handoff.
|
|
||||||
if let Some(owner) = owner {
|
|
||||||
owner.cancel();
|
|
||||||
} else if let Some(canceler) = self.decommission_cancelers.read().await.get(idx).and_then(Option::as_ref) {
|
|
||||||
canceler.cancel();
|
|
||||||
}
|
|
||||||
let movement_gate = self.ctx.data_movement_operation_gate();
|
|
||||||
let _movement_guard = movement_gate.write().await;
|
|
||||||
// Lock order: movement gate, then decommission_cancelers, then pool_meta.
|
|
||||||
// Holding both state locks 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"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
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_save_pool_meta {
|
|
||||||
self.ctx.advance_data_movement_operation_epoch();
|
|
||||||
}
|
|
||||||
|
|
||||||
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
|
||||||
let stage = format!("decommission_cancel for pool {idx}");
|
|
||||||
if let Err(err) =
|
|
||||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())
|
|
||||||
{
|
|
||||||
warn!(
|
|
||||||
event = EVENT_DECOMMISSION_STATE,
|
|
||||||
component = LOG_COMPONENT_ECSTORE,
|
|
||||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
|
||||||
pool_index = idx,
|
|
||||||
state = "terminal_reload_failed",
|
|
||||||
error = %err,
|
|
||||||
"Decommission cancel saved locally but pool meta reload failed"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tracing::instrument(skip(self))]
|
#[tracing::instrument(skip(self))]
|
||||||
@@ -11956,6 +12061,517 @@ mod pools_tests {
|
|||||||
assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none));
|
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 {
|
||||||
|
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 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.data_movement_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_persists_before_signaling_and_quiesces_before_return() {
|
||||||
|
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.data_movement_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
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
tokio::time::timeout(StdDuration::from_secs(1), save_started.notified())
|
||||||
|
.await
|
||||||
|
.expect("cancel should persist while the side effect is still in flight");
|
||||||
|
assert!(save_entered.load(Ordering::SeqCst));
|
||||||
|
assert!(!cancel.is_finished(), "cancel must wait for the injected save");
|
||||||
|
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 fail = tokio::spawn({
|
||||||
|
let store = store.clone();
|
||||||
|
async move { store.decommission_failed(0).await }
|
||||||
|
});
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
assert!(!fail.is_finished(), "fail must serialize behind the pending cancel");
|
||||||
|
|
||||||
|
save_release.notify_one();
|
||||||
|
tokio::time::timeout(StdDuration::from_secs(1), canceler.token().cancelled())
|
||||||
|
.await
|
||||||
|
.expect("the durable cancel must signal the active worker");
|
||||||
|
assert!(!cancel.is_finished(), "cancel must wait for in-flight movement after the durable signal");
|
||||||
|
{
|
||||||
|
let pool_meta = store.pool_meta.read().await;
|
||||||
|
let info = pool_meta.pools[0]
|
||||||
|
.decommission
|
||||||
|
.as_ref()
|
||||||
|
.expect("the durable cancel must be published before quiescence");
|
||||||
|
assert!(info.canceled);
|
||||||
|
assert!(!info.complete);
|
||||||
|
assert!(!info.failed);
|
||||||
|
assert!(info.start_time.is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
drop(side_effect);
|
||||||
|
tokio::time::timeout(StdDuration::from_secs(1), &mut cancel)
|
||||||
|
.await
|
||||||
|
.expect("cancel should finish after in-flight movement quiesces")
|
||||||
|
.expect("cancel task should not panic")
|
||||||
|
.expect("cancel should commit");
|
||||||
|
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]
|
#[tokio::test]
|
||||||
async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() {
|
async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() {
|
||||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||||
|
|||||||
Reference in New Issue
Block a user