diff --git a/crates/filemeta/src/metacache.rs b/crates/filemeta/src/metacache.rs index 462ce6fae..8a3499b58 100644 --- a/crates/filemeta/src/metacache.rs +++ b/crates/filemeta/src/metacache.rs @@ -21,6 +21,7 @@ use arc_swap::ArcSwapOption; use rmp::Marker; use serde::{Deserialize, Serialize}; use std::cmp::Ordering; +use std::collections::{HashMap, HashSet}; use std::str::from_utf8; use std::{ fmt::Debug, @@ -37,8 +38,10 @@ use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; use tokio::spawn; use tokio::sync::Mutex; use tracing::{debug, warn}; +use uuid::Uuid; const SLASH_SEPARATOR: &str = "/"; +pub const MAX_META_CACHE_HEAL_CANDIDATES: usize = 1024; #[derive(Clone, Debug, Default)] pub struct MetadataResolutionParams { @@ -66,6 +69,41 @@ pub struct MetaCacheEntry { pub reusable: bool, } +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub enum MetaCacheHealCandidateKind { + Object, + DeleteMarker, + UnversionedObject, +} + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct MetaCacheHealCandidate { + pub object: String, + pub version_id: Option, + pub kind: MetaCacheHealCandidateKind, + /// Number of raw disk entries that carried this validated version. + pub replica_count: usize, +} + +impl MetaCacheHealCandidate { + pub fn validated_version(&self) -> Option { + match self.kind { + MetaCacheHealCandidateKind::Object | MetaCacheHealCandidateKind::DeleteMarker => self.version_id, + MetaCacheHealCandidateKind::UnversionedObject => None, + } + } + + pub fn is_unversioned(&self) -> bool { + self.kind == MetaCacheHealCandidateKind::UnversionedObject + } +} + +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub struct MetaCacheHealDiscovery { + pub candidates: Vec, + pub unverified_count: usize, +} + impl MetaCacheEntry { pub fn marshal_msg(&self) -> Result> { let mut wr = Vec::new(); @@ -370,6 +408,176 @@ impl MetaCacheEntries { }) } + /// Discover validated object/delete-marker versions and safe unversioned + /// inspection candidates in the raw entries without applying read quorum. + /// This is intentionally separate from [`Self::resolve`]: a sub-quorum + /// version is a valid heal target even though it must not participate in + /// normal reads or writes. + /// + /// The validated list is bounded and deduplicated by object, version id, + /// and metadata kind; each candidate retains the number of raw disk + /// entries that carried it so callers can classify sub-quorum versions. + /// Entries whose xl.meta cannot be decoded are counted separately for + /// discovery accounting; they never become versionless destructive heal + /// requests and do not consume the validated quota. An + /// [`MetaCacheHealCandidateKind::UnversionedObject`] is always consumed by + /// a non-destructive scanner request. + pub fn discover_heal_candidates(&self, bucket: &str, max_candidates: usize) -> MetaCacheHealDiscovery { + let limit = max_candidates.min(MAX_META_CACHE_HEAL_CANDIDATES); + if limit == 0 || bucket.is_empty() { + return MetaCacheHealDiscovery::default(); + } + + let mut discovery = MetaCacheHealDiscovery { + candidates: Vec::::with_capacity(limit.min(self.0.len())), + unverified_count: 0, + }; + let mut seen: HashMap<(String, Option, MetaCacheHealCandidateKind), usize> = + HashMap::with_capacity(limit.min(self.0.len())); + + for entry in self.0.iter().flatten() { + if !valid_heal_candidate_name(bucket, entry) { + continue; + } + + let meta = match FileMeta::load(&entry.metadata) { + Ok(meta) => meta, + Err(_) => { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + }; + let mut entry_seen = HashSet::new(); + + for shallow in meta.versions { + let version = match shallow.parse_version_meta() { + Ok(version) if version.valid() => version, + Ok(_) | Err(_) => { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + }; + if version.free_version() { + continue; + } + + let payload_header = version.header(); + if normalize_version_id(shallow.header.version_id) != normalize_version_id(payload_header.version_id) + || shallow.header.version_type != payload_header.version_type + { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + + let (kind, version_id) = match version.version_type { + VersionType::Object + if version.object.is_some() && version.delete_marker.is_none() && version.legacy_object.is_none() => + { + match version.object.as_ref().and_then(|object| object.version_id) { + Some(id) if !id.is_nil() => (MetaCacheHealCandidateKind::Object, Some(id)), + Some(_) | None => (MetaCacheHealCandidateKind::UnversionedObject, None), + } + } + VersionType::Delete + if version.delete_marker.is_some() && version.object.is_none() && version.legacy_object.is_none() => + { + let Some(id) = version.delete_marker.as_ref().and_then(|marker| marker.version_id) else { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + }; + if id.is_nil() { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + (MetaCacheHealCandidateKind::DeleteMarker, Some(id)) + } + VersionType::Legacy + if version.legacy_object.is_some() && version.object.is_none() && version.delete_marker.is_none() => + { + let Some(legacy) = version.legacy_object.as_ref() else { + continue; + }; + if legacy.version_id.is_empty() { + (MetaCacheHealCandidateKind::UnversionedObject, None) + } else { + let Ok(id) = Uuid::parse_str(&legacy.version_id) else { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + }; + if id.is_nil() { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + (MetaCacheHealCandidateKind::Object, Some(id)) + } + } + _ => { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + }; + + if normalize_version_id(payload_header.version_id) != version_id { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + + // `all_parts=true` is the trust-boundary check for versioned + // candidates. A null/legacy object may still need the old + // non-destructive inspection fallback when its part arrays + // are parseable but incomplete; never use that fallback for + // a candidate carrying a real version id. + let file_info = match version.clone().into_fileinfo(bucket, &entry.name, true) { + Ok(file_info) => file_info, + Err(_) if version_id.is_none() && matches!(kind, MetaCacheHealCandidateKind::UnversionedObject) => { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + match version.into_fileinfo(bucket, &entry.name, false) { + Ok(file_info) => file_info, + Err(_) => continue, + } + } + Err(_) => { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + }; + if file_info.volume != bucket || file_info.name != entry.name { + discovery.unverified_count = discovery.unverified_count.saturating_add(1); + continue; + } + + let candidate = MetaCacheHealCandidate { + object: entry.name.clone(), + version_id, + kind, + replica_count: 1, + }; + let key = (candidate.object.clone(), candidate.version_id, candidate.kind.clone()); + if !entry_seen.contains(&key) { + // Keep per-entry dedupe bounded as well as the global + // candidate union. Once the cap is reached, only keys + // already present in the global map may update replica + // counts; novel versions are accounting-only. + if entry_seen.len() >= limit && !seen.contains_key(&key) { + continue; + } + entry_seen.insert(key.clone()); + } + if let Some(index) = seen.get(&key).copied() { + discovery.candidates[index].replica_count = discovery.candidates[index].replica_count.saturating_add(1); + } else { + seen.insert(key, discovery.candidates.len()); + discovery.candidates.push(candidate); + if discovery.candidates.len() >= limit { + return discovery; + } + } + } + } + + discovery + } + fn resolve_inner(&self, mut params: MetadataResolutionParams, enforce_write_quorum: bool) -> Option { if self.0.is_empty() { debug!( @@ -546,6 +754,28 @@ impl MetaCacheEntries { } } +fn valid_heal_candidate_name(bucket: &str, entry: &MetaCacheEntry) -> bool { + if bucket.is_empty() || entry.name.is_empty() || entry.is_dir() || entry.name.contains('\0') { + return false; + } + + // Validate raw key components without normalizing them. A dot component + // could otherwise escape the bucket when the key is later mapped back to + // a disk path; a final empty component is retained for valid keys ending + // in '/'. + let mut components = entry.name.split('/').peekable(); + while let Some(component) = components.next() { + if component == "." || component == ".." || (component.is_empty() && components.peek().is_some()) { + return false; + } + } + true +} + +fn normalize_version_id(version_id: Option) -> Option { + version_id.filter(|id| !id.is_nil()) +} + #[derive(Debug, Default)] pub struct MetaCacheEntriesSortedResult { pub entries: Option, @@ -991,7 +1221,7 @@ impl Cache { mod tests { use super::*; use crate::test_data::create_real_xlmeta; - use crate::{FileMetaVersion, MetaDeleteMarker, TRANSITION_COMPLETE}; + use crate::{FileMetaVersion, MetaDeleteMarker, MetaObjectV1, MetaObjectV1Erasure, MetaObjectV1Stat, TRANSITION_COMPLETE}; use std::collections::HashMap; use std::io::Cursor; use std::sync::{ @@ -1592,6 +1822,277 @@ mod tests { ); } + #[test] + fn discover_heal_candidates_keeps_sub_quorum_versions_and_deduplicates() { + let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"); + let entries = MetaCacheEntries(vec![ + Some(metacache_entry_single_version(1, now, "one")), + Some(metacache_entry_single_version(2, now, "two")), + Some(metacache_entry_single_version(2, now, "two")), + Some(metacache_entry_single_version(3, now, "three")), + ]); + + let discovery = entries.discover_heal_candidates("bucket", 16); + let ids: std::collections::HashSet = discovery + .candidates + .iter() + .filter_map(|candidate| candidate.version_id) + .collect(); + assert_eq!( + ids, + [Uuid::from_u128(1), Uuid::from_u128(2), Uuid::from_u128(3)] + .into_iter() + .collect() + ); + assert_eq!(discovery.candidates.len(), 3, "duplicate tied versions must be emitted once"); + assert_eq!( + discovery + .candidates + .iter() + .find(|candidate| candidate.version_id == Some(Uuid::from_u128(2))) + .expect("duplicate version should be discovered") + .replica_count, + 2 + ); + } + + #[test] + fn discover_heal_candidates_covers_divergent_quorum_boundaries_n2_n4_n6() { + let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"); + + for (disk_count, quorum) in [(2usize, 1usize), (4, 2), (6, 3)] { + let target_id = Uuid::from_u128(0x1000 + disk_count as u128); + for target_replicas in [quorum.saturating_sub(1), quorum, quorum + 1] { + let entries = (0..disk_count) + .map(|disk| { + let version_id = if disk < target_replicas { + target_id + } else { + Uuid::from_u128(0x2000 + disk as u128) + }; + Some(metacache_entry_single_version(version_id.as_u128(), now, "divergent")) + }) + .collect(); + let discovery = MetaCacheEntries(entries).discover_heal_candidates("bucket", 32); + let target = discovery + .candidates + .iter() + .find(|candidate| candidate.version_id == Some(target_id)); + assert_eq!(target.is_some(), target_replicas > 0, "N={disk_count}, replicas={target_replicas}"); + if let Some(target) = target { + assert_eq!(target.replica_count, target_replicas); + } + } + } + } + + #[test] + fn discover_heal_candidates_separates_delete_markers_and_preserves_unversioned_objects() { + let mut marker_meta = FileMeta::new(); + marker_meta + .add_version(FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + version_id: Some(Uuid::from_u128(99)), + deleted: true, + mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")), + ..Default::default() + }) + .expect("delete marker should be added"); + let marker = MetaCacheEntry { + name: "object".to_string(), + metadata: marker_meta.marshal_msg().expect("delete marker metadata should marshal"), + cached: Some(marker_meta), + reusable: false, + }; + + let unversioned_entry = metacache_entry_with_mod_time( + OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"), + "unversioned", + ); + let discovery = MetaCacheEntries(vec![Some(marker), Some(unversioned_entry)]).discover_heal_candidates("bucket", 16); + assert!(discovery.candidates.iter().any(|candidate| { + candidate.kind == MetaCacheHealCandidateKind::DeleteMarker && candidate.version_id == Some(Uuid::from_u128(99)) + })); + assert!(discovery.candidates.iter().any(|candidate| { + candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none() + })); + } + + #[test] + fn discover_heal_candidates_skips_free_versions() { + let object_id = Uuid::from_u128(100); + let free_id = Uuid::from_u128(101); + let mut meta = FileMeta::new(); + meta.add_version(FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + version_id: Some(object_id), + transition_status: TRANSITION_COMPLETE.to_string(), + transitioned_objname: "remote/object".to_string(), + transition_version_id: Some(Uuid::from_u128(102)), + transition_tier: "WARM".to_string(), + mod_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }) + .expect("transitioned object should be added"); + let mut delete = FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + version_id: Some(object_id), + mod_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }; + delete.set_tier_free_version_id(&free_id.to_string()); + meta.delete_version(&delete).expect("free version should be persisted"); + + let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry { + name: "object".to_string(), + metadata: meta.marshal_msg().expect("free version metadata should marshal"), + cached: Some(meta), + reusable: false, + })]) + .discover_heal_candidates("bucket", 16); + assert!(discovery.candidates.is_empty()); + } + + #[test] + fn discover_heal_candidates_preserves_unversioned_legacy_object() { + let legacy = MetaObjectV1 { + version: "1.0.1".to_string(), + format: "xl".to_string(), + stat: MetaObjectV1Stat { + size: 1, + mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")), + name: "object".to_string(), + ..Default::default() + }, + erasure: MetaObjectV1Erasure { + data_blocks: 4, + parity_blocks: 2, + index: 1, + distribution: vec![1, 2, 3, 4, 5, 6], + ..Default::default() + }, + ..Default::default() + }; + let version = FileMetaVersion { + version_type: VersionType::Legacy, + legacy_object: Some(legacy), + ..Default::default() + }; + let mut meta = FileMeta::new(); + meta.versions + .push(FileMetaShallowVersion::try_from(version).expect("legacy metadata should marshal")); + let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry { + name: "object".to_string(), + metadata: meta.marshal_msg().expect("legacy metadata should marshal"), + cached: Some(meta), + reusable: false, + })]) + .discover_heal_candidates("bucket", 16); + assert!(discovery.candidates.iter().any(|candidate| { + candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none() + })); + } + + #[test] + fn discover_heal_candidates_rejects_nil_and_malformed_metadata_and_is_bounded() { + let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"); + let mut nil = metacache_entry_single_version(1, now, "nil"); + let mut nil_meta = FileMeta::load(&nil.metadata).expect("nil fixture should decode"); + let mut nil_version = nil_meta.versions[0] + .parse_version_meta() + .expect("nil fixture version should decode"); + nil_version.object.as_mut().expect("object fixture").version_id = Some(Uuid::nil()); + nil_meta.versions[0] = FileMetaShallowVersion::try_from(nil_version).expect("nil fixture should marshal"); + nil.metadata = nil_meta.marshal_msg().expect("nil fixture metadata should marshal"); + + let mut mismatched = metacache_entry_single_version(2, now, "mismatched"); + let mut mismatched_meta = FileMeta::load(&mismatched.metadata).expect("mismatched fixture should decode"); + mismatched_meta.versions[0].header.version_id = Some(Uuid::from_u128(200)); + mismatched.metadata = mismatched_meta.marshal_msg().expect("mismatched metadata should marshal"); + + let mut short_parts = metacache_entry_single_version(3, now, "short-parts"); + let mut short_parts_meta = FileMeta::load(&short_parts.metadata).expect("short-parts fixture should decode"); + let mut short_parts_version = short_parts_meta.versions[0] + .parse_version_meta() + .expect("short-parts fixture version should decode"); + let object = short_parts_version.object.as_mut().expect("object fixture"); + object.part_numbers = vec![1]; + object.part_actual_sizes = vec![1]; + object.part_sizes.clear(); + short_parts_meta.versions[0] = FileMetaShallowVersion::try_from(short_parts_version).expect("short-parts should marshal"); + short_parts.metadata = short_parts_meta.marshal_msg().expect("short-parts metadata should marshal"); + + let mut short_unversioned = metacache_entry_with_mod_time(now, "short-unversioned"); + let mut short_unversioned_meta = + FileMeta::load(&short_unversioned.metadata).expect("short-unversioned fixture should decode"); + let mut short_unversioned_version = short_unversioned_meta.versions[0] + .parse_version_meta() + .expect("short-unversioned version should decode"); + let unversioned_object = short_unversioned_version.object.as_mut().expect("unversioned object fixture"); + unversioned_object.part_numbers = vec![1]; + unversioned_object.part_actual_sizes = vec![1]; + unversioned_object.part_sizes.clear(); + short_unversioned_meta.versions[0] = + FileMetaShallowVersion::try_from(short_unversioned_version).expect("short-unversioned should marshal"); + short_unversioned.metadata = short_unversioned_meta + .marshal_msg() + .expect("short-unversioned metadata should marshal"); + + let mut malformed = nil.clone(); + malformed.name = "malformed".to_string(); + malformed.metadata = vec![1, 2, 3]; + + let entries = MetaCacheEntries( + std::iter::once(Some(nil)) + .chain(std::iter::once(Some(mismatched))) + .chain(std::iter::once(Some(short_parts))) + .chain(std::iter::once(Some(short_unversioned))) + .chain(std::iter::once(Some(malformed))) + .chain((0..32).map(|id| Some(metacache_entry_single_version(id + 10, now, "bounded")))) + .collect(), + ); + let discovery = entries.discover_heal_candidates("bucket", 5); + assert!(discovery.candidates.len() <= 5); + assert!( + !discovery + .candidates + .iter() + .any(|candidate| candidate.version_id == Some(Uuid::nil())) + ); + assert!( + !discovery + .candidates + .iter() + .any(|candidate| candidate.version_id == Some(Uuid::from_u128(2))) + ); + assert!( + !discovery + .candidates + .iter() + .any(|candidate| candidate.version_id == Some(Uuid::from_u128(3))) + ); + assert!(discovery.candidates.iter().any(|candidate| { + candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none() + })); + assert!( + discovery.unverified_count >= 1, + "malformed and rejected metadata must remain observable during discovery" + ); + + for invalid_name in ["../object", "object//", "object\0name"] { + let mut entry = metacache_entry_single_version(400, now, invalid_name); + entry.name = invalid_name.to_string(); + let discovery = MetaCacheEntries(vec![Some(entry)]).discover_heal_candidates("bucket", 5); + assert!( + discovery.candidates.is_empty(), + "invalid key should not become a heal candidate: {invalid_name:?}" + ); + } + } + #[test] fn resolve_rejects_partial_latest_and_returns_committed_previous_metadata() { let old_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"); diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 693dd6ebc..4486c247b 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -43,7 +43,7 @@ use rustfs_common::metrics::{ UpdateCurrentPathFn, current_path_updater, global_metrics, }; use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit, trace_subscriber_count}; -use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams}; +use rustfs_filemeta::{MAX_META_CACHE_HEAL_CANDIDATES, MetaCacheEntries, MetaCacheEntry, MetaCacheHealCandidateKind}; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, VersioningConfiguration}; use time::OffsetDateTime; @@ -96,6 +96,10 @@ const METRIC_SCANNER_EXCESS_OBJECT_VERSION_SIZE_TOTAL: &str = "rustfs_scanner_ex 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 METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL: &str = "rustfs_scanner_heal_discovery_candidates_total"; +const METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL: &str = "rustfs_scanner_heal_discovery_sub_quorum_total"; +const METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL: &str = "rustfs_scanner_heal_discovery_unverified_total"; +const METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL: &str = "rustfs_scanner_heal_discovery_queued_total"; const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128; // --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) -- @@ -883,7 +887,7 @@ impl FolderScanner { object: Option, version_id: Option, request: HealChannelRequest, - ) -> Result<(), ScannerError> { + ) -> Result { let candidate_type = pending_scanner_heal_candidate_type(kind); let priority = request.priority; let scan_mode = request.scan_mode.unwrap_or(self.scan_mode); @@ -911,7 +915,7 @@ impl FolderScanner { error = %err, "Scanner deferred heal request after channel error" ); - return Ok(()); + return Ok(HealAdmissionResult::Full); } }; self.update_pending_scanner_heal_after_admission( @@ -923,7 +927,7 @@ impl FolderScanner { result, ); if result.is_admitted() { - return Ok(()); + return Ok(result); } record_high_priority_heal_escalation(candidate_type, priority, result); @@ -944,7 +948,7 @@ impl FolderScanner { state = "high_priority_not_admitted", "Scanner high-priority heal admission failed" ); - Ok(()) + Ok(result) } pub fn set_heal_object_select(&mut self, prob: u32) { @@ -1736,14 +1740,7 @@ impl FolderScanner { break; } - let mut resolver = MetadataResolutionParams { - dir_quorum: self.disks_quorum, - obj_quorum: self.disks_quorum, - bucket: "".to_string(), - strict: false, - ..Default::default() - }; - + let mut previous_bucket = String::new(); for name in abandoned_children { if !self.should_heal().await { break; @@ -1751,7 +1748,7 @@ impl FolderScanner { let (bucket, prefix) = path2_bucket_object(name.as_str()); - if bucket != resolver.bucket { + if bucket != previous_bucket { self.send_required_scanner_heal_request( PendingScannerHealKind::Bucket, bucket.clone(), @@ -1760,10 +1757,9 @@ impl FolderScanner { build_bucket_heal_request(bucket.clone(), HealChannelPriority::High), ) .await?; + previous_bucket = bucket.clone(); } - resolver.bucket = bucket.clone(); - let child_ctx = ctx.child_token(); let (agreed_tx, mut agreed_rx) = mpsc::channel::(1); @@ -1880,6 +1876,7 @@ impl FolderScanner { let mut agreed_closed = false; let mut partial_closed = false; let mut finished_closed = false; + let mut seen_heal_candidates: HashSet<(String, Option, MetaCacheHealCandidateKind)> = HashSet::new(); loop { if agreed_closed && partial_closed && finished_closed { @@ -1904,65 +1901,56 @@ impl FolderScanner { break; } - let Some(entry) = resolve_object_heal_entry(&entries, resolver.clone()) else { - continue; - }; + let discovery = entries.discover_heal_candidates(&bucket, MAX_META_CACHE_HEAL_CANDIDATES); + counter!(METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL) + .increment(u64::try_from(discovery.candidates.len()).unwrap_or(u64::MAX)); + counter!(METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL).increment( + u64::try_from( + discovery + .candidates + .iter() + .filter(|candidate| candidate.replica_count < disks_quorum) + .count(), + ) + .unwrap_or(u64::MAX), + ); + counter!(METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL).increment( + u64::try_from(discovery.unverified_count).unwrap_or(u64::MAX), + ); - (self.update_current_path)(&entry.name).await; - - if entry.is_dir() { - continue; - } - - let fivs = match entry.file_info_versions(&bucket) { - Ok(fivs) => fivs, - Err(e) => { - error!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_FOLDER_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_FOLDER, - bucket = %bucket, - entry = %entry.name, - state = "file_info_versions_failed", - error = %e, - "Scanner list_path_raw failed to resolve file versions" - ); - self.send_required_scanner_heal_request( - PendingScannerHealKind::Object, - bucket.clone(), - Some(entry.name.clone()), - None, - build_object_heal_request( - bucket.clone(), - entry.name.clone(), - None, - self.scan_mode, - HealChannelPriority::High, - ), - ) - .await?; - found_objects = true; + for candidate in discovery.candidates { + let version_id = candidate.validated_version().map(|id| id.to_string()); + let identity = (candidate.object.clone(), version_id.clone(), candidate.kind.clone()); + if seen_heal_candidates.len() >= MAX_META_CACHE_HEAL_CANDIDATES + && !seen_heal_candidates.contains(&identity) + { continue; } - }; - - for fiv in fivs.versions { - 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( + if !seen_heal_candidates.insert(identity) { + continue; + } + let mut request = build_object_heal_request( + bucket.clone(), + candidate.object.clone(), + version_id.clone(), + self.scan_mode, + HealChannelPriority::High, + ); + if candidate.is_unversioned() { + request.remove_corrupted = Some(false); + } + (self.update_current_path)(&candidate.object).await; + let admission = 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(), - version_id, - self.scan_mode, - HealChannelPriority::High, - ), + Some(candidate.object.clone()), + version_id, + request, ) .await?; + if admission.is_admitted() { + counter!(METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL).increment(1); + } found_objects = true; } diff --git a/crates/scanner/src/scanner_folder/item_actions.rs b/crates/scanner/src/scanner_folder/item_actions.rs index 15bd6736b..f5cb3507a 100644 --- a/crates/scanner/src/scanner_folder/item_actions.rs +++ b/crates/scanner/src/scanner_folder/item_actions.rs @@ -13,6 +13,8 @@ // limitations under the License. /// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers. use super::*; +#[cfg(test)] +use rustfs_filemeta::MetadataResolutionParams; /// Cached folder information for scanning #[derive(Clone, Debug)] @@ -88,6 +90,7 @@ pub(super) fn build_object_heal_request( } } +#[cfg(test)] pub(super) fn resolve_object_heal_entry( entries: &MetaCacheEntries, resolver: MetadataResolutionParams, diff --git a/crates/scanner/src/scanner_folder/ledger.rs b/crates/scanner/src/scanner_folder/ledger.rs index 7ac03983c..040821b82 100644 --- a/crates/scanner/src/scanner_folder/ledger.rs +++ b/crates/scanner/src/scanner_folder/ledger.rs @@ -305,13 +305,17 @@ pub(super) fn build_pending_scanner_heal_request(entry: &PendingScannerHeal) -> 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( + let mut request = build_object_heal_request( entry.bucket.clone(), object.clone(), entry.version_id.clone(), entry.scan_mode, HealChannelPriority::High, - ) + ); + if entry.version_id.is_none() { + request.remove_corrupted = Some(false); + } + request }), } } diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 503b0e75f..67eafed19 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -17,7 +17,7 @@ use crate::SCANNER_SLEEPER; use super::*; use crate::storage_api::VersionPurgeStatusType; use crate::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, TierStats, new_disk, storageclass}; -use rustfs_filemeta::{FileInfo, FileMeta}; +use rustfs_filemeta::{FileInfo, FileMeta, MetadataResolutionParams}; use std::io::Write; #[cfg(unix)] use std::os::unix::fs::{PermissionsExt, symlink}; @@ -1121,6 +1121,17 @@ fn test_pending_heal_reconstructs_object_request_with_version() { assert_eq!(request.source, HealRequestSource::Scanner); } +#[test] +fn test_pending_heal_reconstructs_unversioned_request_without_removal() { + let pending = pending_heal(PendingScannerHealKind::Object, "bucket", Some("object"), None, 1, 1); + + let request = build_pending_scanner_heal_request(&pending).expect("unversioned object request should rebuild"); + + assert!(request.object_version_id.is_none()); + assert_eq!(request.remove_corrupted, Some(false)); + assert_eq!(request.recreate_missing, Some(false)); +} + #[test] fn test_pending_heal_retry_candidates_respect_cap_and_order() { let pending: Vec = (0..(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET + 2)) @@ -1323,6 +1334,20 @@ fn metadata_for_object(bucket: &str, object: &str) -> Vec { meta.marshal_msg().expect("test metadata should marshal") } +fn metadata_for_object_version(bucket: &str, object: &str, version_id: Option) -> Vec { + let mut file_info = FileInfo::new(object, 4, 2); + file_info.volume = bucket.to_string(); + file_info.name = object.to_string(); + file_info.version_id = version_id; + file_info.versioned = version_id.is_some(); + file_info.mod_time = Some(OffsetDateTime::now_utc()); + file_info.size = 1; + + let mut meta = FileMeta::new(); + meta.add_version(file_info).expect("test metadata version should be accepted"); + meta.marshal_msg().expect("test metadata should marshal") +} + async fn write_test_object_metadata(root: &std::path::Path, bucket: &str, object: &str) { write_test_object_metadata_bytes(root, bucket, object, &metadata_for_object(bucket, object)).await; } @@ -1726,12 +1751,21 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() { let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone()); let heal_starts = Arc::new(AtomicUsize::new(0)); let heal_starts_clone = heal_starts.clone(); + let healed_versions = Arc::new(Mutex::new(Vec::>::new())); + let healed_versions_clone = healed_versions.clone(); let mut heal_rx = rustfs_common::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests"); let _heal_responder = tokio::spawn(async move { while let Some(command) = heal_rx.recv().await { - if let rustfs_common::heal_channel::HealChannelCommand::Start { response_tx, .. } = command { + if let rustfs_common::heal_channel::HealChannelCommand::Start { + request, response_tx, .. + } = command + { heal_starts_clone.fetch_add(1, Ordering::Relaxed); + healed_versions_clone + .lock() + .expect("heal version capture lock should not be poisoned") + .push(request.object_version_id); let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); } } @@ -1739,13 +1773,18 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() { let bucket = "src-archive"; let object = "snapshots/37b3f20d941e2f5e6d99114d9bb2f3e67a8a2e5c9c4c5a1b0d6e7f8091a2b3c4"; - let metadata = metadata_for_object(bucket, object); - write_test_object_metadata_bytes(&temp_dir, bucket, object, &metadata).await; + let orphan_version = Uuid::from_u128(0x1934); + let shared_version = Uuid::from_u128(0x1935); + let orphan_metadata = metadata_for_object_version(bucket, object, Some(orphan_version)); + let shared_metadata = metadata_for_object_version(bucket, object, Some(shared_version)); + write_test_object_metadata_bytes(&temp_dir, bucket, object, &orphan_metadata).await; + let mut expected_metadata = vec![(temp_dir.join(bucket).join(object).join("xl.meta"), orphan_metadata.clone())]; let mut disks = vec![scanner.local_disk.clone()]; for disk_name in ["disk2", "disk3", "disk4"] { let disk_root = temp_dir.join(disk_name); - write_test_object_metadata_bytes(&disk_root, bucket, object, &metadata).await; + write_test_object_metadata_bytes(&disk_root, bucket, object, &shared_metadata).await; + expected_metadata.push((disk_root.join(bucket).join(object).join("xl.meta"), shared_metadata.clone())); let endpoint = Endpoint::try_from(disk_root.to_string_lossy().as_ref()).expect("failed to create extra disk endpoint"); let disk = new_disk( &endpoint, @@ -1794,8 +1833,29 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() { .new_cache .checked_flatten(bucket) .expect("healed cache must contain canonical child links"); - assert_eq!(root.objects, 1); + // The fixture intentionally exposes two divergent version histories, so + // the scanner keeps both logical versions visible while discovering heals. + assert_eq!(root.objects, 2); assert!(heal_starts.load(Ordering::Relaxed) > 0, "test must execute the heal child-link path"); + let orphan_version_text = orphan_version.to_string(); + assert!( + healed_versions + .lock() + .expect("heal version capture lock should not be poisoned") + .iter() + .any(|version| version.as_deref() == Some(orphan_version_text.as_str())), + "sub-quorum orphan version must be submitted as an exact heal candidate" + ); + for (path, expected) in expected_metadata { + assert_eq!( + tokio::fs::read(&path) + .await + .expect("scanner discovery must not delete metadata"), + expected, + "scanner discovery must not modify candidate metadata: {}", + path.display() + ); + } } #[tokio::test]