Compare commits

..

17 Commits

Author SHA1 Message Date
houseme a20e162d04 Merge branch 'main' into houseme/feat/heal-outcomes-delivery 2026-09-06 03:22:26 +08:00
houseme 0449a9d544 fix(heal): drop test locks before awaits
Limit synchronous mock mutex guards to pre-await scopes in canonical outcome tests.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-06 03:03:48 +08:00
houseme 516d0294a6 Merge branch 'main' into houseme/feat/heal-outcomes-delivery 2026-09-06 03:00:57 +08:00
houseme 09947213fe fix(heal): preserve compatible listing EOF outcomes
Keep truncated heal listings without continuation tokens as complete compatibility EOFs and assert the canonical task outcome.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-06 02:08:30 +08:00
houseme 6b8489db92 Merge branch 'main' into houseme/feat/heal-outcomes-delivery 2026-09-06 01:49:51 +08:00
houseme e9c9505212 Merge branch 'main' into houseme/feat/heal-outcomes-delivery 2026-09-06 01:33:08 +08:00
houseme f5baedc8ba chore(heal): integrate frozen main for outcome delivery
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-06 00:59:50 +08:00
houseme 71b19bd522 fix(heal): preserve cancellation and retry only failed listing pages
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 21:18:01 +08:00
houseme cbd3ff9ad7 feat(heal): record bounded canonical object and task outcomes
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 21:18:01 +08:00
houseme 61210d02d2 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-delivery-base 2026-09-05 21:07:31 +08:00
houseme bdbdca07c8 fix(deps): preserve supported hotpath focus expressions
Keep the profiler runtime before its regex-lite compatibility regression.
Track the opt-in validation required to remove this constraint in backlog.

Refs rustfs/backlog#2302.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:50:46 +08:00
houseme 53efaa2b8f chore(deps): refresh profiling dependencies for the next batch
Update hotpath and its macro crate to the compatible patch release before
the next dependency-ready implementation tasks.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:10:54 +08:00
houseme ec9672a397 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b4-base 2026-09-05 18:54:50 +08:00
houseme cee35f7e54 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:43:37 +08:00
houseme ef7e7afd8c Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:35:05 +08:00
houseme 652ebb12c6 fix(ecstore): remove duplicate local rename implementation
Keep the canonical commit module after concurrent storage changes merged.
The control-write and rollback changes are already present there.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:28:01 +08:00
houseme 38c03d9d5d chore(deps): refresh scanner heal batch dependency baseline
Regenerate compatible lockfile selections before the next implementation
batch. Cargo upgrade leaves direct requirements unchanged.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:22:13 +08:00
12 changed files with 1028 additions and 23 deletions
+9
View File
@@ -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>,
},
#[error("Invalid heal type: {heal_type}")]
InvalidHealType { heal_type: String },
+6
View File
@@ -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),
+6
View File
@@ -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<CompletedHealStatus> {
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<Arc<HealTaskOutcome>>,
pub status: HealTaskStatus,
pub result_items: Vec<HealResultItem>,
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;
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<u64>) -> 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<u64>) ->
};
HealTaskReport {
status: completed.status.clone(),
outcome: completed.outcome.clone(),
result_items,
result_items_truncated: completed.result_items_truncated || lagged,
progress: completed.progress.clone(),
+3
View File
@@ -83,6 +83,7 @@ pub(super) struct CompletedHealStatus {
pub(super) heal_type: HealType,
pub(super) status: HealTaskStatus,
pub(super) progress: Option<HealProgress>,
pub(super) outcome: Option<Arc<HealTaskOutcome>>,
pub(super) retained_bytes: std::sync::OnceLock<usize>,
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::<Self>();
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(),
@@ -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);
+64
View File
@@ -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(),
+1
View File
@@ -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;
+305
View File
@@ -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<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) {
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<HealAbortReason>) {
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::<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);
}
}
+132
View File
@@ -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<tokio::sync::Mutex<Option<Arc<OutcomeFinishTestHook>>>> =
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<RwLock<HealTaskStatus>>,
/// Progress tracking
pub progress: Arc<RwLock<HealProgress>>,
outcome: Arc<RwLock<HealTaskOutcome>>,
/// 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<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 {
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!(
+98 -19
View File
@@ -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<String> = 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;
}
}
+2 -2
View File
@@ -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
+401 -2
View File
@@ -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<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(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<bool>,
list_each_bucket: bool,
fail_second_listing_page: bool,
recoverable_second_page_failures: Mutex<Option<usize>>,
listing_tokens: Mutex<Vec<Option<String>>>,
healed_objects: Mutex<Vec<String>>,
heal_object_calls: Mutex<Vec<String>>,
heal_object_version_ids: Mutex<Vec<Option<String>>>,
@@ -995,12 +1357,41 @@ impl HealStorageAPI for MockStorage {
_include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, 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()],