use std::{
cell::RefCell,
sync::{
Arc,
atomic::{AtomicU64, AtomicUsize, Ordering},
},
};
use serde::{Deserialize, Serialize};
use super::cpu;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct OpStats {
pub fts_postings_bytes: u64,
pub vector_cells_scanned: u64,
pub vector_candidates_scanned: u64,
pub vector_rows_reranked: u64,
pub planned_read_ranges: u64,
pub sql_page_bytes: u64,
pub rows_materialized: u64,
pub kernel_cpu_ns: u64,
}
#[derive(Debug, Default)]
pub(crate) struct OpStatsCollector {
fts_postings_bytes: AtomicU64,
vector_cells_scanned: AtomicU64,
vector_candidates_scanned: AtomicU64,
vector_rows_reranked: AtomicU64,
planned_read_ranges: AtomicU64,
sql_page_bytes: AtomicU64,
rows_materialized: AtomicU64,
kernel_cpu_ns: AtomicU64,
}
impl OpStatsCollector {
pub(crate) fn add_fts_postings_bytes(&self, bytes: u64) {
self.fts_postings_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub(crate) fn add_vector_scan(&self, cells: u64, candidates: u64) {
self.vector_cells_scanned
.fetch_add(cells, Ordering::Relaxed);
self.vector_candidates_scanned
.fetch_add(candidates, Ordering::Relaxed);
}
pub(crate) fn add_vector_rows_reranked(&self, rows: u64) {
self.vector_rows_reranked.fetch_add(rows, Ordering::Relaxed);
}
pub(crate) fn add_planned_read_ranges(&self, ranges: u64) {
self.planned_read_ranges
.fetch_add(ranges, Ordering::Relaxed);
}
pub(crate) fn add_sql_page_bytes(&self, bytes: u64) {
self.sql_page_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub(crate) fn add_rows_materialized(&self, rows: u64) {
self.rows_materialized.fetch_add(rows, Ordering::Relaxed);
}
pub(crate) fn add_kernel_cpu_ns(&self, ns: u64) {
self.kernel_cpu_ns.fetch_add(ns, Ordering::Relaxed);
}
pub(crate) fn snapshot(&self) -> OpStats {
OpStats {
fts_postings_bytes: self.fts_postings_bytes.load(Ordering::Relaxed),
vector_cells_scanned: self.vector_cells_scanned.load(Ordering::Relaxed),
vector_candidates_scanned: self.vector_candidates_scanned.load(Ordering::Relaxed),
vector_rows_reranked: self.vector_rows_reranked.load(Ordering::Relaxed),
planned_read_ranges: self.planned_read_ranges.load(Ordering::Relaxed),
sql_page_bytes: self.sql_page_bytes.load(Ordering::Relaxed),
rows_materialized: self.rows_materialized.load(Ordering::Relaxed),
kernel_cpu_ns: self.kernel_cpu_ns.load(Ordering::Relaxed),
}
}
}
thread_local! {
static CURRENT: RefCell<Option<Arc<OpStatsCollector>>> = const { RefCell::new(None) };
}
static ACTIVE_SCOPES: AtomicUsize = AtomicUsize::new(0);
pub(crate) fn metering_active() -> bool {
ACTIVE_SCOPES.load(Ordering::Relaxed) > 0
}
struct ActiveScopeGuard;
impl Drop for ActiveScopeGuard {
fn drop(&mut self) {
ACTIVE_SCOPES.fetch_sub(1, Ordering::Relaxed);
}
}
struct ScopeGuard {
previous: Option<Arc<OpStatsCollector>>,
}
impl Drop for ScopeGuard {
fn drop(&mut self) {
CURRENT.with(|slot| {
*slot.borrow_mut() = self.previous.take();
});
}
}
pub fn with_op_stats<T>(f: impl FnOnce() -> T) -> (T, OpStats) {
let collector = Arc::new(OpStatsCollector::default());
let previous = CURRENT.with(|slot| slot.borrow_mut().replace(Arc::clone(&collector)));
let _guard = ScopeGuard { previous };
ACTIVE_SCOPES.fetch_add(1, Ordering::Relaxed);
let _active = ActiveScopeGuard;
let value = f();
(value, collector.snapshot())
}
pub(crate) fn current() -> Option<Arc<OpStatsCollector>> {
CURRENT.with(|slot| slot.borrow().clone())
}
pub(crate) fn timed_kernel<T>(
collector: &Option<Arc<OpStatsCollector>>,
f: impl FnOnce() -> T,
) -> T {
if !metering_active() {
return f();
}
let Some(stats) = collector else {
return f();
};
let start = cpu::thread_cpu_ns();
let value = f();
stats.add_kernel_cpu_ns(cpu::thread_cpu_delta_ns(start));
value
}
pub(crate) fn timed_section<R>(f: impl FnOnce() -> R) -> (R, u64) {
let start = metering_active().then(cpu::thread_cpu_ns).flatten();
let out = f();
(out, cpu::thread_cpu_delta_ns(start))
}
pub(crate) fn suppressed<T>(f: impl FnOnce() -> T) -> T {
let previous = CURRENT.with(|slot| slot.borrow_mut().take());
let _guard = ScopeGuard { previous };
f()
}
#[cfg(test)]
mod tests {
use std::panic::catch_unwind;
use super::*;
#[test]
fn scope_collects_and_clears() {
assert!(current().is_none(), "no collector outside a scope");
let (value, stats) = with_op_stats(|| {
let collector = current().expect("collector installed inside the scope");
collector.add_fts_postings_bytes(123);
7u32
});
assert_eq!(value, 7);
assert_eq!(stats.fts_postings_bytes, 123);
assert!(current().is_none(), "scope uninstalls its collector");
}
#[test]
fn nested_scopes_shadow_and_restore() {
let (_, outer) = with_op_stats(|| {
let outer_collector = current().expect("outer collector");
outer_collector.add_fts_postings_bytes(1);
let (_, inner) = with_op_stats(|| {
current()
.expect("inner collector shadows outer")
.add_fts_postings_bytes(10);
});
assert_eq!(inner.fts_postings_bytes, 10);
current()
.expect("outer collector restored")
.add_fts_postings_bytes(2);
});
assert_eq!(
outer.fts_postings_bytes, 3,
"outer scope never absorbs the inner query's counters"
);
}
#[test]
fn a_panicking_scope_still_restores_the_previous_collector() {
let (_, outer) = with_op_stats(|| {
let result = catch_unwind(|| {
let (_, _) = with_op_stats(|| panic!("kernel failure"));
});
assert!(result.is_err());
current()
.expect("outer collector survives the inner panic")
.add_fts_postings_bytes(5);
});
assert_eq!(outer.fts_postings_bytes, 5);
}
}