From bb43ef861a09b57cab1d2e0b9a2cb37a8bc567bb Mon Sep 17 00:00:00 2001 From: cxymds Date: Mon, 21 Sep 2026 15:29:30 +0800 Subject: [PATCH] fix(decommission): use owned data for source capacity (#8038) --- crates/ecstore/src/core/pools.rs | 354 ++++++++++++++++++++++++++++--- 1 file changed, 323 insertions(+), 31 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index ad39321c6..84a1f3c62 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -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 { + 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 { 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 { + 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 { + async fn get_decommission_pool_capacity_info( + &self, + idx: usize, + use_owned_usage: bool, + ) -> Result { 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> { + #[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>, ) -> (usize, usize, usize) { let width = layout.width(); let mut sets: HashMap> = 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::>(); + 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);