fix(heal): atomically persist page progress

This commit is contained in:
马登山
2026-08-22 20:53:55 +08:00
parent db709aee57
commit 69ae9ef671
4 changed files with 136 additions and 33 deletions
+45 -31
View File
@@ -15,7 +15,7 @@
use crate::heal::{ use crate::heal::{
progress::{HealProgress, add_bytes, increment_counter}, progress::{HealProgress, add_bytes, increment_counter},
resume::{ resume::{
CheckpointManager, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key, CheckpointManager, CheckpointObjectOutcome, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key,
replacement_target_identities_match, replacement_target_identities_match,
}, },
storage::{HealStorageAPI, next_heal_listing_token}, 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) { 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); let counter_ok = increment_counter(processed_objects);
completed_in_page = completed_in_page.saturating_add(1); completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_new_versions_total").increment(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) (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
}; };
if !counter_ok || counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager checkpoint_manager
.set_skipped_version_counts(skipped_new, skipped_ilm) .record_object_outcome(
.await?; key,
checkpoint_manager CheckpointObjectOutcome::Processed,
.update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed) *successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
skipped_new,
skipped_ilm,
!counter_ok || counter_unknown,
)
.await?; .await?;
if !counter_ok || counter_unknown { if !counter_ok || counter_unknown {
resume_manager.mark_counter_unknown().await?; resume_manager.mark_counter_unknown().await?;
@@ -943,7 +946,6 @@ impl ErasureSetHealer {
) )
.await? .await?
{ {
checkpoint_manager.add_processed_object(key).await?;
let counter_ok = increment_counter(processed_objects); let counter_ok = increment_counter(processed_objects);
completed_in_page = completed_in_page.saturating_add(1); completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_ilm_expired_total").increment(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) (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
}; };
if !counter_ok || counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager checkpoint_manager
.set_skipped_version_counts(skipped_new, skipped_ilm) .record_object_outcome(
.await?; key,
checkpoint_manager CheckpointObjectOutcome::Processed,
.update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed) *successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
skipped_new,
skipped_ilm,
!counter_ok || counter_unknown,
)
.await?; .await?;
if !counter_ok || counter_unknown { if !counter_ok || counter_unknown {
resume_manager.mark_counter_unknown().await?; resume_manager.mark_counter_unknown().await?;
@@ -1100,11 +1106,10 @@ impl ErasureSetHealer {
while let Some((key, object, version_id, result)) = page_tasks.next().await { while let Some((key, object, version_id, result)) = page_tasks.next().await {
let (object_size, result) = result; let (object_size, result) = result;
let mut telemetry_unknown = false; let mut telemetry_unknown = false;
match result { let checkpoint_outcome = match result {
Ok(true) => { Ok(true) => {
telemetry_unknown |= !increment_counter(successful_objects); telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
checkpoint_manager.add_processed_object(key).await?;
debug!( debug!(
target: "rustfs::heal::erasure_healer", target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE, event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -1117,9 +1122,9 @@ impl ErasureSetHealer {
state = "healed", state = "healed",
"Erasure set object healed" "Erasure set object healed"
); );
CheckpointObjectOutcome::Processed
} }
Ok(false) => { Ok(false) => {
checkpoint_manager.add_processed_object(key).await?;
telemetry_unknown |= !increment_counter(successful_objects); telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
debug!( debug!(
@@ -1134,12 +1139,12 @@ impl ErasureSetHealer {
state = "missing_treated_as_ok", state = "missing_treated_as_ok",
"Erasure set missing object treated as ok" "Erasure set missing object treated as ok"
); );
CheckpointObjectOutcome::Processed
} }
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(Error::TransientSkip { message }) => { Err(Error::TransientSkip { message }) => {
telemetry_unknown |= !increment_counter(skipped_objects); telemetry_unknown |= !increment_counter(skipped_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); 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", { 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, event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL, component = LOG_COMPONENT_HEAL,
@@ -1152,11 +1157,11 @@ impl ErasureSetHealer {
error = %message, error = %message,
"Erasure set object heal skipped due to transient error" "Erasure set object heal skipped due to transient error"
}); });
CheckpointObjectOutcome::Skipped
} }
Err(err) => { Err(err) => {
telemetry_unknown |= !increment_counter(failed_objects); telemetry_unknown |= !increment_counter(failed_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); 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", { demote_to_debug_when!(!take_failure_log_sample(&mut failure_samples_logged), warn, target: "rustfs::heal::erasure_healer", {
event = EVENT_HEAL_ERASURE_OBJECT_STATE, event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL, component = LOG_COMPONENT_HEAL,
@@ -1169,8 +1174,9 @@ impl ErasureSetHealer {
error = %err, error = %err,
"Erasure set object heal failed" "Erasure set object heal failed"
}); });
CheckpointObjectOutcome::Failed
} }
} };
telemetry_unknown |= !increment_counter(processed_objects); telemetry_unknown |= !increment_counter(processed_objects);
completed_in_page += 1; completed_in_page += 1;
@@ -1189,11 +1195,18 @@ impl ErasureSetHealer {
} }
progress.counter_unknown progress.counter_unknown
}; };
if telemetry_unknown || progress_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager 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?; .await?;
if telemetry_unknown || progress_unknown { if telemetry_unknown || progress_unknown {
resume_manager.mark_counter_unknown().await?; resume_manager.mark_counter_unknown().await?;
@@ -1206,12 +1219,13 @@ impl ErasureSetHealer {
*current_object_index = global_obj_idx; *current_object_index = global_obj_idx;
// Persist the authoritative cursor FIRST (points at the next page // Persist the checkpoint ledger and page position before exposing
// boundary), then prune the per-version dedup sets. Both are // the next resume cursor. A crash before cursor publication keeps
// idempotent under crash: heal_object re-heals safely. // the page identities available for exact-once replay.
let next_cursor = if is_truncated { next_token.clone() } else { None }; 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?; 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 // Check if there are more pages
if !is_truncated { if !is_truncated {
break; break;
+1 -1
View File
@@ -31,7 +31,7 @@ mod checkpoint;
mod replacement; mod replacement;
mod utils; mod utils;
pub use checkpoint::{CheckpointManager, ResumeCheckpoint}; pub use checkpoint::{CheckpointManager, CheckpointObjectOutcome, ResumeCheckpoint};
pub(crate) use replacement::replacement_target_identities_match; pub(crate) use replacement::replacement_target_identities_match;
use replacement::replacement_targets_match_identities; use replacement::replacement_targets_match_identities;
pub use replacement::{ pub use replacement::{
+56 -1
View File
@@ -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. /// to the new `compose_key` identities, so a stale checkpoint is discarded.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5; pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
#[derive(Debug, Clone, Copy)]
pub enum CheckpointObjectOutcome {
Processed,
Failed,
Skipped,
}
/// resume checkpoint /// resume checkpoint
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ResumeCheckpoint { pub struct ResumeCheckpoint {
@@ -310,7 +317,7 @@ impl CheckpointManager {
self.save_checkpoint_throttled().await 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<()> { pub async fn complete_page(&self, bucket_index: usize, object_index: usize) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await; let mut checkpoint = self.checkpoint.write().await;
checkpoint.complete_page(bucket_index, object_index); checkpoint.complete_page(bucket_index, object_index);
@@ -318,6 +325,26 @@ impl CheckpointManager {
self.save_checkpoint_throttled().await 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. /// Reset the checkpoint to the start of the scan for a retry, then persist.
pub async fn reset_for_retry(&self) -> Result<()> { pub async fn reset_for_retry(&self) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await; let mut checkpoint = self.checkpoint.write().await;
@@ -352,6 +379,34 @@ impl CheckpointManager {
self.save_checkpoint_if_due().await 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<()> { pub async fn update_progress(&self, successful: u64, failed: u64, skipped: u64, bytes: u64) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await; let mut checkpoint = self.checkpoint.write().await;
checkpoint.update_progress(successful, failed, skipped, bytes); checkpoint.update_progress(successful, failed, skipped, bytes);
+34
View File
@@ -1476,6 +1476,40 @@ fn test_checkpoint_object_sets_dedupe_and_prune() {
assert!(checkpoint.failed_objects.is_empty()); 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] #[test]
fn test_checkpoint_loads_legacy_vec_format() { fn test_checkpoint_loads_legacy_vec_format() {
// Checkpoints written before the HashSet migration stored the object // Checkpoints written before the HashSet migration stored the object