style(ecstore): format decommission supervision tests

This commit is contained in:
overtrue
2026-08-22 13:58:45 +08:00
parent 255943fa43
commit ce4a72869d
+45 -126
View File
@@ -231,10 +231,7 @@ fn bind_missing_decommission_cancelers(
bound bound
} }
fn take_decommission_canceler( fn take_decommission_canceler(cancelers: &mut [Option<DecommissionCanceler>], idx: usize) -> Option<DecommissionCanceler> {
cancelers: &mut [Option<DecommissionCanceler>],
idx: usize,
) -> Option<DecommissionCanceler> {
cancelers.get_mut(idx).and_then(Option::take) cancelers.get_mut(idx).and_then(Option::take)
} }
@@ -244,10 +241,7 @@ fn take_decommission_canceler_for_operation(
owner: &DecommissionCanceler, owner: &DecommissionCanceler,
) -> Option<DecommissionCanceler> { ) -> Option<DecommissionCanceler> {
let slot = cancelers.get_mut(idx)?; let slot = cancelers.get_mut(idx)?;
if slot if slot.as_ref().is_some_and(|canceler| canceler.owns_same_operation(owner)) {
.as_ref()
.is_some_and(|canceler| canceler.owns_same_operation(owner))
{
slot.take() slot.take()
} else { } else {
None None
@@ -283,10 +277,7 @@ fn update_decommission_for_operation<T>(
} }
fn has_active_decommission_canceler(cancelers: &[Option<DecommissionCanceler>]) -> bool { fn has_active_decommission_canceler(cancelers: &[Option<DecommissionCanceler>]) -> bool {
cancelers cancelers.iter().flatten().any(DecommissionCanceler::is_active)
.iter()
.flatten()
.any(DecommissionCanceler::is_active)
} }
fn cancel_decommission_canceler(canceler: Option<DecommissionCanceler>) -> bool { fn cancel_decommission_canceler(canceler: Option<DecommissionCanceler>) -> bool {
@@ -326,9 +317,7 @@ fn ensure_decommission_routines_scheduled(bound_count: usize, expected_count: us
Ok(()) Ok(())
} }
fn guard_decommission_cancelers( fn guard_decommission_cancelers(index_cancelers: Vec<(usize, DecommissionCanceler)>) -> Vec<(usize, DecommissionCancelerGuard)> {
index_cancelers: Vec<(usize, DecommissionCanceler)>,
) -> Vec<(usize, DecommissionCancelerGuard)> {
index_cancelers index_cancelers
.into_iter() .into_iter()
.map(|(idx, canceler)| (idx, DecommissionCancelerGuard::new(canceler))) .map(|(idx, canceler)| (idx, DecommissionCancelerGuard::new(canceler)))
@@ -3028,11 +3017,7 @@ impl ECStore {
} }
} }
async fn decommission_cancel_with_owner( async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
&self,
idx: usize,
owner: Option<&DecommissionCanceler>,
) -> Result<()> {
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
let _start_guard = self.start_gate.lock().await; let _start_guard = self.start_gate.lock().await;
@@ -3059,13 +3044,9 @@ impl ECStore {
ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?; ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?;
let previous_pool_meta = lock.clone(); let previous_pool_meta = lock.clone();
let Some(changed) = update_decommission_for_operation( let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut lock, idx, owner, |pool_meta| {
cancelers.as_slice(), pool_meta.decommission_cancel(idx)
&mut lock, }) else {
idx,
owner,
|pool_meta| pool_meta.decommission_cancel(idx),
) else {
return Ok(()); return Ok(());
}; };
let terminal_canceler = if let Some(owner) = owner { let terminal_canceler = if let Some(owner) = owner {
@@ -3084,9 +3065,7 @@ impl ECStore {
terminal_canceler, terminal_canceler,
) )
}; };
let canceled_worker = terminal_canceler let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active);
.as_ref()
.is_some_and(DecommissionCanceler::is_active);
if !canceled_worker && !already_canceled { if !canceled_worker && !already_canceled {
warn!( warn!(
event = EVENT_DECOMMISSION_STATE, event = EVENT_DECOMMISSION_STATE,
@@ -4045,11 +4024,7 @@ impl ECStore {
} }
#[tracing::instrument(skip(self, canceler))] #[tracing::instrument(skip(self, canceler))]
pub async fn do_decommission_in_routine( pub async fn do_decommission_in_routine(self: &Arc<Self>, canceler: DecommissionCanceler, idx: usize) -> Result<()> {
self: &Arc<Self>,
canceler: DecommissionCanceler,
idx: usize,
) -> Result<()> {
let rx = canceler.token().clone(); let rx = canceler.token().clone();
self.run_decommission_in_routine(rx, idx, &canceler).await self.run_decommission_in_routine(rx, idx, &canceler).await
} }
@@ -4263,11 +4238,7 @@ impl ECStore {
self.decommission_failed_with_owner(idx, Some(owner)).await self.decommission_failed_with_owner(idx, Some(owner)).await
} }
async fn decommission_failed_with_owner( async fn decommission_failed_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
&self,
idx: usize,
owner: Option<&DecommissionCanceler>,
) -> Result<()> {
self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta()) self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta())
.await .await
} }
@@ -4290,13 +4261,11 @@ impl ECStore {
let cancelers = self.decommission_cancelers.read().await; let cancelers = self.decommission_cancelers.read().await;
let mut pool_meta = self.pool_meta.write().await; let mut pool_meta = self.pool_meta.write().await;
let previous_pool_meta = pool_meta.clone(); let previous_pool_meta = pool_meta.clone();
let Some(changed) = update_decommission_for_operation( let Some(changed) =
cancelers.as_slice(), update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| {
&mut pool_meta, pool_meta.decommission_failed(idx)
idx, })
owner, else {
|pool_meta| pool_meta.decommission_failed(idx),
) else {
return Ok(()); return Ok(());
}; };
let terminal_canceler = if let Some(owner) = owner { let terminal_canceler = if let Some(owner) = owner {
@@ -4370,11 +4339,7 @@ impl ECStore {
self.complete_decommission_with_owner(idx, Some(owner)).await self.complete_decommission_with_owner(idx, Some(owner)).await
} }
async fn complete_decommission_with_owner( async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
&self,
idx: usize,
owner: Option<&DecommissionCanceler>,
) -> Result<()> {
ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?; ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?;
let _start_guard = self.start_gate.lock().await; let _start_guard = self.start_gate.lock().await;
@@ -4384,13 +4349,11 @@ impl ECStore {
let cancelers = self.decommission_cancelers.read().await; let cancelers = self.decommission_cancelers.read().await;
let mut pool_meta = self.pool_meta.write().await; let mut pool_meta = self.pool_meta.write().await;
let previous_pool_meta = pool_meta.clone(); let previous_pool_meta = pool_meta.clone();
let Some(changed) = update_decommission_for_operation( let Some(changed) =
cancelers.as_slice(), update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| {
&mut pool_meta, pool_meta.decommission_complete(idx)
idx, })
owner, else {
|pool_meta| pool_meta.decommission_complete(idx),
) else {
return Ok(()); return Ok(());
}; };
let terminal_canceler = if let Some(owner) = owner { let terminal_canceler = if let Some(owner) = owner {
@@ -4627,13 +4590,7 @@ impl ECStore {
let mut cancelers = self.decommission_cancelers.write().await; let mut cancelers = self.decommission_cancelers.write().await;
let pool_meta = self.pool_meta.read().await; let pool_meta = self.pool_meta.read().await;
ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?; ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?;
reserve_decommission_start_cancelers( reserve_decommission_start_cancelers(&pool_meta, &indices, local_indices, rx, cancelers.as_mut_slice())?
&pool_meta,
&indices,
local_indices,
rx,
cancelers.as_mut_slice(),
)?
} else { } else {
let pool_meta = self.pool_meta.read().await; let pool_meta = self.pool_meta.read().await;
ensure_decommission_start_pool_states(&pool_meta, &indices)?; ensure_decommission_start_pool_states(&pool_meta, &indices)?;
@@ -5862,16 +5819,11 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
#[cfg(test)] #[cfg(test)]
mod pools_tests { mod pools_tests {
use super::{ use super::{
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler,
DecommissionCanceler, DecommissionStartPoolState, DecommissionTerminalState, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo,
ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, PoolStatus, apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers,
await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state,
cancel_decommission_canceler,
classify_decommission_terminal_state, count_decommission_item,
count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, 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,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo,
decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available,
ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool,
@@ -5879,33 +5831,26 @@ DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity, ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity,
ensure_decommission_terminal_operation_supported, ensure_local_decommission_pool_leaders, ensure_decommission_terminal_operation_supported, ensure_local_decommission_pool_leaders,
ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, get_by_index, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, get_by_index,
guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
is_decommission_cancel_requested,
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done,
merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result,
pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store, pool_meta_has_active_decommission, require_decommission_store, reserve_decommission_start_cancelers,
resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error, resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result,
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result, resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result,
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result, resolve_decommission_progress_save_result, resolve_decommission_terminal_mark_after_error_result,
resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result, resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
resolve_decommission_update_after_result,
resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta, resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta,
reserve_decommission_start_cancelers, run_decommission_buckets_bounded, run_decommission_listing_with_retry, run_decommission_buckets_bounded, run_decommission_listing_with_retry, should_cleanup_decommission_source_entry,
should_cleanup_decommission_source_entry,
should_continue_decommission_queue, should_count_decommission_version_complete, should_continue_decommission_queue, should_count_decommission_version_complete,
should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal, 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, should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine,
spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler, spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler,
take_decommission_canceler, take_decommission_canceler, touch_decommission_progress, track_decommission_current_object,
touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage, track_decommission_current_object_stage, update_decommission_for_operation, validate_start_decommission_request,
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_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context,
wait_decommission_worker_drain, with_decommission_entry_context,
touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage,
track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request,
}; };
use crate::data_movement; use crate::data_movement;
use crate::disk::endpoint::Endpoint; use crate::disk::endpoint::Endpoint;
@@ -5929,10 +5874,7 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
Arc::new(|_| Box::pin(async {})) Arc::new(|_| Box::pin(async {}))
} }
fn decommission_worker_test_store( fn decommission_worker_test_store(pool_meta: PoolMeta, cancelers: Vec<Option<DecommissionCanceler>>) -> Arc<ECStore> {
pool_meta: PoolMeta,
cancelers: Vec<Option<DecommissionCanceler>>,
) -> Arc<ECStore> {
let ctx = Arc::new(InstanceContext::new()); let ctx = Arc::new(InstanceContext::new());
let endpoint_pools = EndpointServerPools::default(); let endpoint_pools = EndpointServerPools::default();
Arc::new(ECStore { Arc::new(ECStore {
@@ -8590,14 +8532,8 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
let second_parent = CancellationToken::new(); let second_parent = CancellationToken::new();
let mut cancelers = vec![None, None]; let mut cancelers = vec![None, None];
let first = reserve_decommission_start_cancelers( let first = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &first_parent, cancelers.as_mut_slice())
&pool_meta, .expect("first start should reserve its worker");
&[0],
&[0],
&first_parent,
cancelers.as_mut_slice(),
)
.expect("first start should reserve its worker");
pool_meta pool_meta
.decommission( .decommission(
0, 0,
@@ -8609,13 +8545,7 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
) )
.expect("first start should install active metadata"); .expect("first start should install active metadata");
let second = reserve_decommission_start_cancelers( let second = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &second_parent, cancelers.as_mut_slice());
&pool_meta,
&[0],
&[0],
&second_parent,
cancelers.as_mut_slice(),
);
assert!(matches!(second, Err(Error::DecommissionAlreadyRunning))); assert!(matches!(second, Err(Error::DecommissionAlreadyRunning)));
let current = cancelers[0].as_ref().expect("first operation should retain the slot"); let current = cancelers[0].as_ref().expect("first operation should retain the slot");
@@ -8674,11 +8604,7 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
#[test] #[test]
fn test_missing_decommission_worker_prefix_stops_at_active_worker() { fn test_missing_decommission_worker_prefix_stops_at_active_worker() {
let cancelers = vec![ let cancelers = vec![None, Some(DecommissionCanceler::new(CancellationToken::new())), None];
None,
Some(DecommissionCanceler::new(CancellationToken::new())),
None,
];
let missing = missing_decommission_worker_prefix(&[0, 1, 2], cancelers.as_slice()); let missing = missing_decommission_worker_prefix(&[0, 1, 2], cancelers.as_slice());
@@ -8931,10 +8857,7 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
async fn test_decommission_supervisor_failure_cancels_queued_successor() { async fn test_decommission_supervisor_failure_cancels_queued_successor() {
let first = DecommissionCanceler::new(CancellationToken::new()); let first = DecommissionCanceler::new(CancellationToken::new());
let queued = DecommissionCanceler::new(CancellationToken::new()); let queued = DecommissionCanceler::new(CancellationToken::new());
let store = decommission_worker_test_store( let store = decommission_worker_test_store(PoolMeta::default(), vec![Some(first.clone()), Some(queued.clone())]);
PoolMeta::default(),
vec![Some(first.clone()), Some(queued.clone())],
);
let guards = guard_decommission_cancelers(vec![(0, first.clone()), (1, queued.clone())]); let guards = guard_decommission_cancelers(vec![(0, first.clone()), (1, queued.clone())]);
spawn_decommission_index_cancelers(store.clone(), CancellationToken::new(), guards) spawn_decommission_index_cancelers(store.clone(), CancellationToken::new(), guards)
@@ -9021,13 +8944,9 @@ track_decommission_current_object, track_decommission_current_object_stage, vali
..Default::default() ..Default::default()
}; };
let changed = update_decommission_for_operation( let changed = update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, 0, Some(&stale), |pool_meta| {
cancelers.as_slice(), pool_meta.decommission_cancel(0)
&mut pool_meta, });
0,
Some(&stale),
|pool_meta| pool_meta.decommission_cancel(0),
);
assert!(changed.is_none()); assert!(changed.is_none());
assert!( assert!(