mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 20:36:38 +00:00
fix(ecstore): supervise decommission worker exits
This commit is contained in:
+267
-102
@@ -95,6 +95,7 @@ const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4;
|
||||
const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30;
|
||||
const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3;
|
||||
const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
const DECOMMISSION_TERMINAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
|
||||
/// Background decommission walks must tolerate slow object migrations; the
|
||||
/// stall timeout is the drive-health bound, not the total listing duration.
|
||||
const DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
|
||||
@@ -333,6 +334,12 @@ fn guard_decommission_cancelers(
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn await_decommission_worker(idx: usize, worker: tokio::task::JoinHandle<Result<()>>) -> Result<()> {
|
||||
worker
|
||||
.await
|
||||
.map_err(|err| Error::other(format!("decommission worker {idx} task join error: {err}")))?
|
||||
}
|
||||
|
||||
fn reserve_decommission_start_cancelers(
|
||||
pool_meta: &PoolMeta,
|
||||
indices: &[usize],
|
||||
@@ -524,7 +531,7 @@ fn spawn_decommission_index_cancelers(
|
||||
store: Arc<ECStore>,
|
||||
rx: CancellationToken,
|
||||
index_cancelers: Vec<(usize, DecommissionCancelerGuard)>,
|
||||
) {
|
||||
) -> tokio::task::JoinHandle<()> {
|
||||
tokio::spawn(async move {
|
||||
let mut stop_queue = false;
|
||||
|
||||
@@ -532,21 +539,16 @@ fn spawn_decommission_index_cancelers(
|
||||
let canceler = canceler_guard.canceler().clone();
|
||||
if stop_queue || rx.is_cancelled() {
|
||||
canceler.cancel();
|
||||
if let Err(err) = store.decommission_cancel_for_operation(idx, &canceler).await {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "queued_cancel_failed",
|
||||
error = %err,
|
||||
"Failed to cancel queued decommission"
|
||||
);
|
||||
}
|
||||
store.retry_decommission_cancel_for_operation(idx, &canceler).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Err(err) = store.do_decommission_in_routine(canceler, idx).await {
|
||||
let worker = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
async move { store.do_decommission_in_routine(canceler, idx).await }
|
||||
});
|
||||
if let Err(err) = await_decommission_worker(idx, worker).await {
|
||||
error!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
@@ -556,6 +558,7 @@ fn spawn_decommission_index_cancelers(
|
||||
error = %err,
|
||||
"Decommission routine failed"
|
||||
);
|
||||
store.retry_decommission_failed_for_operation(idx, &canceler).await;
|
||||
stop_queue = true;
|
||||
continue;
|
||||
}
|
||||
@@ -565,7 +568,7 @@ fn spawn_decommission_index_cancelers(
|
||||
!should_continue_decommission_queue(&pool_meta, idx)
|
||||
};
|
||||
}
|
||||
});
|
||||
})
|
||||
}
|
||||
|
||||
fn decommission_meta_bucket_options() -> MakeBucketOptions {
|
||||
@@ -2944,10 +2947,7 @@ impl ECStore {
|
||||
}
|
||||
|
||||
async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
let _canceler_guard = DecommissionCancelerGuard::new(owner.clone());
|
||||
let result = self.decommission_cancel_with_owner(idx, Some(owner)).await;
|
||||
self.release_decommission_canceler_slot(idx, owner).await;
|
||||
result
|
||||
self.decommission_cancel_with_owner(idx, Some(owner)).await
|
||||
}
|
||||
|
||||
async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||
@@ -2955,6 +2955,78 @@ impl ECStore {
|
||||
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner);
|
||||
}
|
||||
|
||||
async fn decommission_terminal_retryable_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> bool {
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
if !decommission_canceler_is_owned_by(cancelers.as_slice(), idx, owner) {
|
||||
owner.release();
|
||||
return false;
|
||||
}
|
||||
|
||||
let retryable = {
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
pool_meta
|
||||
.pools
|
||||
.get(idx)
|
||||
.and_then(|pool| pool.decommission.as_ref())
|
||||
.is_some_and(|info| info.has_decommission_state() && !info.complete && !info.failed && !info.canceled)
|
||||
};
|
||||
if !retryable {
|
||||
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner);
|
||||
}
|
||||
retryable
|
||||
}
|
||||
|
||||
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 {
|
||||
return;
|
||||
};
|
||||
if !self.decommission_terminal_retryable_for_operation(idx, owner).await {
|
||||
return;
|
||||
}
|
||||
attempt += 1;
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "terminal_save_retry",
|
||||
terminal = "canceled",
|
||||
attempt,
|
||||
error = %err,
|
||||
"Decommission terminal save will be retried"
|
||||
);
|
||||
tokio::time::sleep(DECOMMISSION_TERMINAL_RETRY_DELAY).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn retry_decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||
let mut attempt = 0usize;
|
||||
loop {
|
||||
let Err(err) = self.decommission_failed_for_operation(idx, owner).await else {
|
||||
return;
|
||||
};
|
||||
if !self.decommission_terminal_retryable_for_operation(idx, owner).await {
|
||||
return;
|
||||
}
|
||||
attempt += 1;
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "terminal_save_retry",
|
||||
terminal = "failed",
|
||||
attempt,
|
||||
error = %err,
|
||||
"Decommission terminal save will be retried"
|
||||
);
|
||||
tokio::time::sleep(DECOMMISSION_TERMINAL_RETRY_DELAY).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn decommission_cancel_with_owner(
|
||||
&self,
|
||||
idx: usize,
|
||||
@@ -2965,8 +3037,8 @@ impl ECStore {
|
||||
|
||||
// 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, canceled_worker) = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
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) {
|
||||
@@ -2995,19 +3067,25 @@ impl ECStore {
|
||||
) else {
|
||||
return Ok(());
|
||||
};
|
||||
let canceler = if let Some(owner) = owner {
|
||||
take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner)
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
take_decommission_canceler(cancelers.as_mut_slice(), idx)
|
||||
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),
|
||||
cancel_decommission_canceler(canceler),
|
||||
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,
|
||||
@@ -3028,6 +3106,10 @@ impl ECStore {
|
||||
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())?;
|
||||
@@ -3039,6 +3121,7 @@ impl ECStore {
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn clear_decommission(&self, idx: usize) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "clear decommission")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
let (should_reload_pool_meta, previous_pool_meta) = {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
@@ -3047,11 +3130,6 @@ impl ECStore {
|
||||
(changed, changed.then_some(previous_pool_meta))
|
||||
};
|
||||
|
||||
{
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx);
|
||||
}
|
||||
|
||||
if should_reload_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;
|
||||
@@ -3060,6 +3138,11 @@ impl ECStore {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
{
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx);
|
||||
}
|
||||
|
||||
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("clear_decommission for pool {idx}");
|
||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?;
|
||||
@@ -3136,12 +3219,29 @@ impl ECStore {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
let indices = {
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
first_resumable_decommission_queue_indices(&pool_meta)
|
||||
.into_iter()
|
||||
.filter(|idx| indices.contains(idx))
|
||||
.collect::<Vec<_>>()
|
||||
};
|
||||
if indices.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let index_cancelers = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
bind_decommission_cancelers(indices.as_slice(), rx, cancelers.as_mut_slice())
|
||||
let missing = missing_decommission_worker_prefix(indices.as_slice(), cancelers.as_slice());
|
||||
if missing.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let bound = bind_missing_decommission_cancelers(missing.as_slice(), rx, cancelers.as_mut_slice());
|
||||
let guards = guard_decommission_cancelers(bound);
|
||||
ensure_decommission_routines_scheduled(guards.len(), missing.len())?;
|
||||
guards
|
||||
};
|
||||
let index_cancelers = guard_decommission_cancelers(index_cancelers);
|
||||
ensure_decommission_routines_scheduled(index_cancelers.len(), indices.len())?;
|
||||
Ok(index_cancelers)
|
||||
}
|
||||
|
||||
@@ -3152,7 +3252,9 @@ impl ECStore {
|
||||
indices: Vec<usize>,
|
||||
) -> Result<()> {
|
||||
let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?;
|
||||
spawn_decommission_index_cancelers(store, rx, index_cancelers);
|
||||
if !index_cancelers.is_empty() {
|
||||
let _ = spawn_decommission_index_cancelers(store, rx, index_cancelers);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -3168,18 +3270,12 @@ impl ECStore {
|
||||
}
|
||||
|
||||
let rx = CancellationToken::new();
|
||||
let index_cancelers = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
let missing = missing_decommission_worker_prefix(indices.as_slice(), cancelers.as_slice());
|
||||
bind_missing_decommission_cancelers(missing.as_slice(), &rx, cancelers.as_mut_slice())
|
||||
};
|
||||
|
||||
let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?;
|
||||
if index_cancelers.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let index_cancelers = guard_decommission_cancelers(index_cancelers);
|
||||
spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers);
|
||||
let _ = spawn_decommission_index_cancelers(self.clone(), rx, index_cancelers);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -3204,7 +3300,7 @@ impl ECStore {
|
||||
let index_cancelers = self
|
||||
.start_decommission_with_routines(indices, &rx, local_indices.as_slice())
|
||||
.await?;
|
||||
spawn_decommission_index_cancelers(store, rx, index_cancelers);
|
||||
let _ = spawn_decommission_index_cancelers(store, rx, index_cancelers);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -3954,11 +4050,7 @@ impl ECStore {
|
||||
idx: usize,
|
||||
) -> Result<()> {
|
||||
let rx = canceler.token().clone();
|
||||
let _canceler_guard = DecommissionCancelerGuard::new(canceler.clone());
|
||||
let result = self.run_decommission_in_routine(rx, idx, &canceler).await;
|
||||
self.release_decommission_canceler_slot(idx, &canceler).await;
|
||||
|
||||
result
|
||||
self.run_decommission_in_routine(rx, idx, &canceler).await
|
||||
}
|
||||
|
||||
async fn run_decommission_in_routine(
|
||||
@@ -4167,10 +4259,7 @@ impl ECStore {
|
||||
}
|
||||
|
||||
async fn decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
let _canceler_guard = DecommissionCancelerGuard::new(owner.clone());
|
||||
let result = self.decommission_failed_with_owner(idx, Some(owner)).await;
|
||||
self.release_decommission_canceler_slot(idx, owner).await;
|
||||
result
|
||||
self.decommission_failed_with_owner(idx, Some(owner)).await
|
||||
}
|
||||
|
||||
async fn decommission_failed_with_owner(
|
||||
@@ -4183,8 +4272,8 @@ impl ECStore {
|
||||
|
||||
// Lock order: decommission_cancelers before pool_meta. Holding both makes
|
||||
// owner validation and the terminal transition one atomic operation.
|
||||
let (should_reload_pool_meta, previous_pool_meta) = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
let (should_reload_pool_meta, previous_pool_meta, terminal_canceler) = {
|
||||
let cancelers = self.decommission_cancelers.read().await;
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
let previous_pool_meta = pool_meta.clone();
|
||||
let Some(changed) = update_decommission_for_operation(
|
||||
@@ -4196,27 +4285,31 @@ impl ECStore {
|
||||
) else {
|
||||
return Ok(());
|
||||
};
|
||||
let canceler = if let Some(owner) = owner {
|
||||
take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner)
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
take_decommission_canceler(cancelers.as_mut_slice(), idx)
|
||||
cancelers.get(idx).and_then(Option::as_ref).cloned()
|
||||
};
|
||||
cancel_decommission_canceler(canceler);
|
||||
(changed, changed.then_some(previous_pool_meta))
|
||||
(changed, changed.then_some(previous_pool_meta), terminal_canceler)
|
||||
};
|
||||
|
||||
if should_reload_pool_meta {
|
||||
if 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 should_reload_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 should_reload_pool_meta {
|
||||
{
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
pool_meta.mark_decommission_progress_saved();
|
||||
}
|
||||
}
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
self.release_decommission_canceler_slot(idx, canceler).await;
|
||||
}
|
||||
if should_reload_pool_meta {
|
||||
if let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("decommission_failed for pool {idx}");
|
||||
if let Some(err) = observe_decommission_terminal_reload_result(
|
||||
@@ -4260,10 +4353,7 @@ impl ECStore {
|
||||
}
|
||||
|
||||
async fn complete_decommission_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
let _canceler_guard = DecommissionCancelerGuard::new(owner.clone());
|
||||
let result = self.complete_decommission_with_owner(idx, Some(owner)).await;
|
||||
self.release_decommission_canceler_slot(idx, owner).await;
|
||||
result
|
||||
self.complete_decommission_with_owner(idx, Some(owner)).await
|
||||
}
|
||||
|
||||
async fn complete_decommission_with_owner(
|
||||
@@ -4276,8 +4366,8 @@ impl ECStore {
|
||||
|
||||
// Lock order: decommission_cancelers before pool_meta. Holding both makes
|
||||
// owner validation and the terminal transition one atomic operation.
|
||||
let (should_reload_pool_meta, previous_pool_meta) = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
let (should_reload_pool_meta, previous_pool_meta, terminal_canceler) = {
|
||||
let cancelers = self.decommission_cancelers.read().await;
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
let previous_pool_meta = pool_meta.clone();
|
||||
let Some(changed) = update_decommission_for_operation(
|
||||
@@ -4289,27 +4379,31 @@ impl ECStore {
|
||||
) else {
|
||||
return Ok(());
|
||||
};
|
||||
let canceler = if let Some(owner) = owner {
|
||||
take_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner)
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
take_decommission_canceler(cancelers.as_mut_slice(), idx)
|
||||
cancelers.get(idx).and_then(Option::as_ref).cloned()
|
||||
};
|
||||
cancel_decommission_canceler(canceler);
|
||||
(changed, changed.then_some(previous_pool_meta))
|
||||
(changed, changed.then_some(previous_pool_meta), terminal_canceler)
|
||||
};
|
||||
|
||||
if should_reload_pool_meta {
|
||||
if 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 should_reload_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 should_reload_pool_meta {
|
||||
{
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
pool_meta.mark_decommission_progress_saved();
|
||||
}
|
||||
}
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
self.release_decommission_canceler_slot(idx, canceler).await;
|
||||
}
|
||||
if should_reload_pool_meta {
|
||||
if let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("complete_decommission for pool {idx}");
|
||||
if let Some(err) = observe_decommission_terminal_reload_result(
|
||||
@@ -5754,12 +5848,11 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
|
||||
#[cfg(test)]
|
||||
mod pools_tests {
|
||||
use super::{
|
||||
DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta,
|
||||
DecommissionCanceler, DecommissionCancelerGuard, DecommissionStartPoolState, DecommissionTerminalState,
|
||||
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo,
|
||||
DecommissionCanceler, DecommissionStartPoolState, DecommissionTerminalState,
|
||||
ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info,
|
||||
PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers,
|
||||
bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler,
|
||||
bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state,
|
||||
await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers,
|
||||
cancel_decommission_canceler,
|
||||
classify_decommission_terminal_state, count_decommission_item,
|
||||
count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options,
|
||||
decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options,
|
||||
@@ -5791,7 +5884,9 @@ DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
|
||||
should_continue_decommission_queue, should_count_decommission_version_complete,
|
||||
should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal,
|
||||
should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine,
|
||||
split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler,
|
||||
spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler,
|
||||
take_decommission_canceler,
|
||||
touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage,
|
||||
update_decommission_for_operation, validate_start_decommission_request, wait_decommission_listing_retry,
|
||||
wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context,
|
||||
wait_decommission_worker_drain, with_decommission_entry_context,
|
||||
@@ -5802,10 +5897,10 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
|
||||
use crate::disk::endpoint::Endpoint;
|
||||
use crate::error::{Error, StorageError};
|
||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
use crate::runtime::instance::InstanceContext;
|
||||
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
|
||||
use rustfs_filemeta::{
|
||||
FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo,
|
||||
};
|
||||
use crate::store::ECStore;
|
||||
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
|
||||
use rustfs_rio::Index;
|
||||
use std::sync::{
|
||||
Arc,
|
||||
@@ -5820,6 +5915,27 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
|
||||
Arc::new(|_| Box::pin(async {}))
|
||||
}
|
||||
|
||||
fn decommission_worker_test_store(
|
||||
pool_meta: PoolMeta,
|
||||
cancelers: Vec<Option<DecommissionCanceler>>,
|
||||
) -> Arc<ECStore> {
|
||||
let ctx = Arc::new(InstanceContext::new());
|
||||
let endpoint_pools = EndpointServerPools::default();
|
||||
Arc::new(ECStore {
|
||||
id: uuid::Uuid::new_v4(),
|
||||
disk_map: std::collections::HashMap::new(),
|
||||
pools: Vec::new(),
|
||||
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, ctx.clone()),
|
||||
pool_meta: tokio::sync::RwLock::new(pool_meta),
|
||||
rebalance_meta: tokio::sync::RwLock::new(None),
|
||||
decommission_cancelers: tokio::sync::RwLock::new(cancelers),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::new(()),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
})
|
||||
}
|
||||
|
||||
fn decommission_test_pool_endpoint(idx: usize, is_local: bool) -> PoolEndpoints {
|
||||
let port = 9000usize + idx;
|
||||
let mut endpoint =
|
||||
@@ -8748,24 +8864,73 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_canceler_guard_releases_operation_on_task_abort() {
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let cancelers = vec![Some(canceler.clone())];
|
||||
async fn test_decommission_supervisor_observes_worker_abort() {
|
||||
let (started_tx, started_rx) = tokio::sync::oneshot::channel();
|
||||
let worker_canceler = canceler.clone();
|
||||
let worker = tokio::spawn(async move {
|
||||
let _guard = DecommissionCancelerGuard::new(worker_canceler);
|
||||
started_tx.send(()).expect("worker start should be observed");
|
||||
std::future::pending::<()>().await;
|
||||
#[allow(unreachable_code)]
|
||||
Ok(())
|
||||
});
|
||||
|
||||
started_rx.await.expect("worker should install its guard");
|
||||
started_rx.await.expect("worker start should be observed");
|
||||
worker.abort();
|
||||
let join_error = worker.await.expect_err("aborted worker should return a join error");
|
||||
let err = await_decommission_worker(3, worker)
|
||||
.await
|
||||
.expect_err("supervisor should observe aborted worker");
|
||||
|
||||
assert!(join_error.is_cancelled());
|
||||
assert!(canceler.is_cancelled());
|
||||
assert!(!has_active_decommission_canceler(cancelers.as_slice()));
|
||||
assert!(err.to_string().contains("decommission worker 3 task join error"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_supervisor_observes_worker_panic() {
|
||||
let worker = tokio::spawn(async move {
|
||||
panic!("injected decommission worker panic");
|
||||
#[allow(unreachable_code)]
|
||||
Ok(())
|
||||
});
|
||||
|
||||
let err = await_decommission_worker(4, worker)
|
||||
.await
|
||||
.expect_err("supervisor should observe panicked worker");
|
||||
|
||||
assert!(err.to_string().contains("decommission worker 4 task join error"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_worker_metadata_missing_releases_owned_slot() {
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let store = decommission_worker_test_store(PoolMeta::default(), vec![Some(canceler.clone())]);
|
||||
canceler.cancel();
|
||||
|
||||
let err = store
|
||||
.do_decommission_in_routine(canceler.clone(), 0)
|
||||
.await
|
||||
.expect_err("missing worker metadata should fail the routine");
|
||||
|
||||
assert!(err.to_string().contains("target pool was not found"));
|
||||
assert!(!canceler.is_active());
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_supervisor_failure_cancels_queued_successor() {
|
||||
let first = DecommissionCanceler::new(CancellationToken::new());
|
||||
let queued = DecommissionCanceler::new(CancellationToken::new());
|
||||
let store = decommission_worker_test_store(
|
||||
PoolMeta::default(),
|
||||
vec![Some(first.clone()), Some(queued.clone())],
|
||||
);
|
||||
let guards = guard_decommission_cancelers(vec![(0, first.clone()), (1, queued.clone())]);
|
||||
|
||||
spawn_decommission_index_cancelers(store.clone(), CancellationToken::new(), guards)
|
||||
.await
|
||||
.expect("decommission supervisor should finish after queued cleanup");
|
||||
|
||||
assert!(!first.is_active());
|
||||
assert!(!queued.is_active());
|
||||
assert!(queued.is_cancelled());
|
||||
assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user