From c14a442586ddcbfdfa2b7305cfc19df307e14811 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 24 Jun 2026 07:59:39 +0800 Subject: [PATCH] fix(storage): harden rebalance and decommission state (#3730) --- .../e2e_test/src/reliant/grpc_lock_server.rs | 21 + crates/ecstore/src/notification_sys.rs | 26 + crates/ecstore/src/pools.rs | 952 +++++++++++++++--- crates/ecstore/src/rebalance/control.rs | 42 +- crates/ecstore/src/rebalance/meta.rs | 37 +- .../src/rebalance/rebalance_unit_tests.rs | 237 ++++- crates/ecstore/src/rebalance/runtime.rs | 21 +- crates/ecstore/src/rpc/peer_rest_client.rs | 86 +- .../src/generated/proto_gen/node_service.rs | 177 ++++ crates/protos/src/node.proto | 30 + .../decommission-compatibility.md | 6 + rustfs/src/admin/handlers/pools.rs | 312 +++++- rustfs/src/admin/handlers/rebalance.rs | 14 +- rustfs/src/admin/route_policy.rs | 6 + rustfs/src/admin/route_registration_test.rs | 2 + rustfs/src/app/admin_usecase.rs | 358 +++++-- rustfs/src/app/mod.rs | 1 + rustfs/src/storage/rpc/node_service.rs | 123 ++- 18 files changed, 2153 insertions(+), 298 deletions(-) diff --git a/crates/e2e_test/src/reliant/grpc_lock_server.rs b/crates/e2e_test/src/reliant/grpc_lock_server.rs index 085d51fe0..01849428f 100644 --- a/crates/e2e_test/src/reliant/grpc_lock_server.rs +++ b/crates/e2e_test/src/reliant/grpc_lock_server.rs @@ -664,6 +664,27 @@ impl NodeService for MinimalLockNodeService { Err(Status::unimplemented("lock-only test server")) } + async fn start_decommission( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("lock-only test server")) + } + + async fn cancel_decommission( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("lock-only test server")) + } + + async fn clear_decommission( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("lock-only test server")) + } + async fn get_metrics( &self, _request: Request, diff --git a/crates/ecstore/src/notification_sys.rs b/crates/ecstore/src/notification_sys.rs index 7d9e0c735..e65174300 100644 --- a/crates/ecstore/src/notification_sys.rs +++ b/crates/ecstore/src/notification_sys.rs @@ -109,6 +109,14 @@ impl NotificationSys { self.all_peer_clients[idx].clone() } + pub fn peer_client_for_grid_host(&self, grid_host: &str) -> Option { + self.all_peer_clients + .iter() + .flatten() + .find(|client| client.grid_host == grid_host) + .cloned() + } + pub async fn delete_policy(&self, policy_name: &str) -> Vec { let mut futures = Vec::with_capacity(self.peer_clients.len()); for client in self.peer_clients.iter() { @@ -1215,6 +1223,24 @@ mod tests { assert!(msg.contains("local save failed")); } + #[test] + fn peer_client_for_grid_host_matches_exact_grid_host() { + let sys = NotificationSys { + peer_clients: Vec::new(), + all_peer_clients: vec![Some(PeerRestClient::new( + "127.0.0.1:9000".to_string().try_into().expect("peer host should parse"), + "http://127.0.0.1:9000".to_string(), + ))], + peer_admin_caches: Vec::new(), + }; + + let client = sys + .peer_client_for_grid_host("http://127.0.0.1:9000") + .expect("matching grid host should return peer client"); + assert_eq!(client.grid_host, "http://127.0.0.1:9000"); + assert!(sys.peer_client_for_grid_host("http://node-b:9000").is_none()); + } + #[test] fn load_rebalance_meta_aggregate_failures_return_error() { let err = aggregate_notification_failures( diff --git a/crates/ecstore/src/pools.rs b/crates/ecstore/src/pools.rs index cffc5470b..6ebad8d86 100644 --- a/crates/ecstore/src/pools.rs +++ b/crates/ecstore/src/pools.rs @@ -25,7 +25,7 @@ use crate::bucket::{ object_lock::objectlock_sys::BucketObjectLockSys, }; use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw}; -use crate::config::com::{CONFIG_PREFIX, read_config, save_config}; +use crate::config::com::{CONFIG_PREFIX, read_config, read_config_no_lock, save_config, save_config_with_opts}; use crate::data_movement; use crate::data_usage::DATA_USAGE_CACHE_NAME; use crate::disk::error::DiskError; @@ -33,8 +33,7 @@ use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; use crate::endpoints::EndpointServerPools; use crate::error::{Error, Result}; use crate::error::{ - StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_operation_canceled, - is_err_version_not_found, + StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, }; use crate::global::resolve_object_store_handle; use crate::object_api::{GetObjectReader, ObjectOptions}; @@ -78,6 +77,10 @@ const LOG_SUBSYSTEM_POOLS: &str = "pools"; const EVENT_DECOMMISSION_STATE: &str = "decommission_state"; const EVENT_DECOMMISSION_BUCKET: &str = "decommission_bucket"; const EVENT_DECOMMISSION_ENTRY: &str = "decommission_entry"; +const DECOMMISSION_STAGE_MIGRATE_OBJECT: &str = "migrate_object"; +const DECOMMISSION_STAGE_CLEANUP_PREFLIGHT: &str = "cleanup_preflight"; +const DECOMMISSION_STAGE_SOURCE_CLEANUP: &str = "source_cleanup"; +const DECOMMISSION_STAGE_ENTRY_FINISHED: &str = "entry_finished"; 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"; @@ -420,22 +423,72 @@ fn invalid_decommission_pool_index_error(pool_count: usize, idx: usize) -> Error Error::other(format!("invalid decommission pool index {idx} for {pool_count} pools")) } -fn ensure_decommission_start_allowed(pool_present: bool, decommission_active: bool, decommission_complete: bool) -> Result<()> { - if !pool_present { - return Err(Error::other("failed to start decommission: target pool was not found")); - } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum DecommissionStartPoolState { + Missing, + Active, + Decommissioning, + Decommissioned, + Blocked, +} - if decommission_active { - return Err(StorageError::DecommissionAlreadyRunning); - } +fn decommission_start_pool_state(pool: Option<&PoolStatus>) -> DecommissionStartPoolState { + let Some(pool) = pool else { + return DecommissionStartPoolState::Missing; + }; + let Some(info) = pool.decommission.as_ref() else { + return DecommissionStartPoolState::Active; + }; - if decommission_complete { - return Err(Error::other("failed to start decommission: target pool decommission is already complete")); + if info.complete { + DecommissionStartPoolState::Decommissioned + } else if info.failed || info.canceled { + DecommissionStartPoolState::Blocked + } else { + DecommissionStartPoolState::Decommissioning + } +} + +fn is_decommission_start_active_pool(pool: &PoolStatus) -> bool { + decommission_start_pool_state(Some(pool)) == DecommissionStartPoolState::Active +} + +fn ensure_decommission_start_allowed(state: DecommissionStartPoolState) -> Result<()> { + match state { + DecommissionStartPoolState::Missing => Err(Error::other("failed to start decommission: target pool was not found")), + DecommissionStartPoolState::Active => Ok(()), + DecommissionStartPoolState::Decommissioning => Err(StorageError::DecommissionAlreadyRunning), + DecommissionStartPoolState::Decommissioned => { + Err(Error::other("failed to start decommission: target pool is already decommissioned")) + } + DecommissionStartPoolState::Blocked => Err(Error::other( + "failed to start decommission: target pool decommission is blocked; clear failed or canceled metadata before starting again", + )), + } +} + +fn ensure_decommission_start_keeps_active_pool(meta: &PoolMeta, indices: &[usize]) -> Result<()> { + let active_count = meta + .pools + .iter() + .filter(|pool| is_decommission_start_active_pool(pool)) + .count(); + if active_count <= indices.len() { + return Err(Error::other( + "failed to start decommission: at least one active pool must remain after decommission start", + )); } Ok(()) } +fn ensure_decommission_start_pool_states(meta: &PoolMeta, indices: &[usize]) -> Result<()> { + for idx in indices.iter().copied() { + ensure_decommission_start_allowed(decommission_start_pool_state(meta.pools.get(idx)))?; + } + ensure_decommission_start_keeps_active_pool(meta, indices) +} + 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)); @@ -507,7 +560,13 @@ fn count_decommission_item(meta: &mut PoolMeta, idx: usize, size: usize, failed: Ok(()) } -fn track_decommission_current_object(meta: &mut PoolMeta, idx: usize, bucket: &str, object: &str) -> Result<()> { +fn track_decommission_current_object_stage( + meta: &mut PoolMeta, + idx: usize, + bucket: &str, + object: &str, + stage: &str, +) -> Result<()> { let pool_count = meta.pools.len(); ensure_valid_decommission_pool_index(pool_count, idx)?; @@ -520,6 +579,27 @@ fn track_decommission_current_object(meta: &mut PoolMeta, idx: usize, bucket: &s info.object = object.to_string(); info.bucket = bucket.to_string(); + info.stage = stage.to_string(); + Ok(()) +} + +fn track_decommission_current_object(meta: &mut PoolMeta, idx: usize, bucket: &str, object: &str) -> Result<()> { + track_decommission_current_object_stage(meta, idx, bucket, object, "") +} + +fn touch_decommission_progress(meta: &mut PoolMeta, idx: usize) -> Result<()> { + let pool_count = meta.pools.len(); + ensure_valid_decommission_pool_index(pool_count, idx)?; + + let Some(pool) = meta.pools.get_mut(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let Some(info) = pool.decommission.as_mut() else { + return Err(decommission_metadata_not_initialized_error("touch decommission progress")); + }; + + pool.last_update = OffsetDateTime::now_utc(); + info.mark_progress_saved(); Ok(()) } @@ -610,6 +690,37 @@ fn load_decommission_entry_versions(entry: &MetaCacheEntry, bucket: &str, stage: .map_err(|err| with_decommission_entry_context(stage, bucket, &entry.name, err)) } +fn empty_decommission_entry_versions(bucket: &str, object: &str) -> FileInfoVersions { + FileInfoVersions { + volume: bucket.to_string(), + name: object.to_string(), + versions: Vec::new(), + ..Default::default() + } +} + +fn resolve_decommission_entry_exact_versions( + result: Result>, + entry: &MetaCacheEntry, + bucket: &str, + stage: &str, +) -> Result { + match result { + Ok(Some(fivs)) => Ok(fivs), + Ok(None) => Ok(empty_decommission_entry_versions(bucket, &entry.name)), + Err(err) => Err(with_decommission_entry_context(stage, bucket, &entry.name, err)), + } +} + +async fn load_decommission_entry_exact_versions( + set: &SetDisks, + entry: &MetaCacheEntry, + bucket: &str, + stage: &str, +) -> Result { + resolve_decommission_entry_exact_versions(set.load_file_info_versions_exact(bucket, &entry.name).await, entry, bucket, stage) +} + fn resolve_decommission_check_after_list_result(list_result: Result<()>, entry_error: Option) -> Result<()> { if let Some(err) = entry_error { Err(err) } else { list_result } } @@ -618,6 +729,24 @@ fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) result.map_err(|err| Error::other(format!("decommission pool meta reload failed during {stage}: {err}"))) } +fn apply_decommission_status_space_info(mut pool_info: PoolStatus, space_info: PoolSpaceInfo) -> PoolStatus { + match pool_info.decommission.as_mut() { + Some(d) => { + d.total_size = space_info.total; + d.current_size = space_info.free; + } + None => { + pool_info.decommission = Some(PoolDecommissionInfo { + total_size: space_info.total, + current_size: space_info.free, + ..Default::default() + }); + } + } + + pool_info +} + fn resolve_start_decommission_pool_meta_reload_result(result: Result<()>) -> Result<()> { resolve_decommission_pool_meta_reload_result(result, "start_decommission") } @@ -689,21 +818,19 @@ fn is_decommission_copy_cleanup_safe_error(err: &Error) -> bool { is_err_object_not_found(err) || is_err_version_not_found(err) } -fn should_cleanup_decommission_source_entry(decommissioned: usize, total_versions: usize, expired: usize) -> bool { - decommissioned.saturating_add(expired) == total_versions +fn is_decommission_target_capacity_error(err: &Error) -> bool { + if matches!(err, Error::DiskFull | Error::StorageFull) { + return true; + } + + let message = err.to_string(); + let disk_full = Error::DiskFull.to_string(); + let storage_full = Error::StorageFull.to_string(); + message.contains(&disk_full) || message.contains(&storage_full) } -fn decommission_start_guard_state(pool: Option<&PoolStatus>) -> (bool, bool, bool) { - if let Some(pool) = pool { - let active = pool - .decommission - .as_ref() - .is_some_and(|info| is_decommission_active(info.complete, info.failed, info.canceled)); - let complete = pool.decommission.as_ref().is_some_and(|info| info.complete); - (true, active, complete) - } else { - (false, false, false) - } +fn should_cleanup_decommission_source_entry(decommissioned: usize, total_versions: usize, expired: usize) -> bool { + decommissioned.saturating_add(expired) == total_versions } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -720,8 +847,8 @@ fn classify_decommission_terminal_state(failed_items_present: bool) -> Decommiss } } -fn should_preserve_decommission_canceled_state(meta_canceled: bool, cancel_signal: bool) -> bool { - meta_canceled || cancel_signal +fn should_preserve_decommission_canceled_state(meta_canceled: bool, _cancel_signal: bool) -> bool { + meta_canceled } fn should_continue_decommission_queue(meta: &PoolMeta, idx: usize) -> bool { @@ -970,6 +1097,10 @@ struct PersistedPoolDecommissionInfo { pub bytes_done: usize, #[serde(rename = "bytesDecommissionedFailed")] pub bytes_failed: usize, + #[serde(rename = "terminalReloadAttemptAt", with = "time::serde::rfc3339::option", default)] + pub terminal_reload_attempt_at: Option, + #[serde(rename = "terminalReloadFailures", default)] + pub terminal_reload_failures: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -1094,10 +1225,13 @@ impl TryFrom for PoolDecommissionInfo { bucket: value.bucket, prefix: value.prefix, object: value.object, + stage: String::new(), items_decommissioned: value.items_decommissioned, items_decommission_failed: value.items_decommission_failed, bytes_done: value.bytes_done, bytes_failed: value.bytes_failed, + terminal_reload_attempt_at: value.terminal_reload_attempt_at, + terminal_reload_failures: value.terminal_reload_failures, progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed), }) } @@ -1122,10 +1256,13 @@ impl TryFrom for PoolDecommissionInfo { bucket: String::new(), prefix: String::new(), object: String::new(), + stage: String::new(), items_decommissioned: value.items_decommissioned, items_decommission_failed: value.items_decommission_failed, bytes_done: value.bytes_done, bytes_failed: value.bytes_failed, + terminal_reload_attempt_at: None, + terminal_reload_failures: Vec::new(), progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed), }) } @@ -1171,6 +1308,8 @@ impl From<&PoolDecommissionInfo> for PersistedPoolDecommissionInfo { items_decommission_failed: value.items_decommission_failed, bytes_done: value.bytes_done, bytes_failed: value.bytes_failed, + terminal_reload_attempt_at: value.terminal_reload_attempt_at, + terminal_reload_failures: value.terminal_reload_failures.clone(), } } } @@ -1238,23 +1377,13 @@ impl PoolMeta { } } - pub async fn load(&mut self, pool: Arc, _pools: Vec>) -> Result<()> { - let data = match read_config(pool, POOL_META_NAME).await { - Ok(data) => { - if data.is_empty() { - return Ok(()); - } else if data.len() <= 4 { - return Err(Error::other("pool metadata load failed: metadata payload is too short")); - } - data - } - Err(err) => { - if err == Error::ConfigNotFound { - return Ok(()); - } - return Err(err); - } - }; + fn load_from_config_data(&mut self, data: Vec) -> Result<()> { + if data.is_empty() { + return Ok(()); + } else if data.len() <= 4 { + return Err(Error::other("pool metadata load failed: metadata payload is too short")); + } + let format = LittleEndian::read_u16(&data[0..2]); if format != POOL_META_FORMAT { return Err(Error::other(format!("pool metadata load failed: unknown format {format}"))); @@ -1275,9 +1404,35 @@ impl PoolMeta { Ok(()) } - pub async fn save(&self, pools: Vec>) -> Result<()> { + pub async fn load(&mut self, pool: Arc, _pools: Vec>) -> Result<()> { + let data = match read_config(pool, POOL_META_NAME).await { + Ok(data) => data, + Err(err) => { + if err == Error::ConfigNotFound { + return Ok(()); + } + return Err(err); + } + }; + self.load_from_config_data(data) + } + + async fn load_no_lock(&mut self, pool: Arc) -> Result<()> { + let data = match read_config_no_lock(pool, POOL_META_NAME).await { + Ok(data) => data, + Err(err) => { + if err == Error::ConfigNotFound { + return Ok(()); + } + return Err(err); + } + }; + self.load_from_config_data(data) + } + + fn encode_config_data(&self) -> Result> { if self.dont_save { - return Ok(()); + return Ok(Vec::new()); } let mut data = Vec::new(); data.write_u16::(POOL_META_FORMAT)?; @@ -1285,7 +1440,14 @@ impl PoolMeta { let mut buf = Vec::new(); PersistedPoolMeta::from(self).serialize(&mut Serializer::new(&mut buf))?; data.write_all(&buf)?; + Ok(data) + } + pub async fn save(&self, pools: Vec>) -> Result<()> { + let data = self.encode_config_data()?; + if data.is_empty() { + return Ok(()); + } for pool in pools { save_config(pool, POOL_META_NAME, data.clone()).await?; } @@ -1293,6 +1455,28 @@ impl PoolMeta { Ok(()) } + async fn save_no_lock(&self, pools: Vec>) -> Result<()> { + let data = self.encode_config_data()?; + if data.is_empty() { + return Ok(()); + } + for pool in pools { + save_config_with_opts( + pool, + POOL_META_NAME, + data.clone(), + &ObjectOptions { + max_parity: true, + no_lock: true, + ..Default::default() + }, + ) + .await?; + } + + Ok(()) + } + pub fn decommission_cancel(&mut self, idx: usize) -> bool { if let Some(stats) = self.pools.get_mut(idx) { if let Some(d) = &stats.decommission { @@ -1304,6 +1488,8 @@ impl PoolMeta { pd.failed = false; pd.complete = false; pd.start_time = None; + pd.terminal_reload_attempt_at = None; + pd.terminal_reload_failures.clear(); stats.decommission = Some(pd); true @@ -1328,6 +1514,8 @@ impl PoolMeta { pd.failed = true; pd.complete = false; pd.start_time = None; + pd.terminal_reload_attempt_at = None; + pd.terminal_reload_failures.clear(); stats.decommission = Some(pd); true @@ -1373,6 +1561,8 @@ impl PoolMeta { pd.canceled = false; pd.failed = false; pd.complete = true; + pd.terminal_reload_attempt_at = None; + pd.terminal_reload_failures.clear(); stats.decommission = Some(pd); true @@ -1394,12 +1584,7 @@ impl PoolMeta { return Err(invalid_decommission_pool_index_error(pool_count, idx)); }; - let decommission_active = pool - .decommission - .as_ref() - .is_some_and(|info| is_decommission_active(info.complete, info.failed, info.canceled)); - let decommission_complete = pool.decommission.as_ref().is_some_and(|info| info.complete); - ensure_decommission_start_allowed(true, decommission_active, decommission_complete)?; + ensure_decommission_start_allowed(decommission_start_pool_state(Some(pool)))?; let previous = pool.decommission.as_ref(); let now = OffsetDateTime::now_utc(); @@ -1417,6 +1602,28 @@ impl PoolMeta { self.set_decommission_state(idx, pi, true) } + pub fn record_decommission_terminal_reload_failure(&mut self, idx: usize, stage: &str, message: String) -> Result { + let pool_count = self.pools.len(); + ensure_valid_decommission_pool_index(pool_count, idx)?; + + let Some(pool) = self.pools.get_mut(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let Some(info) = pool.decommission.as_mut() else { + return Err(decommission_metadata_not_initialized_error("record decommission terminal reload failure")); + }; + + let failure = format!("{stage}: {message}"); + if info.terminal_reload_failures.last().is_some_and(|last| last == &failure) { + return Ok(false); + } + + pool.last_update = OffsetDateTime::now_utc(); + info.terminal_reload_attempt_at = Some(pool.last_update); + info.terminal_reload_failures.push(failure); + Ok(true) + } + pub fn promote_queued_decommission(&mut self, idx: usize) -> bool { if let Some(pool) = self.pools.get_mut(idx) && let Some(info) = pool.decommission.as_mut() @@ -1490,6 +1697,10 @@ impl PoolMeta { } pub fn track_current_bucket_object(&mut self, idx: usize, bucket: String, object: String) { + self.track_current_bucket_object_stage(idx, bucket, object, String::new()); + } + + pub fn track_current_bucket_object_stage(&mut self, idx: usize, bucket: String, object: String, stage: String) { if self.pools.get(idx).is_none_or(|v| v.decommission.is_none()) { return; } @@ -1499,6 +1710,7 @@ impl PoolMeta { { info.object = object; info.bucket = bucket; + info.stage = stage; } } @@ -1647,6 +1859,8 @@ pub struct PoolDecommissionInfo { pub prefix: String, #[serde(skip)] pub object: String, + #[serde(skip)] + pub stage: String, #[serde(rename = "objectsDecommissioned")] pub items_decommissioned: usize, @@ -1656,6 +1870,10 @@ pub struct PoolDecommissionInfo { pub bytes_done: usize, #[serde(rename = "bytesDecommissionedFailed")] pub bytes_failed: usize, + #[serde(rename = "terminalReloadAttemptAt", with = "time::serde::rfc3339::option", default)] + pub terminal_reload_attempt_at: Option, + #[serde(rename = "terminalReloadFailures", default)] + pub terminal_reload_failures: Vec, #[serde(skip)] pub progress_save_item_baseline: usize, } @@ -1929,16 +2147,12 @@ impl ECStore { pool_meta.clone() }; let mut latest_pool_meta = PoolMeta::default(); - latest_pool_meta.load(rebalance_pool, self.pools.clone()).await?; + latest_pool_meta.load_no_lock(rebalance_pool).await?; if latest_pool_meta.pools.is_empty() { latest_pool_meta = current_pool_meta; } - for idx in indices.iter().copied() { - let (pool_present, decommission_active, decommission_complete) = - decommission_start_guard_state(latest_pool_meta.pools.get(idx)); - ensure_decommission_start_allowed(pool_present, decommission_active, decommission_complete)?; - } + ensure_decommission_start_pool_states(&latest_pool_meta, indices)?; let previous_pool_meta = latest_pool_meta.clone(); let first_idx = indices.first().copied(); @@ -1951,7 +2165,7 @@ impl ECStore { latest_pool_meta.queue_buckets(idx, decom_buckets.clone()); } - latest_pool_meta.save(self.pools.clone()).await?; + latest_pool_meta.save_no_lock(self.pools.clone()).await?; { let mut pool_meta = self.pool_meta.write().await; *pool_meta = latest_pool_meta; @@ -1960,24 +2174,18 @@ impl ECStore { Ok(previous_pool_meta) } + async fn ensure_decommission_rebalance_idle_after_refresh(&self) -> Result<()> { + self.load_rebalance_meta().await?; + ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await) + } + pub async fn status(&self, idx: usize) -> Result { let space_info = self.get_decommission_pool_space_info(idx).await?; let pool_meta = self.pool_meta.read().await; - let mut pool_info = get_by_index(pool_meta.pools.as_slice(), idx, "fetch decommission status")?.clone(); - if let Some(d) = pool_info.decommission.as_mut() { - d.total_size = space_info.total; - d.current_size = space_info.free; - } else { - pool_info.decommission = Some(PoolDecommissionInfo { - total_size: space_info.total, - current_size: space_info.free, - ..Default::default() - }); - } - - Ok(pool_info) + let pool_info = get_by_index(pool_meta.pools.as_slice(), idx, "fetch decommission status")?.clone(); + Ok(apply_decommission_status_space_info(pool_info, space_info)) } async fn get_decommission_pool_space_info(&self, idx: usize) -> Result { @@ -2108,6 +2316,19 @@ impl ECStore { Ok(()) } + async fn record_decommission_terminal_reload_failure(&self, idx: usize, stage: &str, err: Error) -> Result<()> { + let changed = { + let mut pool_meta = self.pool_meta.write().await; + pool_meta.record_decommission_terminal_reload_failure(idx, stage, err.to_string())? + }; + + if changed { + self.save_current_pool_meta().await?; + } + + Ok(()) + } + pub async fn is_decommission_running(&self) -> bool { { let cancelers = self.decommission_cancelers.read().await; @@ -2197,7 +2418,7 @@ impl ECStore { ); validate_start_decommission_request(&indices, self.single_pool())?; - ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)?; + self.ensure_decommission_rebalance_idle_after_refresh().await?; let store = require_decommission_store(resolve_object_store_handle(), "start decommission")?; let local_indices = local_decommission_queue_prefix(&self.endpoints(), &indices)?; @@ -2227,6 +2448,38 @@ impl ECStore { Ok(()) } + async fn save_decommission_entry_progress_stage( + &self, + idx: usize, + bucket: &str, + object: &str, + stage: &'static str, + ) -> Result<()> { + { + let mut pool_meta = self.pool_meta.write().await; + track_decommission_current_object_stage(&mut pool_meta, idx, bucket, object, stage) + .map_err(|err| with_decommission_entry_context(stage, bucket, object, err))?; + touch_decommission_progress(&mut pool_meta, idx) + .map_err(|err| with_decommission_entry_context(stage, bucket, object, err))?; + } + + if let Some(err) = resolve_decommission_progress_save_result(self.save_current_pool_meta().await) { + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + bucket = %bucket, + object = %object, + stage, + error = ?err, + "Decommission progress stage save failed" + ); + } + + Ok(()) + } + #[allow(unused_assignments, clippy::too_many_arguments)] #[tracing::instrument(skip(self, set, _worker_permit, lifecycle_config, lock_retention, replication_config))] async fn decommission_entry( @@ -2269,7 +2522,7 @@ impl ECStore { } decommission_cancel_signal_result(rx.is_cancelled())?; - let mut fivs = load_decommission_entry_versions(&entry, &bucket, "file_info_versions")?; + let mut fivs = load_decommission_entry_exact_versions(&set, &entry, &bucket, "file_info_versions").await?; fivs.versions .sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time))); @@ -2350,6 +2603,15 @@ impl ECStore { ignore = true; cleanup_ignored = true; } else { + if is_decommission_target_capacity_error(&err) { + return Err(with_decommission_entry_context( + "delete_marker_copy", + bucket.as_str(), + version.name.as_str(), + err, + )); + } + failure = true; error = Some(err) @@ -2421,6 +2683,15 @@ impl ECStore { break; } + if is_decommission_target_capacity_error(&err) { + return Err(with_decommission_entry_context( + "decommission_tiered_object", + bucket.as_str(), + version.name.as_str(), + err, + )); + } + failure = true; error!("decommission_pool: decommission_tiered_object err {:?}", &err); error = Some(err); @@ -2466,6 +2737,14 @@ impl ECStore { let bucket_name = bucket.clone(); let object_name = rd.object_info.name.clone(); + self.save_decommission_entry_progress_stage( + idx, + bucket_name.as_str(), + object_name.as_str(), + DECOMMISSION_STAGE_MIGRATE_OBJECT, + ) + .await?; + if let Err(err) = self.clone().decommission_object(idx, bucket, rd).await { if is_decommission_copy_cleanup_safe_error(&err) { ignore = true; @@ -2473,6 +2752,15 @@ impl ECStore { break; } + if is_decommission_target_capacity_error(&err) { + return Err(with_decommission_entry_context( + DECOMMISSION_STAGE_MIGRATE_OBJECT, + bucket_name.as_str(), + object_name.as_str(), + err, + )); + } + failure = true; error!("decommission_pool: decommission_object err {:?}", &err); @@ -2536,6 +2824,14 @@ impl ECStore { if should_cleanup_decommission_source_entry(decommissioned, fivs.versions.len(), expired) { decommission_cancel_signal_result(rx.is_cancelled())?; + self.save_decommission_entry_progress_stage( + idx, + bucket.as_str(), + entry.name.as_str(), + DECOMMISSION_STAGE_CLEANUP_PREFLIGHT, + ) + .await?; + data_movement::ensure_source_cleanup_versions_unchanged( set.clone(), bucket.as_str(), @@ -2547,6 +2843,14 @@ impl ECStore { .await .map_err(|err| with_decommission_entry_context("cleanup_preflight", bucket.as_str(), entry.name.as_str(), err))?; + self.save_decommission_entry_progress_stage( + idx, + bucket.as_str(), + entry.name.as_str(), + DECOMMISSION_STAGE_SOURCE_CLEANUP, + ) + .await?; + let cleanup_result = set .delete_object( bucket.as_str(), @@ -2596,6 +2900,9 @@ impl ECStore { } }; + self.save_decommission_entry_progress_stage(idx, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_ENTRY_FINISHED) + .await?; + if should_save_progress { let save_result = self.save_current_pool_meta().await; if let Some(err) = resolve_decommission_progress_save_result(save_result) { @@ -2974,7 +3281,7 @@ impl ECStore { "Decommission background routine failed" ); - if is_err_operation_canceled(&err) || should_preserve_decommission_canceled_state(canceled, rx.is_cancelled()) { + if should_preserve_decommission_canceled_state(canceled, rx.is_cancelled()) { warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -3111,6 +3418,21 @@ impl ECStore { resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), stage.as_str(), ) { + if let Err(record_err) = self + .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) + .await + { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_reload_record_failed", + error = %record_err, + original_error = %err, + "Decommission terminal reload failure record failed" + ); + } warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -3161,6 +3483,21 @@ impl ECStore { resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str()), stage.as_str(), ) { + if let Err(record_err) = self + .record_decommission_terminal_reload_failure(idx, stage.as_str(), err.clone()) + .await + { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "terminal_reload_record_failed", + error = %record_err, + original_error = %err, + "Decommission terminal reload failure record failed" + ); + } warn!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -3287,7 +3624,7 @@ impl ECStore { let indices = dedup_indices(&indices); validate_start_decommission_request(&indices, self.single_pool())?; - ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)?; + self.ensure_decommission_rebalance_idle_after_refresh().await?; ensure_decommission_start_local_leader(&self.endpoints(), &indices)?; for idx in indices.iter().copied() { @@ -3296,11 +3633,7 @@ impl ECStore { { let pool_meta = self.pool_meta.read().await; - for idx in indices.iter().copied() { - let (pool_present, decommission_active, decommission_complete) = - decommission_start_guard_state(pool_meta.pools.get(idx)); - ensure_decommission_start_allowed(pool_present, decommission_active, decommission_complete)?; - } + ensure_decommission_start_pool_states(&pool_meta, &indices)?; } let decom_buckets = self.get_buckets_to_decommission().await?; @@ -3333,7 +3666,7 @@ impl ECStore { } let _start_guard = self.start_gate.lock().await; - ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)?; + self.ensure_decommission_rebalance_idle_after_refresh().await?; let previous_pool_meta = self .save_current_pool_meta_for_decommission_start(&indices, space_infos, decom_buckets) @@ -3658,6 +3991,29 @@ mod tests { assert!(!is_decommission_copy_cleanup_safe_error(&err)); } + #[test] + fn decommission_target_capacity_error_accepts_direct_capacity_errors() { + assert!(is_decommission_target_capacity_error(&Error::DiskFull)); + assert!(is_decommission_target_capacity_error(&Error::StorageFull)); + } + + #[test] + fn decommission_target_capacity_error_accepts_wrapped_capacity_errors() { + let disk_full = Error::other(format!("decommission_object: put_object failed for bucket/object: {}", Error::DiskFull)); + let storage_full = Error::other(format!( + "decommission_object: put_object failed for bucket/object: {}", + Error::StorageFull + )); + + assert!(is_decommission_target_capacity_error(&disk_full)); + assert!(is_decommission_target_capacity_error(&storage_full)); + } + + #[test] + fn decommission_target_capacity_error_rejects_unrelated_errors() { + assert!(!is_decommission_target_capacity_error(&Error::SlowDown)); + } + #[test] fn should_skip_decommission_delete_marker_characterizes_empty_marker_without_replication() { let version = rustfs_filemeta::FileInfo { @@ -3809,6 +4165,8 @@ mod tests { items_decommission_failed: 1, bytes_done: 1024, bytes_failed: 128, + terminal_reload_attempt_at: Some(start_time), + terminal_reload_failures: vec!["complete_decommission: peer node-a failed".to_string()], ..Default::default() }), }], @@ -3838,34 +4196,71 @@ mod tests { assert_eq!(restored_decommission.bucket, "bucket-b"); assert_eq!(restored_decommission.prefix, "prefix"); assert_eq!(restored_decommission.object, "object.txt"); + assert!(restored_decommission.stage.is_empty()); assert_eq!(restored_decommission.items_decommissioned, 7); assert_eq!(restored_decommission.items_decommission_failed, 1); assert_eq!(restored_decommission.bytes_done, 1024); assert_eq!(restored_decommission.bytes_failed, 128); + assert_eq!(restored_decommission.terminal_reload_attempt_at, Some(start_time)); + assert_eq!( + restored_decommission.terminal_reload_failures, + vec!["complete_decommission: peer node-a failed".to_string()] + ); assert!(restored_decommission.queued); assert_eq!(restored_decommission.items_since_last_progress_save(), 0); } #[test] - fn pool_meta_decode_supports_legacy_payload() { + fn pool_meta_records_decommission_terminal_reload_failure_once() { let start_time = OffsetDateTime::now_utc(); - let legacy_meta = PoolMeta { + let mut pool_meta = PoolMeta { version: POOL_META_VERSION, pools: vec![PoolStatus { + id: 1, + cmd_line: "/data/pool1/disk{1...4}".to_string(), + last_update: start_time, + decommission: Some(PoolDecommissionInfo { + complete: true, + ..Default::default() + }), + }], + dont_save: false, + }; + + assert!( + pool_meta + .record_decommission_terminal_reload_failure(0, "complete_decommission", "peer node-a failed".to_string()) + .expect("terminal reload failure should be recorded") + ); + assert!( + !pool_meta + .record_decommission_terminal_reload_failure(0, "complete_decommission", "peer node-a failed".to_string()) + .expect("duplicate terminal reload failure should be ignored") + ); + + let decommission = pool_meta.pools[0].decommission.as_ref().expect("decommission should exist"); + assert!(decommission.terminal_reload_attempt_at.is_some()); + assert_eq!( + decommission.terminal_reload_failures, + vec!["complete_decommission: peer node-a failed".to_string()] + ); + } + + #[test] + fn pool_meta_decode_supports_legacy_payload() { + let start_time = OffsetDateTime::now_utc(); + let legacy_meta = LegacyPoolMeta { + version: POOL_META_VERSION, + pools: vec![LegacyPoolStatus { id: 3, cmd_line: "/legacy/pool".to_string(), last_update: start_time, - decommission: Some(PoolDecommissionInfo { + decommission: Some(LegacyPoolDecommissionInfo { start_time: Some(start_time), items_decommissioned: 9, items_decommission_failed: 2, bytes_done: 2048, bytes_failed: 256, - queued_buckets: vec!["not-persisted".to_string()], - decommissioned_buckets: vec!["not-persisted".to_string()], - bucket: "not-persisted".to_string(), - prefix: "not-persisted".to_string(), - object: "not-persisted".to_string(), ..Default::default() }), }], @@ -4099,13 +4494,13 @@ mod tests { #[test] fn pool_meta_decode_rejects_invalid_legacy_decommission_terminal_state() { let start_time = OffsetDateTime::now_utc(); - let legacy_meta = PoolMeta { + let legacy_meta = LegacyPoolMeta { version: POOL_META_VERSION, - pools: vec![PoolStatus { + pools: vec![LegacyPoolStatus { id: 1, cmd_line: "/legacy/pool".to_string(), last_update: start_time, - decommission: Some(PoolDecommissionInfo { + decommission: Some(LegacyPoolDecommissionInfo { start_time: Some(start_time), complete: true, failed: false, @@ -4113,7 +4508,7 @@ mod tests { ..Default::default() }), }], - dont_save: true, + dont_save: false, }; let mut payload = Vec::new(); @@ -4367,12 +4762,14 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi mod pools_tests { use super::{ DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, - DecommissionTerminalState, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, bind_decommission_cancelers, - bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, - count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, - decommission_start_guard_state, dedup_indices, default_decommission_bucket_concurrency, + DecommissionStartPoolState, DecommissionTerminalState, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, + apply_decommission_status_space_info, bind_decommission_cancelers, bind_missing_decommission_cancelers, + cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item, + decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, + decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available, - ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_local_leader, + 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, @@ -4380,17 +4777,18 @@ mod pools_tests { missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, pool_meta_has_active_decommission, require_decommission_store, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, - resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, - resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result, - resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result, - resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result, - resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result, - resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta, - run_decommission_buckets_bounded, should_cleanup_decommission_source_entry, should_continue_decommission_queue, - should_count_decommission_version_complete, should_preserve_decommission_canceled_state, - should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload, - should_skip_canceled_decommission_routine, split_decommission_buckets, take_and_cancel_decommission_canceler, - take_decommission_canceler, track_decommission_current_object, validate_start_decommission_request, + resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, + resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, + resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result, + resolve_decommission_progress_save_result, resolve_decommission_spawn_failure_result, + resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result, + resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result, + rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, should_cleanup_decommission_source_entry, + should_continue_decommission_queue, should_count_decommission_version_complete, + should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal, + should_retry_decommission_cancel_reload, should_skip_canceled_decommission_routine, split_decommission_buckets, + take_and_cancel_decommission_canceler, take_decommission_canceler, touch_decommission_progress, + track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request, wait_decommission_worker_drain, with_decommission_entry_context, }; use crate::data_movement; @@ -4398,7 +4796,7 @@ mod pools_tests { use crate::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::error::Error; use crate::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; - use rustfs_filemeta::MetaCacheEntry; + use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; use rustfs_rio::Index; use std::sync::{ Arc, @@ -4435,6 +4833,49 @@ mod pools_tests { } } + #[test] + fn test_apply_decommission_status_space_info_adds_idle_pool_usage() { + let status = apply_decommission_status_space_info( + decommission_test_pool_status(0, None), + PoolSpaceInfo { + free: 25, + total: 100, + used: 75, + }, + ); + + let decommission = status.decommission.expect("idle pool status should include usage info"); + assert_eq!(decommission.total_size, 100); + assert_eq!(decommission.current_size, 25); + assert!(decommission.start_time.is_none()); + assert!(!decommission.complete); + assert!(!decommission.failed); + assert!(!decommission.canceled); + } + + #[test] + fn test_apply_decommission_status_space_info_refreshes_active_decommission_sizes() { + let status = apply_decommission_status_space_info( + decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + total_size: 1, + current_size: 1, + ..Default::default() + }), + ), + PoolSpaceInfo { + free: 25, + total: 100, + used: 75, + }, + ); + + let decommission = status.decommission.expect("active decommission info should remain present"); + assert_eq!(decommission.total_size, 100); + assert_eq!(decommission.current_size, 25); + } + #[test] fn test_dedup_indices_removes_duplicates_preserving_order() { assert_eq!(dedup_indices(&[0, 2, 1, 2, 3, 0]), vec![0, 2, 1, 3]); @@ -4664,14 +5105,15 @@ mod pools_tests { active.pools[0].decommission = Some(PoolDecommissionInfo::default()); assert!(active.is_suspended(0)); - let (_, decommission_active, _) = decommission_start_guard_state(active.pools.first()); - assert!(decommission_active); + assert_eq!( + decommission_start_pool_state(active.pools.first()), + DecommissionStartPoolState::Decommissioning + ); rollback_start_decommission_pool_meta(&mut active, previous); assert!(!active.is_suspended(0)); - let (_, decommission_active, _) = decommission_start_guard_state(active.pools.first()); - assert!(!decommission_active); + assert_eq!(decommission_start_pool_state(active.pools.first()), DecommissionStartPoolState::Active); } #[test] @@ -4924,6 +5366,28 @@ mod pools_tests { let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist"); assert_eq!(info.bucket, "bucket-a"); assert_eq!(info.object, "object-a"); + assert!(info.stage.is_empty()); + } + + #[test] + fn test_track_decommission_current_object_stage_updates_stage() { + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(PoolDecommissionInfo::default()), + }], + ..Default::default() + }; + + track_decommission_current_object_stage(&mut meta, 0, "bucket-a", "object-a", "cleanup_preflight") + .expect("valid state should track bucket/object stage"); + + let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist"); + assert_eq!(info.bucket, "bucket-a"); + assert_eq!(info.object, "object-a"); + assert_eq!(info.stage, "cleanup_preflight"); } #[test] @@ -5163,6 +5627,55 @@ mod pools_tests { assert!(message.contains("object obj.txt")); } + #[test] + fn test_resolve_decommission_entry_exact_versions_preserves_full_parts() { + let entry = MetaCacheEntry { + name: "obj.txt".to_string(), + metadata: Vec::new(), + cached: None, + reusable: false, + }; + let fivs = FileInfoVersions { + volume: "bucket-a".to_string(), + name: "obj.txt".to_string(), + versions: vec![FileInfo { + name: "obj.txt".to_string(), + parts: vec![ObjectPartInfo { + number: 1, + etag: "part-etag".to_string(), + size: 128, + actual_size: 128, + ..Default::default() + }], + ..Default::default() + }], + ..Default::default() + }; + + let resolved = resolve_decommission_entry_exact_versions(Ok(Some(fivs)), &entry, "bucket-a", "file_info_versions") + .expect("exact versions should be preserved"); + + assert_eq!(resolved.versions[0].parts.len(), 1); + assert_eq!(resolved.versions[0].parts[0].etag, "part-etag"); + } + + #[test] + fn test_resolve_decommission_entry_exact_versions_uses_empty_when_source_missing() { + let entry = MetaCacheEntry { + name: "obj.txt".to_string(), + metadata: Vec::new(), + cached: None, + reusable: false, + }; + + let resolved = resolve_decommission_entry_exact_versions(Ok(None), &entry, "bucket-a", "file_info_versions") + .expect("missing source metadata should be treated as empty"); + + assert_eq!(resolved.volume, "bucket-a"); + assert_eq!(resolved.name, "obj.txt"); + assert!(resolved.versions.is_empty()); + } + #[test] fn test_resolve_decommission_check_after_list_result_prefers_entry_error() { let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), Some(Error::SlowDown)) @@ -5304,6 +5817,29 @@ mod pools_tests { ); } + #[test] + fn test_touch_decommission_progress_updates_last_update_and_save_baseline() { + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(PoolDecommissionInfo { + items_decommissioned: 3, + items_decommission_failed: 2, + ..Default::default() + }), + }], + ..Default::default() + }; + + touch_decommission_progress(&mut meta, 0).expect("valid decommission progress should be touched"); + + assert!(meta.pools[0].last_update > OffsetDateTime::UNIX_EPOCH); + let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist"); + assert_eq!(info.items_since_last_progress_save(), 0); + } + #[test] fn test_pool_meta_update_after_skips_before_time_and_item_thresholds() { let mut meta = PoolMeta { @@ -5565,7 +6101,8 @@ mod pools_tests { #[test] fn test_ensure_decommission_start_allowed_rejects_missing_pool() { - let err = ensure_decommission_start_allowed(false, false, false).expect_err("missing pool should be invalid"); + let err = + ensure_decommission_start_allowed(DecommissionStartPoolState::Missing).expect_err("missing pool should be invalid"); assert!( err.to_string() .contains("failed to start decommission: target pool was not found") @@ -5574,28 +6111,37 @@ mod pools_tests { #[test] fn test_ensure_decommission_start_allowed_rejects_running_state() { - let err = ensure_decommission_start_allowed(true, true, false).expect_err("active decommission should be rejected"); + let err = ensure_decommission_start_allowed(DecommissionStartPoolState::Decommissioning) + .expect_err("active decommission should be rejected"); assert!(matches!(err, Error::DecommissionAlreadyRunning)); } #[test] fn test_ensure_decommission_start_allowed_rejects_completed_state() { - let err = ensure_decommission_start_allowed(true, false, true).expect_err("completed decommission should be rejected"); - assert!(err.to_string().contains("target pool decommission is already complete")); + let err = ensure_decommission_start_allowed(DecommissionStartPoolState::Decommissioned) + .expect_err("completed decommission should be rejected"); + assert!(err.to_string().contains("target pool is already decommissioned")); } #[test] - fn test_ensure_decommission_start_allowed_allows_failed_or_canceled_state() { - assert!(ensure_decommission_start_allowed(true, false, false).is_ok()); + fn test_ensure_decommission_start_allowed_rejects_blocked_state() { + let err = ensure_decommission_start_allowed(DecommissionStartPoolState::Blocked) + .expect_err("blocked decommission should be rejected"); + assert!(err.to_string().contains("target pool decommission is blocked")); } #[test] - fn test_decommission_start_guard_state_reports_missing_pool() { - assert_eq!(decommission_start_guard_state(None), (false, false, false)); + fn test_ensure_decommission_start_allowed_allows_active_state() { + assert!(ensure_decommission_start_allowed(DecommissionStartPoolState::Active).is_ok()); } #[test] - fn test_decommission_start_guard_state_reports_idle_pool_without_decommission_info() { + fn test_decommission_start_pool_state_reports_missing_pool() { + assert_eq!(decommission_start_pool_state(None), DecommissionStartPoolState::Missing); + } + + #[test] + fn test_decommission_start_pool_state_reports_idle_pool_without_decommission_info() { let pool = PoolStatus { id: 0, cmd_line: "pool-0".to_string(), @@ -5603,11 +6149,11 @@ mod pools_tests { decommission: None, }; - assert_eq!(decommission_start_guard_state(Some(&pool)), (true, false, false)); + assert_eq!(decommission_start_pool_state(Some(&pool)), DecommissionStartPoolState::Active); } #[test] - fn test_decommission_start_guard_state_reports_active_pool_when_not_terminal() { + fn test_decommission_start_pool_state_reports_decommissioning_pool_when_not_terminal() { let pool = PoolStatus { id: 0, cmd_line: "pool-0".to_string(), @@ -5620,11 +6166,11 @@ mod pools_tests { }), }; - assert_eq!(decommission_start_guard_state(Some(&pool)), (true, true, false)); + assert_eq!(decommission_start_pool_state(Some(&pool)), DecommissionStartPoolState::Decommissioning); } #[test] - fn test_decommission_start_guard_state_reports_canceled_pool_as_restartable() { + fn test_decommission_start_pool_state_reports_canceled_pool_as_blocked() { let pool = PoolStatus { id: 0, cmd_line: "pool-0".to_string(), @@ -5637,11 +6183,11 @@ mod pools_tests { }), }; - assert_eq!(decommission_start_guard_state(Some(&pool)), (true, false, false)); + assert_eq!(decommission_start_pool_state(Some(&pool)), DecommissionStartPoolState::Blocked); } #[test] - fn test_decommission_start_guard_state_reports_failed_pool_as_restartable() { + fn test_decommission_start_pool_state_reports_failed_pool_as_blocked() { let pool = PoolStatus { id: 0, cmd_line: "pool-0".to_string(), @@ -5652,11 +6198,11 @@ mod pools_tests { }), }; - assert_eq!(decommission_start_guard_state(Some(&pool)), (true, false, false)); + assert_eq!(decommission_start_pool_state(Some(&pool)), DecommissionStartPoolState::Blocked); } #[test] - fn test_decommission_start_guard_state_reports_completed_pool() { + fn test_decommission_start_pool_state_reports_completed_pool() { let pool = PoolStatus { id: 0, cmd_line: "pool-0".to_string(), @@ -5667,7 +6213,75 @@ mod pools_tests { }), }; - assert_eq!(decommission_start_guard_state(Some(&pool)), (true, false, true)); + assert_eq!(decommission_start_pool_state(Some(&pool)), DecommissionStartPoolState::Decommissioned); + } + + #[test] + fn test_ensure_decommission_start_keeps_active_pool_rejects_last_active_pool() { + let meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }], + ..Default::default() + }; + + let err = ensure_decommission_start_keeps_active_pool(&meta, &[0]).expect_err("last active pool should be rejected"); + + assert!(err.to_string().contains("at least one active pool must remain")); + } + + #[test] + fn test_ensure_decommission_start_pool_states_rejects_blocked_pool() { + let meta = PoolMeta { + pools: vec![ + PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(PoolDecommissionInfo { + failed: true, + ..Default::default() + }), + }, + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + ], + ..Default::default() + }; + + let err = ensure_decommission_start_pool_states(&meta, &[0]).expect_err("blocked pool should be rejected"); + + assert!(err.to_string().contains("target pool decommission is blocked")); + } + + #[test] + fn test_ensure_decommission_start_pool_states_allows_active_pool_with_remaining_active_pool() { + let meta = PoolMeta { + pools: vec![ + PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + ], + ..Default::default() + }; + + assert!(ensure_decommission_start_pool_states(&meta, &[0]).is_ok()); } #[test] @@ -5704,7 +6318,7 @@ mod pools_tests { #[test] fn test_should_preserve_decommission_canceled_state_when_signal_canceled() { - assert!(should_preserve_decommission_canceled_state(false, true)); + assert!(!should_preserve_decommission_canceled_state(false, true)); } #[test] @@ -6059,7 +6673,7 @@ mod pools_tests { } #[test] - fn test_pool_meta_decommission_retry_preserves_completed_bucket_progress() { + fn test_pool_meta_failed_decommission_requires_clear_before_restart() { let mut meta = PoolMeta { pools: vec![PoolStatus { id: 0, @@ -6083,6 +6697,29 @@ mod pools_tests { ..Default::default() }; + let err = meta + .decommission( + 0, + PoolSpaceInfo { + total: 200, + free: 50, + used: 150, + }, + ) + .expect_err("failed decommission should be blocked until cleared"); + assert!(err.to_string().contains("target pool decommission is blocked")); + let blocked = meta.pools[0] + .decommission + .as_ref() + .expect("blocked metadata should remain until clear"); + assert!(blocked.failed); + assert_eq!(blocked.decommissioned_buckets, vec!["bucket-done".to_string()]); + assert_eq!(blocked.items_decommissioned, 7); + assert_eq!(blocked.bytes_done, 1024); + + assert!(meta.clear_decommission(0).expect("failed decommission should clear")); + assert!(meta.pools[0].decommission.is_none()); + meta.decommission( 0, PoolSpaceInfo { @@ -6091,7 +6728,7 @@ mod pools_tests { used: 150, }, ) - .expect("failed decommission should be restartable"); + .expect("cleared decommission should be restartable"); meta.queue_buckets( 0, vec![ @@ -6113,11 +6750,11 @@ mod pools_tests { assert!(!info.failed); assert!(!info.canceled); assert!(!info.complete); - assert_eq!(info.decommissioned_buckets, vec!["bucket-done".to_string()]); - assert_eq!(info.queued_buckets, vec!["bucket-pending".to_string()]); - assert_eq!(info.items_decommissioned, 7); + assert!(info.decommissioned_buckets.is_empty()); + assert_eq!(info.queued_buckets, vec!["bucket-done".to_string(), "bucket-pending".to_string()]); + assert_eq!(info.items_decommissioned, 0); assert_eq!(info.items_decommission_failed, 0); - assert_eq!(info.bytes_done, 1024); + assert_eq!(info.bytes_done, 0); assert_eq!(info.bytes_failed, 0); assert_eq!(info.items_since_last_progress_save(), 0); assert_eq!(info.start_size, 50); @@ -6130,7 +6767,7 @@ mod pools_tests { } #[test] - fn test_pool_meta_queued_decommission_retry_preserves_canceled_completed_buckets() { + fn test_pool_meta_canceled_queued_decommission_requires_clear_before_restart() { let mut meta = PoolMeta { pools: vec![PoolStatus { id: 0, @@ -6147,6 +6784,29 @@ mod pools_tests { ..Default::default() }; + let err = meta + .queue_decommission( + 0, + PoolSpaceInfo { + total: 100, + free: 25, + used: 75, + }, + ) + .expect_err("canceled queued decommission should be blocked until cleared"); + assert!(err.to_string().contains("target pool decommission is blocked")); + let blocked = meta.pools[0] + .decommission + .as_ref() + .expect("blocked metadata should remain until clear"); + assert!(blocked.canceled); + assert_eq!(blocked.decommissioned_buckets, vec!["bucket-done".to_string()]); + assert_eq!(blocked.items_decommissioned, 5); + assert_eq!(blocked.bytes_done, 512); + + assert!(meta.clear_decommission(0).expect("canceled decommission should clear")); + assert!(meta.pools[0].decommission.is_none()); + meta.queue_decommission( 0, PoolSpaceInfo { @@ -6155,7 +6815,7 @@ mod pools_tests { used: 75, }, ) - .expect("canceled queued decommission should be restartable"); + .expect("cleared queued decommission should be restartable"); let info = meta.pools[0] .decommission @@ -6163,9 +6823,9 @@ mod pools_tests { .expect("decommission info should be rebuilt"); assert!(info.queued); assert!(info.start_time.is_none()); - assert_eq!(info.decommissioned_buckets, vec!["bucket-done".to_string()]); - assert_eq!(info.items_decommissioned, 5); - assert_eq!(info.bytes_done, 512); + assert!(info.decommissioned_buckets.is_empty()); + assert_eq!(info.items_decommissioned, 0); + assert_eq!(info.bytes_done, 0); } #[test] diff --git a/crates/ecstore/src/rebalance/control.rs b/crates/ecstore/src/rebalance/control.rs index 9539674ce..625e0e5bc 100644 --- a/crates/ecstore/src/rebalance/control.rs +++ b/crates/ecstore/src/rebalance/control.rs @@ -1,10 +1,10 @@ use super::meta::{ - clone_first_arc, clone_rebalance_pool_stats, defer_bucket_in_rebalance_queue, ensure_valid_rebalance_pool_index, - invalid_rebalance_pool_index_error, is_rebalance_conflicting_with_decommission, mark_rebalance_bucket_done, - merge_rebalance_meta, percent_free_ratio, rebalance_metadata_not_initialized_error, record_rebalance_cleanup_warning_in_meta, - record_rebalance_stop_propagation_snapshot, resolve_next_rebalance_bucket, rollback_rebalance_start_meta_snapshot_for_id, - should_accept_rebalance_stats_update, should_pool_participate, stop_rebalance_meta_snapshot_for_id, - validate_init_rebalance_state, + RebalanceMetaMergeOutcome, clone_first_arc, clone_rebalance_pool_stats, defer_bucket_in_rebalance_queue, + ensure_valid_rebalance_pool_index, invalid_rebalance_pool_index_error, is_rebalance_conflicting_with_decommission, + mark_rebalance_bucket_done, merge_rebalance_meta, percent_free_ratio, rebalance_metadata_not_initialized_error, + record_rebalance_cleanup_warning_in_meta, record_rebalance_stop_propagation_snapshot, resolve_next_rebalance_bucket, + rollback_rebalance_start_meta_snapshot_for_id, should_accept_rebalance_stats_update, should_pool_participate, + stop_rebalance_meta_snapshot_for_id, validate_init_rebalance_state, }; use super::worker::{ rebalance_meta_lock_error, resolve_load_rebalance_stats_update_result, resolve_rebalance_meta_load_result, @@ -27,6 +27,18 @@ use time::OffsetDateTime; use tracing::{debug, info}; use uuid::Uuid; +pub(super) fn validate_rebalance_disk_stats_coverage(disk_stats: &[DiskStat]) -> Result<()> { + for (idx, disk_stat) in disk_stats.iter().enumerate() { + if disk_stat.total_space == 0 { + return Err(Error::other(format!( + "rebalance storage info is incomplete: pool {idx} has no reported capacity" + ))); + } + } + + Ok(()) +} + impl ECStore { pub(super) async fn save_rebalance_meta_with_merge( &self, @@ -50,7 +62,9 @@ impl ECStore { let mut merged = RebalanceMeta::new(); match merged.load_with_opts(pool.clone(), opts.clone()).await { Ok(()) => { - merge_rebalance_meta(&mut merged, local_snapshot); + if merge_rebalance_meta(&mut merged, local_snapshot) == RebalanceMetaMergeOutcome::RejectedActiveConflict { + return Err(Error::RebalanceAlreadyRunning); + } } Err(Error::ConfigNotFound) => { merged = local_snapshot.clone(); @@ -90,6 +104,10 @@ impl ECStore { "Loaded rebalance metadata" ); } else { + { + let mut rebalance_meta = self.rebalance_meta.write().await; + *rebalance_meta = None; + } debug!( event = EVENT_REBALANCE_STATE, component = LOG_COMPONENT_ECSTORE, @@ -191,6 +209,7 @@ impl ECStore { } let percent_free_goal = percent_free_ratio(total_free, total_cap); + validate_rebalance_disk_stats_coverage(&disk_stats)?; let mut pool_stats = Vec::with_capacity(self.pools.len()); @@ -421,6 +440,15 @@ impl ECStore { false } + pub async fn pool_rebalance_status(&self, pool_index: usize) -> (RebalStatus, bool) { + let rebalance_meta = self.rebalance_meta.read().await; + rebalance_meta + .as_ref() + .and_then(|meta| meta.pool_stats.get(pool_index)) + .map(|pool_stat| (pool_stat.info.status, pool_stat.info.stopping)) + .unwrap_or_default() + } + pub async fn current_rebalance_id(&self) -> Option { let rebalance_meta = self.rebalance_meta.read().await; rebalance_meta diff --git a/crates/ecstore/src/rebalance/meta.rs b/crates/ecstore/src/rebalance/meta.rs index c78ee68cc..289a58b41 100644 --- a/crates/ecstore/src/rebalance/meta.rs +++ b/crates/ecstore/src/rebalance/meta.rs @@ -457,7 +457,10 @@ pub(super) fn complete_rebalance_pools_at_goal(meta: &mut RebalanceMeta, now: Of let mut changed = false; for pool_stat in meta.pool_stats.iter_mut() { - if !is_rebalance_pool_started(pool_stat) || has_deferred_rebalance_error(pool_stat) { + if !is_rebalance_pool_started(pool_stat) + || has_deferred_rebalance_error(pool_stat) + || has_rebalance_cleanup_warnings(pool_stat) + { continue; } @@ -481,7 +484,7 @@ pub(super) fn complete_rebalance_pools_with_empty_queue(meta: &mut RebalanceMeta let mut changed = false; for pool_stat in meta.pool_stats.iter_mut() { - if !is_rebalance_pool_started(pool_stat) || !pool_stat.buckets.is_empty() { + if !is_rebalance_pool_started(pool_stat) || !pool_stat.buckets.is_empty() || has_rebalance_cleanup_warnings(pool_stat) { continue; } @@ -501,6 +504,11 @@ pub(super) fn has_deferred_rebalance_error(pool_stat: &RebalanceStats) -> bool { .as_deref() .is_some_and(|last_error| last_error.starts_with(REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX)) } + +pub(super) fn has_rebalance_cleanup_warnings(pool_stat: &RebalanceStats) -> bool { + pool_stat.cleanup_warnings.count > 0 +} + pub(super) fn clone_first_arc(values: &[Arc], err_msg: &str) -> Result> { values.first().cloned().ok_or_else(|| Error::other(err_msg)) } @@ -608,7 +616,7 @@ pub(super) fn validate_init_rebalance_state(decommission_running: bool, current_ if !ensure_rebalance_not_decommissioning(decommission_running) { return Err(Error::DecommissionAlreadyRunning); } - if current_meta.is_some_and(is_rebalance_in_progress) { + if current_meta.is_some_and(|meta| !is_rebalance_meta_replaceable_for_new_id(meta)) { return Err(Error::RebalanceAlreadyRunning); } @@ -812,14 +820,29 @@ pub(super) fn merge_rebalance_pool_stats(remote: &mut RebalanceStats, local: &Re } } -pub(super) fn merge_rebalance_meta(remote: &mut RebalanceMeta, local: &RebalanceMeta) { +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum RebalanceMetaMergeOutcome { + Merged, + Replaced, + RejectedActiveConflict, +} + +pub(super) fn is_rebalance_meta_replaceable_for_new_id(meta: &RebalanceMeta) -> bool { + meta.stopped_at.is_some() || !is_rebalance_in_progress(meta) +} + +pub(super) fn merge_rebalance_meta(remote: &mut RebalanceMeta, local: &RebalanceMeta) -> RebalanceMetaMergeOutcome { if remote.id.is_empty() { *remote = local.clone(); - return; + return RebalanceMetaMergeOutcome::Replaced; } if !local.id.is_empty() && remote.id != local.id { - return; + if is_rebalance_meta_replaceable_for_new_id(remote) { + *remote = local.clone(); + return RebalanceMetaMergeOutcome::Replaced; + } + return RebalanceMetaMergeOutcome::RejectedActiveConflict; } remote.percent_free_goal = local.percent_free_goal; @@ -837,6 +860,8 @@ pub(super) fn merge_rebalance_meta(remote: &mut RebalanceMeta, local: &Rebalance merge_rebalance_pool_stats(remote_pool_stat, local_pool_stat); } } + + RebalanceMetaMergeOutcome::Merged } pub(super) fn mark_started_rebalance_pools_stopped(meta: &mut RebalanceMeta, stop_time: OffsetDateTime) { diff --git a/crates/ecstore/src/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/rebalance/rebalance_unit_tests.rs index 0beb835c1..ba9a56d05 100644 --- a/crates/ecstore/src/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/rebalance/rebalance_unit_tests.rs @@ -13,19 +13,21 @@ // limitations under the License. use super::REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX; +use super::control::validate_rebalance_disk_stats_coverage; use super::meta::{ - RebalanceTerminalEvent, apply_rebalance_save_option, apply_rebalance_terminal_event, apply_stopped_at, - classify_rebalance_terminal_event, clone_arc_by_index, clone_first_arc, clone_rebalance_pool_stats, + RebalanceMetaMergeOutcome, RebalanceTerminalEvent, apply_rebalance_save_option, apply_rebalance_terminal_event, + apply_stopped_at, classify_rebalance_terminal_event, clone_arc_by_index, clone_first_arc, clone_rebalance_pool_stats, complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, defer_bucket_in_rebalance_queue, ensure_rebalance_not_decommissioning, ensure_valid_rebalance_pool_index, first_rebalance_bucket, has_deferred_rebalance_error, is_rebalance_actively_running, is_rebalance_conflicting_with_decommission, - is_rebalance_in_progress, is_rebalance_stopped_terminal_event, mark_rebalance_bucket_done, merge_rebalance_bucket_lists, - merge_rebalance_meta, next_rebal_bucket_from_stat, percent_free_ratio, rebalance_goal_reached, - rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error, rebalance_meta_load_unknown_version_error, - record_rebalance_cleanup_warning_in_meta, remove_rebalanced_buckets_from_queue, resolve_next_rebalance_bucket, - resolve_rebalance_participants, should_accept_rebalance_stats_update, should_ignore_rebalance_data_usage_cache, - should_pool_participate, should_preserve_rebalance_stopped_state, should_skip_start_rebalance, stop_rebalance_meta_snapshot, - stop_rebalance_state, take_bucket_from_rebalance_queue, validate_init_rebalance_state, validate_start_rebalance_state, + is_rebalance_in_progress, is_rebalance_meta_replaceable_for_new_id, is_rebalance_stopped_terminal_event, + mark_rebalance_bucket_done, merge_rebalance_bucket_lists, merge_rebalance_meta, next_rebal_bucket_from_stat, + percent_free_ratio, rebalance_goal_reached, rebalance_meta_load_no_data_error, rebalance_meta_load_unknown_format_error, + rebalance_meta_load_unknown_version_error, record_rebalance_cleanup_warning_in_meta, remove_rebalanced_buckets_from_queue, + resolve_next_rebalance_bucket, resolve_rebalance_participants, should_accept_rebalance_stats_update, + should_ignore_rebalance_data_usage_cache, should_pool_participate, should_preserve_rebalance_stopped_state, + should_skip_start_rebalance, stop_rebalance_meta_snapshot, stop_rebalance_state, take_bucket_from_rebalance_queue, + validate_init_rebalance_state, validate_start_rebalance_state, }; use super::migration::{ MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait, @@ -44,8 +46,8 @@ use super::worker::{ wait_rebalance_listing_retry, with_rebalance_entry_context, }; use super::{ - GetObjectReader, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, RebalanceBucketOutcome, RebalanceCleanupWarnings, - RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats, + DiskStat, GetObjectReader, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, RebalanceBucketOutcome, + RebalanceCleanupWarnings, RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats, }; use crate::data_movement; use crate::data_usage::DATA_USAGE_CACHE_NAME; @@ -1207,7 +1209,7 @@ fn test_merge_rebalance_meta_preserves_updates_from_multiple_pools() { ..Default::default() }; - merge_rebalance_meta(&mut remote, &local); + assert_eq!(merge_rebalance_meta(&mut remote, &local), RebalanceMetaMergeOutcome::Merged); assert_eq!(remote.pool_stats[0].num_versions, 4); assert_eq!(remote.pool_stats[0].object, "remote-object"); @@ -1220,6 +1222,163 @@ fn test_merge_rebalance_meta_preserves_updates_from_multiple_pools() { assert_eq!(remote.pool_stats[1].cleanup_warnings.last_at, Some(warning_at)); } +#[test] +fn test_merge_rebalance_meta_replaces_terminal_metadata_for_new_rebalance() { + let old_completed_at = OffsetDateTime::from_unix_timestamp(1_000).expect("valid old completion timestamp"); + let new_started_at = OffsetDateTime::from_unix_timestamp(2_000).expect("valid new start timestamp"); + let mut remote = RebalanceMeta { + id: "old-rebalance".to_string(), + percent_free_goal: 0.25, + pool_stats: vec![ + RebalanceStats { + buckets: Vec::new(), + rebalanced_buckets: vec!["bucket-a".to_string()], + participating: true, + info: RebalanceInfo { + status: RebalStatus::Completed, + end_time: Some(old_completed_at), + ..Default::default() + }, + num_versions: 7, + bytes: 700, + ..Default::default() + }, + RebalanceStats::default(), + ], + ..Default::default() + }; + let local = RebalanceMeta { + id: "new-rebalance".to_string(), + percent_free_goal: 0.5, + pool_stats: vec![ + RebalanceStats { + buckets: vec!["bucket-a".to_string()], + participating: true, + info: RebalanceInfo { + start_time: Some(new_started_at), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats { + buckets: vec!["bucket-a".to_string()], + participating: true, + info: RebalanceInfo { + start_time: Some(new_started_at), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats::default(), + ], + ..Default::default() + }; + + assert_eq!(merge_rebalance_meta(&mut remote, &local), RebalanceMetaMergeOutcome::Replaced); + + assert_eq!(remote.id, "new-rebalance"); + assert_eq!(remote.percent_free_goal, 0.5); + assert_eq!(remote.pool_stats.len(), 3); + assert_eq!(remote.pool_stats[0].info.status, RebalStatus::Started); + assert_eq!(remote.pool_stats[0].buckets, vec!["bucket-a"]); + assert_eq!(remote.pool_stats[0].rebalanced_buckets, Vec::::new()); + assert_eq!(remote.pool_stats[0].num_versions, 0); + assert_eq!(remote.pool_stats[1].info.status, RebalStatus::Started); + assert!(remote.pool_stats[1].participating); +} + +#[test] +fn test_merge_rebalance_meta_preserves_active_metadata_for_different_rebalance_id() { + let old_started_at = OffsetDateTime::from_unix_timestamp(1_000).expect("valid old start timestamp"); + let new_started_at = OffsetDateTime::from_unix_timestamp(2_000).expect("valid new start timestamp"); + let mut remote = RebalanceMeta { + id: "active-rebalance".to_string(), + percent_free_goal: 0.25, + pool_stats: vec![RebalanceStats { + buckets: vec!["bucket-a".to_string()], + participating: true, + info: RebalanceInfo { + start_time: Some(old_started_at), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + let local = RebalanceMeta { + id: "new-rebalance".to_string(), + percent_free_goal: 0.5, + pool_stats: vec![RebalanceStats { + buckets: vec!["bucket-b".to_string()], + participating: true, + info: RebalanceInfo { + start_time: Some(new_started_at), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + + assert_eq!( + merge_rebalance_meta(&mut remote, &local), + RebalanceMetaMergeOutcome::RejectedActiveConflict + ); + + assert_eq!(remote.id, "active-rebalance"); + assert_eq!(remote.percent_free_goal, 0.25); + assert_eq!(remote.pool_stats.len(), 1); + assert_eq!(remote.pool_stats[0].buckets, vec!["bucket-a"]); + assert_eq!(remote.pool_stats[0].info.start_time, Some(old_started_at)); +} + +#[test] +fn test_merge_rebalance_meta_replaces_stopped_started_metadata_for_new_rebalance() { + let stopped_at = OffsetDateTime::from_unix_timestamp(1_500).expect("valid stop timestamp"); + let new_started_at = OffsetDateTime::from_unix_timestamp(2_000).expect("valid new start timestamp"); + let mut remote = RebalanceMeta { + id: "stopped-rebalance".to_string(), + stopped_at: Some(stopped_at), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + stopping: true, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + let local = RebalanceMeta { + id: "new-rebalance".to_string(), + percent_free_goal: 0.5, + pool_stats: vec![RebalanceStats { + buckets: vec!["bucket-a".to_string()], + participating: true, + info: RebalanceInfo { + start_time: Some(new_started_at), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + + assert!(is_rebalance_meta_replaceable_for_new_id(&remote)); + assert_eq!(merge_rebalance_meta(&mut remote, &local), RebalanceMetaMergeOutcome::Replaced); + + assert_eq!(remote.id, "new-rebalance"); + assert_eq!(remote.stopped_at, None); + assert_eq!(remote.pool_stats[0].info.status, RebalStatus::Started); + assert!(!remote.pool_stats[0].info.stopping); +} + #[test] fn test_merge_rebalance_meta_does_not_overwrite_failed_with_started_stats() { let now = OffsetDateTime::from_unix_timestamp(2_000).unwrap(); @@ -2202,6 +2361,7 @@ fn test_validate_init_rebalance_state_rejects_active_rebalance() { #[test] fn test_validate_init_rebalance_state_allows_terminal_or_missing_rebalance() { + let stopped_at = OffsetDateTime::now_utc(); let completed = RebalanceMeta { pool_stats: vec![RebalanceStats { participating: true, @@ -2213,9 +2373,23 @@ fn test_validate_init_rebalance_state_allows_terminal_or_missing_rebalance() { }], ..Default::default() }; + let stopped_started = RebalanceMeta { + stopped_at: Some(stopped_at), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + stopping: true, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; validate_init_rebalance_state(false, None).expect("missing rebalance meta should allow init"); validate_init_rebalance_state(false, Some(&completed)).expect("terminal rebalance meta should allow init"); + validate_init_rebalance_state(false, Some(&stopped_started)).expect("stopped rebalance meta should allow new init"); } #[tokio::test] @@ -2298,6 +2472,38 @@ fn test_rebalance_goal_not_reached_for_issue_3137_initial_imbalance() { assert!(!rebalance_goal_reached(pool0_free, pool0_capacity, 0, goal)); } +#[test] +fn test_validate_rebalance_disk_stats_coverage_rejects_missing_pool_capacity() { + let disk_stats = vec![ + DiskStat { + total_space: 1_000, + available_space: 100, + }, + DiskStat::default(), + ]; + + let err = validate_rebalance_disk_stats_coverage(&disk_stats) + .expect_err("missing pool capacity should reject rebalance initialization"); + + assert!(err.to_string().contains("pool 1 has no reported capacity")); +} + +#[test] +fn test_validate_rebalance_disk_stats_coverage_accepts_all_pools() { + let disk_stats = vec![ + DiskStat { + total_space: 1_000, + available_space: 100, + }, + DiskStat { + total_space: 2_000, + available_space: 1_500, + }, + ]; + + assert!(validate_rebalance_disk_stats_coverage(&disk_stats).is_ok()); +} + #[test] fn test_complete_rebalance_pools_at_goal_marks_started_participants_completed() { let now = OffsetDateTime::from_unix_timestamp(1_000).unwrap(); @@ -2943,7 +3149,7 @@ fn test_record_rebalance_cleanup_warning_in_meta_preserves_last_error() { } #[test] -fn test_complete_rebalance_pools_with_empty_queue_preserves_cleanup_warnings() { +fn test_complete_rebalance_pools_with_empty_queue_skips_cleanup_warnings() { let warning_at = OffsetDateTime::from_unix_timestamp(9_000).unwrap(); let completed_at = OffsetDateTime::from_unix_timestamp(10_000).unwrap(); let mut meta = RebalanceMeta { @@ -2966,9 +3172,10 @@ fn test_complete_rebalance_pools_with_empty_queue_preserves_cleanup_warnings() { ..Default::default() }; - assert!(complete_rebalance_pools_with_empty_queue(&mut meta, completed_at)); + assert!(!complete_rebalance_pools_with_empty_queue(&mut meta, completed_at)); - assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Completed); + assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started); + assert_eq!(meta.pool_stats[0].info.end_time, None); assert!(meta.pool_stats[0].info.last_error.is_none()); assert_eq!(meta.pool_stats[0].cleanup_warnings.count, 1); assert_eq!(meta.pool_stats[0].cleanup_warnings.last_message.as_deref(), Some("cleanup failed")); diff --git a/crates/ecstore/src/rebalance/runtime.rs b/crates/ecstore/src/rebalance/runtime.rs index bbb5707de..78f91bf22 100644 --- a/crates/ecstore/src/rebalance/runtime.rs +++ b/crates/ecstore/src/rebalance/runtime.rs @@ -1,8 +1,9 @@ use super::meta::{ apply_rebalance_save_option, apply_rebalance_terminal_event, classify_rebalance_terminal_event, clone_first_arc, complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, ensure_valid_rebalance_pool_index, - has_deferred_rebalance_error, is_rebalance_in_progress, rebalance_goal_reached, resolve_rebalance_participants, - should_preserve_rebalance_stopped_state, should_skip_start_rebalance, validate_start_rebalance_state, + has_deferred_rebalance_error, has_rebalance_cleanup_warnings, is_rebalance_in_progress, rebalance_goal_reached, + resolve_rebalance_participants, should_preserve_rebalance_stopped_state, should_skip_start_rebalance, + validate_start_rebalance_state, }; use super::worker::{ resolve_rebalance_bucket_result, resolve_rebalance_meta_save_result, resolve_rebalance_save_task_result, @@ -211,7 +212,20 @@ impl ECStore { if let Some(meta) = rebalance_meta.as_mut() { let meta_stopped = meta.stopped_at.is_some(); if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) { - if should_preserve_rebalance_stopped_state( + if matches!(&terminal_event, super::meta::RebalanceTerminalEvent::Completed { .. }) + && has_rebalance_cleanup_warnings(pool_stat) + { + pool_stat.info.stopping = false; + pool_stat.info.status = RebalStatus::Failed; + pool_stat.info.end_time = Some(now); + pool_stat.info.last_error = Some( + pool_stat + .cleanup_warnings + .last_message + .clone() + .unwrap_or_else(|| "rebalance source cleanup warnings prevented completion".to_string()), + ); + } else if should_preserve_rebalance_stopped_state( meta_stopped, pool_stat.info.status, &terminal_event, @@ -513,6 +527,7 @@ impl ECStore { }; if !has_deferred_rebalance_error(pool_stat) + && !has_rebalance_cleanup_warnings(pool_stat) && rebalance_goal_reached( pool_stat.init_free_space, pool_stat.init_capacity, diff --git a/crates/ecstore/src/rpc/peer_rest_client.rs b/crates/ecstore/src/rpc/peer_rest_client.rs index c7cbb30d9..9173e7f50 100644 --- a/crates/ecstore/src/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/rpc/peer_rest_client.rs @@ -29,12 +29,13 @@ use rustfs_madmin::{ }; use rustfs_protos::evict_failed_connection; use rustfs_protos::proto_gen::node_service::{ - DeleteBucketMetadataRequest, DeletePolicyRequest, DeleteServiceAccountRequest, DeleteUserRequest, GetCpusRequest, - GetLiveEventsRequest, GetMemInfoRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, - GetProcInfoRequest, GetSeLinuxInfoRequest, GetSysConfigRequest, GetSysErrorsRequest, LoadBucketMetadataRequest, - LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, - LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, - ReloadSiteReplicationConfigRequest, ServerInfoRequest, SignalServiceRequest, StartProfilingRequest, StopRebalanceRequest, + CancelDecommissionRequest, ClearDecommissionRequest, DeleteBucketMetadataRequest, DeletePolicyRequest, + DeleteServiceAccountRequest, DeleteUserRequest, GetCpusRequest, GetLiveEventsRequest, GetMemInfoRequest, GetMetricsRequest, + GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, GetSysConfigRequest, + GetSysErrorsRequest, LoadBucketMetadataRequest, LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest, + LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, + LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ServerInfoRequest, + SignalServiceRequest, StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, node_service_client::NodeServiceClient, }; use rustfs_utils::XHost; @@ -972,6 +973,79 @@ impl PeerRestClient { .await } + pub async fn start_decommission(&self, pool_indices: Vec) -> Result<()> { + self.finalize_result( + async { + let pool_indices = pool_indices + .into_iter() + .map(|idx| { + u32::try_from(idx).map_err(|_| Error::other(format!("decommission pool index {idx} exceeds RPC range"))) + }) + .collect::>>()?; + let mut client = self.get_client().await?; + let request = Request::new(StartDecommissionRequest { pool_indices }); + + let response = client.start_decommission(request).await?.into_inner(); + if !response.success { + if let Some(msg) = response.error_info { + return Err(Error::other(msg)); + } + return Err(Error::other("")); + } + + Ok(()) + } + .await, + ) + .await + } + + pub async fn decommission_cancel(&self, pool_index: usize) -> Result<()> { + self.finalize_result( + async { + let pool_index = u32::try_from(pool_index) + .map_err(|_| Error::other(format!("decommission pool index {pool_index} exceeds RPC range")))?; + let mut client = self.get_client().await?; + let request = Request::new(CancelDecommissionRequest { pool_index }); + + let response = client.cancel_decommission(request).await?.into_inner(); + if !response.success { + if let Some(msg) = response.error_info { + return Err(Error::other(msg)); + } + return Err(Error::other("")); + } + + Ok(()) + } + .await, + ) + .await + } + + pub async fn clear_decommission(&self, pool_index: usize) -> Result<()> { + self.finalize_result( + async { + let pool_index = u32::try_from(pool_index) + .map_err(|_| Error::other(format!("decommission pool index {pool_index} exceeds RPC range")))?; + let mut client = self.get_client().await?; + let request = Request::new(ClearDecommissionRequest { pool_index }); + + let response = client.clear_decommission(request).await?.into_inner(); + if !response.success { + if let Some(msg) = response.error_info { + return Err(Error::other(msg)); + } + return Err(Error::other("")); + } + + Ok(()) + } + .await, + ) + .await + } + pub async fn load_transition_tier_config(&self) -> Result<()> { self.finalize_result( async { diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 9f201b054..e3b123893 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -1104,6 +1104,42 @@ pub struct LoadRebalanceMetaResponse { #[prost(string, optional, tag = "2")] pub error_info: ::core::option::Option<::prost::alloc::string::String>, } +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct StartDecommissionRequest { + #[prost(uint32, repeated, tag = "1")] + pub pool_indices: ::prost::alloc::vec::Vec, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct StartDecommissionResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(string, optional, tag = "2")] + pub error_info: ::core::option::Option<::prost::alloc::string::String>, +} +#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] +pub struct CancelDecommissionRequest { + #[prost(uint32, tag = "1")] + pub pool_index: u32, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct CancelDecommissionResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(string, optional, tag = "2")] + pub error_info: ::core::option::Option<::prost::alloc::string::String>, +} +#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ClearDecommissionRequest { + #[prost(uint32, tag = "1")] + pub pool_index: u32, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ClearDecommissionResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(string, optional, tag = "2")] + pub error_info: ::core::option::Option<::prost::alloc::string::String>, +} #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] pub struct LoadTransitionTierConfigRequest {} #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] @@ -2356,6 +2392,51 @@ pub mod node_service_client { .insert(GrpcMethod::new("node_service.NodeService", "LoadRebalanceMeta")); self.inner.unary(req, path, codec).await } + pub async fn start_decommission( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/StartDecommission"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "StartDecommission")); + self.inner.unary(req, path, codec).await + } + pub async fn cancel_decommission( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/CancelDecommission"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "CancelDecommission")); + self.inner.unary(req, path, codec).await + } + pub async fn clear_decommission( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/ClearDecommission"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "ClearDecommission")); + self.inner.unary(req, path, codec).await + } pub async fn load_transition_tier_config( &mut self, request: impl tonic::IntoRequest, @@ -2715,6 +2796,18 @@ pub mod node_service_server { &self, request: tonic::Request, ) -> std::result::Result, tonic::Status>; + async fn start_decommission( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn cancel_decommission( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn clear_decommission( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; async fn load_transition_tier_config( &self, request: tonic::Request, @@ -4927,6 +5020,90 @@ pub mod node_service_server { }; Box::pin(fut) } + "/node_service.NodeService/StartDecommission" => { + #[allow(non_camel_case_types)] + struct StartDecommissionSvc(pub Arc); + impl tonic::server::UnaryService for StartDecommissionSvc { + type Response = super::StartDecommissionResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::start_decommission(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = StartDecommissionSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.NodeService/CancelDecommission" => { + #[allow(non_camel_case_types)] + struct CancelDecommissionSvc(pub Arc); + impl tonic::server::UnaryService for CancelDecommissionSvc { + type Response = super::CancelDecommissionResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::cancel_decommission(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = CancelDecommissionSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.NodeService/ClearDecommission" => { + #[allow(non_camel_case_types)] + struct ClearDecommissionSvc(pub Arc); + impl tonic::server::UnaryService for ClearDecommissionSvc { + type Response = super::ClearDecommissionResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::clear_decommission(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = ClearDecommissionSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } "/node_service.NodeService/LoadTransitionTierConfig" => { #[allow(non_camel_case_types)] struct LoadTransitionTierConfigSvc(pub Arc); diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index e0b6c9778..42bdbcacc 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -782,6 +782,33 @@ message LoadRebalanceMetaResponse { optional string error_info = 2; } +message StartDecommissionRequest { + repeated uint32 pool_indices = 1; +} + +message StartDecommissionResponse { + bool success = 1; + optional string error_info = 2; +} + +message CancelDecommissionRequest { + uint32 pool_index = 1; +} + +message CancelDecommissionResponse { + bool success = 1; + optional string error_info = 2; +} + +message ClearDecommissionRequest { + uint32 pool_index = 1; +} + +message ClearDecommissionResponse { + bool success = 1; + optional string error_info = 2; +} + message LoadTransitionTierConfigRequest {} message LoadTransitionTierConfigResponse { @@ -895,6 +922,9 @@ service NodeService { rpc ReloadPoolMeta(ReloadPoolMetaRequest) returns (ReloadPoolMetaResponse) {}; rpc StopRebalance(StopRebalanceRequest) returns (StopRebalanceResponse) {}; rpc LoadRebalanceMeta(LoadRebalanceMetaRequest) returns (LoadRebalanceMetaResponse) {}; + rpc StartDecommission(StartDecommissionRequest) returns (StartDecommissionResponse) {}; + rpc CancelDecommission(CancelDecommissionRequest) returns (CancelDecommissionResponse) {}; + rpc ClearDecommission(ClearDecommissionRequest) returns (ClearDecommissionResponse) {}; rpc LoadTransitionTierConfig(LoadTransitionTierConfigRequest) returns (LoadTransitionTierConfigResponse) {}; rpc GetLiveEvents(GetLiveEventsRequest) returns (GetLiveEventsResponse) {}; } diff --git a/docs/architecture/decommission-compatibility.md b/docs/architecture/decommission-compatibility.md index 249ee903f..f01d134a9 100644 --- a/docs/architecture/decommission-compatibility.md +++ b/docs/architecture/decommission-compatibility.md @@ -33,6 +33,11 @@ being moved while still allowing a request to contain later targets whose leader are different nodes. Later queued targets are recovered or promoted by the leader that owns that target. +Admin start, cancel, and clear requests may arrive on any cluster node. When the +target pool first endpoint is remote, RustFS forwards the operation over the +authenticated internode RPC channel to that first endpoint. The receiving node +still enforces the local-leader rule before mutating decommission state. + ### Persisted Metadata Shape The queue is persisted in pool metadata and decoded with the rest of @@ -155,6 +160,7 @@ The queued multi-pool contract is guarded by: - `test_contextualized_decommission_start_request_allows_multiple_target_pools` - `test_decommission_start_local_leader_allows_remote_queued_pool` - `test_local_decommission_queue_prefix_stops_at_remote_leader` +- `test_decommission_peer_target_returns_none_for_local_first_endpoint` - `test_pool_meta_queued_decommission_is_not_suspended_until_promoted` - `test_pool_meta_promoted_queued_decommission_can_be_canceled` - `test_first_resumable_decommission_queue_indices_stops_at_failed_or_canceled_state` diff --git a/rustfs/src/admin/handlers/pools.rs b/rustfs/src/admin/handlers/pools.rs index f2c8b2544..bf0d379c7 100644 --- a/rustfs/src/admin/handlers/pools.rs +++ b/rustfs/src/admin/handlers/pools.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use http::{HeaderMap, StatusCode, Uri}; +use http::{HeaderMap, HeaderValue, StatusCode, Uri}; use matchit::Params; use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_utils::{ @@ -26,11 +26,12 @@ use tracing::{error, info, warn}; use crate::{ admin::{ + EndpointServerPools, PeerRestClient, auth::validate_admin_request, router::{AdminOperation, Operation, S3Router}, }, app::admin_usecase::{DefaultAdminUsecase, QueryPoolStatusRequest}, - app::context::{resolve_endpoints_handle, resolve_object_store_handle}, + app::context::{resolve_endpoints_handle, resolve_notification_system, resolve_object_store_handle}, auth::{check_key_valid, get_session_token}, error::ApiError, server::{ADMIN_PREFIX, RemoteAddr}, @@ -213,6 +214,84 @@ fn validate_start_decommission_guards(decommission_running: bool, rebalance_runn Ok(()) } +#[cfg(test)] +fn validate_pool_mutation_leader( + endpoints: &EndpointServerPools, + idx: usize, + operation: &str, + audit: PoolAuditContext<'_>, +) -> s3s::S3Result<()> { + let endpoint = endpoints + .as_ref() + .get(idx) + .and_then(|pool| pool.endpoints.as_ref().first()) + .ok_or_else(|| pool_admin_pool_index_error_with_audit(operation, idx, endpoints.as_ref().len(), audit))?; + + if !endpoint.is_local { + log_pool_request_rejected_with_index_audit( + operation_to_event(operation), + "not_pool_leader", + idx, + endpoints.as_ref().len(), + audit, + ); + return Err(S3Error::with_message( + S3ErrorCode::OperationAborted, + format!("Failed to {operation}: pool {idx} must be handled by its first endpoint {endpoint}"), + )); + } + + Ok(()) +} + +fn decommission_peer_target( + endpoints: &EndpointServerPools, + idx: usize, + operation: &str, + audit: PoolAuditContext<'_>, +) -> s3s::S3Result> { + let endpoint = endpoints + .as_ref() + .get(idx) + .and_then(|pool| pool.endpoints.as_ref().first()) + .ok_or_else(|| pool_admin_pool_index_error_with_audit(operation, idx, endpoints.as_ref().len(), audit))?; + + if endpoint.is_local { + return Ok(None); + } + + let grid_host = endpoint.grid_host(); + let Some(notification_sys) = resolve_notification_system() else { + log_pool_request_rejected_with_index_audit( + operation_to_event(operation), + "notification_sys_not_initialized", + idx, + endpoints.as_ref().len(), + audit, + ); + return Err(S3Error::with_message( + S3ErrorCode::OperationAborted, + format!("Failed to {operation}: target pool first endpoint is not reachable"), + )); + }; + + let Some(client) = notification_sys.peer_client_for_grid_host(&grid_host) else { + log_pool_request_rejected_with_index_audit( + operation_to_event(operation), + "target_peer_not_found", + idx, + endpoints.as_ref().len(), + audit, + ); + return Err(S3Error::with_message( + S3ErrorCode::OperationAborted, + format!("Failed to {operation}: target pool first endpoint is not reachable"), + )); + }; + + Ok(Some(client)) +} + fn contextualize_admin_pool_api_error( err: crate::error::ApiError, operation: &str, @@ -304,6 +383,7 @@ fn operation_to_event(operation: &str) -> &'static str { match operation { "list pools" => "list_pools", "load pool status" => "query_pool_status", + "load decommission status" => "query_decommission_status", "start decommission" => "start_decommission", "cancel decommission" => "cancel_decommission", "clear decommission" => "clear_decommission", @@ -339,6 +419,12 @@ pub fn register_pool_route(r: &mut S3Router) -> std::io::Result< AdminOperation(&StatusPool {}), )?; + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v3/decommission/status").as_str(), + AdminOperation(&StatusDecommission {}), + )?; + r.insert( Method::POST, format!("{}{}", ADMIN_PREFIX, "/v3/pools/decommission").as_str(), @@ -508,6 +594,62 @@ impl Operation for StatusPool { } } +pub struct StatusDecommission {} + +#[async_trait::async_trait] +impl Operation for StatusDecommission { + // GET //decommission/status[?pool=http://server{1...4}/disk{1...4}] + #[tracing::instrument(skip_all)] + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let Some(input_cred) = req.credentials else { + return Err(pool_admin_missing_credentials_error("load decommission status")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![ + Action::AdminAction(AdminAction::ServerInfoAdminAction), + Action::AdminAction(AdminAction::DecommissionAdminAction), + ], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await?; + + let query = parse_status_pool_query(&req.uri).map_err(|_| pool_admin_query_parse_error("load decommission status"))?; + + let usecase = DefaultAdminUsecase::from_global(); + let data = if query.pool.is_empty() { + let status = usecase.execute_list_decommission_status().await.map_err(S3Error::from)?; + serde_json::to_vec(&status) + } else { + let status = usecase + .execute_query_decommission_status(QueryPoolStatusRequest { + pool: query.pool, + by_id: query.by_id.as_str() == "true", + }) + .await + .map_err(S3Error::from)?; + serde_json::to_vec(&status) + } + .map_err(|e| { + log_pool_request_failed!("query_decommission_status", "serialize_decommission_status_failed", e); + S3Error::with_message(S3ErrorCode::InternalError, "serialize decommission status failed") + })?; + + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); + log_pool_response_emitted!("query_decommission_status"); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), header)) + } +} + pub struct StartDecommission {} #[async_trait::async_trait] @@ -572,27 +714,6 @@ impl Operation for StartDecommission { return Err(decommission_admin_not_initialized_error_with_audit("start decommission", audit)); }; - let decommission_running = store.is_decommission_running().await; - let rebalance_running = store.is_rebalance_started().await; - if decommission_running { - log_pool_request_rejected_with_context( - "start_decommission", - "decommission_already_running", - &request_id, - &actor, - &remote_addr, - ); - } else if rebalance_running { - log_pool_request_rejected_with_context( - "start_decommission", - "rebalance_in_progress", - &request_id, - &actor, - &remote_addr, - ); - } - validate_start_decommission_guards(decommission_running, rebalance_running)?; - let query = parse_mutation_pool_query(&req.uri) .map_err(|_| pool_admin_query_parse_error_with_audit("start decommission", audit))?; let is_byid = query.by_id.as_str() == "true"; @@ -631,13 +752,49 @@ impl Operation for StartDecommission { } let pools_indices = parsed_indices; - if !pools_indices.is_empty() { + if let Some(first_idx) = pools_indices.first().copied() + && let Some(client) = decommission_peer_target(&endpoints, first_idx, "start decommission", audit)? + { let pool_context = format!("pools {:?}", &pools_indices); - store - .decommission(ctx.clone(), pools_indices.clone()) + client + .start_decommission(pools_indices.clone()) .await .map_err(ApiError::from) .map_err(|err| contextualize_admin_pool_api_error(err, "start decommission", &pool_context))?; + } else { + store + .load_rebalance_meta() + .await + .map_err(|e| s3_error!(InternalError, "failed to refresh rebalance metadata before decommission start: {}", e))?; + let decommission_running = store.is_decommission_running().await; + let rebalance_running = store.is_rebalance_started().await; + if decommission_running { + log_pool_request_rejected_with_context( + "start_decommission", + "decommission_already_running", + &request_id, + &actor, + &remote_addr, + ); + } else if rebalance_running { + log_pool_request_rejected_with_context( + "start_decommission", + "rebalance_in_progress", + &request_id, + &actor, + &remote_addr, + ); + } + validate_start_decommission_guards(decommission_running, rebalance_running)?; + + if !pools_indices.is_empty() { + let pool_context = format!("pools {:?}", &pools_indices); + store + .decommission(ctx.clone(), pools_indices.clone()) + .await + .map_err(ApiError::from) + .map_err(|err| contextualize_admin_pool_api_error(err, "start decommission", &pool_context))?; + } } info!( @@ -734,15 +891,23 @@ impl Operation for CancelDecommission { return Err(pool_admin_pool_not_found_error_with_audit("cancel decommission", &query.pool, audit)); }; - let Some(store) = resolve_object_store_handle() else { - return Err(decommission_admin_not_initialized_error_with_audit("cancel decommission", audit)); - }; + if let Some(client) = decommission_peer_target(&endpoints, idx, "cancel decommission", audit)? { + client + .decommission_cancel(idx) + .await + .map_err(ApiError::from) + .map_err(|err| contextualize_admin_pool_api_error(err, "cancel decommission", format!("pool {idx}")))?; + } else { + let Some(store) = resolve_object_store_handle() else { + return Err(decommission_admin_not_initialized_error_with_audit("cancel decommission", audit)); + }; - store - .decommission_cancel(idx) - .await - .map_err(ApiError::from) - .map_err(|err| contextualize_admin_pool_api_error(err, "cancel decommission", format!("pool {idx}")))?; + store + .decommission_cancel(idx) + .await + .map_err(ApiError::from) + .map_err(|err| contextualize_admin_pool_api_error(err, "cancel decommission", format!("pool {idx}")))?; + } info!( event = EVENT_ADMIN_RESPONSE_EMITTED, @@ -839,15 +1004,23 @@ impl Operation for ClearDecommission { return Err(pool_admin_pool_not_found_error_with_audit("clear decommission", &query.pool, audit)); }; - let Some(store) = resolve_object_store_handle() else { - return Err(decommission_admin_not_initialized_error_with_audit("clear decommission", audit)); - }; + if let Some(client) = decommission_peer_target(&endpoints, idx, "clear decommission", audit)? { + client + .clear_decommission(idx) + .await + .map_err(ApiError::from) + .map_err(|err| contextualize_admin_pool_api_error(err, "clear decommission", format!("pool {idx}")))?; + } else { + let Some(store) = resolve_object_store_handle() else { + return Err(decommission_admin_not_initialized_error_with_audit("clear decommission", audit)); + }; - store - .clear_decommission(idx) - .await - .map_err(ApiError::from) - .map_err(|err| contextualize_admin_pool_api_error(err, "clear decommission", format!("pool {idx}")))?; + store + .clear_decommission(idx) + .await + .map_err(ApiError::from) + .map_err(|err| contextualize_admin_pool_api_error(err, "clear decommission", format!("pool {idx}")))?; + } info!( event = EVENT_ADMIN_RESPONSE_EMITTED, @@ -870,12 +1043,26 @@ impl Operation for ClearDecommission { mod pools_handler_tests { use super::{ PoolAuditContext, contextualize_admin_pool_api_error, decommission_admin_not_initialized_error_with_audit, - has_duplicate_indices, parse_mutation_pool_query, parse_pool_idx_by_id, parse_status_pool_query, - pool_admin_missing_credentials_error, pool_admin_missing_credentials_error_with_request, + decommission_peer_target, has_duplicate_indices, parse_mutation_pool_query, parse_pool_idx_by_id, + parse_status_pool_query, pool_admin_missing_credentials_error, pool_admin_missing_credentials_error_with_request, pool_admin_pool_index_error_with_audit, pool_admin_pool_not_found_error_with_audit, pool_admin_pool_parse_error_with_audit, pool_admin_query_parse_error, pool_admin_query_parse_error_with_audit, - validate_start_decommission_guards, + validate_pool_mutation_leader, validate_start_decommission_guards, }; + use crate::admin::{Endpoint, EndpointServerPools, Endpoints, PoolEndpoints}; + + fn test_pool_endpoints(is_local: bool) -> EndpointServerPools { + let mut endpoint = Endpoint::try_from("http://127.0.0.1:9000/disk").expect("test endpoint should parse"); + endpoint.is_local = is_local; + EndpointServerPools::from(vec![PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 1, + cmd_line: "http://127.0.0.1:9000/disk".to_string(), + endpoints: Endpoints::from(vec![endpoint]), + platform: String::new(), + }]) + } #[test] fn test_parse_pool_idx_by_id_rejects_non_numeric() { @@ -969,6 +1156,41 @@ mod pools_handler_tests { assert!(validate_start_decommission_guards(false, false).is_ok()); } + #[test] + fn test_validate_pool_mutation_leader_allows_local_first_endpoint() { + let endpoints = test_pool_endpoints(true); + let audit = PoolAuditContext::new("request", "actor", "remote"); + + assert!(validate_pool_mutation_leader(&endpoints, 0, "cancel decommission", audit).is_ok()); + } + + #[test] + fn test_decommission_peer_target_returns_none_for_local_first_endpoint() { + let endpoints = test_pool_endpoints(true); + let audit = PoolAuditContext::new("request", "actor", "remote"); + + let target = decommission_peer_target(&endpoints, 0, "start decommission", audit) + .expect("local first endpoint should resolve without peer lookup"); + + assert!(target.is_none()); + } + + #[test] + fn test_validate_pool_mutation_leader_rejects_remote_first_endpoint() { + let endpoints = test_pool_endpoints(false); + let audit = PoolAuditContext::new("request", "actor", "remote"); + + let err = validate_pool_mutation_leader(&endpoints, 0, "cancel decommission", audit) + .expect_err("remote first endpoint should reject mutation"); + + assert_eq!(err.code(), &s3s::S3ErrorCode::OperationAborted); + assert!( + err.message() + .expect("rejection should include message") + .contains("must be handled by its first endpoint") + ); + } + #[test] fn test_contextualize_admin_pool_api_error_preserves_code_and_adds_pool_context() { let err = crate::error::ApiError { diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index b7597008b..cb9d7548a 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -323,8 +323,8 @@ fn rebalance_used_pct(total: u64, available: u64) -> f64 { (total - bounded_available) as f64 / total as f64 } -fn rebalance_remaining_buckets(buckets: usize, rebalanced_buckets: usize) -> usize { - buckets.saturating_sub(rebalanced_buckets) +fn rebalance_remaining_buckets(buckets: usize, _rebalanced_buckets: usize) -> usize { + buckets } fn rebalance_pool_used(disk_stats: &[DiskStat], idx: usize) -> f64 { @@ -1090,7 +1090,7 @@ mod rebalance_handler_tests { assert_eq!(progress.num_objects, 3); assert_eq!(progress.num_versions, 5); assert_eq!(progress.bytes, 100); - assert_eq!(progress.remaining_buckets, 2); + assert_eq!(progress.remaining_buckets, 3); assert_eq!(progress.bucket, "bucket-b"); assert_eq!(progress.object, "obj-1"); assert_eq!(progress.elapsed, 50); @@ -1157,9 +1157,9 @@ mod rebalance_handler_tests { } #[test] - fn test_rebalance_remaining_buckets_is_saturating_sub() { - assert_eq!(rebalance_remaining_buckets(10, 7), 3); - assert_eq!(rebalance_remaining_buckets(3, 10), 0); + fn test_rebalance_remaining_buckets_uses_pending_queue_len() { + assert_eq!(rebalance_remaining_buckets(10, 7), 10); + assert_eq!(rebalance_remaining_buckets(3, 10), 3); } #[test] @@ -1227,7 +1227,7 @@ mod rebalance_handler_tests { assert_eq!(active.used, 0.5); assert_eq!(active.progress.as_ref().unwrap().bucket, "bucket-b"); assert_eq!(active.progress.as_ref().unwrap().object, "obj-2"); - assert_eq!(active.progress.as_ref().unwrap().remaining_buckets, 1); + assert_eq!(active.progress.as_ref().unwrap().remaining_buckets, 2); let inactive = &statuses[1]; assert_eq!(inactive.id, 1); diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index 772bddc20..870254e98 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -235,6 +235,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ DECOMMISSION, RouteRiskLevel::High, ), + admin( + HttpMethod::Get, + "/rustfs/admin/v3/decommission/status", + DECOMMISSION, + RouteRiskLevel::Sensitive, + ), admin(HttpMethod::Post, "/rustfs/admin/v3/pools/cancel", DECOMMISSION, RouteRiskLevel::High), admin(HttpMethod::Post, "/rustfs/admin/v3/pools/clear", DECOMMISSION, RouteRiskLevel::High), admin(HttpMethod::Post, "/rustfs/admin/v3/rebalance/start", REBALANCE, RouteRiskLevel::High), diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 874d6ff66..6fd9d591d 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -173,6 +173,7 @@ fn expected_admin_route_matrix() -> Vec { admin_route(Method::GET, "/v3/metrics"), admin_route(Method::GET, "/v3/pools/list"), admin_route(Method::GET, "/v3/pools/status"), + admin_route(Method::GET, "/v3/decommission/status"), admin_route(Method::POST, "/v3/pools/decommission"), admin_route(Method::POST, "/v3/pools/cancel"), admin_route(Method::POST, "/v3/pools/clear"), @@ -1036,6 +1037,7 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::GET, &admin_path("/v3/metrics")); assert_route(&router, Method::GET, &admin_path("/v3/pools/list")); + assert_route(&router, Method::GET, &admin_path("/v3/decommission/status")); assert_route(&router, Method::POST, &admin_path("/v3/rebalance/start")); assert_route(&router, Method::GET, &admin_path("/v3/rebalance/status")); assert_route(&router, Method::POST, &admin_path("/v3/heal/")); diff --git a/rustfs/src/app/admin_usecase.rs b/rustfs/src/app/admin_usecase.rs index 0d85c9264..a54938797 100644 --- a/rustfs/src/app/admin_usecase.rs +++ b/rustfs/src/app/admin_usecase.rs @@ -17,7 +17,7 @@ use super::ECStore; use super::EndpointServerPools; use super::get_server_info; -use super::{PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free}; +use super::{PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity, get_total_usable_capacity_free}; use super::{apply_bucket_usage_memory_overlay, load_data_usage_from_backend}; use crate::app::context::{AppContext, get_global_app_context, resolve_object_store_handle_for_context}; use crate::capacity::resolve_admin_used_capacity; @@ -81,6 +81,8 @@ pub struct AdminPoolDecommissionInfo { pub prefix: String, #[serde(rename = "object")] pub object: String, + #[serde(rename = "stage")] + pub stage: String, #[serde(rename = "objectsDecommissioned")] pub items_decommissioned: usize, #[serde(rename = "objectsDecommissionedFailed")] @@ -111,24 +113,58 @@ pub struct AdminPoolStatus { pub used: f64, #[serde(rename = "status")] pub status: String, + #[serde(rename = "decommissionStatus")] + pub decommission_status: String, + #[serde(rename = "rebalanceStatus")] + pub rebalance_status: String, #[serde(rename = "decommissionInfo")] pub decommission: Option, } pub type AdminPoolListItem = AdminPoolStatus; +#[derive(Debug, Clone, serde::Serialize)] +pub struct AdminDecommissionPoolStatus { + #[serde(rename = "id")] + pub id: usize, + #[serde(rename = "cmdline")] + pub cmd_line: String, + #[serde(rename = "status")] + pub status: String, + #[serde(rename = "poolStatus")] + pub pool_status: String, + #[serde(rename = "decommissionInfo")] + pub decommission: Option, +} + +#[derive(Debug, Clone, serde::Serialize)] +pub struct AdminDecommissionStatus { + #[serde(rename = "pools")] + pub pools: Vec, +} + #[derive(Clone, Default)] pub struct DefaultAdminUsecase { context: Option>, } impl DefaultAdminUsecase { - const POOL_STATUS_ACTIVE: &'static str = "active"; const POOL_STATUS_CANCELED: &'static str = "canceled"; const POOL_STATUS_COMPLETE: &'static str = "complete"; const POOL_STATUS_FAILED: &'static str = "failed"; const POOL_STATUS_QUEUED: &'static str = "queued"; const POOL_STATUS_RUNNING: &'static str = "running"; + const POOL_STATUS_UNKNOWN: &'static str = "unknown"; + const POOL_STATE_ACTIVE: &'static str = "active"; + const POOL_STATE_BLOCKED: &'static str = "blocked"; + const POOL_STATE_DECOMMISSIONED: &'static str = "decommissioned"; + const POOL_STATE_DECOMMISSIONING: &'static str = "decommissioning"; + const REBALANCE_STATUS_COMPLETED: &'static str = "completed"; + const REBALANCE_STATUS_FAILED: &'static str = "failed"; + const REBALANCE_STATUS_NONE: &'static str = "none"; + const REBALANCE_STATUS_STARTED: &'static str = "started"; + const REBALANCE_STATUS_STOPPING: &'static str = "stopping"; + const REBALANCE_STATUS_STOPPED: &'static str = "stopped"; #[cfg(test)] pub fn without_context() -> Self { @@ -277,19 +313,19 @@ impl DefaultAdminUsecase { } pub async fn execute_list_pools(&self) -> AdminUsecaseResult> { + let Some(store) = self.object_store() else { + return Err(Self::app_error(S3ErrorCode::InternalError, "Not init")); + }; let pool_statuses = self.execute_list_pool_statuses().await?; - Ok(pool_statuses.into_iter().map(Self::pool_list_item_from_status).collect()) + let mut items = Vec::with_capacity(pool_statuses.len()); + for status in pool_statuses { + let rebalance_status = store.pool_rebalance_status(status.id).await; + items.push(Self::pool_list_item_from_status(status, rebalance_status)); + } + Ok(items) } - pub async fn execute_query_pool_status(&self, req: QueryPoolStatusRequest) -> AdminUsecaseResult { - let Some(endpoints) = self.endpoints() else { - return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); - }; - - if endpoints.legacy() { - return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); - } - + fn resolve_pool_index(&self, req: &QueryPoolStatusRequest, endpoints: &EndpointServerPools) -> AdminUsecaseResult { let has_idx = if req.by_id { Self::parse_pool_idx_by_id(&req.pool, endpoints.as_ref().len()) } else { @@ -301,18 +337,62 @@ impl DefaultAdminUsecase { return Err(Self::app_error_default(S3ErrorCode::InvalidArgument)); }; + Ok(idx) + } + + pub async fn execute_query_pool_status(&self, req: QueryPoolStatusRequest) -> AdminUsecaseResult { + let Some(endpoints) = self.endpoints() else { + return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); + }; + + if endpoints.legacy() { + return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); + } + + let idx = self.resolve_pool_index(&req, &endpoints)?; + let Some(store) = self.object_store() else { return Err(Self::app_error(S3ErrorCode::InternalError, "Not init")); }; - store - .status(idx) - .await - .map(Self::pool_list_item_from_status) - .map_err(ApiError::from) + let status = store.status(idx).await.map_err(ApiError::from)?; + let rebalance_status = store.pool_rebalance_status(idx).await; + Ok(Self::pool_list_item_from_status(status, rebalance_status)) } - fn pool_list_item_from_status(status: PoolStatus) -> AdminPoolListItem { + pub async fn execute_list_decommission_status(&self) -> AdminUsecaseResult { + let pool_statuses = self.execute_list_pool_statuses().await?; + Ok(AdminDecommissionStatus { + pools: pool_statuses + .into_iter() + .map(Self::decommission_pool_status_from_status) + .collect(), + }) + } + + pub async fn execute_query_decommission_status( + &self, + req: QueryPoolStatusRequest, + ) -> AdminUsecaseResult { + let Some(endpoints) = self.endpoints() else { + return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); + }; + + if endpoints.legacy() { + return Err(Self::app_error_default(S3ErrorCode::NotImplemented)); + } + + let idx = self.resolve_pool_index(&req, &endpoints)?; + + let Some(store) = self.object_store() else { + return Err(Self::app_error(S3ErrorCode::InternalError, "Not init")); + }; + + let status = store.status(idx).await.map_err(ApiError::from)?; + Ok(Self::decommission_pool_status_from_status(status)) + } + + fn pool_list_item_from_status(status: PoolStatus, rebalance_status: (RebalStatus, bool)) -> AdminPoolListItem { let PoolStatus { id, cmd_line, @@ -322,6 +402,7 @@ impl DefaultAdminUsecase { let total_size = decommission.as_ref().map(|info| info.total_size).unwrap_or_default(); let current_size = decommission.as_ref().map(|info| info.current_size).unwrap_or_default(); let used_size = total_size.saturating_sub(current_size); + let pool_state = Self::pool_lifecycle_state(decommission.as_ref()); AdminPoolStatus { id, @@ -331,19 +412,64 @@ impl DefaultAdminUsecase { current_size, used_size, used: Self::used_ratio(total_size, used_size), - status: Self::pool_list_status(decommission.as_ref()).to_string(), + status: pool_state.to_string(), + decommission_status: Self::pool_decommission_status(decommission.as_ref()).to_string(), + rebalance_status: Self::pool_rebalance_status(rebalance_status).to_string(), decommission: decommission.map(Self::admin_decommission_info_from_pool), } } - fn pool_list_status(decommission: Option<&PoolDecommissionInfo>) -> &'static str { + fn pool_lifecycle_state(decommission: Option<&PoolDecommissionInfo>) -> &'static str { + match decommission { + Some(info) if info.complete => Self::POOL_STATE_DECOMMISSIONED, + Some(info) if info.failed || info.canceled => Self::POOL_STATE_BLOCKED, + Some(_) => Self::POOL_STATE_DECOMMISSIONING, + None => Self::POOL_STATE_ACTIVE, + } + } + + fn pool_decommission_status(decommission: Option<&PoolDecommissionInfo>) -> &'static str { match decommission { Some(info) if info.complete => Self::POOL_STATUS_COMPLETE, Some(info) if info.failed => Self::POOL_STATUS_FAILED, Some(info) if info.canceled => Self::POOL_STATUS_CANCELED, Some(info) if info.queued => Self::POOL_STATUS_QUEUED, Some(info) if info.start_time.is_some() => Self::POOL_STATUS_RUNNING, - _ => Self::POOL_STATUS_ACTIVE, + Some(_) => Self::POOL_STATUS_UNKNOWN, + None => Self::REBALANCE_STATUS_NONE, + } + } + + fn pool_rebalance_status((status, stopping): (RebalStatus, bool)) -> &'static str { + if stopping { + return Self::REBALANCE_STATUS_STOPPING; + } + + match status { + RebalStatus::None => Self::REBALANCE_STATUS_NONE, + RebalStatus::Started => Self::REBALANCE_STATUS_STARTED, + RebalStatus::Completed => Self::REBALANCE_STATUS_COMPLETED, + RebalStatus::Stopped => Self::REBALANCE_STATUS_STOPPED, + RebalStatus::Failed => Self::REBALANCE_STATUS_FAILED, + } + } + + fn decommission_pool_status_from_status(status: PoolStatus) -> AdminDecommissionPoolStatus { + let PoolStatus { + id, + cmd_line, + decommission, + .. + } = status; + let pool_status = Self::pool_lifecycle_state(decommission.as_ref()).to_string(); + let status = Self::pool_decommission_status(decommission.as_ref()).to_string(); + + AdminDecommissionPoolStatus { + id, + cmd_line, + status, + pool_status, + decommission: decommission.map(Self::admin_decommission_info_from_pool), } } @@ -363,6 +489,7 @@ impl DefaultAdminUsecase { bucket: info.bucket, prefix: info.prefix, object: info.object, + stage: info.stage, items_decommissioned: info.items_decommissioned, items_decommission_failed: info.items_decommission_failed, bytes_done: info.bytes_done, @@ -446,7 +573,7 @@ mod tests { } #[test] - fn admin_pool_list_item_maps_capacity_and_active_status() { + fn admin_pool_list_item_maps_capacity_and_unknown_decommission_status() { let now = OffsetDateTime::UNIX_EPOCH; let pool = PoolStatus { id: 2, @@ -459,24 +586,29 @@ mod tests { }), }; - let item = DefaultAdminUsecase::pool_list_item_from_status(pool); + let item = DefaultAdminUsecase::pool_list_item_from_status(pool, (RebalStatus::None, false)); assert_eq!(item.id, 2); assert_eq!(item.total_size, 1_000); assert_eq!(item.current_size, 250); assert_eq!(item.used_size, 750); assert!((item.used - 0.75).abs() < f64::EPSILON); - assert_eq!(item.status, "active"); + assert_eq!(item.status, "decommissioning"); + assert_eq!(item.decommission_status, "unknown"); + assert_eq!(item.rebalance_status, "none"); } #[test] fn admin_pool_list_item_serializes_admin_api_fields() { - let item = DefaultAdminUsecase::pool_list_item_from_status(PoolStatus { - id: 1, - cmd_line: "pool-1".to_string(), - last_update: OffsetDateTime::UNIX_EPOCH, - decommission: None, - }); + let item = DefaultAdminUsecase::pool_list_item_from_status( + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + (RebalStatus::Completed, false), + ); let value = serde_json::to_value(item).unwrap(); @@ -491,6 +623,8 @@ mod tests { "usedSize": 0, "used": 0.0, "status": "active", + "decommissionStatus": "none", + "rebalanceStatus": "completed", "decommissionInfo": null }) ); @@ -509,7 +643,7 @@ mod tests { }), }; - let item = DefaultAdminUsecase::pool_list_item_from_status(pool); + let item = DefaultAdminUsecase::pool_list_item_from_status(pool, (RebalStatus::None, false)); assert_eq!(item.total_size, 100); assert_eq!(item.current_size, 150); @@ -531,33 +665,40 @@ mod tests { }), }; - let item = DefaultAdminUsecase::pool_list_item_from_status(pool); + let item = DefaultAdminUsecase::pool_list_item_from_status(pool, (RebalStatus::Started, false)); - assert_eq!(item.status, "running"); + assert_eq!(item.status, "decommissioning"); + assert_eq!(item.decommission_status, "running"); + assert_eq!(item.rebalance_status, "started"); } #[test] fn admin_pool_list_item_exposes_queued_decommission_state() { - let item = DefaultAdminUsecase::pool_list_item_from_status(PoolStatus { - id: 3, - cmd_line: "pool-3".to_string(), - last_update: OffsetDateTime::UNIX_EPOCH, - decommission: Some(PoolDecommissionInfo { - queued: true, - queued_buckets: vec!["bucket-a".to_string(), ".rustfs.sys/config".to_string()], - decommissioned_buckets: vec!["bucket-done".to_string()], - bucket: "bucket-a".to_string(), - prefix: "prefix/".to_string(), - object: "object.txt".to_string(), - items_decommissioned: 7, - items_decommission_failed: 1, - bytes_done: 1024, - bytes_failed: 64, - ..Default::default() - }), - }); + let item = DefaultAdminUsecase::pool_list_item_from_status( + PoolStatus { + id: 3, + cmd_line: "pool-3".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(PoolDecommissionInfo { + queued: true, + queued_buckets: vec!["bucket-a".to_string(), ".rustfs.sys/config".to_string()], + decommissioned_buckets: vec!["bucket-done".to_string()], + bucket: "bucket-a".to_string(), + prefix: "prefix/".to_string(), + object: "object.txt".to_string(), + stage: "migrate_object".to_string(), + items_decommissioned: 7, + items_decommission_failed: 1, + bytes_done: 1024, + bytes_failed: 64, + ..Default::default() + }), + }, + (RebalStatus::None, false), + ); - assert_eq!(item.status, "queued"); + assert_eq!(item.status, "decommissioning"); + assert_eq!(item.decommission_status, "queued"); let value = serde_json::to_value(item).expect("admin pool status should serialize"); assert_eq!(value["decommissionInfo"]["queued"], true); assert_eq!( @@ -568,6 +709,7 @@ mod tests { assert_eq!(value["decommissionInfo"]["bucket"], "bucket-a"); assert_eq!(value["decommissionInfo"]["prefix"], "prefix/"); assert_eq!(value["decommissionInfo"]["object"], "object.txt"); + assert_eq!(value["decommissionInfo"]["stage"], "migrate_object"); assert_eq!(value["decommissionInfo"]["objectsDecommissioned"], 7); assert_eq!(value["decommissionInfo"]["objectsDecommissionedFailed"], 1); assert_eq!(value["decommissionInfo"]["bytesDecommissioned"], 1024); @@ -577,28 +719,128 @@ mod tests { #[test] fn admin_pool_list_item_maps_terminal_decommission_statuses() { - let complete = DefaultAdminUsecase::pool_list_status(Some(&PoolDecommissionInfo { + let complete = DefaultAdminUsecase::pool_decommission_status(Some(&PoolDecommissionInfo { complete: true, ..Default::default() })); - let failed = DefaultAdminUsecase::pool_list_status(Some(&PoolDecommissionInfo { + let failed = DefaultAdminUsecase::pool_decommission_status(Some(&PoolDecommissionInfo { failed: true, ..Default::default() })); - let canceled = DefaultAdminUsecase::pool_list_status(Some(&PoolDecommissionInfo { + let canceled = DefaultAdminUsecase::pool_decommission_status(Some(&PoolDecommissionInfo { canceled: true, ..Default::default() })); - let queued = DefaultAdminUsecase::pool_list_status(Some(&PoolDecommissionInfo { + let queued = DefaultAdminUsecase::pool_decommission_status(Some(&PoolDecommissionInfo { queued: true, ..Default::default() })); - let idle = DefaultAdminUsecase::pool_list_status(None); + let idle = DefaultAdminUsecase::pool_decommission_status(None); assert_eq!(complete, "complete"); assert_eq!(failed, "failed"); assert_eq!(canceled, "canceled"); assert_eq!(queued, "queued"); - assert_eq!(idle, "active"); + assert_eq!(idle, "none"); + } + + #[test] + fn admin_pool_list_item_keeps_rebalance_failure_separate_from_pool_state() { + let item = DefaultAdminUsecase::pool_list_item_from_status( + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + (RebalStatus::Failed, false), + ); + + assert_eq!(item.status, "active"); + assert_eq!(item.decommission_status, "none"); + assert_eq!(item.rebalance_status, "failed"); + } + + #[test] + fn admin_pool_list_item_maps_stopping_rebalance_status() { + let item = DefaultAdminUsecase::pool_list_item_from_status( + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: None, + }, + (RebalStatus::Started, true), + ); + + assert_eq!(item.status, "active"); + assert_eq!(item.decommission_status, "none"); + assert_eq!(item.rebalance_status, "stopping"); + } + + #[test] + fn admin_pool_lifecycle_state_distinguishes_decommission_terminal_states() { + let complete = DefaultAdminUsecase::pool_lifecycle_state(Some(&PoolDecommissionInfo { + complete: true, + ..Default::default() + })); + let failed = DefaultAdminUsecase::pool_lifecycle_state(Some(&PoolDecommissionInfo { + failed: true, + ..Default::default() + })); + let canceled = DefaultAdminUsecase::pool_lifecycle_state(Some(&PoolDecommissionInfo { + canceled: true, + ..Default::default() + })); + + assert_eq!(complete, "decommissioned"); + assert_eq!(failed, "blocked"); + assert_eq!(canceled, "blocked"); + } + + #[test] + fn admin_decommission_status_serializes_task_status_and_pool_status() { + let item = DefaultAdminUsecase::decommission_pool_status_from_status(PoolStatus { + id: 3, + cmd_line: "pool-3".to_string(), + last_update: OffsetDateTime::UNIX_EPOCH, + decommission: Some(PoolDecommissionInfo { + failed: true, + ..Default::default() + }), + }); + + let value = serde_json::to_value(item).expect("decommission status should serialize"); + + assert_eq!( + value, + serde_json::json!({ + "id": 3, + "cmdline": "pool-3", + "status": "failed", + "poolStatus": "blocked", + "decommissionInfo": { + "startTime": null, + "startSize": 0, + "totalSize": 0, + "currentSize": 0, + "complete": false, + "failed": true, + "canceled": false, + "queued": false, + "queuedBuckets": [], + "decommissionedBuckets": [], + "bucket": "", + "prefix": "", + "object": "", + "stage": "", + "objectsDecommissioned": 0, + "objectsDecommissionedFailed": 0, + "bytesDecommissioned": 0, + "bytesDecommissionedFailed": 0, + "waitingReason": null + } + }) + ); } } diff --git a/rustfs/src/app/mod.rs b/rustfs/src/app/mod.rs index 260886949..f930e7fa3 100644 --- a/rustfs/src/app/mod.rs +++ b/rustfs/src/app/mod.rs @@ -105,6 +105,7 @@ pub(crate) type ObjectInfo = :: pub(crate) type ObjectOptions = ::ObjectOptions; pub(crate) type PoolDecommissionInfo = ecstore_capacity::PoolDecommissionInfo; pub(crate) type PoolStatus = ecstore_capacity::PoolStatus; +pub(crate) type RebalStatus = crate::storage::ecstore_rebalance::RebalStatus; pub(crate) type StorageError = crate::storage::StorageError; pub(crate) type Error = StorageError; pub(crate) type TierConfigMgr = crate::storage::TierConfigMgr; diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index d924cf662..2e77eed6a 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -13,11 +13,11 @@ // limitations under the License. use super::super::{ - CollectMetricsOpts, DeleteOptions, DiskError, DiskInfoOptions, DiskStore, FileInfoVersions, LocalPeerS3Client, MetricType, - PEER_RESTSIGNAL, PEER_RESTSUB_SYS, ReadMultipleReq, ReadMultipleResp, ReadOptions, SERVICE_SIGNAL_REFRESH_CONFIG, - SERVICE_SIGNAL_RELOAD_DYNAMIC, StorageDiskRpcExt as _, StoragePeerS3ClientExt as _, UpdateMetadataOpts, all_local_disk_path, - collect_local_metrics, find_local_disk_by_ref, get_local_server_property, load_bucket_metadata, - reload_transition_tier_config, resolve_object_store_handle, set_bucket_metadata, + CollectMetricsOpts, DeleteOptions, DiskError, DiskInfoOptions, DiskStore, ECStore, Error, FileInfoVersions, + LocalPeerS3Client, MetricType, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, ReadMultipleReq, ReadMultipleResp, ReadOptions, + SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, StorageDiskRpcExt as _, StoragePeerS3ClientExt as _, + UpdateMetadataOpts, all_local_disk_path, collect_local_metrics, find_local_disk_by_ref, get_local_server_property, + load_bucket_metadata, reload_transition_tier_config, resolve_object_store_handle, set_bucket_metadata, }; use crate::admin::service::{ config::{reload_dynamic_config_runtime_state, reload_runtime_config_snapshot}, @@ -46,6 +46,7 @@ use std::{collections::HashMap, io::Cursor, pin::Pin, sync::Arc}; use tokio::spawn; use tokio::sync::mpsc; use tokio_stream::wrappers::ReceiverStream; +use tokio_util::sync::CancellationToken; use tonic::{Request, Response, Status, Streaming}; use tracing::{debug, error, info, warn}; @@ -140,6 +141,23 @@ fn stop_rebalance_response(result: super::super::Result<()>) -> StopRebalanceRes } } +fn ensure_rpc_decommission_local_leader(store: &ECStore, idx: usize) -> super::super::Result<()> { + let endpoints = store.endpoints(); + let endpoint = endpoints + .as_ref() + .get(idx) + .and_then(|pool| pool.endpoints.as_ref().first()) + .ok_or_else(|| Error::other(format!("invalid decommission pool index {idx} for {} pools", endpoints.as_ref().len())))?; + + if !endpoint.is_local { + return Err(Error::other(format!( + "decommission for pool {idx} must run on the pool first endpoint {endpoint}" + ))); + } + + Ok(()) +} + #[path = "bucket.rs"] mod bucket; #[path = "disk.rs"] @@ -1057,6 +1075,101 @@ impl Node for NodeService { })) } + async fn start_decommission( + &self, + request: Request, + ) -> Result, Status> { + let Some(store) = resolve_object_store_handle() else { + return Ok(Response::new(StartDecommissionResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); + }; + + let mut indices = Vec::with_capacity(request.get_ref().pool_indices.len()); + for idx in request.into_inner().pool_indices { + indices.push( + usize::try_from(idx) + .map_err(|_| Status::invalid_argument(format!("decommission pool index {idx} exceeds local range")))?, + ); + } + + match store.decommission(CancellationToken::new(), indices).await { + Ok(()) => Ok(Response::new(StartDecommissionResponse { + success: true, + error_info: None, + })), + Err(err) => Ok(Response::new(StartDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + })), + } + } + + async fn cancel_decommission( + &self, + request: Request, + ) -> Result, Status> { + let Some(store) = resolve_object_store_handle() else { + return Ok(Response::new(CancelDecommissionResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); + }; + + let idx = usize::try_from(request.into_inner().pool_index) + .map_err(|_| Status::invalid_argument("decommission pool index exceeds local range"))?; + if let Err(err) = ensure_rpc_decommission_local_leader(&store, idx) { + return Ok(Response::new(CancelDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + })); + } + + match store.decommission_cancel(idx).await { + Ok(()) => Ok(Response::new(CancelDecommissionResponse { + success: true, + error_info: None, + })), + Err(err) => Ok(Response::new(CancelDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + })), + } + } + + async fn clear_decommission( + &self, + request: Request, + ) -> Result, Status> { + let Some(store) = resolve_object_store_handle() else { + return Ok(Response::new(ClearDecommissionResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); + }; + + let idx = usize::try_from(request.into_inner().pool_index) + .map_err(|_| Status::invalid_argument("decommission pool index exceeds local range"))?; + if let Err(err) = ensure_rpc_decommission_local_leader(&store, idx) { + return Ok(Response::new(ClearDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + })); + } + + match store.clear_decommission(idx).await { + Ok(()) => Ok(Response::new(ClearDecommissionResponse { + success: true, + error_info: None, + })), + Err(err) => Ok(Response::new(ClearDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + })), + } + } + async fn load_transition_tier_config( &self, _request: Request,