diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 19c99ddfd..1e80cea90 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -231,10 +231,7 @@ fn bind_missing_decommission_cancelers( bound } -fn take_decommission_canceler( - cancelers: &mut [Option], - idx: usize, -) -> Option { +fn take_decommission_canceler(cancelers: &mut [Option], idx: usize) -> Option { cancelers.get_mut(idx).and_then(Option::take) } @@ -244,10 +241,7 @@ fn take_decommission_canceler_for_operation( owner: &DecommissionCanceler, ) -> Option { let slot = cancelers.get_mut(idx)?; - if slot - .as_ref() - .is_some_and(|canceler| canceler.owns_same_operation(owner)) - { + if slot.as_ref().is_some_and(|canceler| canceler.owns_same_operation(owner)) { slot.take() } else { None @@ -283,10 +277,7 @@ fn update_decommission_for_operation( } fn has_active_decommission_canceler(cancelers: &[Option]) -> bool { - cancelers - .iter() - .flatten() - .any(DecommissionCanceler::is_active) + cancelers.iter().flatten().any(DecommissionCanceler::is_active) } fn cancel_decommission_canceler(canceler: Option) -> bool { @@ -326,9 +317,7 @@ fn ensure_decommission_routines_scheduled(bound_count: usize, expected_count: us Ok(()) } -fn guard_decommission_cancelers( - index_cancelers: Vec<(usize, DecommissionCanceler)>, -) -> Vec<(usize, DecommissionCancelerGuard)> { +fn guard_decommission_cancelers(index_cancelers: Vec<(usize, DecommissionCanceler)>) -> Vec<(usize, DecommissionCancelerGuard)> { index_cancelers .into_iter() .map(|(idx, canceler)| (idx, DecommissionCancelerGuard::new(canceler))) @@ -3028,11 +3017,7 @@ impl ECStore { } } - async fn decommission_cancel_with_owner( - &self, - idx: usize, - owner: Option<&DecommissionCanceler>, - ) -> Result<()> { + 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; @@ -3059,13 +3044,9 @@ impl ECStore { 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 { + 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 { @@ -3084,9 +3065,7 @@ impl ECStore { terminal_canceler, ) }; - let canceled_worker = terminal_canceler - .as_ref() - .is_some_and(DecommissionCanceler::is_active); + let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active); if !canceled_worker && !already_canceled { warn!( event = EVENT_DECOMMISSION_STATE, @@ -4045,11 +4024,7 @@ impl ECStore { } #[tracing::instrument(skip(self, canceler))] - pub async fn do_decommission_in_routine( - self: &Arc, - canceler: DecommissionCanceler, - idx: usize, - ) -> Result<()> { + pub async fn do_decommission_in_routine(self: &Arc, canceler: DecommissionCanceler, idx: usize) -> Result<()> { let rx = canceler.token().clone(); self.run_decommission_in_routine(rx, idx, &canceler).await } @@ -4263,11 +4238,7 @@ impl ECStore { self.decommission_failed_with_owner(idx, Some(owner)).await } - async fn decommission_failed_with_owner( - &self, - idx: usize, - owner: Option<&DecommissionCanceler>, - ) -> Result<()> { + async fn decommission_failed_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta()) .await } @@ -4290,13 +4261,11 @@ impl ECStore { 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( - cancelers.as_slice(), - &mut pool_meta, - idx, - owner, - |pool_meta| pool_meta.decommission_failed(idx), - ) else { + let Some(changed) = + update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| { + pool_meta.decommission_failed(idx) + }) + else { return Ok(()); }; let terminal_canceler = if let Some(owner) = owner { @@ -4370,11 +4339,7 @@ impl ECStore { self.complete_decommission_with_owner(idx, Some(owner)).await } - async fn complete_decommission_with_owner( - &self, - idx: usize, - owner: Option<&DecommissionCanceler>, - ) -> Result<()> { + async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?; let _start_guard = self.start_gate.lock().await; @@ -4384,13 +4349,11 @@ impl ECStore { 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( - cancelers.as_slice(), - &mut pool_meta, - idx, - owner, - |pool_meta| pool_meta.decommission_complete(idx), - ) else { + let Some(changed) = + update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| { + pool_meta.decommission_complete(idx) + }) + else { return Ok(()); }; let terminal_canceler = if let Some(owner) = owner { @@ -4627,13 +4590,7 @@ impl ECStore { let mut cancelers = self.decommission_cancelers.write().await; let pool_meta = self.pool_meta.read().await; ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?; - reserve_decommission_start_cancelers( - &pool_meta, - &indices, - local_indices, - rx, - cancelers.as_mut_slice(), - )? + reserve_decommission_start_cancelers(&pool_meta, &indices, local_indices, rx, cancelers.as_mut_slice())? } else { let pool_meta = self.pool_meta.read().await; 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)] mod pools_tests { use super::{ - DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, - DecommissionCanceler, DecommissionStartPoolState, DecommissionTerminalState, - ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, - await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, - cancel_decommission_canceler, - classify_decommission_terminal_state, count_decommission_item, + DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler, + DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, + PoolStatus, apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers, + bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, 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, 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, @@ -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_terminal_operation_supported, ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, get_by_index, - guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, - is_decommission_cancel_requested, + guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, 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, - 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_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_partial_listing_entry, resolve_decommission_pool_meta_reload_result, - resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result, - resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result, - resolve_decommission_update_after_result, + resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result, + resolve_decommission_progress_save_result, resolve_decommission_terminal_mark_after_error_result, + resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result, 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, - should_cleanup_decommission_source_entry, + run_decommission_buckets_bounded, run_decommission_listing_with_retry, should_cleanup_decommission_source_entry, 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, 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, + 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_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::disk::endpoint::Endpoint; @@ -5929,10 +5874,7 @@ 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>, - ) -> Arc { + fn decommission_worker_test_store(pool_meta: PoolMeta, cancelers: Vec>) -> Arc { let ctx = Arc::new(InstanceContext::new()); let endpoint_pools = EndpointServerPools::default(); Arc::new(ECStore { @@ -8590,14 +8532,8 @@ track_decommission_current_object, track_decommission_current_object_stage, vali let second_parent = CancellationToken::new(); let mut cancelers = vec![None, None]; - let first = reserve_decommission_start_cancelers( - &pool_meta, - &[0], - &[0], - &first_parent, - cancelers.as_mut_slice(), - ) - .expect("first start should reserve its worker"); + let first = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &first_parent, cancelers.as_mut_slice()) + .expect("first start should reserve its worker"); pool_meta .decommission( 0, @@ -8609,13 +8545,7 @@ track_decommission_current_object, track_decommission_current_object_stage, vali ) .expect("first start should install active metadata"); - let second = reserve_decommission_start_cancelers( - &pool_meta, - &[0], - &[0], - &second_parent, - cancelers.as_mut_slice(), - ); + let second = reserve_decommission_start_cancelers(&pool_meta, &[0], &[0], &second_parent, cancelers.as_mut_slice()); assert!(matches!(second, Err(Error::DecommissionAlreadyRunning))); 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] fn test_missing_decommission_worker_prefix_stops_at_active_worker() { - let cancelers = vec![ - None, - Some(DecommissionCanceler::new(CancellationToken::new())), - None, - ]; + let cancelers = vec![None, Some(DecommissionCanceler::new(CancellationToken::new())), None]; 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() { 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 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) @@ -9021,13 +8944,9 @@ track_decommission_current_object, track_decommission_current_object_stage, vali ..Default::default() }; - let changed = update_decommission_for_operation( - cancelers.as_slice(), - &mut pool_meta, - 0, - Some(&stale), - |pool_meta| pool_meta.decommission_cancel(0), - ); + let changed = update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, 0, Some(&stale), |pool_meta| { + pool_meta.decommission_cancel(0) + }); assert!(changed.is_none()); assert!(