mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-21 03:46:37 +00:00
refactor(heal): split task.rs per heal kind (#6293)
This commit is contained in:
+6
-3875
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,450 @@
|
||||
// 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.
|
||||
/// bucket/cluster/prefix heal: the recursive bucket-objects sweep and the erasure-set usage baseline
|
||||
use super::*;
|
||||
|
||||
impl HealTask {
|
||||
pub(super) 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}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) 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(())
|
||||
}
|
||||
|
||||
pub(super) 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(())
|
||||
}
|
||||
|
||||
pub(super) 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(())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,506 @@
|
||||
// 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.
|
||||
/// erasure-set heal: drives the ErasureSetHealer across the set's buckets
|
||||
use super::*;
|
||||
|
||||
impl HealTask {
|
||||
pub(super) 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::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::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::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}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,342 @@
|
||||
// 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.
|
||||
/// metadata and erasure-decode heal for a single object version
|
||||
use super::*;
|
||||
|
||||
impl HealTask {
|
||||
pub(super) 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}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) 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}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,447 @@
|
||||
// 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.
|
||||
/// object-level heal: metadata dir canonicalization and missing-object recreation
|
||||
use super::*;
|
||||
|
||||
impl HealTask {
|
||||
// specific heal implementation method
|
||||
#[tracing::instrument(skip(self), fields(bucket = %bucket, object = %object, version_id = ?version_id))]
|
||||
#[hotpath::measure]
|
||||
pub(super) 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}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user