From ec8f05950612adb3f418f6e6bcfdf688af7ef992 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 7 Apr 2026 20:32:32 +0800 Subject: [PATCH] feat: add timing and metrics for query execution phases (#2419) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Shreyan Jain Co-authored-by: Shreyan Jain Co-authored-by: heihutu Co-authored-by: loverustfs Co-authored-by: cxymds Co-authored-by: 安正超 Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: houseme <4829346+houseme@users.noreply.github.com> --- Cargo.lock | 1 + crates/common/src/metrics.rs | 432 ++++++++++--------- crates/ecstore/src/bucket/quota/checker.rs | 7 +- crates/kms/src/api_types.rs | 9 +- crates/s3select-api/Cargo.toml | 1 + crates/s3select-api/src/query/execution.rs | 74 +++- crates/s3select-query/src/execution/query.rs | 34 +- 7 files changed, 309 insertions(+), 249 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8d0de0f17..f7d162eaf 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8434,6 +8434,7 @@ dependencies = [ "futures", "futures-core", "http 1.4.0", + "metrics", "object_store", "parking_lot 0.12.5", "pin-project-lite", diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index 2753397eb..7d41a9dc1 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -22,12 +22,12 @@ use std::{ future::Future, pin::Pin, sync::{ - Arc, OnceLock, + Arc, Mutex, OnceLock, atomic::{AtomicU64, Ordering}, }, time::{Duration, SystemTime}, }; -use tokio::sync::{Mutex, RwLock}; +use tokio::sync::RwLock; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum IlmAction { @@ -110,7 +110,7 @@ pub fn global_metrics() -> &'static Arc { GLOBAL_METRICS.get_or_init(|| Arc::new(Metrics::new())) } -#[derive(Clone, Debug, PartialEq, PartialOrd)] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Metric { // START Realtime metrics, that only records // last minute latencies and total operation count. @@ -188,7 +188,6 @@ impl Metric { if index >= Self::Last as usize { return None; } - // Safe conversion using match instead of unsafe transmute match index { 0 => Some(Self::ReadMetadata), 1 => Some(Self::CheckMissing), @@ -220,16 +219,48 @@ impl Metric { } } -/// Thread-safe wrapper for LastMinuteLatency with atomic operations -#[derive(Default)] +// --------------------------------------------------------------------------- +// LockedLastMinuteLatency +// --------------------------------------------------------------------------- +// +// Uses std::sync::Mutex instead of tokio::sync::Mutex. +// +// Rationale: the critical section is a handful of integer additions inside +// LastMinuteLatency::add_all / get_total — no I/O, no blocking syscalls, no +// awaiting. A std blocking mutex is cheaper (no task-wakeup overhead) and, +// crucially, lets every caller be *synchronous*. That eliminates the need to +// spawn a background task just to record a duration, which was the only reason +// tokio::spawn appeared in log/time/time_size/time_n/time_ilm. +// +// Note: vec![LockedLastMinuteLatency::default(); N] with the old Arc-based +// Clone made every element share the *same* inner mutex — a latent bug where +// all metrics slots wrote to one counter. The new Clone creates a fresh +// independent Mutex per element, matching the intent. + +/// Thread-safe wrapper for LastMinuteLatency backed by a std blocking mutex. pub struct LockedLastMinuteLatency { + // Arc so Clone is cheap *and* each cloned value stays independent (its own + // allocation). We never hand out the Arc to two Metrics at once; the Arc + // is purely a convenience for the Clone impl below. latency: Arc>, } +impl Default for LockedLastMinuteLatency { + fn default() -> Self { + Self::new() + } +} + +// Produce a fresh, independent slot — *not* a shared alias. This is what +// vec![val; N] and #[derive(Clone)] on the parent struct both need. impl Clone for LockedLastMinuteLatency { fn clone(&self) -> Self { + let inner = match self.latency.lock() { + Ok(guard) => guard.clone(), + Err(poisoned) => poisoned.into_inner().clone(), + }; Self { - latency: Arc::clone(&self.latency), + latency: Arc::new(Mutex::new(inner)), } } } @@ -241,14 +272,13 @@ impl LockedLastMinuteLatency { } } - /// Add a duration measurement - pub async fn add(&self, duration: Duration) { - self.add_size(duration, 0).await; + /// Record a duration sample (no size). + pub fn add(&self, duration: Duration) { + self.add_size(duration, 0); } - /// Add a duration measurement with size - pub async fn add_size(&self, duration: Duration, size: u64) { - let mut latency = self.latency.lock().await; + /// Record a duration sample with an associated byte count. + pub fn add_size(&self, duration: Duration, size: u64) { let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) .unwrap_or_default() @@ -259,17 +289,27 @@ impl LockedLastMinuteLatency { total: duration.as_secs(), size, }; - latency.add_all(now, &elem); + + match self.latency.lock() { + Ok(mut guard) => guard.add_all(now, &elem), + Err(poisoned) => poisoned.into_inner().add_all(now, &elem), + } } - /// Get total accumulated metrics for the last minute - pub async fn total(&self) -> AccElem { - let mut latency = self.latency.lock().await; - latency.get_total() + /// Return accumulated totals for the last minute window. + pub fn total(&self) -> AccElem { + match self.latency.lock() { + Ok(mut latency) => latency.get_total(), + Err(poisoned) => poisoned.into_inner().get_total(), + } } } -/// Current path tracker for monitoring active scan paths +// --------------------------------------------------------------------------- +// CurrentPathTracker — unchanged, still uses tokio::sync::RwLock because it +// lives inside async path-update callbacks. +// --------------------------------------------------------------------------- + struct CurrentPathTracker { current_path: Arc>, } @@ -290,17 +330,16 @@ impl CurrentPathTracker { } } -/// Main scanner metrics structure +// --------------------------------------------------------------------------- +// Metrics +// --------------------------------------------------------------------------- + pub struct Metrics { - // All fields must be accessed atomically and aligned. operations: Vec, latency: Vec, actions: Vec, actions_latency: Vec, - // Current paths contains disk -> tracker mappings current_paths: Arc>>>, - - // Cycle information cycle_info: Arc>>, } @@ -331,8 +370,6 @@ const OTEL_SCANNER_CYCLES: &str = "rustfs_scanner_cycles_total"; const OTEL_SCANNER_CYCLE_DURATION_SECONDS: &str = "rustfs_scanner_cycle_duration_seconds"; const OTEL_SCANNER_BUCKET_DRIVE_DURATION_SECONDS: &str = "rustfs_scanner_bucket_drive_duration_seconds"; -/// Emit an OTEL counter increment for the given scanner metric. -/// ScanCycle and ScanBucketDrive are handled by dedicated emit functions with labels. fn emit_otel_counter(metric: usize, count: u64) { match Metric::from_index(metric) { Some(Metric::ScanObject) => { @@ -345,8 +382,6 @@ fn emit_otel_counter(metric: usize, count: u64) { } } -/// Emit OTel metrics for a completed scan cycle. -/// Counter with result label + gauge for last successful cycle duration. pub fn emit_scan_cycle_complete(success: bool, duration: Duration) { let result = if success { "success" } else { "error" }; metrics::counter!(OTEL_SCANNER_CYCLES, "result" => result).increment(1); @@ -355,8 +390,6 @@ pub fn emit_scan_cycle_complete(success: bool, duration: Duration) { } } -/// Emit OTel metrics for a completed bucket-drive scan. -/// Counter with result/bucket/disk labels + histogram for duration. pub fn emit_scan_bucket_drive_complete(success: bool, bucket: &str, disk: &str, duration: Duration) { let result = if success { "success" } else { "error" }; metrics::counter!( @@ -376,226 +409,202 @@ pub fn emit_scan_bucket_drive_complete(success: bool, bucket: &str, disk: &str, impl Metrics { pub fn new() -> Self { - let operations = (0..Metric::Last as usize).map(|_| AtomicU64::new(0)).collect(); - - let latency = (0..Metric::LastRealtime as usize) - .map(|_| LockedLastMinuteLatency::new()) - .collect(); - Self { - operations, - latency, + operations: (0..Metric::Last as usize).map(|_| AtomicU64::new(0)).collect(), + // Each slot gets its own fresh LockedLastMinuteLatency so that + // different metrics never accidentally share state. + latency: (0..Metric::LastRealtime as usize) + .map(|_| LockedLastMinuteLatency::new()) + .collect(), actions: (0..IlmAction::ActionCount as usize).map(|_| AtomicU64::new(0)).collect(), - actions_latency: vec![LockedLastMinuteLatency::default(); IlmAction::ActionCount as usize], + actions_latency: (0..IlmAction::ActionCount as usize) + .map(|_| LockedLastMinuteLatency::new()) + .collect(), current_paths: Arc::new(RwLock::new(HashMap::new())), cycle_info: Arc::new(RwLock::new(None)), } } - /// Log scanner action with custom metadata - compatible with existing usage + // ----------------------------------------------------------------------- + // Metric recording helpers + // + // All of these are now pure sync closures. No tokio::spawn, no heap- + // allocated future, no scheduler overhead — just an atomic increment and a + // std-mutex lock for fewer than ~10 ns of work. + // ----------------------------------------------------------------------- + + /// Return a closure that records one observation of `metric` (with + /// optional caller-supplied metadata). Call it once the operation ends. pub fn log(metric: Metric) -> impl Fn(&HashMap) { - let metric = metric as usize; - let start_time = SystemTime::now(); + let metric_idx = metric as usize; + let start = SystemTime::now(); move |_custom: &HashMap| { - let duration = SystemTime::now().duration_since(start_time).unwrap_or_default(); - - // Update operation count - global_metrics().operations[metric].fetch_add(1, Ordering::Relaxed); - emit_otel_counter(metric, 1); - - // Update latency for realtime metrics (spawn async task for this) - if (metric) < Metric::LastRealtime as usize { - let metric_index = metric; - tokio::spawn(async move { - global_metrics().latency[metric_index].add(duration).await; - }); - } - - // Log trace metrics - if metric as u8 > Metric::StartTrace as u8 { - //debug!(metric = metric.as_str(), duration_ms = duration.as_millis(), "Scanner trace metric"); + let duration = SystemTime::now().duration_since(start).unwrap_or_default(); + global_metrics().operations[metric_idx].fetch_add(1, Ordering::Relaxed); + emit_otel_counter(metric_idx, 1); + if metric_idx < Metric::LastRealtime as usize { + global_metrics().latency[metric_idx].add(duration); } } } - /// Time scanner action with size - returns function that takes size + /// Return a closure that records one observation of `metric` together with + /// a byte count. Call `done(size_bytes)` when the operation ends. pub fn time_size(metric: Metric) -> impl Fn(u64) { - let metric = metric as usize; - let start_time = SystemTime::now(); + let metric_idx = metric as usize; + let start = SystemTime::now(); move |size: u64| { - let duration = SystemTime::now().duration_since(start_time).unwrap_or_default(); - - // Update operation count - global_metrics().operations[metric].fetch_add(1, Ordering::Relaxed); - emit_otel_counter(metric, 1); - - // Update latency for realtime metrics with size (spawn async task) - if (metric) < Metric::LastRealtime as usize { - let metric_index = metric; - tokio::spawn(async move { - global_metrics().latency[metric_index].add_size(duration, size).await; - }); + let duration = SystemTime::now().duration_since(start).unwrap_or_default(); + global_metrics().operations[metric_idx].fetch_add(1, Ordering::Relaxed); + emit_otel_counter(metric_idx, 1); + if metric_idx < Metric::LastRealtime as usize { + global_metrics().latency[metric_idx].add_size(duration, size); } } } - /// Time a scanner action - returns a closure to call when done + /// Return a closure that records one observation of `metric`. + /// Call `done()` when the operation ends. pub fn time(metric: Metric) -> impl Fn() { - let metric = metric as usize; - let start_time = SystemTime::now(); + let metric_idx = metric as usize; + let start = SystemTime::now(); move || { - let duration = SystemTime::now().duration_since(start_time).unwrap_or_default(); - - // Update operation count - global_metrics().operations[metric].fetch_add(1, Ordering::Relaxed); - emit_otel_counter(metric, 1); - - // Update latency for realtime metrics (spawn async task) - if (metric) < Metric::LastRealtime as usize { - let metric_index = metric; - tokio::spawn(async move { - global_metrics().latency[metric_index].add(duration).await; - }); + let duration = SystemTime::now().duration_since(start).unwrap_or_default(); + global_metrics().operations[metric_idx].fetch_add(1, Ordering::Relaxed); + emit_otel_counter(metric_idx, 1); + if metric_idx < Metric::LastRealtime as usize { + global_metrics().latency[metric_idx].add(duration); } } } - /// Time N scanner actions - returns function that takes count, then returns completion function + /// Return a two-stage closure: first call takes an item count, second call + /// (returned closure) fires when the batch of `count` operations ends. pub fn time_n(metric: Metric) -> Box Box + Send + Sync> { - let metric = metric as usize; - let start_time = SystemTime::now(); + let metric_idx = metric as usize; + let start = SystemTime::now(); Box::new(move |count: usize| { Box::new(move || { - let duration = SystemTime::now().duration_since(start_time).unwrap_or_default(); - - // Update operation count - global_metrics().operations[metric].fetch_add(count as u64, Ordering::Relaxed); - emit_otel_counter(metric, count as u64); - - // Update latency for realtime metrics (spawn async task) - if (metric) < Metric::LastRealtime as usize { - let metric_index = metric; - tokio::spawn(async move { - global_metrics().latency[metric_index].add(duration).await; - }); + let duration = SystemTime::now().duration_since(start).unwrap_or_default(); + global_metrics().operations[metric_idx].fetch_add(count as u64, Ordering::Relaxed); + emit_otel_counter(metric_idx, count as u64); + if metric_idx < Metric::LastRealtime as usize { + global_metrics().latency[metric_idx].add(duration); } }) }) } - /// Time ILM action with versions - returns function that takes versions, then returns completion function + /// Return a two-stage closure for ILM actions: first call takes a version + /// count, second call fires when the action completes. pub fn time_ilm(a: IlmAction) -> Box Box + Send + Sync> { - let a_clone = a as usize; - if a_clone == IlmAction::NoneAction as usize || a_clone >= IlmAction::ActionCount as usize { - return Box::new(move |_: u64| Box::new(move || {})); + let a_idx = a as usize; + if a_idx == IlmAction::NoneAction as usize || a_idx >= IlmAction::ActionCount as usize { + return Box::new(|_| Box::new(|| {})); } let start = SystemTime::now(); Box::new(move |versions: u64| { Box::new(move || { - let duration = SystemTime::now().duration_since(start).unwrap_or(Duration::from_secs(0)); - tokio::spawn(async move { - global_metrics().actions[a_clone].fetch_add(versions, Ordering::Relaxed); - global_metrics().actions_latency[a_clone].add(duration).await; - }); + let duration = SystemTime::now().duration_since(start).unwrap_or_default(); + global_metrics().actions[a_idx].fetch_add(versions, Ordering::Relaxed); + global_metrics().actions_latency[a_idx].add(duration); }) }) } - /// Increment time with specific duration - pub async fn inc_time(metric: Metric, duration: Duration) { - let metric = metric as usize; - // Update operation count - global_metrics().operations[metric].fetch_add(1, Ordering::Relaxed); - emit_otel_counter(metric, 1); - - // Update latency for realtime metrics - if (metric) < Metric::LastRealtime as usize { - global_metrics().latency[metric].add(duration).await; + /// Record a single observation of `metric` with a caller-supplied duration. + /// No longer async — nothing inside requires it. + pub fn inc_time(metric: Metric, duration: Duration) { + let metric_idx = metric as usize; + global_metrics().operations[metric_idx].fetch_add(1, Ordering::Relaxed); + emit_otel_counter(metric_idx, 1); + if metric_idx < Metric::LastRealtime as usize { + global_metrics().latency[metric_idx].add(duration); } } - /// Get lifetime operation count for a metric + // ----------------------------------------------------------------------- + // Read-side helpers + // ----------------------------------------------------------------------- + + /// Lifetime operation count for `metric`. pub fn lifetime(&self, metric: Metric) -> u64 { - let metric = metric as usize; - if (metric) >= Metric::Last as usize { + let idx = metric as usize; + if idx >= Metric::Last as usize { return 0; } - self.operations[metric].load(Ordering::Relaxed) + self.operations[idx].load(Ordering::Relaxed) } - /// Get last minute statistics for a metric - pub async fn last_minute(&self, metric: Metric) -> AccElem { - let metric = metric as usize; - if (metric) >= Metric::LastRealtime as usize { + /// Last-minute accumulated stats for a realtime metric. + /// No longer async — LockedLastMinuteLatency::total() is now synchronous. + pub fn last_minute(&self, metric: Metric) -> AccElem { + let idx = metric as usize; + if idx >= Metric::LastRealtime as usize { return AccElem::default(); } - self.latency[metric].total().await + self.latency[idx].total() } - /// Set current cycle information + /// Replace the current cycle record. pub async fn set_cycle(&self, cycle: Option) { *self.cycle_info.write().await = cycle; } - /// Get current cycle information + /// Read the current cycle record. pub async fn get_cycle(&self) -> Option { self.cycle_info.read().await.clone() } - /// Get current active paths + /// Snapshot of every path currently being scanned. pub async fn get_current_paths(&self) -> Vec { - let mut result = Vec::new(); let paths = self.current_paths.read().await; - + let mut result = Vec::with_capacity(paths.len()); for (disk, tracker) in paths.iter() { - let path = tracker.get_path().await; - result.push(format!("{disk}/{path}")); + result.push(format!("{disk}/{}", tracker.get_path().await)); } - result } - /// Get number of active drives + /// Number of drives with an active scan in progress. pub async fn active_drives(&self) -> usize { self.current_paths.read().await.len() } - /// Generate metrics report + /// Build a full metrics report snapshot. pub async fn report(&self) -> M_ScannerMetrics { - let mut metrics = M_ScannerMetrics::default(); + let mut m = M_ScannerMetrics::default(); - // Set cycle information if let Some(cycle) = self.get_cycle().await { - metrics.current_cycle = cycle.current; - metrics.cycles_completed_at = cycle.cycle_completed; - metrics.current_started = cycle.started; + m.current_cycle = cycle.current; + m.cycles_completed_at = cycle.cycle_completed; + m.current_started = cycle.started; } - // Replace default start time with global init time if it's the placeholder if let Some(init_time) = crate::get_global_init_time().await { - metrics.current_started = init_time; + m.current_started = init_time; } - metrics.collected_at = Utc::now(); - metrics.active_paths = self.get_current_paths().await; + m.collected_at = Utc::now(); + m.active_paths = self.get_current_paths().await; - // Lifetime operations + // Lifetime operation counts for i in 0..Metric::Last as usize { let count = self.operations[i].load(Ordering::Relaxed); if count > 0 && let Some(metric) = Metric::from_index(i) { - metrics.life_time_ops.insert(metric.as_str().to_string(), count); + m.life_time_ops.insert(metric.as_str().to_string(), count); } } - // Last minute statistics for realtime metrics + // Last-minute stats for realtime metrics — now plain sync calls for i in 0..Metric::LastRealtime as usize { - let last_min = self.latency[i].total().await; + let last_min = self.latency[i].total(); if last_min.n > 0 && let Some(metric) = Metric::from_index(i) { - metrics.last_minute.actions.insert( + m.last_minute.actions.insert( metric.as_str().to_string(), TimedAction { count: last_min.n, @@ -606,23 +615,23 @@ impl Metrics { } } - // Lifetime ILM operations + // Lifetime ILM counts for i in 0..IlmAction::ActionCount as usize { let count = self.actions[i].load(Ordering::Relaxed); if count > 0 && let Some(action) = IlmAction::from_index(i) { - metrics.life_time_ilm.insert(action.as_str().to_string(), count); + m.life_time_ilm.insert(action.as_str().to_string(), count); } } - // Last minute ILM latency + // Last-minute ILM latency — plain sync calls for i in 0..IlmAction::ActionCount as usize { - let last_min = self.actions_latency[i].total().await; + let last_min = self.actions_latency[i].total(); if last_min.n > 0 && let Some(action) = IlmAction::from_index(i) { - metrics.last_minute.ilm.insert( + m.last_minute.ilm.insert( action.as_str().to_string(), TimedAction { count: last_min.n, @@ -633,56 +642,65 @@ impl Metrics { } } - metrics + m } } -// Type aliases for compatibility with existing code -pub type UpdateCurrentPathFn = Arc Pin + Send>> + Send + Sync>; -pub type CloseDiskFn = Arc Pin + Send>> + Send + Sync>; - -/// Create a current path updater for tracking scan progress -pub fn current_path_updater(disk: &str, initial: &str) -> (UpdateCurrentPathFn, CloseDiskFn) { - let tracker = Arc::new(CurrentPathTracker::new(initial.to_string())); - let disk_name = disk.to_string(); - - // Store the tracker in global metrics - let tracker_clone = Arc::clone(&tracker); - let disk_clone = disk_name.clone(); - tokio::spawn(async move { - global_metrics().current_paths.write().await.insert(disk_clone, tracker_clone); - }); - - let update_fn = { - let tracker = Arc::clone(&tracker); - Arc::new(move |path: &str| -> Pin + Send>> { - let tracker = Arc::clone(&tracker); - let path = path.to_string(); - Box::pin(async move { - tracker.update_path(path).await; - }) - }) - }; - - let done_fn = { - let disk_name = disk_name.clone(); - Arc::new(move || -> Pin + Send>> { - let disk_name = disk_name.clone(); - Box::pin(async move { - global_metrics().current_paths.write().await.remove(&disk_name); - }) - }) - }; - - (update_fn, done_fn) -} - impl Default for Metrics { fn default() -> Self { Self::new() } } +// --------------------------------------------------------------------------- +// Path tracking helpers +// --------------------------------------------------------------------------- + +pub type UpdateCurrentPathFn = Arc Pin + Send>> + Send + Sync>; +pub type CloseDiskFn = Arc Pin + Send>> + Send + Sync>; + +/// Register a new disk in the global path tracker and return two callbacks: +/// one to update the current path and one to deregister the disk when done. +pub fn current_path_updater(disk: &str, initial: &str) -> (UpdateCurrentPathFn, CloseDiskFn) { + let tracker = Arc::new(CurrentPathTracker::new(initial.to_string())); + let disk_name = disk.to_string(); + + let tracker_clone = Arc::clone(&tracker); + let disk_insert = disk_name.clone(); + tokio::spawn(async move { + global_metrics() + .current_paths + .write() + .await + .insert(disk_insert, tracker_clone); + }); + + let update_fn: UpdateCurrentPathFn = { + let tracker = Arc::clone(&tracker); + Arc::new(move |path: &str| { + let tracker = Arc::clone(&tracker); + let path = path.to_string(); + Box::pin(async move { tracker.update_path(path).await }) + }) + }; + + let done_fn: CloseDiskFn = { + let disk = disk_name.clone(); + Arc::new(move || { + let disk = disk.clone(); + Box::pin(async move { + global_metrics().current_paths.write().await.remove(&disk); + }) + }) + }; + + (update_fn, done_fn) +} + +// --------------------------------------------------------------------------- +// CloseDiskGuard +// --------------------------------------------------------------------------- + pub struct CloseDiskGuard(CloseDiskFn); impl CloseDiskGuard { @@ -697,16 +715,10 @@ impl CloseDiskGuard { impl Drop for CloseDiskGuard { fn drop(&mut self) { - // Drop cannot be async, so we spawn the async cleanup task - // The task will run in the background and complete asynchronously if let Ok(handle) = tokio::runtime::Handle::try_current() { let close_fn = self.0.clone(); - handle.spawn(async move { - close_fn().await; - }); - } else { - // If we're not in a tokio runtime context, we can't spawn - // This is a best-effort cleanup, so we just skip it + handle.spawn(async move { close_fn().await }); } + // If there is no runtime we are in a test or shutdown path; skip cleanup. } } diff --git a/crates/ecstore/src/bucket/quota/checker.rs b/crates/ecstore/src/bucket/quota/checker.rs index 8510f52e6..6f1f05845 100644 --- a/crates/ecstore/src/bucket/quota/checker.rs +++ b/crates/ecstore/src/bucket/quota/checker.rs @@ -107,9 +107,10 @@ impl QuotaChecker { }; let duration = start_time.elapsed(); - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaCheck, duration).await; + // inc_time is now a plain fn (not async) — no .await needed. + rustfs_common::metrics::Metrics::inc_time(Metric::QuotaCheck, duration); if !allowed { - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaViolation, duration).await; + rustfs_common::metrics::Metrics::inc_time(Metric::QuotaViolation, duration); } Ok(result) @@ -146,7 +147,7 @@ impl QuotaChecker { .await .map_err(QuotaError::StorageError)?; - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()).await; + rustfs_common::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()); Ok(()) } diff --git a/crates/kms/src/api_types.rs b/crates/kms/src/api_types.rs index b769bd228..c8957ea72 100644 --- a/crates/kms/src/api_types.rs +++ b/crates/kms/src/api_types.rs @@ -122,11 +122,7 @@ pub enum ConfigureKmsRequest { )] VaultKv2(ConfigureVaultKmsRequest), /// Configure with Vault Transit backend - #[serde( - rename = "VaultTransit", - alias = "vault-transit", - alias = "vault_transit" - )] + #[serde(rename = "VaultTransit", alias = "vault-transit", alias = "vault_transit")] VaultTransit(ConfigureVaultTransitKmsRequest), } @@ -472,8 +468,7 @@ mod tests { "mount_path": "transit", "default_key_id": "rustfs-master-key" }); - let request: ConfigureKmsRequest = - serde_json::from_value(raw).expect("vault-transit request should deserialize"); + let request: ConfigureKmsRequest = serde_json::from_value(raw).expect("vault-transit request should deserialize"); let config = request.to_kms_config(); assert_eq!(config.backend, KmsBackend::VaultTransit); let vault = config.vault_transit_config().expect("vault-transit config should be present"); diff --git a/crates/s3select-api/Cargo.toml b/crates/s3select-api/Cargo.toml index b09c770bd..9890975b1 100644 --- a/crates/s3select-api/Cargo.toml +++ b/crates/s3select-api/Cargo.toml @@ -26,6 +26,7 @@ categories = ["web-programming", "development-tools", "asynchronous"] documentation = "https://docs.rs/rustfs-s3select-api/latest/rustfs_s3select_api/" [dependencies] +metrics = { workspace = true } async-trait.workspace = true bytes.workspace = true chrono.workspace = true diff --git a/crates/s3select-api/src/query/execution.rs b/crates/s3select-api/src/query/execution.rs index 865599083..1ed532cce 100644 --- a/crates/s3select-api/src/query/execution.rs +++ b/crates/s3select-api/src/query/execution.rs @@ -18,13 +18,13 @@ use std::sync::Arc; use std::task::{Context, Poll}; use std::time::{Duration, Instant}; -use parking_lot::RwLock; - use async_trait::async_trait; use datafusion::arrow::datatypes::{Schema, SchemaRef}; use datafusion::arrow::record_batch::RecordBatch; use datafusion::physical_plan::SendableRecordBatchStream; use futures::{Stream, StreamExt, TryStreamExt}; +use parking_lot::RwLock; +use tracing::debug; use crate::{QueryError, QueryResult}; @@ -32,6 +32,36 @@ use super::Query; use super::logical_planner::Plan; use super::session::SessionCtx; +pub struct PhaseTimer { + phase_name: &'static str, + start_time: Instant, +} + +impl PhaseTimer { + pub fn new(phase_name: &'static str) -> Self { + Self { + phase_name, + start_time: Instant::now(), + } + } +} + +impl Drop for PhaseTimer { + fn drop(&mut self) { + let duration = self.start_time.elapsed(); + + if !std::thread::panicking() { + metrics::histogram!( + "rustfs.s3select.phase.duration.seconds", + "phase" => self.phase_name + ) + .record(duration.as_secs_f64()); + } + + debug!("Phase '{}' took {:?}", self.phase_name, duration); + } +} + pub type QueryExecutionRef = Arc; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -94,13 +124,13 @@ impl Output { } } - /// Returns the number of records affected by the query operation - /// - /// If it is a select statement, returns the number of rows in the result set - /// - /// -1 means unknown - /// - /// panic! when StreamData's number of records greater than i64::Max + /// Returns the number of records affected by the query operation + /// + /// If it is a select statement, returns the number of rows in the result set + /// + /// -1 means unknown + /// + /// panic! when StreamData's number of records greater than i64::Max pub async fn affected_rows(self) -> i64 { self.num_rows().await as i64 } @@ -138,6 +168,20 @@ pub struct QueryStateMachine { } impl QueryStateMachine { + fn record_phase_timestamp(&self, phase: &'static str, event: &'static str) { + let elapsed = self.start.elapsed(); + metrics::histogram!( + "rustfs.s3select.phase.timestamp.seconds", + "phase" => phase, + "event" => event + ) + .record(elapsed.as_secs_f64()); + } + #[must_use] + pub fn time_phase(&self, phase_name: &'static str) -> PhaseTimer { + PhaseTimer::new(phase_name) + } + pub fn begin(query: Query, session: SessionCtx) -> Self { Self { session, @@ -148,30 +192,30 @@ impl QueryStateMachine { } pub fn begin_analyze(&self) { - // TODO record time + self.record_phase_timestamp("analyze", "start"); self.translate_to(QueryState::RUNNING(RUNNING::ANALYZING)); } pub fn end_analyze(&self) { - // TODO record time + self.record_phase_timestamp("analyze", "end"); } pub fn begin_optimize(&self) { - // TODO record time + self.record_phase_timestamp("optimize", "start"); self.translate_to(QueryState::RUNNING(RUNNING::OPTIMIZING)); } pub fn end_optimize(&self) { - // TODO + self.record_phase_timestamp("optimize", "end"); } pub fn begin_schedule(&self) { - // TODO + self.record_phase_timestamp("schedule", "start"); self.translate_to(QueryState::RUNNING(RUNNING::SCHEDULING)); } pub fn end_schedule(&self) { - // TODO + self.record_phase_timestamp("schedule", "end"); } pub fn finish(&self) { diff --git a/crates/s3select-query/src/execution/query.rs b/crates/s3select-query/src/execution/query.rs index d48cf28de..dfdb2a4c8 100644 --- a/crates/s3select-query/src/execution/query.rs +++ b/crates/s3select-query/src/execution/query.rs @@ -50,21 +50,27 @@ impl SqlQueryExecution { } async fn start(&self) -> QueryResult { - // begin optimize - self.query_state_machine.begin_optimize(); - let physical_plan = self.optimizer.optimize(&self.plan, &self.query_state_machine.session).await?; - self.query_state_machine.end_optimize(); + let physical_plan = { + // Time optimize phase - dropped at end of this block + let _optimize_timer = self.query_state_machine.time_phase("optimize"); + self.query_state_machine.begin_optimize(); + let plan = self.optimizer.optimize(&self.plan, &self.query_state_machine.session).await?; + self.query_state_machine.end_optimize(); + plan + }; - // begin schedule - self.query_state_machine.begin_schedule(); - let stream = self - .scheduler - .schedule(physical_plan.clone(), self.query_state_machine.session.inner().task_ctx()) - .await? - .stream(); - - debug!("Success build result stream."); - self.query_state_machine.end_schedule(); + let stream = { + // Time schedule phase - dropped at end of this block + let _schedule_timer = self.query_state_machine.time_phase("schedule"); + self.query_state_machine.begin_schedule(); + let stream = self + .scheduler + .schedule(physical_plan.clone(), self.query_state_machine.session.inner().task_ctx()) + .await? + .stream(); + self.query_state_machine.end_schedule(); + stream + }; Ok(Output::StreamData(stream)) }