fix(ecstore): bound decommission target gate contention (#8061)

* fix(ecstore): bound decommission target gate contention

Target capacity gate contention during pool decommission escalated a
per-object transient into a durable bucket pause: each contended object
failed its bucket entry, which re-ran the whole bucket listing and amplified
attempts on the same objects.

- centralize the decommission capacity failure classification so gate
  contention, benign contention and fatal failures are decided once
- retry target-gate contention inline (bounded, jittered) before a mutation
  is admitted, covering put, part, complete, new-multipart and abort
- defer contended entries to the end of the round instead of failing the
  bucket entry, and require the deferred set to drain before a set completes
- treat a missing object or version, an overwrite and a superseded upload id
  as benign contention that is neither counted as a failure nor escalated
- use full-jitter exponential backoff, capped, for decommission retries
- expose the capacity pause reason, the waiting reason and a cumulative
  pause count in the admin pool status, plus gate-retry and per-object
  attempt metrics

Related: rustfs/backlog#2644

* fix(ecstore): correct deferred replay and metadata compatibility
This commit is contained in:
cxymds
2026-09-22 19:20:05 +08:00
committed by GitHub
parent 9c30cc8851
commit d0ce2f758b
4 changed files with 1026 additions and 88 deletions
File diff suppressed because it is too large Load Diff
+7
View File
@@ -4253,6 +4253,13 @@ mod decommission_lock_order_tests {
panic!("cancellation probe finished before observing target contention: {result:?}");
}
}
// The mutation absorbs the gate contention inline before the outer wait
// loop ever sees it, and that inline budget is bounded.
assert_eq!(
retry_observer.target_gate_inline_retries(),
crate::core::pools::DECOMMISSION_MUTATION_GATE_MAX_INLINE_ATTEMPTS - 1,
"inline target-gate retries must stop at their bounded attempt budget"
);
tokio::time::sleep(Duration::from_millis(800)).await;
assert_eq!(
retry_observer.target_gate_exact_reloads(),
@@ -92,7 +92,21 @@ When decommission metadata is present, `decommissionInfo` includes:
- progress counters: `objectsDecommissioned`, `objectsDecommissionedFailed`, `bytesDecommissioned`, and `bytesDecommissionedFailed`;
- current location: `bucket`, `prefix`, and `object`;
- queue/history lists: `queuedBuckets` and `decommissionedBuckets`;
- `waitingReason`: `queued` for queued entries and `waiting_for_worker` when metadata exists but no worker has started.
- `waitingReason`: `capacity` while the pool is paused on target capacity,
`queued` for queued entries, and `waiting_for_worker` when metadata exists but
no worker has started;
- `capacityBlockedReason`: the persisted detail for the active capacity pause,
absent once the pause clears.
A capacity pause is reported ahead of the worker states so operators can tell
"contending but progressing" from "genuinely out of target capacity". The existing `pool.bin` layout is unchanged;
cumulative pause history requires a separately versioned persistence contract.
The object-attempt metrics count entry passes and repeated version-copy attempts
within a listing and its deferred replay. Inline gate waits use the gate-retry
counter. A new listing after a durable pause starts a new per-object observation;
the maximum is the highest observation in the process, not a persisted lifetime
attempt count.
This makes queued pools and stalled metadata visible without requiring operators to inspect pool metadata files directly.
+69 -1
View File
@@ -43,6 +43,10 @@ use std::sync::Arc;
use tracing::{debug, error, info, warn};
pub type AdminUsecaseResult<T> = Result<T, ApiError>;
/// `waitingReason` value reported while a pool is paused on target capacity.
const DECOMMISSION_WAITING_REASON_CAPACITY: &str = "capacity";
pub const ADMIN_CLUSTER_SNAPSHOT_ROUTE: &str = "/rustfs/admin/v4/cluster/snapshot";
pub const ADMIN_EXTENSIONS_CATALOG_ROUTE: &str = "/rustfs/admin/v4/extensions/catalog";
pub const ADMIN_RUNTIME_CAPABILITIES_ROUTE: &str = "/rustfs/admin/v4/runtime/capabilities";
@@ -108,6 +112,8 @@ pub struct AdminPoolDecommissionInfo {
pub bytes_failed: usize,
#[serde(rename = "waitingReason")]
pub waiting_reason: Option<String>,
#[serde(rename = "capacityBlockedReason", skip_serializing_if = "Option::is_none")]
pub capacity_blocked_reason: Option<String>,
#[serde(rename = "unresolvedEntries", skip_serializing_if = "Vec::is_empty")]
pub unresolved_entries: Vec<DecommissionUnresolvedEntry>,
}
@@ -622,12 +628,23 @@ impl DefaultAdminUsecase {
bytes_done: info.bytes_done,
bytes_failed: info.bytes_failed,
waiting_reason,
capacity_blocked_reason: info.capacity_blocked_reason,
unresolved_entries: info.unresolved_entries,
}
}
/// Why a pool is not currently making progress.
///
/// A durable capacity pause is reported ahead of the worker states: it is
/// the actionable condition, and `capacityBlockedReason` carries the detail.
fn decommission_waiting_reason(info: &PoolDecommissionInfo) -> Option<&'static str> {
if !info.has_decommission_state() || info.complete || info.failed || info.canceled || info.start_time.is_some() {
if !info.has_decommission_state() || info.complete || info.failed || info.canceled {
return None;
}
if info.capacity_blocked_reason.is_some() {
return Some(DECOMMISSION_WAITING_REASON_CAPACITY);
}
if info.start_time.is_some() {
return None;
}
if info.queued {
@@ -995,6 +1012,57 @@ mod tests {
assert_eq!(item.rebalance_status, "started");
}
/// A durable capacity pause must be distinguishable from "slow but
/// progressing": the pause surfaces its own waiting reason, the persisted
/// detail.
#[test]
fn admin_pool_list_item_exposes_decommission_capacity_pause() {
let item = DefaultAdminUsecase::pool_list_item_from_status(
PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
total_size: 1_000,
current_size: 500,
capacity_blocked_reason: Some("target pool 1 target capacity mutation gate is busy".to_string()),
..Default::default()
}),
},
(RebalStatus::None, false),
);
let value = serde_json::to_value(item).expect("paused pool status should serialize");
assert_eq!(value["decommissionInfo"]["waitingReason"], "capacity");
assert_eq!(
value["decommissionInfo"]["capacityBlockedReason"],
"target pool 1 target capacity mutation gate is busy"
);
}
/// Once the pause clears, `waitingReason` must return to the worker states
/// without retaining a stale capacity reason.
#[test]
fn admin_pool_list_item_clears_capacity_pause() {
let item = DefaultAdminUsecase::pool_list_item_from_status(
PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
}),
},
(RebalStatus::None, false),
);
let value = serde_json::to_value(item).expect("resumed pool status should serialize");
assert!(value["decommissionInfo"]["waitingReason"].is_null());
assert!(value["decommissionInfo"].get("capacityBlockedReason").is_none());
}
#[test]
fn admin_pool_list_item_exposes_queued_decommission_state() {
let item = DefaultAdminUsecase::pool_list_item_from_status(