diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index d461184c4..2edd136de 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -261,7 +261,8 @@ pub mod 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, record_bucket_object_delete_memory, record_bucket_object_version_write_memory, record_bucket_object_write_memory, - record_compression_total_memory, refresh_versioned_bucket_usage_from_object_layer, remove_bucket_usage_from_backend, + 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, }; } diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index f35b4984c..defad2f3f 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -33,7 +33,7 @@ use crate::{ pub use local_snapshot::{LocalUsageSnapshot, read_snapshot as read_local_snapshot, snapshot_path}; use rustfs_data_usage::{ BucketTargetUsageInfo, BucketUsageInfo, CompressionTotalInfo, DataUsageCache, DataUsageEntry, DataUsageInfo, DiskUsageStatus, - SizeSummary, + SizeHistogram, SizeSummary, VersionsHistogram, }; use rustfs_io_metrics::record_system_path_failure; use rustfs_utils::path::SLASH_SEPARATOR; @@ -446,9 +446,11 @@ pub async fn compute_bucket_usage(store: Arc, bucket_name: &str) -> Res let mut marker: Option = None; let mut version_marker: Option = None; let mut object_names: HashSet = HashSet::new(); + let mut object_versions: HashMap = HashMap::new(); let mut versions_count: u64 = 0; let mut total_size: u64 = 0; let mut delete_markers: u64 = 0; + let mut size_histogram = SizeHistogram::default(); loop { let result = store @@ -475,6 +477,8 @@ pub async fn compute_bucket_usage(store: Arc, bucket_name: &str) -> Res let object_size = object.size.max(0) as u64; object_names.insert(object.name.clone()); + *object_versions.entry(object.name.clone()).or_insert(0) += 1; + size_histogram.add(object_size); total_size = total_size.saturating_add(object_size); versions_count = versions_count.saturating_add(1); } @@ -495,18 +499,39 @@ pub async fn compute_bucket_usage(store: Arc, bucket_name: &str) -> Res } let objects_count = object_names.len() as u64; + let mut versions_histogram = VersionsHistogram::default(); + for version_count in object_versions.values() { + versions_histogram.add(*version_count); + } let usage = BucketUsageInfo { size: total_size, objects_count, versions_count, delete_markers_count: delete_markers, + object_size_histogram: size_histogram.to_map(), + object_versions_histogram: versions_histogram.to_map(), ..Default::default() }; Ok(usage) } +pub async fn refresh_bucket_usage_from_object_layer( + store: Arc, + data_usage_info: &mut DataUsageInfo, + bucket: &str, +) -> Result { + let refresh_started_at = SystemTime::now(); + let usage = compute_bucket_usage(store, bucket).await?; + replace_bucket_usage_memory_from_authoritative(bucket, usage.clone(), refresh_started_at).await; + data_usage_info.bucket_sizes.insert(bucket.to_string(), usage.size); + data_usage_info.buckets_usage.insert(bucket.to_string(), usage.clone()); + set_buckets_count_from_usage(data_usage_info); + data_usage_info.calculate_totals(); + Ok(usage) +} + pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc, data_usage_info: &mut DataUsageInfo) { let listed_bucket_names = match store .list_bucket(&BucketOptions { @@ -533,22 +558,14 @@ pub async fn refresh_versioned_bucket_usage_from_object_layer(store: Arc usage, - Err(err) => { - debug!( - bucket = %bucket, - error = %err, - "failed to refresh versioned bucket usage from object layer" - ); - continue; - } - }; - - replace_bucket_usage_memory_from_authoritative(&bucket, usage.clone(), refresh_started_at).await; - data_usage_info.bucket_sizes.insert(bucket.clone(), usage.size); - data_usage_info.buckets_usage.insert(bucket, usage); + if let Err(err) = refresh_bucket_usage_from_object_layer(store.clone(), data_usage_info, &bucket).await { + debug!( + bucket = %bucket, + error = %err, + "failed to refresh versioned bucket usage from object layer" + ); + continue; + } changed = true; } diff --git a/rustfs/src/admin/handlers/account_info.rs b/rustfs/src/admin/handlers/account_info.rs index 24957c652..d3abbb7b8 100644 --- a/rustfs/src/admin/handlers/account_info.rs +++ b/rustfs/src/admin/handlers/account_info.rs @@ -18,11 +18,16 @@ use crate::admin::runtime_sources::{current_action_credentials, current_object_s 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, + refresh_versioned_bucket_usage_from_object_layer, replace_bucket_usage_memory_from_info, +}; use crate::auth::get_condition_values; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use http::{HeaderMap, HeaderValue}; use hyper::{Method, StatusCode}; use matchit::Params; +use rustfs_data_usage::BucketUsageInfo; use rustfs_policy::policy::BucketPolicy; use rustfs_policy::policy::default::DEFAULT_POLICIES; use rustfs_policy::policy::{Args, action::Action, action::S3Action}; @@ -31,6 +36,7 @@ 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)] @@ -57,6 +63,17 @@ fn resolve_bucket_access(can_list_bucket: bool, can_get_bucket_location: bool, c (can_list_bucket || can_get_bucket_location, can_put_object) } +fn apply_usage_to_bucket_access_info(bucket_info: &mut rustfs_madmin::BucketAccessInfo, usage: Option<&BucketUsageInfo>) { + 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(); +} + #[async_trait::async_trait] impl Operation for AccountInfoHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { @@ -201,12 +218,19 @@ 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()) + .await + .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, e.to_string()))?; + refresh_versioned_bucket_usage_from_object_layer(store.clone(), &mut data_usage_info).await; + 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() { let (rd, wr) = is_allow(bucket.name.clone()).await; if rd || wr { // TODO: BucketQuotaSys // TODO: other attributes - account_info.buckets.push(rustfs_madmin::BucketAccessInfo { + let mut bucket_info = rustfs_madmin::BucketAccessInfo { name: bucket.name.clone(), details: Some(rustfs_madmin::BucketDetails { versioning: BucketVersioningSys::enabled(bucket.name.as_str()).await, @@ -216,7 +240,18 @@ impl Operation for AccountInfoHandler { created: bucket.created, 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); } } @@ -267,4 +302,41 @@ mod tests { assert_eq!(resolve_bucket_access(false, true, false), (true, false)); assert_eq!(resolve_bucket_access(false, false, true), (false, true)); } + + #[test] + fn accountinfo_bucket_access_info_uses_data_usage_stats() { + let mut bucket_info = rustfs_madmin::BucketAccessInfo { + name: "agent".to_string(), + ..Default::default() + }; + let usage = BucketUsageInfo { + size: 222 * 1024 * 1024, + objects_count: 109, + object_size_histogram: HashMap::from([("1MiB-10MiB".to_string(), 109)]), + object_versions_histogram: HashMap::from([("SINGLE_VERSION".to_string(), 109)]), + ..Default::default() + }; + + 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); + } + + #[test] + fn accountinfo_bucket_access_info_ignores_missing_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); + } } diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index 5a500d8cd..e30fa6219 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -453,11 +453,40 @@ pub(crate) static ERR_TIER_NOT_FOUND: AdminErrorRef = AdminErrorRef(|| &ecstore_ pub(crate) mod data_usage { use std::sync::Arc; + pub(crate) async fn apply_bucket_usage_memory_overlay(data_usage_info: &mut rustfs_data_usage::DataUsageInfo) { + 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 refresh_versioned_bucket_usage_from_object_layer( + store: Arc, + data_usage_info: &mut rustfs_data_usage::DataUsageInfo, + ) { + crate::storage::storage_api::ecstore_data_usage::refresh_versioned_bucket_usage_from_object_layer(store, data_usage_info) + .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) mod access { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 80982651b..0e8a9211f 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -360,7 +360,8 @@ 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, refresh_versioned_bucket_usage_from_object_layer, remove_bucket_usage_from_backend, + record_bucket_object_write_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, }; }