Skip to main content

lora_executor/
profile.rs

1//! Profile-mode instrumentation for the pull executor.
2//!
3//! When a `MetricsCollector` is installed via [`install_collector`], the
4//! pull builder wraps each operator's [`RowSource`] in a [`MeteredSource`]
5//! that records, for that operator id:
6//!  - `next_calls`: number of `next_row` invocations.
7//!  - `rows`: number of `Some(_)` rows produced.
8//!  - `elapsed_ns`: cumulative wall-clock time spent inside `next_row`,
9//!    *inclusive* of time spent pulling from upstream operators. This
10//!    matches the "operator + descendants" view that is most actionable
11//!    when reading a profile.
12//!
13//! The collector is installed for the lifetime of one query through a
14//! thread-local guard; the executor is single-threaded per-query, so this
15//! is safe and avoids threading an extra parameter through every source
16//! constructor. Outside `profile()` calls the thread-local is `None` and
17//! `wrap_metered` returns its argument unchanged — `execute()` and
18//! streaming reads pay nothing.
19
20use std::cell::RefCell;
21use std::collections::BTreeMap;
22use std::sync::{Arc, Mutex};
23
24use web_time::Instant;
25
26use lora_compiler::physical::PhysicalNodeId;
27
28use crate::errors::ExecResult;
29use crate::pull::RowSource;
30use crate::value::Row;
31
32/// Per-operator metrics collected during a profile run.
33#[derive(Debug, Clone, Default)]
34pub struct OperatorProfile {
35    pub rows: u64,
36    pub elapsed_ns: u64,
37    pub next_calls: u64,
38}
39
40/// Shared accumulator. Pass an `Arc` clone to each `MeteredSource` so
41/// every operator writes into the same map.
42#[derive(Debug, Default)]
43pub struct MetricsCollector {
44    inner: Mutex<BTreeMap<PhysicalNodeId, OperatorProfile>>,
45}
46
47impl MetricsCollector {
48    pub fn new() -> Self {
49        Self::default()
50    }
51
52    pub fn record(&self, op: PhysicalNodeId, elapsed_ns: u64, produced_row: bool) {
53        if let Ok(mut map) = self.inner.lock() {
54            let entry = map.entry(op).or_default();
55            entry.next_calls += 1;
56            entry.elapsed_ns = entry.elapsed_ns.saturating_add(elapsed_ns);
57            if produced_row {
58                entry.rows += 1;
59            }
60        }
61    }
62
63    pub fn snapshot(&self) -> BTreeMap<PhysicalNodeId, OperatorProfile> {
64        self.inner.lock().map(|m| m.clone()).unwrap_or_default()
65    }
66}
67
68thread_local! {
69    static CURRENT: RefCell<Option<Arc<MetricsCollector>>> = const { RefCell::new(None) };
70}
71
72/// RAII guard. While alive, the thread-local current collector is set
73/// to the supplied collector; on drop it is cleared.
74pub struct CollectorGuard {
75    _private: (),
76}
77
78impl CollectorGuard {
79    pub fn install(collector: Arc<MetricsCollector>) -> Self {
80        CURRENT.with(|cell| {
81            *cell.borrow_mut() = Some(collector);
82        });
83        Self { _private: () }
84    }
85}
86
87impl Drop for CollectorGuard {
88    fn drop(&mut self) {
89        CURRENT.with(|cell| {
90            *cell.borrow_mut() = None;
91        });
92    }
93}
94
95/// Wrap `inner` with timing instrumentation if a collector is currently
96/// installed; otherwise return `inner` unchanged.
97pub(crate) fn wrap_metered<'a>(
98    op_id: PhysicalNodeId,
99    inner: Box<dyn RowSource + 'a>,
100) -> Box<dyn RowSource + 'a> {
101    let collector = CURRENT.with(|cell| cell.borrow().clone());
102    match collector {
103        Some(c) => Box::new(MeteredSource {
104            inner,
105            op_id,
106            collector: c,
107        }),
108        None => inner,
109    }
110}
111
112struct MeteredSource<'a> {
113    inner: Box<dyn RowSource + 'a>,
114    op_id: PhysicalNodeId,
115    collector: Arc<MetricsCollector>,
116}
117
118impl<'a> RowSource for MeteredSource<'a> {
119    fn next_row(&mut self) -> ExecResult<Option<Row>> {
120        let t0 = Instant::now();
121        let result = self.inner.next_row();
122        let elapsed_ns = t0.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64;
123        let produced = matches!(&result, Ok(Some(_)));
124        self.collector.record(self.op_id, elapsed_ns, produced);
125        result
126    }
127}