From 86e9362459cd9ea5458f2874c97ea96fdb698eff Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 25 Jun 2026 22:07:38 +0800 Subject: [PATCH] fix(heal): canonicalize scanner object-dir repairs (#3864) --- .../ecstore/src/cache_value/metacache_set.rs | 160 ++++++--- crates/heal/src/heal/task.rs | 324 +++++++++++++++++- crates/scanner/src/scanner_folder.rs | 125 +++++-- 3 files changed, 532 insertions(+), 77 deletions(-) diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index 678e2b9fc..0a122e186 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -29,7 +29,7 @@ use tokio::io::AsyncRead; use tokio::spawn; use tokio::time::timeout; use tokio_util::sync::CancellationToken; -use tracing::{error, warn}; +use tracing::{debug, error, warn}; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_METACACHE: &str = "metacache"; @@ -55,6 +55,10 @@ async fn peek_with_timeout(reader: &mut MetacacheReader } } +fn is_missing_path_error(err: &DiskError) -> bool { + matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) +} + #[cfg(test)] #[derive(Clone)] pub(crate) enum TestReaderBehavior { @@ -197,17 +201,31 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d "metacache_walk_dir_primary_failed", primary_walk_started.elapsed().as_secs_f64() * 1000.0, ); - warn!( - event = EVENT_METACACHE_LISTING, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_METACACHE, - bucket = %opts_clone.bucket, - path = %opts_clone.path, - disk_index = disk_idx, - state = "walk_dir_failed", - error = ?err, - "Metacache walk_dir failed" - ); + if is_missing_path_error(&err) { + debug!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "walk_dir_missing_path", + error = ?err, + "Metacache walk_dir missing path skipped" + ); + } else { + warn!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "walk_dir_failed", + error = ?err, + "Metacache walk_dir failed" + ); + } last_err = Some(err); need_fallback = true; } @@ -232,17 +250,32 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d } let Some(disk) = disk_op else { - warn!( - event = EVENT_METACACHE_LISTING, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_METACACHE, - bucket = %opts_clone.bucket, - path = %opts_clone.path, - disk_index = disk_idx, - state = "fallback_disk_missing", - "Metacache fallback disk missing" - ); let err = last_err.unwrap_or(DiskError::DiskNotFound); + if is_missing_path_error(&err) { + debug!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "fallback_disk_missing_for_path", + error = ?err, + "Metacache fallback disk unavailable for missing path" + ); + } else { + warn!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "fallback_disk_missing", + error = ?err, + "Metacache fallback disk missing" + ); + } record_producer_error(&producer_errs_clone, disk_idx, &err); return Err(err); }; @@ -279,17 +312,31 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d "metacache_walk_dir_fallback_failed", fallback_walk_started.elapsed().as_secs_f64() * 1000.0, ); - error!( - event = EVENT_METACACHE_LISTING, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_METACACHE, - bucket = %opts_clone.bucket, - path = %opts_clone.path, - disk_index = disk_idx, - state = "fallback_walk_dir_failed", - error = ?err, - "Metacache fallback walk_dir failed" - ); + if is_missing_path_error(&err) { + debug!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "fallback_walk_dir_missing_path", + error = ?err, + "Metacache fallback walk_dir missing path skipped" + ); + } else { + error!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %opts_clone.bucket, + path = %opts_clone.path, + disk_index = disk_idx, + state = "fallback_walk_dir_failed", + error = ?err, + "Metacache fallback walk_dir failed" + ); + } last_err = Some(err); } } @@ -560,16 +607,29 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d let merge_started = std::time::Instant::now(); if let Err(err) = revjob.await.map_err(std::io::Error::other)? { rustfs_io_metrics::record_stage_duration("metacache_merge_failed", merge_started.elapsed().as_secs_f64() * 1000.0); - error!( - event = EVENT_METACACHE_LISTING, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_METACACHE, - bucket = %log_bucket, - path = %log_path, - state = "merge_job_failed", - error = ?err, - "Metacache merge job failed" - ); + if is_missing_path_error(&err) { + debug!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %log_bucket, + path = %log_path, + state = "merge_job_missing_path", + error = ?err, + "Metacache merge job missing path skipped" + ); + } else { + error!( + event = EVENT_METACACHE_LISTING, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_METACACHE, + bucket = %log_bucket, + path = %log_path, + state = "merge_job_failed", + error = ?err, + "Metacache merge job failed" + ); + } cancel_rx.cancel(); for job in jobs { job.abort(); @@ -597,8 +657,8 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d match result { Ok(Ok(())) => {} Ok(Err(err)) => { - if matches!(err, DiskError::FileNotFound | DiskError::VolumeNotFound) { - warn!( + if is_missing_path_error(&err) { + debug!( event = EVENT_METACACHE_LISTING, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_METACACHE, @@ -668,6 +728,16 @@ mod tests { assert_eq!(err, DiskError::ErasureReadQuorum); } + #[test] + fn missing_path_error_classification_excludes_actionable_failures() { + assert!(is_missing_path_error(&DiskError::FileNotFound)); + assert!(is_missing_path_error(&DiskError::FileVersionNotFound)); + assert!(is_missing_path_error(&DiskError::VolumeNotFound)); + assert!(!is_missing_path_error(&DiskError::Timeout)); + assert!(!is_missing_path_error(&DiskError::DiskNotFound)); + assert!(!is_missing_path_error(&DiskError::FileAccessDenied)); + } + #[tokio::test] async fn list_path_raw_returns_timeout_when_reader_stalls_before_completion() { let err = list_path_raw( diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index e96171b75..ab21cb55c 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -693,13 +693,36 @@ impl HealTask { "Heal object stage entered" ); self.check_control_flags().await?; - let object_exists = match self.await_with_control(self.storage.object_exists(bucket, object)).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 { warn!( target: "rustfs::heal::task", @@ -968,6 +991,40 @@ impl HealTask { } } + async fn canonicalize_scanner_missing_object_dir(&self, bucket: &str, object: &str) -> Result> { + 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!( @@ -2171,6 +2228,7 @@ mod tests { use crate::heal::storage::{DiskStatus, HealObjectInfo}; use rustfs_madmin::heal_commands::HealResultItem; use rustfs_storage_api::BucketInfo; + use std::collections::HashMap; use std::sync::Mutex; #[derive(Default)] @@ -2178,9 +2236,11 @@ mod tests { listed: Mutex, healed_objects: Mutex>, heal_object_calls: Mutex>, + heal_object_version_ids: Mutex>>, bucket_heal_opts: Mutex>, object_heal_opts: Mutex>, object_exists: Mutex>, + object_exists_by_name: Mutex>, heal_object_outcome: Mutex>, format_no_heal_required: Mutex, listed_prefixes: Mutex>, @@ -2193,6 +2253,13 @@ mod tests { ErrOther(&'static str), } + #[derive(Clone, Copy)] + enum MockObjectExists { + Exists(bool), + TransientSkip(&'static str), + OtherError(&'static str), + } + #[test] fn test_missing_object_dir_heal_result_matches_only_object_level_not_found() { assert!(is_missing_object_dir_heal_result("x.rnd/", &Error::Disk(DiskError::FileNotFound))); @@ -2258,7 +2325,14 @@ mod tests { }]) } - async fn object_exists(&self, _bucket: &str, _object: &str) -> Result { + async fn object_exists(&self, _bucket: &str, object: &str) -> Result { + if let Some(result) = self.object_exists_by_name.lock().unwrap().get(object).copied() { + return match result { + MockObjectExists::Exists(exists) => Ok(exists), + MockObjectExists::TransientSkip(message) => Err(Error::transient_skip(message)), + MockObjectExists::OtherError(message) => Err(Error::other(message)), + }; + } Ok(self.object_exists.lock().unwrap().unwrap_or(true)) } @@ -2274,10 +2348,14 @@ mod tests { &self, bucket: &str, object: &str, - _version_id: Option<&str>, + version_id: Option<&str>, opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { self.heal_object_calls.lock().unwrap().push(object.to_string()); + self.heal_object_version_ids + .lock() + .unwrap() + .push(version_id.map(ToString::to_string)); if let Some(outcome) = self.heal_object_outcome.lock().unwrap().take() { return match outcome { MockHealObjectOutcome::OkWithOtherError(message) => { @@ -2630,6 +2708,246 @@ mod tests { assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]); } + #[tokio::test] + async fn test_heal_scanner_missing_object_dir_canonicalizes_existing_plain_object() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true)); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd/".to_string(), + version_id: Some("version-a".to_string()), + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("scanner object-dir candidate should heal the existing plain object"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd".to_string()]); + assert_eq!( + storage.heal_object_version_ids.lock().unwrap().as_slice(), + [Some("version-a".to_string())] + ); + assert_eq!(storage.healed_objects.lock().unwrap().as_slice(), ["x.rnd".to_string()]); + } + + #[tokio::test] + async fn test_heal_scanner_existing_trailing_slash_object_is_not_canonicalized() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(true)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true)); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd/".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("existing trailing-slash object should keep its exact key"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]); + } + + #[tokio::test] + async fn test_heal_admin_missing_object_dir_does_not_canonicalize_plain_object() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true)); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::ErrOther("File not found"))), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd/".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Admin; + let task = HealTask::from_request(request, storage.clone()); + + let err = task + .execute() + .await + .expect_err("admin object-dir request must not be canonicalized"); + + assert!(matches!(err, Error::TaskExecutionFailed { .. })); + assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]); + } + + #[tokio::test] + async fn test_heal_scanner_canonicalizes_only_one_trailing_slash() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd//".to_string(), MockObjectExists::Exists(false)); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(true)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::Exists(true)); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd//".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("scanner canonicalization should remove only one trailing slash"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["x.rnd/".to_string()]); + } + + #[tokio::test] + async fn test_heal_scanner_trimmed_object_exists_error_is_not_recreated() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::OtherError("backend unavailable")); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd/".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + let err = task + .execute() + .await + .expect_err("trimmed object_exists error must not be treated as missing"); + + assert!(matches!(err, Error::Other(_))); + assert!(storage.heal_object_calls.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn test_heal_scanner_trimmed_object_exists_transient_skip_is_not_recreated() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("x.rnd/".to_string(), MockObjectExists::Exists(false)); + object_exists_by_name.insert("x.rnd".to_string(), MockObjectExists::TransientSkip("backend busy")); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "x.rnd/".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("trimmed object_exists transient skip should complete without recreate"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert!(storage.heal_object_calls.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn test_heal_scanner_empty_trimmed_object_keeps_existing_skip_behavior() { + let mut object_exists_by_name = HashMap::new(); + object_exists_by_name.insert("/".to_string(), MockObjectExists::Exists(false)); + let storage = Arc::new(MockStorage { + object_exists_by_name: Mutex::new(object_exists_by_name), + heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::OkWithOtherError("File not found"))), + ..Default::default() + }); + let mut request = HealRequest::new( + HealType::Object { + bucket: "bucket-a".to_string(), + object: "/".to_string(), + version_id: None, + }, + HealOptions { + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + request.source = HealRequestSource::Scanner; + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("empty canonical object must keep existing scanner skip behavior"); + + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); + assert_eq!(storage.heal_object_calls.lock().unwrap().as_slice(), ["/".to_string()]); + } + #[tokio::test] async fn test_heal_recreate_scanner_synthetic_object_dir_skips_err_not_found() { let storage = Arc::new(MockStorage { diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index d924ef187..fd3b05038 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -474,6 +474,21 @@ fn resolve_object_heal_entry(entries: &MetaCacheEntries, resolver: MetadataResol .cloned() } +fn is_missing_path_disk_error(err: &DiskError) -> bool { + matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) +} + +fn disk_errors_are_only_missing_paths(errs: &[Option]) -> bool { + let mut saw_missing_path = false; + for err in errs.iter().flatten() { + if !is_missing_path_disk_error(err) { + return false; + } + saw_missing_path = true; + } + saw_missing_path +} + fn heal_priority_label(priority: HealChannelPriority) -> &'static str { match priority { HealChannelPriority::Low => "low", @@ -995,13 +1010,18 @@ impl ScannerItem { async fn enqueue_heal(&mut self, oi: &ObjectInfo) { let done_heal = Metrics::time(Metric::HealAbandonedObject); + let object = if oi.name.is_empty() { + self.object_path() + } else { + oi.name.clone() + }; debug!( target: "rustfs::scanner::folder", event = EVENT_SCANNER_HEAL_ADMISSION, component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_HEAL, bucket = %self.bucket, - object = %self.object_path(), + object = %object, version_id = %oi.version_id.unwrap_or_default(), state = "request_started", "Scanner heal admission started" @@ -1021,7 +1041,7 @@ impl ScannerItem { component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_HEAL, bucket = %self.bucket, - object = %self.object_path(), + object = %object, version_id = %oi.version_id.unwrap_or_default(), object_age_secs = age_secs.unwrap_or_default(), cooldown_secs = cooldown.as_secs(), @@ -1036,7 +1056,7 @@ impl ScannerItem { "object", build_object_heal_request( self.bucket.clone(), - self.object_path(), + object.clone(), oi.version_id .and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) }), scan_mode, @@ -1056,7 +1076,7 @@ impl ScannerItem { component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_HEAL, bucket = %self.bucket, - object = %self.object_path(), + object = %object, admission = %describe_heal_admission(result), state = "not_admitted", "Scanner heal admission rejected low-priority request" @@ -1068,7 +1088,7 @@ impl ScannerItem { component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_HEAL, bucket = %self.bucket, - object = %self.object_path(), + object = %object, state = "submit_failed", error = %e, "Scanner heal admission submission failed" @@ -1438,7 +1458,7 @@ impl FolderScanner { let mut dir_reader = match tokio::fs::read_dir(&dir_path).await { Ok(dir_reader) => dir_reader, Err(e) if e.kind() == ErrorKind::NotFound => { - warn!( + debug!( target: "rustfs::scanner::folder", event = EVENT_SCANNER_FOLDER_STATE, component = LOG_COMPONENT_SCANNER, @@ -1458,7 +1478,7 @@ impl FolderScanner { Ok(Some(entry)) => entry, Ok(None) => break, Err(e) if e.kind() == ErrorKind::NotFound => { - warn!( + debug!( target: "rustfs::scanner::folder", event = EVENT_SCANNER_FOLDER_STATE, component = LOG_COMPONENT_SCANNER, @@ -1505,7 +1525,7 @@ impl FolderScanner { let mut entry_type = match entry.file_type().await { Ok(entry_type) => entry_type, Err(e) if e.kind() == ErrorKind::NotFound => { - warn!( + debug!( target: "rustfs::scanner::folder", event = EVENT_SCANNER_FOLDER_STATE, component = LOG_COMPONENT_SCANNER, @@ -1537,7 +1557,7 @@ impl FolderScanner { let metadata = match tokio::fs::metadata(&file_path).await { Ok(metadata) => metadata, Err(e) if e.kind() == ErrorKind::NotFound => { - warn!( + debug!( target: "rustfs::scanner::folder", event = EVENT_SCANNER_FOLDER_STATE, component = LOG_COMPONENT_SCANNER, @@ -2043,17 +2063,31 @@ impl FolderScanner { ) .await { - error!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_FOLDER_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_FOLDER, - bucket = %bucket_clone, - prefix = %prefix_clone, - state = "list_path_failed", - error = %e, - "Scanner list_path_raw failed" - ); + if is_missing_path_disk_error(&e) { + debug!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_FOLDER_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_FOLDER, + bucket = %bucket_clone, + prefix = %prefix_clone, + state = "list_path_missing", + error = %e, + "Scanner list_path_raw missing path skipped" + ); + } else { + error!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_FOLDER_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_FOLDER, + bucket = %bucket_clone, + prefix = %prefix_clone, + state = "list_path_failed", + error = %e, + "Scanner list_path_raw failed" + ); + } } }); @@ -2151,15 +2185,27 @@ impl FolderScanner { finished_closed = true; continue; }; - error!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_FOLDER_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_FOLDER, - state = "list_path_finished_with_errors", - errors = ?errs, - "Scanner list_path_raw finished with disk errors" - ); + if disk_errors_are_only_missing_paths(&errs) { + debug!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_FOLDER_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_FOLDER, + state = "list_path_finished_missing_paths", + errors = ?errs, + "Scanner list_path_raw finished with missing paths" + ); + } else { + error!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_FOLDER_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_FOLDER, + state = "list_path_finished_with_errors", + errors = ?errs, + "Scanner list_path_raw finished with disk errors" + ); + } child_ctx.cancel(); } _ = child_ctx.cancelled() => { @@ -3227,6 +3273,27 @@ mod tests { assert!(entry.is_object_dir()); } + #[test] + fn test_disk_errors_are_only_missing_paths_accepts_missing_mix() { + let errs = vec![ + Some(DiskError::FileNotFound), + None, + Some(DiskError::FileVersionNotFound), + Some(DiskError::VolumeNotFound), + ]; + + assert!(disk_errors_are_only_missing_paths(&errs)); + } + + #[test] + fn test_disk_errors_are_only_missing_paths_rejects_empty_or_actionable_errors() { + assert!(!disk_errors_are_only_missing_paths(&[None, None])); + assert!(!disk_errors_are_only_missing_paths(&[ + Some(DiskError::FileNotFound), + Some(DiskError::Timeout), + ])); + } + #[test] fn test_effective_object_heal_scan_mode_keeps_normal_when_bitrot_disabled() { let now = OffsetDateTime::now_utc();