// Copyright 2024 RustFS Team // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration}; use serde::{Deserialize, Serialize, ser::SerializeMap}; use sha2::{Digest, Sha256}; use std::{ collections::{HashMap, HashSet}, future::Future, sync::{Arc, LazyLock, Once}, time::SystemTime, }; use http::HeaderMap; use metrics::{counter, describe_counter, describe_histogram, histogram}; #[cfg(test)] 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, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeReconciliationEntry, SizeReconciliationScope, SizeSummary, TierAccountingProof, TierStats, UNKNOWN_TIER, UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP, UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP, UnknownTierStats, hash_path, prefix_usage_in_cache, }; use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use tokio::time::{Duration, Instant, sleep, timeout}; use tracing::{debug, warn}; use crate::raw_page_index::{RawEnumerationPageIndex, RawEnumerationPageOwnerStatus}; use crate::storage_api::owner::HTTPPreconditions; use crate::{ BUCKET_META_PREFIX, EcstoreError as Error, EcstoreResult as StorageResult, RUSTFS_META_BUCKET, ReplicationConfig, SCANNER_PUBLICATION_EPOCH_CHANGED, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, StorageError, TRANSITION_COMPLETE, save_config, save_config_with_preconditions, scanner_publication_admission_for_epoch, storageclass, }; use crate::{ScannerConfigObjectDelete, ScannerObjectIO}; // Data usage constants pub const DATA_USAGE_ROOT: &str = SLASH_SEPARATOR; const DATA_USAGE_BLOOM_NAME: &str = ".bloomcycle.bin"; pub const DATA_USAGE_CACHE_NAME: &str = ".usage-cache.bin"; pub(crate) const DATA_USAGE_CACHE_KEY_FORMAT: u16 = 1; const DATA_USAGE_CACHE_SAVE_RETRIES: u32 = 2; const DATA_USAGE_CACHE_BACKUP_SAVE_TIMEOUT_SECS_MAX: u64 = 5; const DATA_USAGE_CACHE_BACKUP_SAVE_RETRIES: u32 = 0; const DATA_USAGE_CACHE_SAVE_RETRY_BACKOFF_MAX: Duration = Duration::from_millis(350); const DATA_USAGE_CACHE_PERSISTENCE_MARGIN: Duration = Duration::from_secs(5); const METRIC_CACHE_SAVE_ATTEMPT_TOTAL: &str = "rustfs_scanner_cache_save_attempt_total"; const METRIC_CACHE_SAVE_TIMEOUT_TOTAL: &str = "rustfs_scanner_cache_save_timeout_total"; const METRIC_CACHE_SAVE_RETRY_TOTAL: &str = "rustfs_scanner_cache_save_retry_total"; const METRIC_CACHE_SAVE_DURATION_SECONDS: &str = "rustfs_scanner_cache_save_duration_seconds"; const METRIC_CACHE_BACKUP_REVISION_FAILURE_TOTAL: &str = "rustfs_scanner_cache_backup_revision_failure_total"; const LOG_COMPONENT_SCANNER: &str = "scanner"; const LOG_SUBSYSTEM_CACHE: &str = "cache"; const EVENT_SCANNER_CACHE_LOAD_STATE: &str = "scanner_cache_load_state"; const EVENT_SCANNER_CACHE_SAVE_STATE: &str = "scanner_cache_save_state"; static CACHE_SAVE_METRICS_ONCE: Once = Once::new(); pub const DATA_USAGE_SCAN_CHECKPOINT_VERSION: u16 = 1; pub const DATA_USAGE_RAW_ENUMERATION_CURSOR_VERSION: u16 = 1; const DATA_USAGE_SCAN_CURSOR_MAX_BYTES: usize = 16 * 1024; #[derive(Clone, Debug, PartialEq, Eq)] pub(crate) enum DataUsageCacheRevision { Missing, Etag(String), } impl DataUsageCacheRevision { pub(crate) fn preconditions(&self) -> HTTPPreconditions { match self { Self::Missing => HTTPPreconditions { if_none_match: Some("*".to_string()), ..Default::default() }, Self::Etag(etag) => HTTPPreconditions { if_match: Some(etag.clone()), ..Default::default() }, } } } pub(crate) async fn read_config_with_revision( store: Arc, path: &str, ) -> StorageResult<(Option>, DataUsageCacheRevision)> { match store .get_object_reader( RUSTFS_META_BUCKET, path, None, HeaderMap::new(), &ObjectOptions { no_lock: true, ..Default::default() }, ) .await { Ok(mut reader) => { let revision = reader .object_info .etag .as_ref() .filter(|etag| !etag.is_empty()) .cloned() .map(DataUsageCacheRevision::Etag) .ok_or_else(|| StorageError::other(format!("scanner config object {path} has no ETag")))?; Ok((Some(reader.read_all().await?), revision)) } Err( Error::ConfigNotFound | Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_), ) => Ok((None, DataUsageCacheRevision::Missing)), Err(err) => Err(err), } } pub(crate) fn usage_floor_primary_read_error_allows_backup(err: &Error) -> bool { match err { Error::FileCorrupt | Error::CorruptedFormat | Error::CorruptedBackend | Error::PartMissingOrCorrupt | Error::LessData | Error::MoreData => true, Error::Io(io_error) => { matches!(io_error.kind(), std::io::ErrorKind::InvalidData | std::io::ErrorKind::UnexpectedEof) || error_chain_has_usage_floor_corruption_signature(io_error) } _ => false, } } fn error_chain_has_usage_floor_corruption_signature(error: &(dyn std::error::Error + 'static)) -> bool { let mut current = Some(error); while let Some(err) = current { let message = err.to_string(); if message.contains("InlineData value out of range") || message.contains("InlineData key out of range") || message.contains("insufficient data for metadata") || message.contains("insufficient data for meta length") || message.contains("insufficient data for CRC") { return true; } current = err.source(); } false } /// Read only the object revision without materializing its body. pub(crate) async fn read_config_revision(store: Arc, path: &str) -> StorageResult { match store .get_object_reader( RUSTFS_META_BUCKET, path, None, HeaderMap::new(), &ObjectOptions { no_lock: true, ..Default::default() }, ) .await { Ok(reader) => reader .object_info .etag .filter(|etag| !etag.is_empty()) .map(DataUsageCacheRevision::Etag) .ok_or_else(|| StorageError::other(format!("scanner config object {path} has no ETag"))), Err( Error::ConfigNotFound | Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_), ) => Ok(DataUsageCacheRevision::Missing), Err(err) => Err(err), } } #[derive(Clone, Debug, PartialEq, Eq)] pub(crate) struct DataUsageCacheRevisions { main: DataUsageCacheRevision, backup: Option, } pub static DATA_USAGE_BUCKET: LazyLock = LazyLock::new(|| format!("{RUSTFS_META_BUCKET}{SLASH_SEPARATOR}{BUCKET_META_PREFIX}")); pub static DATA_USAGE_OBJ_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_OBJECT_NAME}")); /// Durable evidence for recovery of the exact empty usage fence written by /// rc.2/rc.3 bucket cleanup before the first authoritative scanner snapshot. pub static DATA_USAGE_RECOVERY_PATH: LazyLock = LazyLock::new(|| format!("{}.recovery-pending.json", DATA_USAGE_OBJ_NAME_PATH.as_str())); pub static DATA_USAGE_OBSERVED_OBJ_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_OBSERVED_OBJECT_NAME}")); pub static LEGACY_DATA_USAGE_OBJ_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{LEGACY_DATA_USAGE_OBJECT_NAME}")); pub static DATA_USAGE_BLOOM_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_BLOOM_NAME}")); /// Durable companion object for a cycle-state object which cannot be decoded. /// The primary object is deliberately never replaced or deleted by recovery. pub static DATA_USAGE_BLOOM_RECOVERY_PATH: LazyLock = LazyLock::new(|| format!("{}.recovery-required.json", DATA_USAGE_BLOOM_NAME_PATH.as_str())); 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; /// Scanner-side accounting on the shared [`SizeSummary`]. /// /// The type itself lives in `rustfs-data-usage`, which sits below the storage /// layer and cannot see `ObjectInfo`, so this stays an extension trait rather /// than an inherent method (backlog#1828). 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 { fn actions_accounting(&mut self, oi: &ObjectInfo, size: i64, actual_size: i64) { 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()) && size == actual_size { self.versions = self.versions.saturating_add(1); } let logical_size = size.max(0); let size = usize::try_from(logical_size).unwrap_or(usize::MAX); self.total_size = self.total_size.saturating_add(size); let logical_bytes = u64::try_from(logical_size).unwrap_or(u64::MAX); let physical_bytes = u64::try_from(oi.size.max(0)).unwrap_or(0); let mut proof = TierAccountingProof { logical_total: logical_bytes, logical_known: 0, physical_total: physical_bytes, physical_known: 0, overflowed: false, }; if oi.transitioned_object.free_version { proof.logical_known = logical_bytes; proof.physical_known = physical_bytes; self.tier_accounting_proof.saturating_add(proof); return; } let tier = if oi.transitioned_object.status == TRANSITION_COMPLETE { oi.transitioned_object.tier.as_str() } else { oi.storage_class.as_deref().unwrap_or(storageclass::STANDARD) }; let builtin_tier = tier == storageclass::STANDARD || tier == storageclass::RRS; let tier_registry_is_empty = self.tier_stats.is_empty() || (self.tier_stats.len() == 1 && self.tier_stats.contains_key(UNKNOWN_TIER)); let known_tier = tier != UNKNOWN_TIER && (builtin_tier || self.tier_stats.contains_key(tier)); // With no configured tier, retain the historical empty-map shape for // ordinary STANDARD/RRS objects. A non-built-in key is still an // observable unknown and must create only the fixed bucket. if tier_registry_is_empty && known_tier { proof.logical_known = logical_bytes; proof.physical_known = physical_bytes; self.tier_accounting_proof.saturating_add(proof); return; } // Configured tiers and the fixed bucket are normally seeded, so the // hot path can mutate them without allocating a key for every object. // The fallback inserts only when a legacy/no-config summary sees its // first unknown key. let tier_stats = if known_tier { if let Some(stats) = self.tier_stats.get_mut(tier) { stats } else { self.tier_stats.entry(tier.to_owned()).or_default() } } else if let Some(stats) = self.tier_stats.get_mut(UNKNOWN_TIER) { stats } else { self.tier_stats.entry(UNKNOWN_TIER.to_string()).or_default() }; *tier_stats = tier_stats.add(&TierStats { total_size: physical_bytes, num_versions: 1, num_objects: u64::from(oi.is_latest), }); if known_tier { proof.logical_known = logical_bytes; proof.physical_known = physical_bytes; } if !known_tier { self.unknown_tier_stats .record_dimensions(tier, logical_bytes, physical_bytes, 1, u64::from(oi.is_latest)); if self.unknown_tier_stats.counter_overflowed { proof.overflowed = true; } } self.tier_accounting_proof.saturating_add(proof); } 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 ===== #[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum DataUsageScanCheckpointReason { Runtime, Objects, Directories, Unknown, } impl DataUsageScanCheckpointReason { pub fn as_str(self) -> &'static str { match self { Self::Runtime => "runtime", Self::Objects => "objects", Self::Directories => "directories", Self::Unknown => "unknown", } } } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] pub struct DataUsageScanCheckpoint { pub version: u16, pub resume_after: String, pub reason: DataUsageScanCheckpointReason, } impl DataUsageScanCheckpoint { pub fn new(resume_after: String, reason: DataUsageScanCheckpointReason) -> Self { Self { version: DATA_USAGE_SCAN_CHECKPOINT_VERSION, resume_after, reason, } } } /// Durable raw directory-page cursor for a bucket scan. #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct DataUsageRawEnumerationCursor { pub version: u16, pub parent: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub last_entry: Option, pub entries_seen: u64, pub page_digest: [u8; 32], } impl DataUsageRawEnumerationCursor { pub fn new(parent: String, last_entry: Option, entries_seen: u64, page_digest: [u8; 32]) -> Self { Self { version: DATA_USAGE_RAW_ENUMERATION_CURSOR_VERSION, parent, last_entry, entries_seen, page_digest, } } fn is_valid_for_bucket(&self, bucket: &str) -> bool { self.version == DATA_USAGE_RAW_ENUMERATION_CURSOR_VERSION && bucket != DATA_USAGE_ROOT && path_is_in_bucket_scope(bucket, &self.parent) && self.parent.len() <= DATA_USAGE_SCAN_CURSOR_MAX_BYTES && self.page_digest != [0; 32] && match &self.last_entry { Some(last_entry) => { !last_entry.is_empty() && self.entries_seen > 0 && last_entry.len() <= DATA_USAGE_SCAN_CURSOR_MAX_BYTES && !last_entry.contains(SLASH_SEPARATOR) } None => self.entries_seen == 0, } } } fn path_is_in_bucket_scope(bucket: &str, path: &str) -> bool { path == bucket || path .strip_prefix(bucket) .is_some_and(|suffix| suffix.starts_with(SLASH_SEPARATOR)) } /// Durable scope of a bucket checkpoint, independent of namespace mutation counters. #[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct DataUsageScanIdentity { pub version: u16, pub bucket_incarnation: uuid::Uuid, pub set_layout: DataUsageScanPlanDigest, pub publication_epoch: u64, pub tier_registry_generation: u64, pub scan_mode: HealScanMode, } impl Serialize for DataUsageScanIdentity { fn serialize(&self, serializer: S) -> Result { let mut map = serializer.serialize_map(Some(6))?; map.serialize_entry("version", &self.version)?; map.serialize_entry("bucket_incarnation", &self.bucket_incarnation)?; map.serialize_entry("set_layout", &self.set_layout)?; map.serialize_entry("publication_epoch", &self.publication_epoch)?; map.serialize_entry("tier_registry_generation", &self.tier_registry_generation)?; map.serialize_entry("scan_mode", &self.scan_mode)?; map.end() } } impl DataUsageScanIdentity { pub(crate) fn is_valid(&self) -> bool { self.version == 1 && !self.bucket_incarnation.is_nil() && matches!(self.scan_mode, HealScanMode::Normal | HealScanMode::Deep) } } /// A forward coverage sweep may span budgets, but not authorize mixed mutation generations. #[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct DataUsageScanProgress { pub started_plan: DataUsageScanPlanDigest, pub requested_plan: DataUsageScanPlanDigest, } impl Serialize for DataUsageScanProgress { fn serialize(&self, serializer: S) -> Result { let mut map = serializer.serialize_map(Some(2))?; map.serialize_entry("started_plan", &self.started_plan)?; map.serialize_entry("requested_plan", &self.requested_plan)?; map.end() } } #[derive(Clone, Debug, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct DataUsageScanCoverageReceipt { pub through: String, pub digest: [u8; 32], } impl Serialize for DataUsageScanCoverageReceipt { fn serialize(&self, serializer: S) -> Result { let mut map = serializer.serialize_map(Some(2))?; map.serialize_entry("through", &self.through)?; map.serialize_entry("digest", &self.digest)?; map.end() } } struct CheckpointDigestWriter(Sha256); impl std::io::Write for CheckpointDigestWriter { fn write(&mut self, bytes: &[u8]) -> std::io::Result { self.0.update(bytes); Ok(bytes.len()) } fn flush(&mut self) -> std::io::Result<()> { Ok(()) } } #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageEntryInfo { pub name: String, pub parent: String, pub entry: DataUsageEntry, /// Durable bucket incarnation that produced this bucket root. Missing /// values are legacy/unproven and must not authorize cold-bucket reuse. #[serde(default, skip_serializing_if = "Option::is_none")] pub bucket_incarnation: Option, /// Registry generation used to classify this root entry. Older remote /// workers omit it; callers must reject that result when a frozen cycle /// requires generation fencing. #[serde(default, skip_serializing_if = "Option::is_none")] pub tier_registry_generation: Option, } #[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq, Hash)] #[serde(deny_unknown_fields)] pub struct DataUsageCacheSource { pub pool_index: usize, pub set_index: usize, } impl DataUsageCacheSource { pub const fn new(pool_index: usize, set_index: usize) -> Self { Self { pool_index, set_index } } } #[derive(Clone, Copy, Debug, Default, Serialize, Deserialize, PartialEq, Eq, Hash)] #[serde(transparent)] pub struct DataUsageScanPlanDigest(pub [u8; 32]); #[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, Deserialize)] pub struct DataUsageCacheInfo { pub name: String, pub next_cycle: u64, pub last_update: Option, pub skip_healing: bool, pub lifecycle: Option>, pub replication: Option>, #[serde(default)] pub failed_objects: HashMap, #[serde(default)] pub scan_resume_after: Option, #[serde(default)] pub scan_checkpoint: Option, #[serde(default)] pub scan_raw_enumeration_cursor: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub scan_raw_enumeration_page_index: Option, #[serde(default)] pub scan_identity: Option, #[serde(default)] pub scan_progress: Option, #[serde(default)] pub scan_coverage_receipt: Option, #[serde(default)] pub pending_heals: Vec, #[serde(default)] pub object_lock: Option>, #[serde(default)] pub leader_epoch: u64, #[serde(default)] pub source: Option, #[serde(default)] pub snapshot_complete: bool, #[serde(default)] pub scan_plan_digest: Option, /// Full activity and inventory scope of a set scan; only a complete /// snapshot proves coverage. Bucket caches bind this scope into their /// opaque scan plan digest instead. #[serde(default)] pub scan_coverage_digest: Option, #[serde(default)] pub cache_key_format: u16, /// Registry generation used for the completed/partial scan. This is /// process-local audit data; older cache writers omit it. #[serde(default, skip_serializing_if = "Option::is_none")] pub tier_registry_generation: Option, /// 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, /// Whether the entries retained while a set scan was incomplete come /// from a prior complete set snapshot. This is observational input only. #[serde(default)] pub lkg_snapshot_complete: bool, #[serde(default)] pub lkg_next_cycle: Option, #[serde(default)] pub lkg_last_update: Option, #[serde(default)] pub lkg_leader_epoch: Option, #[serde(default)] pub lkg_scan_plan_digest: Option, /// Activity-sensitive identity for same-cycle set snapshot reuse. The /// structural plan remains reusable across ordinary bucket writes. #[serde(default)] pub scan_execution_digest: Option, /// Durable bucket incarnations captured for a complete set aggregate. /// Missing or nil entries are legacy/unproven and cannot authorize /// skipping an unselected bucket in a later scoped set scan. #[serde(default, skip_serializing_if = "HashMap::is_empty")] pub scan_bucket_incarnations: HashMap, } impl Serialize for DataUsageCacheInfo { fn serialize(&self, serializer: S) -> Result where S: serde::Serializer, { // Keep this metadata map-encoded so older readers can ignore fields // appended by newer scanner versions during rolling upgrades. let field_count = 16 + usize::from(self.scan_raw_enumeration_cursor.is_some()) + usize::from(self.scan_raw_enumeration_page_index.is_some()) + usize::from(self.scan_identity.is_some()) + usize::from(self.scan_progress.is_some()) + usize::from(self.scan_coverage_receipt.is_some()) + usize::from(self.scan_coverage_digest.is_some()) + usize::from(self.tier_registry_generation.is_some()) + usize::from(!self.size_reconciliation.is_empty()) + usize::from(self.lkg_snapshot_complete) + usize::from(self.lkg_next_cycle.is_some()) + usize::from(self.lkg_last_update.is_some()) + usize::from(self.lkg_leader_epoch.is_some()) + usize::from(self.lkg_scan_plan_digest.is_some()) + usize::from(self.scan_execution_digest.is_some()) + usize::from(!self.scan_bucket_incarnations.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)?; state.serialize_entry("last_update", &self.last_update)?; state.serialize_entry("skip_healing", &self.skip_healing)?; state.serialize_entry("lifecycle", &self.lifecycle)?; state.serialize_entry("replication", &self.replication)?; state.serialize_entry("failed_objects", &self.failed_objects)?; state.serialize_entry("scan_resume_after", &self.scan_resume_after)?; state.serialize_entry("scan_checkpoint", &self.scan_checkpoint)?; if let Some(cursor) = &self.scan_raw_enumeration_cursor { state.serialize_entry("scan_raw_enumeration_cursor", cursor)?; } if let Some(index) = &self.scan_raw_enumeration_page_index { state.serialize_entry("scan_raw_enumeration_page_index", index)?; } if let Some(identity) = self.scan_identity { state.serialize_entry("scan_identity", &identity)?; } if let Some(progress) = self.scan_progress { state.serialize_entry("scan_progress", &progress)?; } if let Some(receipt) = &self.scan_coverage_receipt { state.serialize_entry("scan_coverage_receipt", receipt)?; } state.serialize_entry("pending_heals", &self.pending_heals)?; state.serialize_entry("object_lock", &self.object_lock)?; state.serialize_entry("source", &self.source)?; state.serialize_entry("snapshot_complete", &self.snapshot_complete)?; state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?; if let Some(coverage) = self.scan_coverage_digest { state.serialize_entry("scan_coverage_digest", &coverage)?; } state.serialize_entry("cache_key_format", &self.cache_key_format)?; if let Some(generation) = self.tier_registry_generation { state.serialize_entry("tier_registry_generation", &generation)?; } if !self.size_reconciliation.is_empty() { state.serialize_entry("size_reconciliation", &self.size_reconciliation)?; } if self.lkg_snapshot_complete { state.serialize_entry("lkg_snapshot_complete", &true)?; } if let Some(next_cycle) = self.lkg_next_cycle { state.serialize_entry("lkg_next_cycle", &next_cycle)?; } if let Some(last_update) = self.lkg_last_update { state.serialize_entry("lkg_last_update", &last_update)?; } if let Some(leader_epoch) = self.lkg_leader_epoch { state.serialize_entry("lkg_leader_epoch", &leader_epoch)?; } if let Some(scan_plan_digest) = self.lkg_scan_plan_digest { state.serialize_entry("lkg_scan_plan_digest", &scan_plan_digest)?; } if let Some(scan_execution_digest) = self.scan_execution_digest { state.serialize_entry("scan_execution_digest", &scan_execution_digest)?; } if !self.scan_bucket_incarnations.is_empty() { state.serialize_entry("scan_bucket_incarnations", &self.scan_bucket_incarnations)?; } state.end() } } /// Data usage cache #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageCache { pub info: DataUsageCacheInfo, pub cache: HashMap, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(crate) enum DataUsageCachePrepareOutcome { Reused, Reset, RejectedNewerCycle, RejectedNewerLeader, } impl DataUsageCache { /// Reconcile tier keys loaded from an older cache against the registry /// frozen for this scan. New metadata is already routed through /// `UNKNOWN_TIER`; this pass handles retired keys that predate that rule. /// Legacy `TierStats` carries physical bytes only, so this migration does /// not manufacture a logical unknown-byte value from that physical total. pub(crate) fn fold_retired_tiers(&mut self, tier_names: &[String]) { let known_tiers = tier_names.iter().map(String::as_str).collect::>(); for entry in self.cache.values_mut() { let Some(tiers) = entry.all_tier_stats.as_mut() else { continue }; let existing_unknown = tiers.tiers.get(UNKNOWN_TIER).cloned().unwrap_or_default(); let companion_present = entry.unknown_tier_stats.as_ref().is_some_and(|stats| !stats.is_empty()); let migrate_existing_unknown = !companion_present; let mut retired = TierStats::default(); let mut retired_key_found = false; if migrate_existing_unknown { retired = retired.add(&existing_unknown); } for (tier, stats) in &tiers.tiers { if tier != UNKNOWN_TIER && tier != storageclass::STANDARD && tier != storageclass::RRS && !known_tiers.contains(tier.as_str()) { retired_key_found = true; retired = retired.add(stats); } } tiers.fold_unknown_tiers(tier_names.iter().map(String::as_str)); if !retired.is_empty() && !companion_present { entry.add_unknown_tier_stats(&UnknownTierStats { // The legacy map stores physical bytes only. Logical // bytes remain zero until a fresh object scan observes // them under the current metadata format. unknown_physical_bytes: retired.total_size, unknown_objects: retired.num_objects, unknown_versions: retired.num_versions, ..Default::default() }); // The legacy tier map has no logical-byte dimension, so a // proof that classified this retired key as known cannot be // repaired safely. Mark it unvalidated and require a fresh // scan rather than guessing a logical subtraction. entry.tier_accounting_proof = None; } else if retired_key_found { // A nonempty companion has no provenance tying it to the // retired map keys. Reject the mixed cache until a fresh scan // reconciles the dimensions instead of double-counting them. entry.tier_accounting_proof = None; } } } /// Prefix-level usage query over this (writer-side) cache; see /// [`prefix_usage_in_cache`] for the semantics /// (rustfs/backlog#1872). pub fn prefix_usage(&self, bucket: &str, prefix: &str, max_entries: usize) -> Option { prefix_usage_in_cache(&self.cache, bucket, prefix, max_entries) } pub(crate) fn prepare_for_scan( &mut self, name: &str, next_cycle: u64, leader_epoch: u64, source: DataUsageCacheSource, scan_plan_digest: DataUsageScanPlanDigest, require_source: bool, ) -> DataUsageCachePrepareOutcome { if self.info.next_cycle > next_cycle { return DataUsageCachePrepareOutcome::RejectedNewerCycle; } if self.info.leader_epoch > leader_epoch { return DataUsageCachePrepareOutcome::RejectedNewerLeader; } let source_matches = self.info.source == Some(source); let plan_matches = self.info.scan_plan_digest == Some(scan_plan_digest); let metadata_is_reusable = self.info.name == name && self.info.leader_epoch == leader_epoch && plan_matches && (source_matches || (!require_source && self.info.source.is_none())) && self.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT; let reusable = metadata_is_reusable && (self.cache.is_empty() || if name == DATA_USAGE_ROOT || self.info.snapshot_complete { self.checked_flatten_complete_scope(name).is_some() } else { self.checked_flatten(name).is_some() }); if !reusable { 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(), 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; self.info.leader_epoch = leader_epoch; self.info.source = Some(source); self.info.scan_plan_digest = Some(scan_plan_digest); self.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT; self.info.snapshot_complete = false; if reusable { DataUsageCachePrepareOutcome::Reused } else { DataUsageCachePrepareOutcome::Reset } } pub(crate) fn prepare_bucket_checkpoint( &mut self, name: &str, next_cycle: u64, leader_epoch: u64, source: DataUsageCacheSource, scan_plan_digest: DataUsageScanPlanDigest, identity: DataUsageScanIdentity, ) -> DataUsageCachePrepareOutcome { if self.info.next_cycle > next_cycle { return DataUsageCachePrepareOutcome::RejectedNewerCycle; } if self.info.leader_epoch > leader_epoch { return DataUsageCachePrepareOutcome::RejectedNewerLeader; } let reusable = identity.is_valid() && name != DATA_USAGE_ROOT && self.info.name == name && self.info.source == Some(source) && self.info.leader_epoch == leader_epoch && self.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT && self.info.scan_identity == Some(identity) && self.info.tier_registry_generation == Some(identity.tier_registry_generation) && (self.cache.is_empty() || self.checked_flatten_complete_scope(name).is_some()); if reusable && self.info.snapshot_complete && self.info.scan_progress.is_none() && self.info.scan_checkpoint.is_none() && self.info.scan_raw_enumeration_cursor.is_none() && self.info.scan_raw_enumeration_page_index.is_none() && self.info.scan_resume_after.is_none() && self.info.scan_coverage_receipt.is_none() && self.info.scan_plan_digest == Some(scan_plan_digest) { return self.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, true); } if !reusable { let keep_debts = self.info.name == name && self .info .scan_identity .is_none_or(|previous| previous.bucket_incarnation == identity.bucket_incarnation); let (pending_heals, size_reconciliation) = if keep_debts { ( std::mem::take(&mut self.info.pending_heals), std::mem::take(&mut self.info.size_reconciliation), ) } else { (Vec::new(), HashMap::new()) }; *self = Self::default(); self.info.pending_heals = pending_heals; self.info.size_reconciliation = size_reconciliation; } if self.validated_raw_enumeration_cursor().is_none() { self.info.scan_raw_enumeration_cursor = None; } if self.validated_raw_enumeration_page_index().is_none() { self.info.scan_raw_enumeration_page_index = None; } let cursor_is_valid = (self.info.scan_checkpoint.is_none() && self.info.scan_raw_enumeration_cursor.is_none() && self.info.scan_raw_enumeration_page_index.is_none() && self.info.scan_resume_after.is_none() && self.info.scan_coverage_receipt.is_none()) || self.validated_scan_frontier().is_some() || self.info.scan_raw_enumeration_cursor.is_some() || self.info.scan_raw_enumeration_page_index.is_some(); if !cursor_is_valid { self.info.scan_progress = None; } self.info.name = name.to_owned(); self.info.next_cycle = next_cycle; self.info.leader_epoch = leader_epoch; self.info.source = Some(source); self.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT; self.info.tier_registry_generation = Some(identity.tier_registry_generation); self.info.scan_identity = Some(identity); self.info.snapshot_complete = false; if let Some(progress) = &mut self.info.scan_progress { progress.requested_plan = scan_plan_digest; } else { self.info.scan_progress = Some(DataUsageScanProgress { started_plan: scan_plan_digest, requested_plan: scan_plan_digest, }); self.info.scan_resume_after = None; self.info.scan_checkpoint = None; self.info.scan_raw_enumeration_cursor = None; self.info.scan_raw_enumeration_page_index = None; self.info.scan_coverage_receipt = None; } // Old readers do not understand coverage sweeps. An absent plan makes // their existing prepare path rebuild instead of promoting mixed data. self.info.scan_plan_digest = None; if reusable { DataUsageCachePrepareOutcome::Reused } else { DataUsageCachePrepareOutcome::Reset } } fn coverage_prefix_digest(&self, through: &str) -> Result<[u8; 32], serde_json::Error> { let mut writer = CheckpointDigestWriter(Sha256::new()); serde_json::to_writer( &mut writer, &( &self.info.name, self.info.scan_identity, self.info.source, self.info.leader_epoch, self.info.cache_key_format, self.info.scan_progress.map(|progress| progress.started_plan), through, ), )?; let mut prefix = self .cache .iter() .filter(|(key, _)| { let ancestor = through .strip_prefix(key.as_str()) .is_some_and(|suffix| suffix.starts_with('/')); let descendant = key.strip_prefix(through).is_some_and(|suffix| suffix.starts_with('/')); (key.as_str() <= through && !ancestor) || descendant }) .collect::>(); prefix.sort_unstable_by_key(|(key, _)| *key); for (key, entry) in prefix { let mut value = serde_json::to_value(entry)?; value.sort_all_objects(); if let Some(children) = value.get_mut("children").and_then(serde_json::Value::as_array_mut) { children.sort_unstable_by(|left, right| left.as_str().cmp(&right.as_str())); } serde_json::to_writer(&mut writer, &(key, value))?; } Ok(writer.0.finalize().into()) } pub(crate) fn validated_scan_frontier(&self) -> Option<&str> { let receipt = self.info.scan_coverage_receipt.as_ref()?; let checkpoint = self.info.scan_checkpoint.as_ref()?; (self.info.scan_progress.is_some() && self.info.scan_identity.is_some_and(|identity| identity.is_valid()) && self.info.source.is_some() && receipt.through.len() <= 16 * 1024 && checkpoint.version == DATA_USAGE_SCAN_CHECKPOINT_VERSION && checkpoint.resume_after == receipt.through && self.info.scan_resume_after.as_deref() == Some(receipt.through.as_str()) && receipt .through .strip_prefix(&self.info.name) .is_some_and(|suffix| suffix.starts_with('/')) && self.find(&receipt.through).is_some() && self.coverage_prefix_digest(&receipt.through).ok() == Some(receipt.digest)) .then_some(receipt.through.as_str()) } pub(crate) fn validated_raw_enumeration_cursor(&self) -> Option<&DataUsageRawEnumerationCursor> { let cursor = self.info.scan_raw_enumeration_cursor.as_ref()?; (self.info.scan_progress.is_some() && self.info.scan_identity.is_some_and(|identity| identity.is_valid()) && self.info.source.is_some() && cursor.is_valid_for_bucket(&self.info.name)) .then_some(cursor) } pub(crate) fn validated_raw_enumeration_page_index(&self) -> Option<&RawEnumerationPageIndex> { let index = self.info.scan_raw_enumeration_page_index.as_ref()?; if self.info.scan_progress.is_none() || !self.info.scan_identity.is_some_and(|identity| identity.is_valid()) || self.info.source.is_none() || index.committed_entries().is_err() || index.indexed_entries().is_err() { return None; } let parent = match index.status() { RawEnumerationPageOwnerStatus::Unsupported => return None, RawEnumerationPageOwnerStatus::Building { parent, .. } | RawEnumerationPageOwnerStatus::Ready { parent, .. } => { parent } }; path_is_in_bucket_scope(&self.info.name, &parent).then_some(index) } /// Seal only the frontier supplied by completed traversal, never a restored cursor. pub(crate) fn seal_scan_frontier(&mut self, frontier: Option<&str>) -> Result<(), serde_json::Error> { if self.info.scan_progress.is_none() { self.info.scan_coverage_receipt = None; return Ok(()); } let frontier = frontier.filter(|path| path.len() <= 16 * 1024 && self.find(path).is_some()); self.info.scan_coverage_receipt = match frontier { Some(through) => Some(DataUsageScanCoverageReceipt { through: through.to_owned(), digest: self.coverage_prefix_digest(through)?, }), None => None, }; self.info.scan_resume_after = frontier.map(str::to_owned); let reason = self .info .scan_checkpoint .as_ref() .map_or(DataUsageScanCheckpointReason::Unknown, |checkpoint| checkpoint.reason); self.info.scan_checkpoint = frontier.map(|through| DataUsageScanCheckpoint::new(through.to_owned(), reason)); Ok(()) } fn ensure_cache_save_metrics_registered() { CACHE_SAVE_METRICS_ONCE.call_once(|| { describe_counter!( METRIC_CACHE_SAVE_ATTEMPT_TOTAL, "Total scanner data usage cache save attempts by result and cache type." ); describe_counter!( METRIC_CACHE_SAVE_TIMEOUT_TOTAL, "Total scanner data usage cache save timeouts by cache type." ); describe_counter!( METRIC_CACHE_SAVE_RETRY_TOTAL, "Total scanner data usage cache save retries by cache type." ); describe_histogram!( METRIC_CACHE_SAVE_DURATION_SECONDS, "Duration of scanner data usage cache save attempts in seconds." ); }); } fn cache_path_type(path: &str) -> &'static str { if path.ends_with(".bkp") { "backup" } else { "main" } } pub fn replace(&mut self, path: &str, parent: &str, e: DataUsageEntry) { let hash = hash_path(path); self.cache.insert(hash.key(), e); if !parent.is_empty() { let parent_hash = hash_path(parent); self.cache.entry(parent_hash.key()).or_default().add_child(&hash); } } pub fn replace_hashed(&mut self, hash: &DataUsageHash, parent: &Option, e: &DataUsageEntry) { self.cache.insert(hash.key(), e.clone()); if let Some(parent) = parent { self.cache.entry(parent.key()).or_default().add_child(hash); } } pub fn find(&self, path: &str) -> Option<&DataUsageEntry> { self.cache.get(&hash_path(path).key()) } pub fn find_children_copy(&mut self, h: DataUsageHash) -> DataUsageHashMap { self.cache.entry(h.string()).or_default().children.clone() } pub fn flatten(&self, root: &DataUsageEntry) -> DataUsageEntry { let mut visited = HashSet::new(); self.flatten_with_guard(root, &mut visited, 0) } pub(crate) fn checked_flatten(&self, path: &str) -> Option { self.checked_flatten_inner(path).map(|(entry, _)| entry) } pub(crate) fn checked_flatten_complete(&self, path: &str) -> Option { self.checked_flatten_inner(path) .filter(|(_, visited)| *visited == self.cache.len()) .map(|(entry, _)| entry) } pub(crate) fn checked_flatten_complete_scope(&self, path: &str) -> Option { if path == DATA_USAGE_ROOT { return self.checked_flatten_complete(path); } let (entry, visited) = self.checked_flatten_inner(path)?; let root_parent_only = { let path_key = hash_path(path).key(); self.cache .get(DATA_USAGE_ROOT) .is_some_and(|root| root_is_parent_only(root, &path_key)) }; let expected_entries = self.cache.len().saturating_sub(usize::from(root_parent_only)); (visited == expected_entries).then_some(entry) } pub(crate) fn has_complete_root_inventory(&self, bucket_keys: &HashSet) -> bool { let Some(root) = self.find(DATA_USAGE_ROOT) else { return false; }; // Set roots only connect bucket entries. Scalar data at the root, an // extra bucket, or an orphan must not disappear during bucket folding. root.children.len() == bucket_keys.len() && bucket_keys.iter().all(|key| root.children.contains(key)) && root.size == 0 && root.objects == 0 && root.versions == 0 && root.delete_markers == 0 && root.failed_objects == 0 && !root.compacted && root.obj_sizes.is_empty() && root.obj_versions.is_empty() && root.replication_stats.is_none() && root.all_tier_stats.is_none() && root.unknown_tier_stats.is_none() && root.tier_accounting_proof.is_none() && self.checked_flatten_complete(DATA_USAGE_ROOT).is_some() } fn checked_flatten_inner(&self, path: &str) -> Option<(DataUsageEntry, usize)> { let root_key = hash_path(path).key(); let (root_key, root) = self.cache.get_key_value(&root_key)?; if root.compacted && !root.children.is_empty() { return None; } let mut visited = HashSet::from([root_key.as_str()]); let mut pending = root.children.iter().map(|child| (child.as_str(), 1usize)).collect::>(); let mut flattened = DataUsageEntry::default(); if !flattened.checked_merge(root) { return None; } flattened.compacted = root.compacted; while let Some((key, depth)) = pending.pop() { if depth > MAX_DATA_USAGE_CACHE_DEPTH || !visited.insert(key) { return None; } let entry = self.cache.get(key)?; if (entry.compacted || depth == MAX_DATA_USAGE_CACHE_DEPTH) && !entry.children.is_empty() { return None; } pending.extend(entry.children.iter().map(|child| (child.as_str(), depth + 1))); if !flattened.checked_merge(entry) { return None; } } Some((flattened, visited.len())) } 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_with_guard(&e, visited, depth + 1); } root.merge(&e); } } root.children.clear(); root } 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()); 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); } } if let Some(parent) = parent { self.cache.entry(parent.key()).or_default().add_child(hash); } } } 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() { need_remove.push(child.clone()); } } self.cache.remove(&hash.string()); if depth >= MAX_DATA_USAGE_CACHE_DEPTH { return; } for child in need_remove { self.delete_recursive_guard(&DataUsageHash(child), visited, depth + 1); } } pub fn size_recursive(&self, path: &str) -> Option { match self.find(path) { Some(root) => { if root.children.is_empty() { return Some(root.clone()); } 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.is_empty()) { flat.replication_stats = None; } Some(flat) } None => None, } } pub fn search_parent(&self, hash: &DataUsageHash) -> Option { let want = hash.key(); if let Some(last_index) = want.rfind('/') && let Some(v) = self.find(&want[0..last_index]) && v.children.contains(&want) { return Some(hash_path(&want[0..last_index])); } for (k, v) in self.cache.iter() { if v.children.contains(&want) { return Some(DataUsageHash(k.clone())); } } None } pub fn is_compacted(&self, hash: &DataUsageHash) -> bool { self.cache.get(&hash.key()).is_some_and(|due| due.compacted) } pub fn force_compact(&mut self, limit: usize) { if self.cache.len() < limit { return; } let top = hash_path(&self.info.name).key(); let Some(top_e) = self.find(&top).cloned() else { return; }; if top_e.children.len() > 250_000 { self.reduce_children_of(&hash_path(&self.info.name), limit, true); } if self.cache.len() <= limit { return; } let mut found = HashSet::new(); found.insert(top); mark(self, &top_e, &mut found); self.cache.retain(|k, _| found.contains(k)); } pub fn reduce_children_of(&mut self, path: &DataUsageHash, limit: usize, compact_self: bool) { let Some(e) = self.cache.get(&path.key()).cloned() else { return; }; if e.compacted { return; } if e.children.len() > limit && compact_self { let mut flat = self.size_recursive(&path.key()).unwrap_or_default(); flat.compacted = true; self.delete_recursive(path); self.replace_hashed(path, &None, &flat); return; } let total = self.total_children_rec(&path.key()); if total < limit { return; } let mut candidates = Vec::new(); let mut remove = total - limit; add(self, path, &mut candidates); candidates.sort_by_key(|a| a.objects); let mut candidate_index = 0; while remove > 0 && candidate_index < candidates.len() { let e = &candidates[candidate_index]; let candidate = e.path.clone(); if candidate == *path && !compact_self { break; } let removing = self.total_children_rec(&candidate.key()); let mut flat = match self.size_recursive(&candidate.key()) { Some(flat) => flat, None => { candidate_index += 1; continue; } }; flat.compacted = true; self.delete_recursive(&candidate); self.replace_hashed(&candidate, &None, &flat); remove = remove.saturating_sub(removing); candidate_index += 1; } } 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() || depth >= MAX_DATA_USAGE_CACHE_DEPTH { return 0; } let mut n = 0; for ch in root.children.iter() { if visited.insert(ch.clone()) { n += 1 + self.total_children_rec_guard(ch, visited, depth + 1); } } n } pub fn merge(&mut self, o: &DataUsageCache) { let Some(mut existing_root) = self.root() else { if o.root().is_none() { return; } *self = o.clone(); return; }; let Some(other_root) = o.root() else { return; }; if o.info.last_update > self.info.last_update { self.info.last_update = o.info.last_update; } existing_root.merge(&other_root); self.cache.insert(hash_path(&self.info.name).key(), existing_root); let root_hash = self.root_hash(); for key in other_root.children.iter() { let Some(entry) = o.cache.get(key) else { continue; }; let flat = o.flatten(entry); if let Some(existing) = self.cache.get_mut(key) { existing.merge(&flat); } else { self.replace_hashed(&DataUsageHash(key.clone()), &Some(root_hash.clone()), &flat); } } } pub fn root_hash(&self) -> DataUsageHash { hash_path(&self.info.name) } pub fn root(&self) -> Option { self.find(&self.info.name).cloned() } /// Convert cache to DataUsageInfo for a specific path pub fn dui(&self, path: &str, buckets: &[String]) -> DataUsageInfo { let e = match self.find(path) { Some(e) => e, None => return DataUsageInfo::default(), }; let flat = self.flatten(e); let mut buckets_usage = HashMap::new(); for bucket_name in buckets.iter() { let e = match self.find(bucket_name) { Some(e) => e, None => continue, }; let flat = self.flatten(e); let mut bui = BucketUsageInfo { size: flat.size as u64, versions_count: flat.versions as u64, objects_count: flat.objects as u64, delete_markers_count: flat.delete_markers as u64, object_size_histogram: flat.obj_sizes.to_map(), object_versions_histogram: flat.obj_versions.to_map(), ..Default::default() }; if let Some(rs) = &flat.replication_stats { bui.replica_size = rs.replica_size; bui.replica_count = rs.replica_count; for (arn, stat) in rs.targets.iter() { bui.replication_info.insert( arn.clone(), BucketTargetUsageInfo { replication_pending_size: stat.pending_size, replicated_size: stat.replicated_size, replication_failed_size: stat.failed_size, replication_pending_count: stat.pending_count, replication_failed_count: stat.failed_count, replicated_count: stat.replicated_count, ..Default::default() }, ); } } buckets_usage.insert(bucket_name.clone(), bui); } DataUsageInfo { last_update: self.info.last_update, objects_total_count: flat.objects as u64, versions_total_count: flat.versions as u64, delete_markers_total_count: flat.delete_markers as u64, objects_total_size: flat.size as u64, tier_stats: flat.all_tier_stats.filter(|tiers| !tiers.is_empty()), unknown_tier_stats: flat.unknown_tier_stats.filter(|stats| !stats.is_empty()), buckets_count: u64::try_from(buckets.len()).unwrap_or(u64::MAX), buckets_usage, ..Default::default() } } pub fn marshal_msg(&self) -> Result, Box> { let mut buf = Vec::new(); self.serialize(&mut rmp_serde::Serializer::new(&mut buf))?; Ok(buf) } pub fn unmarshal(buf: &[u8]) -> Result> { let t: Self = rmp_serde::from_slice(buf)?; Ok(t) } } mod persistence; #[derive(Default, Clone)] struct Inner { objects: usize, path: DataUsageHash, } 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; 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 — // total_children_rec returns 0 for leaves, so `remove` would never decrement. if !e.children.is_empty() { candidates.push(Inner { objects, path: path.clone(), }); } objects } 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() { if !found.insert(k.to_string()) { continue; } if let Some(ch) = duc.cache.get(k) { mark_with_depth(duc, ch, found, depth + 1); } } } fn root_is_parent_only(root: &DataUsageEntry, child: &str) -> bool { root.children.len() == 1 && root.children.contains(child) && root.size == 0 && root.objects == 0 && root.versions == 0 && root.delete_markers == 0 && root.replication_stats.is_none() && !root.compacted && root.failed_objects == 0 } /// Trait for storage-specific operations on DataUsageCache #[async_trait::async_trait] pub trait DataUsageCacheStorage { /// Load data usage cache from backend storage async fn load(store: &dyn std::any::Any, name: &str) -> Result> where Self: Sized; /// Save data usage cache to backend storage async fn save(&self, name: &str) -> Result<(), Box>; } #[cfg(test)] mod tests;