diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 53bfa729a..91bb0dd2a 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -15,8 +15,8 @@ use crate::heal::{ progress::{HealProgress, add_bytes, increment_counter}, resume::{ - CheckpointManager, CheckpointObjectOutcome, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key, - replacement_target_identities_match, + CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, ReplacementTargetIdentity, ResumeManager, + ResumeUtils, compose_key, replacement_target_identities_match, }, storage::{HealStorageAPI, next_heal_listing_token}, task::{demote_to_debug_when, is_missing_object_dir_heal_result, take_failure_log_sample}, @@ -909,17 +909,17 @@ impl ErasureSetHealer { (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown) }; checkpoint_manager - .record_object_outcome( - key, - CheckpointObjectOutcome::Processed, - *successful_objects, - *failed_objects, - *skipped_objects, - bytes_processed, - skipped_new, - skipped_ilm, - !counter_ok || counter_unknown, - ) + .record_object_outcome(CheckpointObjectOutcomeRecord { + object: key, + outcome: CheckpointObjectOutcome::Processed, + successful: *successful_objects, + failed: *failed_objects, + skipped: *skipped_objects, + bytes: bytes_processed, + skipped_new_versions: skipped_new, + skipped_ilm_expired: skipped_ilm, + counter_unknown: !counter_ok || counter_unknown, + }) .await?; if !counter_ok || counter_unknown { resume_manager.mark_counter_unknown().await?; @@ -974,17 +974,17 @@ impl ErasureSetHealer { (progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown) }; checkpoint_manager - .record_object_outcome( - key, - CheckpointObjectOutcome::Processed, - *successful_objects, - *failed_objects, - *skipped_objects, - bytes_processed, - skipped_new, - skipped_ilm, - !counter_ok || counter_unknown, - ) + .record_object_outcome(CheckpointObjectOutcomeRecord { + object: key, + outcome: CheckpointObjectOutcome::Processed, + successful: *successful_objects, + failed: *failed_objects, + skipped: *skipped_objects, + bytes: bytes_processed, + skipped_new_versions: skipped_new, + skipped_ilm_expired: skipped_ilm, + counter_unknown: !counter_ok || counter_unknown, + }) .await?; if !counter_ok || counter_unknown { resume_manager.mark_counter_unknown().await?; @@ -1204,17 +1204,17 @@ impl ErasureSetHealer { (progress.counter_unknown, progress.skipped_new_versions, progress.skipped_ilm_expired) }; checkpoint_manager - .record_object_outcome( - key, - checkpoint_outcome, - *successful_objects, - *failed_objects, - *skipped_objects, - bytes_processed, + .record_object_outcome(CheckpointObjectOutcomeRecord { + object: key, + outcome: checkpoint_outcome, + successful: *successful_objects, + failed: *failed_objects, + skipped: *skipped_objects, + bytes: bytes_processed, skipped_new_versions, skipped_ilm_expired, - telemetry_unknown || progress_unknown, - ) + counter_unknown: telemetry_unknown || progress_unknown, + }) .await?; if telemetry_unknown || progress_unknown { resume_manager.mark_counter_unknown().await?; @@ -2413,17 +2413,17 @@ mod resume_loop_tests { }, ); env.checkpoint - .record_object_outcome( - compose_key("object", Some("v1")), - CheckpointObjectOutcome::Failed, - 0, - 1, - 0, - 0, - 0, - 0, - false, - ) + .record_object_outcome(CheckpointObjectOutcomeRecord { + object: compose_key("object", Some("v1")), + outcome: CheckpointObjectOutcome::Failed, + successful: 0, + failed: 1, + skipped: 0, + bytes: 0, + skipped_new_versions: 0, + skipped_ilm_expired: 0, + counter_unknown: false, + }) .await .unwrap(); env.checkpoint.advance_page(0, 1).await.unwrap(); diff --git a/crates/heal/src/heal/resume.rs b/crates/heal/src/heal/resume.rs index 39a2c898b..34a856817 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, CheckpointObjectOutcome, ResumeCheckpoint}; +pub use checkpoint::{CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, 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 a5657df87..3619356fd 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -41,6 +41,19 @@ pub enum CheckpointObjectOutcome { Skipped, } +#[derive(Debug)] +pub struct CheckpointObjectOutcomeRecord { + pub object: String, + pub outcome: CheckpointObjectOutcome, + pub successful: u64, + pub failed: u64, + pub skipped: u64, + pub bytes: u64, + pub skipped_new_versions: u64, + pub skipped_ilm_expired: u64, + pub counter_unknown: bool, +} + /// resume checkpoint #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ResumeCheckpoint { @@ -389,18 +402,18 @@ impl CheckpointManager { } /// 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<()> { + pub async fn record_object_outcome(&self, record: CheckpointObjectOutcomeRecord) -> Result<()> { + let CheckpointObjectOutcomeRecord { + object, + outcome, + successful, + failed, + skipped, + bytes, + skipped_new_versions, + skipped_ilm_expired, + counter_unknown, + } = record; let mut checkpoint = self.checkpoint.write().await; match outcome { CheckpointObjectOutcome::Processed => checkpoint.add_processed_object(object), diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index b5d03d7b7..ffb663c29 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1483,17 +1483,17 @@ async fn checkpoint_page_commit_keeps_ledger_until_cursor_is_durable() { 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, - ) + .record_object_outcome(CheckpointObjectOutcomeRecord { + object: "bucket/object:v1".to_string(), + outcome: CheckpointObjectOutcome::Processed, + successful: 1, + failed: 0, + skipped: 0, + bytes: 128, + skipped_new_versions: 0, + skipped_ilm_expired: 0, + counter_unknown: false, + }) .await .unwrap(); checkpoint.advance_page(0, 1).await.unwrap();