Compare commits

..

3 Commits

4 changed files with 398 additions and 787 deletions
+132 -716
View File
@@ -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))]
@@ -5254,6 +5143,9 @@ impl ECStore {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
self.ensure_decommission_multipart_uploads_drained(idx, &pool, &buckets)
.await?;
for (set_index, set) in pool.disk_set.iter().enumerate() {
for bucket_info in &buckets {
let mut lifecycle_config = None;
@@ -5367,6 +5259,49 @@ impl ECStore {
Ok(())
}
async fn ensure_decommission_multipart_uploads_drained(
&self,
idx: usize,
pool: &Sets,
buckets: &[DecomBucketInfo],
) -> Result<()> {
let mut bucket_names = buckets
.iter()
.filter(|bucket| bucket.name != RUSTFS_META_BUCKET)
.map(|bucket| bucket.name.as_str())
.collect::<Vec<_>>();
bucket_names.sort_unstable();
bucket_names.dedup();
// Take one bucket fence at a time so cross-bucket COPY cannot form an
// ABBA cycle. Suspension prevents new source uploads after each fence.
for bucket in bucket_names {
let lifecycle_guard = self.acquire_bucket_lifecycle_write_lock(bucket).await?;
if lifecycle_guard.is_lock_lost() {
return Err(Error::other(format!(
"decommission multipart drain lost the bucket lifecycle fence for `{bucket}`"
)));
}
for set in &pool.disk_set {
if let Some(upload_path) = set.first_multipart_upload_path_for_decommission(bucket).await? {
return Err(Error::other(format!(
"pool {idx} still contains multipart upload `{upload_path}` for bucket `{bucket}`; resolve it before retrying decommission"
)));
}
}
}
Ok(())
}
#[cfg(test)]
pub(crate) async fn ensure_decommission_multipart_uploads_drained_for_test(self: &Arc<Self>, idx: usize) -> Result<()> {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
self.ensure_decommission_multipart_uploads_drained(idx, pool.as_ref(), &buckets)
.await
}
#[tracing::instrument(skip(self, rd))]
async fn decommission_object(
self: Arc<Self>,
@@ -9756,525 +9691,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());
+64 -47
View File
@@ -556,6 +556,69 @@ async fn multipart_upload_paths_on_disk(disk: DiskStore, bucket: &str) -> disk::
}
impl SetDisks {
async fn discover_multipart_upload_paths(
&self,
orig_bucket: &str,
error_path: &str,
) -> Result<(Vec<Option<DiskStore>>, Vec<String>, usize)> {
let disks = self.disks.read().await.clone();
if disks.is_empty() {
return Err(Error::ErasureReadQuorum);
}
let discovery_quorum = if self.default_parity_count == 0 {
disks.len()
} else {
(disks.len() / 2).max(1)
};
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
let mut candidate_counts = HashMap::<String, usize>::new();
let mut discovery_tasks = JoinSet::new();
for (index, disk) in disks.iter().enumerate() {
let disk = disk.clone();
let orig_bucket = orig_bucket.to_string();
discovery_tasks.spawn(async move {
let result = match disk {
Some(disk) => multipart_upload_paths_on_disk(disk, &orig_bucket).await,
None => Err(DiskError::DiskNotFound),
};
(index, result)
});
}
while let Some(task_result) = discovery_tasks.join_next().await {
let Ok((index, result)) = task_result else {
continue;
};
match result {
Ok(paths) => {
discovery_errors[index] = None;
for path in paths {
*candidate_counts.entry(path).or_insert(0) += 1;
}
}
Err(err) => discovery_errors[index] = Some(err),
}
}
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
return Err(to_object_err(err.into(), vec![orig_bucket, error_path]));
}
let mut candidate_paths = candidate_counts
.into_iter()
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
.collect::<Vec<_>>();
candidate_paths.sort_unstable();
Ok((disks, candidate_paths, discovery_quorum))
}
pub(crate) async fn first_multipart_upload_path_for_decommission(&self, bucket: &str) -> Result<Option<String>> {
let (_, paths, _) = self
.discover_multipart_upload_paths(bucket, RUSTFS_META_MULTIPART_BUCKET)
.await?;
Ok(paths.into_iter().next())
}
async fn acquire_multipart_upload_read_lock(
&self,
op: &'static str,
@@ -747,53 +810,7 @@ impl SetDisks {
max_uploads: usize,
expected_incarnation_id: Option<Uuid>,
) -> Result<ListMultipartsInfo> {
let disks = self.disks.read().await.clone();
if disks.is_empty() {
return Err(Error::ErasureReadQuorum);
}
let discovery_quorum = if self.default_parity_count == 0 {
disks.len()
} else {
(disks.len() / 2).max(1)
};
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
let mut candidate_counts = HashMap::<String, usize>::new();
let mut discovery_tasks = JoinSet::new();
for (index, disk) in disks.iter().enumerate() {
let disk = disk.clone();
let bucket = bucket.to_string();
discovery_tasks.spawn(async move {
let result = match disk {
Some(disk) => multipart_upload_paths_on_disk(disk, &bucket).await,
None => Err(DiskError::DiskNotFound),
};
(index, result)
});
}
while let Some(task_result) = discovery_tasks.join_next().await {
let Ok((index, result)) = task_result else {
continue;
};
match result {
Ok(paths) => {
discovery_errors[index] = None;
for path in paths {
*candidate_counts.entry(path).or_insert(0) += 1;
}
}
Err(err) => discovery_errors[index] = Some(err),
}
}
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
return Err(to_object_err(err.into(), vec![bucket, prefix]));
}
let candidate_paths = candidate_counts
.into_iter()
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
.collect::<Vec<_>>();
let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix).await?;
let listed_uploads = stream::iter(candidate_paths)
.map(|upload_path| {
let disks = &disks;
+171
View File
@@ -1589,6 +1589,177 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn suspended_decommission_source_multipart_remains_operable_until_drained() {
let temp_dir = tempfile::tempdir().expect("create decommission multipart drain store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-multipart-drain", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decommission-multipart-drain-{}", uuid::Uuid::new_v4());
let complete_object = "complete.bin";
let abort_object = "abort.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create decommission multipart drain bucket");
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
let lifecycle_guard = store
.acquire_bucket_lifecycle_read_lock(&bucket)
.await
.expect("acquire multipart creation lifecycle fence");
let mut upload_opts = ObjectOptions {
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
let complete_upload = store.pools[0]
.new_multipart_upload(&bucket, complete_object, &upload_opts)
.await
.expect("create source upload to complete");
let abort_upload = store.pools[0]
.new_multipart_upload(&bucket, abort_object, &upload_opts)
.await
.expect("create source upload to abort");
drop(lifecycle_guard);
mark_test_pool_decommissioning(&store, 0).await;
let err = store
.ensure_decommission_multipart_uploads_drained_for_test(0)
.await
.expect_err("an unresolved source multipart upload must block final decommission");
let drain_error = err.to_string();
assert!(
drain_error.contains("still contains multipart upload") && drain_error.contains(&bucket),
"the drain error must identify both the upload path and user bucket: {drain_error}"
);
let listed = store
.list_multipart_uploads(&bucket, "", None, None, None, 100)
.await
.expect("list uploads from suspended decommission source");
assert!(
listed
.uploads
.iter()
.any(|upload| upload.upload_id.as_str() == complete_upload.upload_id.as_str()),
"the upload selected before suspension must remain visible"
);
store
.get_multipart_info(&bucket, complete_object, &complete_upload.upload_id, &ObjectOptions::default())
.await
.expect("read upload metadata from suspended decommission source");
let mut part_reader = PutObjReader::from_vec(b"multipart body".to_vec());
let part = store
.put_object_part(
&bucket,
complete_object,
&complete_upload.upload_id,
1,
&mut part_reader,
&ObjectOptions::default(),
)
.await
.expect("write part to suspended decommission source");
let parts = store
.list_object_parts(&bucket, complete_object, &complete_upload.upload_id, None, 100, &ObjectOptions::default())
.await
.expect("list parts from suspended decommission source");
assert_eq!(parts.parts.len(), 1);
assert_eq!(parts.parts[0].etag.as_deref(), part.etag.as_deref());
store
.clone()
.complete_multipart_upload(
&bucket,
complete_object,
&complete_upload.upload_id,
vec![crate::storage_api_contracts::multipart::CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.expect("complete upload on suspended decommission source");
store
.abort_multipart_upload(&bucket, abort_object, &abort_upload.upload_id, &ObjectOptions::default())
.await
.expect("abort upload on suspended decommission source");
store
.ensure_decommission_multipart_uploads_drained_for_test(0)
.await
.expect("final decommission gate should open after all source uploads are resolved");
assert_pool_object_present(&store.pools[0], &bucket, complete_object).await;
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn active_multipart_upload_routes_before_faulted_suspended_source() {
let temp_dir = tempfile::tempdir().expect("create active-first multipart routing store dir");
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "active-first-multipart-routing", &[4, 4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("active-first-multipart-routing-{}", uuid::Uuid::new_v4());
let object = "target-upload.bin";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create active-first multipart routing bucket");
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
let lifecycle_guard = store
.acquire_bucket_lifecycle_read_lock(&bucket)
.await
.expect("acquire multipart creation lifecycle fence");
let mut upload_opts = ObjectOptions {
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
let upload = store.pools[1]
.new_multipart_upload(&bucket, object, &upload_opts)
.await
.expect("create upload in active target pool");
drop(lifecycle_guard);
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
let original_source_disks = {
let mut disks = source_set.disks.write().await;
let original = disks.clone();
disks.fill(None);
original
};
let source_result = store.pools[0]
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
.await;
let routed_result = store
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
.await;
*source_set.disks.write().await = original_source_disks;
assert!(
matches!(&source_result, Err(StorageError::ErasureReadQuorum)),
"the suspended source must expose the injected hard read failure: {source_result:?}"
);
let routed = routed_result.expect("the active target UploadID must be resolved before the faulted suspended source");
assert_eq!(routed.upload_id, upload.upload_id);
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn delete_objects_skips_active_rebalance_source_pool() {
+31 -24
View File
@@ -196,6 +196,25 @@ async fn list_pool_multipart_uploads_for_incarnation(
}
impl ECStore {
async fn existing_multipart_pool_order(&self) -> Vec<usize> {
// A draining source must not hide a valid UploadID in an active target,
// while physical order within each phase preserves fail-closed errors.
let mut active = Vec::with_capacity(self.pools.len());
let mut draining = Vec::new();
for (idx, pool) in self.pools.iter().enumerate() {
if self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
if self.is_suspended(pool.pool_idx).await {
draining.push(idx);
} else {
active.push(idx);
}
}
active.extend(draining);
active
}
#[allow(clippy::too_many_arguments)]
pub async fn list_multipart_uploads_for_bucket_incarnation(
&self,
@@ -290,10 +309,8 @@ impl ECStore {
.await;
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
return match pool
.list_object_parts(bucket, object, upload_id, part_number_marker, max_parts, opts)
.await
@@ -353,10 +370,8 @@ impl ECStore {
let mut common_prefixes = HashSet::new();
let mut source_truncated = false;
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
let res = list_pool_multipart_uploads_for_incarnation(
pool,
bucket,
@@ -523,10 +538,8 @@ impl ECStore {
.await;
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
Ok(res) => return Ok(res),
Err(err) => {
@@ -586,10 +599,8 @@ impl ECStore {
return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await;
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
return match pool.get_multipart_info(bucket, object, upload_id, opts).await {
Ok(res) => Ok(res),
@@ -624,10 +635,8 @@ impl ECStore {
return self.pools[0].abort_multipart_upload(bucket, object, upload_id, opts).await;
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
let err = match pool.abort_multipart_upload(bucket, object, upload_id, opts).await {
Ok(_) => return Ok(()),
@@ -685,10 +694,8 @@ impl ECStore {
.await;
}
for pool in self.pools.iter() {
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
for pool_idx in self.existing_multipart_pool_order().await {
let pool = &self.pools[pool_idx];
let pool = pool.clone();
let err = match pool