diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 30d5565e1..cdb81f72b 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -892,7 +892,7 @@ impl ErasureSetHealer { 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); - let progress = { + let (outcome_record, counter_unknown) = { let mut progress = self.progress.write().await; progress.record_skipped_new_version(); progress.set_current_object(Some(format!("skipped_new: {bucket}/{}", item.name))); @@ -906,22 +906,23 @@ impl ErasureSetHealer { if !counter_ok { progress.mark_unknown(); } - progress.clone() + ( + CheckpointObjectOutcomeRecord { + object: key, + outcome: CheckpointObjectOutcome::Processed, + successful: progress.objects_healed, + failed: progress.objects_failed, + skipped: progress.skipped_objects, + bytes: progress.bytes_processed, + skipped_new_versions: progress.skipped_new_versions, + skipped_ilm_expired: progress.skipped_ilm_expired, + counter_unknown: progress.counter_unknown, + }, + progress.counter_unknown, + ) }; - checkpoint_manager - .record_object_outcome(CheckpointObjectOutcomeRecord { - object: key, - outcome: CheckpointObjectOutcome::Processed, - successful: progress.objects_healed, - failed: progress.objects_failed, - skipped: progress.skipped_objects, - bytes: progress.bytes_processed, - skipped_new_versions: progress.skipped_new_versions, - skipped_ilm_expired: progress.skipped_ilm_expired, - counter_unknown: progress.counter_unknown, - }) - .await?; - if progress.counter_unknown { + checkpoint_manager.record_object_outcome(outcome_record).await?; + if counter_unknown { resume_manager.mark_counter_unknown().await?; } debug!( @@ -957,7 +958,7 @@ impl ErasureSetHealer { 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); - let progress = { + let (outcome_record, counter_unknown) = { let mut progress = self.progress.write().await; progress.record_skipped_ilm_expired(); progress.set_current_object(Some(format!("skipped_ilm: {bucket}/{}", item.name))); @@ -971,22 +972,23 @@ impl ErasureSetHealer { if !counter_ok { progress.mark_unknown(); } - progress.clone() + ( + CheckpointObjectOutcomeRecord { + object: key, + outcome: CheckpointObjectOutcome::Processed, + successful: progress.objects_healed, + failed: progress.objects_failed, + skipped: progress.skipped_objects, + bytes: progress.bytes_processed, + skipped_new_versions: progress.skipped_new_versions, + skipped_ilm_expired: progress.skipped_ilm_expired, + counter_unknown: progress.counter_unknown, + }, + progress.counter_unknown, + ) }; - checkpoint_manager - .record_object_outcome(CheckpointObjectOutcomeRecord { - object: key, - outcome: CheckpointObjectOutcome::Processed, - successful: progress.objects_healed, - failed: progress.objects_failed, - skipped: progress.skipped_objects, - bytes: progress.bytes_processed, - skipped_new_versions: progress.skipped_new_versions, - skipped_ilm_expired: progress.skipped_ilm_expired, - counter_unknown: progress.counter_unknown, - }) - .await?; - if progress.counter_unknown { + checkpoint_manager.record_object_outcome(outcome_record).await?; + if counter_unknown { resume_manager.mark_counter_unknown().await?; } debug!( @@ -1188,7 +1190,7 @@ impl ErasureSetHealer { telemetry_unknown |= !increment_counter(processed_objects); completed_in_page += 1; - let progress = { + let (outcome_record, counter_unknown) = { let mut progress = self.progress.write().await; progress.set_current_object(Some(format!("{bucket}/{object}"))); progress.update_object_progress( @@ -1201,22 +1203,23 @@ impl ErasureSetHealer { if telemetry_unknown { progress.mark_unknown(); } - progress.clone() + ( + CheckpointObjectOutcomeRecord { + object: key, + outcome: checkpoint_outcome, + successful: progress.objects_healed, + failed: progress.objects_failed, + skipped: progress.skipped_objects, + bytes: progress.bytes_processed, + skipped_new_versions: progress.skipped_new_versions, + skipped_ilm_expired: progress.skipped_ilm_expired, + counter_unknown: progress.counter_unknown, + }, + progress.counter_unknown, + ) }; - checkpoint_manager - .record_object_outcome(CheckpointObjectOutcomeRecord { - object: key, - outcome: checkpoint_outcome, - successful: progress.objects_healed, - failed: progress.objects_failed, - skipped: progress.skipped_objects, - bytes: progress.bytes_processed, - skipped_new_versions: progress.skipped_new_versions, - skipped_ilm_expired: progress.skipped_ilm_expired, - counter_unknown: progress.counter_unknown, - }) - .await?; - if progress.counter_unknown { + checkpoint_manager.record_object_outcome(outcome_record).await?; + if counter_unknown { resume_manager.mark_counter_unknown().await?; } diff --git a/crates/heal/src/heal/progress.rs b/crates/heal/src/heal/progress.rs index 7b72b5244..7f4ba3b26 100644 --- a/crates/heal/src/heal/progress.rs +++ b/crates/heal/src/heal/progress.rs @@ -15,6 +15,27 @@ use serde::{Deserialize, Serialize}; use std::time::{Duration, SystemTime}; +pub(crate) fn stable_generation(parts: &[&[u8]]) -> u64 { + let mut hash = 0xcbf29ce484222325u64; + for part in parts { + for byte in (part.len() as u64).to_be_bytes().into_iter().chain(part.iter().copied()) { + hash ^= u64::from(byte); + hash = hash.wrapping_mul(0x100000001b3); + } + } + hash +} + +#[cfg(test)] +mod stable_generation_tests { + use super::stable_generation; + + #[test] + fn stable_generation_has_a_fixed_vector() { + assert_eq!(stable_generation(&[b"rustfs", b"heal", b"42"]), 11_007_672_338_488_385_056); + } +} + pub(crate) fn increment_counter(counter: &mut u64) -> bool { match counter.checked_add(1) { Some(next) => { @@ -847,6 +868,23 @@ mod tests { assert_eq!(aggregate.progress_percentage, 0.0); } + #[test] + fn aggregate_accepts_multiple_sets_from_one_snapshot_generation() { + let progress = |objects_scanned| HealProgress { + kind: HealProgressKind::ObjectSweep, + objects_scanned, + objects_total_count: 10, + progress_state: HealProgressState::Running, + baseline_generation: Some(7), + baseline_known: true, + ..Default::default() + }; + + let aggregate = aggregate_heal_progress([progress(5), progress(3)]).expect("progress should aggregate"); + assert!(aggregate.baseline_known); + assert_eq!(aggregate.baseline_generation, Some(7)); + } + #[test] fn stage_updates_do_not_double_count_object_outcomes() { let mut progress = HealProgress::new(); diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index afa9e913f..a3e7c0cdf 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -19,11 +19,10 @@ use base64::engine::general_purpose::URL_SAFE_NO_PAD; use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use serde::{Deserialize, Serialize}; -use std::collections::hash_map::DefaultHasher; -use std::hash::{Hash, Hasher}; use std::sync::Arc; use tracing::{debug, error, warn}; +use super::progress::stable_generation; use super::storage_api::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached}; use super::storage_api::storage::{ BucketInfo, BucketOperations, DiskSetSelector, HealOperations as _, ListOperations as _, ObjectIO as _, @@ -805,14 +804,36 @@ impl HealStorageAPI for ECStoreHealStorage { } let identity = info.snapshot_identity(); - let mut hasher = DefaultHasher::new(); - identity.last_update.hash(&mut hasher); - identity.scanner_cycle.hash(&mut hasher); - identity.scanner_epoch.hash(&mut hasher); + let mut canonical = Vec::new(); + match identity.last_update { + Some(last_update) => { + canonical.push(1); + canonical.extend_from_slice( + &last_update + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() + .to_be_bytes(), + ); + } + None => canonical.push(0), + } + for value in [identity.scanner_cycle, identity.scanner_epoch] { + match value { + Some(value) => { + canonical.push(1); + canonical.extend_from_slice(&value.to_be_bytes()); + } + None => canonical.push(0), + } + } let mut scope = buckets.to_vec(); scope.sort_unstable(); - scope.hash(&mut hasher); - baseline.generation = Some(hasher.finish()); + for bucket in scope { + canonical.extend_from_slice(&(bucket.len() as u64).to_be_bytes()); + canonical.extend_from_slice(bucket.as_bytes()); + } + baseline.generation = Some(stable_generation(&[&canonical])); Ok(Some(baseline)) } diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index d46e97869..8ef9c90ab 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -13,7 +13,7 @@ // limitations under the License. /// bucket/cluster/prefix heal: the recursive bucket-objects sweep and the erasure-set usage baseline use super::*; -use crate::heal::progress::{add_bytes, increment_counter}; +use crate::heal::progress::{add_bytes, increment_counter, stable_generation}; impl HealTask { pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> { @@ -444,7 +444,7 @@ impl HealTask { Ok(()) } - pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String], set_disk_id: &str) -> Result<()> { + pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String], _set_disk_id: &str) -> Result<()> { if matches!(self.options.scan_mode, HealScanMode::Deep) || matches!(self.source, HealRequestSource::AutoHeal) { return Ok(()); } @@ -463,15 +463,7 @@ impl HealTask { bytes, generation, } = baseline; - let generation = generation.map(|snapshot_generation| { - use std::hash::{Hash, Hasher}; - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - snapshot_generation.hash(&mut hasher); - set_disk_id.hash(&mut hasher); - self.options.pool_index.hash(&mut hasher); - self.options.set_index.hash(&mut hasher); - hasher.finish() - }); + let generation = generation.map(|snapshot_generation| stable_generation(&[&snapshot_generation.to_be_bytes()])); let mut progress = self.progress.write().await; if let Some(generation) = generation { progress.set_total_baseline_with_generation(objects_count, bytes, generation);