From 1e303e5be0519d783ee5099fe44deaa288b7dfbc Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Sun, 28 Jun 2026 07:58:07 +0800 Subject: [PATCH] fix(data-usage): refresh versioned usage state (#3969) fix(data-usage): refresh versioned usage from authoritative state Co-authored-by: Henry Guo Co-authored-by: houseme --- crates/ecstore/src/data_usage/mod.rs | 145 ++++++++++++++++++++++++++- 1 file changed, 143 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index 9d2176b50..73ab206cc 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -14,7 +14,11 @@ pub mod local_snapshot; -use crate::storage_api_contracts::{list::ListOperations as _, object::ObjectIO as _}; +use crate::storage_api_contracts::{ + bucket::{BucketOperations as _, BucketOptions}, + list::ListOperations as _, + object::ObjectIO as _, +}; use crate::{ bucket::{metadata_sys::get_replication_config, versioning::VersioningApi as _, versioning_sys::BucketVersioningSys}, config::com::read_config, @@ -463,7 +467,20 @@ pub async fn compute_bucket_usage(store: Arc, bucket_name: &str) -> Res } pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc, data_usage_info: &mut DataUsageInfo) { - let buckets = data_usage_info.buckets_usage.keys().cloned().collect::>(); + let listed_bucket_names = match store + .list_bucket(&BucketOptions { + no_metadata: true, + ..Default::default() + }) + .await + { + Ok(buckets) => buckets.into_iter().map(|bucket| bucket.name).collect::>(), + Err(err) => { + debug!(error = %err, "failed to list buckets while refreshing versioned bucket usage"); + Vec::new() + } + }; + let buckets = bucket_names_for_versioned_refresh(data_usage_info, listed_bucket_names); let mut changed = false; for bucket in buckets { @@ -475,6 +492,7 @@ pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc usage, Err(err) => { @@ -487,6 +505,7 @@ pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc, +) -> Vec { + let mut buckets = data_usage_info.buckets_usage.keys().cloned().collect::>(); + buckets.extend(listed_bucket_names.into_iter().filter(|bucket| !bucket.is_empty())); + + let mut buckets = buckets.into_iter().collect::>(); + buckets.sort(); + buckets +} + async fn ensure_bucket_usage_cached(bucket: &str) { let cache = memory_cache().read().await; if cache.contains_key(bucket) { @@ -540,6 +571,17 @@ fn bucket_usage_counts_match(left: &BucketUsageInfo, right: &BucketUsageInfo) -> && left.delete_markers_count == right.delete_markers_count } +async fn replace_bucket_usage_memory_from_authoritative(bucket: &str, usage: BucketUsageInfo, refresh_started_at: SystemTime) { + let mut cache = memory_cache().write().await; + if let Some(existing) = cache.get(bucket) + && existing.usage_updated_at > refresh_started_at + { + return; + } + + cache.insert(bucket.to_string(), cached_bucket_usage_from_backend(usage, refresh_started_at)); +} + /// Fast in-memory update for immediate quota and admin usage consistency. pub async fn record_bucket_object_write_memory(bucket: &str, previous_current_size: Option, new_size: u64) { record_bucket_object_write_memory_inner(bucket, previous_current_size, new_size, false).await; @@ -1263,6 +1305,105 @@ mod tests { ); } + #[test] + fn versioned_refresh_bucket_names_include_live_buckets_without_usage_snapshot() { + let mut persisted = DataUsageInfo::default(); + persisted.buckets_usage.insert( + "bucket-a".to_string(), + BucketUsageInfo { + objects_count: 1, + versions_count: 1, + size: 10, + ..Default::default() + }, + ); + + let buckets = bucket_names_for_versioned_refresh(&persisted, vec!["bucket-b".to_string(), "bucket-a".to_string()]); + + assert_eq!(buckets, vec!["bucket-a".to_string(), "bucket-b".to_string()]); + } + + #[tokio::test] + #[serial] + async fn authoritative_versioned_refresh_replaces_stale_dirty_memory() { + clear_usage_memory_cache_for_test().await; + + let old_persisted = data_usage_info_for_test("bucket-a", 0, 0, SystemTime::now() - Duration::from_secs(10)); + replace_bucket_usage_memory_from_info(&old_persisted).await; + record_bucket_object_write_memory("bucket-a", None, 15).await; + + let authoritative = BucketUsageInfo { + objects_count: 1, + versions_count: 1, + delete_markers_count: 1, + size: 10, + ..Default::default() + }; + replace_bucket_usage_memory_from_authoritative("bucket-a", authoritative.clone(), SystemTime::now()).await; + + let mut response = old_persisted.clone(); + response.buckets_usage.insert("bucket-a".to_string(), authoritative); + response.bucket_sizes.insert("bucket-a".to_string(), 10); + response.calculate_totals(); + + apply_bucket_usage_memory_overlay(&mut response).await; + + assert_eq!(response.objects_total_count, 1); + assert_eq!(response.versions_total_count, 1); + assert_eq!(response.delete_markers_total_count, 1); + assert_eq!(response.objects_total_size, 10); + assert_eq!( + response.buckets_usage.get("bucket-a").map(|usage| ( + usage.objects_count, + usage.versions_count, + usage.delete_markers_count, + usage.size + )), + Some((1, 1, 1, 10)) + ); + } + + #[tokio::test] + #[serial] + async fn authoritative_versioned_refresh_preserves_newer_dirty_memory() { + clear_usage_memory_cache_for_test().await; + + let old_persisted = data_usage_info_for_test("bucket-a", 0, 0, SystemTime::now() - Duration::from_secs(10)); + replace_bucket_usage_memory_from_info(&old_persisted).await; + let refresh_started = SystemTime::now() - Duration::from_secs(1); + record_bucket_object_write_memory("bucket-a", None, 15).await; + + replace_bucket_usage_memory_from_authoritative( + "bucket-a", + BucketUsageInfo { + objects_count: 1, + versions_count: 1, + delete_markers_count: 1, + size: 10, + ..Default::default() + }, + refresh_started, + ) + .await; + + let mut response = old_persisted.clone(); + apply_bucket_usage_memory_overlay(&mut response).await; + + assert_eq!(response.objects_total_count, 1); + assert_eq!(response.versions_total_count, 1); + assert_eq!(response.delete_markers_total_count, 0); + assert_eq!(response.objects_total_size, 15); + assert_eq!( + response.buckets_usage.get("bucket-a").map(|usage| ( + usage.objects_count, + usage.versions_count, + usage.delete_markers_count, + usage.size + )), + Some((1, 1, 0, 15)) + ); + } + #[tokio::test] #[serial] async fn scanner_sync_preserves_dirty_delete_marker_with_later_snapshot() {