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");