mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 528e22c150 |
@@ -22,10 +22,9 @@ use tracing::{info, warn};
|
||||
|
||||
const BUCKET: &str = "conditional-put-race-bucket";
|
||||
|
||||
async fn cleanup_object(client: &Client, key: &str) {
|
||||
if let Err(e) = client.delete_object().bucket(BUCKET).key(key).send().await {
|
||||
warn!("Failed to delete object '{}' from bucket '{}' during cleanup: {:?}", key, BUCKET, e);
|
||||
}
|
||||
async fn cleanup_object(client: &Client, key: &str) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
client.delete_object().bucket(BUCKET).key(key).send().await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn conditional_put(
|
||||
@@ -71,14 +70,13 @@ async fn run_race_iteration(
|
||||
test_key: &str,
|
||||
iteration: usize,
|
||||
) -> Result<usize, Box<dyn std::error::Error + Send + Sync>> {
|
||||
cleanup_object(&clients[0], test_key).await;
|
||||
cleanup_object(&clients[0], test_key).await?;
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
|
||||
|
||||
let head_result = clients[0].head_object().bucket(BUCKET).key(test_key).send().await;
|
||||
|
||||
if head_result.is_ok() {
|
||||
warn!("Warning: Object still exists after cleanup, skipping iteration {}", iteration);
|
||||
return Ok(0);
|
||||
match clients[0].head_object().bucket(BUCKET).key(test_key).send().await {
|
||||
Ok(_) => return Err(format!("object still exists after cleanup in iteration {iteration}").into()),
|
||||
Err(error) if error.as_service_error().is_some_and(|error| error.is_not_found()) => {}
|
||||
Err(error) => return Err(format!("failed to verify cleanup in iteration {iteration}: {error:?}").into()),
|
||||
}
|
||||
|
||||
info!("\n=== Iteration {} ===", iteration);
|
||||
@@ -120,14 +118,16 @@ async fn run_race_iteration(
|
||||
|
||||
info!("Result: {} out of {} succeeded", success_count, clients.len());
|
||||
|
||||
if had_error {
|
||||
return Err("one or more conditional PUTs failed unexpectedly".into());
|
||||
}
|
||||
|
||||
if success_count > 1 {
|
||||
info!(">>> RACE CONDITION DETECTED!");
|
||||
} else if success_count == 1 {
|
||||
info!(">>> Correct behavior: exactly 1 writer succeeded.");
|
||||
} else if had_error {
|
||||
return Err("all conditional PUTs failed (e.g. cluster/bucket not ready)".into());
|
||||
} else {
|
||||
info!(">>> Unexpected: no writers succeeded.");
|
||||
return Err("no conditional PUT succeeded".into());
|
||||
}
|
||||
|
||||
Ok(success_count)
|
||||
@@ -167,7 +167,7 @@ async fn test_conditional_put_race_cluster() -> Result<(), Box<dyn std::error::E
|
||||
}
|
||||
}
|
||||
|
||||
cleanup_object(&clients[0], &test_key).await;
|
||||
cleanup_object(&clients[0], &test_key).await?;
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
|
||||
}
|
||||
|
||||
@@ -177,7 +177,7 @@ async fn test_conditional_put_race_cluster() -> Result<(), Box<dyn std::error::E
|
||||
info!("Total iterations: {}", iterations);
|
||||
info!("Correct (1 winner): {}", correct_count);
|
||||
info!("Race conditions: {}", races_detected);
|
||||
info!("Errors (skipped): {}", error_count);
|
||||
info!("Failed iterations: {}", error_count);
|
||||
|
||||
assert_eq!(races_detected, 0, "Race conditions detected: {}/{}", races_detected, iterations);
|
||||
assert_eq!(
|
||||
@@ -185,6 +185,10 @@ async fn test_conditional_put_race_cluster() -> Result<(), Box<dyn std::error::E
|
||||
"{} iteration(s) failed due to errors (e.g. cluster not ready)",
|
||||
error_count
|
||||
);
|
||||
assert_eq!(
|
||||
correct_count, iterations,
|
||||
"only {correct_count}/{iterations} iterations observed exactly one winner"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -201,7 +205,7 @@ async fn test_conditional_put_basic_cluster() -> Result<(), Box<dyn std::error::
|
||||
|
||||
let client = cluster.create_s3_client(0)?;
|
||||
let test_key = "basic-conditional-put";
|
||||
cleanup_object(&client, test_key).await;
|
||||
cleanup_object(&client, test_key).await?;
|
||||
|
||||
let result = client
|
||||
.put_object()
|
||||
@@ -233,6 +237,6 @@ async fn test_conditional_put_basic_cluster() -> Result<(), Box<dyn std::error::
|
||||
assert_eq!(code, "PreconditionFailed");
|
||||
}
|
||||
|
||||
cleanup_object(&client, test_key).await;
|
||||
cleanup_object(&client, test_key).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1120,48 +1120,6 @@ 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<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) {
|
||||
rollback_decommission_pool_meta(pool_meta, previous_pool_meta);
|
||||
}
|
||||
@@ -1612,7 +1570,7 @@ struct PersistedPoolMeta {
|
||||
pub pools: Vec<PersistedPoolStatus>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolStatus {
|
||||
#[serde(rename = "id")]
|
||||
@@ -1625,7 +1583,7 @@ struct PersistedPoolStatus {
|
||||
pub decommission: Option<PersistedPoolDecommissionInfo>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolDecommissionInfo {
|
||||
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
|
||||
@@ -3046,156 +3004,14 @@ impl ECStore {
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn decommission_cancel(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||
pub async fn decommission_cancel(&self, idx: usize) -> Result<()> {
|
||||
self.decommission_cancel_with_owner(idx, None).await
|
||||
}
|
||||
|
||||
async fn decommission_cancel_for_operation(self: &Arc<Self>, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
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 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);
|
||||
@@ -3223,11 +3039,7 @@ impl ECStore {
|
||||
retryable
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -3277,10 +3089,87 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
|
||||
async fn decommission_cancel_with_owner(self: &Arc<Self>, 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
|
||||
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(())
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
@@ -9756,525 +9645,6 @@ 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());
|
||||
|
||||
Reference in New Issue
Block a user