mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 19:55:37 +00:00
Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a20e162d04 | |||
| 0449a9d544 | |||
| 516d0294a6 | |||
| 09947213fe | |||
| 6b8489db92 | |||
| e9c9505212 | |||
| f5baedc8ba | |||
| 71b19bd522 | |||
| cbd3ff9ad7 | |||
| 61210d02d2 | |||
| bdbdca07c8 | |||
| 53efaa2b8f | |||
| ec9672a397 | |||
| cee35f7e54 | |||
| ef7e7afd8c | |||
| 652ebb12c6 | |||
| 38c03d9d5d |
@@ -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 },
|
||||
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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!(
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()],
|
||||
|
||||
Reference in New Issue
Block a user