fix(heal): preserve progress status across nodes

This commit is contained in:
overtrue
2026-08-23 06:49:29 +08:00
parent 37c5fef252
commit e63033fbc7
15 changed files with 467 additions and 190 deletions
@@ -1089,7 +1089,9 @@ impl PeerRestClient {
.await?
.max_decoding_message_size(BACKGROUND_HEAL_STATUS_MAX_MESSAGE_SIZE);
let response = match client
.background_heal_status(Request::new(BackgroundHealStatusRequest::default()))
.background_heal_status(Request::new(BackgroundHealStatusRequest {
protocol_version: rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION,
}))
.await
{
Ok(response) => response.into_inner(),
+32 -32
View File
@@ -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 (skipped_new, skipped_ilm, counter_unknown) = {
let progress = {
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,22 @@ impl ErasureSetHealer {
if !counter_ok {
progress.mark_unknown();
}
(progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
progress.clone()
};
checkpoint_manager
.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,
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 !counter_ok || counter_unknown {
if progress.counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
debug!(
@@ -957,7 +957,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 (skipped_new, skipped_ilm, counter_unknown) = {
let progress = {
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 +971,22 @@ impl ErasureSetHealer {
if !counter_ok {
progress.mark_unknown();
}
(progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
progress.clone()
};
checkpoint_manager
.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,
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 !counter_ok || counter_unknown {
if progress.counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
debug!(
@@ -1188,7 +1188,7 @@ impl ErasureSetHealer {
telemetry_unknown |= !increment_counter(processed_objects);
completed_in_page += 1;
let (progress_unknown, skipped_new_versions, skipped_ilm_expired) = {
let progress = {
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_object_progress(
@@ -1201,22 +1201,22 @@ impl ErasureSetHealer {
if telemetry_unknown {
progress.mark_unknown();
}
(progress.counter_unknown, progress.skipped_new_versions, progress.skipped_ilm_expired)
progress.clone()
};
checkpoint_manager
.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,
counter_unknown: telemetry_unknown || progress_unknown,
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 telemetry_unknown || progress_unknown {
if progress.counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
@@ -1506,8 +1506,8 @@ mod resume_loop_tests {
};
use crate::heal::progress::HealProgress;
use crate::heal::resume::{
CheckpointManager, CheckpointObjectOutcome, RESUME_CHECKPOINT_FILE, ReplacementTargetIdentity, ResumeDeleteFailure,
ResumeManager, ResumeUtils, compose_key,
CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, RESUME_CHECKPOINT_FILE,
ReplacementTargetIdentity, ResumeDeleteFailure, ResumeManager, ResumeUtils, compose_key,
};
use crate::heal::storage::{HealLifecycleExpiryContext, HealListItem, HealObjectInfo, HealStorageAPI};
use crate::heal::storage_api::status::BucketInfo;
+3 -83
View File
@@ -2011,91 +2011,11 @@ impl HealManager {
return None;
}
let mut snapshot = HealProgress::default();
let mut has_object_sweep = false;
let mut all_object_baselines_known = true;
let mut counter_overflow = false;
let mut stage_current = 0_u64;
let mut stage_total = 0_u64;
let mut progresses = Vec::with_capacity(active_tasks.len());
for task in active_tasks {
let progress = task.get_progress().await;
let object_sweep = matches!(progress.kind, crate::heal::progress::HealProgressKind::ObjectSweep);
has_object_sweep |= object_sweep;
if object_sweep {
all_object_baselines_known &= progress.baseline_known;
}
counter_overflow |=
progress.counter_unknown || matches!(progress.progress_state, crate::heal::progress::HealProgressState::Unknown);
match stage_current.checked_add(progress.stage_current) {
Some(sum) => stage_current = sum,
None => counter_overflow = true,
}
match stage_total.checked_add(progress.stage_total) {
Some(sum) => stage_total = sum,
None => counter_overflow = true,
}
for (target, value) in [
(&mut snapshot.objects_scanned, progress.objects_scanned),
(&mut snapshot.objects_healed, progress.objects_healed),
(&mut snapshot.objects_failed, progress.objects_failed),
(&mut snapshot.skipped_objects, progress.skipped_objects),
(&mut snapshot.skipped_new_versions, progress.skipped_new_versions),
(&mut snapshot.skipped_ilm_expired, progress.skipped_ilm_expired),
(&mut snapshot.objects_total_count, progress.objects_total_count),
(&mut snapshot.objects_total_size, progress.objects_total_size),
(&mut snapshot.bytes_processed, progress.bytes_processed),
] {
match target.checked_add(value) {
Some(sum) => *target = sum,
None => counter_overflow = true,
}
}
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
(current, None) => current,
};
snapshot.last_update_time = match (snapshot.last_update_time, progress.last_update_time) {
(Some(current), Some(next)) => Some(current.max(next)),
(None, next) => next,
(current, None) => current,
};
if progress.current_object.is_some() {
snapshot.current_object = progress.current_object;
}
progresses.push(task.get_progress().await);
}
snapshot.kind = if has_object_sweep {
crate::heal::progress::HealProgressKind::ObjectSweep
} else {
crate::heal::progress::HealProgressKind::Stage
};
snapshot.stage_current = stage_current;
snapshot.stage_total = stage_total;
snapshot.baseline_known = has_object_sweep && all_object_baselines_known;
snapshot.progress_state = if counter_overflow {
crate::heal::progress::HealProgressState::Unknown
} else if has_object_sweep && !all_object_baselines_known {
crate::heal::progress::HealProgressState::Indeterminate
} else if has_object_sweep {
crate::heal::progress::HealProgressState::Running
} else if stage_total == 0 {
crate::heal::progress::HealProgressState::Indeterminate
} else {
crate::heal::progress::HealProgressState::Running
};
if counter_overflow {
snapshot.progress_percentage = 0.0;
} else if !has_object_sweep {
snapshot.progress_percentage = if stage_total == 0 {
0.0
} else {
((stage_current as f64 / stage_total as f64) * 100.0).min(99.999)
};
} else {
snapshot.refresh_progress_percentage();
}
snapshot.refresh_estimated_completion_time();
Some(snapshot)
crate::heal::progress::aggregate_heal_progress(progresses)
}
}
+122 -2
View File
@@ -65,8 +65,8 @@ pub enum HealProgressState {
Completed,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
#[serde(default, rename_all = "camelCase")]
pub struct HealProgress {
#[serde(default)]
pub kind: HealProgressKind,
@@ -382,6 +382,101 @@ impl HealProgress {
}
}
pub fn aggregate_heal_progress(progresses: impl IntoIterator<Item = HealProgress>) -> Option<HealProgress> {
let mut snapshot = HealProgress::default();
let mut found = false;
let mut has_object_sweep = false;
let mut all_object_baselines_known = true;
let mut baseline_generation = None;
let mut baseline_generation_consistent = true;
let mut all_ledgers_complete = true;
let mut counter_overflow = false;
for progress in progresses {
found = true;
let object_sweep = matches!(progress.kind, HealProgressKind::ObjectSweep);
has_object_sweep |= object_sweep;
all_ledgers_complete &= progress.ledger_complete;
if object_sweep {
all_object_baselines_known &= progress.baseline_known;
match baseline_generation {
None => baseline_generation = Some(progress.baseline_generation),
Some(generation) => baseline_generation_consistent &= generation == progress.baseline_generation,
}
}
counter_overflow |= progress.counter_unknown || matches!(progress.progress_state, HealProgressState::Unknown);
for (target, value) in [
(&mut snapshot.objects_scanned, progress.objects_scanned),
(&mut snapshot.objects_healed, progress.objects_healed),
(&mut snapshot.objects_failed, progress.objects_failed),
(&mut snapshot.skipped_objects, progress.skipped_objects),
(&mut snapshot.skipped_new_versions, progress.skipped_new_versions),
(&mut snapshot.skipped_ilm_expired, progress.skipped_ilm_expired),
(&mut snapshot.objects_total_count, progress.objects_total_count),
(&mut snapshot.objects_total_size, progress.objects_total_size),
(&mut snapshot.bytes_processed, progress.bytes_processed),
(&mut snapshot.stage_current, progress.stage_current),
(&mut snapshot.stage_total, progress.stage_total),
] {
match target.checked_add(value) {
Some(sum) => *target = sum,
None => {
*target = u64::MAX;
counter_overflow = true;
}
}
}
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
(current, None) => current,
};
snapshot.last_update_time = match (snapshot.last_update_time, progress.last_update_time) {
(Some(current), Some(next)) => Some(current.max(next)),
(None, next) => next,
(current, None) => current,
};
if progress.current_object.is_some() {
snapshot.current_object = progress.current_object;
}
}
if !found {
return None;
}
snapshot.kind = if has_object_sweep {
HealProgressKind::ObjectSweep
} else {
HealProgressKind::Stage
};
snapshot.baseline_known = has_object_sweep && all_object_baselines_known && baseline_generation_consistent;
snapshot.baseline_generation = if snapshot.baseline_known && baseline_generation_consistent {
baseline_generation.flatten()
} else {
None
};
snapshot.ledger_complete = all_ledgers_complete;
snapshot.counter_unknown = counter_overflow;
if counter_overflow {
snapshot.progress_state = HealProgressState::Unknown;
snapshot.progress_percentage = if snapshot.ledger_complete { 100.0 } else { 0.0 };
} else if snapshot.ledger_complete {
snapshot.progress_state = HealProgressState::Completed;
snapshot.progress_percentage = 100.0;
} else if has_object_sweep {
snapshot.refresh_progress_percentage();
} else if snapshot.stage_total == 0 {
snapshot.progress_state = HealProgressState::Indeterminate;
snapshot.progress_percentage = 0.0;
} else {
snapshot.progress_state = HealProgressState::Running;
snapshot.progress_percentage = ((snapshot.stage_current as f64 / snapshot.stage_total as f64) * 100.0).min(99.999);
}
snapshot.refresh_estimated_completion_time();
Some(snapshot)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealStatistics {
/// Total heal tasks
@@ -725,6 +820,31 @@ mod tests {
progress.mark_completed();
assert!(progress.is_completed());
assert_eq!(progress.progress_state, HealProgressState::Unknown);
let aggregate = aggregate_heal_progress([progress]).expect("progress should aggregate");
assert!(aggregate.ledger_complete);
assert!(aggregate.counter_unknown);
assert_eq!(aggregate.progress_state, HealProgressState::Unknown);
assert_eq!(aggregate.progress_percentage, 100.0);
}
#[test]
fn aggregate_rejects_mixed_baseline_generations() {
let progress = |generation| HealProgress {
kind: HealProgressKind::ObjectSweep,
objects_scanned: 5,
objects_total_count: 10,
progress_state: HealProgressState::Running,
baseline_generation: Some(generation),
baseline_known: true,
..Default::default()
};
let aggregate = aggregate_heal_progress([progress(1), progress(2)]).expect("progress should aggregate");
assert!(!aggregate.baseline_known);
assert_eq!(aggregate.baseline_generation, None);
assert_eq!(aggregate.progress_state, HealProgressState::Indeterminate);
assert_eq!(aggregate.progress_percentage, 0.0);
}
#[test]
+7 -8
View File
@@ -29,10 +29,10 @@ use super::{
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
/// Current on-disk schema version for `ResumeCheckpoint`. Schema 5 could
/// persist dedup identities without the aggregate counters needed to restore
/// them safely, so stale checkpoints are discarded and replayed.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6;
#[derive(Debug, Clone, Copy)]
pub enum CheckpointObjectOutcome {
@@ -252,10 +252,9 @@ impl CheckpointManager {
});
}
// A checkpoint from an older schema stored latest-only dedup identities
// that are not comparable to the new per-version `compose_key`
// identities. Discard the stale sets and position, then stamp the
// current schema so the scan restarts cleanly.
// Older checkpoints can contain identities that are not comparable to
// the current keys or lack their corresponding aggregate counters.
// Discard the stale sets and position so the scan restarts cleanly.
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
return Err(Error::TaskExecutionFailed {
message: format!(
+4 -4
View File
@@ -1603,14 +1603,14 @@ async fn test_resumestate_schema_v0_discarded_on_load() {
}
#[tokio::test]
async fn test_checkpoint_schema_v4_discarded_on_load() {
async fn test_checkpoint_schema_v5_discarded_on_load() {
let (temp_dir, disk) = schema_test_disk().await;
// The previous checkpoint schema is unsafe once its paired resume
// state is discarded: retaining either position would skip work.
// Schema v5 can persist failed identities without the aggregate counters
// that make those identities safe to deduplicate after an upgrade.
let task_id = "00000000-0000-4000-8000-000000000002";
let legacy = r#"{
"schema_version": 4,
"schema_version": 5,
"task_id": "00000000-0000-4000-8000-000000000002",
"checkpoint_time": 1700000000,
"current_bucket_index": 2,
+3
View File
@@ -445,6 +445,9 @@ impl HealTask {
}
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(());
}
let baseline = match self
.await_with_control(self.storage.erasure_set_usage_baseline(buckets))
.await
+44
View File
@@ -22,6 +22,7 @@ use std::sync::Mutex;
use tempfile::TempDir;
use super::super::storage_api::status::BucketInfo;
use crate::heal::progress::HealProgressState;
#[tokio::test]
async fn retry_request_carries_remaining_timeout_budget() {
@@ -2154,6 +2155,49 @@ async fn erasure_set_heal_applies_usage_baseline_to_progress() {
assert!((progress.progress_percentage - 25.0).abs() < 0.001);
}
#[tokio::test]
async fn erasure_set_disk_walk_keeps_cluster_usage_baseline_indeterminate() {
for (scan_mode, source) in [
(HealScanMode::Deep, HealRequestSource::Admin),
(HealScanMode::Normal, HealRequestSource::AutoHeal),
] {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline: Mutex::new(Some(HealBucketUsageBaseline {
objects_count: 10,
bytes: 8,
generation: Some(1),
})),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
scan_mode,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = source;
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("erasure set heal should complete");
let progress = task.get_progress().await;
assert!(!progress.baseline_known);
assert_eq!(progress.baseline_generation, None);
assert_eq!(progress.progress_state, HealProgressState::Indeterminate);
}
}
#[tokio::test]
async fn erasure_set_heal_ignores_usage_baseline_errors() {
let temp = TempDir::new().expect("temporary directory should be created");
+1 -1
View File
@@ -19,7 +19,7 @@ pub use error::{Error, Result};
pub use heal::{
HealManager, HealOperationsSnapshot, HealOptions, HealPriority, HealPriorityCounts, HealRequest, HealSourceCounts, HealType,
channel::HealChannelProcessor,
progress::HealProgress,
progress::{HealProgress, aggregate_heal_progress},
resume::{ReplacementRecoveryRecord, ReplacementRecoveryState, ResumeUtils},
};
use rustfs_concurrency::WorkloadAdmissionSnapshotProvider;
@@ -1215,7 +1215,10 @@ pub struct ScannerActivityResponse {
pub dirty_usage_pending: bool,
}
#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
pub struct BackgroundHealStatusRequest {}
pub struct BackgroundHealStatusRequest {
#[prost(uint32, tag = "1")]
pub protocol_version: u32,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct BackgroundHealStatusResponse {
#[prost(bool, tag = "1")]
+19
View File
@@ -170,6 +170,7 @@ pub fn internode_rpc_max_message_size() -> usize {
pub const HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE: usize = heal_control::RESULT_MAX_SIZE + 1024;
pub const HEAL_CONTROL_PROTOCOL_VERSION: u32 = 3;
pub const DYNAMIC_CONFIG_PROTOCOL_VERSION: u32 = 1;
pub const BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION: u32 = 2;
pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-capability-v3\0";
pub const REMOTE_VERSION_STATE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-tier-remote-version-state-capability-v1\0";
pub const CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-cross-pool-fence-capability-v1\0";
@@ -2346,10 +2347,28 @@ pub async fn evict_failed_connection_with_log_level(addr: &str, log_level: Conne
#[cfg(test)]
mod tests {
use super::*;
use prost::Message as _;
use std::sync::Mutex;
static INTERNODE_RPC_MSGPACK_ONLY_ENV_LOCK: Mutex<()> = Mutex::new(());
#[derive(Clone, PartialEq, prost::Message)]
struct BackgroundHealStatusRequestV1 {}
#[test]
fn background_heal_status_request_remains_rolling_upgrade_compatible() {
let current = proto_gen::node_service::BackgroundHealStatusRequest {
protocol_version: BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION,
};
let encoded = current.encode_to_vec();
BackgroundHealStatusRequestV1::decode(encoded.as_slice()).expect("v1 server should ignore the version field");
let encoded = BackgroundHealStatusRequestV1 {}.encode_to_vec();
let decoded = proto_gen::node_service::BackgroundHealStatusRequest::decode(encoded.as_slice())
.expect("v2 server should accept a v1 request");
assert_eq!(decoded.protocol_version, 0);
}
#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
struct CompatPayloadField {
message: &'static str,
+3 -1
View File
@@ -846,7 +846,9 @@ message ScannerActivityResponse {
bool dirty_usage_pending = 9;
}
message BackgroundHealStatusRequest {}
message BackgroundHealStatusRequest {
uint32 protocol_version = 1;
}
message BackgroundHealStatusResponse {
bool success = 1;