From 7148c5d5cd38d5a1540fff3af50abfd170f4898c Mon Sep 17 00:00:00 2001 From: mujunxiang <1948535941@qq.com> Date: Tue, 3 Dec 2024 15:31:43 +0800 Subject: [PATCH] sanner(1) Signed-off-by: mujunxiang <1948535941@qq.com> --- ecstore/src/disk/local.rs | 1 + ecstore/src/disk/mod.rs | 1 + ecstore/src/heal/data_scanner.rs | 18 +++++++-- ecstore/src/heal/data_scanner_metric.rs | 2 + ecstore/src/heal/data_usage_cache.rs | 10 ++--- ecstore/src/metrics_realtime.rs | 50 ++++++++++++++++++++++++- ecstore/src/set_disk.rs | 2 + rustfs/src/admin/handlers.rs | 2 +- 8 files changed, 76 insertions(+), 10 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index c85db05fe..6117df74b 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -2092,6 +2092,7 @@ impl DiskAPI for LocalDisk { ) .await?; data_usage_info.info.last_update = Some(SystemTime::now()); + info!("ns_scanner completed: {data_usage_info:?}"); Ok(data_usage_info) } diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 079d299b0..ab57ce45a 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -351,6 +351,7 @@ impl DiskAPI for Disk { updates: Sender, scan_mode: HealScanMode, ) -> Result { + info!("ns_scanner"); match self { Disk::Local(local_disk) => local_disk.ns_scanner(cache, updates, scan_mode).await, Disk::Remote(remote_disk) => remote_disk.ns_scanner(cache, updates, scan_mode).await, diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index 8562d32df..16ff8cc05 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -86,23 +86,28 @@ lazy_static! { } pub async fn init_data_scanner() { - let mut r = rand::thread_rng(); - let random = r.gen_range(0.0..1.0); tokio::spawn(async move { loop { run_data_scanner().await; + let random = { + let mut r = rand::thread_rng(); + r.gen_range(0.0..1.0) + }; let duration = Duration::from_secs_f64(random * (SCANNER_CYCLE.load(std::sync::atomic::Ordering::SeqCst) as f64)); let sleep_duration = if duration < Duration::new(1, 0) { Duration::new(1, 0) } else { duration }; + + info!("data scanner will sleeping {sleep_duration:?}"); sleep(sleep_duration).await; } }); } async fn run_data_scanner() { + info!("run_data_scanner"); let Some(store) = new_object_layer_fn() else { error!("errServerNotInitialized"); return; @@ -163,8 +168,10 @@ async fn run_data_scanner() { }); let mut res = HashMap::new(); res.insert("cycle".to_string(), cycle_info.current.to_string()); + info!("start ns_scanner"); match store.clone().ns_scanner(tx, cycle_info.current as usize, scan_mode).await { Ok(_) => { + info!("ns_scanner completed"); cycle_info.next += 1; cycle_info.current = 0; cycle_info.cycle_completed.push(Utc::now()); @@ -176,9 +183,14 @@ async fn run_data_scanner() { globalScannerMetrics.write().await.set_cycle(Some(cycle_info.clone())).await; let mut tmp = Vec::new(); tmp.write_u64::(cycle_info.next).unwrap(); - let _ = save_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH, &tmp).await; + if let Ok(data) = cycle_info.marshal_msg(&tmp) { + let _ = save_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH, &data).await; + } else { + let _ = save_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH, &tmp).await; + } } Err(err) => { + info!("ns_scanner failed: {:?}", err); res.insert("error".to_string(), err.to_string()); } } diff --git a/ecstore/src/heal/data_scanner_metric.rs b/ecstore/src/heal/data_scanner_metric.rs index b0f364bdd..77b97c1d5 100644 --- a/ecstore/src/heal/data_scanner_metric.rs +++ b/ecstore/src/heal/data_scanner_metric.rs @@ -14,6 +14,7 @@ use std::{ time::SystemTime, }; use tokio::sync::RwLock; +use tracing::info; use super::data_scanner::{CurrentScannerCycle, UpdateCurrentPathFn}; @@ -148,6 +149,7 @@ impl ScannerMetrics { } pub async fn set_cycle(&mut self, c: Option) { + info!("ScannerMetrics set_cycle {c:?}"); *self.cycle_info.write().await = c; } diff --git a/ecstore/src/heal/data_usage_cache.rs b/ecstore/src/heal/data_usage_cache.rs index 09be08b99..fbc6003b6 100644 --- a/ecstore/src/heal/data_usage_cache.rs +++ b/ecstore/src/heal/data_usage_cache.rs @@ -131,7 +131,7 @@ const OBJECTS_VERSION_COUNT_INTERVALS: [ObjectHistogramInterval; DATA_USAGE_VERS ]; // sizeHistogram is a size histogram. -#[derive(Clone, Serialize, Deserialize)] +#[derive(Clone, Debug, Serialize, Deserialize)] pub struct SizeHistogram(Vec); impl Default for SizeHistogram { @@ -168,7 +168,7 @@ impl SizeHistogram { } // versionsHistogram is a histogram of number of versions in an object. -#[derive(Clone, Serialize, Deserialize)] +#[derive(Clone, Debug, Serialize, Deserialize)] pub struct VersionsHistogram(Vec); impl Default for VersionsHistogram { @@ -238,7 +238,7 @@ impl ReplicationAllStats { } } -#[derive(Clone, Default, Serialize, Deserialize)] +#[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageEntry { pub children: DataUsageHashMap, // These fields do no include any children. @@ -340,7 +340,7 @@ pub struct DataUsageEntryInfo { pub entry: DataUsageEntry, } -#[derive(Clone, Default, Serialize, Deserialize)] +#[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageCacheInfo { pub name: String, pub next_cycle: u32, @@ -367,7 +367,7 @@ pub struct DataUsageCacheInfo { // } // } -#[derive(Clone, Default, Serialize, Deserialize)] +#[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageCache { pub info: DataUsageCacheInfo, pub cache: HashMap, diff --git a/ecstore/src/metrics_realtime.rs b/ecstore/src/metrics_realtime.rs index a219723b6..97b396970 100644 --- a/ecstore/src/metrics_realtime.rs +++ b/ecstore/src/metrics_realtime.rs @@ -4,6 +4,7 @@ use chrono::Utc; use common::globals::{GLOBAL_Local_Node_Name, GLOBAL_Rustfs_Addr}; use madmin::metrics::{DiskIOStats, DiskMetric, RealtimeMetrics}; use serde::{Deserialize, Serialize}; +use tracing::info; use crate::{ admin_server_info::get_local_server_property, @@ -55,8 +56,10 @@ impl MetricType { } pub async fn collect_local_metrics(types: MetricType, opts: &CollectMetricsOpts) -> RealtimeMetrics { + info!("collect_local_metrics"); let mut real_time_metrics = RealtimeMetrics::default(); if types.0 == MetricType::NONE.0 { + info!("types is None, return"); return real_time_metrics; } @@ -75,11 +78,13 @@ pub async fn collect_local_metrics(types: MetricType, opts: &CollectMetricsOpts) } if types.contains(&MetricType::DISK) { + info!("start get disk metrics"); let mut aggr = DiskMetric { collected_at: Utc::now(), ..Default::default() }; for (name, disk) in collect_local_disks_metrics(&opts.disks).await.into_iter() { + info!("got disk metric, name: {name}, metric: {disk:?}"); real_time_metrics.by_disk.insert(name, disk.clone()); aggr.merge(&disk); } @@ -87,10 +92,31 @@ pub async fn collect_local_metrics(types: MetricType, opts: &CollectMetricsOpts) } if types.contains(&MetricType::SCANNER) { + info!("start get scanner metrics"); let metrics = globalScannerMetrics.read().await.report().await; real_time_metrics.aggregated.scanner = Some(metrics); } - RealtimeMetrics::default() + + if types.contains(&MetricType::OS) {} + + if types.contains(&MetricType::BATCH_JOBS) {} + + if types.contains(&MetricType::SITE_RESYNC) {} + + if types.contains(&MetricType::NET) {} + + if types.contains(&MetricType::MEM) {} + + if types.contains(&MetricType::CPU) {} + + if types.contains(&MetricType::RPC) {} + + real_time_metrics + .by_host + .insert(by_host_name.clone(), real_time_metrics.aggregated.clone()); + real_time_metrics.hosts.push(by_host_name); + + real_time_metrics } async fn collect_local_disks_metrics(disks: &HashSet) -> HashMap { @@ -167,3 +193,25 @@ async fn collect_local_disks_metrics(disks: &HashSet) -> HashMap, heal_scan_mode: HealScanMode, ) -> Result<()> { + info!("ns_scanner"); if buckets.is_empty() { return Ok(()); } @@ -2846,6 +2847,7 @@ impl SetDisks { } let _ = join_all(futures).await; let _ = task.await; + info!("ns_scanner completed"); Ok(()) } diff --git a/rustfs/src/admin/handlers.rs b/rustfs/src/admin/handlers.rs index 1419beaed..fcec96c97 100644 --- a/rustfs/src/admin/handlers.rs +++ b/rustfs/src/admin/handlers.rs @@ -431,7 +431,7 @@ impl Operation for MetricsHandler { // todo write resp match serde_json::to_vec(&m) { Ok(re) => { - info!("got metrics, send it to client, m: {m:?}, re: {re:?}"); + info!("got metrics, send it to client, m: {m:?}"); let _ = tx.send(Ok(Bytes::from(re))).await; } Err(e) => {