From 4589148f4803d83a6fa604f0197b65fa69c7c5d5 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 18 Jul 2026 00:09:02 +0800 Subject: [PATCH] fix(admin): serve data usage endpoints from scanner snapshot instead of live listing (#4980) --- crates/ecstore/src/api/mod.rs | 6 +- crates/ecstore/src/data_usage/mod.rs | 73 ++++++++- crates/madmin/src/user.rs | 68 ++++++-- rustfs/src/admin/handlers/account_info.rs | 52 +++--- rustfs/src/admin/handlers/system.rs | 42 +++-- rustfs/src/admin/storage_api.rs | 19 +-- rustfs/src/app/admin_usecase.rs | 43 +---- .../app/data_usage_snapshot_gating_test.rs | 155 ++++++++++++++++++ rustfs/src/app/mod.rs | 2 + rustfs/src/app/storage_api.rs | 31 +--- rustfs/src/storage/storage_api.rs | 12 +- 11 files changed, 370 insertions(+), 133 deletions(-) create mode 100644 rustfs/src/app/data_usage_snapshot_gating_test.rs diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index d0cbfb23c..8c9d26431 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -258,12 +258,14 @@ pub mod config { pub mod data_usage { pub use crate::data_usage::{ - DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, - load_compression_total_from_memory, load_data_usage_from_backend, record_bucket_delete_marker_memory, + DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage, + init_compression_total_memory_from_backend, live_bucket_usage_computations, load_compression_total_from_memory, + load_data_usage_from_backend, load_data_usage_from_backend_cached, record_bucket_delete_marker_memory, record_bucket_object_delete_memory, record_bucket_object_version_write_memory, record_bucket_object_write_memory, record_bucket_object_write_unknown_previous_memory, record_compression_total_memory, refresh_bucket_usage_from_object_layer, refresh_versioned_bucket_usage_from_object_layer, remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, store_compression_total_in_backend, + store_data_usage_in_backend, }; } diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index c27e7a599..a67e69d57 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -41,7 +41,10 @@ use rustfs_utils::path::SLASH_SEPARATOR; use std::{ collections::{HashMap, HashSet, hash_map::Entry}, future::Future, - sync::{Arc, LazyLock, OnceLock}, + sync::{ + Arc, LazyLock, OnceLock, + atomic::{AtomicU64, Ordering}, + }, time::{Duration, SystemTime}, }; use tokio::fs; @@ -76,6 +79,28 @@ static USAGE_MEMORY_CACHE: OnceLock = OnceLock::new(); static USAGE_CACHE_UPDATING: OnceLock = OnceLock::new(); static LIVE_BUCKET_USAGE_CACHE: OnceLock = OnceLock::new(); +/// Cached copy of the last persisted data usage snapshot, served to admin +/// endpoints for up to `DATA_USAGE_CACHE_TTL_SECS` between backend reads. +#[derive(Debug, Clone)] +struct CachedDataUsageSnapshot { + info: DataUsageInfo, + loaded_at: SystemTime, +} + +type DataUsageSnapshotCache = Arc>>; + +static DATA_USAGE_SNAPSHOT_CACHE: OnceLock = OnceLock::new(); + +// Always-on revert detector for rustfs/backlog#1306: one relaxed increment per +// full-bucket version listing is negligible and lets tests prove that admin +// request paths never trigger live listings. +static LIVE_BUCKET_USAGE_COMPUTATIONS: AtomicU64 = AtomicU64::new(0); + +/// Number of live full-bucket usage computations performed by this process. +pub fn live_bucket_usage_computations() -> u64 { + LIVE_BUCKET_USAGE_COMPUTATIONS.load(Ordering::Relaxed) +} + /// Deferred persist thresholds for compression totals: persist after this many /// operations recorded, but no more often than the min interval. const COMPRESSION_PERSIST_BATCH_SIZE: u64 = 100; @@ -115,6 +140,10 @@ fn cache_updating() -> &'static CacheUpdating { USAGE_CACHE_UPDATING.get_or_init(|| Arc::new(RwLock::new(false))) } +fn data_usage_snapshot_cache() -> &'static DataUsageSnapshotCache { + DATA_USAGE_SNAPSHOT_CACHE.get_or_init(|| Arc::new(RwLock::new(None))) +} + fn live_bucket_usage_cache() -> &'static LiveBucketUsageCache { LIVE_BUCKET_USAGE_CACHE.get_or_init(|| { moka::future::Cache::builder() @@ -193,6 +222,12 @@ async fn save_data_usage_in_backend(data_usage_info: DataUsageInfo, store: Arc) -> Result) -> Result { + let ttl = Duration::from_secs(DATA_USAGE_CACHE_TTL_SECS); + + { + let cache = data_usage_snapshot_cache().read().await; + if let Some(cached) = cache.as_ref() + && SystemTime::now().duration_since(cached.loaded_at).unwrap_or_default() < ttl + { + return Ok(cached.info.clone()); + } + } + + // Re-check under the write lock so concurrent expirations trigger a single + // backend read instead of a stampede. + let mut cache = data_usage_snapshot_cache().write().await; + if let Some(cached) = cache.as_ref() + && SystemTime::now().duration_since(cached.loaded_at).unwrap_or_default() < ttl + { + return Ok(cached.info.clone()); + } + + let info = load_data_usage_from_backend(store).await?; + *cache = Some(CachedDataUsageSnapshot { + info: info.clone(), + loaded_at: SystemTime::now(), + }); + Ok(info) +} + /// Aggregate usage information from local disk snapshots. fn merge_snapshot(aggregated: &mut DataUsageInfo, mut snapshot: LocalUsageSnapshot, latest_update: &mut Option) { if let Some(update) = snapshot.last_update @@ -548,6 +618,7 @@ impl BucketUsageAccumulator { type UsageVersionPage = StorageListObjectVersionsInfo; pub async fn compute_bucket_usage(store: Arc, bucket_name: &str) -> Result { + LIVE_BUCKET_USAGE_COMPUTATIONS.fetch_add(1, Ordering::Relaxed); let bucket = bucket_name.to_string(); compute_bucket_usage_with_pages(bucket_name, move |marker, version_marker| { let store = Arc::clone(&store); diff --git a/crates/madmin/src/user.rs b/crates/madmin/src/user.rs index 114dad6e6..fe7de5bc5 100644 --- a/crates/madmin/src/user.rs +++ b/crates/madmin/src/user.rs @@ -406,10 +406,17 @@ pub struct AccountInfo { #[derive(Debug, Serialize, Deserialize, Default)] pub struct BucketAccessInfo { pub name: String, - pub size: u64, - pub objects: u64, - pub object_sizes_histogram: HashMap, - pub object_versions_histogram: HashMap, + // Usage stats are absent (not zero) when no scanner snapshot covers the + // bucket yet, so clients can distinguish "unknown" from "empty" + // (rustfs/backlog#1306). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub size: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub objects: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub object_sizes_histogram: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub object_versions_histogram: Option>, pub details: Option, pub prefix_usage: HashMap, #[serde(rename = "expiration", with = "time::serde::rfc3339::option")] @@ -727,6 +734,43 @@ mod tests { use time::OffsetDateTime; use time::macros::datetime; + /// Wire pin (rustfs/backlog#1306): usage stats without a scanner snapshot + /// must be omitted from the JSON, not serialized as zeros. + #[test] + fn bucket_access_info_omits_absent_usage_stats() { + let info = BucketAccessInfo { + name: "no-snapshot".to_string(), + ..Default::default() + }; + + let value = serde_json::to_value(&info).unwrap(); + let obj = value.as_object().unwrap(); + assert!(!obj.contains_key("size")); + assert!(!obj.contains_key("objects")); + assert!(!obj.contains_key("object_sizes_histogram")); + assert!(!obj.contains_key("object_versions_histogram")); + } + + /// Wire pin (rustfs/backlog#1306): populated usage stats keep their + /// existing snake_case keys and numeric values. + #[test] + fn bucket_access_info_serializes_present_usage_stats() { + let info = BucketAccessInfo { + name: "snapshot".to_string(), + size: Some(1024), + objects: Some(7), + object_sizes_histogram: Some(HashMap::from([("1MiB-10MiB".to_string(), 7)])), + object_versions_histogram: Some(HashMap::from([("SINGLE_VERSION".to_string(), 7)])), + ..Default::default() + }; + + let value = serde_json::to_value(&info).unwrap(); + assert_eq!(value["size"], 1024); + assert_eq!(value["objects"], 7); + assert_eq!(value["object_sizes_histogram"]["1MiB-10MiB"], 7); + assert_eq!(value["object_versions_histogram"]["SINGLE_VERSION"], 7); + } + #[test] fn test_account_status_try_from_invalid() { let result = AccountStatus::try_from("invalid"); @@ -1029,10 +1073,10 @@ mod tests { let bucket_info = BucketAccessInfo { name: "test-bucket".to_string(), - size: 6000000, - objects: 150, - object_sizes_histogram: sizes_histogram, - object_versions_histogram: versions_histogram, + size: Some(6000000), + objects: Some(150), + object_sizes_histogram: Some(sizes_histogram), + object_versions_histogram: Some(versions_histogram), details: Some(BucketDetails { versioning: true, versioning_suspended: false, @@ -1049,10 +1093,10 @@ mod tests { }; assert_eq!(bucket_info.name, "test-bucket"); - assert_eq!(bucket_info.size, 6000000); - assert_eq!(bucket_info.objects, 150); - assert_eq!(bucket_info.object_sizes_histogram.len(), 2); - assert_eq!(bucket_info.object_versions_histogram.len(), 2); + assert_eq!(bucket_info.size, Some(6000000)); + assert_eq!(bucket_info.objects, Some(150)); + assert_eq!(bucket_info.object_sizes_histogram.as_ref().map(HashMap::len), Some(2)); + assert_eq!(bucket_info.object_versions_histogram.as_ref().map(HashMap::len), Some(2)); assert!(bucket_info.details.is_some()); assert_eq!(bucket_info.prefix_usage.len(), 2); assert!(bucket_info.created.is_some()); diff --git a/rustfs/src/admin/handlers/account_info.rs b/rustfs/src/admin/handlers/account_info.rs index d0d4e2a95..829c2ee3f 100644 --- a/rustfs/src/admin/handlers/account_info.rs +++ b/rustfs/src/admin/handlers/account_info.rs @@ -18,10 +18,7 @@ use crate::admin::runtime_sources::{current_action_credentials, object_store_fro use crate::admin::storage_api::bucket::versioning_sys::BucketVersioningSys; use crate::admin::storage_api::contract::admin::StorageAdminApi; use crate::admin::storage_api::contract::bucket::{BucketOperations, BucketOptions}; -use crate::admin::storage_api::data_usage::{ - apply_bucket_usage_memory_overlay, load_data_usage_from_backend, refresh_bucket_usage_from_object_layer, - replace_bucket_usage_memory_from_info, -}; +use crate::admin::storage_api::data_usage::{apply_bucket_usage_memory_overlay, load_data_usage_from_backend_cached}; use crate::admin::storage_api::metadata_sys; use crate::auth::get_condition_values; use crate::server::{ADMIN_PREFIX, RemoteAddr}; @@ -38,7 +35,6 @@ use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error} use serde::Serialize; use std::collections::HashMap; use std::sync::Arc; -use tracing::debug; #[allow(dead_code)] #[derive(Debug, Serialize, Default)] @@ -66,14 +62,16 @@ fn resolve_bucket_access(can_list_bucket: bool, can_get_bucket_location: bool, c } fn apply_usage_to_bucket_access_info(bucket_info: &mut rustfs_madmin::BucketAccessInfo, usage: Option<&BucketUsageInfo>) { + // No snapshot coverage for this bucket: leave the stats absent so clients + // render "unknown" instead of confirmed zeros (rustfs/backlog#1306). let Some(usage) = usage else { return; }; - bucket_info.size = usage.size; - bucket_info.objects = usage.objects_count; - bucket_info.object_sizes_histogram = usage.object_size_histogram.clone(); - bucket_info.object_versions_histogram = usage.object_versions_histogram.clone(); + bucket_info.size = Some(usage.size); + bucket_info.objects = Some(usage.objects_count); + bucket_info.object_sizes_histogram = Some(usage.object_size_histogram.clone()); + bucket_info.object_versions_histogram = Some(usage.object_versions_histogram.clone()); } fn object_lock_config_enabled(config: &ObjectLockConfiguration) -> bool { @@ -255,10 +253,12 @@ impl Operation for AccountInfoHandler { .await .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, e.to_string()))?; - let mut data_usage_info = load_data_usage_from_backend(store.clone()) + // Serve the last persisted scanner snapshot plus the in-memory overlay. + // This request path must never trigger a live full-version listing + // (rustfs/backlog#1306); freshness is owned by the scanner. + let mut data_usage_info = load_data_usage_from_backend_cached(store.clone()) .await .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, e.to_string()))?; - replace_bucket_usage_memory_from_info(&data_usage_info).await; apply_bucket_usage_memory_overlay(&mut data_usage_info).await; for bucket in buckets.iter() { @@ -277,15 +277,6 @@ impl Operation for AccountInfoHandler { access: rustfs_madmin::AccountAccess { read: rd, write: wr }, ..Default::default() }; - // AccountInfo backs Console bucket stats, so prefer object-layer usage over potentially cold scanner snapshots. - if let Err(err) = refresh_bucket_usage_from_object_layer(store.clone(), &mut data_usage_info, &bucket.name).await - { - debug!( - bucket = %bucket.name, - error = %err, - "failed to refresh account info bucket usage from object layer" - ); - } apply_usage_to_bucket_access_info(&mut bucket_info, data_usage_info.buckets_usage.get(&bucket.name)); account_info.buckets.push(bucket_info); } @@ -356,25 +347,28 @@ mod tests { apply_usage_to_bucket_access_info(&mut bucket_info, Some(&usage)); - assert_eq!(bucket_info.size, usage.size); - assert_eq!(bucket_info.objects, usage.objects_count); - assert_eq!(bucket_info.object_sizes_histogram, usage.object_size_histogram); - assert_eq!(bucket_info.object_versions_histogram, usage.object_versions_histogram); + assert_eq!(bucket_info.size, Some(usage.size)); + assert_eq!(bucket_info.objects, Some(usage.objects_count)); + assert_eq!(bucket_info.object_sizes_histogram.as_ref(), Some(&usage.object_size_histogram)); + assert_eq!(bucket_info.object_versions_histogram.as_ref(), Some(&usage.object_versions_histogram)); } + /// Buckets without snapshot coverage keep their stats absent so the wire + /// omits the fields instead of reporting confirmed zeros + /// (rustfs/backlog#1306 contract decision, option A). #[test] - fn accountinfo_bucket_access_info_ignores_missing_usage() { + fn accountinfo_bucket_access_info_leaves_stats_absent_without_usage() { let mut bucket_info = rustfs_madmin::BucketAccessInfo { name: "agent".to_string(), - size: 5, - objects: 2, ..Default::default() }; apply_usage_to_bucket_access_info(&mut bucket_info, None); - assert_eq!(bucket_info.size, 5); - assert_eq!(bucket_info.objects, 2); + assert_eq!(bucket_info.size, None); + assert_eq!(bucket_info.objects, None); + assert_eq!(bucket_info.object_sizes_histogram, None); + assert_eq!(bucket_info.object_versions_histogram, None); } #[test] diff --git a/rustfs/src/admin/handlers/system.rs b/rustfs/src/admin/handlers/system.rs index 43d764845..f1587486d 100644 --- a/rustfs/src/admin/handlers/system.rs +++ b/rustfs/src/admin/handlers/system.rs @@ -795,6 +795,16 @@ impl Operation for RuntimeCapabilitiesHandler { } } +/// Authorization gate for GET datausageinfo: any-of the dedicated admin action +/// OR the bucket listing action. Pinned by a unit test so the gate cannot +/// silently narrow or widen (rustfs/backlog#1306). +fn data_usage_info_gate_actions() -> Vec { + vec![ + Action::AdminAction(AdminAction::DataUsageInfoAdminAction), + Action::S3Action(S3Action::ListBucketAction), + ] +} + #[async_trait::async_trait] impl Operation for DataUsageInfoHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { @@ -807,18 +817,7 @@ impl Operation for DataUsageInfoHandler { check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; let remote_addr = req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)); - validate_admin_request( - &req.headers, - &cred, - owner, - false, - vec![ - Action::AdminAction(AdminAction::DataUsageInfoAdminAction), - Action::S3Action(S3Action::ListBucketAction), - ], - remote_addr, - ) - .await?; + validate_admin_request(&req.headers, &cred, owner, false, data_usage_info_gate_actions(), remote_addr).await?; let usecase = default_admin_usecase(); let info = usecase.execute_query_data_usage_info().await.map_err(S3Error::from)?; @@ -840,7 +839,8 @@ impl Operation for DataUsageInfoHandler { mod tests { use super::{ OBSERVABILITY_SUMMARY_RESOLVED, ServerInfoResponse, TOPOLOGY_SNAPSHOT_NOT_AVAILABLE, TOPOLOGY_SUMMARY_RESOLVED, - build_runtime_capabilities_response, build_runtime_capabilities_summary, system_admin_discovery, + build_runtime_capabilities_response, build_runtime_capabilities_summary, data_usage_info_gate_actions, + system_admin_discovery, }; use crate::admin::runtime_sources::DefaultAdminUsecase; use crate::admin::storage_api::cluster::{ @@ -849,6 +849,22 @@ mod tests { }; use rustfs_concurrency::WorkloadClass; use rustfs_madmin::{InfoMessage, StorageInfo}; + use rustfs_policy::policy::action::{Action, AdminAction, S3Action}; + + /// Authz regression pin (rustfs/backlog#1306): datausageinfo stays an + /// any-of gate over exactly DataUsageInfoAdminAction OR ListBucketAction. + /// The OR semantics of the multi-action loop are covered by the + /// `evaluate_admin_actions` tests in `crate::admin::auth`. + #[test] + fn data_usage_info_gate_keeps_dual_action_or_semantics() { + assert_eq!( + data_usage_info_gate_actions(), + vec![ + Action::AdminAction(AdminAction::DataUsageInfoAdminAction), + Action::S3Action(S3Action::ListBucketAction), + ] + ); + } #[tokio::test] async fn runtime_capabilities_response_reports_missing_topology_before_storage_init() { diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index 01d4ae059..72fbd6a8c 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -456,27 +456,16 @@ pub(crate) mod data_usage { crate::storage::storage_api::ecstore_data_usage::apply_bucket_usage_memory_overlay(data_usage_info).await; } - pub(crate) async fn refresh_bucket_usage_from_object_layer( - store: Arc, - data_usage_info: &mut rustfs_data_usage::DataUsageInfo, - bucket_name: &str, - ) -> Result { - crate::storage::storage_api::ecstore_data_usage::refresh_bucket_usage_from_object_layer( - store, - data_usage_info, - bucket_name, - ) - .await - } - pub(crate) async fn load_data_usage_from_backend( store: Arc, ) -> Result { crate::storage::storage_api::ecstore_data_usage::load_data_usage_from_backend(store).await } - pub(crate) async fn replace_bucket_usage_memory_from_info(data_usage_info: &rustfs_data_usage::DataUsageInfo) { - crate::storage::storage_api::ecstore_data_usage::replace_bucket_usage_memory_from_info(data_usage_info).await; + pub(crate) async fn load_data_usage_from_backend_cached( + store: Arc, + ) -> Result { + crate::storage::storage_api::ecstore_data_usage::load_data_usage_from_backend_cached(store).await } } diff --git a/rustfs/src/app/admin_usecase.rs b/rustfs/src/app/admin_usecase.rs index 001cc5240..678aea387 100644 --- a/rustfs/src/app/admin_usecase.rs +++ b/rustfs/src/app/admin_usecase.rs @@ -19,11 +19,7 @@ use super::storage_api::admin_usecase::capacity::{ PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity, get_total_usable_capacity_free, }; use super::storage_api::admin_usecase::contract::StorageAdminApi; -use super::storage_api::admin_usecase::contract::bucket::{BucketOperations, BucketOptions}; -use super::storage_api::admin_usecase::data_usage::{ - apply_bucket_usage_memory_overlay, load_data_usage_from_backend, refresh_bucket_usage_from_object_layer, - replace_bucket_usage_memory_from_info, -}; +use super::storage_api::admin_usecase::data_usage::{apply_bucket_usage_memory_overlay, load_data_usage_from_backend_cached}; use super::storage_api::admin_usecase::{ECStore, EndpointServerPools}; use crate::app::runtime_sources::{ AppContext, current_app_context, current_endpoints_handle, current_object_store_handle_for_context, @@ -245,13 +241,18 @@ impl DefaultAdminUsecase { return Err(Self::app_error(S3ErrorCode::InternalError, "Not init")); }; - let mut info = load_data_usage_from_backend(store.clone()).await.map_err(|e| { + Self::query_data_usage_info_with_store(store).await + } + + /// Serve the last persisted scanner snapshot plus the in-memory overlay. + /// This request path must never trigger a live full-version listing + /// (rustfs/backlog#1306); freshness is owned by the scanner. + pub(crate) async fn query_data_usage_info_with_store(store: Arc) -> AdminUsecaseResult { + let mut info = load_data_usage_from_backend_cached(store.clone()).await.map_err(|e| { error!("load_data_usage_from_backend failed {:?}", e); Self::app_error(S3ErrorCode::InternalError, "load_data_usage_from_backend failed") })?; - replace_bucket_usage_memory_from_info(&info).await; apply_bucket_usage_memory_overlay(&mut info).await; - Self::refresh_live_bucket_usage_for_data_usage_info(store.clone(), &mut info).await; let storage_info = StorageAdminApi::storage_info(store.as_ref()).await; @@ -321,32 +322,6 @@ impl DefaultAdminUsecase { Ok(info) } - async fn refresh_live_bucket_usage_for_data_usage_info(store: Arc, data_usage_info: &mut DataUsageInfo) { - let buckets = match store - .list_bucket(&BucketOptions { - no_metadata: true, - ..Default::default() - }) - .await - { - Ok(buckets) => buckets, - Err(err) => { - debug!(error = %err, "failed to list buckets while refreshing data usage info"); - return; - } - }; - - for bucket in buckets { - if let Err(err) = refresh_bucket_usage_from_object_layer(store.clone(), data_usage_info, &bucket.name).await { - debug!( - bucket = %bucket.name, - error = %err, - "failed to refresh data usage info bucket usage from object layer" - ); - } - } - } - pub async fn execute_list_pool_statuses(&self) -> AdminUsecaseResult> { let Some(store) = self.object_store() else { return Err(Self::app_error(S3ErrorCode::InternalError, "Not init")); diff --git a/rustfs/src/app/data_usage_snapshot_gating_test.rs b/rustfs/src/app/data_usage_snapshot_gating_test.rs new file mode 100644 index 000000000..d2ab87ff1 --- /dev/null +++ b/rustfs/src/app/data_usage_snapshot_gating_test.rs @@ -0,0 +1,155 @@ +// 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. + +//! Revert detector for rustfs/backlog#1306: the admin data usage endpoint must +//! serve the persisted scanner snapshot plus the in-memory overlay and never +//! trigger a live full-version listing on the request path. +//! +//! The observable is the always-on `live_bucket_usage_computations()` counter +//! incremented by `compute_bucket_usage` (the only entry point into request +//! path full version listings). The test first proves the counter moves when a +//! live computation *does* run (positive control, so a broken counter cannot +//! make the guard pass vacuously), then drives the endpoint use case +//! end-to-end against a pre-seeded snapshot and asserts the counter stays put +//! while the response carries the seeded numbers. + +use super::gating_test_env::shared_gating_ecstore; +use super::storage_api::test::StoragePutObjReader as PutObjReader; +use super::storage_api::test::contract::bucket::{BucketOperations, MakeBucketOptions}; +use super::storage_api::test::contract::object::ObjectIO as _; +use super::storage_api::test::data_usage::{ + compute_bucket_usage, live_bucket_usage_computations, record_bucket_object_write_memory, store_data_usage_in_backend, +}; +use crate::app::admin_usecase::DefaultAdminUsecase; +use rustfs_data_usage::{BucketUsageInfo, DataUsageInfo}; +use serial_test::serial; +use std::collections::HashMap; +use std::time::SystemTime; +use uuid::Uuid; + +const SEEDED_BUCKET_SIZE: u64 = 123_456; +const SEEDED_BUCKET_OBJECTS: u64 = 42; + +fn seeded_data_usage_info(bucket: &str, last_update: SystemTime) -> DataUsageInfo { + let usage = BucketUsageInfo { + size: SEEDED_BUCKET_SIZE, + objects_count: SEEDED_BUCKET_OBJECTS, + versions_count: SEEDED_BUCKET_OBJECTS, + ..Default::default() + }; + + let mut info = DataUsageInfo { + last_update: Some(last_update), + buckets_count: 1, + objects_total_count: SEEDED_BUCKET_OBJECTS, + objects_total_size: SEEDED_BUCKET_SIZE, + ..Default::default() + }; + info.buckets_usage = HashMap::from([(bucket.to_string(), usage)]); + info.bucket_sizes = HashMap::from([(bucket.to_string(), SEEDED_BUCKET_SIZE)]); + info +} + +#[tokio::test] +#[serial] +async fn data_usage_endpoint_serves_snapshot_without_live_listing() { + let ecstore = shared_gating_ecstore().await; + let live_bucket = format!("usage-live-{}", Uuid::new_v4()); + let seeded_bucket = format!("usage-seeded-{}", Uuid::new_v4()); + let overlay_bucket = format!("usage-overlay-{}", Uuid::new_v4()); + + // Positive control: a direct live computation must move the revert + // detector counter, otherwise the "no increment" assertion below would be + // vacuously true with a broken counter. + ecstore + .make_bucket(&live_bucket, &MakeBucketOptions::default()) + .await + .expect("create live bucket"); + for object in ["control-a.bin", "control-b.bin"] { + let mut reader = PutObjReader::from_vec(b"live payload".to_vec()); + ecstore + .put_object(&live_bucket, object, &mut reader, &Default::default()) + .await + .expect("put control object"); + } + + let before_control = live_bucket_usage_computations(); + let control_usage = compute_bucket_usage(ecstore.clone(), &live_bucket) + .await + .expect("live computation over the control bucket"); + assert_eq!(control_usage.objects_count, 2, "control bucket must be fully listed"); + assert!( + live_bucket_usage_computations() > before_control, + "positive control: compute_bucket_usage must increment the live-listing counter" + ); + + // Pre-seed a scanner snapshot for a bucket the endpoint has to serve + // verbatim, then record an in-memory overlay write for another bucket. + ecstore + .make_bucket(&seeded_bucket, &MakeBucketOptions::default()) + .await + .expect("create seeded bucket"); + let seeded_at = SystemTime::now(); + store_data_usage_in_backend(seeded_data_usage_info(&seeded_bucket, seeded_at), ecstore.clone()) + .await + .expect("persist seeded data usage snapshot"); + + record_bucket_object_write_memory(&overlay_bucket, None, 512).await; + + let before_endpoint = live_bucket_usage_computations(); + let info = DefaultAdminUsecase::query_data_usage_info_with_store(ecstore.clone()) + .await + .expect("query data usage info"); + assert_eq!( + live_bucket_usage_computations(), + before_endpoint, + "revert detector: the data usage endpoint must not run live full-version listings" + ); + + // The endpoint must serve the seeded snapshot numbers, not recomputed ones. + assert_eq!(info.last_update, Some(seeded_at), "endpoint must report the snapshot timestamp"); + let seeded_usage = info + .buckets_usage + .get(&seeded_bucket) + .expect("seeded bucket must come from the snapshot"); + assert_eq!(seeded_usage.size, SEEDED_BUCKET_SIZE); + assert_eq!(seeded_usage.objects_count, SEEDED_BUCKET_OBJECTS); + + // The in-memory overlay stays applied on top of the snapshot. + let overlay_usage = info + .buckets_usage + .get(&overlay_bucket) + .expect("overlay bucket must come from the memory overlay"); + assert_eq!(overlay_usage.size, 512); +} + +/// Wire pin for the no-snapshot response shape (rustfs/backlog#1306): a +/// default `DataUsageInfo` must keep serializing `last_update` as `null` with +/// empty bucket maps, so "no snapshot yet" stays distinguishable from real +/// stats and a future `skip_serializing_if`/`now()` fallback trips this test. +#[test] +fn data_usage_info_default_serializes_null_last_update_and_empty_buckets() { + let value = serde_json::to_value(DataUsageInfo::default()).expect("serialize default DataUsageInfo"); + + assert!(value["last_update"].is_null(), "last_update must serialize as null: {value}"); + assert_eq!(value["buckets_count"], 0); + assert!( + value["buckets_usage"].as_object().is_some_and(|map| map.is_empty()), + "buckets_usage must serialize as an empty map: {value}" + ); + assert!( + value["bucket_sizes"].as_object().is_some_and(|map| map.is_empty()), + "bucket_sizes must serialize as an empty map: {value}" + ); +} diff --git a/rustfs/src/app/mod.rs b/rustfs/src/app/mod.rs index 0c0c5cffe..c0ee9f0fe 100644 --- a/rustfs/src/app/mod.rs +++ b/rustfs/src/app/mod.rs @@ -29,6 +29,8 @@ pub(crate) mod storage_api; #[cfg(test)] mod capacity_dirty_scope_test; #[cfg(test)] +mod data_usage_snapshot_gating_test; +#[cfg(test)] mod delete_objects_stat_gating_test; #[cfg(test)] mod gating_test_env; diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index 22e581340..d575f79f9 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -51,28 +51,17 @@ pub(crate) mod data_usage { crate::storage::storage_api::ecstore_data_usage::apply_bucket_usage_memory_overlay(data_usage_info).await; } - pub(crate) async fn load_data_usage_from_backend( + pub(crate) async fn load_data_usage_from_backend_cached( store: Arc, ) -> Result { - crate::storage::storage_api::ecstore_data_usage::load_data_usage_from_backend(store).await + crate::storage::storage_api::ecstore_data_usage::load_data_usage_from_backend_cached(store).await } - pub(crate) async fn refresh_bucket_usage_from_object_layer( - store: Arc, - data_usage_info: &mut rustfs_data_usage::DataUsageInfo, - bucket_name: &str, - ) -> Result { - crate::storage::storage_api::ecstore_data_usage::refresh_bucket_usage_from_object_layer( - store, - data_usage_info, - bucket_name, - ) - .await - } - - pub(crate) async fn replace_bucket_usage_memory_from_info(data_usage_info: &rustfs_data_usage::DataUsageInfo) { - crate::storage::storage_api::ecstore_data_usage::replace_bucket_usage_memory_from_info(data_usage_info).await; - } + // Test-only observables for the rustfs/backlog#1306 revert detector. + #[cfg(test)] + pub(crate) use crate::storage::storage_api::ecstore_data_usage::{ + compute_bucket_usage, live_bucket_usage_computations, store_data_usage_in_backend, + }; pub(crate) async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64, removed_current_object: bool) { crate::storage::storage_api::ecstore_data_usage::record_bucket_object_delete_memory( @@ -936,10 +925,6 @@ pub(crate) mod s3_api { pub(crate) mod admin_usecase { pub(crate) mod contract { - pub(crate) mod bucket { - pub(crate) use super::super::super::storage_contracts::{BucketOperations, BucketOptions}; - } - pub(crate) use super::super::storage_contracts::StorageAdminApi; } @@ -1073,7 +1058,7 @@ pub(crate) mod test { } } - pub(crate) use super::{bucket, ecfs, object_utils, runtime}; + pub(crate) use super::{bucket, data_usage, ecfs, object_utils, runtime}; pub(crate) use crate::storage::storage_api::{ ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader, }; diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index d1afb3b68..ea3f71bbb 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -383,10 +383,14 @@ pub(crate) mod ecstore_config { pub(crate) mod ecstore_data_usage { pub(crate) use rustfs_ecstore::api::data_usage::{ apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, load_data_usage_from_backend, - record_bucket_delete_marker_memory, record_bucket_object_delete_memory, record_bucket_object_version_write_memory, - record_bucket_object_write_memory, record_bucket_object_write_unknown_previous_memory, - refresh_bucket_usage_from_object_layer, remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, - store_compression_total_in_backend, + load_data_usage_from_backend_cached, record_bucket_delete_marker_memory, record_bucket_object_delete_memory, + record_bucket_object_version_write_memory, record_bucket_object_write_memory, + record_bucket_object_write_unknown_previous_memory, remove_bucket_usage_from_backend, store_compression_total_in_backend, + }; + // Test-only observables for the rustfs/backlog#1306 revert detector. + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::data_usage::{ + compute_bucket_usage, live_bucket_usage_computations, store_data_usage_in_backend, }; }