feat: add timing and metrics for query execution phases (#2419)

Signed-off-by: Shreyan Jain <shreyan11@duck.com>
Co-authored-by: Shreyan Jain <shreyan11@duck.com>
Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: loverustfs <hello@rustfs.com>
Co-authored-by: cxymds <Cxymds@qq.com>
Co-authored-by: 安正超 <anzhengchao@gmail.com>
Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com>
Co-authored-by: houseme <4829346+houseme@users.noreply.github.com>
This commit is contained in:
houseme
2026-04-07 20:32:32 +08:00
committed by GitHub
parent 9cd1d6b4d5
commit ec8f059506
7 changed files with 309 additions and 249 deletions
Generated
+1
View File
@@ -8434,6 +8434,7 @@ dependencies = [
"futures",
"futures-core",
"http 1.4.0",
"metrics",
"object_store",
"parking_lot 0.12.5",
"pin-project-lite",
+222 -210
View File
@@ -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<Metrics> {
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<Mutex<LastMinuteLatency>>,
}
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<RwLock<String>>,
}
@@ -290,17 +330,16 @@ impl CurrentPathTracker {
}
}
/// Main scanner metrics structure
// ---------------------------------------------------------------------------
// Metrics
// ---------------------------------------------------------------------------
pub struct Metrics {
// All fields must be accessed atomically and aligned.
operations: Vec<AtomicU64>,
latency: Vec<LockedLastMinuteLatency>,
actions: Vec<AtomicU64>,
actions_latency: Vec<LockedLastMinuteLatency>,
// Current paths contains disk -> tracker mappings
current_paths: Arc<RwLock<HashMap<String, Arc<CurrentPathTracker>>>>,
// Cycle information
cycle_info: Arc<RwLock<Option<CurrentCycle>>>,
}
@@ -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<String, String>) {
let metric = metric as usize;
let start_time = SystemTime::now();
let metric_idx = metric as usize;
let start = SystemTime::now();
move |_custom: &HashMap<String, String>| {
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<dyn Fn(usize) -> Box<dyn Fn() + Send + Sync> + 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<dyn Fn(u64) -> Box<dyn Fn() + Send + Sync> + 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<CurrentCycle>) {
*self.cycle_info.write().await = cycle;
}
/// Get current cycle information
/// Read the current cycle record.
pub async fn get_cycle(&self) -> Option<CurrentCycle> {
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<String> {
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<dyn Fn(&str) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>;
pub type CloseDiskFn = Arc<dyn Fn() -> Pin<Box<dyn Future<Output = ()> + 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<Box<dyn Future<Output = ()> + 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<Box<dyn Future<Output = ()> + 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<dyn Fn(&str) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>;
pub type CloseDiskFn = Arc<dyn Fn() -> Pin<Box<dyn Future<Output = ()> + 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.
}
}
+4 -3
View File
@@ -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(())
}
+2 -7
View File
@@ -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");
+1
View File
@@ -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
+59 -15
View File
@@ -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<dyn QueryExecution>;
#[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) {
+20 -14
View File
@@ -50,21 +50,27 @@ impl SqlQueryExecution {
}
async fn start(&self) -> QueryResult<Output> {
// 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))
}