mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
fix(heal): stabilize progress generations
This commit is contained in:
@@ -892,7 +892,7 @@ impl ErasureSetHealer {
|
|||||||
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);
|
||||||
let progress = {
|
let (outcome_record, counter_unknown) = {
|
||||||
let mut progress = self.progress.write().await;
|
let mut progress = self.progress.write().await;
|
||||||
progress.record_skipped_new_version();
|
progress.record_skipped_new_version();
|
||||||
progress.set_current_object(Some(format!("skipped_new: {bucket}/{}", item.name)));
|
progress.set_current_object(Some(format!("skipped_new: {bucket}/{}", item.name)));
|
||||||
@@ -906,22 +906,23 @@ impl ErasureSetHealer {
|
|||||||
if !counter_ok {
|
if !counter_ok {
|
||||||
progress.mark_unknown();
|
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
|
checkpoint_manager.record_object_outcome(outcome_record).await?;
|
||||||
.record_object_outcome(CheckpointObjectOutcomeRecord {
|
if counter_unknown {
|
||||||
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 {
|
|
||||||
resume_manager.mark_counter_unknown().await?;
|
resume_manager.mark_counter_unknown().await?;
|
||||||
}
|
}
|
||||||
debug!(
|
debug!(
|
||||||
@@ -957,7 +958,7 @@ impl ErasureSetHealer {
|
|||||||
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);
|
||||||
let progress = {
|
let (outcome_record, counter_unknown) = {
|
||||||
let mut progress = self.progress.write().await;
|
let mut progress = self.progress.write().await;
|
||||||
progress.record_skipped_ilm_expired();
|
progress.record_skipped_ilm_expired();
|
||||||
progress.set_current_object(Some(format!("skipped_ilm: {bucket}/{}", item.name)));
|
progress.set_current_object(Some(format!("skipped_ilm: {bucket}/{}", item.name)));
|
||||||
@@ -971,22 +972,23 @@ impl ErasureSetHealer {
|
|||||||
if !counter_ok {
|
if !counter_ok {
|
||||||
progress.mark_unknown();
|
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
|
checkpoint_manager.record_object_outcome(outcome_record).await?;
|
||||||
.record_object_outcome(CheckpointObjectOutcomeRecord {
|
if counter_unknown {
|
||||||
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 {
|
|
||||||
resume_manager.mark_counter_unknown().await?;
|
resume_manager.mark_counter_unknown().await?;
|
||||||
}
|
}
|
||||||
debug!(
|
debug!(
|
||||||
@@ -1188,7 +1190,7 @@ impl ErasureSetHealer {
|
|||||||
|
|
||||||
telemetry_unknown |= !increment_counter(processed_objects);
|
telemetry_unknown |= !increment_counter(processed_objects);
|
||||||
completed_in_page += 1;
|
completed_in_page += 1;
|
||||||
let progress = {
|
let (outcome_record, counter_unknown) = {
|
||||||
let mut progress = self.progress.write().await;
|
let mut progress = self.progress.write().await;
|
||||||
progress.set_current_object(Some(format!("{bucket}/{object}")));
|
progress.set_current_object(Some(format!("{bucket}/{object}")));
|
||||||
progress.update_object_progress(
|
progress.update_object_progress(
|
||||||
@@ -1201,22 +1203,23 @@ impl ErasureSetHealer {
|
|||||||
if telemetry_unknown {
|
if telemetry_unknown {
|
||||||
progress.mark_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
|
checkpoint_manager.record_object_outcome(outcome_record).await?;
|
||||||
.record_object_outcome(CheckpointObjectOutcomeRecord {
|
if counter_unknown {
|
||||||
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 {
|
|
||||||
resume_manager.mark_counter_unknown().await?;
|
resume_manager.mark_counter_unknown().await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -15,6 +15,27 @@
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use std::time::{Duration, SystemTime};
|
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 {
|
pub(crate) fn increment_counter(counter: &mut u64) -> bool {
|
||||||
match counter.checked_add(1) {
|
match counter.checked_add(1) {
|
||||||
Some(next) => {
|
Some(next) => {
|
||||||
@@ -847,6 +868,23 @@ mod tests {
|
|||||||
assert_eq!(aggregate.progress_percentage, 0.0);
|
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]
|
#[test]
|
||||||
fn stage_updates_do_not_double_count_object_outcomes() {
|
fn stage_updates_do_not_double_count_object_outcomes() {
|
||||||
let mut progress = HealProgress::new();
|
let mut progress = HealProgress::new();
|
||||||
|
|||||||
@@ -19,11 +19,10 @@ use base64::engine::general_purpose::URL_SAFE_NO_PAD;
|
|||||||
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
|
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
|
||||||
use rustfs_madmin::heal_commands::HealResultItem;
|
use rustfs_madmin::heal_commands::HealResultItem;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use std::collections::hash_map::DefaultHasher;
|
|
||||||
use std::hash::{Hash, Hasher};
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tracing::{debug, error, warn};
|
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::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached};
|
||||||
use super::storage_api::storage::{
|
use super::storage_api::storage::{
|
||||||
BucketInfo, BucketOperations, DiskSetSelector, HealOperations as _, ListOperations as _, ObjectIO as _,
|
BucketInfo, BucketOperations, DiskSetSelector, HealOperations as _, ListOperations as _, ObjectIO as _,
|
||||||
@@ -805,14 +804,36 @@ impl HealStorageAPI for ECStoreHealStorage {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let identity = info.snapshot_identity();
|
let identity = info.snapshot_identity();
|
||||||
let mut hasher = DefaultHasher::new();
|
let mut canonical = Vec::new();
|
||||||
identity.last_update.hash(&mut hasher);
|
match identity.last_update {
|
||||||
identity.scanner_cycle.hash(&mut hasher);
|
Some(last_update) => {
|
||||||
identity.scanner_epoch.hash(&mut hasher);
|
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();
|
let mut scope = buckets.to_vec();
|
||||||
scope.sort_unstable();
|
scope.sort_unstable();
|
||||||
scope.hash(&mut hasher);
|
for bucket in scope {
|
||||||
baseline.generation = Some(hasher.finish());
|
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))
|
Ok(Some(baseline))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,7 +13,7 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
/// bucket/cluster/prefix heal: the recursive bucket-objects sweep and the erasure-set usage baseline
|
/// bucket/cluster/prefix heal: the recursive bucket-objects sweep and the erasure-set usage baseline
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::heal::progress::{add_bytes, increment_counter};
|
use crate::heal::progress::{add_bytes, increment_counter, stable_generation};
|
||||||
|
|
||||||
impl HealTask {
|
impl HealTask {
|
||||||
pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> {
|
pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> {
|
||||||
@@ -444,7 +444,7 @@ impl HealTask {
|
|||||||
Ok(())
|
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) {
|
if matches!(self.options.scan_mode, HealScanMode::Deep) || matches!(self.source, HealRequestSource::AutoHeal) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
@@ -463,15 +463,7 @@ impl HealTask {
|
|||||||
bytes,
|
bytes,
|
||||||
generation,
|
generation,
|
||||||
} = baseline;
|
} = baseline;
|
||||||
let generation = generation.map(|snapshot_generation| {
|
let generation = generation.map(|snapshot_generation| stable_generation(&[&snapshot_generation.to_be_bytes()]));
|
||||||
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 mut progress = self.progress.write().await;
|
let mut progress = self.progress.write().await;
|
||||||
if let Some(generation) = generation {
|
if let Some(generation) = generation {
|
||||||
progress.set_total_baseline_with_generation(objects_count, bytes, generation);
|
progress.set_total_baseline_with_generation(objects_count, bytes, generation);
|
||||||
|
|||||||
Reference in New Issue
Block a user