From 7eddd1cf833bb46f45f012e8c8dfac5d9ed33f39 Mon Sep 17 00:00:00 2001 From: houseme Date: Fri, 28 Aug 2026 20:09:45 +0800 Subject: [PATCH] fix(heal): skip dangling delete grace failures (#6799) Co-authored-by: heihutu --- crates/ecstore/src/disk/error.rs | 37 +++++++++ crates/ecstore/src/error/mod.rs | 4 + .../src/set_disk/core/io_primitives.rs | 23 +++++- crates/ecstore/src/set_disk/ops/heal.rs | 15 +++- crates/heal/src/error.rs | 20 +++++ crates/heal/src/heal/erasure_healer.rs | 19 ++++- crates/heal/src/heal/task.rs | 28 +++++++ crates/heal/src/heal/task/heal_bucket.rs | 16 +++- crates/heal/src/heal/task/heal_object.rs | 8 ++ crates/heal/src/heal/task/tests.rs | 75 +++++++++++++++++++ 10 files changed, 238 insertions(+), 7 deletions(-) diff --git a/crates/ecstore/src/disk/error.rs b/crates/ecstore/src/disk/error.rs index dff44d78b..5f1d1175b 100644 --- a/crates/ecstore/src/disk/error.rs +++ b/crates/ecstore/src/disk/error.rs @@ -22,6 +22,7 @@ pub type Error = DiskError; pub type Result = core::result::Result; const METACACHE_OUTPUT_STREAM_CLOSED: &str = "metacache output stream closed"; +pub(crate) const HEAL_DANGLING_DELETE_GRACE_MESSAGE: &str = "dangling object deletion deferred by heal grace window"; /// Marker carried by a shard-read `io::Error` when the underlying reader can /// no longer be realigned after a fresh remote open failed. The marker is @@ -33,6 +34,12 @@ pub(crate) struct TerminalReadError { source: DiskError, } +#[derive(Debug)] +struct DanglingDeleteGraceError { + retry_after_secs: i64, + grace_secs: i64, +} + // DiskError == StorageErr #[derive(Debug, thiserror::Error)] pub enum DiskError { @@ -200,6 +207,18 @@ impl StdError for TerminalReadError { } } +impl std::fmt::Display for DanglingDeleteGraceError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "{HEAL_DANGLING_DELETE_GRACE_MESSAGE}; retry_after_secs={}; grace_secs={}", + self.retry_after_secs, self.grace_secs + ) + } +} + +impl StdError for DanglingDeleteGraceError {} + fn classify_internode_missing_error(error: &InternodeHttpError) -> Option { if error.is_remote_file_not_found() { return Some(DiskError::FileNotFound); @@ -253,6 +272,24 @@ impl DiskError { DiskError::Io(std::io::Error::other(error)) } + pub(crate) fn dangling_delete_grace(retry_after_secs: i64, grace_secs: i64) -> Self { + DiskError::other(DanglingDeleteGraceError { + retry_after_secs, + grace_secs, + }) + } + + pub fn is_dangling_delete_grace(&self) -> bool { + matches!(self, DiskError::Io(io_error) if Self::io_error_is_dangling_delete_grace(io_error)) + } + + pub fn io_error_is_dangling_delete_grace(io_error: &io::Error) -> bool { + io_error + .get_ref() + .is_some_and(|source| source.downcast_ref::().is_some()) + || io_error.to_string().contains(HEAL_DANGLING_DELETE_GRACE_MESSAGE) + } + pub(crate) fn metacache_output_stream_closed() -> Self { DiskError::Io(std::io::Error::new(std::io::ErrorKind::BrokenPipe, METACACHE_OUTPUT_STREAM_CLOSED)) } diff --git a/crates/ecstore/src/error/mod.rs b/crates/ecstore/src/error/mod.rs index 706f6e926..9895018f7 100644 --- a/crates/ecstore/src/error/mod.rs +++ b/crates/ecstore/src/error/mod.rs @@ -277,6 +277,10 @@ impl StorageError { | StorageError::NamespaceLockQuorumUnavailable { .. } ) } + + pub fn is_dangling_delete_grace(&self) -> bool { + matches!(self, StorageError::Io(io_error) if DiskError::io_error_is_dangling_delete_grace(io_error)) + } } impl From for StorageError { diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 30b221256..03d29ff86 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -5696,6 +5696,8 @@ impl SetDisks { { let grace = dangling_delete_grace(); if !grace.is_zero() && OffsetDateTime::now_utc() - mod_time < grace { + let elapsed = OffsetDateTime::now_utc() - mod_time; + let retry_after_secs = grace.saturating_sub(elapsed).whole_seconds().max(0); info!( bucket = bucket, object = object, @@ -5703,7 +5705,7 @@ impl SetDisks { grace_secs = grace.whole_seconds(), "skipping dangling-object deletion within grace window" ); - return Err(DiskError::ErasureReadQuorum); + return Err(DiskError::dangling_delete_grace(retry_after_secs, grace.whole_seconds())); } } @@ -6799,6 +6801,7 @@ pub(in crate::set_disk) mod rename_fanout_barrier { #[cfg(test)] mod tests { + use crate::disk::error::HEAL_DANGLING_DELETE_GRACE_MESSAGE; use crate::disk::local::{DurabilityMode, durability_mode_override}; use super::*; @@ -10746,7 +10749,13 @@ mod tests { let object = "object"; let (_dir, disk) = read_multiple_test_disk(bucket, &[]).await; let set = io_primitives_test_set(vec![Some(disk.clone()), None, None], 1).await; - let mut fi = metadata_test_fileinfo(object); + let mut fi = FileInfo::new(object, 2, 1); + fi.volume = bucket.to_string(); + fi.name = object.to_string(); + fi.size = 1; + fi.erasure.index = 1; + fi.metadata.insert("etag".to_string(), "etag-1".to_string()); + fi.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None); fi.mod_time = Some(OffsetDateTime::now_utc()); disk.write_metadata(bucket, bucket, object, fi.clone()) .await @@ -10764,7 +10773,15 @@ mod tests { .await .expect_err("recent dangling metadata must stay protected by grace"); - assert_eq!(err, DiskError::ErasureReadQuorum); + let message = err.to_string(); + assert!( + message.contains(HEAL_DANGLING_DELETE_GRACE_MESSAGE), + "grace-protected dangling cleanup must explain the deferred delete: {message}" + ); + assert!( + message.contains("retry_after_secs="), + "grace-protected dangling cleanup must include retry timing: {message}" + ); disk.read_all(bucket, &path_join_buf(&[object, STORAGE_FORMAT_FILE])) .await .expect("metadata should remain during dangling grace"); diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 171e83fc2..6c5c4f08e 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -3543,7 +3543,20 @@ mod heal_result_report_tests { .await .expect("grace-protected dangling metadata should return a typed heal result"); - assert_eq!(error, Some(DiskError::ErasureReadQuorum)); + let error = error.expect("grace-protected dangling metadata should be reported as deferred"); + assert!( + error.is_dangling_delete_grace(), + "grace-protected dangling metadata should keep a typed deferred-cleanup marker: {error}" + ); + let message = error.to_string(); + assert!( + message.contains("dangling object deletion deferred by heal grace window"), + "grace-protected dangling metadata should explain that cleanup was deferred: {message}" + ); + assert!( + message.contains("retry_after_secs="), + "grace-protected dangling metadata should include retry timing: {message}" + ); assert!( temp_dirs[0] .path() diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 88b43e563..336ba00ad 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -16,6 +16,8 @@ use thiserror::Error; use super::heal::{DiskError, EcstoreError}; +const HEAL_DANGLING_DELETE_GRACE_MESSAGE: &str = "dangling object deletion deferred by heal grace window"; + /// Custom error type for heal operations /// This enum defines various error variants that can occur during /// the execution of heal-related tasks, such as I/O errors, storage errors, @@ -98,6 +100,9 @@ impl Error { // them. Error::Storage(EcstoreError::Lock(lock_err)) => !lock_err.is_fatal(), Error::Storage(err) => { + if err.is_dangling_delete_grace() { + return true; + } err.is_quorum_error() || matches!( err, @@ -110,6 +115,9 @@ impl Error { || is_recoverable_heal_error_message(&err.to_string()) } Error::Disk(err) => { + if err.is_dangling_delete_grace() { + return true; + } matches!( err, DiskError::DiskNotFound @@ -127,6 +135,18 @@ impl Error { _ => false, } } + + pub(crate) fn is_dangling_delete_grace(&self) -> bool { + match self { + Error::Storage(err) => err.is_dangling_delete_grace(), + Error::Disk(err) => err.is_dangling_delete_grace(), + Error::Io(err) => DiskError::io_error_is_dangling_delete_grace(err), + Error::TaskExecutionFailed { message } | Error::Other(message) => { + message.contains(HEAL_DANGLING_DELETE_GRACE_MESSAGE) + } + _ => false, + } + } } /// Documented substring fallback for errors that reach heal with their typed diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 4d1e0b7c5..a6c288dfd 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -159,9 +159,13 @@ impl ErasureSetHealer { /// Classify an error returned by [`HealStorageAPI::heal_object`]. /// - /// Both the inner `Ok((_, Some(err)))` and the outer `Err(err)` produced by - /// `heal_object` wrap `Error::Storage(StorageError)`, so match on that. + /// Most heal object failures wrap `Error::Storage(StorageError)`, while + /// compatibility markers can also arrive through Disk/Io/task wrappers. fn classify_heal_object_error(err: &Error) -> HealObjectOutcome { + if err.is_dangling_delete_grace() { + return HealObjectOutcome::Transient; + } + let Error::Storage(se) = err else { return HealObjectOutcome::Failed; }; @@ -1459,6 +1463,7 @@ mod tests { // genuine object absence, or transient failures get recorded as "healed" and // permanently skipped. use super::{EcstoreError, Error, HealObjectOutcome}; + use crate::heal::DiskError; fn classify(err: EcstoreError) -> HealObjectOutcome { ErasureSetHealer::classify_heal_object_error(&Error::Storage(err)) @@ -1479,6 +1484,16 @@ mod tests { )); } + #[test] + fn dangling_delete_grace_is_transient() { + assert!(matches!( + ErasureSetHealer::classify_heal_object_error(&Error::Disk(DiskError::other( + "dangling object deletion deferred by heal grace window; retry_after_secs=3599; grace_secs=3600" + ))), + HealObjectOutcome::Transient + )); + } + #[test] fn genuine_object_absence_is_absent() { assert!(matches!(classify(EcstoreError::FileNotFound), HealObjectOutcome::Absent)); diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index ffd2ac618..b39d92162 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -682,6 +682,10 @@ impl HealTask { Self::is_data_usage_cache_object(bucket, object) && Self::is_transient_lock_or_timeout_error(err) } + fn is_dangling_delete_grace_error(err: &Error) -> bool { + err.is_dangling_delete_grace() + } + fn is_no_heal_required_error(err: &Error) -> bool { match err { Error::Storage(EcstoreError::NoHealRequired) | Error::Disk(DiskError::NoHealRequired) => true, @@ -746,6 +750,30 @@ impl HealTask { true } + async fn skip_dangling_delete_grace_error(&self, bucket: &str, object: &str, err: &Error) -> bool { + if !Self::is_dangling_delete_grace_error(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 = "dangling_delete_grace_skip", + error = %err, + "Heal object dangling cleanup deferred by grace window" + ); + let mut progress = self.progress.write().await; + progress.set_current_object(Some(format!("skipped: {bucket}/{object}"))); + progress.update_object_progress(1, 0, 0, 1, 0); + progress.update_stage(3, 3); + 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; diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index 6237d246f..27242c2d0 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -359,7 +359,21 @@ impl HealTask { }; if let Some(err) = error { - if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) { + if Self::is_dangling_delete_grace_error(&err) { + telemetry_unknown |= !increment_counter(&mut skipped); + 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 = "dangling_delete_grace_skip", + error = %err, + "Heal bucket object dangling cleanup deferred by grace window" + ); + } else if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) { telemetry_unknown |= !increment_counter(&mut skipped); warn!( target: "rustfs::heal::task", diff --git a/crates/heal/src/heal/task/heal_object.rs b/crates/heal/src/heal/task/heal_object.rs index 037bc3d41..06c46de94 100644 --- a/crates/heal/src/heal/task/heal_object.rs +++ b/crates/heal/src/heal/task/heal_object.rs @@ -169,6 +169,10 @@ impl HealTask { match heal_result { Ok((result, error)) => { if let Some(e) = error { + if self.skip_dangling_delete_grace_error(bucket, object, &e).await { + return Ok(()); + } + if self.skip_data_usage_cache_heal_error(bucket, object, &e).await { return Ok(()); } @@ -257,6 +261,10 @@ impl HealTask { Err(Error::TaskCancelled) => Err(Error::TaskCancelled), Err(Error::TaskTimeout) => Err(Error::TaskTimeout), Err(e) => { + if self.skip_dangling_delete_grace_error(bucket, object, &e).await { + return Ok(()); + } + if self.skip_data_usage_cache_heal_error(bucket, object, &e).await { return Ok(()); } diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index bd2b23606..d4c5309c4 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -705,6 +705,7 @@ fn replacement_identity( enum MockHealObjectOutcome { OkWithOtherError(&'static str), ErrOther(&'static str), + DanglingGraceDeferred, RetryableReadQuorum, RetryableSlowDown, PermanentOther(&'static str), @@ -806,6 +807,12 @@ impl HealStorageAPI for MockStorage { .and_then(VecDeque::pop_front) { return match outcome { + MockHealObjectOutcome::DanglingGraceDeferred => Ok(( + HealResultItem::default(), + Some(Error::Disk(DiskError::other( + "dangling object deletion deferred by heal grace window; retry_after_secs=3599; grace_secs=3600", + ))), + )), MockHealObjectOutcome::RetryableReadQuorum => Err(Error::Storage(EcstoreError::InsufficientReadQuorum( bucket.to_string(), object.to_string(), @@ -820,6 +827,12 @@ impl HealStorageAPI for MockStorage { } if let Some(outcome) = self.heal_object_outcome.lock().unwrap().take() { return match outcome { + MockHealObjectOutcome::DanglingGraceDeferred => Ok(( + HealResultItem::default(), + Some(Error::Disk(DiskError::other( + "dangling object deletion deferred by heal grace window; retry_after_secs=3599; grace_secs=3600", + ))), + )), MockHealObjectOutcome::OkWithOtherError(message) => Ok((HealResultItem::default(), Some(Error::other(message)))), MockHealObjectOutcome::ErrOther(message) | MockHealObjectOutcome::PermanentOther(message) => { Err(Error::other(message)) @@ -1310,6 +1323,31 @@ async fn test_cluster_heal_visits_bucket_objects() { assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); } +#[tokio::test] +async fn object_heal_skips_dangling_delete_grace_without_failing_task() { + let storage = Arc::new(MockStorage { + heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::DanglingGraceDeferred)), + ..Default::default() + }); + let task = HealTask::from_request( + HealRequest::object("bucket-a".to_string(), "recent.txt".to_string(), None), + storage.clone(), + ); + + task.execute() + .await + .expect("grace-protected dangling cleanup should be reported as a skipped object"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert!(storage.healed_objects.lock().unwrap().is_empty()); + let progress = task.get_progress().await; + assert_eq!(progress.current_object.as_deref(), Some("skipped: bucket-a/recent.txt")); + assert_eq!(progress.objects_scanned, 1); + assert_eq!(progress.objects_healed, 0); + assert_eq!(progress.objects_failed, 0); + assert_eq!(progress.skipped_objects, 1); +} + #[tokio::test(start_paused = true)] async fn test_recursive_bucket_heal_retries_only_retryable_objects() { let storage = Arc::new(MockStorage::default()); @@ -1345,6 +1383,43 @@ async fn test_recursive_bucket_heal_retries_only_retryable_objects() { assert_eq!(progress.objects_failed, 0); } +#[tokio::test(start_paused = true)] +async fn recursive_bucket_heal_skips_dangling_delete_grace_without_batch_failure() { + let storage = Arc::new(MockStorage::default()); + storage + .heal_object_outcomes + .lock() + .unwrap() + .insert("object-a".to_string(), VecDeque::from([MockHealObjectOutcome::DanglingGraceDeferred])); + 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("grace-protected dangling cleanup should not fail the bucket heal batch"); + + assert_eq!( + storage.heal_object_calls.lock().unwrap().as_slice(), + ["object-a".to_string(), "object-b".to_string()] + ); + assert_eq!(storage.healed_objects.lock().unwrap().as_slice(), ["object-b".to_string()]); + let progress = task.get_progress().await; + assert_eq!(progress.objects_scanned, 2); + assert_eq!(progress.objects_healed, 1); + assert_eq!(progress.objects_failed, 0); + assert_eq!(progress.skipped_objects, 1); +} + #[tokio::test(start_paused = true)] async fn test_recursive_bucket_heal_reports_typed_exhausted_and_permanent_failures() { let storage = Arc::new(MockStorage::default());