diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index a718ecd04..c5c883f8a 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -739,9 +739,14 @@ pub struct Metrics { scanner_set_scans_queued: AtomicU64, scanner_set_scans_active: AtomicU64, scanner_disk_bucket_scan_states: Mutex>, + scanner_leader_lock_state: RwLock, + scanner_leader_lock_held: AtomicBool, + scanner_leader_lock_last_error: RwLock, + scanner_leader_lock_last_update_unix_secs: AtomicU64, last_scan_cycle_result: AtomicU8, last_scan_cycle_partial_reason: AtomicU8, last_scan_cycle_partial_source: AtomicU8, + last_scan_cycle_end_unix_secs: AtomicU64, last_scan_cycle_duration_millis: AtomicU64, last_scan_cycle_objects_scanned: AtomicU64, last_scan_cycle_directories_scanned: AtomicU64, @@ -1105,6 +1110,16 @@ pub struct ScannerMetricsReport { pub active_paths: Vec, pub current_scan_mode: String, #[serde(default)] + pub leader_lock_state: String, + #[serde(default)] + pub leader_lock_held_by_this_process: bool, + #[serde(default)] + pub leader_lock_last_error: String, + #[serde(default)] + pub leader_lock_last_update_unix_secs: u64, + #[serde(default)] + pub last_cycle_end_unix_secs: u64, + #[serde(default)] pub current_set_scan_concurrency_limit: u64, #[serde(default)] pub current_set_scans_queued: u64, @@ -1678,9 +1693,14 @@ impl Metrics { scanner_set_scans_queued: AtomicU64::new(0), scanner_set_scans_active: AtomicU64::new(0), scanner_disk_bucket_scan_states: Mutex::new(HashMap::new()), + scanner_leader_lock_state: RwLock::new("unknown".to_string()), + scanner_leader_lock_held: AtomicBool::new(false), + scanner_leader_lock_last_error: RwLock::new(String::new()), + scanner_leader_lock_last_update_unix_secs: AtomicU64::new(0), last_scan_cycle_result: AtomicU8::new(SCAN_CYCLE_RESULT_UNKNOWN), last_scan_cycle_partial_reason: AtomicU8::new(ScanCyclePartialReason::Unknown as u8), last_scan_cycle_partial_source: AtomicU8::new(0), + last_scan_cycle_end_unix_secs: AtomicU64::new(0), last_scan_cycle_duration_millis: AtomicU64::new(0), last_scan_cycle_objects_scanned: AtomicU64::new(0), last_scan_cycle_directories_scanned: AtomicU64::new(0), @@ -2280,6 +2300,19 @@ impl Metrics { HealScanMode::from_u8(self.current_scan_mode.load(Ordering::Relaxed)).unwrap_or(HealScanMode::Unknown) } + pub async fn record_scanner_leader_liveness(&self, state: impl Into, held: bool, error: impl Into) { + *self.scanner_leader_lock_state.write().await = state.into(); + *self.scanner_leader_lock_last_error.write().await = error.into(); + self.scanner_leader_lock_held.store(held, Ordering::Relaxed); + self.scanner_leader_lock_last_update_unix_secs + .store(Utc::now().timestamp().max(0) as u64, Ordering::Relaxed); + } + + pub fn record_scanner_cycle_end_time(&self) { + self.last_scan_cycle_end_unix_secs + .store(Utc::now().timestamp().max(0) as u64, Ordering::Relaxed); + } + pub fn record_scan_cycle_complete(&self, success: bool, duration: Duration) { let result = if success { SCAN_CYCLE_RESULT_SUCCESS @@ -2287,6 +2320,7 @@ impl Metrics { self.failed_scan_cycles.fetch_add(1, Ordering::Relaxed); SCAN_CYCLE_RESULT_ERROR }; + self.record_scanner_cycle_end_time(); self.last_scan_cycle_result.store(result, Ordering::Relaxed); self.last_scan_cycle_partial_reason .store(ScanCyclePartialReason::Unknown as u8, Ordering::Relaxed); @@ -2305,6 +2339,7 @@ impl Metrics { reason: ScanCyclePartialReason, source: Option, ) { + self.record_scanner_cycle_end_time(); self.partial_scan_cycles.fetch_add(1, Ordering::Relaxed); match reason { ScanCyclePartialReason::Unknown => &self.partial_scan_cycles_unknown, @@ -2660,6 +2695,11 @@ impl Metrics { .map(|(disk, state)| format!("{disk}/{}", state.path)) .collect(); m.current_scan_mode = self.current_scan_mode().as_str().to_string(); + m.leader_lock_state = self.scanner_leader_lock_state.read().await.clone(); + m.leader_lock_held_by_this_process = self.scanner_leader_lock_held.load(Ordering::Relaxed); + m.leader_lock_last_error = self.scanner_leader_lock_last_error.read().await.clone(); + m.leader_lock_last_update_unix_secs = self.scanner_leader_lock_last_update_unix_secs.load(Ordering::Relaxed); + m.last_cycle_end_unix_secs = self.last_scan_cycle_end_unix_secs.load(Ordering::Relaxed); m.current_set_scan_concurrency_limit = self.scanner_set_scan_concurrency_limit.load(Ordering::Relaxed); m.current_set_scans_queued = self.scanner_set_scans_queued.load(Ordering::Relaxed); m.current_set_scans_active = self.scanner_set_scans_active.load(Ordering::Relaxed); @@ -3671,6 +3711,21 @@ mod tests { assert_eq!(report.current_scan_mode, HealScanMode::Deep.as_str()); } + #[tokio::test] + async fn report_includes_scanner_leader_liveness() { + let metrics = Metrics::new(); + + metrics.record_scanner_leader_liveness("contended", false, "lock busy").await; + metrics.record_scanner_cycle_end_time(); + let report = metrics.report().await; + + assert_eq!(report.leader_lock_state, "contended"); + assert!(!report.leader_lock_held_by_this_process); + assert_eq!(report.leader_lock_last_error, "lock busy"); + assert!(report.leader_lock_last_update_unix_secs > 0); + assert!(report.last_cycle_end_unix_secs > 0); + } + #[tokio::test] async fn report_includes_last_scan_cycle_result() { let metrics = Metrics::new(); diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index c79f0cba3..a95e9a620 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -23,6 +23,7 @@ use std::{ use http::HeaderMap; use metrics::{counter, describe_counter, describe_histogram, histogram}; +use rustfs_common::heal_channel::HealScanMode; #[cfg(test)] use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS; pub use rustfs_data_usage::{ @@ -75,6 +76,8 @@ pub static DATA_USAGE_BLOOM_NAME_PATH: LazyLock = pub static BACKGROUND_HEAL_INFO_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}.background-heal.json")); +const MAX_DATA_USAGE_CACHE_DEPTH: usize = 1024; + #[derive(Clone, Copy, Default, Debug, Serialize, Deserialize, PartialEq)] pub struct TierStats { pub total_size: u64, @@ -191,7 +194,7 @@ impl SizeSummary { } if let Some(tier_stats) = self.tier_stats.get_mut(&tier) { - tier_stats.add(&TierStats::from_object_info(oi)); + *tier_stats = tier_stats.add(&TierStats::from_object_info(oi)); } } } @@ -259,6 +262,31 @@ pub struct DataUsageEntryInfo { pub entry: DataUsageEntry, } +#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum PendingScannerHealKind { + Bucket, + Object, +} + +#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] +pub struct PendingScannerHeal { + pub kind: PendingScannerHealKind, + pub bucket: String, + #[serde(default)] + pub object: Option, + #[serde(default)] + pub version_id: Option, + pub scan_mode: HealScanMode, + pub first_seen: u64, + pub last_attempt: u64, + pub attempts: u32, + #[serde(default)] + pub last_admission_result: String, + #[serde(default)] + pub last_admission_reason: String, +} + /// Data usage cache info #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageCacheInfo { @@ -274,6 +302,8 @@ pub struct DataUsageCacheInfo { pub scan_resume_after: Option, #[serde(default)] pub scan_checkpoint: Option, + #[serde(default)] + pub pending_heals: Vec, } /// Data usage cache @@ -334,12 +364,25 @@ impl DataUsageCache { } pub fn flatten(&self, root: &DataUsageEntry) -> DataUsageEntry { + let mut visited = HashSet::new(); + self.flatten_with_guard(root, &mut visited, 0) + } + + fn flatten_with_guard(&self, root: &DataUsageEntry, visited: &mut HashSet, depth: usize) -> DataUsageEntry { let mut root = root.clone(); + if depth >= MAX_DATA_USAGE_CACHE_DEPTH { + root.children.clear(); + return root; + } + for id in root.children.clone().iter() { + if !visited.insert(id.clone()) { + continue; + } if let Some(e) = self.cache.get(id) { let mut e = e.clone(); if !e.children.is_empty() { - e = self.flatten(&e); + e = self.flatten_with_guard(&e, visited, depth + 1); } root.merge(&e); } @@ -349,13 +392,31 @@ impl DataUsageCache { } pub fn copy_with_children(&mut self, src: &DataUsageCache, hash: &DataUsageHash, parent: &Option) { + let mut visited = HashSet::new(); + self.copy_with_children_guard(src, hash, parent, &mut visited, 0); + } + + fn copy_with_children_guard( + &mut self, + src: &DataUsageCache, + hash: &DataUsageHash, + parent: &Option, + visited: &mut HashSet, + depth: usize, + ) { + if !visited.insert(hash.key()) { + return; + } + if let Some(e) = src.cache.get(&hash.string()) { self.cache.insert(hash.key(), e.clone()); - for ch in e.children.iter() { - if *ch == hash.key() { - return; + if depth < MAX_DATA_USAGE_CACHE_DEPTH { + for ch in e.children.iter() { + if *ch == hash.key() { + continue; + } + self.copy_with_children_guard(src, &DataUsageHash(ch.to_string()), &Some(hash.clone()), visited, depth + 1); } - self.copy_with_children(src, &DataUsageHash(ch.to_string()), &Some(hash.clone())); } if let Some(parent) = parent { self.cache.entry(parent.key()).or_default().add_child(hash); @@ -364,6 +425,15 @@ impl DataUsageCache { } pub fn delete_recursive(&mut self, hash: &DataUsageHash) { + let mut visited = HashSet::new(); + self.delete_recursive_guard(hash, &mut visited, 0); + } + + fn delete_recursive_guard(&mut self, hash: &DataUsageHash, visited: &mut HashSet, depth: usize) { + if !visited.insert(hash.key()) { + return; + } + let mut need_remove = Vec::new(); if let Some(v) = self.cache.get(&hash.string()) { for child in v.children.iter() { @@ -371,8 +441,11 @@ impl DataUsageCache { } } self.cache.remove(&hash.string()); + if depth >= MAX_DATA_USAGE_CACHE_DEPTH { + return; + } for child in need_remove { - self.delete_recursive(&DataUsageHash(child)); + self.delete_recursive_guard(&DataUsageHash(child), visited, depth + 1); } } @@ -382,7 +455,9 @@ impl DataUsageCache { if root.children.is_empty() { return Some(root.clone()); } - let mut flat = self.flatten(root); + let mut visited = HashSet::new(); + visited.insert(hash_path(path).key()); + let mut flat = self.flatten_with_guard(root, &mut visited, 0); if flat.replication_stats.as_ref().is_some_and(|stats| stats.empty()) { flat.replication_stats = None; } @@ -488,16 +563,24 @@ impl DataUsageCache { } pub fn total_children_rec(&self, path: &str) -> usize { + let mut visited = HashSet::new(); + visited.insert(hash_path(path).key()); + self.total_children_rec_guard(path, &mut visited, 0) + } + + fn total_children_rec_guard(&self, path: &str, visited: &mut HashSet, depth: usize) -> usize { let Some(root) = self.find(path) else { return 0; }; - if root.children.is_empty() { + if root.children.is_empty() || depth >= MAX_DATA_USAGE_CACHE_DEPTH { return 0; } - let mut n = root.children.len(); + let mut n = 0; for ch in root.children.iter() { - n += self.total_children_rec(ch); + if visited.insert(ch.clone()) { + n += 1 + self.total_children_rec_guard(ch, visited, depth + 1); + } } n } @@ -938,13 +1021,29 @@ struct Inner { } fn add(data_usage_cache: &DataUsageCache, path: &DataUsageHash, candidates: &mut Vec) -> usize { + let mut visited = HashSet::new(); + visited.insert(path.key()); + add_with_guard(data_usage_cache, path, candidates, &mut visited, 0) +} + +fn add_with_guard( + data_usage_cache: &DataUsageCache, + path: &DataUsageHash, + candidates: &mut Vec, + visited: &mut HashSet, + depth: usize, +) -> usize { let e = match data_usage_cache.cache.get(&path.key()) { Some(e) => e, None => return 0, }; let mut objects = e.objects; - for ch in e.children.iter() { - objects += add(data_usage_cache, &DataUsageHash(ch.clone()), candidates); + if depth < MAX_DATA_USAGE_CACHE_DEPTH { + for ch in e.children.iter() { + if visited.insert(ch.clone()) { + objects += add_with_guard(data_usage_cache, &DataUsageHash(ch.clone()), candidates, visited, depth + 1); + } + } } // Collect internal nodes (with children) as compaction candidates. // Leaf nodes have no children to remove, so compacting them is a no-op — @@ -959,10 +1058,20 @@ fn add(data_usage_cache: &DataUsageCache, path: &DataUsageHash, candidates: &mut } fn mark(duc: &DataUsageCache, entry: &DataUsageEntry, found: &mut HashSet) { + mark_with_depth(duc, entry, found, 0); +} + +fn mark_with_depth(duc: &DataUsageCache, entry: &DataUsageEntry, found: &mut HashSet, depth: usize) { + if depth >= MAX_DATA_USAGE_CACHE_DEPTH { + return; + } + for k in entry.children.iter() { - found.insert(k.to_string()); + if !found.insert(k.to_string()) { + continue; + } if let Some(ch) = duc.cache.get(k) { - mark(duc, ch, found); + mark_with_depth(duc, ch, found, depth + 1); } } } @@ -1082,6 +1191,37 @@ mod tests { assert_eq!(summary.total_size, 0); } + #[test] + fn size_summary_actions_accounting_accumulates_tier_stats() { + let mut summary = SizeSummary::new(); + summary + .tier_stats + .insert(storageclass::STANDARD.to_string(), TierStats::default()); + + let object = ObjectInfo { + storage_class: Some(storageclass::STANDARD.to_string()), + size: 10, + is_latest: true, + ..Default::default() + }; + + summary.actions_accounting(&object, 10, 10); + summary.actions_accounting(&object, 10, 10); + + let stats = summary + .tier_stats + .get(storageclass::STANDARD) + .expect("standard tier stats should remain present"); + assert_eq!( + *stats, + TierStats { + total_size: 20, + num_versions: 2, + num_objects: 2, + } + ); + } + #[test] fn test_data_usage_entry_merge_sums_failed_objects() { let mut left = DataUsageEntry { @@ -1132,6 +1272,7 @@ mod tests { assert_eq!(decoded.next_cycle, 7); assert!(decoded.scan_resume_after.is_none()); assert!(decoded.scan_checkpoint.is_none()); + assert!(decoded.pending_heals.is_empty()); } #[test] @@ -1169,6 +1310,7 @@ mod tests { assert_eq!(decoded.failed_objects.get("bad-object"), Some(&11)); assert!(decoded.scan_resume_after.is_none()); assert!(decoded.scan_checkpoint.is_none()); + assert!(decoded.pending_heals.is_empty()); } #[test] @@ -1267,6 +1409,89 @@ mod tests { assert!(dst.cache.contains_key(&root_hash.key())); } + #[test] + fn test_data_usage_cache_recursive_helpers_tolerate_cycles() { + let root_hash = hash_path("bucket"); + let child_hash = hash_path("bucket/a"); + + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + ..Default::default() + }, + ..Default::default() + }; + cache.replace_hashed(&root_hash, &None, &DataUsageEntry::default()); + cache.replace_hashed( + &child_hash, + &Some(root_hash.clone()), + &DataUsageEntry { + objects: 2, + size: 20, + ..Default::default() + }, + ); + cache.cache.entry(child_hash.key()).or_default().add_child(&root_hash); + + assert_eq!(cache.total_children_rec("bucket"), 1); + + let flat = cache.size_recursive("bucket").expect("cyclic cache should still flatten"); + assert_eq!(flat.objects, 2); + assert_eq!(flat.size, 20); + assert!(flat.children.is_empty()); + + let mut copied = DataUsageCache { + info: cache.info.clone(), + ..Default::default() + }; + copied.copy_with_children(&cache, &root_hash, &None); + assert!(copied.cache.contains_key(&root_hash.key())); + assert!(copied.cache.contains_key(&child_hash.key())); + + copied.delete_recursive(&root_hash); + assert!(!copied.cache.contains_key(&root_hash.key())); + assert!(!copied.cache.contains_key(&child_hash.key())); + } + + #[test] + fn test_data_usage_cache_flatten_does_not_count_root_twice_in_cycle() { + let root_hash = hash_path("bucket"); + let child_hash = hash_path("bucket/a"); + + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + ..Default::default() + }, + ..Default::default() + }; + cache.replace_hashed( + &root_hash, + &None, + &DataUsageEntry { + objects: 1, + size: 10, + ..Default::default() + }, + ); + cache.replace_hashed( + &child_hash, + &Some(root_hash.clone()), + &DataUsageEntry { + objects: 2, + size: 20, + ..Default::default() + }, + ); + cache.cache.entry(child_hash.key()).or_default().add_child(&root_hash); + + let flat = cache.size_recursive("bucket").expect("cyclic cache should still flatten"); + + assert_eq!(flat.objects, 3); + assert_eq!(flat.size, 30); + assert!(flat.children.is_empty()); + } + #[test] fn test_find_children_copy_preserves_missing_entry_behavior() { let mut cache = DataUsageCache::default(); diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index e00847aad..9ddd95753 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -59,6 +59,7 @@ const EVENT_SCANNER_LOCK_STATE: &str = "scanner_lock_state"; const EVENT_SCANNER_PERSIST_STATE: &str = "scanner_persist_state"; const EVENT_SCANNER_RUNTIME_CONFIG: &str = "scanner_runtime_config"; const EVENT_SCANNER_BACKGROUND_HEAL_STATE: &str = "scanner_background_heal_state"; +const METRIC_SCANNER_LEADER_LOCK_TOTAL: &str = "rustfs_scanner_leader_lock_total"; #[cfg(test)] const ENV_SCANNER_START_DELAY_SECS_DEPRECATED: &str = "RUSTFS_DATA_SCANNER_START_DELAY_SECS"; @@ -77,6 +78,14 @@ fn scanner_cycle_budget_config() -> ScannerCycleBudgetConfig { resolve_scanner_runtime_config().cycle_budget } +fn record_scanner_leader_lock_state(state: &'static str) { + metrics::counter!( + METRIC_SCANNER_LEADER_LOCK_TOTAL, + "state" => state + ) + .increment(1); +} + #[cfg(test)] fn scanner_cycle_max_duration() -> Option { resolve_scanner_runtime_config().cycle_budget.max_duration @@ -550,6 +559,18 @@ fn background_heal_info_for_scan_complete(mut info: BackgroundHealInfo, scan_mod Some(info) } +fn background_heal_info_for_scan_result( + info: BackgroundHealInfo, + scan_mode: HealScanMode, + success: bool, +) -> Option { + if !success { + return None; + } + + background_heal_info_for_scan_complete(info, scan_mode) +} + fn retain_recent_cycle_completions(cycle_completed: &mut Vec>) { let keep = data_usage_update_dir_cycles() as usize; if cycle_completed.len() > keep { @@ -780,9 +801,10 @@ async fn run_data_scanner_cycle(ctx: &CancellationToken, storeapi: &Arc "Scanner cycle failed" ); emit_scan_cycle_complete(false, cycle_start.elapsed()); - if let Some(new_heal_info) = background_heal_info_for_scan_complete(background_heal_info.clone(), scan_mode) { + if let Some(new_heal_info) = background_heal_info_for_scan_result(background_heal_info.clone(), scan_mode, false) { save_background_heal_info(storeapi.clone(), new_heal_info).await; } + mark_scan_cycle_idle(cycle_info).await; return; } if cycle_budget.budget_elapsed() && !ctx.is_cancelled() { @@ -813,7 +835,7 @@ async fn run_data_scanner_cycle(ctx: &CancellationToken, storeapi: &Arc done_cycle(); global_metrics().finish_scan_cycle_work(cycle_work_start); emit_scan_cycle_complete(true, cycle_start.elapsed()); - if let Some(new_heal_info) = background_heal_info_for_scan_complete(background_heal_info.clone(), scan_mode) { + if let Some(new_heal_info) = background_heal_info_for_scan_result(background_heal_info.clone(), scan_mode, true) { save_background_heal_info(storeapi.clone(), new_heal_info).await; } @@ -873,6 +895,8 @@ pub async fn run_data_scanner(ctx: CancellationToken, storeapi: Arc) -> let _guard = match storeapi.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock").await { Ok(ns_lock) => match ns_lock.get_write_lock_quiet(get_lock_acquire_timeout()).await { Ok(guard) => { + record_scanner_leader_lock_state("acquired"); + global_metrics().record_scanner_leader_liveness("acquired", true, "").await; debug!( target: "rustfs::scanner", event = EVENT_SCANNER_LOCK_STATE, @@ -885,6 +909,10 @@ pub async fn run_data_scanner(ctx: CancellationToken, storeapi: Arc) -> guard } Err(e) => { + record_scanner_leader_lock_state("contended"); + global_metrics() + .record_scanner_leader_liveness("contended", false, e.to_string()) + .await; debug!( target: "rustfs::scanner", event = EVENT_SCANNER_LOCK_STATE, @@ -899,6 +927,10 @@ pub async fn run_data_scanner(ctx: CancellationToken, storeapi: Arc) -> } }, Err(e) => { + record_scanner_leader_lock_state("create_failed"); + global_metrics() + .record_scanner_leader_liveness("create_failed", false, e.to_string()) + .await; error!( target: "rustfs::scanner", event = EVENT_SCANNER_LOCK_STATE, @@ -965,6 +997,7 @@ pub async fn run_data_scanner(ctx: CancellationToken, storeapi: Arc) -> } global_metrics().set_cycle(None).await; + global_metrics().record_scanner_leader_liveness("stopped", false, "").await; debug!( target: "rustfs::scanner", @@ -993,6 +1026,14 @@ impl Drop for ScannerScanModeGuard { } } +fn stale_data_usage_update_reason(incoming: &DataUsageInfo, existing: &DataUsageInfo) -> Option<&'static str> { + match (incoming.last_update, existing.last_update) { + (Some(new_ts), Some(existing_ts)) if new_ts <= existing_ts => Some("older_or_equal_last_update"), + (None, Some(_)) => Some("missing_incoming_last_update"), + _ => None, + } +} + /// Store data usage info in backend. Will store all objects sent on the receiver until closed. #[instrument(skip(ctx, storeapi))] pub async fn store_data_usage_in_backend( @@ -1010,8 +1051,7 @@ pub async fn store_data_usage_in_backend( if let Ok(buf) = read_config(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await && let Ok(existing) = serde_json::from_slice::(&buf) - && let (Some(new_ts), Some(existing_ts)) = (data_usage_info.last_update, existing.last_update) - && new_ts <= existing_ts + && let Some(reason) = stale_data_usage_update_reason(&data_usage_info, &existing) { debug!( target: "rustfs::scanner", @@ -1019,8 +1059,9 @@ pub async fn store_data_usage_in_backend( component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_RUNTIME, path = %DATA_USAGE_OBJ_NAME_PATH.as_str(), - incoming_last_update = ?new_ts, - existing_last_update = ?existing_ts, + incoming_last_update = ?data_usage_info.last_update, + existing_last_update = ?existing.last_update, + reason = reason, state = "skip_stale_update", "Scanner stale data usage update skipped" ); @@ -1445,6 +1486,45 @@ mod tests { assert_eq!(saved.last_update, Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20))); } + #[tokio::test] + async fn test_store_data_usage_in_backend_rejects_untimestamped_stale_snapshot() { + let store = Arc::new(MemoryConfigStore::default()); + let (sender, receiver) = mpsc::channel(2); + let ctx = CancellationToken::new(); + + let timestamped = DataUsageInfo { + last_update: Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)), + buckets_count: 2, + ..Default::default() + }; + let untimestamped = DataUsageInfo { + last_update: None, + buckets_count: 1, + ..Default::default() + }; + + sender + .send(timestamped) + .await + .expect("timestamped usage snapshot should enqueue"); + sender + .send(untimestamped) + .await + .expect("untimestamped usage snapshot should enqueue"); + drop(sender); + + store_data_usage_in_backend(ctx, store.clone(), receiver).await; + + let objects = store.objects.lock().await; + let saved = objects + .get(&memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str())) + .expect("data usage config should be saved"); + let saved = serde_json::from_slice::(saved).expect("saved usage snapshot should decode"); + + assert_eq!(saved.buckets_count, 2); + assert_eq!(saved.last_update, Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20))); + } + #[test] #[serial] fn test_cycle_interval_prefers_explicit_cycle_override() { @@ -1733,6 +1813,18 @@ mod tests { assert!(background_heal_info_for_scan_complete(info, HealScanMode::Normal).is_none()); } + #[test] + #[serial] + fn test_background_heal_info_for_failed_scan_preserves_deep_mode() { + let info = BackgroundHealInfo { + bitrot_start_time: Some(Utc::now()), + bitrot_start_cycle: 7, + current_scan_mode: HealScanMode::Deep, + }; + + assert!(background_heal_info_for_scan_result(info, HealScanMode::Deep, false).is_none()); + } + #[test] fn test_retain_recent_cycle_completions_keeps_last_entries() { let base = Utc::now(); diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 3c83c0042..19b7b7501 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -21,19 +21,21 @@ use std::time::{Duration, Instant, SystemTime}; use crate::ReplTargetSizeSummary; use crate::data_usage_define::{ DATA_USAGE_SCAN_CHECKPOINT_VERSION, DataUsageCache, DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageScanCheckpoint, - DataUsageScanCheckpointReason, SizeSummary, hash_path, + DataUsageScanCheckpointReason, PendingScannerHeal, PendingScannerHealKind, SizeSummary, hash_path, }; use crate::error::ScannerError; use crate::runtime_config::{ scanner_alert_excess_folders, scanner_alert_excess_version_size, scanner_alert_excess_versions, scanner_yield_every_n_objects, }; use crate::scanner_budget::{ScannerCycleBudget, ScannerCycleBudgetReason}; -use crate::scanner_io::{SCANNER_SKIP_FILE_ERROR, ScannerIODisk as _}; +use crate::scanner_io::{ + SCANNER_SKIP_FILE_ERROR, ScannerIODisk as _, is_scanner_metadata_corrupt_error, is_scanner_metadata_transient_error, +}; use crate::sleeper::DynamicSleeper; use metrics::{counter, describe_counter}; use rustfs_common::heal_channel::{ - HEAL_DELETE_DANGLING, HealAdmissionResult, HealChannelPriority, HealChannelRequest, HealRequestSource, HealScanMode, - send_heal_request_with_admission, + HEAL_DELETE_DANGLING, HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest, + HealRequestSource, HealScanMode, send_heal_request_with_admission, }; use rustfs_common::metrics::{ IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn, @@ -87,6 +89,10 @@ const METRIC_SCANNER_INLINE_HEAL_TOTAL: &str = "rustfs_scanner_inline_heal_total const METRIC_SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL: &str = "rustfs_scanner_excess_object_versions_total"; const METRIC_SCANNER_EXCESS_OBJECT_VERSION_SIZE_TOTAL: &str = "rustfs_scanner_excess_object_version_size_total"; const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders_total"; +const METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL: &str = "rustfs_scanner_pending_heal_prune_total"; +const METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL: &str = "rustfs_scanner_pending_heal_malformed_total"; +const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128; +const MAX_PENDING_SCANNER_HEALS_PER_BUCKET: usize = 10_000; static SCANNER_INLINE_HEAL_WARN_ONCE: Once = Once::new(); static SCANNER_INLINE_HEAL_METRICS_ONCE: Once = Once::new(); @@ -433,6 +439,13 @@ pub struct CachedFolder { /// Type alias for get size function pub type GetSizeFn = Box Result + Send + Sync>; +#[derive(Debug, PartialEq, Eq)] +enum GetSizeFailureAction { + Skip, + RecordFailed, + HealMetadata { object: String }, +} + fn build_bucket_heal_request(bucket: String, priority: HealChannelPriority) -> HealChannelRequest { HealChannelRequest { bucket, @@ -463,6 +476,62 @@ fn build_object_heal_request( } } +fn pending_scanner_heal_candidate_type(kind: PendingScannerHealKind) -> &'static str { + match kind { + PendingScannerHealKind::Bucket => "bucket", + PendingScannerHealKind::Object => "object", + } +} + +fn pending_scanner_heal_matches( + entry: &PendingScannerHeal, + kind: PendingScannerHealKind, + bucket: &str, + object: Option<&str>, + version_id: Option<&str>, +) -> bool { + entry.kind == kind && entry.bucket == bucket && entry.object.as_deref() == object && entry.version_id.as_deref() == version_id +} + +fn pending_scanner_heal_identity(entry: &PendingScannerHeal) -> (u8, &str, Option<&str>, Option<&str>) { + let kind = match entry.kind { + PendingScannerHealKind::Bucket => 0, + PendingScannerHealKind::Object => 1, + }; + (kind, entry.bucket.as_str(), entry.object.as_deref(), entry.version_id.as_deref()) +} + +fn sort_pending_scanner_heals_for_retry(entries: &mut [PendingScannerHeal]) { + entries.sort_by(|a, b| { + a.last_attempt + .cmp(&b.last_attempt) + .then_with(|| a.attempts.cmp(&b.attempts)) + .then_with(|| pending_scanner_heal_identity(a).cmp(&pending_scanner_heal_identity(b))) + }); +} + +fn pending_scanner_heal_retry_candidates(pending_heals: &[PendingScannerHeal], bucket: &str) -> Vec { + let mut entries: Vec = pending_heals.iter().filter(|entry| entry.bucket == bucket).cloned().collect(); + sort_pending_scanner_heals_for_retry(&mut entries); + entries.truncate(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET); + entries +} + +fn build_pending_scanner_heal_request(entry: &PendingScannerHeal) -> Option { + match entry.kind { + PendingScannerHealKind::Bucket => Some(build_bucket_heal_request(entry.bucket.clone(), HealChannelPriority::High)), + PendingScannerHealKind::Object => entry.object.as_ref().map(|object| { + build_object_heal_request( + entry.bucket.clone(), + object.clone(), + entry.version_id.clone(), + entry.scan_mode, + HealChannelPriority::High, + ) + }), + } +} + fn resolve_object_heal_entry(entries: &MetaCacheEntries, resolver: MetadataResolutionParams) -> Option { if let Some(entry) = entries.resolve(resolver) { return entry.is_object().then_some(entry); @@ -588,36 +657,6 @@ async fn send_scanner_heal_request( } } -async fn send_required_scanner_heal_request( - candidate_type: &'static str, - bucket: &str, - object: Option<&str>, - request: HealChannelRequest, -) -> Result<(), ScannerError> { - let priority = request.priority; - let result = send_scanner_heal_request(candidate_type, request).await?; - if result.is_admitted() { - return Ok(()); - } - - record_high_priority_heal_escalation(candidate_type, priority, result); - error!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_HEAL_ADMISSION, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_HEAL, - candidate_type, - bucket, - object = object.unwrap_or(""), - priority = heal_priority_label(priority), - admission = result.result_label(), - reason = result.reason_label(), - state = "high_priority_not_admitted", - "Scanner high-priority heal admission failed" - ); - Err(build_high_priority_heal_admission_error(candidate_type, bucket, object, priority, result)) -} - /// Scanner item representing a file during scanning #[derive(Clone, Debug)] pub struct ScannerItem { @@ -660,6 +699,12 @@ impl ScannerItem { self.object_name = split.last().unwrap_or(&"").to_string(); } + fn metadata_object_path(&self) -> String { + let mut item = self.clone(); + item.transform_meta_dir(); + item.object_path() + } + pub async fn apply_actions( &mut self, object_infos: Vec, @@ -1148,6 +1193,38 @@ impl ScannerItem { } } +fn classify_get_size_failure(item: &ScannerItem, err: &StorageError) -> GetSizeFailureAction { + if matches!(err, StorageError::Io(io) if io.to_string() == SCANNER_SKIP_FILE_ERROR) { + return GetSizeFailureAction::Skip; + } + + if is_scanner_metadata_corrupt_error(err) { + return GetSizeFailureAction::HealMetadata { + object: item.metadata_object_path(), + }; + } + + if is_scanner_metadata_transient_error(err) { + return GetSizeFailureAction::RecordFailed; + } + + GetSizeFailureAction::RecordFailed +} + +fn data_usage_root_has_progress(root: &DataUsageEntry) -> bool { + !root.children.is_empty() + || root.size > 0 + || root.objects > 0 + || root.versions > 0 + || root.delete_markers > 0 + || root.failed_objects > 0 + || root.replication_stats.is_some() +} + +fn partial_cache_is_useful(root: &DataUsageEntry, pending_heals_changed: bool) -> bool { + data_usage_root_has_progress(root) || pending_heals_changed +} + /// Folder scanner for scanning directory structures pub struct FolderScanner { root: String, @@ -1175,6 +1252,7 @@ pub struct FolderScanner { budget: Arc, skip_heal: Arc, local_disk: Arc, + pending_heals_changed: bool, } impl FolderScanner { @@ -1214,6 +1292,118 @@ impl FolderScanner { } } + fn sync_pending_heals(&mut self) { + self.update_cache.info.pending_heals = self.new_cache.info.pending_heals.clone(); + self.pending_heals_changed = true; + } + + fn clear_pending_scanner_heal( + &mut self, + kind: PendingScannerHealKind, + bucket: &str, + object: Option<&str>, + version_id: Option<&str>, + ) { + let before = self.new_cache.info.pending_heals.len(); + self.new_cache + .info + .pending_heals + .retain(|entry| !pending_scanner_heal_matches(entry, kind, bucket, object, version_id)); + if self.new_cache.info.pending_heals.len() != before { + self.sync_pending_heals(); + } + } + + fn record_pending_scanner_heal( + &mut self, + kind: PendingScannerHealKind, + bucket: &str, + object: Option<&str>, + version_id: Option<&str>, + scan_mode: HealScanMode, + result: HealAdmissionResult, + ) { + let now = Self::now_secs(); + if let Some(entry) = self + .new_cache + .info + .pending_heals + .iter_mut() + .find(|entry| pending_scanner_heal_matches(entry, kind, bucket, object, version_id)) + { + entry.last_attempt = now; + entry.attempts = entry.attempts.saturating_add(1); + entry.last_admission_result = result.result_label().to_string(); + entry.last_admission_reason = result.reason_label().to_string(); + self.sync_pending_heals(); + return; + } + + self.new_cache.info.pending_heals.push(PendingScannerHeal { + kind, + bucket: bucket.to_string(), + object: object.map(ToOwned::to_owned), + version_id: version_id.map(ToOwned::to_owned), + scan_mode, + first_seen: now, + last_attempt: now, + attempts: 1, + last_admission_result: result.result_label().to_string(), + last_admission_reason: result.reason_label().to_string(), + }); + self.prune_pending_scanner_heals(); + self.sync_pending_heals(); + } + + fn prune_pending_scanner_heals(&mut self) { + let len = self.new_cache.info.pending_heals.len(); + if len <= MAX_PENDING_SCANNER_HEALS_PER_BUCKET { + return; + } + + sort_pending_scanner_heals_for_retry(&mut self.new_cache.info.pending_heals); + let remove_count = len.saturating_sub(MAX_PENDING_SCANNER_HEALS_PER_BUCKET); + self.new_cache.info.pending_heals.drain(..remove_count); + counter!( + METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL, + "bucket" => self.new_cache.info.name.clone() + ) + .increment(u64::try_from(remove_count).unwrap_or(u64::MAX)); + warn!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_HEAL_ADMISSION, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_HEAL, + bucket = %self.new_cache.info.name, + pruned = remove_count, + remaining = self.new_cache.info.pending_heals.len(), + state = "pending_heal_pruned", + "Scanner pending heal ledger pruned oldest entries" + ); + } + + fn update_pending_scanner_heal_after_admission( + &mut self, + kind: PendingScannerHealKind, + bucket: &str, + object: Option<&str>, + version_id: Option<&str>, + scan_mode: HealScanMode, + result: HealAdmissionResult, + ) { + match result { + HealAdmissionResult::Accepted | HealAdmissionResult::Merged => { + self.clear_pending_scanner_heal(kind, bucket, object, version_id); + } + HealAdmissionResult::Full | HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull) => { + self.record_pending_scanner_heal(kind, bucket, object, version_id, scan_mode, result); + } + HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped) => { + self.clear_pending_scanner_heal(kind, bucket, object, version_id); + } + } + } + fn prune_failed_objects_cache(&mut self) { let ttl = self.failed_object_ttl_secs; if ttl == 0 { @@ -1312,6 +1502,95 @@ impl FolderScanner { true } + async fn send_required_scanner_heal_request( + &mut self, + kind: PendingScannerHealKind, + bucket: String, + object: Option, + version_id: Option, + request: HealChannelRequest, + ) -> Result<(), ScannerError> { + let candidate_type = pending_scanner_heal_candidate_type(kind); + let priority = request.priority; + let scan_mode = request.scan_mode.unwrap_or(self.scan_mode); + let result = send_scanner_heal_request(candidate_type, request).await?; + self.update_pending_scanner_heal_after_admission( + kind, + &bucket, + object.as_deref(), + version_id.as_deref(), + scan_mode, + result, + ); + if result.is_admitted() { + return Ok(()); + } + + record_high_priority_heal_escalation(candidate_type, priority, result); + let admission_error = + build_high_priority_heal_admission_error(candidate_type, &bucket, object.as_deref(), priority, result); + error!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_HEAL_ADMISSION, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_HEAL, + candidate_type, + bucket = %bucket, + object = object.as_deref().unwrap_or(""), + priority = heal_priority_label(priority), + admission = result.result_label(), + reason = result.reason_label(), + error = %admission_error, + state = "high_priority_not_admitted", + "Scanner high-priority heal admission failed" + ); + Ok(()) + } + + async fn retry_pending_scanner_heals(&mut self) -> Result<(), ScannerError> { + if !self.should_heal().await { + return Ok(()); + } + + let bucket = self.new_cache.info.name.clone(); + for pending in pending_scanner_heal_retry_candidates(&self.new_cache.info.pending_heals, &bucket) { + if !self.should_heal().await { + break; + } + + let Some(request) = build_pending_scanner_heal_request(&pending) else { + self.clear_pending_scanner_heal(pending.kind, &pending.bucket, None, pending.version_id.as_deref()); + counter!( + METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL, + "bucket" => pending.bucket.clone(), + "type" => pending_scanner_heal_candidate_type(pending.kind).to_string() + ) + .increment(1); + warn!( + target: "rustfs::scanner::folder", + event = EVENT_SCANNER_HEAL_ADMISSION, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_HEAL, + bucket = %pending.bucket, + state = "pending_heal_malformed", + "Scanner dropped malformed pending heal entry" + ); + continue; + }; + + self.send_required_scanner_heal_request( + pending.kind, + pending.bucket.clone(), + pending.object.clone(), + pending.version_id.clone(), + request, + ) + .await?; + } + + Ok(()) + } + /// Set heal object select probability pub fn set_heal_object_select(&mut self, prob: u32) { self.heal_object_select = prob; @@ -1677,9 +1956,9 @@ impl FolderScanner { let sz = match self.local_disk.get_size(item.clone()).await { Ok(sz) => sz, Err(e) => { - let is_skip_file = matches!(e, StorageError::Io(ref io) if io.to_string() == SCANNER_SKIP_FILE_ERROR); + let failure_action = classify_get_size_failure(&item, &e); - if !is_skip_file { + if failure_action != GetSizeFailureAction::Skip { // Track failed objects to prevent infinite retry loops into.failed_objects += 1; self.record_failed(&item.path); @@ -1699,6 +1978,23 @@ impl FolderScanner { } } + if let GetSizeFailureAction::HealMetadata { object } = failure_action { + self.send_required_scanner_heal_request( + PendingScannerHealKind::Object, + item.bucket.clone(), + Some(object.clone()), + None, + build_object_heal_request( + item.bucket.clone(), + object.clone(), + None, + self.scan_mode, + HealChannelPriority::High, + ), + ) + .await?; + } + timer.sleep().await; continue; } @@ -1976,9 +2272,10 @@ impl FolderScanner { let (bucket, prefix) = path2_bucket_object(name.as_str()); if bucket != resolver.bucket { - send_required_scanner_heal_request( - "bucket", - &bucket, + self.send_required_scanner_heal_request( + PendingScannerHealKind::Bucket, + bucket.clone(), + None, None, build_bucket_heal_request(bucket.clone(), HealChannelPriority::High), ) @@ -2145,10 +2442,11 @@ impl FolderScanner { error = %e, "Scanner list_path_raw failed to resolve file versions" ); - send_required_scanner_heal_request( - "object", - &bucket, - Some(&entry.name), + self.send_required_scanner_heal_request( + PendingScannerHealKind::Object, + bucket.clone(), + Some(entry.name.clone()), + None, build_object_heal_request( bucket.clone(), entry.name.clone(), @@ -2164,14 +2462,16 @@ impl FolderScanner { }; for fiv in fivs.versions { - send_required_scanner_heal_request( - "object", - &bucket, - Some(&entry.name), + let version_id = fiv.version_id.and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) }); + self.send_required_scanner_heal_request( + PendingScannerHealKind::Object, + bucket.clone(), + Some(entry.name.clone()), + version_id.clone(), build_object_heal_request( bucket.clone(), entry.name.clone(), - fiv.version_id.and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) }), + version_id, self.scan_mode, HealChannelPriority::High, ), @@ -2399,6 +2699,7 @@ pub async fn scan_data_folder( budget: budget.clone(), skip_heal, local_disk, + pending_heals_changed: false, }; // Check if context is cancelled @@ -2406,6 +2707,8 @@ pub async fn scan_data_folder( return Err(ScannerError::Other("Operation cancelled".to_string())); } + scanner.retry_pending_scanner_heals().await?; + // Read top level in bucket let mut root = DataUsageEntry::default(); let folder = CachedFolder { @@ -2434,22 +2737,21 @@ pub async fn scan_data_folder( } Err(e) => { if ctx.is_cancelled() { - let root_has_progress = !root.children.is_empty() - || root.size > 0 - || root.objects > 0 - || root.versions > 0 - || root.delete_markers > 0 - || root.failed_objects > 0 - || root.replication_stats.is_some(); + let root_has_progress = data_usage_root_has_progress(&root); + let pending_heals_changed = scanner.pending_heals_changed; let new_cache = scanner.as_mut_new_cache(); if root_has_progress { new_cache.replace_hashed(&hash_path(&cache.info.name), &None, &root); } - if new_cache.root().is_some() { - new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); + if partial_cache_is_useful(&root, pending_heals_changed) { + if new_cache.root().is_some() { + new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); + } new_cache.info.last_update = Some(SystemTime::now()); new_cache.info.next_cycle = cache.info.next_cycle; - set_scan_checkpoint(new_cache, checkpoint_reason_from_budget(budget.reason())); + if root_has_progress { + set_scan_checkpoint(new_cache, checkpoint_reason_from_budget(budget.reason())); + } close_disk().await; return Err(ScannerError::PartialCache(Box::new(new_cache.clone()))); } @@ -2515,6 +2817,7 @@ mod tests { budget: ScannerCycleBudget::new(&CancellationToken::new(), Default::default()), skip_heal: Arc::new(AtomicBool::new(false)), local_disk: disk, + pending_heals_changed: false, }; (scanner, temp_dir) @@ -2578,6 +2881,61 @@ mod tests { assert!(!scanner.should_skip_failed("path2")); } + #[test] + fn test_classify_get_size_failure_marks_metadata_heal_object_path() { + let temp_dir = std::env::temp_dir(); + let file_type = std::fs::metadata(&temp_dir) + .expect("temp dir metadata should be readable") + .file_type(); + let item = ScannerItem { + path: temp_dir.join("bucket/dir/object/xl.meta").to_string_lossy().to_string(), + bucket: "bucket".to_string(), + prefix: "dir/object".to_string(), + object_name: "xl.meta".to_string(), + file_type, + lifecycle: None, + replication: None, + heal_enabled: false, + heal_bitrot: false, + debug: false, + }; + let err = StorageError::other(format!("{}: corrupt metadata", crate::scanner_io::SCANNER_METADATA_CORRUPT_ERROR)); + + let action = classify_get_size_failure(&item, &err); + + assert_eq!( + action, + GetSizeFailureAction::HealMetadata { + object: "dir/object".to_string() + } + ); + } + + #[test] + fn test_classify_get_size_failure_records_transient_metadata_error_without_heal() { + let temp_dir = std::env::temp_dir(); + let file_type = std::fs::metadata(&temp_dir) + .expect("temp dir metadata should be readable") + .file_type(); + let item = ScannerItem { + path: temp_dir.join("bucket/dir/object/xl.meta").to_string_lossy().to_string(), + bucket: "bucket".to_string(), + prefix: "dir/object".to_string(), + object_name: "xl.meta".to_string(), + file_type, + lifecycle: None, + replication: None, + heal_enabled: false, + heal_bitrot: false, + debug: false, + }; + let err = StorageError::other(format!("{}: temporary read failure", crate::scanner_io::SCANNER_METADATA_TRANSIENT_ERROR)); + + let action = classify_get_size_failure(&item, &err); + + assert_eq!(action, GetSizeFailureAction::RecordFailed); + } + #[test] fn test_should_account_replication_stats_only_for_live_object_versions() { let live = ObjectInfo::default(); @@ -3130,6 +3488,199 @@ mod tests { assert_eq!(request.recreate_missing, Some(false)); } + fn pending_heal( + kind: PendingScannerHealKind, + bucket: &str, + object: Option<&str>, + version_id: Option<&str>, + last_attempt: u64, + attempts: u32, + ) -> PendingScannerHeal { + PendingScannerHeal { + kind, + bucket: bucket.to_string(), + object: object.map(ToOwned::to_owned), + version_id: version_id.map(ToOwned::to_owned), + scan_mode: HealScanMode::Deep, + first_seen: 1, + last_attempt, + attempts, + last_admission_result: "full".to_string(), + last_admission_reason: "none".to_string(), + } + } + + #[test] + fn test_pending_heal_reconstructs_bucket_request() { + let pending = pending_heal(PendingScannerHealKind::Bucket, "bucket", None, None, 1, 1); + + let request = build_pending_scanner_heal_request(&pending).expect("bucket request should rebuild"); + + assert_eq!(request.bucket, "bucket"); + assert_eq!(request.priority, HealChannelPriority::High); + assert_eq!(request.source, HealRequestSource::Scanner); + assert_eq!(request.recreate_missing, Some(false)); + assert!(request.object_prefix.is_none()); + } + + #[test] + fn test_pending_heal_reconstructs_object_request_with_version() { + let pending = pending_heal(PendingScannerHealKind::Object, "bucket", Some("path/to/object"), Some("version-a"), 1, 1); + + let request = build_pending_scanner_heal_request(&pending).expect("object request should rebuild"); + + assert_eq!(request.bucket, "bucket"); + assert_eq!(request.object_prefix.as_deref(), Some("path/to/object")); + assert_eq!(request.object_version_id.as_deref(), Some("version-a")); + assert_eq!(request.scan_mode, Some(HealScanMode::Deep)); + assert_eq!(request.priority, HealChannelPriority::High); + assert_eq!(request.source, HealRequestSource::Scanner); + } + + #[test] + fn test_pending_heal_retry_candidates_respect_cap_and_order() { + let pending: Vec = (0..(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET + 2)) + .map(|idx| { + pending_heal( + PendingScannerHealKind::Object, + "bucket", + Some(&format!("object-{idx:03}")), + None, + idx as u64, + 1, + ) + }) + .collect(); + + let candidates = pending_scanner_heal_retry_candidates(&pending, "bucket"); + + assert_eq!(candidates.len(), MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET); + assert_eq!(candidates.first().and_then(|entry| entry.object.as_deref()), Some("object-000")); + assert_eq!(candidates.last().and_then(|entry| entry.object.as_deref()), Some("object-127")); + } + + #[tokio::test] + async fn test_pending_heal_queue_full_deduplicates_object_entry() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard::new(u64::MAX, usize::MAX, &mut scanner, temp_dir); + scanner.new_cache.info.name = "bucket".to_string(); + scanner.update_cache.info.name = "bucket".to_string(); + + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + Some("version-a"), + HealScanMode::Deep, + HealAdmissionResult::Full, + ); + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + Some("version-a"), + HealScanMode::Deep, + HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull), + ); + + assert_eq!(scanner.new_cache.info.pending_heals.len(), 1); + let pending = &scanner.new_cache.info.pending_heals[0]; + assert_eq!(pending.object.as_deref(), Some("object")); + assert_eq!(pending.version_id.as_deref(), Some("version-a")); + assert_eq!(pending.attempts, 2); + assert_eq!(pending.last_admission_result, "dropped"); + assert_eq!(pending.last_admission_reason, "queue_full"); + assert_eq!(scanner.update_cache.info.pending_heals, scanner.new_cache.info.pending_heals); + assert!(scanner.pending_heals_changed); + } + + #[tokio::test] + async fn test_pending_heal_admitted_results_clear_matching_entry() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard::new(u64::MAX, usize::MAX, &mut scanner, temp_dir); + + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + None, + HealScanMode::Normal, + HealAdmissionResult::Full, + ); + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + None, + HealScanMode::Normal, + HealAdmissionResult::Accepted, + ); + + assert!(scanner.new_cache.info.pending_heals.is_empty()); + + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Bucket, + "bucket", + None, + None, + HealScanMode::Normal, + HealAdmissionResult::Full, + ); + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Bucket, + "bucket", + None, + None, + HealScanMode::Normal, + HealAdmissionResult::Merged, + ); + + assert!(scanner.new_cache.info.pending_heals.is_empty()); + } + + #[tokio::test] + async fn test_pending_heal_policy_dropped_clears_without_creating_entry() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard::new(u64::MAX, usize::MAX, &mut scanner, temp_dir); + + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + None, + HealScanMode::Normal, + HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped), + ); + assert!(scanner.new_cache.info.pending_heals.is_empty()); + + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + None, + HealScanMode::Normal, + HealAdmissionResult::Full, + ); + scanner.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + "bucket", + Some("object"), + None, + HealScanMode::Normal, + HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped), + ); + + assert!(scanner.new_cache.info.pending_heals.is_empty()); + } + + #[test] + fn test_partial_cache_is_useful_when_pending_heals_changed() { + let root = DataUsageEntry::default(); + + assert!(!partial_cache_is_useful(&root, false)); + assert!(partial_cache_is_useful(&root, true)); + } + fn metadata_for_object(bucket: &str, object: &str) -> Vec { let mut meta = FileMeta::new(); meta.add_version(FileInfo { diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 7f4c4e5ce..1274e4784 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -29,7 +29,7 @@ use rustfs_config::{ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCU use rustfs_filemeta::FileMeta; use rustfs_utils::path::path_join_buf; use s3s::dto::{BucketLifecycleConfiguration, ReplicationConfiguration}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::path::Path; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::{LazyLock, Mutex as StdMutex, MutexGuard}; @@ -52,6 +52,8 @@ use crate::{ }; pub(crate) const SCANNER_SKIP_FILE_ERROR: &str = "skip file"; +pub(crate) const SCANNER_METADATA_CORRUPT_ERROR: &str = "scanner metadata corrupt"; +pub(crate) const SCANNER_METADATA_TRANSIENT_ERROR: &str = "scanner metadata transient"; const LOG_COMPONENT_SCANNER: &str = "scanner"; const LOG_SUBSYSTEM_IO: &str = "io"; const EVENT_SCANNER_DISK_BUCKET_STATE: &str = "scanner_disk_bucket_state"; @@ -70,17 +72,43 @@ const METRIC_SCANNER_DISK_BUCKET_SCANS_QUEUED: &str = "rustfs_scanner_disk_bucke pub type DirtyUsageBuckets = HashMap; +pub(crate) fn is_scanner_metadata_corrupt_error(err: &StorageError) -> bool { + matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_CORRUPT_ERROR)) +} + +pub(crate) fn is_scanner_metadata_transient_error(err: &StorageError) -> bool { + matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_TRANSIENT_ERROR)) +} + +fn scanner_metadata_corrupt_error(reason: impl std::fmt::Display, bucket: &str, object_path: &str) -> StorageError { + StorageError::other(format!( + "{SCANNER_METADATA_CORRUPT_ERROR}: {reason}, bucket={bucket}, object_path={object_path}" + )) +} + +fn scanner_metadata_transient_error(reason: impl std::fmt::Display, bucket: &str, object_path: &str) -> StorageError { + StorageError::other(format!( + "{SCANNER_METADATA_TRANSIENT_ERROR}: {reason}, bucket={bucket}, object_path={object_path}" + )) +} + #[derive(Clone)] pub struct ScannerBucketScanPlan { buckets: Vec, dirty_usage_buckets: Arc, + failed_dirty_buckets: Arc>>, } impl ScannerBucketScanPlan { - fn new(buckets: Vec, dirty_usage_buckets: Arc) -> Self { + fn new( + buckets: Vec, + dirty_usage_buckets: Arc, + failed_dirty_buckets: Arc>>, + ) -> Self { Self { buckets, dirty_usage_buckets, + failed_dirty_buckets, } } } @@ -165,6 +193,32 @@ fn clear_dirty_usage_buckets(snapshot: &DirtyUsageBuckets) { .record_scanner_dirty_usage_cycle_clear(usize_to_u64_saturated(cleared_buckets), usize_to_u64_saturated(pending_buckets)); } +fn dirty_usage_buckets_excluding_failed(snapshot: &DirtyUsageBuckets, failed_buckets: &HashSet) -> DirtyUsageBuckets { + snapshot + .iter() + .filter(|(bucket, _)| !failed_buckets.contains(*bucket)) + .map(|(bucket, generation)| (bucket.clone(), *generation)) + .collect() +} + +fn should_clear_dirty_usage_snapshot( + result_ok: bool, + completed_all_sets: bool, + budget_elapsed: bool, + dirty_buckets: &DirtyUsageBuckets, + failed_buckets: &HashSet, +) -> Option { + if result_ok && completed_all_sets && !budget_elapsed { + return Some(dirty_usage_buckets_excluding_failed(dirty_buckets, failed_buckets)); + } + + None +} + +async fn record_failed_dirty_bucket(failed_buckets: &Arc>>, bucket: &str) { + failed_buckets.lock().await.insert(bucket.to_string()); +} + fn dirty_usage_snapshot_covers_current(snapshot: &DirtyUsageBuckets) -> bool { let dirty_buckets = dirty_usage_buckets(); dirty_buckets.iter().all(|(bucket, generation)| { @@ -664,6 +718,7 @@ impl ScannerIO for ECStore { let set_scan_limit = scanner_max_concurrent_set_scans(total_results); let dirty_usage_buckets = Arc::new(snapshot_dirty_usage_buckets(&all_buckets)); + let failed_dirty_buckets = Arc::new(Mutex::new(HashSet::::new())); record_set_scan_concurrency_limit(set_scan_limit); debug!( target: "rustfs::scanner::io", @@ -717,7 +772,8 @@ impl ScannerIO for ECStore { }); wait_futs.push(receiver_fut); - let scan_plan = ScannerBucketScanPlan::new(all_buckets.clone(), dirty_usage_buckets.clone()); + let scan_plan = + ScannerBucketScanPlan::new(all_buckets.clone(), dirty_usage_buckets.clone(), failed_dirty_buckets.clone()); // Spawn task to run the scanner let scanner_fut = tokio::spawn(async move { let permit_wait = child_token_clone.clone(); @@ -861,8 +917,15 @@ impl ScannerIO for ECStore { let results = results_mutex.lock().await.clone(); let completed_all_sets = results.iter().all(|result| result.info.last_update.is_some()); let result = finalize_nsscanner_result(&results, first_err); - if result.is_ok() && completed_all_sets && !budget.budget_elapsed() { - clear_dirty_usage_buckets(&dirty_usage_buckets); + let failed_buckets = failed_dirty_buckets.lock().await.clone(); + if let Some(clear_snapshot) = should_clear_dirty_usage_snapshot( + result.is_ok(), + completed_all_sets, + budget.budget_elapsed(), + &dirty_usage_buckets, + &failed_buckets, + ) { + clear_dirty_usage_buckets(&clear_snapshot); } result } @@ -883,6 +946,7 @@ impl ScannerIOCache for SetDisks { let ScannerBucketScanPlan { buckets, dirty_usage_buckets, + failed_dirty_buckets, } = scan_plan; let pool_label = self.pool_index.to_string(); let set_label = self.set_index.to_string(); @@ -952,6 +1016,7 @@ impl ScannerIOCache for SetDisks { } if let Err(e) = bucket_tx.send(bucket.clone()).await { + record_failed_dirty_bucket(&failed_dirty_buckets, &bucket.name).await; error!( target: "rustfs::scanner::io", event = EVENT_SCANNER_SET_STATE, @@ -1012,6 +1077,7 @@ impl ScannerIOCache for SetDisks { let active_disk_bucket_scans_clone = active_disk_bucket_scans.clone(); let pool_label_clone = pool_label.clone(); let set_label_clone = set_label.clone(); + let failed_dirty_buckets_clone = failed_dirty_buckets.clone(); futs.push(tokio::spawn(async move { loop { let Some(bucket) = bucket_rx_mutex_clone.lock().await.recv().await else { @@ -1119,6 +1185,7 @@ impl ScannerIOCache for SetDisks { { Ok(scan_outcome) => scan_outcome, Err(e) => { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; if ctx_clone.is_cancelled() { debug!( target: "rustfs::scanner::io", @@ -1170,6 +1237,7 @@ impl ScannerIOCache for SetDisks { cache = match scan_outcome { ScannerDiskScanOutcome::Complete(cache) => cache, ScannerDiskScanOutcome::Partial(cache) => { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; let done_save = Metrics::time(Metric::SaveUsage); let partial_saved = match cache.save(store_clone_clone.clone(), cache_name.as_str()).await { Ok(()) => true, @@ -1233,6 +1301,7 @@ impl ScannerIOCache for SetDisks { ); if let Err(e) = send_cache_root_entry_info(&bucket_result_tx_clone_clone, &cache).await { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; error!( target: "rustfs::scanner::io", event = EVENT_SCANNER_DATA_USAGE_STREAM, @@ -1247,6 +1316,7 @@ impl ScannerIOCache for SetDisks { let done_save = Metrics::time(Metric::SaveUsage); if let Err(e) = cache.save(store_clone_clone.clone(), &cache_name).await { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; error!( target: "rustfs::scanner::io", event = EVENT_SCANNER_CACHE_PERSIST_STATE, @@ -1325,17 +1395,19 @@ impl ScannerIODisk for Disk { return Err(StorageError::other(SCANNER_SKIP_FILE_ERROR.to_string())); } Err(e) => { - return Err(StorageError::other(format!( - "failed to read metadata: {e}, bucket={}, object_path={}", + return Err(scanner_metadata_transient_error( + format!("failed to read metadata: {e}"), &item.bucket, - &item.object_path() - ))); + &item.object_path(), + )); } }; item.transform_meta_dir(); - let meta = FileMeta::load(&data)?; + let meta = FileMeta::load(&data).map_err(|e| { + scanner_metadata_corrupt_error(format!("failed to load metadata: {e}"), &item.bucket, &item.object_path()) + })?; let fivs = match meta.get_file_info_versions(item.bucket.as_str(), item.object_path().as_str(), false) { Ok(versions) => versions, Err(e) => { @@ -1350,7 +1422,11 @@ impl ScannerIODisk for Disk { error = %e, "Scanner disk bucket failed to resolve file info versions" ); - return Err(StorageError::other(SCANNER_SKIP_FILE_ERROR.to_string())); + return Err(scanner_metadata_corrupt_error( + format!("failed to resolve file info versions: {e}"), + &item.bucket, + &item.object_path(), + )); } }; @@ -1594,6 +1670,38 @@ mod tests { clear_dirty_usage_buckets_for_tests(); } + #[test] + #[serial] + fn dirty_usage_clear_excludes_failed_buckets() { + clear_dirty_usage_buckets_for_tests(); + record_dirty_usage_bucket("photos"); + record_dirty_usage_bucket("videos"); + let buckets = vec![bucket_info("photos"), bucket_info("videos")]; + let snapshot = snapshot_dirty_usage_buckets(&buckets); + let failed_buckets = HashSet::from(["videos".to_string()]); + let clear_snapshot = dirty_usage_buckets_excluding_failed(&snapshot, &failed_buckets); + + clear_dirty_usage_buckets(&clear_snapshot); + + let dirty_buckets = dirty_usage_buckets(); + assert!(!dirty_buckets.contains_key("photos")); + assert!(dirty_buckets.contains_key("videos")); + drop(dirty_buckets); + clear_dirty_usage_buckets_for_tests(); + } + + #[test] + fn dirty_usage_clear_plan_excludes_cache_save_failures() { + let snapshot = DirtyUsageBuckets::from([("photos".to_string(), 1), ("videos".to_string(), 2)]); + let failed_buckets = HashSet::from(["videos".to_string()]); + + let clear_snapshot = should_clear_dirty_usage_snapshot(true, true, false, &snapshot, &failed_buckets) + .expect("successful completed cycle should produce a clear snapshot"); + + assert!(clear_snapshot.contains_key("photos")); + assert!(!clear_snapshot.contains_key("videos")); + } + #[test] #[serial] fn clear_dirty_usage_bucket_removes_deleted_bucket_marker() { @@ -1804,6 +1912,61 @@ mod tests { let _ = tokio::fs::remove_dir_all(&temp_dir).await; } + #[tokio::test] + async fn get_size_marks_corrupt_metadata_for_heal() { + let temp_dir = std::env::temp_dir().join(format!("rustfs-scanner-corrupt-meta-{}", Uuid::new_v4())); + let bucket = "bucket"; + let object = "object"; + let object_dir = temp_dir.join(bucket).join(object); + let metadata_path = object_dir.join(STORAGE_FORMAT_FILE); + + tokio::fs::create_dir_all(&object_dir) + .await + .expect("failed to create object directory"); + tokio::fs::write(&metadata_path, b"not-valid-filemeta") + .await + .expect("failed to write corrupt metadata"); + + let endpoint = Endpoint::try_from(temp_dir.to_string_lossy().as_ref()).expect("failed to create endpoint"); + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("failed to open local disk"); + + let relative_path = metadata_path.to_string_lossy().to_string(); + let (_, scanner_path) = path2_bucket_object_with_base_path(temp_dir.to_string_lossy().as_ref(), relative_path.as_str()); + let file_type = tokio::fs::metadata(&metadata_path) + .await + .expect("failed to stat metadata") + .file_type(); + + let item = ScannerItem { + path: scanner_path, + bucket: bucket.to_string(), + prefix: object.to_string(), + object_name: STORAGE_FORMAT_FILE.to_string(), + file_type, + lifecycle: None, + replication: None, + heal_enabled: false, + heal_bitrot: false, + debug: false, + }; + + let err = disk + .get_size(item) + .await + .expect_err("corrupt metadata should be surfaced as scanner-heal work"); + assert!(is_scanner_metadata_corrupt_error(&err)); + + let _ = tokio::fs::remove_dir_all(&temp_dir).await; + } + #[tokio::test] async fn get_size_counts_delete_markers_separately_from_versions() { let temp_dir = std::env::temp_dir().join(format!("rustfs-scanner-versioned-usage-{}", Uuid::new_v4())); diff --git a/rustfs/src/admin/handlers/scanner.rs b/rustfs/src/admin/handlers/scanner.rs index 166ca0bc1..1a308ecec 100644 --- a/rustfs/src/admin/handlers/scanner.rs +++ b/rustfs/src/admin/handlers/scanner.rs @@ -17,6 +17,8 @@ use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::current_scanner_metrics_report; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, RemoteAddr}; +use crate::startup_background::{ENV_SCANNER_ENABLED, scanner_enabled_from_env}; +use chrono::Utc; use http::{HeaderMap, HeaderValue}; use hyper::{Method, StatusCode}; use matchit::Params; @@ -31,10 +33,63 @@ const JSON_CONTENT_TYPE: &str = "application/json"; #[derive(Debug, Serialize)] struct ScannerStatusResponse { + enabled: bool, + disabled_reason: Option, + freshness: ScannerFreshnessStatus, metrics: ScannerMetricsReport, runtime_config: rustfs_scanner::runtime_config::ScannerRuntimeConfigStatus, } +#[derive(Debug, Serialize)] +struct ScannerFreshnessStatus { + state: &'static str, + last_cycle_end_unix_secs: u64, + max_expected_age_seconds: u64, + reason: Option<&'static str>, +} + +fn scanner_disabled_reason(enabled: bool) -> Option { + (!enabled).then(|| format!("disabled by {ENV_SCANNER_ENABLED}")) +} + +fn scanner_freshness_status( + metrics: &ScannerMetricsReport, + runtime_config: &rustfs_scanner::runtime_config::ScannerRuntimeConfigStatus, +) -> ScannerFreshnessStatus { + const FRESHNESS_MULTIPLIER: u64 = 2; + + let max_expected_age_seconds = runtime_config + .cycle_interval_seconds + .value + .saturating_mul(FRESHNESS_MULTIPLIER); + if metrics.last_cycle_end_unix_secs == 0 { + return ScannerFreshnessStatus { + state: "unknown", + last_cycle_end_unix_secs: 0, + max_expected_age_seconds, + reason: Some("no completed cycle recorded"), + }; + } + + let now = Utc::now().timestamp().max(0) as u64; + let age = now.saturating_sub(metrics.last_cycle_end_unix_secs); + if max_expected_age_seconds > 0 && age > max_expected_age_seconds { + return ScannerFreshnessStatus { + state: "stale", + last_cycle_end_unix_secs: metrics.last_cycle_end_unix_secs, + max_expected_age_seconds, + reason: Some("last cycle is older than freshness window"), + }; + } + + ScannerFreshnessStatus { + state: "fresh", + last_cycle_end_unix_secs: metrics.last_cycle_end_unix_secs, + max_expected_age_seconds, + reason: None, + } +} + pub fn register_scanner_route(r: &mut S3Router) -> std::io::Result<()> { r.insert( Method::GET, @@ -84,9 +139,16 @@ pub struct ScannerStatusHandler {} impl Operation for ScannerStatusHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let _cred = validate_scanner_status_request(&req).await?; + let enabled = scanner_enabled_from_env(); + let metrics = current_scanner_metrics_report().await; + let runtime_config = rustfs_scanner::scanner_runtime_config_status(); + let freshness = scanner_freshness_status(&metrics, &runtime_config); let response = ScannerStatusResponse { - metrics: current_scanner_metrics_report().await, - runtime_config: rustfs_scanner::scanner_runtime_config_status(), + enabled, + disabled_reason: scanner_disabled_reason(enabled), + freshness, + metrics, + runtime_config, }; let body = serde_json::to_vec(&response).map_err(|err| { S3Error::with_message(S3ErrorCode::InternalError, format!("failed to encode scanner status: {err}")) @@ -95,3 +157,44 @@ impl Operation for ScannerStatusHandler { json_response(body) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn scanner_disabled_reason_reports_startup_env_key() { + assert_eq!(scanner_disabled_reason(true), None); + assert_eq!(scanner_disabled_reason(false), Some(format!("disabled by {ENV_SCANNER_ENABLED}"))); + } + + #[test] + fn scanner_freshness_reports_unknown_without_cycle_end() { + let metrics = ScannerMetricsReport::default(); + let mut runtime_config = rustfs_scanner::scanner_runtime_config_status(); + runtime_config.cycle_interval_seconds.value = 60; + + let freshness = scanner_freshness_status(&metrics, &runtime_config); + + assert_eq!(freshness.state, "unknown"); + assert_eq!(freshness.last_cycle_end_unix_secs, 0); + assert_eq!(freshness.max_expected_age_seconds, 120); + assert_eq!(freshness.reason, Some("no completed cycle recorded")); + } + + #[test] + fn scanner_freshness_reports_stale_after_window() { + let metrics = ScannerMetricsReport { + last_cycle_end_unix_secs: Utc::now().timestamp().max(0) as u64 - 121, + ..Default::default() + }; + let mut runtime_config = rustfs_scanner::scanner_runtime_config_status(); + runtime_config.cycle_interval_seconds.value = 60; + + let freshness = scanner_freshness_status(&metrics, &runtime_config); + + assert_eq!(freshness.state, "stale"); + assert_eq!(freshness.max_expected_age_seconds, 120); + assert_eq!(freshness.reason, Some("last cycle is older than freshness window")); + } +} diff --git a/rustfs/src/app/context.rs b/rustfs/src/app/context.rs index 5ddddfad5..19215ecb5 100644 --- a/rustfs/src/app/context.rs +++ b/rustfs/src/app/context.rs @@ -845,6 +845,10 @@ mod tests { }; let endpoint_pools = EndpointServerPools(vec![pool_endpoints]); + if let Some(store) = crate::storage::storage_api::ecstore_global::new_object_layer_fn() { + return (temp_dir, store, endpoint_pools); + } + init_local_disks(endpoint_pools.clone()).await.expect("test local disks"); let store = ECStore::new( "127.0.0.1:0".parse().expect("test addr"), diff --git a/rustfs/src/startup_background.rs b/rustfs/src/startup_background.rs index 3ba720d59..fc4889040 100644 --- a/rustfs/src/startup_background.rs +++ b/rustfs/src/startup_background.rs @@ -22,18 +22,22 @@ use rustfs_utils::get_env_bool_with_aliases; use std::{io::Result, sync::Arc}; use tracing::{debug, info}; -const ENV_SCANNER_ENABLED: &str = "RUSTFS_SCANNER_ENABLED"; -const ENV_SCANNER_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_SCANNER"; +pub(crate) const ENV_SCANNER_ENABLED: &str = "RUSTFS_SCANNER_ENABLED"; +pub(crate) const ENV_SCANNER_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_SCANNER"; const ENV_HEAL_ENABLED: &str = "RUSTFS_HEAL_ENABLED"; const ENV_HEAL_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_HEAL"; const LOG_COMPONENT_MAIN: &str = "main"; const LOG_SUBSYSTEM_STARTUP: &str = "startup"; const EVENT_BACKGROUND_SERVICES_CONFIGURED: &str = "background_services_configured"; +pub(crate) fn scanner_enabled_from_env() -> bool { + get_env_bool_with_aliases(ENV_SCANNER_ENABLED, &[ENV_SCANNER_ENABLED_DEPRECATED], true) +} + pub(crate) async fn init_background_service_runtime(store: Arc) -> Result { let _ = create_ahm_services_cancel_token(); - let enable_scanner = get_env_bool_with_aliases(ENV_SCANNER_ENABLED, &[ENV_SCANNER_ENABLED_DEPRECATED], true); + let enable_scanner = scanner_enabled_from_env(); let enable_heal = get_env_bool_with_aliases(ENV_HEAL_ENABLED, &[ENV_HEAL_ENABLED_DEPRECATED], true); info!( diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 4fc9bf9c1..5f09d0ae3 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -375,6 +375,8 @@ pub(crate) mod ecstore_event { } pub(crate) mod ecstore_global { + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::global::new_object_layer_fn; pub(crate) use rustfs_ecstore::api::global::{ GLOBAL_BOOT_TIME, GLOBAL_TierConfigMgr, get_global_bucket_monitor, get_global_deployment_id, get_global_endpoints_opt, get_global_lock_client, get_global_lock_clients, get_global_region, get_global_tier_config_mgr, global_rustfs_port, diff --git a/rustfs/src/workload_admission.rs b/rustfs/src/workload_admission.rs index 319b6a39a..32e13d10e 100644 --- a/rustfs/src/workload_admission.rs +++ b/rustfs/src/workload_admission.rs @@ -27,7 +27,6 @@ const REPAIR_QUEUE_BACKLOG_PRESENT: &str = "repair queue has pending work"; const REPLICATION_RUNTIME_NOT_INITIALIZED: &str = "replication runtime not initialized"; const REPLICATION_QUEUE_BACKLOG_PRESENT: &str = "replication queue has pending work"; const REPLICATION_QUEUE_STATS_UNAVAILABLE: &str = "replication queue stats unavailable"; -const SCANNER_ADMISSION_DISABLED: &str = "scanner admission disabled because max concurrent set scans is zero"; const SCANNER_ADMISSION_SATURATED: &str = "scanner active work reached configured set-scan limit"; const SCANNER_ACTIVITY_IDLE_OR_NOT_INITIALIZED: &str = "scanner activity idle or not initialized"; const STORAGE_CONCURRENCY_PROVIDER_MISSING_FOREGROUND_READ: &str = @@ -114,9 +113,8 @@ pub fn scanner_workload_admission_snapshot() -> WorkloadAdmissionSnapshot { } fn scanner_workload_admission_snapshot_from_activity(active: u64, limit: usize) -> WorkloadAdmissionSnapshot { - let state = if limit == 0 { - AdmissionState::Disabled - } else if usize::try_from(active).ok().is_some_and(|active| active >= limit) { + let effective_limit = if limit == 0 { None } else { Some(limit) }; + let state = if effective_limit.is_some_and(|limit| usize::try_from(active).ok().is_some_and(|active| active >= limit)) { AdmissionState::Saturated } else if active > 0 { AdmissionState::Open @@ -127,11 +125,10 @@ fn scanner_workload_admission_snapshot_from_activity(active: u64, limit: usize) let snapshot = WorkloadAdmissionSnapshot::new(WorkloadClass::Scanner, state).with_counts( Some(u64_to_usize_saturated(active)), None, - Some(limit), + effective_limit, ); match state { - AdmissionState::Disabled => snapshot.with_reason(SCANNER_ADMISSION_DISABLED), AdmissionState::Saturated => snapshot.with_reason(SCANNER_ADMISSION_SATURATED), AdmissionState::Unknown => snapshot.with_reason(SCANNER_ACTIVITY_IDLE_OR_NOT_INITIALIZED), _ => snapshot, @@ -302,13 +299,24 @@ mod tests { } #[test] - fn scanner_snapshot_reports_disabled_when_set_scan_limit_is_zero() { + fn scanner_snapshot_treats_zero_set_scan_limit_as_topology_derived() { let snapshot = scanner_workload_admission_snapshot_from_activity(0, 0); assert_eq!(snapshot.class, WorkloadClass::Scanner); - assert_eq!(snapshot.state, AdmissionState::Disabled); - assert_eq!(snapshot.limit, Some(0)); - assert_eq!(snapshot.reason.as_deref(), Some(SCANNER_ADMISSION_DISABLED)); + assert_eq!(snapshot.state, AdmissionState::Unknown); + assert_eq!(snapshot.limit, None); + assert_eq!(snapshot.reason.as_deref(), Some(SCANNER_ACTIVITY_IDLE_OR_NOT_INITIALIZED)); + } + + #[test] + fn scanner_snapshot_with_topology_derived_limit_reports_active_work_open() { + let snapshot = scanner_workload_admission_snapshot_from_activity(4, 0); + + assert_eq!(snapshot.class, WorkloadClass::Scanner); + assert_eq!(snapshot.state, AdmissionState::Open); + assert_eq!(snapshot.active, Some(4)); + assert_eq!(snapshot.limit, None); + assert_eq!(snapshot.reason, None); } #[test]