fix(scanner): unify unknown metadata size accounting (#6394)

* fix(scanner): unify unknown metadata size accounting

* fix(scanner): preserve restore expiry semantics

* fix(ci): resolve ecstore clippy warnings

* fix(scanner): close lifecycle review gaps

---------

Signed-off-by: houseme <housemecn@gmail.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
cxymds
2026-08-23 17:28:52 +08:00
committed by GitHub
parent e196a134cc
commit 76a863b3ea
7 changed files with 1212 additions and 81 deletions
+50 -6
View File
@@ -29,8 +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, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME,
PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path,
prefix_usage_in_cache,
PrefixUsageEntry, 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};
@@ -201,6 +201,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 {
@@ -234,6 +238,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 =====
@@ -353,6 +385,10 @@ pub struct DataUsageCacheInfo {
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
#[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<String, SizeReconciliationEntry>,
/// Whether the entries retained while a set scan was incomplete come
/// from a prior complete set snapshot. This is observational input only.
#[serde(default)]
@@ -374,7 +410,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(21))?;
let field_count = 21 + 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)?;
@@ -391,6 +428,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.serialize_entry("lkg_snapshot_complete", &self.lkg_snapshot_complete)?;
state.serialize_entry("lkg_next_cycle", &self.lkg_next_cycle)?;
state.serialize_entry("lkg_last_update", &self.lkg_last_update)?;
@@ -454,14 +494,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;