fix(ecstore): precheck decommission capacity (#4196)

This commit is contained in:
cxymds
2026-07-02 22:36:29 +08:00
committed by GitHub
parent cf33bcc742
commit 5404ab155b
+169 -10
View File
@@ -89,6 +89,7 @@ const DECOMMISSION_PROGRESS_SAVE_INTERVAL: Duration = Duration::seconds(30);
const DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD: usize = 1000;
const DECOMMISSION_BUCKET_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_BUCKET_CONCURRENCY";
const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4;
const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30;
pub const POOL_META_NAME: &str = "pool.bin";
pub const POOL_META_FORMAT: u16 = 1;
@@ -499,6 +500,38 @@ fn ensure_decommission_start_pool_states(meta: &PoolMeta, indices: &[usize]) ->
ensure_decommission_start_keeps_active_pool(meta, indices)
}
fn decommission_target_capacity_required(source_used: usize) -> usize {
source_used
.saturating_mul(100 + DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT)
.div_ceil(100)
}
fn ensure_decommission_start_target_capacity(
meta: &PoolMeta,
indices: &[usize],
space_infos: &[(usize, PoolSpaceInfo)],
) -> Result<()> {
let mut source_used = 0usize;
let mut target_free = 0usize;
for (idx, info) in space_infos {
if indices.contains(idx) {
source_used = source_used.saturating_add(info.used);
} else if meta.pools.get(*idx).is_some_and(is_decommission_start_active_pool) {
target_free = target_free.saturating_add(info.free);
}
}
let required = decommission_target_capacity_required(source_used);
if target_free < required {
return Err(Error::other(format!(
"failed to start decommission: insufficient target pool capacity: required {required} bytes available {target_free} bytes for {source_used} bytes used in decommission pools with {DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT}% overhead"
)));
}
Ok(())
}
fn ensure_valid_decommission_pool_index(pool_count: usize, idx: usize) -> Result<()> {
if idx >= pool_count {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
@@ -2035,7 +2068,7 @@ impl PoolDecommissionInfo {
}
}
#[derive(Debug)]
#[derive(Debug, Clone, Copy)]
pub struct PoolSpaceInfo {
pub free: usize,
pub total: usize,
@@ -2322,6 +2355,14 @@ impl ECStore {
}
}
async fn get_decommission_all_pool_space_infos(&self) -> Result<Vec<(usize, PoolSpaceInfo)>> {
let mut space_infos = Vec::with_capacity(self.pools.len());
for idx in 0..self.pools.len() {
space_infos.push((idx, self.get_decommission_pool_space_info(idx).await?));
}
Ok(space_infos)
}
#[tracing::instrument(skip(self))]
pub async fn decommission_cancel(&self, idx: usize) -> Result<()> {
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
@@ -3798,15 +3839,23 @@ impl ECStore {
}
}
let mut space_infos = Vec::with_capacity(indices.len());
for idx in indices.iter().copied() {
let pi = self.get_decommission_pool_space_info(idx).await?;
space_infos.push((idx, pi));
}
let _start_guard = self.start_gate.lock().await;
self.ensure_decommission_rebalance_idle_after_refresh().await?;
let all_space_infos = self.get_decommission_all_pool_space_infos().await?;
{
let pool_meta = self.pool_meta.read().await;
ensure_decommission_start_pool_states(&pool_meta, &indices)?;
ensure_decommission_start_target_capacity(&pool_meta, &indices, &all_space_infos)?;
}
let mut space_infos = Vec::with_capacity(indices.len());
for (idx, pi) in all_space_infos.iter().copied() {
if indices.contains(&idx) {
space_infos.push((idx, pi));
}
}
let previous_pool_meta = self
.save_current_pool_meta_for_decommission_start(&indices, space_infos, decom_buckets)
.await?;
@@ -4900,9 +4949,10 @@ mod pools_tests {
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_start_local_leader, ensure_decommission_start_pool_states,
ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_terminal_operation_supported,
ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices,
get_by_index, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
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,
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, require_decommission_store, resolve_decommission_bucket_done_save_result,
@@ -6517,6 +6567,115 @@ mod pools_tests {
assert!(err.to_string().contains("at least one active pool must remain"));
}
#[test]
fn test_ensure_decommission_start_target_capacity_allows_sufficient_free_space() {
let meta = PoolMeta {
pools: vec![decommission_test_pool_status(0, None), decommission_test_pool_status(1, None)],
..Default::default()
};
let space_infos = vec![
(
0,
PoolSpaceInfo {
free: 100,
total: 1_000,
used: 900,
},
),
(
1,
PoolSpaceInfo {
free: 1_170,
total: 2_000,
used: 830,
},
),
];
assert!(ensure_decommission_start_target_capacity(&meta, &[0], &space_infos).is_ok());
}
#[test]
fn test_ensure_decommission_start_target_capacity_rejects_insufficient_free_space() {
let meta = PoolMeta {
pools: vec![decommission_test_pool_status(0, None), decommission_test_pool_status(1, None)],
..Default::default()
};
let space_infos = vec![
(
0,
PoolSpaceInfo {
free: 100,
total: 1_000,
used: 900,
},
),
(
1,
PoolSpaceInfo {
free: 1_169,
total: 2_000,
used: 831,
},
),
];
let err = ensure_decommission_start_target_capacity(&meta, &[0], &space_infos)
.expect_err("target free capacity below 130% of source used should be rejected");
assert!(err.to_string().contains("insufficient target pool capacity"));
assert!(err.to_string().contains("required 1170 bytes available 1169 bytes"));
}
#[test]
fn test_ensure_decommission_start_target_capacity_ignores_non_active_target_pool() {
let meta = PoolMeta {
pools: vec![
decommission_test_pool_status(0, None),
decommission_test_pool_status(
1,
Some(PoolDecommissionInfo {
complete: true,
..Default::default()
}),
),
decommission_test_pool_status(2, None),
],
..Default::default()
};
let space_infos = vec![
(
0,
PoolSpaceInfo {
free: 100,
total: 1_000,
used: 900,
},
),
(
1,
PoolSpaceInfo {
free: 10_000,
total: 10_000,
used: 0,
},
),
(
2,
PoolSpaceInfo {
free: 1_169,
total: 2_000,
used: 831,
},
),
];
let err = ensure_decommission_start_target_capacity(&meta, &[0], &space_infos)
.expect_err("completed pools must not contribute target free capacity");
assert!(err.to_string().contains("required 1170 bytes available 1169 bytes"));
}
#[test]
fn test_ensure_decommission_start_pool_states_rejects_blocked_pool() {
let meta = PoolMeta {