fix(decommission): use owned data for source capacity (#8038)

This commit is contained in:
cxymds
2026-09-21 15:29:30 +08:00
committed by GitHub
parent 4047ed6c8d
commit bb43ef861a
+323 -31
View File
@@ -41,7 +41,7 @@ use crate::config::com::{
};
use crate::data_movement;
use crate::data_movement::backpressure::{self, DataMovementOperation};
use crate::data_usage::DATA_USAGE_CACHE_NAME;
use crate::data_usage::{DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, load_data_usage_cache};
use crate::disk::error::DiskError;
use crate::disk::{BUCKET_META_PREFIX, DiskAPI, RUSTFS_META_BUCKET};
use crate::error::{Error, Result};
@@ -925,6 +925,29 @@ fn capacity_target_physical_bytes(data_bytes: usize, layout: DecommissionErasure
Ok(capacity_mul_div_ceil(data_bytes, layout.width(), layout.data))
}
async fn decommission_owned_physical_usage(sets: &Sets, layout: DecommissionErasureLayout) -> HashMap<i32, usize> {
let mut usage = HashMap::with_capacity(sets.disk_set.len());
for (set_index, set) in sets.disk_set.iter().enumerate() {
let Ok(cache) = load_data_usage_cache(set.as_ref(), DATA_USAGE_CACHE_NAME).await else {
continue;
};
let data_usage = cache.dui(DATA_USAGE_ROOT, &[]);
if !data_usage.is_complete_bucket_usage_snapshot() {
continue;
}
let Ok(logical_bytes) = usize::try_from(data_usage.objects_total_size) else {
continue;
};
let Ok(set_index) = i32::try_from(set_index) else {
continue;
};
// The data-usage snapshot excludes unrelated files on a shared mount.
// Convert logical object bytes to the source set's physical EC footprint.
usage.insert(set_index, capacity_mul_div_ceil(logical_bytes, layout.width(), layout.data));
}
usage
}
fn worst_decommission_target_layout(targets: &[DecommissionPoolCapacityInfo]) -> Result<DecommissionErasureLayout> {
targets
.iter()
@@ -1861,6 +1884,132 @@ fn ensure_external_decommission_target_admission(
Err(Error::SlowDown)
}
fn grow_decommission_target_reservation(
meta: &mut PoolMeta,
capacity_infos: &[DecommissionPoolCapacityInfo],
source_pool_index: usize,
target_pool_index: usize,
expected_data_bytes: usize,
mutation_id: uuid::Uuid,
now: OffsetDateTime,
) -> Result<bool> {
let original = meta.clone();
let result = (|| {
let target_capacity = capacity_infos
.iter()
.find(|capacity| capacity.pool_index == target_pool_index)
.ok_or_else(|| {
decommission_capacity_blocked_error("target capacity snapshot is missing during reservation growth")
})?;
let pool_count = meta.pools.len();
let pool = meta
.pools
.get_mut(source_pool_index)
.ok_or_else(|| invalid_decommission_pool_index_error(pool_count, source_pool_index))?;
let reservation = pool
.decommission
.as_mut()
.and_then(|info| info.capacity_reservation.as_mut())
.filter(|reservation| reservation.active())
.ok_or_else(|| decommission_capacity_blocked_error("active reservation disappeared during reservation growth"))?;
let target_layout = reservation
.targets
.iter()
.find(|target| target.pool_index == target_pool_index)
.map(|target| target.layout)
.unwrap_or(target_capacity.layout);
let data_bytes = expected_data_bytes.max(1);
let required_target_bytes = capacity_target_physical_bytes(data_bytes, target_layout)?;
let required_peak = required_target_bytes.saturating_mul(1usize.saturating_add(reservation.temporary_copies));
let existing_target = reservation
.targets
.iter()
.find(|target| target.pool_index == target_pool_index);
if let Some(target) = existing_target
&& target.pending_physical_bytes > 0
&& target.pending_mutation_id != Some(mutation_id)
{
return Err(decommission_capacity_blocked_error(format!(
"source pool {source_pool_index} target pool {target_pool_index} has an unresolved target capacity intent"
)));
}
let existing_remaining = existing_target
.map(|target| target.remaining_reserved_physical_bytes(reservation.temporary_copies))
.unwrap_or_default();
if required_peak <= existing_remaining {
return Ok(false);
}
if reservation.source_data_equivalent_bytes > 0 {
return Err(decommission_capacity_blocked_error(format!(
"source pool {source_pool_index} reservation estimate is exhausted"
)));
}
if let Some(target) = existing_target
&& target_capacity.physical_free <= target.physical_free_at_reservation
{
return Err(decommission_capacity_blocked_error(format!(
"source pool {source_pool_index} target pool {target_pool_index} reservation is exhausted"
)));
}
let current_peak = reservation.remaining_peak_physical_bytes();
let additional_peak = required_peak.saturating_sub(existing_remaining);
let additional_predicted = additional_peak.div_ceil(1usize.saturating_add(reservation.temporary_copies));
let additional_source_data = capacity_source_data_equivalent(additional_predicted, target_layout)?;
reservation.source_data_equivalent_bytes = reservation
.source_data_equivalent_bytes
.saturating_add(additional_source_data);
reservation.source_physical_bytes = reservation
.source_physical_bytes
.saturating_add(capacity_target_physical_bytes(additional_source_data, reservation.source_layout)?);
reservation.predicted_physical_bytes = reservation
.predicted_physical_bytes
.saturating_add(additional_predicted)
.max(capacity_target_physical_bytes(reservation.source_data_equivalent_bytes, target_layout)?);
reservation.temporary_physical_bytes = reservation
.predicted_physical_bytes
.saturating_mul(reservation.temporary_copies);
reservation.peak_physical_bytes = reservation
.predicted_physical_bytes
.saturating_add(reservation.temporary_physical_bytes);
let allocation_growth = reservation.remaining_peak_physical_bytes().saturating_sub(current_peak);
if let Some(target) = reservation
.targets
.iter_mut()
.find(|target| target.pool_index == target_pool_index)
{
target.reserved_physical_bytes = target.reserved_physical_bytes.saturating_add(allocation_growth);
target.physical_total_at_reservation = target.physical_total_at_reservation.max(target_capacity.physical_total);
target.physical_free_at_reservation = target.physical_free_at_reservation.max(target_capacity.physical_free);
} else {
reservation.targets.push(DecommissionCapacityTarget {
pool_index: target_pool_index,
layout: target_layout,
physical_total_at_reservation: target_capacity.physical_total,
physical_free_at_reservation: target_capacity.physical_free,
reserved_physical_bytes: allocation_growth,
consumed_physical_bytes: 0,
observed_physical_bytes: 0,
inflight_physical_bytes: 0,
pending_physical_bytes: 0,
pending_mutation_id: None,
temporary_mutations: Vec::new(),
});
}
renew_decommission_capacity_reservation(reservation, now, true);
pool.last_update = pool.last_update.max(now);
Ok(true)
})();
if result.is_ok() {
if let Err(err) = ensure_decommission_capacity_reservations_available(meta, capacity_infos, "reservation growth") {
*meta = original;
return Err(err);
}
} else {
*meta = original;
}
result
}
fn ensure_decommission_target_owner_admission(
meta: &PoolMeta,
owner: DecommissionCapacityOwner,
@@ -11402,10 +11551,11 @@ impl ECStore {
}
})
.filter(|reservation| {
reservation
.targets
.iter()
.any(|target| target.pool_index == target_pool_index)
!temporary_release
|| reservation
.targets
.iter()
.any(|target| target.pool_index == target_pool_index)
})
.map(|reservation| (owner, reservation.model_version))
});
@@ -11464,10 +11614,11 @@ impl ECStore {
} else {
reservation.admits_owner(owner, OffsetDateTime::now_utc())
}) && reservation.model_version == model_version
&& reservation
.targets
.iter()
.any(|target| target.pool_index == target_pool_index)
&& (!temporary_release
|| reservation
.targets
.iter()
.any(|target| target.pool_index == target_pool_index))
})
.is_some();
if !owner_current {
@@ -11476,9 +11627,22 @@ impl ECStore {
));
}
let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?;
if !temporary_release {
let reservation_grew = if temporary_release {
false
} else {
let data_bytes = expected_data_bytes.unwrap_or(1).max(1);
let grew = grow_decommission_target_reservation(
&mut snapshot,
&capacity_infos,
source_pool_index,
target_pool_index,
data_bytes,
mutation_id,
OffsetDateTime::now_utc(),
)?;
ensure_decommission_capacity_reservations_available(&snapshot, &capacity_infos, "mutation")?;
}
grew
};
let target_layout = snapshot
.pools
.get(source_pool_index)
@@ -11544,7 +11708,7 @@ impl ECStore {
OffsetDateTime::now_utc(),
)?
};
if pending_added > 0 {
if reservation_grew || pending_added > 0 {
let outcome = snapshot
.save_no_lock_armed(self.pools.clone(), &mut save_guard, write_guard.lock_lost_signal(), &[source_pool_index])
.await?;
@@ -12386,7 +12550,11 @@ impl ECStore {
Ok(())
}
async fn get_decommission_pool_capacity_info(&self, idx: usize) -> Result<DecommissionPoolCapacityInfo> {
async fn get_decommission_pool_capacity_info(
&self,
idx: usize,
use_owned_usage: bool,
) -> Result<DecommissionPoolCapacityInfo> {
if let Some(sets) = self.pools.get(idx) {
let mut info = sets.storage_info_snapshot().await;
info.backend = StorageAdminApi::backend_info(self).await;
@@ -12418,8 +12586,13 @@ impl ECStore {
layout.data, layout.parity
)));
}
let owned_physical_used = if use_owned_usage {
Some(decommission_owned_physical_usage(sets, layout).await)
} else {
None
};
let (physical_total, physical_free, physical_used) =
decommission_physical_pool_capacity(&info.disks, idx, layout, space);
decommission_physical_pool_capacity(&info.disks, idx, layout, space, owned_physical_used.as_ref());
Ok(DecommissionPoolCapacityInfo {
pool_index: idx,
@@ -12442,7 +12615,20 @@ impl ECStore {
let mut capacity_infos = Vec::with_capacity(self.pools.len());
for idx in 0..self.pools.len() {
capacity_infos.push(self.get_decommission_pool_capacity_info(idx).await?);
capacity_infos.push(self.get_decommission_pool_capacity_info(idx, false).await?);
}
Ok(capacity_infos)
}
async fn get_decommission_all_pool_capacity_infos_with_owned_usage(&self) -> Result<Vec<DecommissionPoolCapacityInfo>> {
#[cfg(any(test, feature = "test-util"))]
if let Some(capacity_infos) = take_decommission_capacity_info_override_for_test(self.id) {
return Ok(capacity_infos);
}
let mut capacity_infos = Vec::with_capacity(self.pools.len());
for idx in 0..self.pools.len() {
capacity_infos.push(self.get_decommission_pool_capacity_info(idx, true).await?);
}
Ok(capacity_infos)
}
@@ -12544,15 +12730,37 @@ impl ECStore {
// different target so one mutation never holds two target gates.
discard_decommission_capacity_target_permit_except(self.id, owner, None);
}
reservation
if let Some((pool_index, _)) = reservation
.targets
.iter()
.filter_map(candidate)
.max_by_key(|(_, remaining)| *remaining)
.map(|(pool_index, _)| pool_index)
{
return Ok(pool_index);
}
if reservation.source_data_equivalent_bytes > 0 {
return Err(decommission_capacity_blocked_error(format!(
"source pool {} has no target allocation for {expected_data_bytes} data bytes",
owner.source_pool_index
)));
}
let capacities = self.get_decommission_all_pool_capacity_infos().await?;
let active_sources = active_decommission_source_indices(&pool_meta);
capacities
.into_iter()
.filter(|capacity| {
capacity.pool_index != owner.source_pool_index
&& !active_sources.contains(&capacity.pool_index)
&& pool_meta
.pools
.get(capacity.pool_index)
.is_some_and(is_decommission_start_active_pool)
})
.max_by_key(|capacity| capacity.physical_free)
.map(|capacity| capacity.pool_index)
.ok_or_else(|| {
decommission_capacity_blocked_error(format!(
"source pool {} has no target allocation for {expected_data_bytes} data bytes",
"source pool {} has no target pool for {expected_data_bytes} data bytes",
owner.source_pool_index
))
})
@@ -13048,7 +13256,7 @@ impl ECStore {
let capacity_infos = if self.pools.is_empty() {
Vec::new()
} else {
self.get_decommission_all_pool_capacity_infos().await?
self.get_decommission_all_pool_capacity_infos_with_owned_usage().await?
};
let target_fence_proof = crate::services::notification_sys::acquire_decommission_target_fence_fleet_proof();
let mut pool_meta = self.pool_meta.write().await;
@@ -15355,7 +15563,7 @@ impl ECStore {
.is_some_and(|reservation| !reservation.active())
};
if needs_capacity_reservation {
let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?;
let capacity_infos = self.get_decommission_all_pool_capacity_infos_with_owned_usage().await?;
let mut pool_meta = self.pool_meta.write().await;
let version = pool_meta.version;
pool_meta.version = POOL_META_VERSION;
@@ -16219,7 +16427,7 @@ impl ECStore {
self.ensure_decommission_rebalance_idle_after_refresh_under_start_gate()
.await?;
let all_capacity_infos = self.get_decommission_all_pool_capacity_infos().await?;
let all_capacity_infos = self.get_decommission_all_pool_capacity_infos_with_owned_usage().await?;
let target_fence_proof_available =
crate::services::notification_sys::acquire_decommission_target_fence_fleet_proof().is_some();
// Signal cancellation before waiting for the movement writer so active
@@ -22653,6 +22861,7 @@ fn decommission_physical_pool_capacity(
pool_index: usize,
layout: DecommissionErasureLayout,
logical: PoolSpaceInfo,
owned_physical_used: Option<&HashMap<i32, usize>>,
) -> (usize, usize, usize) {
let width = layout.width();
let mut sets: HashMap<i32, Vec<&rustfs_madmin::Disk>> = HashMap::new();
@@ -22692,13 +22901,20 @@ fn decommission_physical_pool_capacity(
.map(|disk| disk.available_space as usize)
.min()
.unwrap_or_default();
let max_used = set_disks
.iter()
.map(|disk| (disk.used_space as usize).max((disk.total_space as usize).saturating_sub(disk.available_space as usize)))
.max()
.unwrap_or_default();
let physical_used_for_set = owned_physical_used
.and_then(|owned| owned.get(&set_disks[0].set_index).copied())
.unwrap_or_else(|| {
set_disks
.iter()
.map(|disk| {
(disk.used_space as usize).max((disk.total_space as usize).saturating_sub(disk.available_space as usize))
})
.max()
.unwrap_or_default()
.saturating_mul(width)
});
physical_total = physical_total.saturating_add(min_total.saturating_mul(width));
physical_used = physical_used.saturating_add(max_used.saturating_mul(width));
physical_used = physical_used.saturating_add(physical_used_for_set);
if set_disks.len() >= width {
physical_free = physical_free.saturating_add(min_free.saturating_mul(width));
}
@@ -22881,6 +23097,8 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
#[cfg(test)]
mod pools_tests {
use std::collections::HashMap;
use super::DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF;
use super::persist_v3_pool_meta_for_test;
use super::record_decommission_entry_error;
@@ -22914,10 +23132,11 @@ mod pools_tests {
ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity,
ensure_decommission_terminal_operation_supported, ensure_decommission_unresolved_verification_disk_count,
ensure_local_decommission_pool_leaders, ensure_pool_meta_write_fence, ensure_valid_decommission_pool_index, get_by_index,
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_decommission_durable_ilm_receipts, merge_pool_meta_updates_for_save, merge_pool_status_refresh,
missing_decommission_worker_prefix, next_decommission_capacity_generation, observe_decommission_terminal_reload_result,
grow_decommission_target_reservation, 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_decommission_durable_ilm_receipts,
merge_pool_meta_updates_for_save, merge_pool_status_refresh, missing_decommission_worker_prefix,
next_decommission_capacity_generation, observe_decommission_terminal_reload_result,
parse_decommission_durable_ilm_receipt_path, pool_meta_has_active_decommission, publish_pool_meta_updates,
read_pool_meta_replica, reconcile_decommission_meta_buckets, reconcile_decommission_unresolved_entries_for_completion,
record_decommission_unresolved_entry, recover_decommission_capacity_reservations,
@@ -28472,6 +28691,44 @@ mod pools_tests {
}
}
#[test]
fn decommission_capacity_reservation_grows_after_stale_empty_snapshot() {
let now = OffsetDateTime::UNIX_EPOCH + Duration::minutes(1);
let layout = DecommissionErasureLayout { data: 3, parity: 1 };
let capacity_infos = vec![
DecommissionPoolCapacityInfo::for_test(0, layout, 0, 0, 0),
DecommissionPoolCapacityInfo::for_test(1, layout, 64, 64, 0),
];
let mut meta = PoolMeta {
version: POOL_META_VERSION,
pools: vec![decommission_test_pool_status(0, None), decommission_test_pool_status(1, None)],
..Default::default()
};
meta.decommission(0, capacity_infos[0].space).unwrap();
reserve_decommission_start_target_capacity(
&mut meta,
&[0],
&capacity_infos,
uuid::Uuid::new_v4(),
1,
now,
DECOMMISSION_CAPACITY_MODEL_VERSION,
)
.expect("an empty scanner snapshot should still start decommission");
let grew = grow_decommission_target_reservation(&mut meta, &capacity_infos, 0, 1, 1, uuid::Uuid::new_v4(), now)
.expect("the reservation should grow against current target capacity");
assert!(grew);
let reservation = meta.pools[0]
.decommission
.as_ref()
.and_then(|info| info.capacity_reservation.as_ref())
.unwrap();
assert_eq!(reservation.source_data_equivalent_bytes, 2);
assert_eq!(reservation.predicted_physical_bytes, 3);
assert_eq!(reservation.targets.len(), 1);
assert_eq!(reservation.targets[0].reserved_physical_bytes, 6);
}
#[test]
fn decommission_capacity_model_reserves_versions_delete_markers_and_temporary_copies_from_physical_usage() {
let reservation = build_decommission_capacity_reservation(
@@ -29028,11 +29285,46 @@ mod pools_tests {
total: 0,
used: 0,
},
None,
);
assert_eq!(capacity, (400, 40, 360));
}
#[test]
fn decommission_physical_capacity_prefers_rustfs_owned_usage_over_filesystem_usage() {
let disks = [10_u64, 20, 30, 40]
.into_iter()
.enumerate()
.map(|(disk_index, available_space)| rustfs_madmin::Disk {
endpoint: format!("http://node-{disk_index}"),
drive_path: format!("/disk-{disk_index}"),
state: "ok".to_string(),
total_space: 100,
used_space: 100 - available_space,
available_space,
pool_index: 0,
set_index: 0,
disk_index: disk_index as i32,
..Default::default()
})
.collect::<Vec<_>>();
let owned = HashMap::from([(0, 12)]);
let capacity = decommission_physical_pool_capacity(
&disks,
0,
DecommissionErasureLayout { data: 2, parity: 2 },
PoolSpaceInfo {
free: 0,
total: 0,
used: 0,
},
Some(&owned),
);
assert_eq!(capacity, (400, 40, 12));
}
#[test]
fn decommission_terminal_transitions_release_capacity_reservations() {
let now = OffsetDateTime::UNIX_EPOCH + Duration::hours(1);