From 2c3e68ad89a0c6b47a5968883bbf4a3ab9aa0f59 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 22 Aug 2026 12:09:21 +0800 Subject: [PATCH 1/4] ci: let feature validation jobs finish (#6364) --- .github/workflows/ci.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6a0d29604..3a06e4667 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -400,7 +400,7 @@ jobs: if: github.event_name != 'pull_request' || github.event.action != 'closed' needs: [ quick-checks ] runs-on: sm-standard-4 - timeout-minutes: 45 + timeout-minutes: 90 env: FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true" steps: @@ -440,7 +440,7 @@ jobs: if: github.event_name != 'pull_request' || github.event.action != 'closed' needs: [ quick-checks ] runs-on: sm-standard-4 - timeout-minutes: 60 + timeout-minutes: 90 env: FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true" steps: @@ -470,7 +470,7 @@ jobs: if: github.event_name != 'pull_request' || github.event.action != 'closed' needs: [ quick-checks ] runs-on: sm-standard-4 - timeout-minutes: 60 + timeout-minutes: 90 strategy: # On a PR, one failing protocol leg is enough to know the PR is not ready, # so stop the sibling leg instead of paying another ~40 minutes for it. From 2e600290791174592ff3546448572f954ac60b60 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 22 Aug 2026 13:49:23 +0800 Subject: [PATCH 2/4] perf(ecstore): throttle decommission checkpoints (#6356) --- crates/ecstore/src/core/pools.rs | 404 +++++++++++++++++++++++-------- 1 file changed, 306 insertions(+), 98 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index fb983dc0c..77f4d4fea 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -89,6 +89,7 @@ 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_PROGRESS_SAVE_RETRY_BACKOFF: Duration = Duration::seconds(1); const DECOMMISSION_BUCKET_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_BUCKET_CONCURRENCY"; const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4; const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30; @@ -638,22 +639,6 @@ fn track_decommission_current_object(meta: &mut PoolMeta, idx: usize, bucket: &s 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(()) -} - fn resolve_decommission_update_after_result(result: Result) -> Result { result.map_err(|err| Error::other(format!("decommission metadata update failed: {err}"))) } @@ -1483,6 +1468,7 @@ impl TryFrom for PoolDecommissionInfo { 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), + progress_save_retry_after: None, }) } } @@ -1514,6 +1500,7 @@ impl TryFrom for PoolDecommissionInfo { terminal_reload_attempt_at: None, terminal_reload_failures: Vec::new(), progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed), + progress_save_retry_after: None, }) } } @@ -1627,6 +1614,82 @@ impl PoolMeta { } } + fn decommission_progress_checkpoint( + &self, + idx: usize, + duration: Duration, + now: OffsetDateTime, + ) -> Result> { + let pool_count = self.pools.len(); + ensure_valid_decommission_pool_index(pool_count, idx)?; + + let Some(pool) = self.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let Some(info) = pool.decommission.as_ref() else { + return Err(decommission_metadata_not_initialized_error("update decommission metadata timestamp")); + }; + + if info.progress_save_retry_after.is_some_and(|retry_after| now < retry_after) { + return Ok(None); + } + + let time_threshold_reached = now.unix_timestamp() - pool.last_update.unix_timestamp() >= duration.whole_seconds(); + let item_threshold_reached = info.items_since_last_progress_save() >= DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD; + if !time_threshold_reached && !item_threshold_reached { + return Ok(None); + } + + Ok(Some(DecommissionProgressCheckpoint { + start_time: info.start_time, + queued: info.queued, + counted_items: info.counted_items(), + checkpoint_at: now, + })) + } + + fn commit_decommission_progress_checkpoint(&mut self, idx: usize, checkpoint: DecommissionProgressCheckpoint) -> bool { + let Some(pool) = self.pools.get_mut(idx) else { + return false; + }; + let Some(info) = pool.decommission.as_mut() else { + return false; + }; + + if info.start_time != checkpoint.start_time + || info.queued != checkpoint.queued + || !is_decommission_active(info.complete, info.failed, info.canceled) + { + return false; + } + + info.progress_save_item_baseline = info.progress_save_item_baseline.max(checkpoint.counted_items); + info.progress_save_retry_after = None; + pool.last_update = pool.last_update.max(checkpoint.checkpoint_at); + true + } + + fn defer_decommission_progress_checkpoint( + &mut self, + idx: usize, + checkpoint: DecommissionProgressCheckpoint, + retry_after: OffsetDateTime, + ) { + let Some(pool) = self.pools.get_mut(idx) else { + return; + }; + let Some(info) = pool.decommission.as_mut() else { + return; + }; + + if info.start_time == checkpoint.start_time + && info.queued == checkpoint.queued + && is_decommission_active(info.complete, info.failed, info.canceled) + { + info.progress_save_retry_after = Some(retry_after); + } + } + fn load_from_config_data(&mut self, data: Vec) -> Result<()> { if data.is_empty() { return Ok(()); @@ -1987,30 +2050,9 @@ impl PoolMeta { } pub fn update_after(&mut self, idx: usize, duration: Duration) -> Result { - let pool_count = self.pools.len(); - ensure_valid_decommission_pool_index(pool_count, idx)?; - - let (last_update, item_threshold_reached) = match self.pools.get(idx) { - Some(pool) if let Some(info) = pool.decommission.as_ref() => ( - pool.last_update, - info.items_since_last_progress_save() >= DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, - ), - Some(_) => { - return Err(decommission_metadata_not_initialized_error("update decommission metadata timestamp")); - } - None => return Err(invalid_decommission_pool_index_error(pool_count, idx)), - }; - let now = OffsetDateTime::now_utc(); - - if now.unix_timestamp() - last_update.unix_timestamp() >= duration.whole_seconds() || item_threshold_reached { - let Some(pool) = self.pools.get_mut(idx) else { - return Err(invalid_decommission_pool_index_error(pool_count, idx)); - }; - pool.last_update = now; - return Ok(true); - } - - Ok(false) + Ok(self + .decommission_progress_checkpoint(idx, duration, OffsetDateTime::now_utc())? + .is_some()) } pub fn validate(&self, pools: Vec>) -> Result { @@ -2151,6 +2193,16 @@ pub struct PoolDecommissionInfo { pub terminal_reload_failures: Vec, #[serde(skip)] pub progress_save_item_baseline: usize, + #[serde(skip)] + pub progress_save_retry_after: Option, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +struct DecommissionProgressCheckpoint { + start_time: Option, + queued: bool, + counted_items: usize, + checkpoint_at: OffsetDateTime, } impl PoolDecommissionInfo { @@ -2185,6 +2237,7 @@ impl PoolDecommissionInfo { fn mark_progress_saved(&mut self) { self.progress_save_item_baseline = self.counted_items(); + self.progress_save_retry_after = None; } pub fn bucket_push(&mut self, bucket: &DecomBucketInfo) { @@ -2489,6 +2542,40 @@ impl ECStore { snapshot.save(self.pools.clone()).await } + async fn save_decommission_progress_checkpoint(&self, idx: usize) -> Result { + // Lock order: save gate, then the short pool metadata read/write sections. Peer + // reloads are intentionally performed by the caller after both locks are released. + let _save_guard = self.pool_meta_save_gate.lock().await; + let (snapshot, checkpoint) = { + let pool_meta = self.pool_meta.read().await; + let Some(checkpoint) = pool_meta.decommission_progress_checkpoint( + idx, + DECOMMISSION_PROGRESS_SAVE_INTERVAL, + OffsetDateTime::now_utc(), + )? + else { + return Ok(false); + }; + + let mut snapshot = pool_meta.clone(); + let Some(pool) = snapshot.pools.get_mut(idx) else { + return Err(invalid_decommission_pool_index_error(snapshot.pools.len(), idx)); + }; + pool.last_update = checkpoint.checkpoint_at; + (snapshot, checkpoint) + }; + + if let Err(err) = snapshot.save(self.pools.clone()).await { + let retry_after = OffsetDateTime::now_utc() + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF; + let mut pool_meta = self.pool_meta.write().await; + pool_meta.defer_decommission_progress_checkpoint(idx, checkpoint, retry_after); + return Err(err); + } + + let mut pool_meta = self.pool_meta.write().await; + Ok(pool_meta.commit_decommission_progress_checkpoint(idx, checkpoint)) + } + async fn save_current_pool_meta_for_decommission_start( &self, indices: &[usize], @@ -2871,7 +2958,7 @@ impl ECStore { Ok(()) } - async fn save_decommission_entry_progress_stage( + async fn track_decommission_entry_progress_stage( &self, idx: usize, bucket: &str, @@ -2882,22 +2969,6 @@ impl ECStore { 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(()) @@ -3165,7 +3236,7 @@ impl ECStore { let bucket_name = bucket.clone(); let object_name = rd.object_info.name.clone(); - self.save_decommission_entry_progress_stage( + self.track_decommission_entry_progress_stage( idx, bucket_name.as_str(), object_name.as_str(), @@ -3259,7 +3330,7 @@ impl ECStore { } decommission_cancel_signal_result(rx.is_cancelled())?; - self.save_decommission_entry_progress_stage( + self.track_decommission_entry_progress_stage( idx, bucket.as_str(), entry.name.as_str(), @@ -3267,7 +3338,7 @@ impl ECStore { ) .await?; - self.save_decommission_entry_progress_stage( + self.track_decommission_entry_progress_stage( idx, bucket.as_str(), entry.name.as_str(), @@ -3334,34 +3405,42 @@ impl ECStore { } }; - self.save_decommission_entry_progress_stage(idx, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_ENTRY_FINISHED) - .await?; + self.track_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) { - warn!( - event = EVENT_DECOMMISSION_ENTRY, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_POOLS, - pool_index = idx, - bucket = %bucket, - object = %entry.name, - state = "progress_save_failed", - error = %err, - "Decommission progress save failed; continuing and will retry at the next checkpoint" - ); - } else { - let mut pool_meta = self.pool_meta.write().await; - pool_meta.mark_decommission_progress_saved(); - if let Some(notification_sys) = runtime_sources::notification_sys() - && let Err(err) = resolve_decommission_entry_reload_result( - notification_sys.reload_pool_meta().await, - bucket.as_str(), - entry.name.as_str(), - ) - { - warn!("{err}"); + match self.save_decommission_progress_checkpoint(idx).await { + Ok(true) => { + if let Some(notification_sys) = runtime_sources::notification_sys() + && let Err(err) = resolve_decommission_entry_reload_result( + notification_sys.reload_pool_meta().await, + bucket.as_str(), + entry.name.as_str(), + ) + { + warn!("{err}"); + } + } + Ok(false) => {} + Err(err) => { + if let Some(err) = resolve_decommission_progress_save_result(Err(err)) { + warn!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + bucket = %bucket, + object = %entry.name, + state = "progress_save_failed", + error = %err, + "Decommission progress save failed; continuing and will retry at the next checkpoint" + ); + } } } } @@ -5264,11 +5343,11 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi #[cfg(test)] mod pools_tests { use super::{ - DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, - DecommissionStartPoolState, DecommissionTerminalState, ListCallback, 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_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, + DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, 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_keeps_active_pool, @@ -5293,9 +5372,8 @@ mod pools_tests { should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload, should_retry_decommission_listing, 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_listing_retry, wait_decommission_worker_drain, - with_decommission_entry_context, + track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request, + wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context, }; use crate::data_movement; use crate::disk::endpoint::Endpoint; @@ -6538,7 +6616,7 @@ mod pools_tests { } #[test] - fn test_touch_decommission_progress_updates_last_update_and_save_baseline() { + fn test_track_decommission_stage_does_not_advance_checkpoint_state() { let mut meta = PoolMeta { pools: vec![PoolStatus { id: 0, @@ -6553,11 +6631,13 @@ mod pools_tests { ..Default::default() }; - touch_decommission_progress(&mut meta, 0).expect("valid decommission progress should be touched"); + track_decommission_current_object_stage(&mut meta, 0, "bucket", "object", "migrate_object") + .expect("valid decommission progress should be tracked"); - assert!(meta.pools[0].last_update > OffsetDateTime::UNIX_EPOCH); + assert_eq!(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); + assert_eq!(info.items_since_last_progress_save(), 5); + assert_eq!(info.stage, "migrate_object"); } #[test] @@ -6632,6 +6712,134 @@ mod pools_tests { assert_eq!(info.items_since_last_progress_save(), 1); } + #[test] + fn test_pool_meta_update_after_does_not_advance_last_update_before_save() { + let last_update = OffsetDateTime::UNIX_EPOCH; + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update, + decommission: Some(PoolDecommissionInfo { + start_time: Some(last_update), + items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, + ..Default::default() + }), + }], + ..Default::default() + }; + + assert!( + meta.update_after(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL) + .expect("item threshold should request a checkpoint") + ); + assert_eq!(meta.pools[0].last_update, last_update); + } + + #[test] + fn test_decommission_progress_checkpoint_commits_exact_snapshot_watermark() { + let start_time = OffsetDateTime::UNIX_EPOCH; + let checkpoint_at = start_time + Duration::seconds(30); + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: start_time, + decommission: Some(PoolDecommissionInfo { + start_time: Some(start_time), + items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, + ..Default::default() + }), + }], + ..Default::default() + }; + + let checkpoint = meta + .decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at) + .expect("valid decommission state should produce a checkpoint") + .expect("item threshold should produce a checkpoint"); + meta.count_item(0, 1, false); + + assert!(meta.commit_decommission_progress_checkpoint(0, checkpoint)); + let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist"); + assert_eq!(info.progress_save_item_baseline, checkpoint.counted_items); + assert_eq!(info.items_since_last_progress_save(), 1); + assert_eq!(meta.pools[0].last_update, checkpoint_at); + } + + #[test] + fn test_decommission_progress_checkpoint_backoff_does_not_advance_baseline() { + let start_time = OffsetDateTime::UNIX_EPOCH; + let checkpoint_at = start_time + Duration::seconds(30); + let retry_after = checkpoint_at + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF; + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: start_time, + decommission: Some(PoolDecommissionInfo { + start_time: Some(start_time), + items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, + ..Default::default() + }), + }], + ..Default::default() + }; + + let checkpoint = meta + .decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at) + .expect("valid decommission state should produce a checkpoint") + .expect("item threshold should produce a checkpoint"); + meta.defer_decommission_progress_checkpoint(0, checkpoint, retry_after); + + assert!( + meta.decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at) + .expect("retry backoff check should succeed") + .is_none() + ); + assert_eq!(meta.pools[0].last_update, start_time); + assert_eq!( + meta.pools[0] + .decommission + .as_ref() + .expect("decommission info should exist") + .progress_save_item_baseline, + 0 + ); + } + + #[test] + fn test_decommission_progress_checkpoint_count_scales_with_threshold() { + let start_time = OffsetDateTime::UNIX_EPOCH; + let checkpoint_at = start_time; + let mut meta = PoolMeta { + pools: vec![PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: start_time, + decommission: Some(PoolDecommissionInfo { + start_time: Some(start_time), + ..Default::default() + }), + }], + ..Default::default() + }; + let mut checkpoint_count = 0; + + for _ in 0..(DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD * 10) { + meta.count_item(0, 1, false); + if let Some(checkpoint) = meta + .decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at) + .expect("valid decommission state should produce a checkpoint") + { + checkpoint_count += 1; + assert!(meta.commit_decommission_progress_checkpoint(0, checkpoint)); + } + } + + assert_eq!(checkpoint_count, 10); + } + #[test] fn test_ensure_decommission_not_rebalancing_rejects_running_rebalance() { let err = ensure_decommission_not_rebalancing(true).expect_err("rebalance running should be rejected"); From a34310a58f0e75a865172b1082c5d139889e5359 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 22 Aug 2026 14:12:28 +0800 Subject: [PATCH 3/4] fix(ecstore): fail closed on unresolved decommission entries (#6367) * fix(ecstore): fail closed on unresolved decommission entries * perf(ecstore): avoid successful listing name clone --- crates/ecstore/src/core/pools.rs | 225 +++++++++++++++++++++++++++---- 1 file changed, 201 insertions(+), 24 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 77f4d4fea..04a94cc05 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -36,7 +36,8 @@ use crate::disk::error::DiskError; use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; 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_version_not_found, + StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_operation_canceled, + is_err_version_not_found, }; use crate::layout::endpoints::EndpointServerPools; use crate::object_api::{GetObjectReader, ObjectOptions}; @@ -758,7 +759,76 @@ async fn load_decommission_entry_exact_versions( } 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 } + match list_result { + Ok(()) => entry_error.map_or(Ok(()), Err), + Err(list_err) => resolve_decommission_listing_error(Some(list_err), entry_error).map_or(Ok(()), Err), + } +} + +fn resolve_decommission_listing_error(listing_error: Option, entry_error: Option) -> Option { + match (listing_error, entry_error) { + (Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&listing_error) => Some(entry_error), + (Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&entry_error) => Some(listing_error), + (Some(listing_error), _) => Some(listing_error), + (None, entry_error) => entry_error, + } +} + +fn decommission_unresolved_listing_error( + bucket: &str, + prefix: &str, + candidate: Option<&str>, + candidate_count: usize, + disk_error_count: usize, + pool_index: usize, + set_index: usize, +) -> Error { + let location = candidate.unwrap_or(prefix); + Error::other(format!( + "decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))" + )) +} + +fn resolve_decommission_partial_listing_entry( + entries: MetaCacheEntries, + resolver: MetadataResolutionParams, + bucket: &str, + prefix: &str, + disk_error_count: usize, + pool_index: usize, + set_index: usize, +) -> Result { + let candidate_count = entries.as_ref().iter().flatten().count(); + if let Some(entry) = entries.resolve(resolver) { + return Ok(entry); + } + + let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next(); + Err(decommission_unresolved_listing_error( + bucket, + prefix, + candidate, + candidate_count, + disk_error_count, + pool_index, + set_index, + )) +} + +async fn record_decommission_entry_error( + entry_error: &Arc>>, + rx: &CancellationToken, + err: Error, +) { + if rx.is_cancelled() { + return; + } + + let mut first_err = entry_error.lock().await; + if first_err.is_none() && !rx.is_cancelled() { + *first_err = Some(err); + rx.cancel(); + } } fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> { @@ -3617,6 +3687,7 @@ impl ECStore { let rx_clone = rx.clone(); let bi = bi.clone(); let set_id = set_idx; + let listing_entry_error = entry_error.clone(); let worker = tokio::spawn(async move { let _listing_permit = listing_permit; run_decommission_listing_with_retry( @@ -3630,7 +3701,11 @@ impl ECStore { let set = set.clone(); let rx = rx_clone.clone(); let bucket = bi.clone(); - async move { set.list_objects_to_decommission(rx, bucket, callback).await } + let entry_error = listing_entry_error.clone(); + async move { + set.list_objects_to_decommission(rx, bucket, callback, entry_error.clone(), idx, set_id) + .await + } }, ) .await @@ -3660,11 +3735,7 @@ impl ECStore { wait_decommission_worker_drain(&workers, worker_limit).await?; - if let Some(err) = listing_worker_error { - return Err(err); - } - - if let Some(err) = entry_error.lock().await.clone() { + if let Some(err) = resolve_decommission_listing_error(listing_worker_error, entry_error.lock().await.clone()) { return Err(err); } @@ -4270,7 +4341,7 @@ impl ECStore { let buckets = self.get_buckets_to_decommission().await?; let pool = self.pools[idx].clone(); - for set in &pool.disk_set { + for (set_index, set) in pool.disk_set.iter().enumerate() { for bucket_info in &buckets { let mut lifecycle_config = None; let mut object_lock_config = None; @@ -4365,7 +4436,7 @@ impl ECStore { }); let list_result = set - .list_objects_to_decommission(callback_rx, bucket_info.clone(), callback) + .list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index) .await; let entry_error = entry_error.lock().await.clone(); resolve_decommission_check_after_list_result(list_result, entry_error)?; @@ -5100,12 +5171,15 @@ mod tests { pub type ListCallback = Arc BoxFuture<'static, ()> + Send + Sync + 'static>; impl SetDisks { - #[tracing::instrument(skip(self, rx, cb_func))] + #[tracing::instrument(skip(self, rx, cb_func, entry_error))] async fn list_objects_to_decommission( self: &Arc, rx: CancellationToken, bucket_info: DecomBucketInfo, cb_func: ListCallback, + entry_error: Arc>>, + pool_index: usize, + set_index: usize, ) -> Result<()> { let (disks, _) = self.get_online_disks_with_healing(false).await; ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?; @@ -5120,6 +5194,12 @@ impl SetDisks { }; let cb1 = cb_func.clone(); + let unresolved_error = entry_error.clone(); + let unresolved_rx = rx.clone(); + let unresolved_bucket = bucket_info.name.clone(); + let unresolved_prefix = bucket_info.prefix.clone(); + let unresolved_pool_index = pool_index; + let unresolved_set_index = set_index; list_path_raw( rx, @@ -5132,20 +5212,51 @@ impl SetDisks { skip_walkdir_total_timeout: true, walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT), agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + partial: Some(Box::new(move |entries: MetaCacheEntries, errs: &[Option]| { let resolver = resolver.clone(); let cb_func = cb_func.clone(); - match entries.resolve(resolver) { - Some(entry) => { + let bucket = unresolved_bucket.clone(); + let prefix = unresolved_prefix.clone(); + let unresolved_error = unresolved_error.clone(); + let unresolved_rx = unresolved_rx.clone(); + let pool_index = unresolved_pool_index; + let set_index = unresolved_set_index; + let disk_error_count = errs.iter().flatten().count(); + if unresolved_rx.is_cancelled() { + return Box::pin(async {}); + } + + match resolve_decommission_partial_listing_entry( + entries, + resolver, + &bucket, + &prefix, + disk_error_count, + pool_index, + set_index, + ) { + Ok(entry) => { warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name); Box::pin(async move { cb_func(entry).await; }) } - None => { - warn!("decommission_pool: list_objects_to_decommission get none"); - Box::pin(async {}) - } + Err(err) => Box::pin(async move { + if unresolved_rx.is_cancelled() { + return; + } + warn!( + event = EVENT_DECOMMISSION_BUCKET, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + bucket = %bucket, + prefix = %prefix, + state = "unresolved_entry", + error = %err, + "Decommission listing failed closed on unresolved metadata" + ); + record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await; + }), } })), ..Default::default() @@ -5153,6 +5264,10 @@ impl SetDisks { ) .await?; + if let Some(err) = entry_error.lock().await.clone() { + return Err(err); + } + Ok(()) } } @@ -5358,11 +5473,12 @@ mod pools_tests { has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, - pool_meta_has_active_decommission, require_decommission_store, resolve_decommission_bucket_done_save_result, - resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, - resolve_decommission_entry_cleanup_delete_result, 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, + pool_meta_has_active_decommission, record_decommission_entry_error, 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_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error, + resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, + resolve_decommission_partial_listing_entry, 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, @@ -5380,7 +5496,9 @@ mod pools_tests { use crate::error::{Error, StorageError}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; - use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; + use rustfs_filemeta::{ + FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo, + }; use rustfs_rio::Index; use std::sync::{ Arc, @@ -6399,6 +6517,65 @@ mod pools_tests { assert!(matches!(err, Error::SlowDown)); } + #[test] + fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() { + let err = resolve_decommission_partial_listing_entry( + MetaCacheEntries(vec![None]), + MetadataResolutionParams { + dir_quorum: 2, + obj_quorum: 2, + bucket: "bucket-a".to_string(), + ..Default::default() + }, + "bucket-a", + "prefix/", + 1, + 2, + 3, + ) + .expect_err("unresolved partial listing must fail closed"); + + let message = err.to_string(); + assert!(message.contains("decommission listing could not resolve metadata")); + assert!(message.contains("bucket-a/prefix/")); + assert!(message.contains("pool 2 set 3")); + assert!(message.contains("1 disk error(s)")); + } + + #[tokio::test] + async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() { + let entry_error = Arc::new(tokio::sync::Mutex::new(None)); + let rx = CancellationToken::new(); + + record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await; + record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await; + + assert!(rx.is_cancelled()); + assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown))); + } + + #[tokio::test] + async fn test_record_decommission_entry_error_ignores_already_canceled_listing() { + let entry_error = Arc::new(tokio::sync::Mutex::new(None)); + let rx = CancellationToken::new(); + rx.cancel(); + + record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await; + + assert!(entry_error.lock().await.is_none()); + } + + #[test] + fn test_resolve_decommission_listing_error_preserves_real_listing_failure() { + let err = resolve_decommission_listing_error(Some(Error::SlowDown), Some(Error::OperationCanceled)) + .expect("listing failure should be returned"); + assert!(matches!(err, Error::SlowDown)); + + let err = resolve_decommission_listing_error(Some(Error::OperationCanceled), Some(Error::SlowDown)) + .expect("entry failure should be returned"); + assert!(matches!(err, Error::SlowDown)); + } + #[test] fn test_resolve_decommission_check_after_list_result_returns_list_result_without_entry_error() { let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), None) From 7143697a5f43118f523bc6ebdee39ea18a70cb8d Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 22 Aug 2026 14:42:10 +0800 Subject: [PATCH 4/4] fix(rpc): bound internode concurrency under multipart load (#6368) --- crates/protos/src/lib.rs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index a9e59b8a0..ccfb8197d 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -2106,6 +2106,9 @@ pub enum ChannelClass { Bulk, } +// Keep multiplexed unary RPCs below h2's per-connection small-frame budget. +const INTERNODE_RPC_CONCURRENCY_LIMIT: usize = 64; + /// Whether control/bulk channel isolation is enabled (env-gated, default off for safe rollout). fn channel_isolation_enabled() -> bool { rustfs_utils::get_env_bool( @@ -2188,6 +2191,7 @@ async fn build_channel(dial_addr: &str, cache_key: &str) -> Result