// 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. use std::{collections::HashMap, sync::Arc}; use crate::{bucket::metadata_sys::get_replication_config, config::com::read_config, store::ECStore}; use rustfs_common::data_usage::{BucketTargetUsageInfo, DataUsageCache, DataUsageEntry, DataUsageInfo, SizeSummary}; use rustfs_utils::path::SLASH_SEPARATOR; use tracing::{error, warn}; use crate::error::Error; // Data usage storage constants pub const DATA_USAGE_ROOT: &str = SLASH_SEPARATOR; const DATA_USAGE_OBJ_NAME: &str = ".usage.json"; const DATA_USAGE_BLOOM_NAME: &str = ".bloomcycle.bin"; pub const DATA_USAGE_CACHE_NAME: &str = ".usage-cache.bin"; // Data usage storage paths lazy_static::lazy_static! { pub static ref DATA_USAGE_BUCKET: String = format!("{}{}{}", crate::disk::RUSTFS_META_BUCKET, SLASH_SEPARATOR, crate::disk::BUCKET_META_PREFIX ); pub static ref DATA_USAGE_OBJ_NAME_PATH: String = format!("{}{}{}", crate::disk::BUCKET_META_PREFIX, SLASH_SEPARATOR, DATA_USAGE_OBJ_NAME ); pub static ref DATA_USAGE_BLOOM_NAME_PATH: String = format!("{}{}{}", crate::disk::BUCKET_META_PREFIX, SLASH_SEPARATOR, DATA_USAGE_BLOOM_NAME ); } /// Store data usage info to backend storage pub async fn store_data_usage_in_backend(data_usage_info: DataUsageInfo, store: Arc) -> Result<(), Error> { let data = serde_json::to_vec(&data_usage_info).map_err(|e| Error::other(format!("Failed to serialize data usage info: {e}")))?; // Save to backend using the same mechanism as original code crate::config::com::save_config(store, &DATA_USAGE_OBJ_NAME_PATH, data) .await .map_err(Error::other)?; Ok(()) } /// Load data usage info from backend storage pub async fn load_data_usage_from_backend(store: Arc) -> Result { let buf: Vec = match read_config(store, &DATA_USAGE_OBJ_NAME_PATH).await { Ok(data) => data, Err(e) => { error!("Failed to read data usage info from backend: {}", e); if e == crate::error::Error::ConfigNotFound { return Ok(DataUsageInfo::default()); } return Err(Error::other(e)); } }; let mut data_usage_info: DataUsageInfo = serde_json::from_slice(&buf).map_err(|e| Error::other(format!("Failed to deserialize data usage info: {e}")))?; warn!("Loaded data usage info from backend {:?}", &data_usage_info); // Handle backward compatibility like original code if data_usage_info.buckets_usage.is_empty() { data_usage_info.buckets_usage = data_usage_info .bucket_sizes .iter() .map(|(bucket, &size)| { ( bucket.clone(), rustfs_common::data_usage::BucketUsageInfo { size, ..Default::default() }, ) }) .collect(); } if data_usage_info.bucket_sizes.is_empty() { data_usage_info.bucket_sizes = data_usage_info .buckets_usage .iter() .map(|(bucket, bui)| (bucket.clone(), bui.size)) .collect(); } for (bucket, bui) in &data_usage_info.buckets_usage { if bui.replicated_size_v1 > 0 || bui.replication_failed_count_v1 > 0 || bui.replication_failed_size_v1 > 0 || bui.replication_pending_count_v1 > 0 { if let Ok((cfg, _)) = get_replication_config(bucket).await { if !cfg.role.is_empty() { data_usage_info.replication_info.insert( cfg.role.clone(), BucketTargetUsageInfo { replication_failed_size: bui.replication_failed_size_v1, replication_failed_count: bui.replication_failed_count_v1, replicated_size: bui.replicated_size_v1, replication_pending_count: bui.replication_pending_count_v1, replication_pending_size: bui.replication_pending_size_v1, ..Default::default() }, ); } } } } Ok(data_usage_info) } /// Create a data usage cache entry from size summary pub fn create_cache_entry_from_summary(summary: &SizeSummary) -> DataUsageEntry { let mut entry = DataUsageEntry::default(); entry.add_sizes(summary); entry } /// Convert data usage cache to DataUsageInfo pub fn cache_to_data_usage_info(cache: &DataUsageCache, path: &str, buckets: &[crate::store_api::BucketInfo]) -> DataUsageInfo { let e = match cache.find(path) { Some(e) => e, None => return DataUsageInfo::default(), }; let flat = cache.flatten(&e); let mut buckets_usage = HashMap::new(); for bucket in buckets.iter() { let e = match cache.find(&bucket.name) { Some(e) => e, None => continue, }; let flat = cache.flatten(&e); let mut bui = rustfs_common::data_usage::BucketUsageInfo { size: flat.size as u64, versions_count: flat.versions as u64, objects_count: flat.objects as u64, delete_markers_count: flat.delete_markers as u64, object_size_histogram: flat.obj_sizes.to_map(), object_versions_histogram: flat.obj_versions.to_map(), ..Default::default() }; if let Some(rs) = &flat.replication_stats { bui.replica_size = rs.replica_size; bui.replica_count = rs.replica_count; for (arn, stat) in rs.targets.iter() { bui.replication_info.insert( arn.clone(), BucketTargetUsageInfo { replication_pending_size: stat.pending_size, replicated_size: stat.replicated_size, replication_failed_size: stat.failed_size, replication_pending_count: stat.pending_count, replication_failed_count: stat.failed_count, replicated_count: stat.replicated_count, ..Default::default() }, ); } } buckets_usage.insert(bucket.name.clone(), bui); } DataUsageInfo { last_update: cache.info.last_update, objects_total_count: flat.objects as u64, versions_total_count: flat.versions as u64, delete_markers_total_count: flat.delete_markers as u64, objects_total_size: flat.size as u64, buckets_count: e.children.len() as u64, buckets_usage, ..Default::default() } } // Helper functions for DataUsageCache operations pub async fn load_data_usage_cache(store: &crate::set_disk::SetDisks, name: &str) -> crate::error::Result { use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; use crate::store_api::{ObjectIO, ObjectOptions}; use http::HeaderMap; use rand::Rng; use std::path::Path; use std::time::Duration; use tokio::time::sleep; let mut d = DataUsageCache::default(); let mut retries = 0; while retries < 5 { let path = Path::new(BUCKET_META_PREFIX).join(name); match store .get_object_reader( RUSTFS_META_BUCKET, path.to_str().unwrap(), None, HeaderMap::new(), &ObjectOptions { no_lock: true, ..Default::default() }, ) .await { Ok(mut reader) => { if let Ok(info) = DataUsageCache::unmarshal(&reader.read_all().await?) { d = info } break; } Err(err) => match err { crate::error::Error::FileNotFound | crate::error::Error::VolumeNotFound => { match store .get_object_reader( RUSTFS_META_BUCKET, name, None, HeaderMap::new(), &ObjectOptions { no_lock: true, ..Default::default() }, ) .await { Ok(mut reader) => { if let Ok(info) = DataUsageCache::unmarshal(&reader.read_all().await?) { d = info } break; } Err(_) => match err { crate::error::Error::FileNotFound | crate::error::Error::VolumeNotFound => { break; } _ => {} }, } } _ => { break; } }, } retries += 1; let dur = { let mut rng = rand::rng(); rng.random_range(0..1_000) }; sleep(Duration::from_millis(dur)).await; } Ok(d) } pub async fn save_data_usage_cache(cache: &DataUsageCache, name: &str) -> crate::error::Result<()> { use crate::config::com::save_config; use crate::disk::BUCKET_META_PREFIX; use crate::new_object_layer_fn; use std::path::Path; let Some(store) = new_object_layer_fn() else { return Err(crate::error::Error::other("errServerNotInitialized")); }; let buf = cache.marshal_msg().map_err(crate::error::Error::other)?; let buf_clone = buf.clone(); let store_clone = store.clone(); let name = Path::new(BUCKET_META_PREFIX).join(name).to_string_lossy().to_string(); let name_clone = name.clone(); tokio::spawn(async move { let _ = save_config(store_clone, &format!("{}{}", &name_clone, ".bkp"), buf_clone).await; }); save_config(store, &name, buf).await?; Ok(()) }