mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1ac09383ad |
+716
-132
@@ -1120,6 +1120,48 @@ fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta:
|
||||
*pool_meta = previous_pool_meta;
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DecommissionCancelCommit {
|
||||
previous_start_time: Option<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);
|
||||
}
|
||||
@@ -1570,7 +1612,7 @@ struct PersistedPoolMeta {
|
||||
pub pools: Vec<PersistedPoolStatus>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolStatus {
|
||||
#[serde(rename = "id")]
|
||||
@@ -1583,7 +1625,7 @@ struct PersistedPoolStatus {
|
||||
pub decommission: Option<PersistedPoolDecommissionInfo>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolDecommissionInfo {
|
||||
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
|
||||
@@ -3004,14 +3046,156 @@ impl ECStore {
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn decommission_cancel(&self, idx: usize) -> Result<()> {
|
||||
pub async fn decommission_cancel(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||
self.decommission_cancel_with_owner(idx, None).await
|
||||
}
|
||||
|
||||
async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
async fn decommission_cancel_for_operation(self: &Arc<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);
|
||||
@@ -3039,7 +3223,11 @@ impl ECStore {
|
||||
retryable
|
||||
}
|
||||
|
||||
async fn retry_decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||
async fn retry_decommission_cancel_for_operation(self: &Arc<Self>, idx: usize, owner: &DecommissionCanceler) {
|
||||
if !self.decommission_terminal_retryable_for_operation(idx, owner).await {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut attempt = 0usize;
|
||||
loop {
|
||||
let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else {
|
||||
@@ -3089,87 +3277,10 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
|
||||
async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
// Lock order: decommission_cancelers before pool_meta. Holding both makes
|
||||
// owner validation and the terminal transition one atomic operation.
|
||||
let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta, terminal_canceler) = {
|
||||
let cancelers = self.decommission_cancelers.read().await;
|
||||
let mut lock = self.pool_meta.write().await;
|
||||
let mut already_canceled = false;
|
||||
let (pool_present, decommission_present, terminal) = if let Some(pool) = lock.pools.get(idx) {
|
||||
if let Some(info) = pool.decommission.as_ref() {
|
||||
already_canceled = info.canceled;
|
||||
(
|
||||
true,
|
||||
info.has_decommission_state(),
|
||||
should_reject_decommission_cancel_as_terminal(info.complete, info.failed),
|
||||
)
|
||||
} else {
|
||||
(true, false, false)
|
||||
}
|
||||
} else {
|
||||
(false, false, false)
|
||||
};
|
||||
|
||||
ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?;
|
||||
let previous_pool_meta = lock.clone();
|
||||
let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut lock, idx, owner, |pool_meta| {
|
||||
pool_meta.decommission_cancel(idx)
|
||||
}) else {
|
||||
return Ok(());
|
||||
};
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
cancelers.get(idx).and_then(Option::as_ref).cloned()
|
||||
};
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
canceler.cancel();
|
||||
}
|
||||
(
|
||||
changed,
|
||||
should_retry_decommission_cancel_reload(changed, already_canceled),
|
||||
already_canceled,
|
||||
changed.then_some(previous_pool_meta),
|
||||
terminal_canceler,
|
||||
)
|
||||
};
|
||||
let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active);
|
||||
if !canceled_worker && !already_canceled {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "cancel_skipped",
|
||||
reason = "no_active_canceler",
|
||||
"Decommission cancel skipped"
|
||||
);
|
||||
}
|
||||
|
||||
self.wait_for_decommission_side_effects().await;
|
||||
|
||||
if should_save_pool_meta && let Err(err) = self.save_current_pool_meta().await {
|
||||
if let Some(previous_pool_meta) = previous_pool_meta {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta);
|
||||
}
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
self.release_decommission_canceler_slot(idx, canceler).await;
|
||||
}
|
||||
|
||||
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("decommission_cancel for pool {idx}");
|
||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
async fn decommission_cancel_with_owner(self: &Arc<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
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
@@ -5143,9 +5254,6 @@ 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;
|
||||
@@ -5259,49 +5367,6 @@ 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>,
|
||||
@@ -9691,6 +9756,525 @@ mod pools_tests {
|
||||
assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_save_failure_preserves_generation_and_token_until_retry() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
version: super::POOL_META_VERSION,
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
|
||||
let err = store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), |_| async { Err(Error::Timeout) })
|
||||
.await
|
||||
.expect_err("injected pool metadata timeout should fail cancel");
|
||||
assert!(matches!(err, Error::Timeout));
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("failed cancel must retain decommission metadata");
|
||||
assert_eq!(info.start_time, Some(generation));
|
||||
assert!(!info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok());
|
||||
}
|
||||
{
|
||||
let cancelers = store.decommission_cancelers.read().await;
|
||||
let current = cancelers[0].as_ref().expect("failed cancel must retain the worker owner");
|
||||
assert!(current.owns_same_operation(&canceler));
|
||||
assert!(current.is_active());
|
||||
}
|
||||
assert!(!canceler.is_cancelled());
|
||||
assert!(store.decommission_terminal_retryable_for_operation(0, &canceler).await);
|
||||
|
||||
let persisted = Arc::new(std::sync::Mutex::new(None));
|
||||
let persisted_for_save = persisted.clone();
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
let data = snapshot.encode_config_data()?;
|
||||
*persisted_for_save
|
||||
.lock()
|
||||
.expect("persisted snapshot lock should not be poisoned") = Some(data);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("cancel retry should commit");
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("committed cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
}
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let repeated_save_called = Arc::new(AtomicBool::new(false));
|
||||
let repeated_save_called_by_closure = repeated_save_called.clone();
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, None, move |_| async move {
|
||||
repeated_save_called_by_closure.store(true, Ordering::SeqCst);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("repeated cancel should be idempotent");
|
||||
assert!(!repeated_save_called.load(Ordering::SeqCst));
|
||||
|
||||
let persisted = persisted
|
||||
.lock()
|
||||
.expect("persisted snapshot lock should not be poisoned")
|
||||
.take()
|
||||
.expect("successful retry should capture persisted bytes");
|
||||
let mut restarted = PoolMeta::default();
|
||||
restarted
|
||||
.load_from_config_data(persisted)
|
||||
.expect("a restarted process should decode the committed cancel");
|
||||
assert!(
|
||||
restarted.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
assert!(first_resumable_decommission_queue_indices(&restarted).is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_serializes_reload_until_local_commit() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
version: super::POOL_META_VERSION,
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
stage: "migrate_object".to_string(),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let (persisted_tx, persisted_rx) = tokio::sync::oneshot::channel();
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
|
||||
let cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let save_release = save_release.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
persisted_tx
|
||||
.send(snapshot.encode_config_data()?)
|
||||
.map_err(|_| Error::other("failed to expose saved cancel snapshot"))?;
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
let persisted = persisted_rx.await.expect("save should expose the canceled snapshot");
|
||||
let reload_started = Arc::new(tokio::sync::Notify::new());
|
||||
let reload = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let reload_started = reload_started.clone();
|
||||
async move {
|
||||
let mut reloaded = PoolMeta::default();
|
||||
reloaded.load_from_config_data(persisted)?;
|
||||
reload_started.notify_one();
|
||||
*store.pool_meta.write().await = reloaded;
|
||||
Ok::<(), Error>(())
|
||||
}
|
||||
});
|
||||
reload_started.notified().await;
|
||||
assert!(!reload.is_finished(), "peer reload must wait for cancel publication");
|
||||
|
||||
save_release.notify_one();
|
||||
cancel
|
||||
.await
|
||||
.expect("cancel task should not panic")
|
||||
.expect("cancel should commit before releasing the reload");
|
||||
reload
|
||||
.await
|
||||
.expect("reload task should not panic")
|
||||
.expect("reload should install the saved cancel after publication");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("reloaded cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
assert!(info.stage.is_empty(), "the persisted snapshot should have been decoded before commit");
|
||||
drop(pool_meta);
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_transaction_survives_caller_abort_after_durable_save() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let persisted = Arc::new(std::sync::Mutex::new(None));
|
||||
let (durable_tx, durable_rx) = tokio::sync::oneshot::channel();
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
let cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let persisted = persisted.clone();
|
||||
let save_release = save_release.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
assert!(
|
||||
snapshot.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
*persisted.lock().expect("persisted cancel lock should not be poisoned") =
|
||||
Some(snapshot.encode_config_data()?);
|
||||
durable_tx
|
||||
.send(())
|
||||
.map_err(|_| Error::other("failed to report durable cancel snapshot"))?;
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
durable_rx.await.expect("save hook should report the durable cancel snapshot");
|
||||
cancel.abort();
|
||||
let join_err = cancel.await.expect_err("caller cancel future should be aborted");
|
||||
assert!(join_err.is_cancelled());
|
||||
assert!(
|
||||
store.pool_meta.try_read().is_err(),
|
||||
"the detached transaction must retain its state guard"
|
||||
);
|
||||
assert!(
|
||||
store.decommission_cancelers.try_read().is_err(),
|
||||
"the detached transaction must retain its owner guard"
|
||||
);
|
||||
|
||||
save_release.notify_one();
|
||||
tokio::time::timeout(StdDuration::from_secs(1), canceler.token().cancelled())
|
||||
.await
|
||||
.expect("detached cancel transaction should terminate the old token");
|
||||
|
||||
let persisted = persisted
|
||||
.lock()
|
||||
.expect("persisted cancel lock should not be poisoned")
|
||||
.take()
|
||||
.expect("save should capture the durable cancel snapshot");
|
||||
let mut durable = PoolMeta::default();
|
||||
durable
|
||||
.load_from_config_data(persisted)
|
||||
.expect("durable cancel snapshot should decode");
|
||||
assert!(
|
||||
durable.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let side_effect_ran = Arc::new(AtomicBool::new(false));
|
||||
let operation_gate = store.ctx.decommission_operation_gate();
|
||||
let result = run_decommission_side_effect(canceler.token(), &operation_gate, {
|
||||
let side_effect_ran = side_effect_ran.clone();
|
||||
move || async move {
|
||||
side_effect_ran.store(true, Ordering::SeqCst);
|
||||
Ok(())
|
||||
}
|
||||
})
|
||||
.await;
|
||||
assert!(matches!(result, Err(Error::OperationCanceled)));
|
||||
assert!(!side_effect_ran.load(Ordering::SeqCst));
|
||||
assert!(
|
||||
store.pool_meta.read().await.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && !info.failed)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_quiesces_and_publishes_only_after_save() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let operation_gate = store.ctx.decommission_operation_gate();
|
||||
let side_effect = operation_gate.read().await;
|
||||
let save_started = Arc::new(tokio::sync::Notify::new());
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
let save_entered = Arc::new(AtomicBool::new(false));
|
||||
|
||||
let mut cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let save_started = save_started.clone();
|
||||
let save_release = save_release.clone();
|
||||
let save_entered = save_entered.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_| async move {
|
||||
save_entered.store(true, Ordering::SeqCst);
|
||||
save_started.notify_one();
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
|
||||
let mut cancel_holds_start_gate = false;
|
||||
for _ in 0..100 {
|
||||
if store.start_gate.try_lock().is_err() {
|
||||
cancel_holds_start_gate = true;
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
assert!(cancel_holds_start_gate, "cancel should reach the operation-gate wait");
|
||||
assert!(!cancel.is_finished(), "cancel must wait for the in-flight side effect");
|
||||
assert!(!save_entered.load(Ordering::SeqCst), "cancel must quiesce side effects before saving");
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok());
|
||||
}
|
||||
assert!(!canceler.is_cancelled());
|
||||
|
||||
drop(side_effect);
|
||||
tokio::time::timeout(StdDuration::from_secs(1), save_started.notified())
|
||||
.await
|
||||
.expect("cancel should reach the injected save after side effects quiesce");
|
||||
assert!(
|
||||
store.pool_meta_save_gate.try_lock().is_err(),
|
||||
"the cancel save must exclude stale full-document saves until publication"
|
||||
);
|
||||
assert!(
|
||||
store.pool_meta.try_read().is_err(),
|
||||
"the active generation must stay write-locked through persistence"
|
||||
);
|
||||
assert!(!canceler.is_cancelled(), "the token must remain live until persistence commits");
|
||||
|
||||
let mut complete = tokio::spawn({
|
||||
let store = store.clone();
|
||||
async move { store.complete_decommission(0).await }
|
||||
});
|
||||
let mut fail = tokio::spawn({
|
||||
let store = store.clone();
|
||||
async move { store.decommission_failed(0).await }
|
||||
});
|
||||
tokio::task::yield_now().await;
|
||||
assert!(!complete.is_finished(), "complete must serialize behind the pending cancel");
|
||||
assert!(!fail.is_finished(), "fail must serialize behind the pending cancel");
|
||||
|
||||
save_release.notify_one();
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut cancel)
|
||||
.await
|
||||
.expect("cancel should finish after persistence commits")
|
||||
.expect("cancel task should not panic")
|
||||
.expect("cancel should commit");
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut complete)
|
||||
.await
|
||||
.expect("complete should finish after cancel commits")
|
||||
.expect("complete task should not panic")
|
||||
.expect("stale complete should be a no-op");
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut fail)
|
||||
.await
|
||||
.expect("fail should finish after cancel commits")
|
||||
.expect("fail task should not panic")
|
||||
.expect("stale fail should be a no-op");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
drop(pool_meta);
|
||||
assert!(canceler.is_cancelled());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_decommission_cancel_commit_rejects_a_replaced_generation() {
|
||||
let old_generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let new_generation = old_generation + Duration::seconds(1);
|
||||
let mut canceled = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(old_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let previous_last_update = canceled.pools[0].last_update;
|
||||
assert!(canceled.decommission_cancel(0));
|
||||
let canceled_pool = canceled.pools.remove(0);
|
||||
let commit = super::DecommissionCancelCommit {
|
||||
previous_start_time: Some(old_generation),
|
||||
previous_queued: false,
|
||||
previous_last_update,
|
||||
canceled_pool,
|
||||
};
|
||||
let mut current = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(new_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let err = super::commit_decommission_cancel(&mut current, 0, commit)
|
||||
.expect_err("an old cancel must not publish over a replacement generation");
|
||||
assert!(err.to_string().contains("operation generation changed"));
|
||||
let info = current.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("replacement generation should remain present");
|
||||
assert_eq!(info.start_time, Some(new_generation));
|
||||
assert!(!info.canceled);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_rejects_stale_retry_after_queued_replacement() {
|
||||
let old_generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(old_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let queued_replacement = Arc::new(std::sync::Mutex::new(None));
|
||||
let queued_replacement_for_save = queued_replacement.clone();
|
||||
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
let saved_cancel = snapshot.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("saved snapshot should contain decommission metadata");
|
||||
assert!(saved_cancel.canceled);
|
||||
assert!(saved_cancel.start_time.is_none());
|
||||
let mut queued_replacement = snapshot.clone();
|
||||
let replacement = queued_replacement
|
||||
.pools
|
||||
.get_mut(0)
|
||||
.expect("cancel snapshot should contain the pool");
|
||||
replacement.last_update += Duration::seconds(1);
|
||||
let info = replacement
|
||||
.decommission
|
||||
.as_mut()
|
||||
.expect("cancel snapshot should contain decommission metadata");
|
||||
info.canceled = false;
|
||||
info.queued = true;
|
||||
*queued_replacement_for_save
|
||||
.lock()
|
||||
.expect("queued replacement lock should not be poisoned") = Some(queued_replacement);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("cancel should commit before a queued replacement is installed");
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let queued_replacement = queued_replacement
|
||||
.lock()
|
||||
.expect("queued replacement lock should not be poisoned")
|
||||
.take()
|
||||
.expect("save should prepare the queued replacement");
|
||||
*store.pool_meta.write().await = queued_replacement;
|
||||
|
||||
let replacement_revision = {
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("queued replacement should remain present");
|
||||
assert!(info.queued);
|
||||
assert!(!info.canceled);
|
||||
assert!(info.start_time.is_none());
|
||||
pool_meta.pools[0].last_update
|
||||
};
|
||||
|
||||
store.retry_decommission_cancel_for_operation(0, &canceler).await;
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("stale retry must preserve the queued replacement");
|
||||
assert_eq!(pool_meta.pools[0].last_update, replacement_revision);
|
||||
assert!(info.queued);
|
||||
assert!(!info.canceled);
|
||||
assert!(info.start_time.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() {
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
|
||||
@@ -556,69 +556,6 @@ 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,
|
||||
@@ -810,7 +747,53 @@ impl SetDisks {
|
||||
max_uploads: usize,
|
||||
expected_incarnation_id: Option<Uuid>,
|
||||
) -> Result<ListMultipartsInfo> {
|
||||
let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix).await?;
|
||||
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 listed_uploads = stream::iter(candidate_paths)
|
||||
.map(|upload_path| {
|
||||
let disks = &disks;
|
||||
|
||||
@@ -1589,177 +1589,6 @@ 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() {
|
||||
|
||||
@@ -196,25 +196,6 @@ 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,
|
||||
@@ -309,8 +290,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
return match pool
|
||||
.list_object_parts(bucket, object, upload_id, part_number_marker, max_parts, opts)
|
||||
.await
|
||||
@@ -370,8 +353,10 @@ impl ECStore {
|
||||
let mut common_prefixes = HashSet::new();
|
||||
let mut source_truncated = false;
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let res = list_pool_multipart_uploads_for_incarnation(
|
||||
pool,
|
||||
bucket,
|
||||
@@ -538,8 +523,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
|
||||
Ok(res) => return Ok(res),
|
||||
Err(err) => {
|
||||
@@ -599,8 +586,10 @@ impl ECStore {
|
||||
return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
return match pool.get_multipart_info(bucket, object, upload_id, opts).await {
|
||||
Ok(res) => Ok(res),
|
||||
@@ -635,8 +624,10 @@ impl ECStore {
|
||||
return self.pools[0].abort_multipart_upload(bucket, object, upload_id, opts).await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let err = match pool.abort_multipart_upload(bucket, object, upload_id, opts).await {
|
||||
Ok(_) => return Ok(()),
|
||||
@@ -694,8 +685,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let pool = pool.clone();
|
||||
let err = match pool
|
||||
|
||||
Reference in New Issue
Block a user