From 23b85b792e60aae3af81ca632d171033c93d90f9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Sat, 22 Aug 2026 22:58:05 +0800 Subject: [PATCH] fix(scanner): unify unknown metadata size accounting --- crates/data-usage/src/data_usage.rs | 90 +- crates/scanner/src/data_usage_define.rs | 55 +- crates/scanner/src/data_usage_define/tests.rs | 39 + crates/scanner/src/scanner_folder.rs | 118 ++- .../src/scanner_folder/item_actions.rs | 834 ++++++++++++++++-- crates/scanner/src/scanner_folder/tests.rs | 60 ++ 6 files changed, 1120 insertions(+), 76 deletions(-) diff --git a/crates/data-usage/src/data_usage.rs b/crates/data-usage/src/data_usage.rs index 9c4692a21..757743b69 100644 --- a/crates/data-usage/src/data_usage.rs +++ b/crates/data-usage/src/data_usage.rs @@ -307,6 +307,38 @@ pub struct DiskUsageStatus { pub snapshot_exists: bool, } +/// A bounded reconciliation record for an object whose logical size could not +/// be trusted at the scanner boundary. The scanner persists these records in +/// its cache; keeping the model here avoids a second, incompatible accounting +/// representation in storage-facing crates. +#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct SizeReconciliationEntry { + /// Stable object/version identity key (not a metrics label). + pub key: String, + pub bucket: String, + pub object: String, + #[serde(default)] + pub version_id: Option, + #[serde(default)] + pub generation: Option, + /// Structured reason label; raw metadata values must never be stored here. + pub reason: String, + #[serde(default)] + pub physical_size: Option, + #[serde(default)] + pub first_seen: u64, + #[serde(default)] + pub attempts: u32, +} + +/// Object scope refreshed by one scanner pass. Existing debts in this scope +/// are removed before the pass's unresolved records are inserted. +#[derive(Debug, Default, Clone, PartialEq, Eq)] +pub struct SizeReconciliationScope { + pub bucket: String, + pub object: String, +} + /// Size summary for a single object or group of objects #[derive(Debug, Default, Clone)] pub struct SizeSummary { @@ -336,6 +368,16 @@ pub struct SizeSummary { pub repl_target_stats: HashMap, /// Per-tier accounting, keyed by storage class or remote tier name pub tier_stats: HashMap, + /// Size-resolution debts observed while scanning this summary. + pub size_reconciliation: Vec, + /// True when the per-object summary exceeded its bounded debt buffer. + /// Callers must retain prior ledger entries rather than treating the + /// partial list as a complete refresh. + pub size_reconciliation_truncated: bool, + /// Object scopes refreshed by this summary. They let the durable ledger + /// remove versions that resolved without allocating one key per healthy + /// version on the hot path. + pub reconciliation_scopes: Vec, } /// Replication target size summary @@ -830,7 +872,8 @@ impl DataUsageEntry { /// /// The canonical wire format is written by the hand-written map-encoded /// `Serialize` on the scanner-side `DataUsageCacheInfo` -/// (`crates/scanner/src/data_usage_define.rs`), which carries 16 fields. +/// (`crates/scanner/src/data_usage_define.rs`), which carries the original 16 +/// fields plus an optional reconciliation field. /// This type decodes only the shared subset and is deliberately not /// `Serialize`: a derived (array) encoding of this 6-field subset would /// corrupt the cache for scanner readers, so no write path may exist here. @@ -1774,6 +1817,51 @@ impl SizeSummary { entry.pending_count = entry.pending_count.saturating_add(stats.pending_count); entry.failed_count = entry.failed_count.saturating_add(stats.failed_count); } + + for entry in &other.size_reconciliation { + self.record_size_reconciliation(entry.clone()); + } + self.size_reconciliation_truncated |= other.size_reconciliation_truncated; + for scope in &other.reconciliation_scopes { + self.record_reconciliation_scope(&scope.bucket, &scope.object); + } + } + + /// Add one reconciliation debt, coalescing repeated observations in the + /// same object summary. The scanner cache applies its own larger bound. + pub fn record_size_reconciliation(&mut self, entry: SizeReconciliationEntry) { + const MAX_SUMMARY_RECONCILIATION_ENTRIES: usize = 1024; + if let Some(existing) = self.size_reconciliation.iter_mut().find(|value| value.key == entry.key) { + existing.reason = entry.reason; + existing.physical_size = entry.physical_size; + existing.generation = entry.generation; + existing.version_id = entry.version_id; + return; + } + if self.size_reconciliation.len() < MAX_SUMMARY_RECONCILIATION_ENTRIES { + self.size_reconciliation.push(entry); + } else { + self.size_reconciliation_truncated = true; + } + } + + /// Mark one object scope as refreshed. Duplicate scopes are suppressed so + /// merging summaries remains bounded and deterministic. + pub fn record_reconciliation_scope(&mut self, bucket: &str, object: &str) { + if !self + .reconciliation_scopes + .iter() + .any(|scope| scope.bucket == bucket && scope.object == object) + { + if self.reconciliation_scopes.len() >= 1024 { + self.size_reconciliation_truncated = true; + return; + } + self.reconciliation_scopes.push(SizeReconciliationScope { + bucket: bucket.to_string(), + object: object.to_string(), + }); + } } } diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index c6ecdd489..02e473c63 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -29,7 +29,8 @@ use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS; pub use rustfs_data_usage::{ AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME, DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry, - PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path, prefix_usage_in_cache, + PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeReconciliationEntry, SizeReconciliationScope, SizeSummary, + TierStats, hash_path, prefix_usage_in_cache, }; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use tokio::time::{Duration, Instant, sleep, timeout}; @@ -192,6 +193,10 @@ const MAX_DATA_USAGE_CACHE_DEPTH: usize = 1024; pub trait ScannerSizeSummaryExt { /// Fold one object's contribution into the summary, including its tier. fn actions_accounting(&mut self, oi: &ObjectInfo, size: i64, actual_size: i64); + /// Fold counters and physical tier usage for an object whose metadata is + /// valid but whose logical size is currently unavailable. Logical totals + /// stay unchanged. + fn actions_accounting_unknown(&mut self, oi: &ObjectInfo); } impl ScannerSizeSummaryExt for SizeSummary { @@ -225,6 +230,34 @@ impl ScannerSizeSummaryExt for SizeSummary { }); } } + + fn actions_accounting_unknown(&mut self, oi: &ObjectInfo) { + if oi.delete_marker { + self.delete_markers = self.delete_markers.saturating_add(1); + return; + } + + if oi.version_id.is_some_and(|v| !v.is_nil()) { + self.versions = self.versions.saturating_add(1); + } + + if oi.transitioned_object.free_version { + return; + } + + let tier = if oi.transitioned_object.status == TRANSITION_COMPLETE { + oi.transitioned_object.tier.clone() + } else { + oi.storage_class.clone().unwrap_or_else(|| storageclass::STANDARD.to_string()) + }; + if let Some(tier_stats) = self.tier_stats.get_mut(&tier) { + *tier_stats = tier_stats.add(&TierStats { + total_size: u64::try_from(oi.size).unwrap_or(0), + num_versions: 1, + num_objects: u64::from(oi.is_latest), + }); + } + } } // ===== Cache-related data structures ===== @@ -344,6 +377,10 @@ pub struct DataUsageCacheInfo { pub scan_plan_digest: Option, #[serde(default)] pub cache_key_format: u16, + /// Bounded durable debts for versions whose logical size was not trusted. + /// The map key is an identity key, never a user-controlled metric label. + #[serde(default)] + pub size_reconciliation: HashMap, } impl Serialize for DataUsageCacheInfo { @@ -353,7 +390,8 @@ impl Serialize for DataUsageCacheInfo { { // Keep this metadata map-encoded so older readers can ignore fields // appended by newer scanner versions during rolling upgrades. - let mut state = serializer.serialize_map(Some(16))?; + let field_count = 16 + usize::from(!self.size_reconciliation.is_empty()); + let mut state = serializer.serialize_map(Some(field_count))?; state.serialize_entry("name", &self.name)?; state.serialize_entry("next_cycle", &self.next_cycle)?; state.serialize_entry("leader_epoch", &self.leader_epoch)?; @@ -370,6 +408,9 @@ impl Serialize for DataUsageCacheInfo { state.serialize_entry("snapshot_complete", &self.snapshot_complete)?; state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?; state.serialize_entry("cache_key_format", &self.cache_key_format)?; + if !self.size_reconciliation.is_empty() { + state.serialize_entry("size_reconciliation", &self.size_reconciliation)?; + } state.end() } } @@ -428,14 +469,18 @@ impl DataUsageCache { self.checked_flatten(name).is_some() }); if !reusable { - let pending_heals = if self.info.name == name { - std::mem::take(&mut self.info.pending_heals) + let (pending_heals, size_reconciliation) = if self.info.name == name { + ( + std::mem::take(&mut self.info.pending_heals), + std::mem::take(&mut self.info.size_reconciliation), + ) } else { - Vec::new() + (Vec::new(), HashMap::new()) }; *self = Self::default(); self.info.name = name.to_string(); self.info.pending_heals = pending_heals; + self.info.size_reconciliation = size_reconciliation; } self.info.next_cycle = next_cycle; diff --git a/crates/scanner/src/data_usage_define/tests.rs b/crates/scanner/src/data_usage_define/tests.rs index ed3fcd544..3c400b000 100644 --- a/crates/scanner/src/data_usage_define/tests.rs +++ b/crates/scanner/src/data_usage_define/tests.rs @@ -673,6 +673,34 @@ fn size_summary_actions_accounting_accumulates_tier_stats() { ); } +#[test] +fn size_summary_unknown_accounting_keeps_physical_tier_and_version_only() { + let mut summary = SizeSummary::default(); + summary + .tier_stats + .insert(storageclass::STANDARD.to_string(), TierStats::default()); + let object = ObjectInfo { + size: 12, + storage_class: Some(storageclass::STANDARD.to_string()), + version_id: Some(uuid::Uuid::new_v4()), + is_latest: true, + ..Default::default() + }; + + summary.actions_accounting_unknown(&object); + + assert_eq!(summary.total_size, 0, "unknown logical size must not become zero or physical bytes"); + assert_eq!(summary.versions, 1); + assert_eq!( + summary.tier_stats.get(storageclass::STANDARD), + Some(&TierStats { + total_size: 12, + num_versions: 1, + num_objects: 1, + }) + ); +} + #[test] fn test_data_usage_entry_merge_sums_failed_objects() { let mut left = DataUsageEntry { @@ -1079,6 +1107,16 @@ fn data_usage_cache_prepare_for_scan_preserves_pending_heal_only_progress() { scan_plan_digest: Some(TEST_PLAN_DIGEST), cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, pending_heals: vec![pending_heal.clone()], + size_reconciliation: HashMap::from([( + "size-key".to_string(), + SizeReconciliationEntry { + key: "size-key".to_string(), + bucket: "bucket".to_string(), + object: "prefix/object".to_string(), + reason: "invalid_declared_size".to_string(), + ..Default::default() + }, + )]), ..Default::default() }, ..Default::default() @@ -1088,6 +1126,7 @@ fn data_usage_cache_prepare_for_scan_preserves_pending_heal_only_progress() { assert_eq!(outcome, DataUsageCachePrepareOutcome::Reused); assert_eq!(cache.info.pending_heals, vec![pending_heal]); + assert!(cache.info.size_reconciliation.contains_key("size-key")); assert!(cache.cache.is_empty()); assert!(!cache.info.snapshot_complete); } diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 693dd6ebc..01c1b8d12 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -20,8 +20,9 @@ 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, PendingScannerHeal, PendingScannerHealKind, ScannerSizeSummaryExt, SizeSummary, hash_path, + DATA_USAGE_SCAN_CHECKPOINT_VERSION, DataUsageCache, DataUsageCacheInfo, DataUsageEntry, DataUsageHash, DataUsageHashMap, + DataUsageScanCheckpoint, DataUsageScanCheckpointReason, PendingScannerHeal, PendingScannerHealKind, ScannerSizeSummaryExt, + SizeReconciliationEntry, SizeSummary, hash_path, }; use crate::error::ScannerError; use crate::runtime_config::{ @@ -97,6 +98,9 @@ const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders 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_SIZE_RECONCILIATION_ENTRIES_PER_BUCKET: usize = 10_000; +const MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET: usize = 8 * 1024 * 1024; +const MAX_SIZE_RECONCILIATION_AGE_SECS: u64 = 7 * 24 * 60 * 60; // --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) -- // @@ -364,7 +368,7 @@ impl PendingScannerAccounting<'_> { fn apply(self, size_summary: &mut SizeSummary, cumulative_size: &mut i64, queued: bool) { let size = if queued { self.expired_size } else { self.retained_size }; size_summary.actions_accounting(self.object, size, self.retained_size); - *cumulative_size += size; + *cumulative_size = cumulative_size.saturating_add(size); } } @@ -675,6 +679,54 @@ pub struct FolderScanner { list_path_raw_options_observer: Option>, } +fn size_reconciliation_entry_bytes(entry: &SizeReconciliationEntry) -> usize { + entry.key.len() + + entry.bucket.len() + + entry.object.len() + + entry.version_id.as_deref().map_or(0, str::len) + + entry.generation.as_deref().map_or(0, str::len) + + entry.reason.len() + + std::mem::size_of::() + + std::mem::size_of::() +} + +fn prune_size_reconciliation(info: &mut DataUsageCacheInfo, now: u64) { + info.size_reconciliation.retain(|key, entry| { + if entry.first_seen == 0 || entry.first_seen > now { + entry.first_seen = now; + } + key == &entry.key + && entry.key.len() <= 4096 + && entry.bucket.len() <= 512 + && entry.object.len() <= 512 + && entry.version_id.as_deref().is_none_or(|value| value.len() <= 64) + && entry.generation.as_deref().is_none_or(|value| value.len() <= 64) + && entry.reason.len() <= 64 + && now.saturating_sub(entry.first_seen) <= MAX_SIZE_RECONCILIATION_AGE_SECS + }); + + while info.size_reconciliation.len() > MAX_SIZE_RECONCILIATION_ENTRIES_PER_BUCKET + || info + .size_reconciliation + .values() + .map(size_reconciliation_entry_bytes) + .sum::() + > MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET + { + let oldest = info + .size_reconciliation + .iter() + .min_by(|(left_key, left), (right_key, right)| { + left.first_seen.cmp(&right.first_seen).then_with(|| left_key.cmp(right_key)) + }) + .map(|(key, _)| key.clone()); + let Some(oldest) = oldest else { + break; + }; + info.size_reconciliation.remove(&oldest); + } +} + impl FolderScanner { fn now_secs() -> u64 { SystemTime::now() @@ -748,6 +800,55 @@ impl FolderScanner { } } + /// Apply the per-object size-resolution ledger updates in one place. The + /// scanner cache is the durable boundary; both working copies are updated + /// so an incremental publication cannot lose a debt or its resolution. + fn apply_size_reconciliation(&mut self, summary: &SizeSummary) { + let now = Self::now_secs(); + // Keep an unresolved identity in place while refreshing its object + // scope. This lets repeated observations increment `attempts`; only + // debts absent from the current pass are considered resolved. + let current_keys = summary + .size_reconciliation + .iter() + .map(|entry| entry.key.clone()) + .collect::>(); + for info in [&mut self.new_cache.info, &mut self.update_cache.info] { + prune_size_reconciliation(info, now); + + if !summary.size_reconciliation_truncated { + for scope in &summary.reconciliation_scopes { + let scope_bucket = item_actions::bounded_reconciliation_field(&scope.bucket); + let scope_object = item_actions::bounded_reconciliation_field(&scope.object); + info.size_reconciliation.retain(|key, entry| { + entry.bucket != scope_bucket || entry.object != scope_object || current_keys.contains(key) + }); + } + } + + for incoming in &summary.size_reconciliation { + if let Some(existing) = info.size_reconciliation.get_mut(&incoming.key) { + existing.reason = incoming.reason.clone(); + existing.physical_size = incoming.physical_size; + existing.generation = incoming.generation.clone(); + existing.version_id = incoming.version_id.clone(); + existing.attempts = existing.attempts.saturating_add(1); + continue; + } + + if size_reconciliation_entry_bytes(incoming) > MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET { + continue; + } + + let mut entry = incoming.clone(); + entry.first_seen = now; + entry.attempts = 1; + info.size_reconciliation.insert(entry.key.clone(), entry); + } + prune_size_reconciliation(info, now); + } + } + fn record_scan_resume_hint(&mut self, folder: &str) { self.new_cache.info.scan_resume_after = Some(folder.to_string()); self.update_cache.info.scan_resume_after = Some(folder.to_string()); @@ -1426,6 +1527,7 @@ impl FolderScanner { abandoned_children.remove(&path_join_buf(&[&item.bucket, &item.object_path()])); apply_scanner_size_summary(into, &sz); + self.apply_size_reconciliation(&sz); into.objects += 1; object_count += 1; self.budget.record_object_scanned(); @@ -2194,6 +2296,10 @@ pub async fn scan_data_folder( list_path_raw_options_observer: None, }; + let now = FolderScanner::now_secs(); + prune_size_reconciliation(&mut scanner.new_cache.info, now); + prune_size_reconciliation(&mut scanner.update_cache.info, now); + // Check if context is cancelled if ctx.is_cancelled() { return Err(ScannerError::Other("Operation cancelled".to_string())); @@ -2217,7 +2323,9 @@ pub async fn scan_data_folder( 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; - let unresolved_objects = root.failed_objects > 0 || !new_cache.info.failed_objects.is_empty(); + let unresolved_objects = root.failed_objects > 0 + || !new_cache.info.failed_objects.is_empty() + || !new_cache.info.size_reconciliation.is_empty(); new_cache.info.snapshot_complete = !unresolved_objects; let had_scan_checkpoint = cache.info.scan_checkpoint.is_some() || new_cache.info.scan_checkpoint.is_some(); new_cache.info.scan_resume_after = None; @@ -2245,7 +2353,7 @@ pub async fn scan_data_folder( if root_has_progress { new_cache.replace_hashed(&root_hash, &None, &root); } - if partial_cache_is_useful(&root, pending_heals_changed) { + if partial_cache_is_useful(&root, pending_heals_changed) || !new_cache.info.size_reconciliation.is_empty() { if new_cache.root().is_some() { new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); } diff --git a/crates/scanner/src/scanner_folder/item_actions.rs b/crates/scanner/src/scanner_folder/item_actions.rs index 15bd6736b..d49a08e3b 100644 --- a/crates/scanner/src/scanner_folder/item_actions.rs +++ b/crates/scanner/src/scanner_folder/item_actions.rs @@ -13,6 +13,7 @@ // limitations under the License. /// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers. use super::*; +use sha2::{Digest as _, Sha256}; /// Cached folder information for scanning #[derive(Clone, Debug)] @@ -32,6 +33,259 @@ pub(super) enum GetSizeFailureAction { HealMetadata { object: String }, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) enum SizeResolutionReason { + CompressedSizeUnknown, + InvalidPhysicalSize, + UnsupportedCompression, + InvalidObjectSize, + InvalidPartSize, + InvalidDeclaredSize, + SizeOverflowOrMismatch, +} + +impl SizeResolutionReason { + fn as_str(self) -> &'static str { + match self { + Self::CompressedSizeUnknown => "compressed_size_unknown", + Self::InvalidPhysicalSize => "invalid_physical_size", + Self::UnsupportedCompression => "unsupported_compression", + Self::InvalidObjectSize => "invalid_object_size", + Self::InvalidPartSize => "invalid_part_size", + Self::InvalidDeclaredSize => "invalid_declared_size", + Self::SizeOverflowOrMismatch => "size_overflow_or_mismatch", + } + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(super) enum SizeResolution { + Known { logical: i64, physical: i64 }, + Unknown { physical: i64, reason: SizeResolutionReason }, + Corrupt { physical: i64, reason: SizeResolutionReason }, +} + +impl SizeResolution { + fn known_size(&self) -> Option { + match self { + Self::Known { logical, .. } => Some(*logical), + Self::Unknown { .. } | Self::Corrupt { .. } => None, + } + } +} + +fn size_reconciliation_key(oi: &ObjectInfo, reason: SizeResolutionReason) -> String { + let version = oi + .version_id + .filter(|version| !version.is_nil()) + .map(|version| version.to_string()) + .unwrap_or_default(); + let generation = oi + .data_dir + .filter(|generation| !generation.is_nil()) + .map(|generation| generation.to_string()) + .unwrap_or_default(); + // Length-prefix each component so an object key containing the separator + // cannot alias another identity. S3 keys are bounded in normal operation; + // oversized persisted values use a digest so a corrupt metadata record + // cannot grow the ledger without bound. + fn component(value: &str) -> String { + const MAX_COMPONENT_LEN: usize = 512; + if value.len() <= MAX_COMPONENT_LEN { + return format!("{}:{}", value.len(), value); + } + let digest = Sha256::digest(value.as_bytes()); + let digest = hex_simd::encode_to_string(digest, hex_simd::AsciiCase::Lower); + format!("hash:{}:{}", value.len(), digest) + } + format!( + "{}|{}|{}|{}|{}", + component(&oi.bucket), + component(&oi.name), + component(&version), + component(&generation), + component(reason.as_str()) + ) +} + +pub(super) fn bounded_reconciliation_field(value: &str) -> String { + const MAX_FIELD_LEN: usize = 512; + if value.len() <= MAX_FIELD_LEN { + return value.to_string(); + } + let digest = hex_simd::encode_to_string(Sha256::digest(value.as_bytes()), hex_simd::AsciiCase::Lower); + let prefix_len = MAX_FIELD_LEN - 65; + let prefix = value + .char_indices() + .take_while(|(offset, ch)| offset.saturating_add(ch.len_utf8()) <= prefix_len) + .map(|(_, ch)| ch) + .collect::(); + format!("{}~{}", prefix, digest) +} + +fn record_size_resolution(summary: &mut SizeSummary, oi: &ObjectInfo, resolution: &SizeResolution) { + match resolution { + SizeResolution::Known { .. } => {} + SizeResolution::Unknown { physical, reason } | SizeResolution::Corrupt { physical, reason } => { + summary.record_size_reconciliation(SizeReconciliationEntry { + key: size_reconciliation_key(oi, *reason), + bucket: bounded_reconciliation_field(&oi.bucket), + object: bounded_reconciliation_field(&oi.name), + version_id: oi + .version_id + .filter(|version| !version.is_nil()) + .map(|version| version.to_string()), + generation: oi + .data_dir + .filter(|generation| !generation.is_nil()) + .map(|generation| generation.to_string()), + reason: reason.as_str().to_string(), + physical_size: u64::try_from(*physical).ok(), + first_seen: 0, + attempts: 0, + }); + } + } +} + +/// Resolve the size metadata once at the scanner trust boundary. A compressed +/// -1 sentinel is valid legacy metadata, but it cannot participate in normal +/// logical-size accounting or size-filtered lifecycle rules. +pub(super) fn resolve_size(oi: &ObjectInfo) -> SizeResolution { + let physical = oi.size; + if physical < 0 { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::InvalidPhysicalSize, + }; + } + + let compressed = match oi.compression_read_plan() { + Ok((_, _, compressed)) => compressed, + Err(_) => { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::UnsupportedCompression, + }; + } + }; + + if oi.actual_size < -1 || (oi.actual_size == -1 && !compressed) { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::InvalidObjectSize, + }; + } + + // Match ObjectInfo::get_actual_size: a positive in-memory value is the + // authoritative decoded size. Stale declared/part metadata must not turn + // an otherwise valid object into a false corruption report. + if oi.actual_size > 0 { + return SizeResolution::Known { + logical: oi.actual_size, + physical, + }; + } + + if oi + .parts + .iter() + .any(|part| part.actual_size < -1 || (part.actual_size < 0 && !compressed)) + { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::InvalidPartSize, + }; + } + + let declared = rustfs_utils::http::get_str(&oi.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE); + let declared = match declared { + Some(value) if value.is_empty() => { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::InvalidDeclaredSize, + }; + } + Some(value) => match value.parse::() { + Ok(value) if value >= 0 => Some(value), + _ => { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::InvalidDeclaredSize, + }; + } + }, + None => None, + }; + + let logical = match oi.get_actual_size() { + Ok(size) if size == -1 && compressed && declared.is_none() => { + return SizeResolution::Unknown { + physical, + reason: SizeResolutionReason::CompressedSizeUnknown, + }; + } + Ok(size) if size >= 0 => size, + Ok(_) | Err(_) => { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::SizeOverflowOrMismatch, + }; + } + }; + + if compressed && logical == 0 && physical != 0 && oi.parts.is_empty() && declared.is_none() { + return SizeResolution::Corrupt { + physical, + reason: SizeResolutionReason::SizeOverflowOrMismatch, + }; + } + + SizeResolution::Known { logical, physical } +} + +fn resolve_sizes(object_infos: &[ObjectInfo]) -> Vec { + object_infos.iter().map(resolve_size).collect() +} + +fn lifecycle_rule_has_size_filter(lifecycle: &BucketLifecycleConfiguration, rule_id: &str) -> bool { + let filter_has_size = |filter: &s3s::dto::LifecycleRuleFilter| { + filter.object_size_greater_than.is_some() + || filter.object_size_less_than.is_some() + || filter + .and + .as_ref() + .is_some_and(|and| and.object_size_greater_than.is_some() || and.object_size_less_than.is_some()) + }; + lifecycle + .rules + .iter() + .find(|rule| rule.id.as_deref().unwrap_or_default() == rule_id) + .and_then(|rule| rule.filter.as_ref()) + .is_some_and(filter_has_size) +} + +fn lifecycle_event_allowed(resolution: &SizeResolution, event: &Event, lifecycle: &BucketLifecycleConfiguration) -> bool { + match resolution { + // Corrupt metadata cannot safely authorize a destructive action, even + // when the evaluator happened to produce a time-only event. + SizeResolution::Corrupt { .. } => false, + // A valid-but-unknown logical size may still execute lifecycle + // actions whose rule is independent of object-size predicates. The + // evaluator has already selected the rule; only that rule's filter + // can make the missing logical value action-critical. + SizeResolution::Unknown { .. } => !lifecycle_rule_has_size_filter(lifecycle, &event.rule_id), + SizeResolution::Known { .. } => true, + } +} + +/// A successful newer-noncurrent batch consumes both known and unresolved +/// versions from the retained-version alert count. The two accounting paths +/// are separate because only known sizes can contribute byte totals. +fn remaining_versions_after_queued_noncurrent(remaining_versions: usize, known_count: usize, unknown_count: usize) -> usize { + remaining_versions.saturating_sub(known_count.saturating_add(unknown_count)) +} + /// How the corrupt-metadata branch records the repair after attempting an /// MRF intent (backlog#1894 axis A). #[derive(Debug, PartialEq, Eq)] @@ -319,34 +573,45 @@ impl ScannerItem { "Scanner lifecycle evaluation started" ); + let resolved_sizes = resolve_sizes(&object_infos); + if let Some(first) = object_infos.first() { + size_summary.record_reconciliation_scope( + &bounded_reconciliation_field(&first.bucket), + &bounded_reconciliation_field(&first.name), + ); + } + for (oi, resolution) in object_infos.iter().zip(resolved_sizes.iter()) { + record_size_resolution(size_summary, oi, resolution); + } + let has_corrupt_size = resolved_sizes + .iter() + .any(|resolution| matches!(resolution, SizeResolution::Corrupt { .. })); + // `versioning_config` is resolved once per object by the caller // (`get_size`) and handed in; only `prefix_enabled` is consulted here. - let Some(lifecycle) = self.lifecycle.as_ref() else { - let mut cumulative_size = 0; - for oi in object_infos.iter() { - let actual_size = match oi.get_actual_size() { - Ok(size) => size, - Err(_) => { - warn!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_LIFECYCLE_ACTION, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_LIFECYCLE, - bucket = %self.bucket, - object = %oi.name, - state = "size_lookup_failed", - "Scanner lifecycle action used fallback size" - ); + let Some(lifecycle) = self.lifecycle.clone() else { + let mut cumulative_size: i64 = 0; + for (oi, resolved_size) in object_infos.iter().zip(resolved_sizes.iter()) { + let accounting_size = match resolved_size { + SizeResolution::Known { logical, .. } => *logical, + // A valid compressed legacy sentinel has no logical size, + // but heal and replication still need to run. The + // physical size is only an input to those operations; it + // is not folded into the logical total below. + SizeResolution::Unknown { physical, .. } => { + self.heal_actions(oi, *physical, size_summary).await; + size_summary.actions_accounting_unknown(oi); continue; } + SizeResolution::Corrupt { .. } => continue, }; - let size = self.heal_actions(oi, actual_size, size_summary).await; + let size = self.heal_actions(oi, accounting_size, size_summary).await; - size_summary.actions_accounting(oi, size, actual_size); + size_summary.actions_accounting(oi, size, accounting_size); - cumulative_size += size; + cumulative_size = cumulative_size.saturating_add(size); } self.alert_excessive_versions(object_infos.len(), cumulative_size); @@ -400,25 +665,108 @@ impl ScannerItem { let mut to_delete_objs: Vec = Vec::new(); let mut noncurrent_events: Vec = Vec::new(); let mut noncurrent_accounting: Vec> = Vec::new(); + let mut noncurrent_unknown: Vec<&ObjectInfo> = Vec::new(); let mut cumulative_size = 0; let mut remaining_versions = object_infos.len(); 'eventLoop: { for (i, event) in events.iter().enumerate() { let oi = &object_infos[i]; - let actual_size = match oi.get_actual_size() { - Ok(size) => size, - Err(_) => { - warn!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_LIFECYCLE_ACTION, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_LIFECYCLE, - bucket = %self.bucket, - object = %oi.name, - state = "size_lookup_failed", - "Scanner lifecycle action used fallback size" - ); - 0 + let known_size = resolved_sizes[i].known_size(); + if has_corrupt_size + && matches!( + event.action, + IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction + ) + { + // An all-version delete would also remove a corrupt + // sibling that could not be reconciled safely. + continue; + } + if !lifecycle_event_allowed(&resolved_sizes[i], event, &lifecycle) { + // An unknown logical size must not make an otherwise + // non-destructive scan disappear from heal/physical-tier + // accounting. Size-filtered or deferred events remain + // pending, so retain the version-only physical counters. + if let SizeResolution::Unknown { physical, .. } = &resolved_sizes[i] { + self.heal_actions(oi, *physical, size_summary).await; + size_summary.actions_accounting_unknown(oi); + } + continue; + } + let actual_size = match known_size { + Some(size) => size, + None => { + match event.action { + IlmAction::DeleteAction + | IlmAction::DeleteRestoredAction + | IlmAction::DeleteRestoredVersionAction + | IlmAction::DeleteAllVersionsAction + | IlmAction::DelMarkerDeleteAllVersionsAction => { + let done_ilm = Metrics::time_ilm(event.action); + let trace_started_at = trace_start_instant(); + let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await; + emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at); + if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) { + done_ilm(1)(); + if event.action == IlmAction::DeleteAllVersionsAction + || event.action == IlmAction::DelMarkerDeleteAllVersionsAction + { + remaining_versions = 0; + } + } else if matches!( + event.action, + IlmAction::DeleteAction + | IlmAction::DeleteRestoredAction + | IlmAction::DeleteRestoredVersionAction + ) { + size_summary.actions_accounting_unknown(oi); + } else { + size_summary.actions_accounting_unknown(oi); + for (j, retained) in object_infos.iter().enumerate().skip(i + 1) { + match &resolved_sizes[j] { + SizeResolution::Known { logical, .. } => PendingScannerAccounting { + object: retained, + retained_size: *logical, + expired_size: 0, + } + .apply(size_summary, &mut cumulative_size, false), + SizeResolution::Unknown { .. } => { + size_summary.actions_accounting_unknown(retained); + } + SizeResolution::Corrupt { .. } => {} + } + } + } + } + IlmAction::DeleteVersionAction => { + if let Some(opt) = object_opts.get(i) { + to_delete_objs.push(ObjectToDelete { + object_name: opt.name.clone(), + version_id: opt.version_id, + ..Default::default() + }); + noncurrent_events.push(event.clone()); + noncurrent_unknown.push(oi); + } + } + IlmAction::TransitionAction | IlmAction::TransitionVersionAction => { + let trace_started_at = trace_start_instant(); + let queued = apply_transition_rule(event, &LcEventSrc::Scanner, oi).await; + emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at); + if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) { + let done_ilm = Metrics::time_ilm(event.action); + done_ilm(1)(); + } + size_summary.actions_accounting_unknown(oi); + } + IlmAction::NoneAction | IlmAction::ActionCount => { + if let SizeResolution::Unknown { physical, .. } = &resolved_sizes[i] { + self.heal_actions(oi, *physical, size_summary).await; + } + size_summary.actions_accounting_unknown(oi); + } + } + continue; } }; @@ -446,36 +794,24 @@ impl ScannerItem { done_ilm(1)(); remaining_versions = 0; } else { - PendingScannerAccounting { - object: oi, - retained_size: actual_size, - expired_size: 0, - } - .apply(size_summary, &mut cumulative_size, false); - for retained in object_infos.iter().skip(i + 1) { - let retained_size = match retained.get_actual_size() { - Ok(size) => size, - Err(_) => { - warn!( - target: "rustfs::scanner::folder", - event = EVENT_SCANNER_LIFECYCLE_ACTION, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_LIFECYCLE, - bucket = %self.bucket, - object = %retained.name, - state = "size_lookup_failed", - "Scanner lifecycle action used fallback size" - ); - 0 - } - }; + if let Some(actual_size) = known_size { PendingScannerAccounting { - object: retained, - retained_size, + object: oi, + retained_size: actual_size, expired_size: 0, } .apply(size_summary, &mut cumulative_size, false); } + for (j, retained) in object_infos.iter().enumerate().skip(i + 1) { + if let Some(retained_size) = resolved_sizes[j].known_size() { + PendingScannerAccounting { + object: retained, + retained_size, + expired_size: 0, + } + .apply(size_summary, &mut cumulative_size, false); + } + } } break 'eventLoop; } @@ -511,11 +847,13 @@ impl ScannerItem { version_id: opt.version_id, ..Default::default() }); - noncurrent_accounting.push(PendingScannerAccounting { - object: oi, - retained_size: actual_size, - expired_size: 0, - }); + if let Some(actual_size) = known_size { + noncurrent_accounting.push(PendingScannerAccounting { + object: oi, + retained_size: actual_size, + expired_size: 0, + }); + } account_now = false; } noncurrent_events.push(event.clone()); @@ -548,7 +886,7 @@ impl ScannerItem { if account_now { size_summary.actions_accounting(oi, size, actual_size); - cumulative_size += size; + cumulative_size = cumulative_size.saturating_add(size); } } } @@ -576,11 +914,20 @@ impl ScannerItem { } if record_scanner_ilm_action_if_queued(global_metrics(), action, count, queued) { done_ilm(count)(); - remaining_versions = remaining_versions.saturating_sub(noncurrent_accounting.len()); + remaining_versions = remaining_versions_after_queued_noncurrent( + remaining_versions, + noncurrent_accounting.len(), + noncurrent_unknown.len(), + ); } for pending in noncurrent_accounting { pending.apply(size_summary, &mut cumulative_size, queued); } + if !queued { + for object in noncurrent_unknown { + size_summary.actions_accounting_unknown(object); + } + } } self.alert_excessive_versions(remaining_versions, cumulative_size); } @@ -929,4 +1276,361 @@ mod tests { assert_eq!(item.object_name, "object"); assert_eq!(item.object_path(), "object"); } + + #[test] + fn size_resolution_rejects_negative_overflow_and_unknown_compression() { + let compressed = |actual_size: i64, declared: Option<&str>| { + let mut user_defined = HashMap::new(); + rustfs_utils::http::insert_str(&mut user_defined, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + if let Some(declared) = declared { + rustfs_utils::http::insert_str(&mut user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, declared.to_string()); + } + ObjectInfo { + size: 12, + actual_size, + user_defined: Arc::new(user_defined), + ..Default::default() + } + }; + + let normal = ObjectInfo { + size: 12, + actual_size: 10, + ..Default::default() + }; + assert_eq!( + resolve_size(&normal), + SizeResolution::Known { + logical: 10, + physical: 12 + } + ); + + let stale_declared_metadata = ObjectInfo { + size: 12, + actual_size: 10, + user_defined: Arc::new(HashMap::from([("x-rustfs-internal-actual-size".to_string(), "not-a-size".to_string())])), + parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo { + actual_size: -2, + ..Default::default() + }]), + ..Default::default() + }; + assert_eq!( + resolve_size(&stale_declared_metadata), + SizeResolution::Known { + logical: 10, + physical: 12 + } + ); + + assert_eq!( + resolve_size(&compressed(0, Some("9"))), + SizeResolution::Known { + logical: 9, + physical: 12 + } + ); + assert_eq!( + resolve_size(&compressed(-1, None)), + SizeResolution::Unknown { + physical: 12, + reason: SizeResolutionReason::CompressedSizeUnknown, + } + ); + assert!(matches!( + resolve_size(&compressed(0, Some("not-a-size"))), + SizeResolution::Corrupt { + reason: SizeResolutionReason::InvalidDeclaredSize, + .. + } + )); + assert!(matches!( + resolve_size(&ObjectInfo { + size: 12, + actual_size: -2, + ..Default::default() + }), + SizeResolution::Corrupt { .. } + )); + assert!(matches!(resolve_size(&compressed(0, Some("-1"))), SizeResolution::Corrupt { .. })); + assert!(matches!(resolve_size(&compressed(0, Some(""))), SizeResolution::Corrupt { .. })); + + let unsupported = { + let mut object = compressed(0, None); + let mut metadata = (*object.user_defined).clone(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "unsupported".to_string()); + object.user_defined = Arc::new(metadata); + object + }; + assert!(matches!(resolve_size(&unsupported), SizeResolution::Corrupt { .. })); + + let invalid_part = { + let mut object = compressed(0, None); + object.parts = Arc::new(vec![rustfs_filemeta::ObjectPartInfo { + size: 12, + actual_size: -2, + ..Default::default() + }]); + object + }; + assert!(matches!(resolve_size(&invalid_part), SizeResolution::Corrupt { .. })); + + let overflow = { + let mut object = compressed(0, None); + object.parts = Arc::new(vec![ + rustfs_filemeta::ObjectPartInfo { + size: 1, + actual_size: i64::MAX, + ..Default::default() + }, + rustfs_filemeta::ObjectPartInfo { + size: 1, + actual_size: 1, + ..Default::default() + }, + ]); + object + }; + assert!(matches!(resolve_size(&overflow), SizeResolution::Corrupt { .. })); + + let mismatch = compressed(0, None); + assert!(matches!(resolve_size(&mismatch), SizeResolution::Corrupt { .. })); + assert_eq!( + resolve_size(&ObjectInfo { + size: 0, + actual_size: 0, + ..Default::default() + }), + SizeResolution::Known { logical: 0, physical: 0 } + ); + } + + #[test] + fn size_resolution_records_and_replays_one_identity() { + let version_id = uuid::Uuid::new_v4(); + let generation = uuid::Uuid::new_v4(); + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, "not-a-number".to_string()); + let corrupt = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 12, + version_id: Some(version_id), + data_dir: Some(generation), + user_defined: Arc::new(metadata), + ..Default::default() + }; + + let mut summary = SizeSummary::default(); + let resolution = resolve_size(&corrupt); + record_size_resolution(&mut summary, &corrupt, &resolution); + record_size_resolution(&mut summary, &corrupt, &resolution); + assert_eq!(summary.size_reconciliation.len(), 1); + assert_eq!(summary.size_reconciliation[0].reason, "invalid_declared_size"); + assert_eq!(summary.size_reconciliation[0].physical_size, Some(12)); + + let known = ObjectInfo { + actual_size: 12, + user_defined: Arc::new(HashMap::new()), + ..corrupt.clone() + }; + record_size_resolution(&mut summary, &known, &resolve_size(&known)); + summary.record_reconciliation_scope(&known.bucket, &known.name); + assert_eq!(summary.reconciliation_scopes.len(), 1); + assert_eq!(summary.reconciliation_scopes[0].bucket, "bucket"); + } + + #[test] + fn malformed_size_has_same_ilm_accounting() { + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, "invalid".to_string()); + let object = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 12, + user_defined: Arc::new(metadata), + ..Default::default() + }; + let resolution = resolve_size(&object); + let mut without_ilm = SizeSummary::default(); + let mut with_ilm = SizeSummary::default(); + record_size_resolution(&mut without_ilm, &object, &resolution); + record_size_resolution(&mut with_ilm, &object, &resolution); + assert_eq!(without_ilm.size_reconciliation, with_ilm.size_reconciliation); + assert_eq!(without_ilm.total_size, 0); + assert_eq!(with_ilm.total_size, 0); + assert!(without_ilm.tier_stats.is_empty()); + assert!(with_ilm.tier_stats.is_empty()); + } + + #[test] + fn size_resolution_parses_once_per_version() { + let objects = vec![ + ObjectInfo { + bucket: "bucket".to_string(), + name: "one".to_string(), + size: 1, + actual_size: 1, + ..Default::default() + }, + ObjectInfo { + bucket: "bucket".to_string(), + name: "two".to_string(), + size: 2, + actual_size: -2, + ..Default::default() + }, + ]; + let resolutions = resolve_sizes(&objects); + assert_eq!(resolutions.len(), objects.len()); + assert!(matches!(resolutions[0], SizeResolution::Known { logical: 1, .. })); + assert!(matches!(resolutions[1], SizeResolution::Corrupt { .. })); + } + + #[test] + fn queued_unknown_noncurrent_versions_are_removed_from_alert_count() { + assert_eq!(remaining_versions_after_queued_noncurrent(3, 1, 2), 0); + assert_eq!(remaining_versions_after_queued_noncurrent(7, 2, 1), 4); + assert_eq!(remaining_versions_after_queued_noncurrent(usize::MAX, usize::MAX, usize::MAX), 0); + } + + #[test] + fn malformed_size_blocks_size_dependent_transition_but_allows_time_only_expiry() { + let size_filtered = BucketLifecycleConfiguration { + rules: vec![s3s::dto::LifecycleRule { + status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED), + expiration: None, + abort_incomplete_multipart_upload: None, + del_marker_expiration: None, + id: Some("size".to_string()), + filter: Some(s3s::dto::LifecycleRuleFilter { + object_size_greater_than: Some(1), + ..Default::default() + }), + noncurrent_version_expiration: None, + noncurrent_version_transitions: None, + prefix: None, + transitions: None, + }], + ..Default::default() + }; + let unknown = SizeResolution::Unknown { + physical: 12, + reason: SizeResolutionReason::CompressedSizeUnknown, + }; + let size_event = Event { + action: IlmAction::DeleteAction, + rule_id: "size".to_string(), + ..Default::default() + }; + assert!(!lifecycle_event_allowed(&unknown, &size_event, &size_filtered)); + assert!(!lifecycle_event_allowed( + &unknown, + &Event { + action: IlmAction::TransitionAction, + rule_id: "size".to_string(), + ..Default::default() + }, + &size_filtered + )); + let mixed_filters = BucketLifecycleConfiguration { + rules: vec![ + size_filtered.rules[0].clone(), + s3s::dto::LifecycleRule { + status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED), + expiration: None, + abort_incomplete_multipart_upload: None, + del_marker_expiration: None, + id: Some("time".to_string()), + filter: None, + noncurrent_version_expiration: None, + noncurrent_version_transitions: None, + prefix: None, + transitions: None, + }, + ], + ..Default::default() + }; + assert!(lifecycle_event_allowed( + &unknown, + &Event { + action: IlmAction::DeleteAction, + rule_id: "time".to_string(), + ..Default::default() + }, + &mixed_filters + )); + assert!(lifecycle_event_allowed( + &unknown, + &Event { + action: IlmAction::TransitionAction, + ..Default::default() + }, + &BucketLifecycleConfiguration::default() + )); + assert!(!lifecycle_event_allowed( + &SizeResolution::Corrupt { + physical: 12, + reason: SizeResolutionReason::InvalidDeclaredSize, + }, + &Event { + action: IlmAction::DeleteAction, + ..Default::default() + }, + &BucketLifecycleConfiguration::default() + )); + assert!(lifecycle_event_allowed( + &SizeResolution::Known { + logical: 10, + physical: 12, + }, + &Event { + action: IlmAction::DeleteAllVersionsAction, + ..Default::default() + }, + &BucketLifecycleConfiguration::default() + )); + assert!(lifecycle_event_allowed( + &unknown, + &Event { + action: IlmAction::DeleteAction, + rule_id: "time-only".to_string(), + ..Default::default() + }, + &BucketLifecycleConfiguration::default() + )); + } + + #[tokio::test] + async fn long_object_size_reconciliation_scope_uses_bounded_identity() { + let object_name = "o".repeat(600); + let mut item = scanner_item_with_prefix(""); + item.object_name = object_name.clone(); + + let mut metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + let object = ObjectInfo { + bucket: item.bucket.clone(), + name: object_name.clone(), + size: 12, + actual_size: -1, + version_id: Some(uuid::Uuid::new_v4()), + user_defined: Arc::new(metadata), + ..Default::default() + }; + let mut summary = SizeSummary::default(); + item.apply_actions(vec![object], None, VersioningConfiguration::default(), &mut summary) + .await; + + let bounded_bucket = bounded_reconciliation_field(&item.bucket); + let bounded_object = bounded_reconciliation_field(&object_name); + assert_eq!(summary.reconciliation_scopes[0].bucket, bounded_bucket); + assert_eq!(summary.reconciliation_scopes[0].object, bounded_object); + assert_eq!(summary.size_reconciliation[0].object, bounded_object); + assert_eq!(summary.versions, 1); + assert_eq!(summary.total_size, 0); + } } diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 503b0e75f..4e8573909 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -388,6 +388,66 @@ async fn test_record_failed_ttl_zero_noop() { assert!(!scanner.should_skip_failed("path2")); } +#[tokio::test] +async fn malformed_size_reconciliation_replays_after_restart() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir); + + let entry = SizeReconciliationEntry { + key: "1:b|6:object|0:|0:".to_string(), + bucket: "b".to_string(), + object: "object".to_string(), + reason: "invalid_declared_size".to_string(), + physical_size: Some(12), + ..Default::default() + }; + let mut summary = SizeSummary::default(); + summary.record_size_reconciliation(entry.clone()); + summary.record_reconciliation_scope("b", "object"); + scanner.apply_size_reconciliation(&summary); + scanner.apply_size_reconciliation(&summary); + + assert_eq!(scanner.new_cache.info.size_reconciliation.len(), 1); + assert_eq!(scanner.update_cache.info.size_reconciliation.len(), 1); + assert_eq!(scanner.new_cache.info.size_reconciliation[&entry.key].attempts, 2); + + let encoded = rmp_serde::to_vec_named(&scanner.new_cache.info).expect("size ledger should encode"); + let decoded: crate::data_usage_define::DataUsageCacheInfo = + rmp_serde::from_slice(&encoded).expect("size ledger should decode"); + assert_eq!(decoded.size_reconciliation.len(), 1); + assert_eq!(decoded.size_reconciliation[&entry.key].reason, "invalid_declared_size"); + + let mut resolved = SizeSummary::default(); + resolved.record_reconciliation_scope("b", "object"); + scanner.apply_size_reconciliation(&resolved); + assert!(scanner.new_cache.info.size_reconciliation.is_empty()); + assert!(scanner.update_cache.info.size_reconciliation.is_empty()); +} + +#[tokio::test] +async fn malformed_size_reconciliation_clears_bounded_long_object_scope() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir); + let long_object = "o".repeat(600); + let bounded_object = item_actions::bounded_reconciliation_field(&long_object); + let entry = SizeReconciliationEntry { + key: "long-object-key".to_string(), + bucket: "b".to_string(), + object: bounded_object, + reason: "invalid_declared_size".to_string(), + ..Default::default() + }; + let mut summary = SizeSummary::default(); + summary.record_size_reconciliation(entry); + scanner.apply_size_reconciliation(&summary); + assert_eq!(scanner.new_cache.info.size_reconciliation.len(), 1); + + let mut resolved = SizeSummary::default(); + resolved.record_reconciliation_scope("b", &long_object); + scanner.apply_size_reconciliation(&resolved); + assert!(scanner.new_cache.info.size_reconciliation.is_empty()); +} + #[test] fn test_classify_get_size_failure_marks_metadata_heal_object_path() { let temp_dir = std::env::temp_dir();