diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index e3c1194d8..ae470d017 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -7135,27 +7135,20 @@ impl ECStore { self.reserve_decommission_routines(rx, indices.as_slice()).await } - pub(crate) async fn spawn_decommission_routines( - &self, - store: Arc, - rx: CancellationToken, - indices: Vec, - ) -> Result<()> { - let index_cancelers = self.reserve_decommission_routines(&rx, indices.as_slice()).await?; - if !index_cancelers.is_empty() { - drop(spawn_decommission_index_cancelers( - store, - rx, - index_cancelers, - Arc::new(Semaphore::new(decommission_entry_concurrency_limit())), - )); - } - - Ok(()) + pub async fn spawn_missing_local_decommission_routines(self: &Arc) -> Result<()> { + self.spawn_missing_local_decommission_routines_with_token(CancellationToken::new()) + .await } - pub async fn spawn_missing_local_decommission_routines(self: &Arc) -> Result<()> { - let rx = CancellationToken::new(); + pub(crate) async fn has_active_local_decommission_worker(&self) -> bool { + let cancelers = self.decommission_cancelers.read().await; + has_active_decommission_canceler(cancelers.as_slice()) + } + + pub(crate) async fn spawn_missing_local_decommission_routines_with_token( + self: &Arc, + rx: CancellationToken, + ) -> Result<()> { let endpoints = self.endpoints(); let index_cancelers = self.reserve_missing_local_decommission_routines(&rx, &endpoints).await?; if index_cancelers.is_empty() { @@ -11132,10 +11125,16 @@ mod tests { assert!(replica_state.repair_write_safe); let worker_cancel = CancellationToken::new(); - store - .spawn_decommission_routines(Arc::clone(&store), worker_cancel.clone(), vec![0]) + let index_cancelers = store + .reserve_decommission_routines(&worker_cancel, &[0]) .await .expect("the committed activation should admit its decommission worker"); + drop(spawn_decommission_index_cancelers( + Arc::clone(&store), + worker_cancel.clone(), + index_cancelers, + Arc::new(Semaphore::new(decommission_entry_concurrency_limit())), + )); let admitted_cancel = store.decommission_cancelers.read().await[0] .clone() .expect("the admitted decommission worker should have a cancellation token"); diff --git a/crates/ecstore/src/error/mod.rs b/crates/ecstore/src/error/mod.rs index 9895018f7..e1bc6b739 100644 --- a/crates/ecstore/src/error/mod.rs +++ b/crates/ecstore/src/error/mod.rs @@ -948,10 +948,6 @@ pub fn is_err_data_movement_overwrite(err: &Error) -> bool { matches!(err, &StorageError::DataMovementOverwriteErr(_, _, _)) } -pub fn is_err_decommission_running(err: &Error) -> bool { - matches!(err, &StorageError::DecommissionAlreadyRunning) -} - #[allow(dead_code, reason = "predicate asserted by this file's tests (backlog#1823)")] pub fn is_err_rebalance_running(err: &Error) -> bool { matches!(err, &StorageError::RebalanceAlreadyRunning) @@ -1347,9 +1343,6 @@ mod tests { #[test] fn test_error_running_state_helpers() { - assert!(is_err_decommission_running(&StorageError::DecommissionAlreadyRunning)); - assert!(!is_err_decommission_running(&StorageError::RebalanceAlreadyRunning)); - assert!(is_err_rebalance_running(&StorageError::RebalanceAlreadyRunning)); assert!(!is_err_rebalance_running(&StorageError::DecommissionAlreadyRunning)); assert!(is_err_operation_canceled(&StorageError::OperationCanceled)); diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 7e0c24a1c..fdaaf6d88 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -17,18 +17,17 @@ use crate::core::pools::{ PoolMetaReplicaState, PoolMetaWriteState, load_pool_meta_identity_observing, local_decommission_queue_prefix, persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission, }; -use crate::error::is_err_decommission_running; use crate::runtime::instance::InstanceContext; use crate::runtime::sources as runtime_sources; use crate::storage_api_contracts::object::EcstoreObjectIO; use rustfs_config::server_config::KVS; use rustfs_credentials::{RPC_SECRET_REQUIRED_OPERATOR_MESSAGE, try_get_rpc_token}; +use std::future::Future; use tracing::{debug, error, info, warn}; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_STORE_INIT: &str = "store_init"; const EVENT_DECOMMISSION_RESUME_RETRY: &str = "decommission_resume_retry"; -const EVENT_DECOMMISSION_RESUME_FAILED: &str = "decommission_resume_failed"; const EVENT_STORE_FORMAT_RETRY: &str = "store_format_retry"; const EVENT_ECSTORE_INIT_STATUS: &str = "ecstore_init_status"; const EVENT_STORE_RPC_SECRET_PREFLIGHT_FAILED: &str = "store_rpc_secret_preflight_failed"; @@ -96,16 +95,13 @@ fn preflight_startup_rpc_secret_with( } } -const LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES: usize = 6; const LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY: Duration = Duration::from_secs(60 * 3); const LOCAL_DECOMMISSION_RESUME_RETRY_DELAY: Duration = Duration::from_secs(30); +const LOCAL_DECOMMISSION_WATCHDOG_INTERVAL: Duration = Duration::from_secs(30); +const LOCAL_DECOMMISSION_WATCHDOG_MAX_RETRY_DELAY: Duration = Duration::from_secs(60 * 5); const REBALANCE_INITIAL_RESUME_DELAY: Duration = Duration::from_secs(10); const REBALANCE_RESUME_RETRY_DELAY: Duration = Duration::from_secs(10); -fn should_retry_local_decommission_resume(err: &Error, attempt: usize) -> bool { - matches!(err, Error::ConfigNotFound) && attempt < LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES -} - fn should_retry_format_load(err: &Error) -> bool { !matches!(err, Error::CorruptedFormat) } @@ -118,14 +114,6 @@ fn should_defer_rebalance_auto_start(distributed: bool, fleet_proof_available: b distributed && !fleet_proof_available } -fn should_schedule_local_decommission_resume( - pool_indices: &[usize], - pool_meta_replica_state: PoolMetaReplicaState, - pool_meta_write_safe: bool, -) -> bool { - !pool_indices.is_empty() && pool_meta_replica_state.repair_write_safe && pool_meta_write_safe -} - async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: Duration) -> bool { tokio::select! { _ = rx.cancelled() => false, @@ -133,6 +121,13 @@ async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: } } +fn local_decommission_watchdog_retry_delay(consecutive_failures: u32) -> Duration { + let exponent = consecutive_failures.saturating_sub(1).min(4); + LOCAL_DECOMMISSION_RESUME_RETRY_DELAY + .saturating_mul(1_u32 << exponent) + .min(LOCAL_DECOMMISSION_WATCHDOG_MAX_RETRY_DELAY) +} + async fn wait_for_rebalance_resume_delay(rx: &CancellationToken, delay: Duration) -> bool { tokio::select! { _ = rx.cancelled() => false, @@ -235,71 +230,63 @@ where Ok(committed) } -async fn resume_local_decommission_after_init(store: Arc, rx: CancellationToken, pool_indices: Vec) { - for attempt in 0..=LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES { +async fn run_local_decommission_watchdog(rx: CancellationToken, mut reconcile: F) +where + F: FnMut() -> Fut, + Fut: Future>, +{ + let mut consecutive_failures = 0_u32; + + loop { if rx.is_cancelled() { return; } - let result = if pool_indices.len() > 1 { - store - .spawn_decommission_routines(store.clone(), rx.clone(), pool_indices.clone()) - .await - } else { - store.decommission(rx.clone(), pool_indices.clone()).await - }; - - match result { - Ok(()) => return, - Err(err) if is_err_decommission_running(&err) => { - if let Err(spawn_err) = store - .spawn_decommission_routines(store.clone(), rx.clone(), pool_indices.clone()) - .await - { - error!( - event = EVENT_DECOMMISSION_RESUME_FAILED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_STORE_INIT, - pool_indices = ?pool_indices, - error = %spawn_err, - reason = "spawn_workers_failed", - "Failed to resume decommission workers" - ); - } - return; + let delay = match reconcile().await { + Ok(()) => { + consecutive_failures = 0; + LOCAL_DECOMMISSION_WATCHDOG_INTERVAL } - Err(err) if should_retry_local_decommission_resume(&err, attempt) => { + Err(err) => { + consecutive_failures = consecutive_failures.saturating_add(1); + let retry_delay = local_decommission_watchdog_retry_delay(consecutive_failures); warn!( event = EVENT_DECOMMISSION_RESUME_RETRY, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_STORE_INIT, - pool_indices = ?pool_indices, - retry_count = attempt + 1, - retry_limit = LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES + 1, + consecutive_failures, + retry_delay_secs = retry_delay.as_secs(), error = %err, - "Retrying decommission resume after missing config" + "Retrying decommission worker recovery" ); - tokio::select! { - _ = rx.cancelled() => return, - _ = tokio::time::sleep(LOCAL_DECOMMISSION_RESUME_RETRY_DELAY) => {} - } - } - Err(err) => { - error!( - event = EVENT_DECOMMISSION_RESUME_FAILED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_STORE_INIT, - pool_indices = ?pool_indices, - error = %err, - reason = "resume_failed", - "Failed to resume decommission" - ); - return; + retry_delay } + }; + + if !wait_for_local_decommission_resume_delay(&rx, delay).await { + return; } } } +async fn supervise_local_decommission_after_init(store: Arc, rx: CancellationToken) { + run_local_decommission_watchdog(rx.clone(), || { + let store = store.clone(); + let worker_rx = rx.clone(); + async move { + store + .ensure_pool_meta_side_effects_safe("decommission worker recovery blocked while pool metadata requires recovery") + .await?; + if store.has_active_local_decommission_worker().await { + return Ok(()); + } + store.refresh_pool_status_meta().await?; + store.spawn_missing_local_decommission_routines_with_token(worker_rx).await + } + }) + .await; +} + async fn resume_rebalance_after_init(store: Arc, rx: CancellationToken) { if !wait_for_rebalance_resume_delay(&rx, REBALANCE_INITIAL_RESUME_DELAY).await { return; @@ -729,31 +716,32 @@ impl ECStore { } let local_pool_indices = local_decommission_queue_prefix(&endpoints, &pool_indices)?; + let has_local_decommission_leadership = endpoints.as_ref().iter().any(pool_first_endpoint_is_local); let pool_meta_write_safe = self .ensure_pool_meta_side_effects_safe("decommission resume blocked while pool metadata requires recovery") .await .is_ok(); - if should_schedule_local_decommission_resume(&local_pool_indices, pool_meta_replica_state, pool_meta_write_safe) { - let store = self.clone(); - let decommission_rx = rx.clone(); - - tokio::spawn(async move { - if !wait_for_local_decommission_resume_delay(&decommission_rx, LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY).await { - return; - } - resume_local_decommission_after_init(store, decommission_rx, local_pool_indices).await; - }); - } else if !local_pool_indices.is_empty() { - error!( - event = EVENT_DECOMMISSION_RESUME_FAILED, + if !pool_meta_replica_state.repair_write_safe || !pool_meta_write_safe { + warn!( + event = EVENT_DECOMMISSION_RESUME_RETRY, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_STORE_INIT, state = "blocked", pool_indices = ?local_pool_indices, reason = "pool_meta_write_blocked", - "Decommission resume blocked until pool metadata replicas are readable and consistent" + "Decommission watchdog waiting for pool metadata recovery" ); } + if has_local_decommission_leadership { + let store = self.clone(); + let decommission_rx = rx.clone(); + tokio::spawn(async move { + if !wait_for_local_decommission_resume_delay(&decommission_rx, LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY).await { + return; + } + supervise_local_decommission_after_init(store, decommission_rx).await; + }); + } runtime_sources::init_bucket_monitor_for_current_endpoints(); crate::bucket::bucket_target_sys::BucketTargetSys::get().start_heartbeat(); @@ -786,12 +774,12 @@ impl ECStore { #[cfg(test)] mod tests { use super::{ - LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, PoolMetaWriteState, establish_pool_meta_bootstrap_identity_if_proven, - load_pool_meta_for_startup, persist_pool_meta_for_startup_if_safe, pool_first_endpoint_is_local, - pool_meta_has_active_decommission, preflight_startup_rpc_secret_with, resolve_startup_pool_defaults_with, - resolve_store_init_stage_result, save_validated_pool_meta_for_startup, should_auto_start_rebalance_after_init, - should_defer_rebalance_auto_start, should_retry_format_load, should_retry_local_decommission_resume, - wait_for_local_decommission_resume_delay, + LOCAL_DECOMMISSION_RESUME_RETRY_DELAY, LOCAL_DECOMMISSION_WATCHDOG_MAX_RETRY_DELAY, PoolMetaWriteState, + establish_pool_meta_bootstrap_identity_if_proven, load_pool_meta_for_startup, local_decommission_watchdog_retry_delay, + persist_pool_meta_for_startup_if_safe, pool_first_endpoint_is_local, pool_meta_has_active_decommission, + preflight_startup_rpc_secret_with, resolve_startup_pool_defaults_with, resolve_store_init_stage_result, + run_local_decommission_watchdog, save_validated_pool_meta_for_startup, should_auto_start_rebalance_after_init, + should_defer_rebalance_auto_start, should_retry_format_load, wait_for_local_decommission_resume_delay, }; #[cfg(feature = "test-util")] use crate::disk::DiskAPI; @@ -1481,15 +1469,6 @@ mod tests { assert!(err.to_string().contains("cannot overwrite an unreadable replica")); assert!(!valid.wrote_without_lock.load(Ordering::SeqCst)); assert!(!unreadable.wrote_without_lock.load(Ordering::SeqCst)); - assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, true)); - assert!(!super::should_schedule_local_decommission_resume( - &[0], - crate::core::pools::PoolMetaReplicaState { - needs_repair: false, - repair_write_safe: true, - }, - false, - )); } #[tokio::test] @@ -1517,7 +1496,6 @@ mod tests { .contains("restart after all replicas are readable and consistent") ); assert!(!repaired.wrote_without_lock.load(Ordering::SeqCst)); - assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, false)); } #[test] @@ -1551,21 +1529,66 @@ mod tests { } #[test] - fn test_should_retry_local_decommission_resume_accepts_config_not_found_before_retry_limit() { - assert!(should_retry_local_decommission_resume(&StorageError::ConfigNotFound, 0)); + fn test_local_decommission_watchdog_retry_delay_is_bounded() { + assert_eq!(local_decommission_watchdog_retry_delay(1), LOCAL_DECOMMISSION_RESUME_RETRY_DELAY); + assert_eq!( + local_decommission_watchdog_retry_delay(u32::MAX), + LOCAL_DECOMMISSION_WATCHDOG_MAX_RETRY_DELAY + ); } - #[test] - fn test_should_retry_local_decommission_resume_rejects_config_not_found_at_retry_limit() { - assert!(!should_retry_local_decommission_resume( - &StorageError::ConfigNotFound, - LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES - )); + #[tokio::test(start_paused = true)] + async fn test_local_decommission_watchdog_retries_general_failures_until_cancelled() { + let rx = CancellationToken::new(); + let attempts = Arc::new(AtomicUsize::new(0)); + let task = tokio::spawn(run_local_decommission_watchdog(rx.clone(), { + let attempts = attempts.clone(); + let rx = rx.clone(); + move || { + let attempts = attempts.clone(); + let rx = rx.clone(); + async move { + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err(StorageError::SlowDown) + } else { + rx.cancel(); + Ok(()) + } + } + } + })); + + tokio::task::yield_now().await; + assert_eq!(attempts.load(Ordering::SeqCst), 1); + tokio::time::advance(LOCAL_DECOMMISSION_RESUME_RETRY_DELAY).await; + task.await.expect("watchdog task should exit after cancellation"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); } - #[test] - fn test_should_retry_local_decommission_resume_rejects_non_config_errors() { - assert!(!should_retry_local_decommission_resume(&StorageError::SlowDown, 0)); + #[tokio::test(start_paused = true)] + async fn test_local_decommission_watchdog_rescans_after_success() { + let rx = CancellationToken::new(); + let attempts = Arc::new(AtomicUsize::new(0)); + let task = tokio::spawn(run_local_decommission_watchdog(rx.clone(), { + let attempts = attempts.clone(); + let rx = rx.clone(); + move || { + let attempts = attempts.clone(); + let rx = rx.clone(); + async move { + if attempts.fetch_add(1, Ordering::SeqCst) == 1 { + rx.cancel(); + } + Ok(()) + } + } + })); + + tokio::task::yield_now().await; + assert_eq!(attempts.load(Ordering::SeqCst), 1); + tokio::time::advance(super::LOCAL_DECOMMISSION_WATCHDOG_INTERVAL).await; + task.await.expect("watchdog task should exit after cancellation"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); } #[test]