mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 03:59:14 +00:00
feat(heal): record bounded canonical object and task outcomes
Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -447,6 +447,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"running".to_string(),
|
"running".to_string(),
|
||||||
None,
|
None,
|
||||||
@@ -463,6 +464,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"running".to_string(),
|
"running".to_string(),
|
||||||
Some(format!("heal task retrying after recoverable failure, attempt {retry_attempt}: {error}")),
|
Some(format!("heal task retrying after recoverable failure, attempt {retry_attempt}: {error}")),
|
||||||
@@ -479,6 +481,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"finished".to_string(),
|
"finished".to_string(),
|
||||||
None,
|
None,
|
||||||
@@ -495,6 +498,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"stopped".to_string(),
|
"stopped".to_string(),
|
||||||
Some("heal task cancelled".to_string()),
|
Some("heal task cancelled".to_string()),
|
||||||
@@ -511,6 +515,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"stopped".to_string(),
|
"stopped".to_string(),
|
||||||
Some("heal task timed out".to_string()),
|
Some("heal task timed out".to_string()),
|
||||||
@@ -527,6 +532,7 @@ impl HealChannelProcessor {
|
|||||||
progress,
|
progress,
|
||||||
next_seq,
|
next_seq,
|
||||||
min_seq,
|
min_seq,
|
||||||
|
..
|
||||||
}) => (
|
}) => (
|
||||||
"stopped".to_string(),
|
"stopped".to_string(),
|
||||||
Some(error),
|
Some(error),
|
||||||
|
|||||||
@@ -13,6 +13,7 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use crate::heal::{
|
use crate::heal::{
|
||||||
|
outcome::HealTaskOutcome,
|
||||||
progress::{HealProgress, HealStatistics},
|
progress::{HealProgress, HealStatistics},
|
||||||
resume::{ReplacementPhase, ResumeGc, ResumeManager, ResumeState, ResumeUtils},
|
resume::{ReplacementPhase, ResumeGc, ResumeManager, ResumeState, ResumeUtils},
|
||||||
storage::HealStorageAPI,
|
storage::HealStorageAPI,
|
||||||
@@ -185,6 +186,7 @@ fn record_displaced_terminal(
|
|||||||
request: &HealRequest,
|
request: &HealRequest,
|
||||||
) -> Arc<CompletedHealStatus> {
|
) -> Arc<CompletedHealStatus> {
|
||||||
let terminal = Arc::new(CompletedHealStatus {
|
let terminal = Arc::new(CompletedHealStatus {
|
||||||
|
outcome: None,
|
||||||
progress: None,
|
progress: None,
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
heal_type: request.heal_type.clone(),
|
heal_type: request.heal_type.clone(),
|
||||||
@@ -268,6 +270,7 @@ async fn publish_completed_heal(
|
|||||||
|
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct HealTaskReport {
|
pub struct HealTaskReport {
|
||||||
|
pub outcome: Option<Arc<HealTaskOutcome>>,
|
||||||
pub status: HealTaskStatus,
|
pub status: HealTaskStatus,
|
||||||
pub result_items: Vec<HealResultItem>,
|
pub result_items: Vec<HealResultItem>,
|
||||||
pub result_items_truncated: bool,
|
pub result_items_truncated: bool,
|
||||||
@@ -285,6 +288,7 @@ async fn active_task_report(task: &HealTask, since: Option<u64>) -> HealTaskRepo
|
|||||||
let window = task.get_result_items_since(since).await;
|
let window = task.get_result_items_since(since).await;
|
||||||
HealTaskReport {
|
HealTaskReport {
|
||||||
status: task.get_status().await,
|
status: task.get_status().await,
|
||||||
|
outcome: Some(Arc::new(task.get_outcome().await)),
|
||||||
result_items: window.items,
|
result_items: window.items,
|
||||||
// The legacy flag stays set once anything was evicted; a lagging
|
// The legacy flag stays set once anything was evicted; a lagging
|
||||||
// incremental cursor additionally marks this response truncated so
|
// incremental cursor additionally marks this response truncated so
|
||||||
@@ -298,6 +302,7 @@ async fn active_task_report(task: &HealTask, since: Option<u64>) -> HealTaskRepo
|
|||||||
|
|
||||||
fn empty_task_report(status: HealTaskStatus) -> HealTaskReport {
|
fn empty_task_report(status: HealTaskStatus) -> HealTaskReport {
|
||||||
HealTaskReport {
|
HealTaskReport {
|
||||||
|
outcome: None,
|
||||||
status,
|
status,
|
||||||
result_items: Vec::new(),
|
result_items: Vec::new(),
|
||||||
result_items_truncated: false,
|
result_items_truncated: false,
|
||||||
@@ -325,6 +330,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option<u64>) ->
|
|||||||
};
|
};
|
||||||
HealTaskReport {
|
HealTaskReport {
|
||||||
status: completed.status.clone(),
|
status: completed.status.clone(),
|
||||||
|
outcome: completed.outcome.clone(),
|
||||||
result_items,
|
result_items,
|
||||||
result_items_truncated: completed.result_items_truncated || lagged,
|
result_items_truncated: completed.result_items_truncated || lagged,
|
||||||
progress: completed.progress.clone(),
|
progress: completed.progress.clone(),
|
||||||
|
|||||||
@@ -83,6 +83,7 @@ pub(super) struct CompletedHealStatus {
|
|||||||
pub(super) heal_type: HealType,
|
pub(super) heal_type: HealType,
|
||||||
pub(super) status: HealTaskStatus,
|
pub(super) status: HealTaskStatus,
|
||||||
pub(super) progress: Option<HealProgress>,
|
pub(super) progress: Option<HealProgress>,
|
||||||
|
pub(super) outcome: Option<Arc<HealTaskOutcome>>,
|
||||||
pub(super) retained_bytes: std::sync::OnceLock<usize>,
|
pub(super) retained_bytes: std::sync::OnceLock<usize>,
|
||||||
pub(super) result_items_truncated: bool,
|
pub(super) result_items_truncated: bool,
|
||||||
pub(super) completed_at: SystemTime,
|
pub(super) completed_at: SystemTime,
|
||||||
@@ -105,6 +106,7 @@ impl CompletedHealStatus {
|
|||||||
fn measure_retained_bytes(&self) -> usize {
|
fn measure_retained_bytes(&self) -> usize {
|
||||||
let mut bytes = size_of::<Self>();
|
let mut bytes = size_of::<Self>();
|
||||||
let mut add = |amount: usize| bytes = bytes.saturating_add(amount);
|
let mut add = |amount: usize| bytes = bytes.saturating_add(amount);
|
||||||
|
add(self.outcome.as_ref().map_or(0, |outcome| outcome.retained_bytes()));
|
||||||
match &self.heal_type {
|
match &self.heal_type {
|
||||||
HealType::Cluster => {}
|
HealType::Cluster => {}
|
||||||
HealType::Bucket { bucket } => add(bucket.capacity()),
|
HealType::Bucket { bucket } => add(bucket.capacity()),
|
||||||
@@ -209,6 +211,7 @@ impl CompletedHealStatus {
|
|||||||
heal_type: task.heal_type.clone(),
|
heal_type: task.heal_type.clone(),
|
||||||
status,
|
status,
|
||||||
progress: Some(task.get_progress().await),
|
progress: Some(task.get_progress().await),
|
||||||
|
outcome: Some(Arc::new(task.get_outcome().await)),
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
result_items_truncated: task.result_items_truncated(),
|
result_items_truncated: task.result_items_truncated(),
|
||||||
completed_at: SystemTime::now(),
|
completed_at: SystemTime::now(),
|
||||||
|
|||||||
@@ -298,6 +298,7 @@ impl HealManager {
|
|||||||
if cancelled_completion {
|
if cancelled_completion {
|
||||||
completed_status = HealTaskStatus::Cancelled;
|
completed_status = HealTaskStatus::Cancelled;
|
||||||
completed_status_entry.status = HealTaskStatus::Cancelled;
|
completed_status_entry.status = HealTaskStatus::Cancelled;
|
||||||
|
completed_status_entry.outcome = Some(Arc::new(task.get_outcome().await));
|
||||||
}
|
}
|
||||||
let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. });
|
let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. });
|
||||||
let successful_completion = matches!(completed_status, HealTaskStatus::Completed);
|
let successful_completion = matches!(completed_status, HealTaskStatus::Completed);
|
||||||
|
|||||||
@@ -103,6 +103,7 @@ struct MockStorage;
|
|||||||
|
|
||||||
fn completed_retention_fixture(completed_at: SystemTime) -> CompletedHealStatus {
|
fn completed_retention_fixture(completed_at: SystemTime) -> CompletedHealStatus {
|
||||||
CompletedHealStatus {
|
CompletedHealStatus {
|
||||||
|
outcome: None,
|
||||||
heal_type: HealType::Cluster,
|
heal_type: HealType::Cluster,
|
||||||
status: HealTaskStatus::Completed,
|
status: HealTaskStatus::Completed,
|
||||||
progress: Some(HealProgress {
|
progress: Some(HealProgress {
|
||||||
@@ -325,6 +326,10 @@ async fn completed_retention_cancel_wins_over_a_prepared_retry_snapshot() {
|
|||||||
for token in [&task_id, &alias] {
|
for token in [&task_id, &alias] {
|
||||||
let report = manager.get_task_report(token).await.expect("cancelled token retained");
|
let report = manager.get_task_report(token).await.expect("cancelled token retained");
|
||||||
assert_eq!(report.status, HealTaskStatus::Cancelled);
|
assert_eq!(report.status, HealTaskStatus::Cancelled);
|
||||||
|
assert_eq!(
|
||||||
|
report.outcome.as_ref().expect("cancelled outcome retained").execution,
|
||||||
|
crate::heal::outcome::HealExecutionOutcome::Aborted(crate::heal::outcome::HealAbortReason::Cancelled)
|
||||||
|
);
|
||||||
assert_eq!(report.progress.expect("frozen progress").objects_scanned, 1);
|
assert_eq!(report.progress.expect("frozen progress").objects_scanned, 1);
|
||||||
}
|
}
|
||||||
assert!(!manager.retrying_heals.lock().await.contains_key(&task_id));
|
assert!(!manager.retrying_heals.lock().await.contains_key(&task_id));
|
||||||
@@ -391,6 +396,7 @@ async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_han
|
|||||||
.expect("scheduler archives terminal");
|
.expect("scheduler archives terminal");
|
||||||
assert!(!manager.active_heals.lock().await.contains_key(&task_id));
|
assert!(!manager.active_heals.lock().await.contains_key(&task_id));
|
||||||
let expected = task.get_progress().await;
|
let expected = task.get_progress().await;
|
||||||
|
let expected_outcome = task.get_outcome().await;
|
||||||
for token in [&task_id, &alias] {
|
for token in [&task_id, &alias] {
|
||||||
assert_eq!(manager.get_task_progress(token).await.expect("terminal progress query"), expected);
|
assert_eq!(manager.get_task_progress(token).await.expect("terminal progress query"), expected);
|
||||||
let report = manager
|
let report = manager
|
||||||
@@ -398,6 +404,7 @@ async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_han
|
|||||||
.await
|
.await
|
||||||
.expect("terminal token remains queryable at handoff");
|
.expect("terminal token remains queryable at handoff");
|
||||||
assert_eq!(report.progress.as_ref(), Some(&expected));
|
assert_eq!(report.progress.as_ref(), Some(&expected));
|
||||||
|
assert_eq!(report.outcome.as_deref(), Some(&expected_outcome));
|
||||||
assert!(report.result_items.is_empty());
|
assert!(report.result_items.is_empty());
|
||||||
match outcome {
|
match outcome {
|
||||||
"success" => assert_eq!(report.status, HealTaskStatus::Completed),
|
"success" => assert_eq!(report.status, HealTaskStatus::Completed),
|
||||||
@@ -1976,6 +1983,7 @@ async fn insert_retrying_request(manager: &HealManager, request: HealRequest) ->
|
|||||||
task_id,
|
task_id,
|
||||||
Arc::new(CompletedHealStatus {
|
Arc::new(CompletedHealStatus {
|
||||||
progress: None,
|
progress: None,
|
||||||
|
outcome: None,
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
heal_type: request.heal_type,
|
heal_type: request.heal_type,
|
||||||
status: HealTaskStatus::Retrying {
|
status: HealTaskStatus::Retrying {
|
||||||
@@ -2700,6 +2708,7 @@ async fn test_retrying_completion_outranks_the_queue_for_the_same_id() {
|
|||||||
task_id.clone(),
|
task_id.clone(),
|
||||||
Arc::new(CompletedHealStatus {
|
Arc::new(CompletedHealStatus {
|
||||||
progress: None,
|
progress: None,
|
||||||
|
outcome: None,
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
heal_type: request.heal_type.clone(),
|
heal_type: request.heal_type.clone(),
|
||||||
status: HealTaskStatus::Retrying {
|
status: HealTaskStatus::Retrying {
|
||||||
@@ -2737,6 +2746,7 @@ async fn test_get_task_status_reads_recent_completed_status() {
|
|||||||
"completed-token".to_string(),
|
"completed-token".to_string(),
|
||||||
Arc::new(CompletedHealStatus {
|
Arc::new(CompletedHealStatus {
|
||||||
progress: None,
|
progress: None,
|
||||||
|
outcome: None,
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
heal_type: HealType::Bucket {
|
heal_type: HealType::Bucket {
|
||||||
bucket: "bucket".to_string(),
|
bucket: "bucket".to_string(),
|
||||||
@@ -2768,6 +2778,7 @@ async fn test_get_task_report_for_path_reads_completed_items() {
|
|||||||
"completed-token".to_string(),
|
"completed-token".to_string(),
|
||||||
Arc::new(CompletedHealStatus {
|
Arc::new(CompletedHealStatus {
|
||||||
progress: None,
|
progress: None,
|
||||||
|
outcome: None,
|
||||||
retained_bytes: std::sync::OnceLock::new(),
|
retained_bytes: std::sync::OnceLock::new(),
|
||||||
heal_type: HealType::Object {
|
heal_type: HealType::Object {
|
||||||
bucket: "bucket".to_string(),
|
bucket: "bucket".to_string(),
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ pub mod channel;
|
|||||||
pub mod erasure_healer;
|
pub mod erasure_healer;
|
||||||
pub mod manager;
|
pub mod manager;
|
||||||
pub mod mrf_queue;
|
pub mod mrf_queue;
|
||||||
|
pub mod outcome;
|
||||||
pub mod progress;
|
pub mod progress;
|
||||||
pub(crate) mod replacement_readiness;
|
pub(crate) mod replacement_readiness;
|
||||||
pub mod resume;
|
pub mod resume;
|
||||||
|
|||||||
@@ -0,0 +1,300 @@
|
|||||||
|
// Copyright 2026 RustFS Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
//! Execution results are separate from repair responsibility. A legacy
|
||||||
|
//! successful storage call supplies no authoritative repair receipt.
|
||||||
|
|
||||||
|
use std::{collections::VecDeque, time::SystemTime};
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
const MAX_OUTCOME_ITEMS: usize = 128;
|
||||||
|
const MAX_OUTCOME_BYTES: usize = 64 * 1024;
|
||||||
|
const MAX_OUTCOME_DETAIL_BYTES: usize = 1024;
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum HealObjectKind {
|
||||||
|
Object,
|
||||||
|
Metadata,
|
||||||
|
Decode,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct HealObjectIdentity {
|
||||||
|
pub kind: HealObjectKind,
|
||||||
|
pub bucket: String,
|
||||||
|
pub object: String,
|
||||||
|
/// The requested version; None remains unresolved, never an absence proof.
|
||||||
|
pub version_id: Option<String>,
|
||||||
|
pub bucket_incarnation_id: Option<Uuid>,
|
||||||
|
pub pool_index: Option<usize>,
|
||||||
|
pub set_index: Option<usize>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum HealDeferredReason {
|
||||||
|
DanglingDeleteGrace,
|
||||||
|
TransientUsageCache,
|
||||||
|
TransientExistenceCheck,
|
||||||
|
Deadline,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum HealFailureClass {
|
||||||
|
Recoverable,
|
||||||
|
RetryExhausted,
|
||||||
|
Permanent,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub enum HealObjectDisposition {
|
||||||
|
/// The legacy storage response does not prove the requested check or commit.
|
||||||
|
Unknown,
|
||||||
|
Repaired,
|
||||||
|
VerifiedHealthy,
|
||||||
|
AuthoritativelyAbsent,
|
||||||
|
Deferred {
|
||||||
|
reason: HealDeferredReason,
|
||||||
|
retry_not_before: Option<SystemTime>,
|
||||||
|
},
|
||||||
|
Failed(HealFailureClass),
|
||||||
|
Cancelled,
|
||||||
|
DryRunObserved,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct HealObjectOutcome {
|
||||||
|
pub identity: HealObjectIdentity,
|
||||||
|
pub disposition: HealObjectDisposition,
|
||||||
|
pub detail: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl HealObjectOutcome {
|
||||||
|
fn retained_bytes(&self) -> usize {
|
||||||
|
size_of::<Self>()
|
||||||
|
.saturating_add(self.identity.bucket.capacity())
|
||||||
|
.saturating_add(self.identity.object.capacity())
|
||||||
|
.saturating_add(self.identity.version_id.as_ref().map_or(0, String::capacity))
|
||||||
|
.saturating_add(self.detail.as_ref().map_or(0, String::capacity))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||||
|
pub enum HealTraversalCoverage {
|
||||||
|
#[default]
|
||||||
|
Unknown,
|
||||||
|
Partial,
|
||||||
|
Complete,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum HealAbortReason {
|
||||||
|
Cancelled,
|
||||||
|
Deadline,
|
||||||
|
Untraversable,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||||
|
pub enum HealExecutionOutcome {
|
||||||
|
#[default]
|
||||||
|
Pending,
|
||||||
|
Running,
|
||||||
|
Completed,
|
||||||
|
CompletedWithErrors,
|
||||||
|
Aborted(HealAbortReason),
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||||
|
pub struct HealOutcomeCounters {
|
||||||
|
pub processed: u64,
|
||||||
|
pub healed: u64,
|
||||||
|
pub unchanged: u64,
|
||||||
|
/// Deferred, cancelled, dry-run and unverified results remain unresolved.
|
||||||
|
pub skipped: u64,
|
||||||
|
pub failed: u64,
|
||||||
|
pub unknown: u64,
|
||||||
|
pub attempt_failures: u64,
|
||||||
|
pub overflowed: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||||
|
pub struct HealTaskOutcome {
|
||||||
|
pub execution: HealExecutionOutcome,
|
||||||
|
pub coverage: HealTraversalCoverage,
|
||||||
|
pub counters: HealOutcomeCounters,
|
||||||
|
/// A bounded diagnostic window, not a complete responsibility ledger.
|
||||||
|
pub objects: VecDeque<HealObjectOutcome>,
|
||||||
|
pub objects_truncated: bool,
|
||||||
|
retained_object_bytes: usize,
|
||||||
|
untraversable: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl HealTaskOutcome {
|
||||||
|
pub(crate) fn start(&mut self) {
|
||||||
|
self.execution = HealExecutionOutcome::Running;
|
||||||
|
self.coverage = HealTraversalCoverage::Partial;
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn attempt_failed(&mut self) {
|
||||||
|
self.counters.overflowed |= !super::progress::increment_counter(&mut self.counters.attempt_failures);
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn mark_untraversable(&mut self) {
|
||||||
|
self.untraversable = true;
|
||||||
|
self.coverage = HealTraversalCoverage::Partial;
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn finish(&mut self, abort: Option<HealAbortReason>) {
|
||||||
|
let abort = abort.or(self.untraversable.then_some(HealAbortReason::Untraversable));
|
||||||
|
self.execution = match abort {
|
||||||
|
Some(reason) => HealExecutionOutcome::Aborted(reason),
|
||||||
|
None if self.counters.failed > 0 => HealExecutionOutcome::CompletedWithErrors,
|
||||||
|
None => HealExecutionOutcome::Completed,
|
||||||
|
};
|
||||||
|
self.coverage = if abort.is_none() && !self.counters.overflowed {
|
||||||
|
HealTraversalCoverage::Complete
|
||||||
|
} else {
|
||||||
|
HealTraversalCoverage::Partial
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn record(&mut self, mut item: HealObjectOutcome) {
|
||||||
|
use super::progress::increment_counter;
|
||||||
|
let counters = &mut self.counters;
|
||||||
|
counters.overflowed |= !increment_counter(&mut counters.processed);
|
||||||
|
let counter = match item.disposition {
|
||||||
|
HealObjectDisposition::Repaired => &mut counters.healed,
|
||||||
|
HealObjectDisposition::VerifiedHealthy | HealObjectDisposition::AuthoritativelyAbsent => &mut counters.unchanged,
|
||||||
|
HealObjectDisposition::Failed(_) => &mut counters.failed,
|
||||||
|
HealObjectDisposition::Unknown => {
|
||||||
|
counters.overflowed |= !increment_counter(&mut counters.unknown);
|
||||||
|
&mut counters.skipped
|
||||||
|
}
|
||||||
|
_ => &mut counters.skipped,
|
||||||
|
};
|
||||||
|
counters.overflowed |= !increment_counter(counter);
|
||||||
|
if let Some(detail) = &mut item.detail {
|
||||||
|
let mut end = detail.len().min(MAX_OUTCOME_DETAIL_BYTES);
|
||||||
|
while !detail.is_char_boundary(end) {
|
||||||
|
end -= 1;
|
||||||
|
}
|
||||||
|
self.objects_truncated |= end < detail.len();
|
||||||
|
detail.truncate(end);
|
||||||
|
detail.shrink_to_fit();
|
||||||
|
}
|
||||||
|
let bytes = item.retained_bytes();
|
||||||
|
if bytes > MAX_OUTCOME_BYTES {
|
||||||
|
self.objects_truncated = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
while self.objects.len() >= MAX_OUTCOME_ITEMS || self.retained_object_bytes.saturating_add(bytes) > MAX_OUTCOME_BYTES {
|
||||||
|
let Some(oldest) = self.objects.pop_front() else { break };
|
||||||
|
self.retained_object_bytes = self.retained_object_bytes.saturating_sub(oldest.retained_bytes());
|
||||||
|
self.objects_truncated = true;
|
||||||
|
}
|
||||||
|
self.retained_object_bytes = self.retained_object_bytes.saturating_add(bytes);
|
||||||
|
self.objects.push_back(item);
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn retained_bytes(&self) -> usize {
|
||||||
|
size_of::<Self>()
|
||||||
|
.saturating_add(self.retained_object_bytes)
|
||||||
|
.saturating_add(self.objects.capacity().saturating_mul(size_of::<HealObjectOutcome>()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod canonical_outcome_tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
fn item(disposition: HealObjectDisposition) -> HealObjectOutcome {
|
||||||
|
HealObjectOutcome {
|
||||||
|
identity: HealObjectIdentity {
|
||||||
|
kind: HealObjectKind::Object,
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
object: "object".to_string(),
|
||||||
|
version_id: None,
|
||||||
|
bucket_incarnation_id: None,
|
||||||
|
pool_index: None,
|
||||||
|
set_index: None,
|
||||||
|
},
|
||||||
|
disposition,
|
||||||
|
detail: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn canonical_outcome_categories_have_one_terminal_count() {
|
||||||
|
let mut outcome = HealTaskOutcome::default();
|
||||||
|
for disposition in [
|
||||||
|
HealObjectDisposition::Unknown,
|
||||||
|
HealObjectDisposition::Repaired,
|
||||||
|
HealObjectDisposition::VerifiedHealthy,
|
||||||
|
HealObjectDisposition::AuthoritativelyAbsent,
|
||||||
|
HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::DanglingDeleteGrace,
|
||||||
|
retry_not_before: None,
|
||||||
|
},
|
||||||
|
HealObjectDisposition::Failed(HealFailureClass::Permanent),
|
||||||
|
HealObjectDisposition::Cancelled,
|
||||||
|
HealObjectDisposition::DryRunObserved,
|
||||||
|
] {
|
||||||
|
outcome.record(item(disposition));
|
||||||
|
}
|
||||||
|
let c = &outcome.counters;
|
||||||
|
assert_eq!((c.processed, c.healed, c.unchanged, c.skipped, c.failed, c.unknown), (8, 1, 2, 4, 1, 1));
|
||||||
|
assert_eq!(c.processed, c.healed + c.unchanged + c.skipped + c.failed);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn canonical_outcome_window_count_bytes_and_oversize_keep_total_counts() {
|
||||||
|
let mut outcome = HealTaskOutcome::default();
|
||||||
|
for _ in 0..MAX_OUTCOME_ITEMS {
|
||||||
|
outcome.record(item(HealObjectDisposition::Unknown));
|
||||||
|
}
|
||||||
|
assert_eq!(outcome.objects.len(), MAX_OUTCOME_ITEMS);
|
||||||
|
assert!(!outcome.objects_truncated);
|
||||||
|
outcome.record(item(HealObjectDisposition::Unknown));
|
||||||
|
assert_eq!(outcome.objects.len(), MAX_OUTCOME_ITEMS);
|
||||||
|
assert!(outcome.objects_truncated);
|
||||||
|
let mut oversized = item(HealObjectDisposition::Failed(HealFailureClass::Permanent));
|
||||||
|
oversized.identity.object = "x".repeat(MAX_OUTCOME_BYTES);
|
||||||
|
outcome.record(oversized);
|
||||||
|
assert_eq!(outcome.counters.processed, u64::try_from(MAX_OUTCOME_ITEMS + 2).expect("bounded count"));
|
||||||
|
assert_eq!(outcome.counters.failed, 1);
|
||||||
|
assert!(outcome.retained_object_bytes <= MAX_OUTCOME_BYTES);
|
||||||
|
for _ in 0..MAX_OUTCOME_ITEMS {
|
||||||
|
let mut failed = item(HealObjectDisposition::Failed(HealFailureClass::Permanent));
|
||||||
|
failed.detail = Some("\u{4fee}".repeat(MAX_OUTCOME_DETAIL_BYTES));
|
||||||
|
outcome.record(failed);
|
||||||
|
}
|
||||||
|
assert!(outcome.retained_object_bytes <= MAX_OUTCOME_BYTES);
|
||||||
|
assert!(outcome.objects.iter().all(|item| {
|
||||||
|
item.detail
|
||||||
|
.as_ref()
|
||||||
|
.is_none_or(|detail| detail.len() <= MAX_OUTCOME_DETAIL_BYTES)
|
||||||
|
}));
|
||||||
|
assert!(outcome.objects.len() < MAX_OUTCOME_ITEMS);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn canonical_outcome_counter_overflow_cannot_claim_complete_coverage() {
|
||||||
|
let mut outcome = HealTaskOutcome::default();
|
||||||
|
outcome.counters.processed = u64::MAX;
|
||||||
|
outcome.record(item(HealObjectDisposition::Unknown));
|
||||||
|
outcome.finish(None);
|
||||||
|
assert!(outcome.counters.overflowed);
|
||||||
|
assert_eq!(outcome.counters.processed, u64::MAX);
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Partial);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -15,6 +15,10 @@
|
|||||||
use crate::heal::{
|
use crate::heal::{
|
||||||
DiskError, EcstoreError, ErasureSetHealer, HealDiskExt as _,
|
DiskError, EcstoreError, ErasureSetHealer, HealDiskExt as _,
|
||||||
erasure_healer::target_outcomes_complete,
|
erasure_healer::target_outcomes_complete,
|
||||||
|
outcome::{
|
||||||
|
HealAbortReason, HealDeferredReason, HealFailureClass, HealObjectDisposition, HealObjectIdentity, HealObjectKind,
|
||||||
|
HealObjectOutcome, HealTaskOutcome,
|
||||||
|
},
|
||||||
progress::HealProgress,
|
progress::HealProgress,
|
||||||
resume::{
|
resume::{
|
||||||
CheckpointManager, ReplacementPhase, ReplacementTargetIdentity, ResumeManager, replacement_target_identities_match,
|
CheckpointManager, ReplacementPhase, ReplacementTargetIdentity, ResumeManager, replacement_target_identities_match,
|
||||||
@@ -394,6 +398,7 @@ pub struct HealTask {
|
|||||||
pub status: Arc<RwLock<HealTaskStatus>>,
|
pub status: Arc<RwLock<HealTaskStatus>>,
|
||||||
/// Progress tracking
|
/// Progress tracking
|
||||||
pub progress: Arc<RwLock<HealProgress>>,
|
pub progress: Arc<RwLock<HealProgress>>,
|
||||||
|
outcome: Arc<RwLock<HealTaskOutcome>>,
|
||||||
/// Result items collected from storage heal calls, each stamped with a
|
/// Result items collected from storage heal calls, each stamped with a
|
||||||
/// monotonically increasing sequence number for incremental consumption
|
/// monotonically increasing sequence number for incremental consumption
|
||||||
/// (the client passes the last seen seq back and receives only newer
|
/// (the client passes the last seen seq back and receives only newer
|
||||||
@@ -460,6 +465,7 @@ impl HealTask {
|
|||||||
result_items_truncated: Arc::new(AtomicBool::new(false)),
|
result_items_truncated: Arc::new(AtomicBool::new(false)),
|
||||||
batch_failure: Arc::new(RwLock::new(None)),
|
batch_failure: Arc::new(RwLock::new(None)),
|
||||||
batch_failure_recorded: Arc::new(AtomicBool::new(false)),
|
batch_failure_recorded: Arc::new(AtomicBool::new(false)),
|
||||||
|
outcome: Arc::new(RwLock::new(HealTaskOutcome::default())),
|
||||||
created_at: request.created_at,
|
created_at: request.created_at,
|
||||||
enqueued_at: request.enqueued_at,
|
enqueued_at: request.enqueued_at,
|
||||||
started_at: Arc::new(RwLock::new(None)),
|
started_at: Arc::new(RwLock::new(None)),
|
||||||
@@ -507,6 +513,66 @@ impl HealTask {
|
|||||||
self.heal_type.kind_label()
|
self.heal_type.kind_label()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn get_outcome(&self) -> HealTaskOutcome {
|
||||||
|
self.outcome.read().await.clone()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn outcome_identity(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
version_id: Option<&str>,
|
||||||
|
pool_index: Option<usize>,
|
||||||
|
set_index: Option<usize>,
|
||||||
|
) -> HealObjectIdentity {
|
||||||
|
HealObjectIdentity {
|
||||||
|
kind: match self.heal_type {
|
||||||
|
HealType::Metadata { .. } => HealObjectKind::Metadata,
|
||||||
|
HealType::ECDecode { .. } => HealObjectKind::Decode,
|
||||||
|
_ => HealObjectKind::Object,
|
||||||
|
},
|
||||||
|
bucket: bucket.to_owned(),
|
||||||
|
object: object.to_owned(),
|
||||||
|
version_id: version_id.map(ToOwned::to_owned),
|
||||||
|
bucket_incarnation_id: None,
|
||||||
|
pool_index,
|
||||||
|
set_index,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn single_object_identity(&self) -> Option<HealObjectIdentity> {
|
||||||
|
let (bucket, object, version) = match &self.heal_type {
|
||||||
|
HealType::Object {
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
version_id,
|
||||||
|
}
|
||||||
|
| HealType::ECDecode {
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
version_id,
|
||||||
|
} => (bucket, object, version_id.as_deref()),
|
||||||
|
HealType::Metadata { bucket, object } => (bucket, object, None),
|
||||||
|
_ => return None,
|
||||||
|
};
|
||||||
|
Some(self.outcome_identity(bucket, object, version, self.options.pool_index, self.options.set_index))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn record_deferred_object(&self, reason: HealDeferredReason) {
|
||||||
|
if let Some(identity) = self.single_object_identity() {
|
||||||
|
let mut outcome = self.outcome.write().await;
|
||||||
|
outcome.attempt_failed();
|
||||||
|
outcome.record(HealObjectOutcome {
|
||||||
|
identity,
|
||||||
|
disposition: HealObjectDisposition::Deferred {
|
||||||
|
reason,
|
||||||
|
retry_not_before: None,
|
||||||
|
},
|
||||||
|
detail: None,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn has_batch_failure(&self) -> bool {
|
pub(crate) fn has_batch_failure(&self) -> bool {
|
||||||
self.batch_failure_recorded.load(Ordering::Acquire)
|
self.batch_failure_recorded.load(Ordering::Acquire)
|
||||||
}
|
}
|
||||||
@@ -634,6 +700,7 @@ impl HealTask {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn skip_due_to_transient_object_exists(&self, bucket: &str, object: &str, err: &Error) -> Result<()> {
|
async fn skip_due_to_transient_object_exists(&self, bucket: &str, object: &str, err: &Error) -> Result<()> {
|
||||||
|
self.record_deferred_object(HealDeferredReason::TransientExistenceCheck).await;
|
||||||
warn!(
|
warn!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
event = EVENT_HEAL_OBJECT_RESULT,
|
event = EVENT_HEAL_OBJECT_RESULT,
|
||||||
@@ -733,6 +800,8 @@ impl HealTask {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
self.record_deferred_object(HealDeferredReason::TransientUsageCache).await;
|
||||||
|
|
||||||
warn!(
|
warn!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
event = EVENT_HEAL_OBJECT_RESULT,
|
event = EVENT_HEAL_OBJECT_RESULT,
|
||||||
@@ -755,6 +824,8 @@ impl HealTask {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
self.record_deferred_object(HealDeferredReason::DanglingDeleteGrace).await;
|
||||||
|
|
||||||
warn!(
|
warn!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
event = EVENT_HEAL_OBJECT_RESULT,
|
event = EVENT_HEAL_OBJECT_RESULT,
|
||||||
@@ -801,6 +872,7 @@ impl HealTask {
|
|||||||
#[tracing::instrument(skip(self), fields(task_id = %self.id, heal_type = ?self.heal_type))]
|
#[tracing::instrument(skip(self), fields(task_id = %self.id, heal_type = ?self.heal_type))]
|
||||||
#[hotpath::measure]
|
#[hotpath::measure]
|
||||||
pub async fn execute(&self) -> Result<()> {
|
pub async fn execute(&self) -> Result<()> {
|
||||||
|
self.outcome.write().await.start();
|
||||||
// update status and timestamps atomically to avoid race conditions
|
// update status and timestamps atomically to avoid race conditions
|
||||||
let now = SystemTime::now();
|
let now = SystemTime::now();
|
||||||
let start_instant = Instant::now();
|
let start_instant = Instant::now();
|
||||||
@@ -860,6 +932,43 @@ impl HealTask {
|
|||||||
HealType::ErasureSet { buckets, set_disk_id } => self.heal_erasure_set(buckets.clone(), set_disk_id.clone()).await,
|
HealType::ErasureSet { buckets, set_disk_id } => self.heal_erasure_set(buckets.clone(), set_disk_id.clone()).await,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
{
|
||||||
|
let mut outcome = self.outcome.write().await;
|
||||||
|
if outcome.counters.processed == 0
|
||||||
|
&& let Some(identity) = self.single_object_identity()
|
||||||
|
{
|
||||||
|
let disposition = match &result {
|
||||||
|
Ok(()) if self.options.dry_run => HealObjectDisposition::DryRunObserved,
|
||||||
|
Ok(()) => HealObjectDisposition::Unknown,
|
||||||
|
Err(Error::TaskCancelled) => HealObjectDisposition::Cancelled,
|
||||||
|
Err(Error::TaskTimeout) => HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::Deadline,
|
||||||
|
retry_not_before: None,
|
||||||
|
},
|
||||||
|
Err(error) => {
|
||||||
|
outcome.attempt_failed();
|
||||||
|
HealObjectDisposition::Failed(if error.is_recoverable_heal() {
|
||||||
|
HealFailureClass::Recoverable
|
||||||
|
} else {
|
||||||
|
HealFailureClass::Permanent
|
||||||
|
})
|
||||||
|
}
|
||||||
|
};
|
||||||
|
outcome.record(HealObjectOutcome {
|
||||||
|
identity,
|
||||||
|
disposition,
|
||||||
|
detail: result.as_ref().err().map(ToString::to_string),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
let abort = match &result {
|
||||||
|
Err(Error::TaskCancelled) => Some(HealAbortReason::Cancelled),
|
||||||
|
Err(Error::TaskTimeout) => Some(HealAbortReason::Deadline),
|
||||||
|
Err(_) if !self.has_batch_failure() && !self.heal_type.is_per_object() => Some(HealAbortReason::Untraversable),
|
||||||
|
_ => None,
|
||||||
|
};
|
||||||
|
outcome.finish(abort);
|
||||||
|
}
|
||||||
|
|
||||||
// update completed time and status
|
// update completed time and status
|
||||||
{
|
{
|
||||||
let mut completed_at = self.completed_at.write().await;
|
let mut completed_at = self.completed_at.write().await;
|
||||||
@@ -944,6 +1053,7 @@ impl HealTask {
|
|||||||
|
|
||||||
pub async fn cancel(&self) -> Result<()> {
|
pub async fn cancel(&self) -> Result<()> {
|
||||||
self.cancel_token.cancel();
|
self.cancel_token.cancel();
|
||||||
|
self.outcome.write().await.finish(Some(HealAbortReason::Cancelled));
|
||||||
let mut status = self.status.write().await;
|
let mut status = self.status.write().await;
|
||||||
*status = HealTaskStatus::Cancelled;
|
*status = HealTaskStatus::Cancelled;
|
||||||
debug!(
|
debug!(
|
||||||
|
|||||||
@@ -214,6 +214,7 @@ impl HealTask {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
failed = failed.saturating_add(1);
|
failed = failed.saturating_add(1);
|
||||||
|
self.outcome.write().await.mark_untraversable();
|
||||||
if err.is_recoverable_heal() {
|
if err.is_recoverable_heal() {
|
||||||
retryable = retryable.saturating_add(1);
|
retryable = retryable.saturating_add(1);
|
||||||
} else {
|
} else {
|
||||||
@@ -260,6 +261,7 @@ impl HealTask {
|
|||||||
|
|
||||||
#[hotpath::measure]
|
#[hotpath::measure]
|
||||||
async fn heal_bucket_objects(&self, bucket: &str, prefix: &str) -> Result<()> {
|
async fn heal_bucket_objects(&self, bucket: &str, prefix: &str) -> Result<()> {
|
||||||
|
let previous_progress = self.get_progress().await;
|
||||||
let mut scanned = 0u64;
|
let mut scanned = 0u64;
|
||||||
let mut healed = 0u64;
|
let mut healed = 0u64;
|
||||||
let mut failed = 0u64;
|
let mut failed = 0u64;
|
||||||
@@ -338,6 +340,14 @@ impl HealTask {
|
|||||||
self.check_control_flags().await?;
|
self.check_control_flags().await?;
|
||||||
let mut telemetry_unknown = false;
|
let mut telemetry_unknown = false;
|
||||||
let object = item.name.as_str();
|
let object = item.name.as_str();
|
||||||
|
let identity =
|
||||||
|
self.outcome_identity(bucket, object, item.version_id.as_deref(), heal_opts.pool, heal_opts.set);
|
||||||
|
let mut disposition = if heal_opts.dry_run {
|
||||||
|
HealObjectDisposition::DryRunObserved
|
||||||
|
} else {
|
||||||
|
HealObjectDisposition::Unknown
|
||||||
|
};
|
||||||
|
let mut detail = None;
|
||||||
{
|
{
|
||||||
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}")));
|
||||||
@@ -380,7 +390,31 @@ impl HealTask {
|
|||||||
};
|
};
|
||||||
|
|
||||||
if let Some(err) = error {
|
if let Some(err) = error {
|
||||||
|
match err {
|
||||||
|
Error::TaskCancelled | Error::TaskTimeout => {
|
||||||
|
let disposition = if matches!(err, Error::TaskCancelled) {
|
||||||
|
HealObjectDisposition::Cancelled
|
||||||
|
} else {
|
||||||
|
HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::Deadline,
|
||||||
|
retry_not_before: None,
|
||||||
|
}
|
||||||
|
};
|
||||||
|
self.outcome.write().await.record(HealObjectOutcome {
|
||||||
|
identity,
|
||||||
|
disposition,
|
||||||
|
detail: None,
|
||||||
|
});
|
||||||
|
return Err(err);
|
||||||
|
}
|
||||||
|
_ => self.outcome.write().await.attempt_failed(),
|
||||||
|
}
|
||||||
|
detail = Some(err.to_string());
|
||||||
if Self::is_dangling_delete_grace_error(&err) {
|
if Self::is_dangling_delete_grace_error(&err) {
|
||||||
|
disposition = HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::DanglingDeleteGrace,
|
||||||
|
retry_not_before: None,
|
||||||
|
};
|
||||||
telemetry_unknown |= !increment_counter(&mut skipped);
|
telemetry_unknown |= !increment_counter(&mut skipped);
|
||||||
warn!(
|
warn!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
@@ -395,6 +429,10 @@ impl HealTask {
|
|||||||
"Heal bucket object dangling cleanup deferred by grace window"
|
"Heal bucket object dangling cleanup deferred by grace window"
|
||||||
);
|
);
|
||||||
} else if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) {
|
} else if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) {
|
||||||
|
disposition = HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::TransientUsageCache,
|
||||||
|
retry_not_before: None,
|
||||||
|
};
|
||||||
telemetry_unknown |= !increment_counter(&mut skipped);
|
telemetry_unknown |= !increment_counter(&mut skipped);
|
||||||
warn!(
|
warn!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
@@ -425,6 +463,11 @@ impl HealTask {
|
|||||||
);
|
);
|
||||||
retry.push(item);
|
retry.push(item);
|
||||||
} else {
|
} else {
|
||||||
|
disposition = HealObjectDisposition::Failed(if err.is_recoverable_heal() {
|
||||||
|
HealFailureClass::RetryExhausted
|
||||||
|
} else {
|
||||||
|
HealFailureClass::Permanent
|
||||||
|
});
|
||||||
telemetry_unknown |= !increment_counter(&mut failed);
|
telemetry_unknown |= !increment_counter(&mut failed);
|
||||||
if err.is_recoverable_heal() {
|
if err.is_recoverable_heal() {
|
||||||
retryable_failed = retryable_failed.saturating_add(1);
|
retryable_failed = retryable_failed.saturating_add(1);
|
||||||
@@ -459,8 +502,20 @@ impl HealTask {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
self.outcome.write().await.record(HealObjectOutcome {
|
||||||
|
identity,
|
||||||
|
disposition,
|
||||||
|
detail,
|
||||||
|
});
|
||||||
|
|
||||||
let mut progress = self.progress.write().await;
|
let mut progress = self.progress.write().await;
|
||||||
progress.update_object_progress(scanned, healed, failed, skipped, bytes);
|
progress.update_object_progress(
|
||||||
|
previous_progress.objects_scanned.saturating_add(scanned),
|
||||||
|
previous_progress.objects_healed.saturating_add(healed),
|
||||||
|
previous_progress.objects_failed.saturating_add(failed),
|
||||||
|
previous_progress.skipped_objects.saturating_add(skipped),
|
||||||
|
previous_progress.bytes_processed.saturating_add(bytes),
|
||||||
|
);
|
||||||
if telemetry_unknown {
|
if telemetry_unknown {
|
||||||
progress.mark_unknown();
|
progress.mark_unknown();
|
||||||
}
|
}
|
||||||
@@ -475,7 +530,7 @@ impl HealTask {
|
|||||||
|
|
||||||
continuation_token = next_heal_listing_token(bucket, prefix, next_token, is_truncated)?;
|
continuation_token = next_heal_listing_token(bucket, prefix, next_token, is_truncated)?;
|
||||||
if continuation_token.is_none() {
|
if continuation_token.is_none() {
|
||||||
// Truncated but no continuation token: end of listing.
|
self.outcome.write().await.mark_untraversable();
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -261,8 +261,8 @@ impl HealTask {
|
|||||||
update_parity: true,
|
update_parity: true,
|
||||||
no_lock: self.options.no_lock,
|
no_lock: self.options.no_lock,
|
||||||
read_repair: false,
|
read_repair: false,
|
||||||
pool: None,
|
pool: self.options.pool_index,
|
||||||
set: None,
|
set: self.options.set_index,
|
||||||
};
|
};
|
||||||
|
|
||||||
let heal_result = self
|
let heal_result = self
|
||||||
|
|||||||
@@ -14,6 +14,294 @@
|
|||||||
|
|
||||||
use super::super::{DiskOption, DiskStore, Endpoint, new_disk};
|
use super::super::{DiskOption, DiskStore, Endpoint, new_disk};
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
|
mod canonical_outcome {
|
||||||
|
use super::*;
|
||||||
|
use crate::heal::outcome::{HealExecutionOutcome, HealTraversalCoverage};
|
||||||
|
|
||||||
|
fn bucket_task(storage: Arc<MockStorage>) -> HealTask {
|
||||||
|
HealTask::from_request(
|
||||||
|
HealRequest::new(
|
||||||
|
HealType::Bucket {
|
||||||
|
bucket: "bucket-a".to_string(),
|
||||||
|
},
|
||||||
|
HealOptions {
|
||||||
|
recursive: true,
|
||||||
|
timeout: None,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
HealPriority::Normal,
|
||||||
|
),
|
||||||
|
storage,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn listing_failure_preserves_processed_objects_and_partial_coverage() {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
fail_second_listing_page: true,
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let task = bucket_task(storage);
|
||||||
|
task.execute().await.expect_err("second page cannot be traversed");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.execution, HealExecutionOutcome::Aborted(HealAbortReason::Untraversable));
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Partial);
|
||||||
|
assert_eq!(outcome.counters.processed, 1);
|
||||||
|
assert_eq!(outcome.objects[0].identity.object, "object-a");
|
||||||
|
assert_eq!(task.get_progress().await.objects_scanned, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn cluster_preserves_cumulative_progress_across_buckets() {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
list_each_bucket: true,
|
||||||
|
listed_buckets: Mutex::new(Some(vec!["bucket-a".to_string(), "bucket-b".to_string()])),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let task = HealTask::from_request(
|
||||||
|
HealRequest::new(
|
||||||
|
HealType::Cluster,
|
||||||
|
HealOptions {
|
||||||
|
recursive: true,
|
||||||
|
timeout: None,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
HealPriority::Normal,
|
||||||
|
),
|
||||||
|
storage,
|
||||||
|
);
|
||||||
|
task.execute().await.expect("both buckets complete");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.counters.processed, 4);
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Complete);
|
||||||
|
let progress = task.get_progress().await;
|
||||||
|
assert_eq!((progress.objects_scanned, progress.objects_healed), (4, 4));
|
||||||
|
assert_eq!(
|
||||||
|
outcome
|
||||||
|
.objects
|
||||||
|
.iter()
|
||||||
|
.filter(|item| item.identity.bucket == "bucket-b")
|
||||||
|
.count(),
|
||||||
|
2
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(start_paused = true)]
|
||||||
|
async fn exhausted_object_does_not_abort_other_objects_or_erase_counts() {
|
||||||
|
let storage = Arc::new(MockStorage::default());
|
||||||
|
storage.heal_object_outcomes.lock().expect("outcomes").insert(
|
||||||
|
"object-a".to_string(),
|
||||||
|
(0..4).map(|_| MockHealObjectOutcome::RetryableReadQuorum).collect(),
|
||||||
|
);
|
||||||
|
let task = bucket_task(storage.clone());
|
||||||
|
task.execute().await.expect_err("legacy adapter retains batch failure");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.execution, HealExecutionOutcome::CompletedWithErrors);
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Complete);
|
||||||
|
assert_eq!((outcome.counters.processed, outcome.counters.failed, outcome.counters.unknown), (2, 1, 1));
|
||||||
|
assert_eq!(outcome.counters.attempt_failures, 4);
|
||||||
|
let failed = outcome
|
||||||
|
.objects
|
||||||
|
.iter()
|
||||||
|
.find(|item| item.identity.object == "object-a")
|
||||||
|
.expect("failed object");
|
||||||
|
assert_eq!(failed.disposition, HealObjectDisposition::Failed(HealFailureClass::RetryExhausted));
|
||||||
|
let calls = storage.heal_object_calls.lock().expect("calls");
|
||||||
|
assert_eq!(calls.iter().filter(|object| object.as_str() == "object-b").count(), 1);
|
||||||
|
let progress = task.get_progress().await;
|
||||||
|
assert_eq!((progress.objects_scanned, progress.objects_healed, progress.objects_failed), (2, 1, 1));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(start_paused = true)]
|
||||||
|
async fn retry_success_counts_one_terminal_outcome() {
|
||||||
|
let storage = Arc::new(MockStorage::default());
|
||||||
|
storage
|
||||||
|
.heal_object_outcomes
|
||||||
|
.lock()
|
||||||
|
.expect("outcomes")
|
||||||
|
.insert("object-a".to_string(), VecDeque::from([MockHealObjectOutcome::RetryableReadQuorum]));
|
||||||
|
let task = bucket_task(storage);
|
||||||
|
task.execute().await.expect("retry should recover");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.execution, HealExecutionOutcome::Completed);
|
||||||
|
assert_eq!(outcome.counters.processed, 2);
|
||||||
|
assert_eq!(outcome.counters.failed, 0);
|
||||||
|
assert_eq!(outcome.counters.attempt_failures, 1);
|
||||||
|
assert_eq!(
|
||||||
|
outcome
|
||||||
|
.objects
|
||||||
|
.iter()
|
||||||
|
.filter(|item| item.identity.object == "object-a")
|
||||||
|
.count(),
|
||||||
|
1
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
outcome.counters.processed,
|
||||||
|
outcome.counters.healed + outcome.counters.unchanged + outcome.counters.skipped + outcome.counters.failed
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn mixed_grace_and_legacy_success_keep_distinct_dispositions() {
|
||||||
|
let storage = Arc::new(MockStorage::default());
|
||||||
|
storage
|
||||||
|
.heal_object_outcomes
|
||||||
|
.lock()
|
||||||
|
.expect("outcomes")
|
||||||
|
.insert("object-a".to_string(), VecDeque::from([MockHealObjectOutcome::DanglingGraceDeferred]));
|
||||||
|
let task = bucket_task(storage);
|
||||||
|
task.execute().await.expect("grace permits traversal completion");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Complete);
|
||||||
|
assert_eq!(outcome.counters.processed, 2);
|
||||||
|
assert_eq!(outcome.counters.healed, 0, "legacy result is not a repair receipt");
|
||||||
|
assert!(matches!(
|
||||||
|
outcome.objects[0].disposition,
|
||||||
|
HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::DanglingDeleteGrace,
|
||||||
|
..
|
||||||
|
}
|
||||||
|
));
|
||||||
|
assert_eq!(outcome.objects[1].disposition, HealObjectDisposition::Unknown);
|
||||||
|
assert!(
|
||||||
|
outcome
|
||||||
|
.objects
|
||||||
|
.iter()
|
||||||
|
.all(|item| item.identity.bucket_incarnation_id.is_none())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
task.get_progress().await.objects_healed,
|
||||||
|
1,
|
||||||
|
"legacy display count remains distinct from proof"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn grace_single_object_is_completed_but_deferred() {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::DanglingGraceDeferred)),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let task = HealTask::from_request(HealRequest::object("bucket-a".to_string(), "recent.txt".to_string(), None), storage);
|
||||||
|
task.execute().await.expect("grace is deferred");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(task.get_status().await, HealTaskStatus::Completed);
|
||||||
|
assert_eq!(outcome.counters.processed, 1);
|
||||||
|
assert!(matches!(
|
||||||
|
outcome.objects[0].disposition,
|
||||||
|
HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::DanglingDeleteGrace,
|
||||||
|
..
|
||||||
|
}
|
||||||
|
));
|
||||||
|
assert_eq!(outcome.counters.attempt_failures, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn dry_run_and_transient_existence_do_not_prove_repair() {
|
||||||
|
for transient in [false, true] {
|
||||||
|
let storage = Arc::new(MockStorage::default());
|
||||||
|
if transient {
|
||||||
|
storage
|
||||||
|
.object_exists_by_name
|
||||||
|
.lock()
|
||||||
|
.expect("existence fixture")
|
||||||
|
.insert("object".to_string(), MockObjectExists::TransientSkip("retry later"));
|
||||||
|
}
|
||||||
|
let mut request = HealRequest::object("bucket-a".to_string(), "object".to_string(), None);
|
||||||
|
request.options.dry_run = !transient;
|
||||||
|
let task = HealTask::from_request(request, storage);
|
||||||
|
task.execute().await.expect("observation may complete");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.counters.healed, 0);
|
||||||
|
if transient {
|
||||||
|
assert!(matches!(
|
||||||
|
outcome.objects[0].disposition,
|
||||||
|
HealObjectDisposition::Deferred {
|
||||||
|
reason: HealDeferredReason::TransientExistenceCheck,
|
||||||
|
..
|
||||||
|
}
|
||||||
|
));
|
||||||
|
} else {
|
||||||
|
assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::DryRunObserved);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn untraversable_bucket_does_not_claim_complete_cluster_coverage() {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
listed_buckets: Mutex::new(Some(vec!["bucket-a".to_string(), "bucket-b".to_string()])),
|
||||||
|
bucket_heal_errors: Mutex::new(HashMap::from([("bucket-a".to_string(), VecDeque::from(["metadata unavailable"]))])),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let task = HealTask::from_request(
|
||||||
|
HealRequest::new(
|
||||||
|
HealType::Cluster,
|
||||||
|
HealOptions {
|
||||||
|
recursive: true,
|
||||||
|
timeout: None,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
HealPriority::Normal,
|
||||||
|
),
|
||||||
|
storage.clone(),
|
||||||
|
);
|
||||||
|
task.execute().await.expect_err("structural bucket error");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.execution, HealExecutionOutcome::Aborted(HealAbortReason::Untraversable));
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Partial);
|
||||||
|
assert_eq!(
|
||||||
|
storage.bucket_heal_calls.lock().expect("bucket calls").as_slice(),
|
||||||
|
["bucket-a", "bucket-b"]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(start_paused = true)]
|
||||||
|
async fn cancellation_and_deadline_leave_partial_coverage() {
|
||||||
|
for cancel in [false, true] {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
block_heal_object: Mutex::new(true),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let mut request = HealRequest::object("bucket-a".to_string(), "object".to_string(), None);
|
||||||
|
request.options.timeout = Some(Duration::from_secs(1));
|
||||||
|
let task = HealTask::from_request(request, storage);
|
||||||
|
if cancel {
|
||||||
|
task.cancel().await.expect("cancel request");
|
||||||
|
}
|
||||||
|
task.execute().await.expect_err("control interruption");
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
assert_eq!(outcome.coverage, HealTraversalCoverage::Partial);
|
||||||
|
assert_eq!(
|
||||||
|
outcome.execution,
|
||||||
|
HealExecutionOutcome::Aborted(if cancel {
|
||||||
|
HealAbortReason::Cancelled
|
||||||
|
} else {
|
||||||
|
HealAbortReason::Deadline
|
||||||
|
})
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn decode_keeps_the_requested_pool_and_set() {
|
||||||
|
let storage = Arc::new(MockStorage::default());
|
||||||
|
let mut request = HealRequest::ec_decode("bucket-a".to_string(), "object".to_string(), Some("version-a".to_string()));
|
||||||
|
request.options.pool_index = Some(2);
|
||||||
|
request.options.set_index = Some(3);
|
||||||
|
let task = HealTask::from_request(request, storage.clone());
|
||||||
|
task.execute().await.expect("decode fixture");
|
||||||
|
let options = storage.object_heal_opts.lock().expect("storage options");
|
||||||
|
assert_eq!((options[0].pool, options[0].set), (Some(2), Some(3)));
|
||||||
|
let outcome = task.get_outcome().await;
|
||||||
|
let identity = &outcome.objects[0].identity;
|
||||||
|
assert_eq!((identity.pool_index, identity.set_index), (Some(2), Some(3)));
|
||||||
|
assert_eq!(identity.version_id.as_deref(), Some("version-a"));
|
||||||
|
assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::Unknown);
|
||||||
|
}
|
||||||
|
}
|
||||||
use crate::heal::storage::{HealListItem, HealObjectInfo};
|
use crate::heal::storage::{HealListItem, HealObjectInfo};
|
||||||
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, TraceSubscription, TraceVal, subscribe_trace_events};
|
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, TraceSubscription, TraceVal, subscribe_trace_events};
|
||||||
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos};
|
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos};
|
||||||
@@ -582,6 +870,8 @@ async fn verified_recovery_keeps_state_when_marker_clear_fails() {
|
|||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
struct MockStorage {
|
struct MockStorage {
|
||||||
listed: Mutex<bool>,
|
listed: Mutex<bool>,
|
||||||
|
list_each_bucket: bool,
|
||||||
|
fail_second_listing_page: bool,
|
||||||
healed_objects: Mutex<Vec<String>>,
|
healed_objects: Mutex<Vec<String>>,
|
||||||
heal_object_calls: Mutex<Vec<String>>,
|
heal_object_calls: Mutex<Vec<String>>,
|
||||||
heal_object_version_ids: Mutex<Vec<Option<String>>>,
|
heal_object_version_ids: Mutex<Vec<Option<String>>>,
|
||||||
@@ -995,12 +1285,19 @@ impl HealStorageAPI for MockStorage {
|
|||||||
_include_lifecycle_object_info: bool,
|
_include_lifecycle_object_info: bool,
|
||||||
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
|
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
|
||||||
self.listed_prefixes.lock().unwrap().push(prefix.to_string());
|
self.listed_prefixes.lock().unwrap().push(prefix.to_string());
|
||||||
|
if self.fail_second_listing_page {
|
||||||
|
return if continuation_token.is_none() {
|
||||||
|
Ok((vec![heal_item("object-a")], Some("next-page".to_string()), true))
|
||||||
|
} else {
|
||||||
|
Err(Error::other("listing unavailable"))
|
||||||
|
};
|
||||||
|
}
|
||||||
if *self.truncate_without_token.lock().unwrap() {
|
if *self.truncate_without_token.lock().unwrap() {
|
||||||
return Ok((vec![heal_item("object-a")], None, true));
|
return Ok((vec![heal_item("object-a")], None, true));
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut listed = self.listed.lock().unwrap();
|
let mut listed = self.listed.lock().unwrap();
|
||||||
if continuation_token.is_none() && !*listed {
|
if continuation_token.is_none() && (!*listed || self.list_each_bucket) {
|
||||||
*listed = true;
|
*listed = true;
|
||||||
let objects = if bucket == RUSTFS_META_BUCKET {
|
let objects = if bucket == RUSTFS_META_BUCKET {
|
||||||
vec![
|
vec![
|
||||||
|
|||||||
Reference in New Issue
Block a user