fix(scanner): fence unknown tier accounting (#6396)

This commit is contained in:
cxymds
2026-08-23 19:28:43 +08:00
committed by GitHub
parent 14e3eb787d
commit b2e60be647
29 changed files with 2827 additions and 259 deletions
+199 -6
View File
@@ -100,13 +100,14 @@ where
let _ = tokio::time::timeout(SCANNER_CACHE_LOCK_LOSS_SHUTDOWN_TIMEOUT, scan).await;
}
pub(crate) fn current_cache_root_entry(
pub(crate) fn current_cache_root_entry_with_generation(
cache: &DataUsageCache,
name: &str,
source: DataUsageCacheSource,
next_cycle: u64,
leader_epoch: u64,
scan_plan_digest: DataUsageScanPlanDigest,
tier_registry_generation: Option<u64>,
) -> std::result::Result<Option<DataUsageEntryInfo>, ScannerError> {
let metadata_is_current = cache.info.name == name
&& cache.info.source == Some(source)
@@ -115,7 +116,8 @@ pub(crate) fn current_cache_root_entry(
&& cache.info.last_update.is_some()
&& cache.info.next_cycle == next_cycle
&& cache.info.leader_epoch == leader_epoch
&& cache.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT;
&& cache.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT
&& tier_registry_generation.is_none_or(|generation| cache.info.tier_registry_generation == Some(generation));
if !metadata_is_current {
return Ok(None);
}
@@ -131,6 +133,13 @@ pub(crate) enum DataUsageCacheScanState {
},
}
#[derive(Clone, Copy, Debug, Default)]
pub(crate) struct DataUsageCacheReuseOptions {
pub(crate) require_source: bool,
pub(crate) tier_registry_generation: Option<u64>,
}
#[cfg(test)]
pub(crate) fn current_cache_root_or_prepare(
cache: &mut DataUsageCache,
name: &str,
@@ -140,11 +149,51 @@ pub(crate) fn current_cache_root_or_prepare(
scan_plan_digest: DataUsageScanPlanDigest,
require_source: bool,
) -> DataUsageCacheScanState {
match current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest) {
current_cache_root_or_prepare_with_generation(
cache,
name,
source,
next_cycle,
leader_epoch,
scan_plan_digest,
DataUsageCacheReuseOptions {
require_source,
tier_registry_generation: None,
},
)
}
pub(crate) fn current_cache_root_or_prepare_with_generation(
cache: &mut DataUsageCache,
name: &str,
source: DataUsageCacheSource,
next_cycle: u64,
leader_epoch: u64,
scan_plan_digest: DataUsageScanPlanDigest,
options: DataUsageCacheReuseOptions,
) -> DataUsageCacheScanState {
if options.tier_registry_generation.is_some_and(|generation| {
cache.info.next_cycle <= next_cycle
&& cache.info.leader_epoch <= leader_epoch
&& cache.info.tier_registry_generation != Some(generation)
}) {
// Make prepare_for_scan take its reset path so an entry classified by
// an older registry cannot be reused under the new cycle generation.
cache.info.scan_plan_digest = None;
}
match current_cache_root_entry_with_generation(
cache,
name,
source,
next_cycle,
leader_epoch,
scan_plan_digest,
options.tier_registry_generation,
) {
Ok(Some(root)) => DataUsageCacheScanState::Current(Box::new(root)),
current => DataUsageCacheScanState::Prepared {
invalid_current: current.err(),
outcome: cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, require_source),
outcome: cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, options.require_source),
},
}
}
@@ -159,7 +208,7 @@ pub(super) fn cache_snapshot_is_current(
scan_plan_digest: DataUsageScanPlanDigest,
) -> bool {
matches!(
current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest),
current_cache_root_entry_with_generation(cache, name, source, next_cycle, leader_epoch, scan_plan_digest, None),
Ok(Some(_))
)
}
@@ -168,6 +217,7 @@ pub(super) fn completed_data_usage_info(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
cancelled: bool,
@@ -183,6 +233,19 @@ pub(super) fn completed_data_usage_info(
return None;
}
// A generation is comparable across nodes because it is derived from the
// frozen registry names. Cycle and leader fencing remain separate cache
// metadata. Legacy peers omit the generation; an all-legacy result remains
// readable, but mixing legacy and new (or two new generations) would make
// the per-tier accounting ambiguous.
let registry_generation = results.first()?.info.tier_registry_generation;
if results.iter().any(|result| match registry_generation {
Some(generation) => result.info.tier_registry_generation != Some(generation),
None => result.info.tier_registry_generation.is_some(),
}) {
return None;
}
if results.iter().any(|result| result.root().is_none()) {
return None;
}
@@ -200,9 +263,16 @@ pub(super) fn completed_data_usage_info(
if !total.checked_merge(&merged) {
return None;
}
if !tier_accounting_proof_is_publishable(&merged, registry_generation, tier_registry_names) {
return None;
}
bucket_entries.insert(bucket.clone(), merged);
}
if !tier_accounting_proof_is_publishable(&total, registry_generation, tier_registry_names) {
return None;
}
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
let buckets_usage = bucket_entries
.iter()
@@ -220,6 +290,7 @@ pub(super) fn completed_data_usage_info(
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
objects_total_size: u64::try_from(total.size).ok()?,
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
unknown_tier_stats: total.unknown_tier_stats.filter(|stats| !stats.is_empty()),
buckets_count: u64::try_from(all_buckets.len()).ok()?,
bucket_sizes,
buckets_usage,
@@ -229,6 +300,109 @@ pub(super) fn completed_data_usage_info(
Some((data_usage_info, merged_last_update))
}
fn tier_accounting_proof_is_publishable(
entry: &DataUsageEntry,
registry_generation: Option<u64>,
tier_registry_names: &[String],
) -> bool {
let has_scalar_usage = entry.size > 0
|| entry.objects > 0
|| entry.versions > 0
|| entry.delete_markers > 0
|| entry.failed_objects > 0
|| !entry.obj_sizes.is_empty()
|| !entry.obj_versions.is_empty()
|| entry.replication_stats.as_ref().is_some_and(|stats| !stats.is_empty());
let has_tier_accounted_data = entry
.all_tier_stats
.as_ref()
.is_some_and(|stats| stats.tiers.values().any(|tier| !tier.is_empty()))
|| entry.unknown_tier_stats.as_ref().is_some_and(|stats| {
stats.counter_overflowed
|| stats.unknown_bytes > 0
|| stats.unknown_physical_bytes > 0
|| stats.unknown_objects > 0
|| stats.unknown_versions > 0
});
let Some(proof) = entry.tier_accounting_proof else {
return !has_scalar_usage && !has_tier_accounted_data;
};
if proof.overflowed
|| entry
.unknown_tier_stats
.as_ref()
.is_some_and(|stats| stats.counter_overflowed)
|| u64::try_from(entry.size).ok() != Some(proof.logical_total)
{
return false;
}
let unknown_logical = entry.unknown_tier_stats.as_ref().map_or(0, |stats| stats.unknown_bytes);
let unknown_physical = entry
.unknown_tier_stats
.as_ref()
.map_or(0, |stats| stats.unknown_physical_bytes);
if proof
.logical_known
.checked_add(unknown_logical)
.is_none_or(|total| total != proof.logical_total)
|| proof
.physical_known
.checked_add(unknown_physical)
.is_none_or(|total| total != proof.physical_total)
{
return false;
}
if registry_generation.is_some()
&& entry.all_tier_stats.as_ref().is_some_and(|stats| {
stats.tiers.keys().any(|tier| {
tier != crate::UNKNOWN_TIER
&& tier != crate::storageclass::STANDARD
&& tier != crate::storageclass::RRS
&& !tier_registry_names.iter().any(|allowed| allowed == tier)
})
})
{
return false;
}
if !has_tier_accounted_data {
return true;
}
let Some(tiers) = entry.all_tier_stats.as_ref() else {
return false;
};
let map_unknown_physical = tiers.tiers.get(crate::UNKNOWN_TIER).map_or(0, |stats| stats.total_size);
let companion_unknown_physical = entry
.unknown_tier_stats
.as_ref()
.map_or(0, |stats| stats.unknown_physical_bytes);
if map_unknown_physical != companion_unknown_physical {
return false;
}
// A no-configuration scan intentionally stores only UNKNOWN_TIER after
// the first unknown object; STANDARD/RRS remain absent to preserve the
// historical empty-map shape. In that shape the scalar proof is the sole
// source of known physical bytes. Configured registries seed at least one
// non-UNKNOWN key, whose map total must match the proof.
let has_known_tier_map = tiers.tiers.keys().any(|tier| tier.as_str() != crate::UNKNOWN_TIER);
if !has_known_tier_map {
return true;
}
let Some(known_tier_physical_total) = tiers
.tiers
.iter()
.filter(|(tier, _)| tier.as_str() != crate::UNKNOWN_TIER)
.map(|(_, stats)| stats)
.try_fold(0_u64, |total, stats| total.checked_add(stats.total_size))
else {
return false;
};
proof.physical_known == known_tier_physical_total
}
/// Build a non-authoritative view from the set snapshots that completed this
/// cycle plus compatible per-set last-known-good caches. The caller must
/// persist this result only on the observational object; a missing set is
@@ -238,6 +412,7 @@ pub(super) fn observational_data_usage_info(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
tier_registry_names: &[String],
expected_plan_digest: DataUsageScanPlanDigest,
scanner_cycle: u64,
leader_epoch: u64,
@@ -316,6 +491,13 @@ pub(super) fn observational_data_usage_info(
if usable.is_empty() {
return None;
}
let registry_generation = usable.first()?.0.info.tier_registry_generation;
if usable.iter().any(|(result, _)| match registry_generation {
Some(generation) => result.info.tier_registry_generation != Some(generation),
None => result.info.tier_registry_generation.is_some(),
}) {
return None;
}
let mut total = DataUsageEntry::default();
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
@@ -337,6 +519,15 @@ pub(super) fn observational_data_usage_info(
}
}
}
if bucket_entries
.values()
.any(|entry| !tier_accounting_proof_is_publishable(entry, registry_generation, tier_registry_names))
{
return None;
}
if !tier_accounting_proof_is_publishable(&total, registry_generation, tier_registry_names) {
return None;
}
let merged_last_update = merged_last_update?;
let buckets_usage = bucket_entries
.iter()
@@ -352,6 +543,7 @@ pub(super) fn observational_data_usage_info(
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
objects_total_size: u64::try_from(total.size).ok()?,
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
unknown_tier_stats: total.unknown_tier_stats.filter(|stats| !stats.is_empty()),
buckets_count: u64::try_from(buckets_usage.len()).ok()?,
bucket_sizes: buckets_usage
.iter()
@@ -463,13 +655,14 @@ pub(super) async fn persist_and_publish_cache_snapshot(
return None;
}
if matches!(
current_cache_root_entry(
current_cache_root_entry_with_generation(
&persisted,
DATA_USAGE_ROOT,
source,
cache_snapshot.info.next_cycle,
cache_snapshot.info.leader_epoch,
scan_plan_digest,
cache_snapshot.info.tier_registry_generation,
),
Ok(Some(_))
) {
+21 -5
View File
@@ -31,6 +31,7 @@ impl ScannerIOCache for SetDisks {
all_buckets,
digest: scan_plan_digest,
leader_epoch,
tier_registry_generation,
publication_epoch,
dirty_usage_buckets,
bucket_failures,
@@ -70,6 +71,7 @@ impl ScannerIOCache for SetDisks {
next_cycle: want_cycle,
last_update: Some(now),
leader_epoch,
tier_registry_generation: Some(tier_registry_generation),
source: Some(source),
snapshot_complete: true,
scan_plan_digest: Some(scan_plan_digest),
@@ -267,7 +269,14 @@ impl ScannerIOCache for SetDisks {
record_disk_bucket_scans_active(0, &pool_label, &set_label);
let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone());
let old_lkg = old_cache.info.snapshot_complete.then(|| {
// Fence a stale set aggregate before copying entries into per-bucket work caches.
if old_cache.info.next_cycle <= want_cycle
&& old_cache.info.leader_epoch <= leader_epoch
&& old_cache.info.tier_registry_generation != Some(tier_registry_generation)
{
old_cache.info.scan_plan_digest = None;
}
let old_lkg = old_cache.info.snapshot_complete.then_some({
(
old_cache.info.next_cycle,
old_cache.info.last_update,
@@ -333,6 +342,7 @@ impl ScannerIOCache for SetDisks {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: want_cycle,
leader_epoch,
tier_registry_generation: Some(tier_registry_generation),
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
@@ -394,8 +404,9 @@ impl ScannerIOCache for SetDisks {
};
let mut cache = cache_mutex_clone.lock().await;
apply_bucket_result_to_cache(&mut cache, result, SystemTime::now());
completed_bucket_count_clone.fetch_add(1, Ordering::Relaxed);
if apply_bucket_result_to_cache(&mut cache, result, SystemTime::now()) {
completed_bucket_count_clone.fetch_add(1, Ordering::Relaxed);
}
}
}
}
@@ -527,6 +538,7 @@ impl ScannerIOCache for SetDisks {
session_id: remote_session_id,
session_sequence: request_sequence,
scan_plan_digest: bucket_scan_plan_digest,
tier_registry_generation,
skip_healing: healing,
scan_mode,
},
@@ -742,14 +754,17 @@ impl ScannerIOCache for SetDisks {
continue;
}
};
let scan_state = current_cache_root_or_prepare(
let scan_state = current_cache_root_or_prepare_with_generation(
&mut cache,
&bucket.name,
source,
want_cycle,
leader_epoch,
bucket_scan_plan_digest,
require_cache_source,
DataUsageCacheReuseOptions {
require_source: require_cache_source,
tier_registry_generation: Some(tier_registry_generation),
},
);
let outcome = match scan_state {
DataUsageCacheScanState::Current(root) => {
@@ -1237,6 +1252,7 @@ impl ScannerIOCache for SetDisks {
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.tier_registry_generation = Some(tier_registry_generation);
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
+12
View File
@@ -48,6 +48,7 @@ impl ScannerIOCycle for ECStore {
scan_mode: HealScanMode,
) -> Result<ScannerCycleResult> {
let child_token = ctx.child_token();
let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch);
// Check the local pool metadata before listing buckets. A failed or
// canceled decommission remains suspended after its worker exits, so
@@ -140,6 +141,8 @@ impl ScannerIOCycle for ECStore {
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before));
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
let cache_cycle_floor = Arc::new(AtomicU64::new(want_cycle));
let tier_registry = runtime_tier_registry_for_cycle(want_cycle, leader_epoch).await;
let tier_registry_generation = tier_registry.generation;
if all_buckets.is_empty() {
reset_set_scan_gauges();
@@ -172,6 +175,9 @@ impl ScannerIOCycle for ECStore {
{
return Ok(ScannerCycleResult::new(status, None).with_publication_epoch(publication_epoch));
}
if status == ScannerCycleStatus::Complete {
complete_tier_registry_cycle(want_cycle, leader_epoch);
}
let dirty_usage_clear =
(status == ScannerCycleStatus::Complete).then(|| dirty_usage_snapshot.buckets.as_ref().clone());
let remote_dirty_usage_acknowledgements = if status == ScannerCycleStatus::Complete {
@@ -266,6 +272,7 @@ impl ScannerIOCycle for ECStore {
all_buckets: Arc::clone(&all_buckets),
digest: scan_plan_digest,
leader_epoch,
tier_registry_generation,
publication_epoch,
dirty_usage_buckets: dirty_usage_snapshot.buckets.clone(),
bucket_failures: bucket_failures.clone(),
@@ -399,6 +406,7 @@ impl ScannerIOCycle for ECStore {
&results,
&expected_sources,
&all_bucket_names,
&tier_registry.names,
bucket_plan_complete,
budget_elapsed,
ctx.is_cancelled(),
@@ -410,6 +418,7 @@ impl ScannerIOCycle for ECStore {
&results,
&expected_sources,
&all_bucket_names,
&tier_registry.names,
scan_plan_digest,
want_cycle,
leader_epoch,
@@ -441,6 +450,9 @@ impl ScannerIOCycle for ECStore {
&failed_buckets,
);
result?;
if cycle_status == ScannerCycleStatus::Complete {
complete_tier_registry_cycle(want_cycle, leader_epoch);
}
let remote_dirty_usage_acknowledgements = if cycle_status == ScannerCycleStatus::Complete {
crate::scanner::scanner_dirty_usage_acknowledgements(&activity_before)
} else {
+27 -15
View File
@@ -13,29 +13,38 @@
// limitations under the License.
/// ScannerIODisk implementation for Disk: get_size and the per-disk bucket scan.
use super::*;
use crate::UNKNOWN_TIER;
///
/// Seed [`SizeSummary::tier_stats`] from the cached tier-name list.
///
/// Preserves the original seeding semantics: with no tiers configured the map
/// stays completely empty (STANDARD/RRS are not seeded either); otherwise the
/// standard storage classes are seeded alongside every configured tier so
/// per-object accounting always finds its tier key.
/// Preserves the original no-tier shape: with no tiers configured the map
/// stays completely empty (STANDARD/RRS/UNKNOWN are not seeded either).
/// Otherwise the standard storage classes and one fixed unknown bucket are
/// seeded alongside every configured tier so per-object accounting never
/// inserts an untrusted metadata key.
pub(super) fn tier_stats_template(tier_names: &[String]) -> HashMap<String, TierStats> {
let mut tier_stats = HashMap::with_capacity(tier_names.len() + 2);
let mut tier_stats = HashMap::with_capacity(tier_names.len() + 3);
for tier_name in tier_names {
tier_stats.insert(tier_name.clone(), TierStats::default());
if tier_name != UNKNOWN_TIER {
tier_stats.insert(tier_name.clone(), TierStats::default());
}
}
if !tier_stats.is_empty() {
tier_stats.insert(storageclass::STANDARD.to_string(), TierStats::default());
tier_stats.insert(storageclass::RRS.to_string(), TierStats::default());
tier_stats.insert(UNKNOWN_TIER.to_string(), TierStats::default());
}
tier_stats
}
#[async_trait::async_trait]
impl ScannerIODisk for Disk {
async fn get_size(&self, mut item: ScannerItem) -> Result<SizeSummary> {
async fn get_size(&self, item: ScannerItem) -> Result<SizeSummary> {
self.get_size_with_tier_names(item, &runtime_tier_names().await).await
}
async fn get_size_with_tier_names(&self, mut item: ScannerItem, tier_names: &[String]) -> Result<SizeSummary> {
let done_object = Metrics::time(Metric::ScanObject);
if !is_xl_meta_path(&item.path) {
@@ -105,12 +114,13 @@ impl ScannerIODisk for Disk {
.map(|v| ObjectInfo::from_file_info(v, item.bucket.as_str(), object_path.as_str(), versioned))
.collect::<Vec<ObjectInfo>>();
let mut size_summary = SizeSummary::default();
// Tier names come from the process-wide TTL cache; seeding from them
// replaces the per-object clone of every full TierConfig.
let tier_names = runtime_tier_names().await;
size_summary.tier_stats = tier_stats_template(&tier_names);
// The caller supplies one registry snapshot for the whole folder scan;
// seeding from it prevents a TTL refresh from mixing generations in a
// single result.
let mut size_summary = SizeSummary {
tier_stats: tier_stats_template(tier_names),
..Default::default()
};
let lock_config = object_lock_config_for_scanner_item(&item).await;
@@ -120,12 +130,14 @@ impl ScannerIODisk for Disk {
// `object_infos`.
global_metrics().record_scanner_versions_scanned(object_infos.len() as u64);
item.apply_actions(object_infos, lock_config, versioning_config, &mut size_summary)
item.apply_actions(object_infos, lock_config, versioning_config, tier_names, &mut size_summary)
.await;
if !free_version_infos.is_empty() {
for oi in free_version_infos {
enqueue_runtime_free_version(oi).await;
if ScannerItem::tier_is_known(&oi, tier_names) {
enqueue_runtime_free_version(oi).await;
}
}
}
@@ -13,7 +13,8 @@
// limitations under the License.
use super::*;
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage};
use crate::data_usage_define::{UNKNOWN_TIER, UnknownTierStats, hash_path};
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccountingProof};
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
@@ -89,6 +90,11 @@ fn completed_root_cache(bucket: &str, objects: usize, update_secs: u64, source:
DataUsageEntry {
objects,
size: objects.saturating_mul(10),
tier_accounting_proof: Some(TierAccountingProof {
logical_total: u64::try_from(objects.saturating_mul(10)).unwrap_or(u64::MAX),
logical_known: u64::try_from(objects.saturating_mul(10)).unwrap_or(u64::MAX),
..Default::default()
}),
..Default::default()
},
);
@@ -102,7 +108,7 @@ fn completed_data_usage_info_for_test(
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
let expected_sources = results.iter().filter_map(|result| result.info.source).collect::<HashSet<_>>();
completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled)
completed_data_usage_info(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled)
}
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
@@ -130,9 +136,10 @@ fn partial_usage_is_observational_not_authoritative_for_quota() {
let expected = HashSet::from([current_source, stalled_source]);
assert!(
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, true, false, false).is_none()
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, &[], true, false, false)
.is_none()
);
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
.expect("a completed set should produce an observational view");
assert!(observed.usage_snapshot_partial);
assert!(!observed.usage_snapshot_complete);
@@ -157,7 +164,7 @@ fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
current.info.next_cycle = 8;
current.info.leader_epoch = 3;
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
.expect("complete set data is a valid observational baseline");
assert_eq!(observed.objects_total_size, 30);
assert_eq!(observed.usage_snapshot_set_states[0].complete, true);
@@ -170,7 +177,7 @@ fn negative_delta_waits_for_set_reconciliation() {
let mut stalled = lkg_root_cache("bucket", 4, source);
stalled.info.lkg_scan_plan_digest = Some(DataUsageScanPlanDigest([9; 32]));
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
@@ -220,7 +227,7 @@ fn old_set_completion_cannot_overwrite_new_aggregate() {
old.info.next_cycle = 7;
old.info.leader_epoch = 2;
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
@@ -231,7 +238,7 @@ fn usage_aggregate_survives_restart_and_leader_failover() {
lkg.info.lkg_leader_epoch = Some(4);
lkg.info.lkg_next_cycle = Some(9);
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, TEST_PLAN_DIGEST, 10, 5)
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 10, 5)
.expect("compatible LKG should survive a leader change");
assert_eq!(observed.usage_snapshot_set_states[0].scanner_epoch, Some(4));
assert_eq!(observed.objects_total_size, 50);
@@ -250,11 +257,11 @@ fn usage_aggregate_cost_is_linear_in_set_count() {
cache.info.leader_epoch = 3;
results.push(cache);
}
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
.expect("all set snapshots should aggregate");
assert_eq!(observed.objects_total_count, 32);
let reversed = results.iter().rev().cloned().collect::<Vec<_>>();
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
.expect("reordered set snapshots should aggregate");
assert_eq!(observed.usage_snapshot_set_states, reversed_observed.usage_snapshot_set_states);
}
@@ -276,11 +283,21 @@ fn completed_data_usage_info_publishes_tier_stats_across_sets() {
let mut first_set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let mut tiered = DataUsageEntry::default();
tiered.add_tier_sizes(&warm(100, 2, 1));
tiered.tier_accounting_proof = Some(TierAccountingProof {
physical_total: 100,
physical_known: 100,
..Default::default()
});
first_set.replace("bucket-b", DATA_USAGE_ROOT, tiered);
let mut second_set = completed_root_cache("bucket-b", 2, 20, DataUsageCacheSource::new(1, 0));
let mut tiered = DataUsageEntry::default();
tiered.add_tier_sizes(&warm(50, 1, 1));
tiered.tier_accounting_proof = Some(TierAccountingProof {
physical_total: 50,
physical_known: 50,
..Default::default()
});
second_set.replace("bucket-a", DATA_USAGE_ROOT, tiered);
let (data_usage_info, _) = completed_data_usage_info_for_test(&[first_set, second_set], &all_buckets, false, false)
@@ -299,6 +316,269 @@ fn completed_data_usage_info_publishes_tier_stats_across_sets() {
);
}
#[test]
fn completed_data_usage_info_rejects_logical_proof_mismatch() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_tier_sizes(&HashMap::from([(
"WARM".to_string(),
TierStats {
total_size: 10,
num_versions: 1,
num_objects: 1,
},
)]));
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 10,
logical_known: 9,
physical_total: 10,
physical_known: 10,
..Default::default()
});
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_logical_total_size_mismatch() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.size = 11;
entry.add_tier_sizes(&HashMap::from([(
"WARM".to_string(),
TierStats {
total_size: 10,
num_versions: 1,
num_objects: 1,
},
)]));
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 10,
logical_known: 10,
physical_total: 10,
physical_known: 10,
..Default::default()
});
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_physical_proof_mismatch() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_tier_sizes(&HashMap::from([(
"WARM".to_string(),
TierStats {
total_size: 10,
num_versions: 1,
num_objects: 1,
},
)]));
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 10,
logical_known: 10,
physical_total: 9,
physical_known: 9,
..Default::default()
});
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_unknown_physical_double_accounting() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_tier_sizes(&HashMap::from([(
UNKNOWN_TIER.to_string(),
TierStats {
total_size: 10,
num_versions: 1,
num_objects: 1,
},
)]));
entry.add_unknown_tier_stats(&UnknownTierStats {
unknown_physical_bytes: 9,
..Default::default()
});
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 10,
logical_known: 10,
physical_total: 10,
physical_known: 10,
..Default::default()
});
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_unknown_counter_overflow() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_unknown_tier_stats(&UnknownTierStats {
counter_overflowed: true,
..Default::default()
});
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_accepts_no_tier_standard_empty_map_with_proof() {
let all_buckets = vec!["bucket-a".to_string()];
let set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let (info, _) = completed_data_usage_info_for_test(&[set], &all_buckets, false, false)
.expect("no-tier STANDARD/RRS usage does not require a tier map");
assert!(info.tier_stats.is_none());
}
#[test]
fn completed_data_usage_info_accepts_no_tier_unknown_and_standard_shape() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.size = 13;
entry.add_tier_sizes(&HashMap::from([(
UNKNOWN_TIER.to_string(),
TierStats {
total_size: 3,
num_versions: 1,
num_objects: 1,
},
)]));
entry.add_unknown_tier_stats(&UnknownTierStats {
unknown_bytes: 9,
unknown_physical_bytes: 3,
unknown_objects: 1,
unknown_versions: 1,
..Default::default()
});
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 13,
logical_known: 4,
physical_total: 7,
physical_known: 4,
..Default::default()
});
let (info, _) = completed_data_usage_info_for_test(&[set], &all_buckets, false, false)
.expect("no-tier STANDARD plus UNKNOWN should remain publishable");
assert_eq!(
info.tier_stats.expect("unknown bucket should be retained").tiers[UNKNOWN_TIER].total_size,
3
);
}
#[test]
fn completed_data_usage_info_accepts_unknown_only_with_current_registry_generation() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
set.info.tier_registry_generation = Some(7);
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.size = 13;
entry.add_tier_sizes(&HashMap::from([(
UNKNOWN_TIER.to_string(),
TierStats {
total_size: 3,
num_versions: 1,
num_objects: 1,
},
)]));
entry.add_unknown_tier_stats(&UnknownTierStats {
unknown_bytes: 9,
unknown_physical_bytes: 3,
unknown_objects: 1,
unknown_versions: 1,
..Default::default()
});
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 13,
logical_known: 4,
physical_total: 7,
physical_known: 4,
..Default::default()
});
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some()
);
}
#[test]
fn completed_data_usage_info_rejects_non_registry_tier_in_current_generation() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
set.info.tier_registry_generation = Some(7);
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_tier_sizes(&HashMap::from([
(
"WARM".to_string(),
TierStats {
total_size: 4,
num_versions: 1,
num_objects: 1,
},
),
(
"RETIRED".to_string(),
TierStats {
total_size: 6,
num_versions: 1,
num_objects: 1,
},
),
]));
entry.tier_accounting_proof = Some(TierAccountingProof {
logical_total: 10,
logical_known: 10,
physical_total: 10,
physical_known: 10,
..Default::default()
});
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none()
);
}
#[test]
fn completed_data_usage_info_rejects_legacy_proof_missing_when_tier_accounted() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.add_tier_sizes(&HashMap::from([(
"WARM".to_string(),
TierStats {
total_size: 10,
num_versions: 1,
num_objects: 1,
},
)]));
entry.tier_accounting_proof = None;
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_legacy_proof_missing_for_scalar_usage() {
let all_buckets = vec!["bucket-a".to_string()];
let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry");
entry.tier_accounting_proof = None;
assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_omits_tier_stats_without_tiered_objects() {
let all_buckets = vec!["bucket-a".to_string()];
@@ -310,6 +590,51 @@ fn completed_data_usage_info_omits_tier_stats_without_tiered_objects() {
assert!(data_usage_info.tier_stats.is_none());
}
#[test]
fn completed_data_usage_info_rejects_legacy_and_new_tier_generations_mixed() {
let all_buckets = vec!["bucket-a".to_string()];
let legacy = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0));
let mut current = completed_root_cache("bucket-a", 1, 20, DataUsageCacheSource::new(1, 0));
current.info.tier_registry_generation = Some(42);
assert!(
completed_data_usage_info_for_test(&[legacy, current], &all_buckets, false, false).is_none(),
"legacy and generation-tagged sets must not publish a mixed snapshot"
);
}
#[test]
fn current_cache_root_with_new_tier_generation_resets_old_cache() {
let source = DataUsageCacheSource::new(0, 0);
let mut cache = completed_root_cache("bucket-a", 1, 10, source);
cache.info.tier_registry_generation = Some(1);
let state = current_cache_root_or_prepare_with_generation(
&mut cache,
DATA_USAGE_ROOT,
source,
0,
0,
TEST_PLAN_DIGEST,
DataUsageCacheReuseOptions {
require_source: false,
tier_registry_generation: Some(2),
},
);
assert!(matches!(
state,
DataUsageCacheScanState::Prepared {
outcome: DataUsageCachePrepareOutcome::Reset,
..
}
));
assert!(cache.cache.is_empty(), "old-generation entries must not be reused");
assert_eq!(cache.info.tier_registry_generation, None);
assert_eq!(cache.info.scan_plan_digest, Some(TEST_PLAN_DIGEST));
assert!(!cache.info.snapshot_complete);
}
#[test]
fn completed_data_usage_info_requires_every_set_before_publish() {
let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string(), "bucket-empty".to_string()];
@@ -434,6 +759,13 @@ fn completed_data_usage_info_flattens_nested_bucket_entries() {
replica_size: 2048,
replica_count: 2,
}),
tier_accounting_proof: Some(TierAccountingProof {
logical_total: 2048,
logical_known: 2048,
physical_total: 2048,
physical_known: 2048,
..Default::default()
}),
..Default::default()
};
nested.obj_sizes.add(2048);
@@ -501,7 +833,8 @@ fn completed_data_usage_info_requires_exact_topology_sources() {
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]);
assert!(
completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, true, false, false).is_none()
completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false)
.is_none()
);
}
@@ -511,7 +844,7 @@ fn completed_data_usage_info_rejects_incomplete_bucket_plan() {
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, false, false, false).is_none());
assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none());
}
#[test]
+44 -5
View File
@@ -23,7 +23,7 @@ use crate::storage_api::owner::{
use crate::storage_api::scan::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions, ObjectIO as _};
use crate::{
DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerObjectOptions,
ScannerPutObjReader, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
ScannerPutObjReader, UNKNOWN_TIER, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
init_local_disks_with_instance_ctx, new_disk, path2_bucket_object_with_base_path,
};
use rustfs_filemeta::FileInfo;
@@ -1030,8 +1030,8 @@ fn is_xl_meta_path_accepts_forward_separator() {
fn tier_stats_template_seeds_tiers_and_standard_classes() {
let template = tier_stats_template(&["WARM".to_string(), "COLD".to_string()]);
assert_eq!(template.len(), 4);
for tier in ["WARM", "COLD", storageclass::STANDARD, storageclass::RRS] {
assert_eq!(template.len(), 5);
for tier in ["WARM", "COLD", storageclass::STANDARD, storageclass::RRS, UNKNOWN_TIER] {
assert_eq!(template.get(tier), Some(&TierStats::default()), "missing seed for tier {tier}");
}
}
@@ -1372,7 +1372,7 @@ fn apply_bucket_result_to_cache_updates_bucket_entry() {
);
let update_time = SystemTime::now();
apply_bucket_result_to_cache(
assert!(apply_bucket_result_to_cache(
&mut cache,
DataUsageEntryInfo {
name: "bucket".to_string(),
@@ -1382,12 +1382,51 @@ fn apply_bucket_result_to_cache_updates_bucket_entry() {
objects: 2,
..Default::default()
},
tier_registry_generation: None,
},
update_time,
);
));
assert_eq!(cache.info.last_update, Some(update_time));
let entry = cache.find("bucket").expect("bucket entry should remain present");
assert_eq!(entry.size, 10);
assert_eq!(entry.objects, 2);
}
#[test]
fn apply_bucket_result_to_cache_rejects_a_different_tier_generation() {
let mut cache = DataUsageCache {
info: DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
tier_registry_generation: Some(7),
..Default::default()
},
..Default::default()
};
cache.replace(
"bucket",
DATA_USAGE_ROOT,
DataUsageEntry {
size: 3,
..Default::default()
},
);
let applied = apply_bucket_result_to_cache(
&mut cache,
DataUsageEntryInfo {
name: "bucket".to_string(),
parent: DATA_USAGE_ROOT.to_string(),
entry: DataUsageEntry {
size: 11,
..Default::default()
},
tier_registry_generation: Some(8),
},
SystemTime::now(),
);
assert!(!applied);
assert_eq!(cache.find("bucket").map(|entry| entry.size), Some(3));
assert!(cache.info.last_update.is_none());
}