From 69ae9ef671016292ca61aed13734637f2b19d1c5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Sat, 22 Aug 2026 20:53:55 +0800 Subject: [PATCH] fix(heal): atomically persist page progress --- crates/heal/src/heal/erasure_healer.rs | 76 ++++++++++++++--------- crates/heal/src/heal/resume.rs | 2 +- crates/heal/src/heal/resume/checkpoint.rs | 57 ++++++++++++++++- crates/heal/src/heal/resume/tests.rs | 34 ++++++++++ 4 files changed, 136 insertions(+), 33 deletions(-) diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 1eb6cc4c2..b26889df6 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -15,7 +15,7 @@ use crate::heal::{ progress::{HealProgress, add_bytes, increment_counter}, resume::{ - CheckpointManager, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key, + CheckpointManager, CheckpointObjectOutcome, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key, replacement_target_identities_match, }, storage::{HealStorageAPI, next_heal_listing_token}, @@ -881,7 +881,6 @@ impl ErasureSetHealer { } if should_skip_new_version(item.mod_time_unix_nanos, started_at_secs) { - checkpoint_manager.add_processed_object(key).await?; let counter_ok = increment_counter(processed_objects); completed_in_page = completed_in_page.saturating_add(1); counter!("rustfs_heal_skipped_new_versions_total").increment(1); @@ -901,14 +900,18 @@ impl ErasureSetHealer { } (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown) }; - if !counter_ok || counter_unknown { - checkpoint_manager.mark_counter_unknown().await?; - } checkpoint_manager - .set_skipped_version_counts(skipped_new, skipped_ilm) - .await?; - checkpoint_manager - .update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed) + .record_object_outcome( + key, + CheckpointObjectOutcome::Processed, + *successful_objects, + *failed_objects, + *skipped_objects, + bytes_processed, + skipped_new, + skipped_ilm, + !counter_ok || counter_unknown, + ) .await?; if !counter_ok || counter_unknown { resume_manager.mark_counter_unknown().await?; @@ -943,7 +946,6 @@ impl ErasureSetHealer { ) .await? { - checkpoint_manager.add_processed_object(key).await?; let counter_ok = increment_counter(processed_objects); completed_in_page = completed_in_page.saturating_add(1); counter!("rustfs_heal_skipped_ilm_expired_total").increment(1); @@ -963,14 +965,18 @@ impl ErasureSetHealer { } (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown) }; - if !counter_ok || counter_unknown { - checkpoint_manager.mark_counter_unknown().await?; - } checkpoint_manager - .set_skipped_version_counts(skipped_new, skipped_ilm) - .await?; - checkpoint_manager - .update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed) + .record_object_outcome( + key, + CheckpointObjectOutcome::Processed, + *successful_objects, + *failed_objects, + *skipped_objects, + bytes_processed, + skipped_new, + skipped_ilm, + !counter_ok || counter_unknown, + ) .await?; if !counter_ok || counter_unknown { resume_manager.mark_counter_unknown().await?; @@ -1100,11 +1106,10 @@ impl ErasureSetHealer { while let Some((key, object, version_id, result)) = page_tasks.next().await { let (object_size, result) = result; let mut telemetry_unknown = false; - match result { + let checkpoint_outcome = match result { Ok(true) => { telemetry_unknown |= !increment_counter(successful_objects); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); - checkpoint_manager.add_processed_object(key).await?; debug!( target: "rustfs::heal::erasure_healer", event = EVENT_HEAL_ERASURE_OBJECT_STATE, @@ -1117,9 +1122,9 @@ impl ErasureSetHealer { state = "healed", "Erasure set object healed" ); + CheckpointObjectOutcome::Processed } Ok(false) => { - checkpoint_manager.add_processed_object(key).await?; telemetry_unknown |= !increment_counter(successful_objects); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); debug!( @@ -1134,12 +1139,12 @@ impl ErasureSetHealer { state = "missing_treated_as_ok", "Erasure set missing object treated as ok" ); + CheckpointObjectOutcome::Processed } Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), Err(Error::TransientSkip { message }) => { telemetry_unknown |= !increment_counter(skipped_objects); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); - checkpoint_manager.add_skipped_object(key).await?; demote_to_debug_when!(!take_failure_log_sample(&mut transient_skip_samples_logged), warn, target: "rustfs::heal::erasure_healer", { event = EVENT_HEAL_ERASURE_OBJECT_STATE, component = LOG_COMPONENT_HEAL, @@ -1152,11 +1157,11 @@ impl ErasureSetHealer { error = %message, "Erasure set object heal skipped due to transient error" }); + CheckpointObjectOutcome::Skipped } Err(err) => { telemetry_unknown |= !increment_counter(failed_objects); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); - checkpoint_manager.add_failed_object(key).await?; demote_to_debug_when!(!take_failure_log_sample(&mut failure_samples_logged), warn, target: "rustfs::heal::erasure_healer", { event = EVENT_HEAL_ERASURE_OBJECT_STATE, component = LOG_COMPONENT_HEAL, @@ -1169,8 +1174,9 @@ impl ErasureSetHealer { error = %err, "Erasure set object heal failed" }); + CheckpointObjectOutcome::Failed } - } + }; telemetry_unknown |= !increment_counter(processed_objects); completed_in_page += 1; @@ -1189,11 +1195,18 @@ impl ErasureSetHealer { } progress.counter_unknown }; - if telemetry_unknown || progress_unknown { - checkpoint_manager.mark_counter_unknown().await?; - } checkpoint_manager - .update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed) + .record_object_outcome( + key, + checkpoint_outcome, + *successful_objects, + *failed_objects, + *skipped_objects, + bytes_processed, + 0, + 0, + telemetry_unknown || progress_unknown, + ) .await?; if telemetry_unknown || progress_unknown { resume_manager.mark_counter_unknown().await?; @@ -1206,12 +1219,13 @@ impl ErasureSetHealer { *current_object_index = global_obj_idx; - // Persist the authoritative cursor FIRST (points at the next page - // boundary), then prune the per-version dedup sets. Both are - // idempotent under crash: heal_object re-heals safely. + // Persist the checkpoint ledger and page position before exposing + // the next resume cursor. A crash before cursor publication keeps + // the page identities available for exact-once replay. let next_cursor = if is_truncated { next_token.clone() } else { None }; + checkpoint_manager.advance_page(bucket_index, *current_object_index).await?; resume_manager.set_resume_cursor(next_cursor.clone()).await?; - checkpoint_manager.complete_page(bucket_index, *current_object_index).await?; + checkpoint_manager.prune_completed_page().await?; // Check if there are more pages if !is_truncated { break; diff --git a/crates/heal/src/heal/resume.rs b/crates/heal/src/heal/resume.rs index b156ef6dd..f0b0adf6f 100644 --- a/crates/heal/src/heal/resume.rs +++ b/crates/heal/src/heal/resume.rs @@ -31,7 +31,7 @@ mod checkpoint; mod replacement; mod utils; -pub use checkpoint::{CheckpointManager, ResumeCheckpoint}; +pub use checkpoint::{CheckpointManager, CheckpointObjectOutcome, ResumeCheckpoint}; pub(crate) use replacement::replacement_target_identities_match; use replacement::replacement_targets_match_identities; pub use replacement::{ diff --git a/crates/heal/src/heal/resume/checkpoint.rs b/crates/heal/src/heal/resume/checkpoint.rs index 40089cadd..38fd2448b 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -34,6 +34,13 @@ const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state"; /// to the new `compose_key` identities, so a stale checkpoint is discarded. pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5; +#[derive(Debug, Clone, Copy)] +pub enum CheckpointObjectOutcome { + Processed, + Failed, + Skipped, +} + /// resume checkpoint #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ResumeCheckpoint { @@ -310,7 +317,7 @@ impl CheckpointManager { self.save_checkpoint_throttled().await } - /// Advance past a completed page and prune the per-object sets, then persist. + /// Persist a completed page position while retaining its identities. pub async fn complete_page(&self, bucket_index: usize, object_index: usize) -> Result<()> { let mut checkpoint = self.checkpoint.write().await; checkpoint.complete_page(bucket_index, object_index); @@ -318,6 +325,26 @@ impl CheckpointManager { self.save_checkpoint_throttled().await } + /// Persist the page position while retaining identities until the resume + /// cursor is durable. + pub async fn advance_page(&self, bucket_index: usize, object_index: usize) -> Result<()> { + let mut checkpoint = self.checkpoint.write().await; + checkpoint.update_position(bucket_index, object_index); + drop(checkpoint); + self.save_checkpoint().await + } + + /// Remove the previous page's dedup identities only after its resume cursor + /// has been durably exposed. + pub async fn prune_completed_page(&self) -> Result<()> { + let mut checkpoint = self.checkpoint.write().await; + checkpoint.processed_objects.clear(); + checkpoint.skipped_objects.clear(); + checkpoint.failed_objects.clear(); + drop(checkpoint); + self.save_checkpoint().await + } + /// Reset the checkpoint to the start of the scan for a retry, then persist. pub async fn reset_for_retry(&self) -> Result<()> { let mut checkpoint = self.checkpoint.write().await; @@ -352,6 +379,34 @@ impl CheckpointManager { self.save_checkpoint_if_due().await } + /// Atomically persist an object's dedup identity with its aggregate result. + pub async fn record_object_outcome( + &self, + object: String, + outcome: CheckpointObjectOutcome, + successful: u64, + failed: u64, + skipped: u64, + bytes: u64, + skipped_new_versions: u64, + skipped_ilm_expired: u64, + counter_unknown: bool, + ) -> Result<()> { + let mut checkpoint = self.checkpoint.write().await; + match outcome { + CheckpointObjectOutcome::Processed => checkpoint.add_processed_object(object), + CheckpointObjectOutcome::Failed => checkpoint.add_failed_object(object), + CheckpointObjectOutcome::Skipped => checkpoint.add_skipped_object(object), + } + checkpoint.update_progress(successful, failed, skipped, bytes); + checkpoint.set_skipped_version_counts(skipped_new_versions, skipped_ilm_expired); + if counter_unknown { + checkpoint.mark_counter_unknown(); + } + drop(checkpoint); + self.save_checkpoint_if_due().await + } + pub async fn update_progress(&self, successful: u64, failed: u64, skipped: u64, bytes: u64) -> Result<()> { let mut checkpoint = self.checkpoint.write().await; checkpoint.update_progress(successful, failed, skipped, bytes); diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index 945290b84..b5d03d7b7 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1476,6 +1476,40 @@ fn test_checkpoint_object_sets_dedupe_and_prune() { assert!(checkpoint.failed_objects.is_empty()); } +#[tokio::test] +async fn checkpoint_page_commit_keeps_ledger_until_cursor_is_durable() { + let (_temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let checkpoint = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap(); + + checkpoint + .record_object_outcome( + "bucket/object:v1".to_string(), + CheckpointObjectOutcome::Processed, + 1, + 0, + 0, + 128, + 0, + 0, + false, + ) + .await + .unwrap(); + checkpoint.advance_page(0, 1).await.unwrap(); + + let reloaded = CheckpointManager::load_from_disk(disk.clone(), &task_id).await.unwrap(); + let snapshot = reloaded.get_checkpoint().await; + assert_eq!(snapshot.current_object_index, 1); + assert_eq!(snapshot.successful_objects, 1); + assert_eq!(snapshot.processed_bytes, 128); + assert!(snapshot.processed_objects.contains("bucket/object:v1")); + + checkpoint.prune_completed_page().await.unwrap(); + let reloaded = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap(); + assert!(reloaded.get_checkpoint().await.processed_objects.is_empty()); +} + #[test] fn test_checkpoint_loads_legacy_vec_format() { // Checkpoints written before the HashSet migration stored the object