From cefdf4d6f9399438942c45e6aefcc1aeb1d09080 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Fri, 11 Apr 2025 04:08:20 +0000 Subject: [PATCH] fix datascanner Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/disk/local.rs | 33 +++++- ecstore/src/heal/data_scanner.rs | 152 ++++++++++++++++++++++----- ecstore/src/heal/data_usage_cache.rs | 5 +- ecstore/src/set_disk.rs | 6 +- ecstore/src/store.rs | 5 +- 5 files changed, 161 insertions(+), 40 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 93681d91f..24e5782fe 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -11,6 +11,8 @@ use super::{ }; use crate::bitrot::bitrot_verify; use crate::bucket::metadata_sys::{self}; +use crate::bucket::versioning::VersioningApi; +use crate::bucket::versioning_sys::BucketVersioningSys; use crate::cache_value::cache::{Cache, Opts, UpdateFn}; use crate::disk::error::{ convert_access_error, is_err_os_not_exist, is_sys_err_handle_invalid, is_sys_err_invalid_arg, is_sys_err_is_dir, @@ -20,7 +22,9 @@ use crate::disk::os::{check_path_length, is_empty_dir}; use crate::disk::STORAGE_FORMAT_FILE; use crate::file_meta::{get_file_info, read_xl_meta_no_data, FileInfoOpts}; use crate::global::{GLOBAL_IsErasureSD, GLOBAL_RootDiskThreshold}; -use crate::heal::data_scanner::{has_active_rules, scan_data_folder, ScannerItem, ShouldSleepFn, SizeSummary}; +use crate::heal::data_scanner::{ + lc_has_active_rules, rep_has_active_rules, scan_data_folder, ScannerItem, ShouldSleepFn, SizeSummary, +}; use crate::heal::data_scanner_metric::{ScannerMetric, ScannerMetrics}; use crate::heal::data_usage_cache::{DataUsageCache, DataUsageEntry}; use crate::heal::error::{ERR_IGNORE_FILE_CONTRIB, ERR_SKIP_FILE}; @@ -2295,6 +2299,7 @@ impl DiskAPI for LocalDisk { Ok(info) } + #[tracing::instrument(level = "info", skip_all)] async fn ns_scanner( &self, cache: &DataUsageCache, @@ -2308,18 +2313,30 @@ impl DiskAPI for LocalDisk { // must befor metadata_sys let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + let mut cache = cache.clone(); + // Check if the current bucket has a configured lifecycle policy + if let Ok((lc, _)) = metadata_sys::get_lifecycle_config(&cache.info.name).await { + if lc_has_active_rules(&lc, "") { + cache.info.life_cycle = Some(lc); + } + } + // Check if the current bucket has replication configuration if let Ok((rcfg, _)) = metadata_sys::get_replication_config(&cache.info.name).await { - if has_active_rules(&rcfg, "", true) { + if rep_has_active_rules(&rcfg, "", true) { // TODO: globalBucketTargetSys } } + let vcfg = match BucketVersioningSys::get(&cache.info.name).await { + Ok(vcfg) => Some(vcfg), + Err(_) => None, + }; + let loc = self.get_disk_location(); let disks = store.get_disks(loc.pool_idx.unwrap(), loc.disk_idx.unwrap()).await?; let disk = Arc::new(LocalDisk::new(&self.endpoint(), false).await?); let disk_clone = disk.clone(); - let mut cache = cache.clone(); cache.info.updates = Some(updates.clone()); let mut data_usage_info = scan_data_folder( &disks, @@ -2328,8 +2345,9 @@ impl DiskAPI for LocalDisk { Box::new(move |item: &ScannerItem| { let mut item = item.clone(); let disk = disk_clone.clone(); + let vcfg = vcfg.clone(); Box::pin(async move { - if item.path.ends_with(&format!("{}{}", SLASH_SEPARATOR, STORAGE_FORMAT_FILE)) { + if !item.path.ends_with(&format!("{}{}", SLASH_SEPARATOR, STORAGE_FORMAT_FILE)) { return Err(Error::from_string(ERR_SKIP_FILE)); } let stop_fn = ScannerMetrics::log(ScannerMetric::ScanObject); @@ -2370,7 +2388,12 @@ impl DiskAPI for LocalDisk { } }; - let versioned = false; + let versioned = if let Some(vcfg) = vcfg.as_ref() { + vcfg.versioned(item.object_path().to_str().unwrap_or_default()) + } else { + false + }; + let mut obj_deleted = false; for info in obj_infos.iter() { let done = ScannerMetrics::time(ScannerMetric::ApplyVersion); diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index 6701f7ba8..315b21fef 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -18,7 +18,10 @@ use super::{ data_usage_cache::{DataUsageCache, DataUsageEntry, DataUsageHash}, heal_commands::{HealScanMode, HEAL_DEEP_SCAN, HEAL_NORMAL_SCAN}, }; -use crate::heal::data_usage::DATA_USAGE_ROOT; +use crate::{ + bucket::{versioning::VersioningApi, versioning_sys::BucketVersioningSys}, + heal::data_usage::DATA_USAGE_ROOT, +}; use crate::{ cache_value::metacache_set::{list_path_raw, ListPathRawOptions}, config::{ @@ -49,7 +52,7 @@ use common::error::{Error, Result}; use lazy_static::lazy_static; use rand::Rng; use rmp_serde::{Deserializer, Serializer}; -use s3s::dto::{ReplicationConfiguration, ReplicationRuleStatus}; +use s3s::dto::{BucketLifecycleConfiguration, ExpirationStatus, LifecycleRule, ReplicationConfiguration, ReplicationRuleStatus}; use serde::{Deserialize, Serialize}; use tokio::{ sync::{ @@ -412,7 +415,7 @@ pub struct ScannerItem { pub prefix: String, pub object_name: String, pub replication: Option, - // todo: lifecycle + pub lifecycle: Option, // typ: fs::Permissions, pub heal: Heal, pub debug: bool, @@ -435,7 +438,7 @@ impl ScannerItem { pub async fn apply_versions_actions(&self, fivs: &[FileInfo]) -> Result> { let obj_infos = self.apply_newer_noncurrent_version_limit(fivs).await?; - if obj_infos.len() >= >::try_into(SCANNER_EXCESS_OBJECT_VERSIONS.load(Ordering::SeqCst)).unwrap() { + if obj_infos.len() >= SCANNER_EXCESS_OBJECT_VERSIONS.load(Ordering::SeqCst) as usize { // todo } @@ -444,9 +447,7 @@ impl ScannerItem { cumulative_size += obj_info.size; } - if cumulative_size - >= >::try_into(SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL_SIZE.load(Ordering::SeqCst)).unwrap() - { + if cumulative_size >= SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL_SIZE.load(Ordering::SeqCst) as usize { //todo } @@ -454,22 +455,31 @@ impl ScannerItem { } pub async fn apply_newer_noncurrent_version_limit(&self, fivs: &[FileInfo]) -> Result> { - let done = ScannerMetrics::time(ScannerMetric::ApplyNonCurrent); - let mut object_infos = Vec::new(); - for info in fivs.iter() { - object_infos.push(info.to_object_info(&self.bucket, &self.object_path().to_string_lossy(), false)); + // let done = ScannerMetrics::time(ScannerMetric::ApplyNonCurrent); + let versioned = match BucketVersioningSys::get(&self.bucket).await { + Ok(vcfg) => vcfg.versioned(self.object_path().to_str().unwrap_or_default()), + Err(_) => false, + }; + let mut object_infos = Vec::with_capacity(fivs.len()); + + if self.lifecycle.is_none() { + for info in fivs.iter() { + object_infos.push(info.to_object_info(&self.bucket, &self.object_path().to_string_lossy(), versioned)); + } + return Ok(object_infos); } - done().await; + + // done().await; Ok(object_infos) } - pub async fn apply_actions(&self, _oi: &ObjectInfo, _size_s: &SizeSummary) -> (bool, usize) { + pub async fn apply_actions(&self, oi: &ObjectInfo, _size_s: &SizeSummary) -> (bool, usize) { let done = ScannerMetrics::time(ScannerMetric::Ilm); //todo: lifecycle done().await; - (false, 0) + (false, oi.size) } } @@ -548,6 +558,7 @@ impl FolderScanner { true } + #[tracing::instrument(level = "info", skip_all)] async fn scan_folder(&mut self, folder: &CachedFolder, into: &mut DataUsageEntry) -> Result<()> { let this_hash = hash_path(&folder.name); let was_compacted = into.compacted; @@ -561,8 +572,18 @@ impl FolderScanner { let (_, prefix) = path_to_bucket_object_with_base_path(&self.root, &folder.name); // Todo: lifeCycle + let active_life_cycle = if let Some(lc) = self.old_cache.info.life_cycle.as_ref() { + if lc_has_active_rules(lc, &prefix) { + self.old_cache.info.life_cycle.clone() + } else { + None + } + } else { + None + }; + let replication_cfg = if self.old_cache.info.replication.is_some() - && has_active_rules(self.old_cache.info.replication.as_ref().unwrap(), &prefix, true) + && rep_has_active_rules(self.old_cache.info.replication.as_ref().unwrap(), &prefix, true) { self.old_cache.info.replication.clone() } else { @@ -593,7 +614,7 @@ impl FolderScanner { continue; } - if !sub_path.is_dir() { + if sub_path.is_dir() { let h = hash_path(ent_name.to_str().unwrap()); if h == this_hash { continue; @@ -638,6 +659,7 @@ impl FolderScanner { .unwrap_or_default(), debug: self.data_usage_scanner_debug, replication: replication_cfg.clone(), + lifecycle: active_life_cycle.clone(), heal: Heal::default(), }; @@ -681,13 +703,11 @@ impl FolderScanner { } let should_compact = self.new_cache.info.name != folder.name - && existing_folders.len() + new_folders.len() - >= >::try_into(DATA_SCANNER_COMPACT_AT_FOLDERS).unwrap() - || existing_folders.len() + new_folders.len() - >= >::try_into(DATA_SCANNER_FORCE_COMPACT_AT_FOLDERS).unwrap(); + && (existing_folders.len() + new_folders.len() >= DATA_SCANNER_COMPACT_AT_FOLDERS as usize + || existing_folders.len() + new_folders.len() >= DATA_SCANNER_FORCE_COMPACT_AT_FOLDERS as usize); let total_folders = existing_folders.len() + new_folders.len(); - if total_folders > >::try_into(SCANNER_EXCESS_FOLDERS.load(Ordering::SeqCst)).unwrap() { + if total_folders > SCANNER_EXCESS_FOLDERS.load(Ordering::SeqCst) as usize { let _prefix_name = format!("{}/", folder.name.trim_end_matches('/')); // todo: notification } @@ -935,7 +955,7 @@ impl FolderScanner { let _ = list_path_raw(rx, lopts).await; if *found_objs.read().await { - let this = CachedFolder { + let this: CachedFolder = CachedFolder { name: k.clone(), parent: this_hash.clone(), object_heal_prob_div: 1, @@ -951,7 +971,7 @@ impl FolderScanner { if !into.compacted && self.new_cache.info.name != folder.name { let mut flat = self.new_cache.size_recursive(&this_hash.key()).unwrap_or_default(); flat.compacted = true; - let compact = if flat.objects < >::try_into(DATA_SCANNER_COMPACT_LEAST_OBJECT).unwrap() { + let compact = if flat.objects < DATA_SCANNER_COMPACT_LEAST_OBJECT as usize { true } else { // Compact if we only have objects as children... @@ -996,6 +1016,7 @@ impl FolderScanner { Ok(()) } + #[tracing::instrument(level = "info", skip_all)] async fn send_update(&mut self) { if SystemTime::now().duration_since(self.last_update).unwrap() < Duration::from_secs(60) { return; @@ -1007,8 +1028,9 @@ impl FolderScanner { } } +#[tracing::instrument(level = "info", skip(into, folder_scanner))] async fn scan(folder: &CachedFolder, into: &mut DataUsageEntry, folder_scanner: &mut FolderScanner) { - let mut dst = if !into.compacted { + let mut dst = if into.compacted { DataUsageEntry::default() } else { into.clone() @@ -1028,7 +1050,74 @@ async fn scan(folder: &CachedFolder, into: &mut DataUsageEntry, folder_scanner: } } -pub fn has_active_rules(config: &ReplicationConfiguration, prefix: &str, recursive: bool) -> bool { +fn lc_get_prefix(rule: &LifecycleRule) -> String { + if let Some(p) = &rule.prefix { + return p.to_string(); + } else if let Some(filter) = &rule.filter { + if let Some(p) = &filter.prefix { + return p.to_string(); + } else if let Some(and) = &filter.and { + if let Some(p) = &and.prefix { + return p.to_string(); + } + } + } + + "".into() +} + +pub fn lc_has_active_rules(config: &BucketLifecycleConfiguration, prefix: &str) -> bool { + if config.rules.is_empty() { + return false; + } + + for rule in config.rules.iter() { + if rule.status == ExpirationStatus::from_static(ExpirationStatus::DISABLED) { + continue; + } + let rule_prefix = lc_get_prefix(rule); + if !prefix.is_empty() && !rule_prefix.is_empty() { + if !prefix.starts_with(&rule_prefix) && !rule_prefix.starts_with(prefix) { + continue; + } + } + + if let Some(e) = &rule.noncurrent_version_expiration { + if let Some(true) = e.noncurrent_days.map(|d| d > 0) { + return true; + } + if let Some(true) = e.newer_noncurrent_versions.map(|d| d > 0) { + return true; + } + } + + if rule.noncurrent_version_transitions.is_some() { + return true; + } + if let Some(true) = rule.expiration.as_ref().map(|e| e.date.is_some()) { + return true; + } + + if let Some(true) = rule.expiration.as_ref().map(|e| e.days.is_some()) { + return true; + } + + if let Some(Some(true)) = rule.expiration.as_ref().map(|e| e.expired_object_delete_marker.map(|m| m)) { + return true; + } + + if let Some(true) = rule.transitions.as_ref().map(|t| !t.is_empty()) { + return true; + } + + if rule.transitions.is_some() { + return true; + } + } + false +} + +pub fn rep_has_active_rules(config: &ReplicationConfiguration, prefix: &str, recursive: bool) -> bool { if config.rules.is_empty() { return false; } @@ -1086,8 +1175,14 @@ pub async fn scan_data_folder( root: base_path, get_size: get_size_fn, old_cache: cache.clone(), - new_cache: DataUsageCache::default(), - update_cache: DataUsageCache::default(), + new_cache: DataUsageCache { + info: cache.info.clone(), + ..Default::default() + }, + update_cache: DataUsageCache { + info: cache.info.clone(), + ..Default::default() + }, data_usage_scanner_debug: false, heal_object_select: 0, scan_mode: heal_scan_mode, @@ -1115,8 +1210,7 @@ pub async fn scan_data_folder( if s.scan_folder(&folder, &mut root).await.is_err() { close_disk().await; } - s.new_cache - .force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN.try_into().unwrap()); + s.new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN as usize); s.new_cache.info.last_update = Some(SystemTime::now()); s.new_cache.info.next_cycle = cache.info.next_cycle; close_disk().await; diff --git a/ecstore/src/heal/data_usage_cache.rs b/ecstore/src/heal/data_usage_cache.rs index cca9d79f3..de5b73734 100644 --- a/ecstore/src/heal/data_usage_cache.rs +++ b/ecstore/src/heal/data_usage_cache.rs @@ -10,7 +10,7 @@ use http::HeaderMap; use path_clean::PathClean; use rand::Rng; use rmp_serde::Serializer; -use s3s::dto::ReplicationConfiguration; +use s3s::dto::{BucketLifecycleConfiguration, ReplicationConfiguration}; use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet}; use std::hash::{DefaultHasher, Hash, Hasher}; @@ -347,7 +347,8 @@ pub struct DataUsageCacheInfo { pub next_cycle: u32, pub last_update: Option, pub skip_healing: bool, - // todo: life_cycle + #[serde(skip)] + pub life_cycle: Option, // pub life_cycle: #[serde(skip)] pub updates: Option>, diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 04b3254a4..63ad26d89 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -3015,7 +3015,7 @@ impl SetDisks { }); // Calc usage let before = cache.info.last_update; - let cache = match disk.clone().ns_scanner(&cache, tx, heal_scan_mode, None).await { + let mut cache = match disk.clone().ns_scanner(&cache, tx, heal_scan_mode, None).await { Ok(cache) => cache, Err(_) => { if cache.info.last_update > before { @@ -3025,6 +3025,9 @@ impl SetDisks { continue; } }; + + cache.info.updates = None; + let _ = task.await; let mut root = DataUsageEntry::default(); if let Some(r) = cache.root() { root = cache.flatten(&r); @@ -3041,7 +3044,6 @@ impl SetDisks { entry: root, }) .await; - let _ = task.await; let _ = cache.save(&cache_name.to_string_lossy()).await; } } diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 1f57b704c..793e31a26 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -740,7 +740,7 @@ impl ECStore { let cancel_clone = cancel.clone(); let all_buckets_clone = all_buckets.clone(); futures.push(async move { - let (tx, mut rx) = mpsc::channel(100); + let (tx, mut rx) = mpsc::channel(1); let task = tokio::spawn(async move { loop { match rx.recv().await { @@ -951,6 +951,7 @@ impl ECStore { } } +#[tracing::instrument(level = "info", skip(all_buckets, updates))] async fn update_scan( all_merged: Arc>, results: Arc>>, @@ -972,7 +973,7 @@ async fn update_scan( } w.merge(info); } - if w.info.last_update > *last_update && w.root().is_none() { + if (last_update.is_none() || w.info.last_update > *last_update) && w.root().is_some() { let _ = updates.send(w.dui(&w.info.name, &all_buckets)).await; *last_update = w.info.last_update; }