From 55cbfe256b39652a0c2418478ac971be309a1678 Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 25 Jun 2026 14:52:24 +0800 Subject: [PATCH] fix(heal): skip scanner object-dir recreate misses (#3839) --- crates/heal/src/heal/task.rs | 266 ++++++++++++++++++++++++++++++++++- 1 file changed, 263 insertions(+), 3 deletions(-) diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 77033d94c..e96171b75 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -13,7 +13,7 @@ // limitations under the License. use crate::heal::{ - DiskError, ErasureSetHealer, + DiskError, EcstoreError, ErasureSetHealer, progress::HealProgress, storage::{HealStorageAPI, next_heal_listing_token}, }; @@ -97,8 +97,17 @@ impl HealType { } } +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) && matches!(err, Error::Disk(DiskError::FileNotFound | DiskError::FileVersionNotFound)) + object.ends_with(SLASH_SEPARATOR) && is_object_level_not_found_error(err) } /// Heal priority @@ -459,6 +468,30 @@ impl HealTask { 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))] pub async fn execute(&self) -> Result<()> { // update status and timestamps atomically to avoid race conditions @@ -969,6 +1002,10 @@ impl HealTask { { 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, @@ -1010,6 +1047,10 @@ impl HealTask { 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, @@ -2136,14 +2177,35 @@ mod tests { struct MockStorage { listed: Mutex, healed_objects: Mutex>, + heal_object_calls: Mutex>, bucket_heal_opts: Mutex>, object_heal_opts: Mutex>, + object_exists: Mutex>, + heal_object_outcome: Mutex>, format_no_heal_required: Mutex, listed_prefixes: Mutex>, truncate_without_token: Mutex, include_object_dir_candidate: Mutex, } + enum MockHealObjectOutcome { + OkWithOtherError(&'static str), + ErrOther(&'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> { @@ -2197,7 +2259,7 @@ mod tests { } async fn object_exists(&self, _bucket: &str, _object: &str) -> Result { - Ok(true) + Ok(self.object_exists.lock().unwrap().unwrap_or(true)) } async fn get_object_size(&self, _bucket: &str, _object: &str) -> Result> { @@ -2215,6 +2277,15 @@ mod tests { _version_id: Option<&str>, opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { + self.heal_object_calls.lock().unwrap().push(object.to_string()); + 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) => Err(Error::other(message)), + }; + } if bucket == RUSTFS_META_BUCKET && object == format!("{BUCKET_META_PREFIX}/{DATA_USAGE_CACHE_NAME}") { return Ok(( HealResultItem::default(), @@ -2528,6 +2599,195 @@ mod tests { 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_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_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_erasure_set_heal_continues_after_format_no_heal_required() { let storage = Arc::new(MockStorage::default());