diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index 33ac33c38..580b01246 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -447,6 +447,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "running".to_string(), None, @@ -463,6 +464,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "running".to_string(), Some(format!("heal task retrying after recoverable failure, attempt {retry_attempt}: {error}")), @@ -479,6 +481,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "finished".to_string(), None, @@ -495,6 +498,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "stopped".to_string(), Some("heal task cancelled".to_string()), @@ -511,6 +515,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "stopped".to_string(), Some("heal task timed out".to_string()), @@ -527,6 +532,7 @@ impl HealChannelProcessor { progress, next_seq, min_seq, + .. }) => ( "stopped".to_string(), Some(error), diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 971590fdb..c90a18729 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -13,6 +13,7 @@ // limitations under the License. use crate::heal::{ + outcome::HealTaskOutcome, progress::{HealProgress, HealStatistics}, resume::{ReplacementPhase, ResumeGc, ResumeManager, ResumeState, ResumeUtils}, storage::HealStorageAPI, @@ -185,6 +186,7 @@ fn record_displaced_terminal( request: &HealRequest, ) -> Arc { let terminal = Arc::new(CompletedHealStatus { + outcome: None, progress: None, retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type.clone(), @@ -268,6 +270,7 @@ async fn publish_completed_heal( #[derive(Debug, Clone)] pub struct HealTaskReport { + pub outcome: Option>, pub status: HealTaskStatus, pub result_items: Vec, pub result_items_truncated: bool, @@ -285,6 +288,7 @@ async fn active_task_report(task: &HealTask, since: Option) -> HealTaskRepo let window = task.get_result_items_since(since).await; HealTaskReport { status: task.get_status().await, + outcome: Some(Arc::new(task.get_outcome().await)), result_items: window.items, // The legacy flag stays set once anything was evicted; a lagging // incremental cursor additionally marks this response truncated so @@ -298,6 +302,7 @@ async fn active_task_report(task: &HealTask, since: Option) -> HealTaskRepo fn empty_task_report(status: HealTaskStatus) -> HealTaskReport { HealTaskReport { + outcome: None, status, result_items: Vec::new(), result_items_truncated: false, @@ -325,6 +330,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option) -> }; HealTaskReport { status: completed.status.clone(), + outcome: completed.outcome.clone(), result_items, result_items_truncated: completed.result_items_truncated || lagged, progress: completed.progress.clone(), diff --git a/crates/heal/src/heal/manager/queue.rs b/crates/heal/src/heal/manager/queue.rs index aceabe42a..a47cd43e2 100644 --- a/crates/heal/src/heal/manager/queue.rs +++ b/crates/heal/src/heal/manager/queue.rs @@ -83,6 +83,7 @@ pub(super) struct CompletedHealStatus { pub(super) heal_type: HealType, pub(super) status: HealTaskStatus, pub(super) progress: Option, + pub(super) outcome: Option>, pub(super) retained_bytes: std::sync::OnceLock, pub(super) result_items_truncated: bool, pub(super) completed_at: SystemTime, @@ -105,6 +106,7 @@ impl CompletedHealStatus { fn measure_retained_bytes(&self) -> usize { let mut bytes = size_of::(); 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 { HealType::Cluster => {} HealType::Bucket { bucket } => add(bucket.capacity()), @@ -209,6 +211,7 @@ impl CompletedHealStatus { heal_type: task.heal_type.clone(), status, progress: Some(task.get_progress().await), + outcome: Some(Arc::new(task.get_outcome().await)), retained_bytes: std::sync::OnceLock::new(), result_items_truncated: task.result_items_truncated(), completed_at: SystemTime::now(), diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index feb0c1231..704a7cbeb 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -298,6 +298,7 @@ impl HealManager { if cancelled_completion { completed_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 successful_completion = matches!(completed_status, HealTaskStatus::Completed); diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index aa1c1293f..cb5a52c1b 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -103,6 +103,7 @@ struct MockStorage; fn completed_retention_fixture(completed_at: SystemTime) -> CompletedHealStatus { CompletedHealStatus { + outcome: None, heal_type: HealType::Cluster, status: HealTaskStatus::Completed, progress: Some(HealProgress { @@ -325,6 +326,10 @@ async fn completed_retention_cancel_wins_over_a_prepared_retry_snapshot() { for token in [&task_id, &alias] { let report = manager.get_task_report(token).await.expect("cancelled token retained"); 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!(!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"); assert!(!manager.active_heals.lock().await.contains_key(&task_id)); let expected = task.get_progress().await; + let expected_outcome = task.get_outcome().await; for token in [&task_id, &alias] { assert_eq!(manager.get_task_progress(token).await.expect("terminal progress query"), expected); let report = manager @@ -398,6 +404,7 @@ async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_han .await .expect("terminal token remains queryable at handoff"); assert_eq!(report.progress.as_ref(), Some(&expected)); + assert_eq!(report.outcome.as_deref(), Some(&expected_outcome)); assert!(report.result_items.is_empty()); match outcome { "success" => assert_eq!(report.status, HealTaskStatus::Completed), @@ -1976,6 +1983,7 @@ async fn insert_retrying_request(manager: &HealManager, request: HealRequest) -> task_id, Arc::new(CompletedHealStatus { progress: None, + outcome: None, retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type, status: HealTaskStatus::Retrying { @@ -2700,6 +2708,7 @@ async fn test_retrying_completion_outranks_the_queue_for_the_same_id() { task_id.clone(), Arc::new(CompletedHealStatus { progress: None, + outcome: None, retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type.clone(), status: HealTaskStatus::Retrying { @@ -2737,6 +2746,7 @@ async fn test_get_task_status_reads_recent_completed_status() { "completed-token".to_string(), Arc::new(CompletedHealStatus { progress: None, + outcome: None, retained_bytes: std::sync::OnceLock::new(), heal_type: HealType::Bucket { bucket: "bucket".to_string(), @@ -2768,6 +2778,7 @@ async fn test_get_task_report_for_path_reads_completed_items() { "completed-token".to_string(), Arc::new(CompletedHealStatus { progress: None, + outcome: None, retained_bytes: std::sync::OnceLock::new(), heal_type: HealType::Object { bucket: "bucket".to_string(), diff --git a/crates/heal/src/heal/mod.rs b/crates/heal/src/heal/mod.rs index c918bea49..5f17c8cd8 100644 --- a/crates/heal/src/heal/mod.rs +++ b/crates/heal/src/heal/mod.rs @@ -16,6 +16,7 @@ pub mod channel; pub mod erasure_healer; pub mod manager; pub mod mrf_queue; +pub mod outcome; pub mod progress; pub(crate) mod replacement_readiness; pub mod resume; diff --git a/crates/heal/src/heal/outcome.rs b/crates/heal/src/heal/outcome.rs new file mode 100644 index 000000000..36eab0225 --- /dev/null +++ b/crates/heal/src/heal/outcome.rs @@ -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, + pub bucket_incarnation_id: Option, + pub pool_index: Option, + pub set_index: Option, +} + +#[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, + }, + Failed(HealFailureClass), + Cancelled, + DryRunObserved, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct HealObjectOutcome { + pub identity: HealObjectIdentity, + pub disposition: HealObjectDisposition, + pub detail: Option, +} + +impl HealObjectOutcome { + fn retained_bytes(&self) -> usize { + size_of::() + .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, + 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) { + 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::() + .saturating_add(self.retained_object_bytes) + .saturating_add(self.objects.capacity().saturating_mul(size_of::())) + } +} + +#[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); + } +} diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index c4c1d902b..1d72bfb25 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -15,6 +15,10 @@ use crate::heal::{ DiskError, EcstoreError, ErasureSetHealer, HealDiskExt as _, erasure_healer::target_outcomes_complete, + outcome::{ + HealAbortReason, HealDeferredReason, HealFailureClass, HealObjectDisposition, HealObjectIdentity, HealObjectKind, + HealObjectOutcome, HealTaskOutcome, + }, progress::HealProgress, resume::{ CheckpointManager, ReplacementPhase, ReplacementTargetIdentity, ResumeManager, replacement_target_identities_match, @@ -394,6 +398,7 @@ pub struct HealTask { pub status: Arc>, /// Progress tracking pub progress: Arc>, + outcome: Arc>, /// Result items collected from storage heal calls, each stamped with a /// monotonically increasing sequence number for incremental consumption /// (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)), batch_failure: Arc::new(RwLock::new(None)), batch_failure_recorded: Arc::new(AtomicBool::new(false)), + outcome: Arc::new(RwLock::new(HealTaskOutcome::default())), created_at: request.created_at, enqueued_at: request.enqueued_at, started_at: Arc::new(RwLock::new(None)), @@ -507,6 +513,66 @@ impl HealTask { 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, + set_index: Option, + ) -> 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 { + 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 { 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<()> { + self.record_deferred_object(HealDeferredReason::TransientExistenceCheck).await; warn!( target: "rustfs::heal::task", event = EVENT_HEAL_OBJECT_RESULT, @@ -733,6 +800,8 @@ impl HealTask { return false; } + self.record_deferred_object(HealDeferredReason::TransientUsageCache).await; + warn!( target: "rustfs::heal::task", event = EVENT_HEAL_OBJECT_RESULT, @@ -755,6 +824,8 @@ impl HealTask { return false; } + self.record_deferred_object(HealDeferredReason::DanglingDeleteGrace).await; + warn!( target: "rustfs::heal::task", 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))] #[hotpath::measure] pub async fn execute(&self) -> Result<()> { + self.outcome.write().await.start(); // update status and timestamps atomically to avoid race conditions let now = SystemTime::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, }; + { + 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 { let mut completed_at = self.completed_at.write().await; @@ -944,6 +1053,7 @@ impl HealTask { pub async fn cancel(&self) -> Result<()> { self.cancel_token.cancel(); + self.outcome.write().await.finish(Some(HealAbortReason::Cancelled)); let mut status = self.status.write().await; *status = HealTaskStatus::Cancelled; debug!( diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index d44df549e..241aaba5a 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -214,6 +214,7 @@ impl HealTask { continue; } failed = failed.saturating_add(1); + self.outcome.write().await.mark_untraversable(); if err.is_recoverable_heal() { retryable = retryable.saturating_add(1); } else { @@ -260,6 +261,7 @@ impl HealTask { #[hotpath::measure] async fn heal_bucket_objects(&self, bucket: &str, prefix: &str) -> Result<()> { + let previous_progress = self.get_progress().await; let mut scanned = 0u64; let mut healed = 0u64; let mut failed = 0u64; @@ -338,6 +340,14 @@ impl HealTask { self.check_control_flags().await?; let mut telemetry_unknown = false; 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; progress.set_current_object(Some(format!("{bucket}/{object}"))); @@ -380,7 +390,31 @@ impl HealTask { }; 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) { + disposition = HealObjectDisposition::Deferred { + reason: HealDeferredReason::DanglingDeleteGrace, + retry_not_before: None, + }; telemetry_unknown |= !increment_counter(&mut skipped); warn!( target: "rustfs::heal::task", @@ -395,6 +429,10 @@ impl HealTask { "Heal bucket object dangling cleanup deferred by grace window" ); } 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); warn!( target: "rustfs::heal::task", @@ -425,6 +463,11 @@ impl HealTask { ); retry.push(item); } else { + disposition = HealObjectDisposition::Failed(if err.is_recoverable_heal() { + HealFailureClass::RetryExhausted + } else { + HealFailureClass::Permanent + }); telemetry_unknown |= !increment_counter(&mut failed); if err.is_recoverable_heal() { retryable_failed = retryable_failed.saturating_add(1); @@ -459,8 +502,20 @@ impl HealTask { continue; } + self.outcome.write().await.record(HealObjectOutcome { + identity, + disposition, + detail, + }); + 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 { progress.mark_unknown(); } @@ -475,7 +530,7 @@ impl HealTask { continuation_token = next_heal_listing_token(bucket, prefix, next_token, is_truncated)?; if continuation_token.is_none() { - // Truncated but no continuation token: end of listing. + self.outcome.write().await.mark_untraversable(); break; } } diff --git a/crates/heal/src/heal/task/heal_metadata.rs b/crates/heal/src/heal/task/heal_metadata.rs index fe2dc98ac..fad5a403e 100644 --- a/crates/heal/src/heal/task/heal_metadata.rs +++ b/crates/heal/src/heal/task/heal_metadata.rs @@ -261,8 +261,8 @@ impl HealTask { update_parity: true, no_lock: self.options.no_lock, read_repair: false, - pool: None, - set: None, + pool: self.options.pool_index, + set: self.options.set_index, }; let heal_result = self diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 1734d0c7a..9693f0ef0 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -14,6 +14,294 @@ use super::super::{DiskOption, DiskStore, Endpoint, new_disk}; use super::*; + +mod canonical_outcome { + use super::*; + use crate::heal::outcome::{HealExecutionOutcome, HealTraversalCoverage}; + + fn bucket_task(storage: Arc) -> 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 rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, TraceSubscription, TraceVal, subscribe_trace_events}; use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos}; @@ -582,6 +870,8 @@ async fn verified_recovery_keeps_state_when_marker_clear_fails() { #[derive(Default)] struct MockStorage { listed: Mutex, + list_each_bucket: bool, + fail_second_listing_page: bool, healed_objects: Mutex>, heal_object_calls: Mutex>, heal_object_version_ids: Mutex>>, @@ -995,12 +1285,19 @@ impl HealStorageAPI for MockStorage { _include_lifecycle_object_info: bool, ) -> Result<(Vec, Option, bool)> { 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() { return Ok((vec![heal_item("object-a")], None, true)); } 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; let objects = if bucket == RUSTFS_META_BUCKET { vec![