From 2e3c777307f6fa6acfbc6a85ba986799c5ee9300 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 11 Jun 2026 19:16:32 +0800 Subject: [PATCH] fix(heal): skip transient usage cache heal errors (#3359) --- crates/heal/src/heal/task.rs | 135 ++++++++++++++++++++++++++++++++--- 1 file changed, 127 insertions(+), 8 deletions(-) diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 46ece118d..727529d61 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -16,6 +16,10 @@ use crate::heal::{ErasureSetHealer, progress::HealProgress, storage::HealStorage use crate::{Error, Result}; use metrics::{counter, histogram}; use rustfs_common::heal_channel::{HealOpts, HealScanMode}; +use rustfs_ecstore::{ + data_usage::DATA_USAGE_CACHE_NAME, + disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}, +}; use rustfs_madmin::heal_commands::HealResultItem; use serde::{Deserialize, Serialize}; use std::{ @@ -320,6 +324,37 @@ impl HealTask { 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) + } + + 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!("Skipping data usage cache heal for {}/{} due to transient error: {}", bucket, object, err); + let mut progress = self.progress.write().await; + progress.update_progress(3, 3, 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 @@ -486,6 +521,10 @@ impl HealTask { 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(()); + } + // Check if this is a "File not found" error during delete operations let error_msg = format!("{e}"); if error_msg.contains("File not found") || error_msg.contains("not found") { @@ -542,6 +581,10 @@ impl HealTask { 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(()); + } + // Check if this is a "File not found" error during delete operations let error_msg = format!("{e}"); if error_msg.contains("File not found") || error_msg.contains("not found") { @@ -758,12 +801,20 @@ impl HealTask { self.record_result_item(result).await; } Ok((_, Some(err))) => { - failed += 1; - warn!("Failed to heal object {}/{}: {}", bucket, object, err); + if Self::should_skip_data_usage_cache_heal_error(bucket, &object, &err) { + warn!("Skipping data usage cache heal for {}/{} due to transient error: {}", bucket, object, err); + } else { + failed += 1; + warn!("Failed to heal object {}/{}: {}", bucket, object, err); + } } Err(err) => { - failed += 1; - warn!("Failed to heal object {}/{}: {}", bucket, object, err); + if Self::should_skip_data_usage_cache_heal_error(bucket, &object, &err) { + warn!("Skipping data usage cache heal for {}/{} due to transient error: {}", bucket, object, err); + } else { + failed += 1; + warn!("Failed to heal object {}/{}: {}", bucket, object, err); + } } }, Err(err) => { @@ -1229,7 +1280,8 @@ mod tests { use super::*; use crate::heal::storage::DiskStatus; use rustfs_ecstore::{ - disk::{DiskStore, endpoint::Endpoint}, + data_usage::DATA_USAGE_CACHE_NAME, + disk::{BUCKET_META_PREFIX, DiskStore, RUSTFS_META_BUCKET, endpoint::Endpoint}, store_api::{BucketInfo, ObjectInfo}, }; use rustfs_madmin::heal_commands::HealResultItem; @@ -1306,13 +1358,21 @@ mod tests { async fn heal_object( &self, - _bucket: &str, + bucket: &str, object: &str, _version_id: Option<&str>, opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { self.healed_objects.lock().unwrap().push(object.to_string()); self.object_heal_opts.lock().unwrap().push(*opts); + 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", + )), + )); + } Ok(( HealResultItem { object_size: 1, @@ -1337,14 +1397,22 @@ mod tests { async fn list_objects_for_heal_page( &self, - _bucket: &str, + bucket: &str, _prefix: &str, continuation_token: Option<&str>, ) -> Result<(Vec, Option, bool)> { let mut listed = self.listed.lock().unwrap(); if continuation_token.is_none() && !*listed { *listed = true; - Ok((vec!["object-a".to_string(), "object-b".to_string()], None, false)) + let objects = if bucket == RUSTFS_META_BUCKET { + vec![ + format!("{BUCKET_META_PREFIX}/{DATA_USAGE_CACHE_NAME}"), + format!("{BUCKET_META_PREFIX}/bucket-metadata.bin"), + ] + } else { + vec!["object-a".to_string(), "object-b".to_string()] + }; + Ok((objects, None, false)) } else { Ok((Vec::new(), None, false)) } @@ -1422,4 +1490,55 @@ mod tests { assert!(object_opts.iter().all(|opts| opts.recreate)); assert!(object_opts.iter().all(|opts| opts.scan_mode == HealScanMode::Deep)); } + + #[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); + } }