diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 336ba00ad..7fa8d7717 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -54,6 +54,15 @@ pub enum Error { #[error("Heal task execution failed: {message}")] TaskExecutionFailed { message: String }, + /// The current page already exhausted its local retry budget. Retrying + /// the enclosing bucket would replay pages whose results were counted. + #[error("Heal listing failed for bucket {bucket}: {source}")] + HealListingFailed { + bucket: String, + #[source] + source: Box, + }, + #[error("Invalid heal type: {heal_type}")] InvalidHealType { heal_type: String }, 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..d8b3e3b16 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 { @@ -287,6 +288,59 @@ pub(super) async fn pause_completed_retention_before_publish(task_id: &str, stat } } +#[tokio::test] +async fn canonical_outcome_cancel_wins_before_worker_finalizes_success() { + use crate::heal::outcome::{HealAbortReason, HealExecutionOutcome}; + use crate::heal::task::{OUTCOME_FINISH_TEST_HOOK, OutcomeFinishTestHook}; + let bucket = "canonical-outcome-cancel-before-finish"; + let manager = HealManager::new(Arc::new(MockStorage), None); + let request = HealRequest::object(bucket.to_string(), "object".to_string(), None); + let task_id = request.id.clone(); + let duplicate = HealRequest::object(bucket.to_string(), "object".to_string(), None); + let alias = duplicate.id.clone(); + let retention_hook = Arc::new(CompletedRetentionHook::default()); + { + let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await; + hooks.insert(bucket.to_string(), retention_hook.clone()); + hooks.insert(task_id.clone(), retention_hook.clone()); + } + let finish_hook = Arc::new(OutcomeFinishTestHook { + task_id: task_id.clone(), + reached: Notify::new(), + release: Notify::new(), + }); + *OUTCOME_FINISH_TEST_HOOK.lock().await = Some(finish_hook.clone()); + manager.submit_heal_request(request).await.expect("admit original"); + manager.submit_heal_request(duplicate).await.expect("admit alias"); + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(5), retention_hook.started.notified()) + .await + .expect("storage started"); + retention_hook.execute.notify_one(); + tokio::time::timeout(Duration::from_secs(5), finish_hook.reached.notified()) + .await + .expect("storage returned before outcome finalization"); + manager.cancel_task(&alias).await.expect("cancel wins publication"); + finish_hook.release.notify_one(); + tokio::time::timeout(Duration::from_secs(5), retention_hook.handoff.notified()) + .await + .expect("scheduler completes cancelled handoff"); + 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("frozen outcome").execution, + HealExecutionOutcome::Aborted(HealAbortReason::Cancelled) + ); + } + retention_hook.finish.notify_one(); + *OUTCOME_FINISH_TEST_HOOK.lock().await = None; + COMPLETED_RETENTION_HOOKS + .lock() + .await + .retain(|key, _| key != bucket && key != &task_id); +} + #[tokio::test] async fn completed_retention_cancel_wins_over_a_prepared_retry_snapshot() { let bucket = "completed-retention-retry-cancel"; @@ -325,6 +379,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 +449,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 +457,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 +2036,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 +2761,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 +2799,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 +2831,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..ada74376d --- /dev/null +++ b/crates/heal/src/heal/outcome.rs @@ -0,0 +1,305 @@ +// 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) { + if self.execution != HealExecutionOutcome::Aborted(HealAbortReason::Cancelled) { + 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) { + if self.execution == HealExecutionOutcome::Aborted(HealAbortReason::Cancelled) { + return; + } + 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..a1b09c3b7 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, @@ -43,6 +47,26 @@ use uuid::Uuid; use super::{BUCKET_META_PREFIX, DATA_USAGE_CACHE_NAME, RUSTFS_META_BUCKET}; +#[cfg(test)] +pub(crate) struct OutcomeFinishTestHook { + pub(crate) task_id: String, + pub(crate) reached: tokio::sync::Notify, + pub(crate) release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(crate) static OUTCOME_FINISH_TEST_HOOK: std::sync::LazyLock>>> = + std::sync::LazyLock::new(|| tokio::sync::Mutex::new(None)); + +#[cfg(test)] +async fn pause_outcome_finish(task_id: &str) { + let hook = OUTCOME_FINISH_TEST_HOOK.lock().await.clone(); + if let Some(hook) = hook.filter(|hook| hook.task_id == task_id) { + hook.reached.notify_one(); + hook.release.notified().await; + } +} + const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_TASK: &str = "task"; const LOG_SUBSYSTEM_OBJECT: &str = "object"; @@ -394,6 +418,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 +485,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 +533,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 +720,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 +820,8 @@ impl HealTask { return false; } + self.record_deferred_object(HealDeferredReason::TransientUsageCache).await; + warn!( target: "rustfs::heal::task", event = EVENT_HEAL_OBJECT_RESULT, @@ -755,6 +844,8 @@ impl HealTask { return false; } + self.record_deferred_object(HealDeferredReason::DanglingDeleteGrace).await; + warn!( target: "rustfs::heal::task", event = EVENT_HEAL_OBJECT_RESULT, @@ -801,6 +892,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 +952,45 @@ impl HealTask { HealType::ErasureSet { buckets, set_disk_id } => self.heal_erasure_set(buckets.clone(), set_disk_id.clone()).await, }; + #[cfg(test)] + pause_outcome_finish(&self.id).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 +1075,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..c6225c844 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; @@ -304,23 +306,47 @@ impl HealTask { let mut continuation_token: Option = None; loop { self.check_control_flags().await?; - let (objects, next_token, is_truncated) = if let Some(set_disk_id) = set_disk_id.as_deref() { - self.await_with_control(self.storage.list_versions_for_heal_page_disk_walk( - set_disk_id, - bucket, - prefix, - continuation_token.as_deref(), - false, - )) - .await? - } else { - self.await_with_control(self.storage.list_objects_for_heal_page( - bucket, - prefix, - continuation_token.as_deref(), - false, - )) - .await? + let mut listing_attempt = 0; + let (objects, next_token, is_truncated) = loop { + let page = if let Some(set_disk_id) = set_disk_id.as_deref() { + self.await_with_control(self.storage.list_versions_for_heal_page_disk_walk( + set_disk_id, + bucket, + prefix, + continuation_token.as_deref(), + false, + )) + .await + } else { + self.await_with_control(self.storage.list_objects_for_heal_page( + bucket, + prefix, + continuation_token.as_deref(), + false, + )) + .await + }; + match page { + Ok(page) => break page, + Err(error @ (Error::TaskCancelled | Error::TaskTimeout)) => return Err(error), + Err(error) => { + self.outcome.write().await.attempt_failed(); + if error.is_recoverable_heal() && listing_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES { + listing_attempt += 1; + self.await_with_control(async { + tokio::time::sleep(self.bucket_object_retry_delay(listing_attempt)).await; + Ok(()) + }) + .await?; + continue; + } + self.outcome.write().await.mark_untraversable(); + return Err(Error::HealListingFailed { + bucket: bucket.to_string(), + source: Box::new(error), + }); + } + } }; let mut pending = objects; @@ -338,6 +364,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 +414,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 +453,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 +487,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 +526,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 +554,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. + // Truncated without a continuation token is a compatibility EOF. 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..8cabf1214 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -14,6 +14,364 @@ 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(start_paused = true)] + async fn cluster_retries_only_the_failed_listing_page() { + let storage = Arc::new(MockStorage { + recoverable_second_page_failures: Mutex::new(Some(1)), + ..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("second-page retry succeeds"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.execution, HealExecutionOutcome::Completed); + assert_eq!(outcome.coverage, HealTraversalCoverage::Complete); + assert_eq!(outcome.counters.processed, 2); + assert_eq!(outcome.counters.attempt_failures, 1); + assert_eq!(task.get_progress().await.objects_scanned, 2); + assert_eq!( + storage.heal_object_calls.lock().expect("object calls").as_slice(), + ["object-a", "object-b"] + ); + assert_eq!( + storage.listing_tokens.lock().expect("listing tokens").as_slice(), + [None, Some("second".to_string()), Some("second".to_string())] + ); + } + + #[tokio::test(start_paused = true)] + async fn exhausted_listing_page_cannot_restart_the_bucket() { + let storage = Arc::new(MockStorage { + recoverable_second_page_failures: Mutex::new(Some(4)), + ..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("listing page budget exhausted"); + 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.counters.attempt_failures, 4); + assert_eq!(task.get_progress().await.objects_scanned, 1); + assert_eq!(storage.heal_object_calls.lock().expect("object calls").as_slice(), ["object-a"]); + assert_eq!(storage.bucket_heal_calls.lock().expect("bucket calls").as_slice(), ["bucket-a"]); + } + + #[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 object_b_calls = { + let calls = storage.heal_object_calls.lock().expect("calls"); + calls.iter().filter(|object| object.as_str() == "object-b").count() + }; + assert_eq!(object_b_calls, 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 pool_and_set = { + let options = storage.object_heal_opts.lock().expect("storage options"); + (options[0].pool, options[0].set) + }; + assert_eq!(pool_and_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 +940,10 @@ 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, + recoverable_second_page_failures: Mutex>, + listing_tokens: Mutex>>, healed_objects: Mutex>, heal_object_calls: Mutex>, heal_object_version_ids: Mutex>>, @@ -995,12 +1357,41 @@ impl HealStorageAPI for MockStorage { _include_lifecycle_object_info: bool, ) -> Result<(Vec, Option, bool)> { self.listed_prefixes.lock().unwrap().push(prefix.to_string()); + self.listing_tokens + .lock() + .expect("listing tokens") + .push(continuation_token.map(ToOwned::to_owned)); + if let Some(remaining) = self + .recoverable_second_page_failures + .lock() + .expect("listing failures") + .as_mut() + { + if continuation_token.is_none() { + return Ok((vec![heal_item("object-a")], Some("second".to_string()), true)); + } + if *remaining > 0 { + *remaining -= 1; + return Err(Error::Storage(EcstoreError::InsufficientReadQuorum( + bucket.to_string(), + "page".to_string(), + ))); + } + return Ok((vec![heal_item("object-b")], None, false)); + } + 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![ @@ -1393,6 +1784,8 @@ async fn test_recursive_bucket_heal_skips_object_dir_candidates() { #[tokio::test] async fn test_recursive_bucket_heal_treats_missing_continuation_token_as_end() { + use crate::heal::outcome::{HealExecutionOutcome, HealTraversalCoverage}; + // A version listing can report the final page as truncated with no // continuation token. That is treated as end-of-listing (not an error), // so the returned page is healed and the pass terminates cleanly instead @@ -1414,10 +1807,16 @@ async fn test_recursive_bucket_heal_treats_missing_continuation_token_as_end() { ); let task = HealTask::from_request(request, storage.clone()); - task.heal_bucket("bucket-a") + task.execute() .await .expect("truncated-without-token must terminate cleanly, not loop or error"); + assert_eq!(task.get_status().await, HealTaskStatus::Completed); + let outcome = task.get_outcome().await; + assert_eq!(outcome.execution, HealExecutionOutcome::Completed); + assert_eq!(outcome.coverage, HealTraversalCoverage::Complete); + assert_eq!(outcome.counters.processed, 1); + assert_eq!( storage.healed_objects.lock().unwrap().as_slice(), ["object-a".to_string()],