diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 84a1f3c62..636b3c04e 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -135,6 +135,7 @@ const DECOMMISSION_CAPACITY_RELEASE_FAILED: &str = "failed"; const DECOMMISSION_CAPACITY_RELEASE_COMPLETED: &str = "completed"; pub(crate) const DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX: &str = "decommission/capacity-target"; const DECOMMISSION_CAPACITY_TARGET_LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250); +const DECOMMISSION_CAPACITY_TARGET_GATE_MAX_ATTEMPTS: usize = 12; const DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_PREFIX: &str = "target pool "; const DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_SUFFIX: &str = " target capacity mutation gate is busy"; const METRIC_DECOMMISSION_CAPACITY_CONFLICTS_TOTAL: &str = "rustfs_decommission_capacity_conflicts_total"; @@ -158,6 +159,10 @@ const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration:: const DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT: usize = 100; const DECOMMISSION_TERMINAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1); const DECOMMISSION_CANCEL_TARGET_LOCK_MAX_ATTEMPTS: usize = 3; + +fn decommission_capacity_target_gate_retry_exhausted(attempt: usize) -> bool { + attempt.saturating_add(1) >= DECOMMISSION_CAPACITY_TARGET_GATE_MAX_ATTEMPTS +} const DECOMMISSION_DURABLE_ILM_RECEIPT_ROOT: &str = "decommission/ilm-receipts"; const DECOMMISSION_DURABLE_ILM_MANIFEST_ROOT: &str = "decommission/ilm-manifests"; const DECOMMISSION_DURABLE_ILM_RECEIPT_SCHEMA: &str = "v2"; @@ -13957,6 +13962,20 @@ impl ECStore { match self.acquire_decommission_capacity_target_guard(target_pool_index).await { Ok(guard) => break guard, Err(err) if is_decommission_capacity_target_gate_busy(&err) => { + if decommission_capacity_target_gate_retry_exhausted(wait_attempt) { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + target_pool_index, + attempts = wait_attempt.saturating_add(1), + state = "capacity_gate_retry_exhausted", + error = %err, + "Decommission target capacity gate remained busy; pausing for supervised retry" + ); + return Err(err); + } wait_attempt = wait_attempt.saturating_add(1); } Err(err) => return Err(err), @@ -23162,8 +23181,9 @@ mod pools_tests { with_decommission_entry_context, }; use super::{ - DecommissionCapacityAdmission, DecommissionCapacityOwner, DecommissionCapacityReleaseProof, - DecommissionCapacityReservation, DecommissionCapacityTemporaryMutation, decommission_capacity_mutation_id, + DECOMMISSION_CAPACITY_TARGET_GATE_MAX_ATTEMPTS, DecommissionCapacityAdmission, DecommissionCapacityOwner, + DecommissionCapacityReleaseProof, DecommissionCapacityReservation, DecommissionCapacityTemporaryMutation, + decommission_capacity_mutation_id, decommission_capacity_target_gate_retry_exhausted, ensure_decommission_target_owner_admission, ensure_exact_delete_capacity_namespace_fences, ensure_external_decommission_target_admission, is_decommission_capacity_blocked_error, plan_exact_delete_capacity_reconciliations, record_decommission_target_consumption, release_decommission_target_inflight, @@ -23270,6 +23290,17 @@ mod pools_tests { identity: StdMutex, String)>>, } + #[test] + fn decommission_target_gate_retry_is_bounded() { + assert!(!decommission_capacity_target_gate_retry_exhausted(0)); + assert!(!decommission_capacity_target_gate_retry_exhausted( + DECOMMISSION_CAPACITY_TARGET_GATE_MAX_ATTEMPTS - 2 + )); + assert!(decommission_capacity_target_gate_retry_exhausted( + DECOMMISSION_CAPACITY_TARGET_GATE_MAX_ATTEMPTS - 1 + )); + } + #[tokio::test] async fn pool_meta_recovery_reconciles_prepare_and_partial_commit_without_format_downgrade() { for previous_version in [POOL_META_V1_VERSION, POOL_META_VERSION, POOL_META_GENERATION_VERSION] { diff --git a/crates/ecstore/src/services/rebalance/runtime.rs b/crates/ecstore/src/services/rebalance/runtime.rs index 2795c4a00..39302d592 100644 --- a/crates/ecstore/src/services/rebalance/runtime.rs +++ b/crates/ecstore/src/services/rebalance/runtime.rs @@ -7,8 +7,9 @@ use super::meta::{ validate_start_rebalance_state, }; use super::worker::{ - resolve_rebalance_bucket_result, resolve_rebalance_meta_save_result, resolve_rebalance_save_task_result, - resolve_rebalance_terminal_error, send_rebalance_done_signal, + rebalance_max_attempts, resolve_rebalance_bucket_result, resolve_rebalance_meta_save_result, + resolve_rebalance_save_task_result, resolve_rebalance_terminal_error, retry_rebalance_metadata_access, + send_rebalance_done_signal, }; use super::{ EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_STATE, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, @@ -449,13 +450,14 @@ impl ECStore { }; if terminal_state_present { - if let Err(err) = store - .save_rebalance_stats_inner( + if let Err(err) = retry_rebalance_metadata_access(None, rebalance_max_attempts(), || { + store.save_rebalance_stats_inner( pool_index, RebalSaveOpt::Stats, Some(save_rebalance_id.as_ref()), ) - .await + }) + .await { let mut rebalance_meta = store.rebalance_meta.write().await; *rebalance_meta = previous_meta; @@ -475,9 +477,10 @@ impl ECStore { } if !terminal_state_saved - && let Err(err) = store - .save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref()) - .await + && let Err(err) = retry_rebalance_metadata_access(None, rebalance_max_attempts(), || { + store.save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref()) + }) + .await { let wrapped = Error::other(format!("rebalance save_task stats save failed for pool {pool_index}: {err}")); error!("{} err: {:?}", msg, wrapped); diff --git a/crates/ecstore/src/services/rebalance/worker.rs b/crates/ecstore/src/services/rebalance/worker.rs index 7b6dc85f0..00ffa7ce3 100644 --- a/crates/ecstore/src/services/rebalance/worker.rs +++ b/crates/ecstore/src/services/rebalance/worker.rs @@ -109,7 +109,13 @@ pub(super) fn resolve_rebalance_save_task_result( } pub(super) fn resolve_rebalance_meta_save_result(result: Result<()>, stage: &str) -> Result<()> { - result.map_err(|err| Error::other(format!("rebalance meta save failed during {stage}: {err}"))) + // Keep the source error reachable: the metadata retry policy classifies + // transient lock timeouts by inspecting the source chain, so collapsing the + // failure into a plain string here would make that retry a no-op. + result.map_err(|err| { + let rendered = format!("rebalance meta save failed during {stage}: {err}"); + crate::data_movement::data_movement_context_error(rendered, err) + }) } pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError, mode: &'static str) -> Error { @@ -739,6 +745,29 @@ mod error_source_tests { } } + #[tokio::test] + async fn rebalance_metadata_retry_engages_for_wrapped_meta_save_lock_timeout() { + let mut attempts = 0; + let result = retry_rebalance_metadata_access(None, 3, || { + attempts += 1; + std::future::ready(resolve_rebalance_meta_save_result( + Err(rebalance_meta_lock_error( + rustfs_lock::LockError::timeout(".rustfs.sys/rebalance.bin@latest", Duration::from_secs(5)), + "write", + )), + "save_rebalance_stats for pool 0 opt Stats", + )) + }) + .await; + + assert_eq!(attempts, 3, "a wrapped meta save lock timeout must stay retryable"); + let err = result.expect_err("persistent lock contention must not become success"); + assert!(matches!( + rebalance_error_source(&err), + Error::Lock(rustfs_lock::LockError::Timeout { .. }) + )); + } + #[tokio::test] async fn rebalance_metadata_retry_cancels_a_pending_attempt() { struct DropProbe(Arc); diff --git a/scripts/error-other-format-baseline.txt b/scripts/error-other-format-baseline.txt index 5939a7027..6c39b9f26 100644 --- a/scripts/error-other-format-baseline.txt +++ b/scripts/error-other-format-baseline.txt @@ -50,7 +50,7 @@ 1|crates/ecstore/src/services/rebalance/entry.rs 8|crates/ecstore/src/services/rebalance/meta.rs 8|crates/ecstore/src/services/rebalance/runtime.rs -18|crates/ecstore/src/services/rebalance/worker.rs +17|crates/ecstore/src/services/rebalance/worker.rs 33|crates/ecstore/src/services/tier/tier.rs 1|crates/ecstore/src/services/tier/tier_config.rs 1|crates/ecstore/src/services/tier/warm_backend.rs