Files
rustfs/crates/heal/src/heal/task.rs
T

5026 lines
194 KiB
Rust

// Copyright 2024 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.
use crate::heal::{
DiskError, EcstoreError, ErasureSetHealer, HealDiskExt as _,
erasure_healer::target_outcomes_complete,
progress::HealProgress,
resume::{
CheckpointManager, ReplacementPhase, ReplacementTargetIdentity, ResumeManager, replacement_target_identities_match,
},
storage::{HealBucketUsageBaseline, HealStorageAPI, next_heal_listing_token},
};
use crate::{Error, Result};
use metrics::{counter, histogram};
use rustfs_common::heal_channel::{HealOpts, HealRequestSource, HealScanMode};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit};
use rustfs_madmin::heal_commands::HealResultItem;
use rustfs_utils::path::SLASH_SEPARATOR;
use serde::{Deserialize, Serialize};
use std::{
future::Future,
sync::{
Arc,
atomic::{AtomicBool, AtomicU64, Ordering},
},
time::{Duration, Instant, SystemTime},
};
use tokio::sync::RwLock;
use tracing::{debug, error, info, warn};
use uuid::Uuid;
use super::{BUCKET_META_PREFIX, DATA_USAGE_CACHE_NAME, RUSTFS_META_BUCKET};
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_TASK: &str = "task";
const LOG_SUBSYSTEM_OBJECT: &str = "object";
const EVENT_HEAL_TASK_STATE: &str = "heal_task_state";
const EVENT_HEAL_OBJECT_STAGE: &str = "heal_object_stage";
const EVENT_HEAL_OBJECT_MISSING: &str = "heal_object_missing";
const MAX_RETAINED_HEAL_RESULT_ITEMS: usize = 1024;
const EVENT_HEAL_OBJECT_RESULT: &str = "heal_object_result";
const MAX_BUCKET_OBJECT_HEAL_RETRIES: u32 = 3;
const MAX_BUCKET_FAILURE_LOG_SAMPLES: u64 = 5;
/// Emits at `$level`, demoted to `debug!` when `$demote` is true. Keeps
/// per-object heal work — Object/Metadata/MRF/ECDecode tasks queued per
/// object by MRF/autoheal/scanner loops, and per-object sweep failures past
/// a sample cap — from amplifying into one info!/warn!/error! line per
/// object during mass recovery (rustfs/rustfs#5716). Aggregate task kinds
/// and foreground (admin/internal) requests keep operator-visible levels;
/// metrics and end-of-sweep summaries carry the aggregate signal for the
/// demoted paths.
macro_rules! demote_to_debug_when {
($demote:expr, $level:ident, target: $target:expr, { $($fields:tt)* }) => {
if $demote {
tracing::debug!(target: $target, $($fields)*);
} else {
tracing::$level!(target: $target, $($fields)*);
}
};
}
pub(crate) use demote_to_debug_when;
const EVENT_HEAL_BUCKET_STAGE: &str = "heal_bucket_stage";
const EVENT_HEAL_BUCKET_RESULT: &str = "heal_bucket_result";
const EVENT_HEAL_METADATA_STAGE: &str = "heal_metadata_stage";
const EVENT_HEAL_METADATA_RESULT: &str = "heal_metadata_result";
const EVENT_HEAL_MRF_STAGE: &str = "heal_mrf_stage";
const EVENT_HEAL_MRF_RESULT: &str = "heal_mrf_result";
const EVENT_HEAL_EC_DECODE_STAGE: &str = "heal_ec_decode_stage";
const EVENT_HEAL_EC_DECODE_RESULT: &str = "heal_ec_decode_result";
const EVENT_HEAL_ERASURE_SET_STAGE: &str = "heal_erasure_set_stage";
const EVENT_HEAL_ERASURE_SET_RESULT: &str = "heal_erasure_set_result";
/// Heal type
#[derive(Debug, Clone)]
pub enum HealType {
/// Cluster heal
Cluster,
/// Object heal
Object {
bucket: String,
object: String,
version_id: Option<String>,
},
/// Bucket heal
Bucket { bucket: String },
/// Prefix heal
Prefix { bucket: String, prefix: String },
/// Erasure Set heal (includes disk format repair)
ErasureSet { buckets: Vec<String>, set_disk_id: String },
/// Metadata heal
Metadata { bucket: String, object: String },
/// MRF heal
MRF { meta_path: String },
/// EC decode heal
ECDecode {
bucket: String,
object: String,
version_id: Option<String>,
},
}
impl HealType {
fn log_kind(&self) -> &'static str {
match self {
Self::Cluster => "cluster",
Self::Object { .. } => "object",
Self::Bucket { .. } => "bucket",
Self::Prefix { .. } => "prefix",
Self::ErasureSet { .. } => "erasure_set",
Self::Metadata { .. } => "metadata",
Self::MRF { .. } => "mrf",
Self::ECDecode { .. } => "ec_decode",
}
}
/// Task kinds enqueued at per-object granularity (MRF, autoheal, scanner,
/// read-repair loops). Their lifecycle and admission logs stay at `debug!`
/// so a recovery loop queuing hundreds of thousands of object heal tasks
/// cannot amplify into per-object `info!`/`warn!` lines; aggregate kinds
/// (cluster/bucket/prefix/erasure-set) keep operator-visible levels.
pub(crate) fn is_per_object(&self) -> bool {
matches!(
self,
Self::Object { .. } | Self::Metadata { .. } | Self::MRF { .. } | Self::ECDecode { .. }
)
}
}
fn is_object_level_not_found_error(err: &Error) -> bool {
match err {
Error::Disk(DiskError::FileNotFound | DiskError::FileVersionNotFound) => true,
Error::Storage(EcstoreError::FileNotFound | EcstoreError::FileVersionNotFound) => true,
Error::Other(message) => matches!(message.as_str(), "File not found" | "File version not found"),
_ => false,
}
}
pub(crate) fn is_missing_object_dir_heal_result(object: &str, err: &Error) -> bool {
object.ends_with(SLASH_SEPARATOR) && is_object_level_not_found_error(err)
}
/// Sample cap for per-object failure logs during a sweep: returns true (and
/// consumes a sample slot) for the first [`MAX_BUCKET_FAILURE_LOG_SAMPLES`]
/// calls, false afterwards so callers demote the remaining occurrences to
/// `debug!`. Aggregate failed/skipped counts still surface in end-of-sweep
/// summaries.
pub(crate) fn take_failure_log_sample(samples_logged: &mut u64) -> bool {
if *samples_logged < MAX_BUCKET_FAILURE_LOG_SAMPLES {
*samples_logged = samples_logged.saturating_add(1);
true
} else {
false
}
}
/// Heal priority
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub enum HealPriority {
/// Low priority
Low = 0,
/// Normal priority
#[default]
Normal = 1,
/// High priority
High = 2,
/// Urgent priority
Urgent = 3,
}
impl HealPriority {
fn as_str(self) -> &'static str {
match self {
Self::Low => "low",
Self::Normal => "normal",
Self::High => "high",
Self::Urgent => "urgent",
}
}
}
/// Heal options
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealOptions {
/// Scan mode
pub scan_mode: HealScanMode,
/// Whether to remove corrupted data
pub remove_corrupted: bool,
/// Whether to recreate
pub recreate_missing: bool,
/// Whether to update parity
pub update_parity: bool,
/// Whether to recursively process
pub recursive: bool,
/// Whether to dry run
pub dry_run: bool,
/// Whether to skip namespace locking
#[serde(default)]
pub no_lock: bool,
/// Aggregate execution timeout across recoverable manager retries
pub timeout: Option<Duration>,
/// pool index
pub pool_index: Option<usize>,
/// set index
pub set_index: Option<usize>,
}
impl Default for HealOptions {
fn default() -> Self {
Self {
scan_mode: HealScanMode::Normal,
remove_corrupted: false,
recreate_missing: true,
update_parity: true,
recursive: false,
dry_run: false,
no_lock: false,
timeout: None,
pool_index: None,
set_index: None,
}
}
}
/// Heal task status
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub enum HealTaskStatus {
/// Pending
Pending,
/// Running
Running,
/// Retrying after a recoverable failure
Retrying { error: String, retry_attempt: u32 },
/// Completed
Completed,
/// Failed
Failed { error: String },
/// Cancelled
Cancelled,
/// Timeout
Timeout,
}
#[derive(Debug)]
pub(crate) struct BatchHealFailure {
pub(crate) scope: String,
pub(crate) failed: u64,
pub(crate) retryable: u64,
pub(crate) permanent: u64,
pub(crate) first_object: String,
pub(crate) first_error: String,
}
impl std::fmt::Display for BatchHealFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"Heal batch failed for {}: {} failed ({} retryable, {} permanent); first failure at {}: {}",
self.scope, self.failed, self.retryable, self.permanent, self.first_object, self.first_error
)
}
}
/// Heal request
#[derive(Debug, Clone)]
pub struct HealRequest {
/// Request ID
pub id: String,
/// Heal type
pub heal_type: HealType,
/// Heal options
pub options: HealOptions,
/// Priority
pub priority: HealPriority,
/// Origin of the request for operational status.
pub source: HealRequestSource,
/// Whether this request should bypass queue admission dedup/full policies.
pub force_start: bool,
/// Number of recoverable retry attempts already scheduled for this request.
pub retry_attempts: u32,
/// Endpoints of the disks being rebuilt by an erasure-set heal. Used to
/// write per-disk healing markers so `DiskInfo.healing` reflects reality;
/// empty when the trigger doesn't know the specific disks (admin API,
/// unclean-shutdown verification).
pub heal_endpoints: Vec<String>,
/// Created time
pub created_at: SystemTime,
/// Queue admission time used for scheduler delay metrics
pub enqueued_at: SystemTime,
}
impl HealRequest {
pub fn new(heal_type: HealType, options: HealOptions, priority: HealPriority) -> Self {
let now = SystemTime::now();
Self {
id: Uuid::new_v4().to_string(),
heal_type,
options,
priority,
source: HealRequestSource::Internal,
force_start: false,
retry_attempts: 0,
heal_endpoints: Vec::new(),
created_at: now,
enqueued_at: now,
}
}
pub fn object(bucket: String, object: String, version_id: Option<String>) -> Self {
Self::new(
HealType::Object {
bucket,
object,
version_id,
},
HealOptions::default(),
HealPriority::Normal,
)
}
pub fn bucket(bucket: String) -> Self {
Self::new(HealType::Bucket { bucket }, HealOptions::default(), HealPriority::Normal)
}
pub fn metadata(bucket: String, object: String) -> Self {
Self::new(HealType::Metadata { bucket, object }, HealOptions::default(), HealPriority::High)
}
pub fn ec_decode(bucket: String, object: String, version_id: Option<String>) -> Self {
Self::new(
HealType::ECDecode {
bucket,
object,
version_id,
},
HealOptions::default(),
HealPriority::Urgent,
)
}
}
/// Heal task
/// Incremental view over a task's retained result items (HS-06).
///
/// `next_seq` is the cursor a client should pass on its next poll; `min_seq`
/// is the oldest sequence still retained; `lagged` means the client's cursor
/// fell behind `min_seq` and items were skipped — the client should restart
/// from `min_seq`.
#[derive(Debug, Clone)]
pub struct HealResultWindow {
pub items: Vec<HealResultItem>,
pub next_seq: u64,
pub min_seq: u64,
pub lagged: bool,
}
pub struct HealTask {
/// Task ID
pub id: String,
/// Heal type
pub heal_type: HealType,
/// Heal options
pub options: HealOptions,
/// Priority inherited from the request
pub priority: HealPriority,
/// Origin inherited from the request
pub source: HealRequestSource,
/// Number of recoverable retry attempts already scheduled for this task.
pub retry_attempts: u32,
/// Endpoints of the disks being rebuilt (see `HealRequest::heal_endpoints`).
pub heal_endpoints: Vec<String>,
/// Durable resume anchor injected by the manager for an existing automatic
/// replacement generation.
replacement_resume_endpoint: Option<String>,
/// Task status
pub status: Arc<RwLock<HealTaskStatus>>,
/// Progress tracking
pub progress: Arc<RwLock<HealProgress>>,
/// 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
/// items; see `get_result_items_since`).
pub result_items: Arc<RwLock<Vec<(u64, HealResultItem)>>>,
/// Next sequence number to assign; starts at 1.
next_item_seq: Arc<AtomicU64>,
/// Sequence number of the oldest item still inside the retention window;
/// equals `next_item_seq` while the window is empty.
min_available_seq: Arc<AtomicU64>,
result_items_truncated: Arc<AtomicBool>,
batch_failure: Arc<RwLock<Option<BatchHealFailure>>>,
batch_failure_recorded: Arc<AtomicBool>,
/// Created time
pub created_at: SystemTime,
/// Queue admission time
pub enqueued_at: SystemTime,
/// Started time
pub started_at: Arc<RwLock<Option<SystemTime>>>,
/// Completed time
pub completed_at: Arc<RwLock<Option<SystemTime>>>,
/// Task start instant for timeout calculation (monotonic)
task_start_instant: Arc<RwLock<Option<Instant>>>,
/// Cancel token
pub cancel_token: tokio_util::sync::CancellationToken,
/// Storage layer interface
pub storage: Arc<dyn HealStorageAPI>,
}
impl HealTask {
async fn verify_replacement_identity_fence(
&self,
expected_identities: &[ReplacementTargetIdentity],
set_disk_id: &str,
stage: &str,
) -> Result<()> {
let actual_identities = self
.await_with_control(self.storage.replacement_target_identities(&self.heal_endpoints))
.await?;
if replacement_target_identities_match(expected_identities, &actual_identities) {
return Ok(());
}
Err(Error::TaskExecutionFailed {
message: format!("Replacement target changed during {stage} for automatic heal {set_disk_id}"),
})
}
pub fn from_request(request: HealRequest, storage: Arc<dyn HealStorageAPI>) -> Self {
Self {
id: request.id,
heal_type: request.heal_type,
options: request.options,
priority: request.priority,
source: request.source,
retry_attempts: request.retry_attempts,
heal_endpoints: request.heal_endpoints,
replacement_resume_endpoint: None,
status: Arc::new(RwLock::new(HealTaskStatus::Pending)),
progress: Arc::new(RwLock::new(HealProgress::new())),
result_items: Arc::new(RwLock::new(Vec::new())),
next_item_seq: Arc::new(AtomicU64::new(1)),
min_available_seq: Arc::new(AtomicU64::new(1)),
result_items_truncated: Arc::new(AtomicBool::new(false)),
batch_failure: Arc::new(RwLock::new(None)),
batch_failure_recorded: Arc::new(AtomicBool::new(false)),
created_at: request.created_at,
enqueued_at: request.enqueued_at,
started_at: Arc::new(RwLock::new(None)),
completed_at: Arc::new(RwLock::new(None)),
task_start_instant: Arc::new(RwLock::new(None)),
cancel_token: tokio_util::sync::CancellationToken::new(),
storage,
}
}
pub fn retry_request(&self) -> HealRequest {
HealRequest {
id: self.id.clone(),
heal_type: self.heal_type.clone(),
options: self.options.clone(),
priority: self.priority,
source: self.source,
force_start: false,
retry_attempts: self.retry_attempts.saturating_add(1),
heal_endpoints: self.heal_endpoints.clone(),
created_at: self.created_at,
enqueued_at: SystemTime::now(),
}
}
pub(crate) async fn retry_request_with_remaining_timeout(&self) -> Result<HealRequest> {
let mut request = self.retry_request();
if self.options.timeout.is_some() {
request.options.timeout = self.remaining_timeout().await?;
}
Ok(request)
}
pub(crate) fn from_replacement_recovery_request(
request: HealRequest,
storage: Arc<dyn HealStorageAPI>,
replacement_resume_endpoint: Option<String>,
) -> Self {
let mut task = Self::from_request(request, storage);
task.replacement_resume_endpoint = replacement_resume_endpoint;
task
}
pub fn metric_type_label(&self) -> &'static str {
match &self.heal_type {
HealType::Cluster => "cluster",
HealType::Object { .. } => "object",
HealType::Bucket { .. } => "bucket",
HealType::Prefix { .. } => "prefix",
HealType::ErasureSet { .. } => "erasure_set",
HealType::Metadata { .. } => "metadata",
HealType::MRF { .. } => "mrf",
HealType::ECDecode { .. } => "ec_decode",
}
}
pub(crate) fn has_batch_failure(&self) -> bool {
self.batch_failure_recorded.load(Ordering::Acquire)
}
pub(crate) async fn record_batch_failure(&self, failure: BatchHealFailure) -> Error {
self.batch_failure_recorded.store(true, Ordering::Release);
let message = failure.to_string();
*self.batch_failure.write().await = Some(failure);
Error::TaskExecutionFailed { message }
}
async fn take_batch_failure(&self) -> Option<BatchHealFailure> {
self.batch_failure.write().await.take()
}
pub fn metric_set_label(&self) -> String {
match &self.heal_type {
HealType::ErasureSet { set_disk_id, .. } => set_disk_id.clone(),
_ => match (self.options.pool_index, self.options.set_index) {
(Some(pool), Some(set)) => format!("pool_{pool}_set_{set}"),
_ => "global".to_string(),
},
}
}
fn emit_trace_task_state(&self, state: &'static str, duration: Duration, error: Option<&Error>) {
trace_emit(|| {
let mut event = TraceEvent::new(TraceKind::Heal, TraceFunc::HealTask)
.with_duration(duration)
.with_attr("task_id", self.id.as_str())
.with_attr("heal_type", self.heal_type.log_kind())
.with_attr("state", state)
.with_attr("source", self.source.as_str())
.with_attr("priority", self.priority.as_str())
.with_attr("retry_attempts", u64::from(self.retry_attempts))
.with_attr("dry_run", self.options.dry_run);
event = match &self.heal_type {
HealType::Cluster => event,
HealType::Object {
bucket,
object,
version_id,
} => {
let event = event.with_bucket(bucket.as_str()).with_object(object.as_str());
match version_id {
Some(version_id) => event.with_attr("version_id", version_id.as_str()),
None => event,
}
}
HealType::Bucket { bucket } => event.with_bucket(bucket.as_str()),
HealType::Prefix { bucket, prefix } => event.with_bucket(bucket.as_str()).with_object(prefix.as_str()),
HealType::ErasureSet { buckets, set_disk_id } => {
let bucket_count = u64::try_from(buckets.len()).unwrap_or(u64::MAX);
event
.with_attr("set_disk_id", set_disk_id.as_str())
.with_attr("bucket_count", bucket_count)
}
HealType::Metadata { bucket, object } => event.with_bucket(bucket.as_str()).with_object(object.as_str()),
HealType::ECDecode {
bucket,
object,
version_id,
} => {
let event = event.with_bucket(bucket.as_str()).with_object(object.as_str());
match version_id {
Some(version_id) => event.with_attr("version_id", version_id.as_str()),
None => event,
}
}
HealType::MRF { meta_path } => event.with_object(meta_path.as_str()),
};
match error {
Some(error) => event.with_attr("error", error.to_string()),
None => event,
}
});
}
async fn remaining_timeout(&self) -> Result<Option<Duration>> {
if let Some(total) = self.options.timeout {
let start_instant = { *self.task_start_instant.read().await };
if let Some(started_at) = start_instant {
let elapsed = started_at.elapsed();
if elapsed >= total {
return Err(Error::TaskTimeout);
}
return Ok(Some(total - elapsed));
}
Ok(Some(total))
} else {
Ok(None)
}
}
async fn check_control_flags(&self) -> Result<()> {
if self.cancel_token.is_cancelled() {
return Err(Error::TaskCancelled);
}
// Only interested in propagating an error if the timeout has expired;
// the actual Duration value is not needed here
let _ = self.remaining_timeout().await?;
Ok(())
}
async fn await_with_control<F, T>(&self, fut: F) -> Result<T>
where
F: Future<Output = Result<T>> + Send,
T: Send,
{
let cancel_token = self.cancel_token.clone();
if let Some(remaining) = self.remaining_timeout().await? {
if remaining.is_zero() {
return Err(Error::TaskTimeout);
}
let mut fut = Box::pin(fut);
tokio::select! {
_ = cancel_token.cancelled() => Err(Error::TaskCancelled),
_ = tokio::time::sleep(remaining) => Err(Error::TaskTimeout),
result = &mut fut => result,
}
} else {
tokio::select! {
_ = cancel_token.cancelled() => Err(Error::TaskCancelled),
result = fut => result,
}
}
}
async fn skip_due_to_transient_object_exists(&self, bucket: &str, object: &str, err: &Error) -> Result<()> {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "transient_skip",
error = %err,
"Heal object skipped due to transient existence check error"
);
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("skipped: {bucket}/{object}")));
progress.update_progress(0, 1, 0, 0);
Ok(())
}
fn is_data_usage_cache_object(bucket: &str, object: &str) -> bool {
bucket == RUSTFS_META_BUCKET
&& object
.strip_prefix(BUCKET_META_PREFIX)
.and_then(|suffix| suffix.strip_prefix('/'))
.is_some_and(|name| name.contains(DATA_USAGE_CACHE_NAME))
}
fn is_transient_lock_or_timeout_error(err: &Error) -> bool {
let message = err.to_string().to_ascii_lowercase();
message.contains("lock acquisition timeout")
|| message.contains("lock acquisition failed")
|| message.contains("timed out")
|| message.contains("deadline has elapsed")
}
fn should_skip_data_usage_cache_heal_error(bucket: &str, object: &str, err: &Error) -> bool {
Self::is_data_usage_cache_object(bucket, object) && Self::is_transient_lock_or_timeout_error(err)
}
fn is_no_heal_required_error(err: &Error) -> bool {
match err {
Error::Storage(EcstoreError::NoHealRequired) | Error::Disk(DiskError::NoHealRequired) => true,
Error::Other(message) => matches!(message.as_str(), "No heal required" | "No healing is required"),
_ => matches!(err.to_string().as_str(), "No heal required" | "No healing is required"),
}
}
fn is_object_not_found_heal_error(err: &Error) -> bool {
match err {
Error::Disk(DiskError::FileNotFound | DiskError::FileVersionNotFound) => true,
Error::Storage(
EcstoreError::FileNotFound
| EcstoreError::FileVersionNotFound
| EcstoreError::ObjectNotFound(_, _)
| EcstoreError::VersionNotFound(_, _, _),
) => true,
Error::Other(message) => {
message.contains("File not found")
|| message.contains("file not found")
|| message.contains("File version not found")
|| message.contains("file version not found")
|| message.contains("Object not found")
|| message.contains("object not found")
}
_ => false,
}
}
fn bucket_object_retry_delay(&self, retry_attempt: u32) -> Duration {
let base = Duration::from_secs(2_u64.saturating_pow(retry_attempt.clamp(1, MAX_BUCKET_OBJECT_HEAL_RETRIES)));
let jitter_seed = self
.id
.bytes()
.fold(0_u64, |acc, byte| acc.wrapping_mul(31).wrapping_add(u64::from(byte)));
Duration::from_millis(jitter_seed % 500).saturating_add(base)
}
fn should_return_typed_heal_error(err: &Error) -> bool {
matches!(err, Error::Storage(_) | Error::Disk(_))
}
async fn skip_data_usage_cache_heal_error(&self, bucket: &str, object: &str, err: &Error) -> bool {
if !Self::should_skip_data_usage_cache_heal_error(bucket, object, err) {
return false;
}
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "data_usage_cache_transient_skip",
error = %err,
"Heal object skipped for data usage cache after transient error"
);
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
true
}
async fn skip_scanner_synthetic_object_dir_missing(&self, bucket: &str, object: &str, err: &Error) -> bool {
if self.source != HealRequestSource::Scanner || !is_missing_object_dir_heal_result(object, err) {
return false;
}
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
source = self.source.as_str(),
result = "synthetic_object_dir_missing",
error = %err,
"Heal recreate skipped scanner synthetic object-dir candidate after object-level not-found"
);
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("skipped: {bucket}/{object}")));
progress.update_progress(4, 4, 0, 0);
true
}
#[tracing::instrument(skip(self), fields(task_id = %self.id, heal_type = ?self.heal_type))]
#[hotpath::measure]
pub async fn execute(&self) -> Result<()> {
// update status and timestamps atomically to avoid race conditions
let now = SystemTime::now();
let start_instant = Instant::now();
let queue_delay = now.duration_since(self.enqueued_at).unwrap_or_default();
let type_label = self.metric_type_label().to_string();
let set_label = self.metric_set_label();
{
let mut status = self.status.write().await;
let mut started_at = self.started_at.write().await;
let mut task_start_instant = self.task_start_instant.write().await;
*status = HealTaskStatus::Running;
*started_at = Some(now);
*task_start_instant = Some(start_instant);
}
histogram!(
"rustfs_heal_queue_delay_seconds",
"type" => type_label.clone(),
"set" => set_label.clone()
)
.record(queue_delay.as_secs_f64());
counter!(
"rustfs_heal_task_start_total",
"type" => type_label,
"set" => set_label
)
.increment(1);
demote_to_debug_when!(self.heal_type.is_per_object(), info, target: "rustfs::heal::task", {
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "started",
queue_delay = ?queue_delay,
"Heal task started"
});
self.emit_trace_task_state("started", Duration::ZERO, None);
let result = match &self.heal_type {
HealType::Cluster => self.heal_cluster().await,
HealType::Object {
bucket,
object,
version_id,
} => self.heal_object(bucket, object, version_id.as_deref()).await,
HealType::Bucket { bucket } => self.heal_bucket(bucket).await,
HealType::Prefix { bucket, prefix } => self.heal_prefix(bucket, prefix).await,
HealType::Metadata { bucket, object } => self.heal_metadata(bucket, object).await,
HealType::MRF { meta_path } => self.heal_mrf(meta_path).await,
HealType::ECDecode {
bucket,
object,
version_id,
} => self.heal_ec_decode(bucket, object, version_id.as_deref()).await,
HealType::ErasureSet { buckets, set_disk_id } => self.heal_erasure_set(buckets.clone(), set_disk_id.clone()).await,
};
// update completed time and status
{
let mut completed_at = self.completed_at.write().await;
*completed_at = Some(SystemTime::now());
}
match &result {
Ok(_) => {
let mut status = self.status.write().await;
*status = HealTaskStatus::Completed;
demote_to_debug_when!(self.heal_type.is_per_object(), info, target: "rustfs::heal::task", {
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "completed",
"Heal task completed"
});
}
Err(Error::TaskCancelled) => {
let mut status = self.status.write().await;
*status = HealTaskStatus::Cancelled;
info!(
target: "rustfs::heal::task",
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "cancelled",
"Heal task cancelled"
);
}
Err(Error::TaskTimeout) => {
let mut status = self.status.write().await;
*status = HealTaskStatus::Timeout;
demote_to_debug_when!(self.heal_type.is_per_object(), warn, target: "rustfs::heal::task", {
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "timed_out",
"Heal task timed out"
});
}
Err(e) => {
let mut status = self.status.write().await;
*status = HealTaskStatus::Failed { error: e.to_string() };
// Per-object failures are already logged with full object
// context by the heal_* implementations and terminally by the
// scheduler's task_failed error!; this generic duplicate would
// multiply every failed object by the retry count.
demote_to_debug_when!(self.heal_type.is_per_object(), error, target: "rustfs::heal::task", {
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "failed",
error = %e,
"Heal task failed"
});
}
}
let terminal_state = match &result {
Ok(_) => "completed",
Err(Error::TaskCancelled) => "cancelled",
Err(Error::TaskTimeout) => "timed_out",
Err(_) => "failed",
};
self.emit_trace_task_state(terminal_state, start_instant.elapsed(), result.as_ref().err());
result
}
pub async fn cancel(&self) -> Result<()> {
self.cancel_token.cancel();
let mut status = self.status.write().await;
*status = HealTaskStatus::Cancelled;
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_TASK_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
heal_type = self.heal_type.log_kind(),
state = "cancelled",
source = "manual",
"Heal task cancellation requested"
);
Ok(())
}
pub async fn get_status(&self) -> HealTaskStatus {
self.status.read().await.clone()
}
pub async fn get_progress(&self) -> HealProgress {
self.progress.read().await.clone()
}
pub async fn get_result_items(&self) -> Vec<HealResultItem> {
self.result_items.read().await.iter().map(|(_, item)| item.clone()).collect()
}
/// Sequence-stamped retained window, used when archiving a completed
/// task so incremental cursors survive the transition (HS-06).
pub async fn get_seqed_result_items(&self) -> Vec<(u64, HealResultItem)> {
self.result_items.read().await.clone()
}
/// Incremental result window (HS-06): `since = None` returns the full
/// retained window (legacy snapshot semantics); `since = Some(seq)`
/// returns only items stamped with a sequence greater than `seq`.
/// `lagged` warns that the caller's cursor fell behind the window start
/// and items were skipped (the response carries `min_seq` as the catch-up
/// cursor).
pub async fn get_result_items_since(&self, since: Option<u64>) -> HealResultWindow {
let result_items = self.result_items.read().await;
let next_seq = self.next_item_seq.load(Ordering::Relaxed);
let min_seq = self.min_available_seq.load(Ordering::Relaxed);
let mut lagged = false;
let items = match since {
None => result_items.iter().map(|(_, item)| item.clone()).collect::<Vec<_>>(),
Some(cursor) => {
if cursor + 1 < min_seq {
lagged = true;
}
result_items
.iter()
.filter(|(seq, _)| *seq > cursor)
.map(|(_, item)| item.clone())
.collect::<Vec<_>>()
}
};
HealResultWindow {
items,
next_seq,
min_seq,
lagged,
}
}
pub fn result_items_truncated(&self) -> bool {
self.result_items_truncated.load(Ordering::Relaxed)
}
async fn record_result_item(&self, result: HealResultItem) {
let seq = self.next_item_seq.fetch_add(1, Ordering::Relaxed);
let mut result_items = self.result_items.write().await;
if result_items.len() < MAX_RETAINED_HEAL_RESULT_ITEMS {
result_items.push((seq, result));
} else {
// Slide the window: the oldest item leaves and the cursor for the
// oldest still-available item moves forward with it.
result_items.remove(0);
self.min_available_seq
.store(result_items.first().map_or(seq, |(oldest, _)| *oldest), Ordering::Relaxed);
result_items.push((seq, result));
self.result_items_truncated.store(true, Ordering::Relaxed);
}
}
// specific heal implementation method
#[tracing::instrument(skip(self), fields(bucket = %bucket, object = %object, version_id = ?version_id))]
#[hotpath::measure]
async fn heal_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
version_id = ?version_id,
stage = "start",
"Heal object started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(0, 4, 0, 0);
}
// Step 1: Check if object exists and get metadata
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
stage = "check_existence",
"Heal object stage entered"
);
self.check_control_flags().await?;
let mut object_exists = match self.await_with_control(self.storage.object_exists(bucket, object)).await {
Ok(exists) => exists,
Err(err @ Error::TransientSkip { .. }) => {
return self.skip_due_to_transient_object_exists(bucket, object, &err).await;
}
Err(err) => return Err(err),
};
let canonicalized_object = if !object_exists {
match self.canonicalize_scanner_missing_object_dir(bucket, object).await {
Ok(canonicalized_object) => canonicalized_object,
Err(err @ Error::TransientSkip { .. }) => {
return self.skip_due_to_transient_object_exists(bucket, object, &err).await;
}
Err(err) => return Err(err),
}
} else {
None
};
let object = if let Some(canonicalized_object) = canonicalized_object.as_deref() {
object_exists = true;
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{canonicalized_object}")));
}
canonicalized_object
} else {
object
};
if !object_exists {
// Background loops (scanner/MRF/autoheal/read-repair) routinely
// race object deletion, so a missing target is per-object noise
// for them; only foreground admin/internal requests keep the warn.
let background_source = !matches!(self.source, HealRequestSource::Admin | HealRequestSource::Internal);
demote_to_debug_when!(background_source, warn, target: "rustfs::heal::task", {
event = EVENT_HEAL_OBJECT_MISSING,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
source = self.source.as_str(),
recreate_missing = self.options.recreate_missing,
"Heal target object is missing"
});
if self.options.recreate_missing {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
stage = "recreate_missing",
"Heal object recreate requested"
);
return self.recreate_missing_object(bucket, object, version_id).await;
} else if self.source == HealRequestSource::Scanner {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
stage = "scanner_missing_probe",
"Heal scanner missing object will be checked by storage layer"
);
} else {
return Err(Error::TaskExecutionFailed {
message: format!("Object not found: {bucket}/{object}"),
});
}
}
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
}
// Step 2: directly call ecstore to perform heal
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
stage = "heal_with_ecstore",
dry_run = self.options.dry_run,
remove_corrupted = self.options.remove_corrupted,
update_parity = self.options.update_parity,
"Heal object stage entered"
);
let heal_opts = HealOpts {
recursive: self.options.recursive,
dry_run: self.options.dry_run,
remove: self.options.remove_corrupted,
recreate: self.options.recreate_missing,
scan_mode: self.options.scan_mode,
update_parity: self.options.update_parity,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
let heal_result = self
.await_with_control(self.storage.heal_object(bucket, object, version_id, &heal_opts))
.await;
match heal_result {
Ok((result, error)) => {
if let Some(e) = error {
if self.skip_data_usage_cache_heal_error(bucket, object, &e).await {
return Ok(());
}
if Self::is_object_not_found_heal_error(&e) {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "treated_as_deleted",
"Heal missing object treated as deleted"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
return Ok(());
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal object operation failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
if Self::should_return_typed_heal_error(&e) {
return Err(e);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal object {bucket}/{object}: {e}"),
});
}
// Step 3: Verify heal result
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
stage = "verify_result",
"Heal object stage entered"
);
let object_size = result.object_size as u64;
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
object_size = object_size,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
result = "ok",
"Heal object repaired"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
if self.skip_data_usage_cache_heal_error(bucket, object, &e).await {
return Ok(());
}
if Self::is_object_not_found_heal_error(&e) {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "treated_as_deleted",
"Heal missing object treated as deleted"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
return Ok(());
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal object operation failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
if Self::should_return_typed_heal_error(&e) {
Err(e)
} else {
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal object {bucket}/{object}: {e}"),
})
}
}
}
}
async fn canonicalize_scanner_missing_object_dir(&self, bucket: &str, object: &str) -> Result<Option<String>> {
if self.source != HealRequestSource::Scanner {
return Ok(None);
}
let Some(candidate) = object.strip_suffix(SLASH_SEPARATOR) else {
return Ok(None);
};
if candidate.is_empty() {
return Ok(None);
}
match self.await_with_control(self.storage.object_exists(bucket, candidate)).await {
Ok(true) => {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object = %candidate,
canonicalized_from = %object,
stage = "canonicalize_scanner_object_dir",
result = "canonicalized",
"Heal scanner object-dir candidate canonicalized"
);
Ok(Some(candidate.to_string()))
}
Ok(false) => Ok(None),
Err(err) => Err(err),
}
}
/// Recreate missing object (for EC decode scenarios)
async fn recreate_missing_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
version_id = ?version_id,
stage = "recreate_missing",
"Heal object recreate started"
);
// Use ecstore's heal_object with recreate option
let heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: false,
recreate: true,
scan_mode: HealScanMode::Deep,
update_parity: true,
no_lock: self.options.no_lock,
pool: None,
set: None,
};
match self
.await_with_control(self.storage.heal_object(bucket, object, version_id, &heal_opts))
.await
{
Ok((result, error)) => {
if let Some(e) = error {
if self.skip_scanner_synthetic_object_dir_missing(bucket, object, &e).await {
return Ok(());
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "recreate_failed",
error = %e,
"Heal object recovery failed"
);
return Err(Error::TaskExecutionFailed {
message: format!("Failed to recreate missing object {bucket}/{object}: {e}"),
});
}
let object_size = result.object_size as u64;
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
object_size,
result = "recreated",
"Heal object recreated"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
if self.skip_scanner_synthetic_object_dir_missing(bucket, object, &e).await {
return Ok(());
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_OBJECT,
task_id = %self.id,
bucket,
object,
result = "recreate_failed",
error = %e,
"Heal object recovery failed"
);
Err(Error::TaskExecutionFailed {
message: format!("Failed to recreate missing object {bucket}/{object}: {e}"),
})
}
}
}
async fn heal_bucket(&self, bucket: &str) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
stage = "start",
recursive = self.options.recursive,
"Heal bucket started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("bucket: {bucket}")));
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if bucket exists
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
stage = "check_existence",
"Heal bucket stage entered"
);
self.check_control_flags().await?;
let bucket_exists = self.await_with_control(self.storage.get_bucket_info(bucket)).await?.is_some();
if !bucket_exists {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
result = "missing",
"Heal bucket failed because the bucket does not exist"
);
return Err(Error::TaskExecutionFailed {
message: format!("Bucket not found: {bucket}"),
});
}
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform bucket heal using ecstore
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
stage = "heal_with_ecstore",
dry_run = self.options.dry_run,
"Heal bucket stage entered"
);
let heal_opts = HealOpts {
recursive: self.options.recursive,
dry_run: self.options.dry_run,
remove: if self.options.recursive {
false
} else {
self.options.remove_corrupted
},
recreate: self.options.recreate_missing,
scan_mode: self.options.scan_mode,
update_parity: self.options.update_parity,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
let heal_result = self.await_with_control(self.storage.heal_bucket(bucket, &heal_opts)).await;
match heal_result {
Ok(result) => {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
recursive = self.options.recursive,
result = "ok",
"Heal bucket completed"
);
self.record_result_item(result).await;
if self.options.recursive {
self.heal_bucket_objects(bucket, "").await?;
}
if !self.options.recursive {
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
result = "failed",
error = %e,
"Heal bucket failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal bucket {bucket}: {e}"),
})
}
}
}
async fn heal_cluster(&self) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
stage = "cluster_recursive",
"Heal cluster started"
);
let bucket_infos = self.await_with_control(self.storage.list_buckets()).await?;
let mut failed = 0_u64;
let mut retryable = 0_u64;
let mut permanent = 0_u64;
let mut first_object = None;
let mut first_error = None;
for bucket_info in bucket_infos {
self.check_control_flags().await?;
let mut retry_attempt = 0_u32;
loop {
match self.heal_bucket(&bucket_info.name).await {
Ok(()) => break,
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
Err(err) => {
if let Some(failure) = self.take_batch_failure().await {
failed = failed.saturating_add(failure.failed);
retryable = retryable.saturating_add(failure.retryable);
permanent = permanent.saturating_add(failure.permanent);
first_object.get_or_insert(failure.first_object);
first_error.get_or_insert(failure.first_error);
break;
}
if err.is_recoverable_heal() && retry_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES {
retry_attempt = retry_attempt.saturating_add(1);
self.await_with_control(async {
tokio::time::sleep(self.bucket_object_retry_delay(retry_attempt)).await;
Ok(())
})
.await?;
continue;
}
failed = failed.saturating_add(1);
if err.is_recoverable_heal() {
retryable = retryable.saturating_add(1);
} else {
permanent = permanent.saturating_add(1);
}
first_object.get_or_insert(bucket_info.name.clone());
first_error.get_or_insert_with(|| err.to_string());
break;
}
}
}
}
if failed > 0 {
let failure = BatchHealFailure {
scope: "cluster".to_string(),
failed,
retryable,
permanent,
first_object: first_object.unwrap_or_default(),
first_error: first_error.unwrap_or_default(),
};
return Err(self.record_batch_failure(failure).await);
}
Ok(())
}
async fn heal_prefix(&self, bucket: &str, prefix: &str) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
prefix,
stage = "prefix_recursive",
"Heal prefix started"
);
self.heal_bucket_objects(bucket, prefix).await
}
#[hotpath::measure]
async fn heal_bucket_objects(&self, bucket: &str, prefix: &str) -> Result<()> {
let mut continuation_token: Option<String> = None;
let mut scanned = 0u64;
let mut healed = 0u64;
let mut failed = 0u64;
let mut retryable_failed = 0u64;
let mut permanent_failed = 0u64;
let mut bytes = 0u64;
let mut first_failed_object = None;
let mut first_error = None;
let mut failure_samples_logged = 0_u64;
let heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: self.options.remove_corrupted,
recreate: self.options.recreate_missing,
scan_mode: self.options.scan_mode,
update_parity: self.options.update_parity,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
loop {
self.check_control_flags().await?;
let (objects, next_token, is_truncated) = self
.await_with_control(
self.storage
.list_objects_for_heal_page(bucket, prefix, continuation_token.as_deref(), false),
)
.await?;
let mut pending = objects;
let mut retry_attempt = 0_u32;
while !pending.is_empty() {
if retry_attempt > 0 {
self.await_with_control(async {
tokio::time::sleep(self.bucket_object_retry_delay(retry_attempt)).await;
Ok(())
})
.await?;
}
let mut retry = Vec::with_capacity(pending.len());
for item in pending {
self.check_control_flags().await?;
let object = item.name.as_str();
if retry_attempt == 0 {
scanned = scanned.saturating_add(1);
}
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(scanned, healed, failed, bytes);
}
let error = match self
.await_with_control(
self.storage
.heal_object(bucket, object, item.version_id.as_deref(), &heal_opts),
)
.await
{
Ok((result, None)) => {
healed = healed.saturating_add(1);
bytes = bytes.saturating_add(u64::try_from(result.object_size).unwrap_or_default());
self.record_result_item(result).await;
None
}
Ok((_, Some(err))) if is_missing_object_dir_heal_result(object, &err) => {
healed = healed.saturating_add(1);
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "object_dir_not_found_skipped",
"Heal bucket object-dir candidate skipped after not-found result"
);
None
}
Ok((_, Some(err))) | Err(err) => Some(err),
};
if let Some(err) = error {
if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "transient_skip",
error = %err,
"Heal bucket object repair skipped due to transient metadata error"
);
} else if err.is_recoverable_heal() && retry_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
retry_attempt = retry_attempt.saturating_add(1),
error = %err,
result = "object_retry_scheduled",
"Heal bucket object retry scheduled"
);
retry.push(item);
} else {
failed = failed.saturating_add(1);
if err.is_recoverable_heal() {
retryable_failed = retryable_failed.saturating_add(1);
} else {
permanent_failed = permanent_failed.saturating_add(1);
}
first_failed_object.get_or_insert_with(|| object.to_string());
first_error.get_or_insert_with(|| err.to_string());
if take_failure_log_sample(&mut failure_samples_logged) {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
retry_attempt,
error = %err,
result = "object_failed",
"Heal bucket object repair failed"
);
}
}
}
let mut progress = self.progress.write().await;
progress.update_progress(scanned, healed, failed, bytes);
}
pending = retry;
retry_attempt = retry_attempt.saturating_add(1);
}
if !is_truncated {
break;
}
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.
break;
}
}
if failed > 0 {
let failure = BatchHealFailure {
scope: format!("bucket:{bucket}"),
failed,
retryable: retryable_failed,
permanent: permanent_failed,
first_object: first_failed_object.unwrap_or_default(),
first_error: first_error.unwrap_or_default(),
};
return Err(self.record_batch_failure(failure).await);
}
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
prefix,
scanned,
healed,
failed,
bytes_processed = bytes,
result = "recursive_ok",
"Heal bucket recursive pass completed"
);
Ok(())
}
async fn apply_erasure_set_usage_baseline(&self, buckets: &[String]) -> Result<()> {
let baseline = match self
.await_with_control(self.storage.erasure_set_usage_baseline(buckets))
.await
{
Ok(Some(baseline)) => baseline,
Ok(None) => return Ok(()),
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(_) => return Ok(()),
};
let HealBucketUsageBaseline { objects_count, bytes } = baseline;
let mut progress = self.progress.write().await;
progress.set_total_baseline(objects_count, bytes);
Ok(())
}
async fn heal_metadata(&self, bucket: &str, object: &str) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
stage = "start",
"Heal metadata started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("metadata: {bucket}/{object}")));
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if object exists
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
stage = "check_existence",
"Heal metadata stage entered"
);
self.check_control_flags().await?;
let object_exists = match self.await_with_control(self.storage.object_exists(bucket, object)).await {
Ok(exists) => exists,
Err(err @ Error::TransientSkip { .. }) => {
return self.skip_due_to_transient_object_exists(bucket, object, &err).await;
}
Err(err) => return Err(err),
};
if !object_exists {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "missing",
"Heal metadata failed because object is missing"
);
return Err(Error::TaskExecutionFailed {
message: format!("Object not found: {bucket}/{object}"),
});
}
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform metadata heal using ecstore
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
stage = "heal_with_ecstore",
"Heal metadata stage entered"
);
let heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: false,
recreate: false,
scan_mode: HealScanMode::Deep,
update_parity: false,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
let heal_result = self
.await_with_control(self.storage.heal_object(bucket, object, None, &heal_opts))
.await;
match heal_result {
Ok((result, error)) => {
if let Some(e) = error {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal metadata failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal metadata {bucket}/{object}: {e}"),
});
}
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
result = "ok",
"Heal metadata repaired"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
self.record_result_item(result).await;
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_METADATA_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal metadata failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal metadata {bucket}/{object}: {e}"),
})
}
}
}
async fn heal_mrf(&self, meta_path: &str) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_MRF_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
meta_path,
stage = "start",
"Heal MRF started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("mrf: {meta_path}")));
progress.update_progress(0, 2, 0, 0);
}
// Parse meta_path to extract bucket and object
let parts: Vec<&str> = meta_path.split('/').collect();
if parts.len() < 2 {
return Err(Error::TaskExecutionFailed {
message: format!("Invalid meta path format: {meta_path}"),
});
}
let bucket = parts[0];
let object = parts[1..].join("/");
// Step 1: Perform MRF heal using ecstore
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_MRF_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
meta_path,
bucket,
object = %object,
stage = "heal_with_ecstore",
"Heal MRF stage entered"
);
let heal_opts = HealOpts {
recursive: true,
dry_run: self.options.dry_run,
remove: self.options.remove_corrupted,
recreate: self.options.recreate_missing,
scan_mode: HealScanMode::Deep,
update_parity: true,
no_lock: self.options.no_lock,
pool: None,
set: None,
};
let heal_result = self
.await_with_control(self.storage.heal_object(bucket, &object, None, &heal_opts))
.await;
match heal_result {
Ok((result, error)) => {
if let Some(e) = error {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_MRF_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
meta_path,
bucket,
object = %object,
result = "failed",
error = %e,
"Heal MRF failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(2, 2, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal MRF {meta_path}: {e}"),
});
}
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_MRF_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
meta_path,
bucket,
object = %object,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
result = "ok",
"Heal MRF repaired"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(2, 2, 0, 0);
}
self.record_result_item(result).await;
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_MRF_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
meta_path,
bucket,
object = %object,
result = "failed",
error = %e,
"Heal MRF failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(2, 2, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal MRF {meta_path}: {e}"),
})
}
}
}
async fn heal_ec_decode(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
version_id = ?version_id,
stage = "start",
"Heal EC decode started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("ec_decode: {bucket}/{object}")));
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if object exists
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
stage = "check_existence",
"Heal EC decode stage entered"
);
self.check_control_flags().await?;
let object_exists = match self.await_with_control(self.storage.object_exists(bucket, object)).await {
Ok(exists) => exists,
Err(err @ Error::TransientSkip { .. }) => {
return self.skip_due_to_transient_object_exists(bucket, object, &err).await;
}
Err(err) => return Err(err),
};
if !object_exists {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "missing",
"Heal EC decode failed because object is missing"
);
return Err(Error::TaskExecutionFailed {
message: format!("Object not found: {bucket}/{object}"),
});
}
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform EC decode heal using ecstore
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
stage = "heal_with_ecstore",
"Heal EC decode stage entered"
);
let heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: false,
recreate: true,
scan_mode: HealScanMode::Deep,
update_parity: true,
no_lock: self.options.no_lock,
pool: None,
set: None,
};
let heal_result = self
.await_with_control(self.storage.heal_object(bucket, object, version_id, &heal_opts))
.await;
match heal_result {
Ok((result, error)) => {
if let Some(e) = error {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal EC decode failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal EC decode {bucket}/{object}: {e}"),
});
}
let object_size = result.object_size as u64;
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
object_size,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
result = "ok",
"Heal EC decode repaired"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_EC_DECODE_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
bucket,
object,
result = "failed",
error = %e,
"Heal EC decode failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal EC decode {bucket}/{object}: {e}"),
})
}
}
}
async fn heal_erasure_set(&self, buckets: Vec<String>, set_disk_id: String) -> Result<()> {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
bucket_count = buckets.len(),
stage = "start",
"Heal erasure set started"
);
// update progress
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("erasure_set: {} ({} buckets)", set_disk_id, buckets.len())));
progress.update_progress(0, 4, 0, 0);
}
let is_auto_replacement = matches!(self.source, HealRequestSource::AutoHeal) && !self.heal_endpoints.is_empty();
let replacement_resume_disk = if is_auto_replacement {
let mut requested_targets = self.heal_endpoints.clone();
requested_targets.sort_unstable();
requested_targets.dedup();
let selection = self
.await_with_control(
self.storage
.get_replacement_resume_disk(&set_disk_id, &self.id, &self.heal_endpoints),
)
.await?;
let disk = match selection {
crate::heal::storage::ReplacementResumeDisk::Existing(disk) => {
if let Some(anchor) = &self.replacement_resume_endpoint
&& disk.endpoint().to_string() != *anchor
{
return Err(Error::TaskExecutionFailed {
message: format!("Replacement resume anchor changed for automatic heal {set_disk_id}"),
});
}
Some(disk)
}
crate::heal::storage::ReplacementResumeDisk::Fresh => {
if self.replacement_resume_endpoint.is_some() {
return Err(Error::TaskExecutionFailed {
message: format!("Replacement resume anchor is unavailable for automatic heal {set_disk_id}"),
});
}
None
}
};
if let Some(disk) = disk.as_ref()
&& ResumeManager::has_replacement_intent(disk, &self.id).await
{
let resume_manager = ResumeManager::load_replacement_intent(disk.clone(), &self.id).await?;
let state = resume_manager.get_state().await;
if state.completed
&& matches!(state.replacement_phase, ReplacementPhase::CleanupPending)
&& state.set_disk_id == set_disk_id
&& state.replacement_targets == requested_targets
&& state.replacement_generation.as_deref() == Some(self.id.as_str())
{
resume_manager.ensure_replacement_completion_proof().await?;
if CheckpointManager::has_checkpoint(disk, &self.id).await {
CheckpointManager::load_from_disk(disk.clone(), &self.id)
.await?
.cleanup()
.await?;
}
resume_manager.cleanup().await?;
return Ok(());
}
}
disk
} else {
None
};
if is_auto_replacement
&& !self
.await_with_control(self.storage.replacement_targets_ready(&self.heal_endpoints))
.await?
{
return Err(Error::TaskExecutionFailed {
message: format!("Replacement target is no longer ready for automatic heal {set_disk_id}"),
});
}
let replacement_resume_disk = if is_auto_replacement {
Some(match replacement_resume_disk {
Some(disk) => disk,
None => {
self.await_with_control(self.storage.get_disk_for_resume_excluding(&set_disk_id, &self.heal_endpoints))
.await?
}
})
} else {
None
};
let mut buckets = if buckets.is_empty() {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
stage = "list_buckets",
"Heal erasure set bucket list resolved"
);
let bucket_infos = self.await_with_control(self.storage.list_buckets()).await?;
bucket_infos.into_iter().map(|info| info.name).collect()
} else {
buckets
};
// Persist automatic replacement intent on a surviving disk before the
// first target format write. A task retry keeps this id; a newly
// admitted blank replacement gets a fresh id and cannot reuse cursor
// progress from an older disk at the same endpoint.
let replacement_resume = if is_auto_replacement {
let identities = self
.await_with_control(self.storage.replacement_target_identities(&self.heal_endpoints))
.await?;
let disk = replacement_resume_disk.clone().ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Replacement resume disk is missing for automatic heal {set_disk_id}"),
})?;
let manager = ResumeManager::new_replacement_intent(
disk.clone(),
self.id.clone(),
set_disk_id.clone(),
buckets.clone(),
self.heal_endpoints.clone(),
identities.clone(),
)
.await?;
buckets = manager.get_state().await.replacement_buckets;
Some((disk, manager, identities))
} else {
None
};
self.apply_erasure_set_usage_baseline(&buckets).await?;
let healing_marker = format!("{set_disk_id}:{}", self.id);
if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() {
let state = resume_manager.get_state().await;
if state.completed && matches!(state.replacement_phase, ReplacementPhase::Verified) {
resume_manager.ensure_replacement_completion_proof().await?;
super::clear_healing_markers_after_verified(&self.heal_endpoints, &healing_marker).await?;
resume_manager.mark_replacement_cleanup_pending().await?;
if CheckpointManager::has_checkpoint(disk, &self.id).await {
CheckpointManager::load_from_disk(disk.clone(), &self.id)
.await?
.cleanup()
.await?;
}
resume_manager.cleanup().await?;
return Ok(());
}
}
// Step 1: Perform disk format heal using ecstore
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
stage = "heal_format",
"Heal erasure set stage entered"
);
if is_auto_replacement {
let Some((_, _, expected_identities)) = replacement_resume.as_ref() else {
return Err(Error::TaskExecutionFailed {
message: format!("Replacement intent is missing for automatic heal {set_disk_id}"),
});
};
self.verify_replacement_identity_fence(expected_identities, &set_disk_id, "format")
.await?;
}
let format_result = if is_auto_replacement {
let pool_index = self.options.pool_index.ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Missing pool scope for automatic replacement heal {set_disk_id}"),
})?;
let set_index = self.options.set_index.ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Missing set scope for automatic replacement heal {set_disk_id}"),
})?;
self.await_with_control(self.storage.heal_replacement_format(
self.options.dry_run,
pool_index,
set_index,
&self.heal_endpoints,
))
.await
} else {
self.await_with_control(self.storage.heal_format(self.options.dry_run)).await
};
match format_result {
Ok((result, error)) => {
if let Some(e) = error {
if Self::is_no_heal_required_error(&e) {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
result = "format_noop",
"Heal erasure set format repair skipped because no format heal was required"
);
} else {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
result = "format_failed",
error = %e,
"Heal erasure set failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
});
}
} else {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
drives_healed = result.drives_healed(),
drives_total = result.drives_reported(),
result = "format_ok",
"Heal erasure set format repaired"
);
}
if !self.options.dry_run && !target_outcomes_complete(&result, &self.heal_endpoints) {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to verify formatted replacement targets for {set_disk_id}"),
});
}
if let Some((_, replacement_resume, expected_identities)) = &replacement_resume {
let identities = self
.await_with_control(self.storage.replacement_target_identities(&self.heal_endpoints))
.await?;
if !replacement_target_identities_match(expected_identities, &identities) {
return Err(Error::TaskExecutionFailed {
message: format!("Replacement target changed after format for automatic heal {set_disk_id}"),
});
}
replacement_resume.mark_replacement_rebuilding(identities).await?;
}
}
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
result = "format_failed",
error = %e,
"Heal erasure set failed"
);
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
});
}
}
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 4, 0, 0);
}
// The rebuilt disks are formatted now: mark them as healing so
// DiskInfo.healing reflects the rebuild until it completes.
super::set_healing_markers(&self.heal_endpoints, &healing_marker).await?;
// Step 2: Get disk for resume functionality
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
stage = "resolve_resume_disk",
"Heal erasure set stage entered"
);
let replacement_target_identities = replacement_resume.as_ref().map(|(_, _, identities)| identities.clone());
let disk = match replacement_resume.as_ref() {
Some((disk, _, _)) => disk.clone(),
None => {
self.await_with_control(self.storage.get_disk_for_resume(&set_disk_id))
.await?
}
};
{
let mut progress = self.progress.write().await;
progress.update_progress(2, 4, 0, 0);
}
// Step 3: Heal bucket structure
// Check control flags before each iteration to ensure timely cancellation.
let bucket_heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: false,
recreate: self.options.recreate_missing,
scan_mode: self.options.scan_mode,
update_parity: self.options.update_parity,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
for bucket in buckets.iter() {
// Check control flags before starting each bucket heal
self.check_control_flags().await?;
if let Some(expected_identities) = replacement_target_identities.as_ref() {
self.verify_replacement_identity_fence(expected_identities, &set_disk_id, "bucket prepass")
.await?;
}
let heal_result = self
.await_with_control(self.storage.heal_bucket(bucket, &bucket_heal_opts))
.await;
match heal_result {
Ok(result) => {
self.record_result_item(result).await;
}
Err(err) => {
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
bucket,
result = "bucket_failed",
error = %err,
"Heal erasure set bucket prepass failed"
);
return Err(err);
}
}
}
// Create erasure set healer with resume support
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
stage = "build_resumable_healer",
"Heal erasure set stage entered"
);
let heal_opts = HealOpts {
recursive: self.options.recursive,
dry_run: self.options.dry_run,
remove: self.options.remove_corrupted,
recreate: self.options.recreate_missing,
scan_mode: self.options.scan_mode,
update_parity: self.options.update_parity,
no_lock: self.options.no_lock,
pool: self.options.pool_index,
set: self.options.set_index,
};
let erasure_healer = ErasureSetHealer::new(
self.storage.clone(),
self.progress.clone(),
self.cancel_token.clone(),
disk,
heal_opts,
self.source,
)
.with_replacement_targets(self.heal_endpoints.clone(), is_auto_replacement.then(|| self.id.clone()))
.with_replacement_identity_fence(replacement_target_identities.clone());
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 4, 0, 0);
}
// Step 4: Execute erasure set heal with resume
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_STAGE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
stage = "execute_resumable_heal",
"Heal erasure set stage entered"
);
let result = self
.await_with_control(erasure_healer.heal_erasure_set(&buckets, &set_disk_id))
.await;
// Keep the markers on failure: the resume state also persists, and the
// next run of this set heal re-marks and eventually clears them.
let result = match result {
Ok(()) => {
if let Some(expected_identities) = replacement_target_identities.as_ref() {
self.verify_replacement_identity_fence(expected_identities, &set_disk_id, "marker completion")
.await?;
}
super::clear_healing_markers_after_verified(&self.heal_endpoints, &healing_marker).await?;
if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() {
resume_manager.mark_replacement_cleanup_pending().await?;
if CheckpointManager::has_checkpoint(disk, &self.id).await {
CheckpointManager::load_from_disk(disk.clone(), &self.id)
.await?
.cleanup()
.await?;
}
resume_manager.cleanup().await?;
}
Ok(())
}
Err(err) => Err(err),
};
{
let mut progress = self.progress.write().await;
let bytes_processed = progress.bytes_processed;
progress.update_progress(4, 4, 0, bytes_processed);
}
match result {
Ok(_) => {
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
bucket_count = buckets.len(),
result = "ok",
"Heal erasure set repaired"
);
Ok(())
}
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_TASK,
task_id = %self.id,
set_disk_id,
result = "failed",
error = %e,
"Heal erasure set failed"
);
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal erasure set {set_disk_id}: {e}"),
})
}
}
}
}
impl std::fmt::Debug for HealTask {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HealTask")
.field("id", &self.id)
.field("heal_type", &self.heal_type)
.field("options", &self.options)
.field("created_at", &self.created_at)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::super::{DiskOption, DiskStore, Endpoint, HealDiskExt as _, new_disk};
use super::*;
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};
use std::collections::{HashMap, VecDeque};
use std::sync::Mutex;
use tempfile::TempDir;
use super::super::storage_api::status::BucketInfo;
#[tokio::test]
async fn retry_request_carries_remaining_timeout_budget() {
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage::default());
let mut request = HealRequest::bucket("bucket".to_string());
request.options.timeout = Some(Duration::from_secs(100));
let task = HealTask::from_request(request, storage.clone());
*task.task_start_instant.write().await = Some(Instant::now() - Duration::from_secs(40));
let retry = task
.retry_request_with_remaining_timeout()
.await
.expect("first retry should retain the unused timeout budget");
let first_remaining = retry.options.timeout.expect("configured timeout should remain present");
assert!(first_remaining <= Duration::from_secs(60));
assert!(first_remaining > Duration::from_secs(59));
let retry_task = HealTask::from_request(retry, storage);
*retry_task.task_start_instant.write().await = Some(Instant::now() - Duration::from_secs(20));
let second_retry = retry_task
.retry_request_with_remaining_timeout()
.await
.expect("second retry should retain only the unused aggregate budget");
let second_remaining = second_retry
.options
.timeout
.expect("configured timeout should remain present");
assert!(second_remaining <= Duration::from_secs(40));
assert!(second_remaining > Duration::from_secs(39));
}
#[test]
fn format_result_requires_every_requested_target_to_be_ok() {
let result = HealResultItem {
after: Infos {
drives: vec![
HealDriveInfo {
endpoint: "disk-a".to_string(),
state: "ok".to_string(),
..Default::default()
},
HealDriveInfo {
endpoint: "disk-b".to_string(),
state: "missing".to_string(),
..Default::default()
},
],
},
..Default::default()
};
assert!(target_outcomes_complete(&result, &["disk-a".to_string()]));
assert!(!target_outcomes_complete(&result, &["disk-a".to_string(), "disk-b".to_string()]));
assert!(!target_outcomes_complete(&result, &["disk-c".to_string()]));
}
#[tokio::test]
async fn automatic_replacement_uses_target_scoped_format() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
replacement_targets_ready: Mutex::new(true),
resume_disk: Mutex::new(Some(disk)),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: Vec::new(),
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..Default::default()
},
HealPriority::Low,
);
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect_err("the mock has no local replacement marker target");
assert_eq!(
*storage.global_format_calls.lock().unwrap(),
0,
"automatic replacement must not call global format"
);
assert_eq!(
storage.replacement_format_calls.lock().unwrap().as_slice(),
&[(0, 0, vec!["replacement-a".to_string()])],
"automatic replacement must pass the exact pool, set, and target"
);
}
#[tokio::test]
async fn automatic_replacement_persists_intent_before_format() {
let storage = Arc::new(MockStorage {
replacement_targets_ready: Mutex::new(true),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
HealTask::from_request(request, storage.clone())
.execute()
.await
.expect_err("intent persistence needs a healthy non-target disk");
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"format must not start before the durable replacement intent exists"
);
}
#[tokio::test]
async fn recovered_replacement_never_uses_a_fresh_resume_disk() {
let storage = Arc::new(MockStorage {
replacement_targets_ready: Mutex::new(true),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
let error = HealTask::from_replacement_recovery_request(request, storage.clone(), Some("survivor-a".to_string()))
.execute()
.await
.expect_err("a durable recovery must not fall back to another resume disk");
assert!(error.to_string().contains("resume anchor is unavailable"));
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"an unavailable durable anchor must block formatting before any write"
);
assert!(!*storage.listed.lock().unwrap(), "an unavailable durable anchor must not list buckets");
}
#[tokio::test]
async fn automatic_replacement_rejects_a_new_identity_after_format() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let disk = make_resume_disk(&temp).await;
let first_identity = replacement_identity("replacement-a", "device-a", "filesystem-a");
let second_identity = replacement_identity("replacement-a", "device-b", "filesystem-b");
let storage = Arc::new(MockStorage {
replacement_targets_ready: Mutex::new(true),
replacement_target_identity_sequences: Mutex::new(VecDeque::from([
vec![first_identity.clone()],
vec![first_identity.clone()],
vec![second_identity],
])),
resume_disk: Mutex::new(Some(disk.clone())),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
let task = HealTask::from_request(request, storage.clone());
let error = task
.execute()
.await
.expect_err("a remounted target after format must fail closed");
assert!(error.to_string().contains("changed after format"));
assert_eq!(storage.replacement_format_calls.lock().unwrap().len(), 1);
assert!(storage.bucket_heal_calls.lock().unwrap().is_empty());
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
let state = ResumeManager::load_replacement_intent(disk, &task.id)
.await
.expect("durable replacement intent should remain available")
.get_state()
.await;
assert_eq!(state.replacement_phase, crate::heal::resume::ReplacementPhase::Intent);
assert_eq!(state.replacement_target_identities, vec![first_identity]);
}
#[tokio::test]
async fn automatic_replacement_reuses_an_existing_non_target_resume_anchor() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let anchor = make_resume_disk(&temp).await;
let task_id = crate::heal::resume::ResumeUtils::generate_task_id();
let identity = ReplacementTargetIdentity {
endpoint: "replacement-a".to_string(),
canonical_path: "/replacement/replacement-a".to_string(),
physical_device_ids: vec!["replacement-a".to_string()],
filesystem_identity: "identity-replacement-a".to_string(),
};
ResumeManager::new_replacement_intent(
anchor.clone(),
task_id.clone(),
"pool_0_set_0".to_string(),
vec!["bucket-a".to_string()],
vec!["replacement-a".to_string()],
vec![identity],
)
.await
.expect("existing intent should be stored on the non-target anchor");
let storage = Arc::new(MockStorage {
replacement_targets_ready: Mutex::new(true),
replacement_resume_disk: Mutex::new(Some(anchor.clone())),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.id = task_id.clone();
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
HealTask::from_request(request, storage.clone())
.execute()
.await
.expect_err("the test has no mounted marker target after format");
assert_eq!(
storage.replacement_format_calls.lock().unwrap().len(),
1,
"an existing non-target anchor must be reused instead of falling back to a fresh anchor"
);
assert!(
storage.resume_disk.lock().unwrap().is_none(),
"the fresh resume-anchor fallback must remain unused"
);
let state = ResumeManager::load_replacement_intent(anchor, &task_id)
.await
.expect("the existing non-target anchor should retain the generation")
.get_state()
.await;
assert_eq!(state.replacement_phase, ReplacementPhase::Rebuilding);
}
#[tokio::test]
async fn automatic_replacement_defers_before_bucket_listing_when_target_is_unready() {
let storage = Arc::new(MockStorage::default());
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: Vec::new(),
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..Default::default()
},
HealPriority::Low,
);
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
HealTask::from_request(request, storage.clone())
.execute()
.await
.expect_err("an unsafe replacement must defer before any scan work");
assert!(!*storage.listed.lock().unwrap(), "unsafe targets must not list buckets");
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"unsafe targets must not format"
);
}
#[tokio::test]
async fn cleanup_pending_recovery_skips_target_readiness_and_format() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let anchor = make_resume_disk(&temp).await;
let task_id = crate::heal::resume::ResumeUtils::generate_task_id();
let identity = replacement_identity("replacement-a", "device-a", "filesystem-a");
let resume_manager = ResumeManager::new_replacement_intent(
anchor.clone(),
task_id.clone(),
"pool_0_set_0".to_string(),
vec!["bucket-a".to_string()],
vec!["replacement-a".to_string()],
vec![identity],
)
.await
.expect("terminal replacement state should persist on the survivor anchor");
resume_manager
.mark_replacement_completed_and_verified()
.await
.expect("terminal replacement proof should persist before cleanup");
resume_manager
.mark_replacement_cleanup_pending()
.await
.expect("failed cleanup must retain a cleanup-pending state");
let storage = Arc::new(MockStorage {
replacement_resume_disk: Mutex::new(Some(anchor.clone())),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.id = task_id.clone();
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
HealTask::from_replacement_recovery_request(request, storage.clone(), Some(anchor.endpoint().to_string()))
.execute()
.await
.expect("cleanup-pending recovery must not require a mounted replacement target");
assert!(
!ResumeManager::has_resume_state(&anchor, &task_id).await,
"terminal cleanup must remove the retained resume state"
);
assert_eq!(*storage.global_format_calls.lock().unwrap(), 0);
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"terminal cleanup must not format replacement targets"
);
assert!(storage.bucket_heal_calls.lock().unwrap().is_empty());
assert!(!*storage.listed.lock().unwrap());
}
#[tokio::test]
async fn cleanup_pending_recovery_removes_checkpoint_without_rebuild_work() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let anchor = make_resume_disk(&temp).await;
let task_id = crate::heal::resume::ResumeUtils::generate_task_id();
let identity = replacement_identity("replacement-a", "device-a", "filesystem-a");
let resume_manager = ResumeManager::new_replacement_intent(
anchor.clone(),
task_id.clone(),
"pool_0_set_0".to_string(),
vec!["bucket-a".to_string()],
vec!["replacement-a".to_string()],
vec![identity],
)
.await
.expect("terminal replacement state should persist on the survivor anchor");
resume_manager
.mark_replacement_completed_and_verified()
.await
.expect("terminal replacement proof should persist before cleanup");
resume_manager
.mark_replacement_cleanup_pending()
.await
.expect("failed cleanup must retain a cleanup-pending state");
CheckpointManager::new(anchor.clone(), task_id.clone())
.await
.expect("checkpoint fixture should persist");
assert!(
CheckpointManager::has_checkpoint(&anchor, &task_id).await,
"checkpoint fixture must exist before restart cleanup"
);
let storage = Arc::new(MockStorage {
replacement_resume_disk: Mutex::new(Some(anchor.clone())),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.id = task_id.clone();
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec!["replacement-a".to_string()];
HealTask::from_replacement_recovery_request(request, storage.clone(), Some(anchor.endpoint().to_string()))
.execute()
.await
.expect("cleanup-pending recovery must finish terminal cleanup");
assert!(
!CheckpointManager::has_checkpoint(&anchor, &task_id).await,
"terminal cleanup must remove the retained checkpoint"
);
assert!(
!ResumeManager::has_resume_state(&anchor, &task_id).await,
"terminal cleanup must remove the retained resume state"
);
assert_eq!(*storage.global_format_calls.lock().unwrap(), 0);
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"terminal checkpoint cleanup must not format replacement targets"
);
assert!(storage.bucket_heal_calls.lock().unwrap().is_empty());
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
assert!(!*storage.listed.lock().unwrap());
}
#[tokio::test]
async fn verified_recovery_keeps_state_when_marker_clear_fails() {
let temp = TempDir::new().expect("temporary resume disk directory should be created");
let anchor = make_resume_disk(&temp).await;
let task_id = crate::heal::resume::ResumeUtils::generate_task_id();
let target = format!("replacement-marker-missing-{task_id}");
let identity = replacement_identity(&target, &target, &format!("identity-{target}"));
let resume_manager = ResumeManager::new_replacement_intent(
anchor.clone(),
task_id.clone(),
"pool_0_set_0".to_string(),
vec!["bucket-a".to_string()],
vec![target.clone()],
vec![identity],
)
.await
.expect("verified replacement state should persist on the survivor anchor");
resume_manager
.mark_replacement_completed_and_verified()
.await
.expect("verified state must persist proof before marker cleanup");
let storage = Arc::new(MockStorage {
replacement_resume_disk: Mutex::new(Some(anchor.clone())),
replacement_targets_ready: Mutex::new(true),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
pool_index: Some(0),
set_index: Some(0),
..HealOptions::default()
},
HealPriority::Low,
);
request.id = task_id.clone();
request.source = HealRequestSource::AutoHeal;
request.heal_endpoints = vec![target];
let error = HealTask::from_replacement_recovery_request(request, storage.clone(), Some(anchor.endpoint().to_string()))
.execute()
.await
.expect_err("marker clear failure must keep the durable terminal state retryable");
assert!(error.to_string().contains("healing marker target is unavailable"));
let state = ResumeManager::load_replacement_intent(anchor.clone(), &task_id)
.await
.expect("verified state must remain for retry after marker clear failure")
.get_state()
.await;
assert!(state.completed);
assert_eq!(state.replacement_phase, ReplacementPhase::Verified);
assert_eq!(*storage.global_format_calls.lock().unwrap(), 0);
assert!(
storage.replacement_format_calls.lock().unwrap().is_empty(),
"marker cleanup retry must not format replacement targets again"
);
assert!(storage.bucket_heal_calls.lock().unwrap().is_empty());
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
assert!(!*storage.listed.lock().unwrap());
}
#[derive(Default)]
struct MockStorage {
listed: Mutex<bool>,
healed_objects: Mutex<Vec<String>>,
heal_object_calls: Mutex<Vec<String>>,
heal_object_version_ids: Mutex<Vec<Option<String>>>,
bucket_heal_opts: Mutex<Vec<HealOpts>>,
object_heal_opts: Mutex<Vec<HealOpts>>,
object_exists: Mutex<Option<bool>>,
object_exists_by_name: Mutex<HashMap<String, MockObjectExists>>,
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
format_no_heal_required: Mutex<bool>,
global_format_calls: Mutex<u32>,
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
replacement_targets_ready: Mutex<bool>,
replacement_target_identity_sequences: Mutex<VecDeque<Vec<crate::heal::resume::ReplacementTargetIdentity>>>,
listed_prefixes: Mutex<Vec<String>>,
truncate_without_token: Mutex<bool>,
include_object_dir_candidate: Mutex<bool>,
listed_buckets: Mutex<Option<Vec<String>>>,
bucket_heal_errors: Mutex<HashMap<String, VecDeque<&'static str>>>,
bucket_heal_calls: Mutex<Vec<String>>,
block_heal_object: Mutex<bool>,
resume_disk: Mutex<Option<DiskStore>>,
replacement_resume_disk: Mutex<Option<DiskStore>>,
usage_baseline: Mutex<Option<HealBucketUsageBaseline>>,
usage_baseline_error: Mutex<bool>,
}
#[test]
fn per_object_heal_types_are_classified_for_log_demotion() {
assert!(
HealType::Object {
bucket: "b".to_string(),
object: "o".to_string(),
version_id: None,
}
.is_per_object()
);
assert!(
HealType::Metadata {
bucket: "b".to_string(),
object: "o".to_string(),
}
.is_per_object()
);
assert!(
HealType::MRF {
meta_path: "p".to_string(),
}
.is_per_object()
);
assert!(
HealType::ECDecode {
bucket: "b".to_string(),
object: "o".to_string(),
version_id: None,
}
.is_per_object()
);
assert!(!HealType::Cluster.is_per_object());
assert!(!HealType::Bucket { bucket: "b".to_string() }.is_per_object());
assert!(
!HealType::Prefix {
bucket: "b".to_string(),
prefix: "p".to_string(),
}
.is_per_object()
);
assert!(
!HealType::ErasureSet {
buckets: Vec::new(),
set_disk_id: "s".to_string(),
}
.is_per_object()
);
}
#[test]
fn failure_log_sampling_caps_at_max_samples() {
let mut samples_logged = 0_u64;
for _ in 0..MAX_BUCKET_FAILURE_LOG_SAMPLES {
assert!(take_failure_log_sample(&mut samples_logged));
}
assert!(!take_failure_log_sample(&mut samples_logged));
assert!(!take_failure_log_sample(&mut samples_logged));
assert_eq!(samples_logged, MAX_BUCKET_FAILURE_LOG_SAMPLES);
}
#[tokio::test]
async fn execute_emits_heal_trace_task_state() {
let mut trace = subscribe_trace_events();
let storage = Arc::new(MockStorage::default());
let task = HealTask::from_request(
HealRequest::object("bucket-a".to_string(), "object-a".to_string(), Some("version-a".to_string())),
storage,
);
task.execute().await.expect("mock object heal should complete");
let started = recv_trace_task_state(&mut trace, &task.id, "started").await;
assert_eq!(started.kind, TraceKind::Heal);
assert_eq!(started.func, TraceFunc::HealTask);
assert_eq!(started.bucket.as_deref(), Some("bucket-a"));
assert_eq!(started.object.as_deref(), Some("object-a"));
assert_eq!(trace_attr_string(&started, "heal_type").as_deref(), Some("object"));
assert_eq!(trace_attr_string(&started, "source").as_deref(), Some("internal"));
assert_eq!(trace_attr_string(&started, "version_id").as_deref(), Some("version-a"));
let completed = recv_trace_task_state(&mut trace, &task.id, "completed").await;
assert_eq!(completed.kind, TraceKind::Heal);
assert_eq!(completed.func, TraceFunc::HealTask);
assert_eq!(trace_attr_string(&completed, "state").as_deref(), Some("completed"));
}
async fn recv_trace_task_state(trace: &mut TraceSubscription, task_id: &str, state: &str) -> TraceEvent {
for _ in 0..32 {
let event = tokio::time::timeout(Duration::from_secs(1), trace.recv())
.await
.expect("trace event should arrive")
.expect("trace bus should stay open");
if trace_attr_string(&event, "task_id").as_deref() == Some(task_id)
&& trace_attr_string(&event, "state").as_deref() == Some(state)
{
return (*event).clone();
}
}
panic!("expected trace state {state} for task {task_id}");
}
fn trace_attr_string(event: &TraceEvent, key: &str) -> Option<String> {
event.attrs.iter().find_map(|attr| {
if attr.key != key {
return None;
}
Some(match &attr.value {
TraceVal::Bool(value) => value.to_string(),
TraceVal::U64(value) => value.to_string(),
TraceVal::I64(value) => value.to_string(),
TraceVal::Str(value) => value.to_string(),
})
})
}
/// Build a latest, non-delete-marker heal list item with no version id.
fn heal_item(name: &str) -> HealListItem {
HealListItem {
name: name.to_string(),
version_id: None,
mod_time_unix_nanos: None,
lifecycle_object_info: None,
is_delete_marker: false,
}
}
fn replacement_identity(
endpoint: &str,
physical_device_id: &str,
filesystem_identity: &str,
) -> crate::heal::resume::ReplacementTargetIdentity {
crate::heal::resume::ReplacementTargetIdentity {
endpoint: endpoint.to_string(),
canonical_path: format!("/replacement/{endpoint}"),
physical_device_ids: vec![physical_device_id.to_string()],
filesystem_identity: filesystem_identity.to_string(),
}
}
enum MockHealObjectOutcome {
OkWithOtherError(&'static str),
ErrOther(&'static str),
RetryableReadQuorum,
RetryableSlowDown,
PermanentOther(&'static str),
}
#[derive(Clone, Copy)]
enum MockObjectExists {
Exists(bool),
TransientSkip(&'static str),
OtherError(&'static str),
}
#[test]
fn test_missing_object_dir_heal_result_matches_only_object_level_not_found() {
assert!(is_missing_object_dir_heal_result("x.rnd/", &Error::Disk(DiskError::FileNotFound)));
assert!(is_missing_object_dir_heal_result("x.rnd/", &Error::Disk(DiskError::FileVersionNotFound)));
assert!(is_missing_object_dir_heal_result("x.rnd/", &Error::Storage(EcstoreError::FileNotFound)));
assert!(is_missing_object_dir_heal_result(
"x.rnd/",
&Error::Other("File version not found".to_string())
));
assert!(!is_missing_object_dir_heal_result("x.rnd/", &Error::Other("Disk not found".to_string())));
assert!(!is_missing_object_dir_heal_result("x.rnd", &Error::Disk(DiskError::FileNotFound)));
}
#[async_trait::async_trait]
impl HealStorageAPI for MockStorage {
async fn get_object_meta(&self, _bucket: &str, _object: &str) -> Result<Option<HealObjectInfo>> {
Ok(None)
}
async fn ec_decode_rebuild(&self, _bucket: &str, _object: &str) -> Result<Vec<u8>> {
Ok(Vec::new())
}
async fn get_bucket_info(&self, bucket: &str) -> Result<Option<BucketInfo>> {
Ok(Some(BucketInfo {
name: bucket.to_string(),
..Default::default()
}))
}
async fn erasure_set_usage_baseline(&self, _buckets: &[String]) -> Result<Option<HealBucketUsageBaseline>> {
if *self.usage_baseline_error.lock().unwrap() {
return Err(Error::Other("usage baseline unavailable".to_string()));
}
Ok(*self.usage_baseline.lock().unwrap())
}
async fn list_buckets(&self) -> Result<Vec<BucketInfo>> {
let buckets = self
.listed_buckets
.lock()
.unwrap()
.clone()
.unwrap_or_else(|| vec!["bucket-a".to_string()]);
Ok(buckets
.into_iter()
.map(|name| BucketInfo {
name,
..Default::default()
})
.collect())
}
async fn object_exists(&self, _bucket: &str, object: &str) -> Result<bool> {
if let Some(result) = self.object_exists_by_name.lock().unwrap().get(object).copied() {
return match result {
MockObjectExists::Exists(exists) => Ok(exists),
MockObjectExists::TransientSkip(message) => Err(Error::transient_skip(message)),
MockObjectExists::OtherError(message) => Err(Error::other(message)),
};
}
Ok(self.object_exists.lock().unwrap().unwrap_or(true))
}
async fn heal_object(
&self,
bucket: &str,
object: &str,
version_id: Option<&str>,
opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
self.heal_object_calls.lock().unwrap().push(object.to_string());
self.heal_object_version_ids
.lock()
.unwrap()
.push(version_id.map(ToString::to_string));
self.object_heal_opts.lock().unwrap().push(*opts);
let block_heal_object = *self.block_heal_object.lock().unwrap();
if block_heal_object {
std::future::pending::<()>().await;
}
if let Some(outcome) = self
.heal_object_outcomes
.lock()
.unwrap()
.get_mut(object)
.and_then(VecDeque::pop_front)
{
return match outcome {
MockHealObjectOutcome::RetryableReadQuorum => Err(Error::Storage(EcstoreError::InsufficientReadQuorum(
bucket.to_string(),
object.to_string(),
))),
MockHealObjectOutcome::RetryableSlowDown => {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::SlowDown))))
}
MockHealObjectOutcome::PermanentOther(message) => Err(Error::other(message)),
MockHealObjectOutcome::OkWithOtherError(message) => {
Ok((HealResultItem::default(), Some(Error::other(message))))
}
MockHealObjectOutcome::ErrOther(message) => Err(Error::other(message)),
};
}
if let Some(outcome) = self.heal_object_outcome.lock().unwrap().take() {
return match outcome {
MockHealObjectOutcome::OkWithOtherError(message) => {
Ok((HealResultItem::default(), Some(Error::other(message))))
}
MockHealObjectOutcome::ErrOther(message) | MockHealObjectOutcome::PermanentOther(message) => {
Err(Error::other(message))
}
MockHealObjectOutcome::RetryableReadQuorum => Err(Error::Storage(EcstoreError::InsufficientReadQuorum(
bucket.to_string(),
object.to_string(),
))),
MockHealObjectOutcome::RetryableSlowDown => {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::SlowDown))))
}
};
}
if bucket == RUSTFS_META_BUCKET && object == format!("{BUCKET_META_PREFIX}/{DATA_USAGE_CACHE_NAME}") {
return Ok((
HealResultItem::default(),
Some(Error::other(
"Lock error: Lock acquisition timeout for resource '.rustfs.sys/buckets/.usage-cache.bin@latest' after 5s",
)),
));
}
if object == "object-dir/" {
return Ok((HealResultItem::default(), Some(Error::Disk(DiskError::FileNotFound))));
}
self.healed_objects.lock().unwrap().push(object.to_string());
Ok((
HealResultItem {
object_size: 1,
..Default::default()
},
None,
))
}
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
self.bucket_heal_calls.lock().unwrap().push(bucket.to_string());
self.bucket_heal_opts.lock().unwrap().push(*opts);
if let Some(message) = self
.bucket_heal_errors
.lock()
.unwrap()
.get_mut(bucket)
.and_then(VecDeque::pop_front)
{
return Err(Error::other(message));
}
Ok(HealResultItem::default())
}
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
*self.global_format_calls.lock().unwrap() += 1;
let no_heal_required = *self.format_no_heal_required.lock().unwrap();
if no_heal_required {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired))))
} else {
Ok((HealResultItem::default(), None))
}
}
async fn heal_replacement_format(
&self,
_dry_run: bool,
pool_index: usize,
set_index: usize,
targets: &[String],
) -> Result<(HealResultItem, Option<Error>)> {
self.replacement_format_calls
.lock()
.unwrap()
.push((pool_index, set_index, targets.to_vec()));
Ok((
HealResultItem {
after: Infos {
drives: targets
.iter()
.map(|endpoint| HealDriveInfo {
endpoint: endpoint.clone(),
state: "ok".to_string(),
..Default::default()
})
.collect(),
},
..Default::default()
},
None,
))
}
async fn replacement_targets_ready(&self, _targets: &[String]) -> Result<bool> {
Ok(*self.replacement_targets_ready.lock().unwrap())
}
async fn list_objects_for_heal_page(
&self,
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
self.listed_prefixes.lock().unwrap().push(prefix.to_string());
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 {
*listed = true;
let objects = if bucket == RUSTFS_META_BUCKET {
vec![
heal_item(&format!("{BUCKET_META_PREFIX}/{DATA_USAGE_CACHE_NAME}")),
heal_item(&format!("{BUCKET_META_PREFIX}/bucket-metadata.bin")),
]
} else if prefix == "logs/" {
vec![heal_item("logs/object-a"), heal_item("logs/object-b")]
} else if *self.include_object_dir_candidate.lock().unwrap() {
vec![heal_item("object-a"), heal_item("object-dir/"), heal_item("object-b")]
} else {
vec![heal_item("object-a"), heal_item("object-b")]
};
Ok((objects, None, false))
} else {
Ok((Vec::new(), None, false))
}
}
async fn get_disk_for_resume(&self, _set_disk_id: &str) -> Result<DiskStore> {
self.resume_disk
.lock()
.unwrap()
.clone()
.ok_or_else(|| Error::other("not implemented in tests"))
}
async fn get_disk_for_resume_excluding(&self, set_disk_id: &str, _excluded_targets: &[String]) -> Result<DiskStore> {
self.get_disk_for_resume(set_disk_id).await
}
async fn get_replacement_resume_disk(
&self,
_set_disk_id: &str,
_task_id: &str,
_excluded_targets: &[String],
) -> Result<crate::heal::storage::ReplacementResumeDisk> {
if let Some(disk) = self.replacement_resume_disk.lock().unwrap().clone() {
return Ok(crate::heal::storage::ReplacementResumeDisk::Existing(disk));
}
Ok(crate::heal::storage::ReplacementResumeDisk::Fresh)
}
async fn replacement_target_identities(
&self,
targets: &[String],
) -> Result<Vec<crate::heal::resume::ReplacementTargetIdentity>> {
if !*self.replacement_targets_ready.lock().unwrap() {
return Err(Error::other("replacement target is not ready"));
}
if let Some(identities) = self.replacement_target_identity_sequences.lock().unwrap().pop_front() {
return Ok(identities);
}
Ok(targets
.iter()
.map(|endpoint| crate::heal::resume::ReplacementTargetIdentity {
endpoint: endpoint.clone(),
canonical_path: format!("/replacement/{endpoint}"),
physical_device_ids: vec![endpoint.clone()],
filesystem_identity: format!("identity-{endpoint}"),
})
.collect())
}
}
#[tokio::test]
async fn scoped_object_heal_slowdown_is_not_treated_as_deleted() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(true)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::RetryableSlowDown)),
..Default::default()
});
let task = HealTask::from_request(
HealRequest::new(
HealType::Object {
bucket: "bucket".to_string(),
object: "object".to_string(),
version_id: None,
},
HealOptions {
pool_index: Some(0),
set_index: Some(1),
..Default::default()
},
HealPriority::Normal,
),
storage.clone(),
);
let err = task.execute().await.expect_err("SlowDown must fail the current heal attempt");
assert!(matches!(err, Error::Storage(EcstoreError::SlowDown)));
assert!(matches!(task.get_status().await, HealTaskStatus::Failed { .. }));
let opts = storage
.object_heal_opts
.lock()
.expect("heal options lock should be available");
assert_eq!(opts[0].pool, Some(0));
assert_eq!(opts[0].set, Some(1));
}
async fn make_resume_disk(temp: &TempDir) -> DiskStore {
let disk_path = temp.path().join("test_disk");
std::fs::create_dir_all(&disk_path).expect("test disk directory should be created");
let endpoint = Endpoint::try_from(disk_path.to_string_lossy().as_ref()).expect("test disk endpoint should be valid");
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("test disk should initialize");
let metadata_volume = disk.make_volume(RUSTFS_META_BUCKET).await;
assert!(
matches!(metadata_volume, Ok(()) | Err(DiskError::VolumeExists)),
"metadata volume should exist: {metadata_volume:?}"
);
disk
}
#[tokio::test]
async fn test_recursive_bucket_heal_visits_objects() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.heal_bucket("bucket-a")
.await
.expect("recursive bucket heal should succeed");
assert_eq!(
storage.healed_objects.lock().unwrap().as_slice(),
["object-a".to_string(), "object-b".to_string()]
);
let progress = task.get_progress().await;
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 2);
let result_items = task.get_result_items().await;
assert_eq!(result_items.len(), 3);
assert_eq!(result_items.iter().filter(|item| item.object_size == 1).count(), 2);
}
#[tokio::test]
async fn result_items_are_bounded_and_report_truncation() {
let storage = Arc::new(MockStorage::default());
let task = HealTask::from_request(HealRequest::bucket("bucket-a".to_string()), storage);
for _ in 0..=MAX_RETAINED_HEAL_RESULT_ITEMS {
task.record_result_item(HealResultItem::default()).await;
}
assert_eq!(task.get_result_items().await.len(), MAX_RETAINED_HEAL_RESULT_ITEMS);
assert!(task.result_items_truncated());
}
// HS-06 (backlog#1870): incremental result windows.
#[tokio::test]
async fn result_items_seq_is_monotonic_and_incremental_slices_work() {
let storage = Arc::new(MockStorage::default());
let task = HealTask::from_request(HealRequest::bucket("bucket-a".to_string()), storage);
for round in 0..5u64 {
let item = HealResultItem {
object_size: round as usize,
..Default::default()
};
task.record_result_item(item).await;
}
let full = task.get_result_items_since(None).await;
assert_eq!(full.items.len(), 5, "None keeps the full-snapshot semantics");
assert_eq!(full.next_seq, 6, "next_seq is one past the last assigned");
assert_eq!(full.min_seq, 1, "nothing was evicted yet");
assert!(!full.lagged);
// Incremental: only items newer than the cursor.
let incremental = task.get_result_items_since(Some(3)).await;
assert_eq!(
incremental.items.iter().map(|item| item.object_size).collect::<Vec<_>>(),
vec![3, 4],
"only sequences greater than the cursor are returned"
);
assert_eq!(incremental.next_seq, 6);
// A cursor at the head is not lagging.
assert!(!task.get_result_items_since(Some(0)).await.lagged);
}
#[tokio::test]
async fn result_items_window_slide_moves_min_seq_and_flags_lagging_cursors() {
let storage = Arc::new(MockStorage::default());
let task = HealTask::from_request(HealRequest::bucket("bucket-a".to_string()), storage);
// Fill the window completely, then push two more items: seq 1 and 2
// are evicted by the slide.
for _ in 0..(MAX_RETAINED_HEAL_RESULT_ITEMS + 2) {
task.record_result_item(HealResultItem::default()).await;
}
let full = task.get_result_items_since(None).await;
assert_eq!(full.items.len(), MAX_RETAINED_HEAL_RESULT_ITEMS);
assert_eq!(full.min_seq, 3, "each evicted head item moved the oldest-available cursor");
assert!(task.result_items_truncated());
// A client still polling from before the eviction is lagging.
let lagging = task.get_result_items_since(Some(0)).await;
assert!(lagging.lagged, "a cursor behind min_seq must be flagged");
assert_eq!(lagging.min_seq, 3, "the response tells the client where to restart");
// A cursor inside the window is fine.
assert!(!task.get_result_items_since(Some(3)).await.lagged);
// The lagging client restarts from min_seq and gets the full window.
let catch_up = task.get_result_items_since(Some(3)).await;
assert_eq!(catch_up.items.len(), MAX_RETAINED_HEAL_RESULT_ITEMS - 1);
assert!(!catch_up.lagged);
}
#[tokio::test]
async fn test_recursive_bucket_heal_skips_object_dir_candidates() {
let storage = Arc::new(MockStorage {
include_object_dir_candidate: Mutex::new(true),
..Default::default()
});
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.heal_bucket("bucket-a")
.await
.expect("recursive bucket heal should skip object-dir candidates");
assert_eq!(
storage.healed_objects.lock().unwrap().as_slice(),
["object-a".to_string(), "object-b".to_string()]
);
let progress = task.get_progress().await;
assert_eq!(progress.objects_scanned, 3);
assert_eq!(progress.objects_healed, 3);
assert_eq!(progress.objects_failed, 0);
}
#[tokio::test]
async fn test_recursive_bucket_heal_treats_missing_continuation_token_as_end() {
// 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
// of erroring or looping forever.
let storage = Arc::new(MockStorage {
truncate_without_token: Mutex::new(true),
..Default::default()
});
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.heal_bucket("bucket-a")
.await
.expect("truncated-without-token must terminate cleanly, not loop or error");
assert_eq!(
storage.healed_objects.lock().unwrap().as_slice(),
["object-a".to_string()],
"the returned page is healed exactly once and the scan ends"
);
}
#[tokio::test]
async fn test_cluster_heal_visits_bucket_objects() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Cluster,
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.execute().await.expect("cluster heal should visit bucket objects");
assert_eq!(
storage.healed_objects.lock().unwrap().as_slice(),
["object-a".to_string(), "object-b".to_string()]
);
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
}
#[tokio::test(start_paused = true)]
async fn test_recursive_bucket_heal_retries_only_retryable_objects() {
let storage = Arc::new(MockStorage::default());
storage
.heal_object_outcomes
.lock()
.unwrap()
.insert("object-a".to_string(), VecDeque::from([MockHealObjectOutcome::RetryableReadQuorum]));
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.heal_bucket("bucket-a")
.await
.expect("retryable object failure should be retried within the listing page");
assert_eq!(
storage.heal_object_calls.lock().unwrap().as_slice(),
["object-a".to_string(), "object-b".to_string(), "object-a".to_string()]
);
let progress = task.get_progress().await;
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 2);
assert_eq!(progress.objects_failed, 0);
}
#[tokio::test(start_paused = true)]
async fn test_recursive_bucket_heal_reports_typed_exhausted_and_permanent_failures() {
let storage = Arc::new(MockStorage::default());
storage.heal_object_outcomes.lock().unwrap().insert(
"object-a".to_string(),
VecDeque::from([
MockHealObjectOutcome::RetryableReadQuorum,
MockHealObjectOutcome::RetryableReadQuorum,
MockHealObjectOutcome::RetryableReadQuorum,
MockHealObjectOutcome::RetryableReadQuorum,
]),
);
storage.heal_object_outcomes.lock().unwrap().insert(
"object-b".to_string(),
VecDeque::from([MockHealObjectOutcome::PermanentOther("invalid metadata")]),
);
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
let err = task
.heal_bucket("bucket-a")
.await
.expect_err("exhausted retryable and permanent failures must fail the task");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
let failure = task
.take_batch_failure()
.await
.expect("batch failure details should be retained on the task");
assert_eq!(failure.failed, 2);
assert_eq!(failure.retryable, 1);
assert_eq!(failure.permanent, 1);
assert_eq!(failure.first_object, "object-b");
let calls = storage.heal_object_calls.lock().unwrap();
assert_eq!(calls.iter().filter(|object| object.as_str() == "object-a").count(), 4);
assert_eq!(calls.iter().filter(|object| object.as_str() == "object-b").count(), 1);
}
#[tokio::test]
async fn test_cluster_heal_continues_after_bucket_failure() {
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 request = HealRequest::new(
HealType::Cluster,
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
let err = task.execute().await.expect_err("cluster task must report the failed bucket");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
let failure = task
.take_batch_failure()
.await
.expect("cluster failure details should be retained on the task");
assert_eq!(failure.failed, 1);
assert_eq!(
storage.bucket_heal_calls.lock().unwrap().as_slice(),
["bucket-a".to_string(), "bucket-b".to_string()]
);
}
#[tokio::test(start_paused = true)]
async fn test_cluster_heal_retries_only_recoverable_bucket() {
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(["lock acquisition timeout"]),
)])),
..Default::default()
});
let request = HealRequest::new(
HealType::Cluster,
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("cluster task should retry a recoverable bucket failure");
assert_eq!(
storage.bucket_heal_calls.lock().unwrap().as_slice(),
["bucket-a".to_string(), "bucket-a".to_string(), "bucket-b".to_string()]
);
}
#[tokio::test]
async fn test_recursive_bucket_heal_does_not_remove_bucket_metadata() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Bucket {
bucket: "bucket-a".to_string(),
},
HealOptions {
recursive: true,
remove_corrupted: true,
recreate_missing: true,
scan_mode: HealScanMode::Deep,
no_lock: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.heal_bucket("bucket-a")
.await
.expect("recursive bucket heal should succeed");
let bucket_opts = storage.bucket_heal_opts.lock().unwrap();
assert_eq!(bucket_opts.len(), 1);
assert!(!bucket_opts[0].remove);
assert!(bucket_opts[0].recreate);
assert_eq!(bucket_opts[0].scan_mode, HealScanMode::Deep);
assert!(bucket_opts[0].no_lock);
let object_opts = storage.object_heal_opts.lock().unwrap();
assert_eq!(object_opts.len(), 2);
assert!(object_opts.iter().all(|opts| opts.remove));
assert!(object_opts.iter().all(|opts| opts.recreate));
assert!(object_opts.iter().all(|opts| opts.scan_mode == HealScanMode::Deep));
assert!(object_opts.iter().all(|opts| opts.no_lock));
}
#[tokio::test]
async fn test_prefix_heal_lists_and_repairs_objects_under_prefix() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Prefix {
bucket: "bucket-a".to_string(),
prefix: "logs/".to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("prefix heal should scan and repair objects under the prefix");
assert_eq!(storage.listed_prefixes.lock().unwrap().as_slice(), ["logs/".to_string()]);
assert_eq!(
storage.healed_objects.lock().unwrap().as_slice(),
["logs/object-a".to_string(), "logs/object-b".to_string()]
);
}
#[tokio::test]
async fn test_data_usage_cache_lock_timeout_does_not_fail_object_heal() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Object {
bucket: RUSTFS_META_BUCKET.to_string(),
object: format!("{BUCKET_META_PREFIX}/{DATA_USAGE_CACHE_NAME}"),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.execute()
.await
.expect("data usage cache lock timeout should be skipped during heal");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
}
#[tokio::test]
async fn test_data_usage_cache_lock_timeout_does_not_fail_recursive_bucket_heal() {
let storage = Arc::new(MockStorage::default());
let request = HealRequest::new(
HealType::Bucket {
bucket: RUSTFS_META_BUCKET.to_string(),
},
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.execute()
.await
.expect("recursive bucket heal should skip transient data usage cache lock timeouts");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
let progress = task.get_progress().await;
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_failed, 0);
}
#[tokio::test]
async fn test_heal_recreate_scanner_synthetic_object_dir_skips_ok_not_found_error() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::OkWithOtherError("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("scanner synthetic object-dir missing result should be skipped");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]);
}
#[tokio::test]
async fn test_heal_scanner_missing_object_dir_canonicalizes_existing_plain_object() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: Some("version-a".to_string()),
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("scanner object-dir candidate should heal the existing plain object");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd".to_string()]);
assert_eq!(
storage.heal_object_version_ids.lock().unwrap().as_slice(),
[Some("version-a".to_string())]
);
assert_eq!(storage.healed_objects.lock().unwrap().as_slice(), ["x.rnd".to_string()]);
}
#[tokio::test]
async fn test_heal_scanner_existing_trailing_slash_object_is_not_canonicalized() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(true));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("existing trailing-slash object should keep its exact key");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]);
}
#[tokio::test]
async fn test_heal_admin_missing_object_dir_does_not_canonicalize_plain_object() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Admin;
let task = HealTask::from_request(request, storage.clone());
let err = task
.execute()
.await
.expect_err("admin object-dir request must not be canonicalized");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]);
}
#[tokio::test]
async fn test_heal_scanner_canonicalizes_only_one_trailing_slash() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd//".to_string(), MockObjectExists::Exists(false));
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(true));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd//".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("scanner canonicalization should remove only one trailing slash");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]);
}
#[tokio::test]
async fn test_heal_scanner_trimmed_object_exists_error_is_not_recreated() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::OtherError("backend unavailable"));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
let err = task
.execute()
.await
.expect_err("trimmed object_exists error must not be treated as missing");
assert!(matches!(err, Error::Other(_)));
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
}
#[tokio::test]
async fn test_heal_scanner_trimmed_object_exists_transient_skip_is_not_recreated() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false));
object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::TransientSkip("backend busy"));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("trimmed object_exists transient skip should complete without recreate");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
}
#[tokio::test]
async fn test_heal_scanner_empty_trimmed_object_keeps_existing_skip_behavior() {
let mut object_exists_by_name = HashMap::new();
object_exists_by_name.insert("/".to_string(), MockObjectExists::Exists(false));
let storage = Arc::new(MockStorage {
object_exists_by_name: Mutex::new(object_exists_by_name),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::OkWithOtherError("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("empty canonical object must keep existing scanner skip behavior");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["/".to_string()]);
}
#[tokio::test]
async fn test_heal_recreate_scanner_synthetic_object_dir_skips_err_not_found() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage);
task.execute()
.await
.expect("scanner synthetic object-dir missing error should be skipped");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
}
#[tokio::test]
async fn test_heal_recreate_scanner_non_dir_not_found_fails() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage);
let err = task
.execute()
.await
.expect_err("scanner non-dir missing object should still fail recreate");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
assert!(matches!(task.get_status().await, HealTaskStatus::Failed { .. }));
}
#[tokio::test]
async fn test_heal_scanner_missing_object_without_recreate_probes_storage() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: false,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("scanner missing object should be checked by storage");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd".to_string()]);
assert!(!storage.object_heal_opts.lock().unwrap()[0].recreate);
}
#[tokio::test]
async fn test_heal_scanner_missing_object_without_recreate_treats_not_found_as_stale() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: false,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("scanner confirmed-not-found object should be treated as stale");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd".to_string()]);
}
#[tokio::test]
async fn test_heal_recreate_scanner_synthetic_object_dir_disk_not_found_fails() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("Disk not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage);
let err = task
.execute()
.await
.expect_err("scanner synthetic object-dir disk-not-found should not be skipped");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
assert!(matches!(task.get_status().await, HealTaskStatus::Failed { .. }));
}
#[tokio::test]
async fn test_heal_recreate_admin_synthetic_object_dir_not_found_fails() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Admin;
let task = HealTask::from_request(request, storage);
let err = task
.execute()
.await
.expect_err("admin object-dir not-found recreate should not be skipped");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
assert!(matches!(task.get_status().await, HealTaskStatus::Failed { .. }));
}
#[tokio::test]
async fn test_heal_recreate_existing_trailing_slash_object_records_normal_result() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(true)),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd/".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Scanner;
let task = HealTask::from_request(request, storage);
task.execute()
.await
.expect("existing trailing-slash object should follow normal heal path");
assert!(matches!(task.get_status().await, HealTaskStatus::Completed));
let result_items = task.get_result_items().await;
assert_eq!(result_items.len(), 1);
assert_eq!(result_items[0].object_size, 1);
}
#[tokio::test]
async fn test_heal_failure_with_remove_corrupted_propagates_remove_flag() {
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(true)),
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::OkWithOtherError(
"can not reconstruct data: not enough available shards (need 12, have 11)",
))),
..Default::default()
});
let request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "x.rnd".to_string(),
version_id: None,
},
HealOptions {
remove_corrupted: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
let err = task.execute().await.expect_err("heal failure should still be reported");
assert!(matches!(err, Error::TaskExecutionFailed { .. }));
assert!(storage.object_heal_opts.lock().unwrap()[0].remove);
}
#[tokio::test]
async fn test_erasure_set_heal_continues_after_format_no_heal_required() {
let storage = Arc::new(MockStorage::default());
*storage.format_no_heal_required.lock().unwrap() = true;
let request = HealRequest::new(
HealType::ErasureSet {
buckets: Vec::new(),
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
let err = task
.heal_erasure_set(Vec::new(), "pool_0_set_0".to_string())
.await
.expect_err("test mock should fail after format when resolving resume disk");
assert!(
err.to_string().contains("not implemented in tests"),
"erasure-set heal should continue past NoHealRequired format result, got: {err}"
);
}
#[tokio::test]
async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
bucket_heal_errors: Mutex::new(HashMap::from([(
"bucket-a".to_string(),
VecDeque::from(["injected bucket prepass failure"]),
)])),
resume_disk: Mutex::new(Some(disk)),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage.clone());
let error = task
.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect_err("bucket prepass failure must stop the erasure-set heal");
assert!(error.to_string().contains("injected bucket prepass failure"));
assert_eq!(storage.bucket_heal_calls.lock().unwrap().as_slice(), ["bucket-a".to_string()]);
assert!(storage.object_heal_opts.lock().unwrap().is_empty());
}
#[tokio::test]
async fn erasure_set_heal_applies_usage_baseline_to_progress() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline: Mutex::new(Some(HealBucketUsageBaseline {
objects_count: 10,
bytes: 8,
})),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("erasure set heal should complete");
let progress = task.get_progress().await;
assert_eq!(progress.objects_total_count, 10);
assert_eq!(progress.objects_total_size, 8);
assert_eq!(progress.bytes_processed, 2);
assert!((progress.progress_percentage - 25.0).abs() < 0.001);
}
#[tokio::test]
async fn erasure_set_heal_ignores_usage_baseline_errors() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline_error: Mutex::new(true),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("usage baseline failures should not fail erasure set heal");
let progress = task.get_progress().await;
assert_eq!(progress.objects_total_count, 0);
assert_eq!(progress.objects_total_size, 0);
}
#[tokio::test]
async fn resumable_erasure_set_execution_is_cancelled_while_object_heal_is_pending() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
block_heal_object: Mutex::new(true),
resume_disk: Mutex::new(Some(disk)),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
no_lock: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = Arc::new(HealTask::from_request(request, storage.clone()));
let execution = tokio::spawn({
let task = task.clone();
async move { task.execute().await }
});
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if !storage.object_heal_opts.lock().unwrap().is_empty() {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("resumable object heal should start");
task.cancel().await.expect("task cancellation should succeed");
let result = tokio::time::timeout(Duration::from_secs(1), execution)
.await
.expect("cancellation should interrupt the pending resumable heal")
.expect("task execution should join");
assert!(matches!(result, Err(Error::TaskCancelled)));
assert!(storage.bucket_heal_opts.lock().unwrap().iter().all(|opts| opts.no_lock));
assert!(storage.object_heal_opts.lock().unwrap().iter().all(|opts| opts.no_lock));
}
}