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};
23use std::time::Instant;
24
25use lora_compiler::physical::PhysicalNodeId;
26
27use crate::errors::ExecResult;
28use crate::pull::RowSource;
29use crate::value::Row;
30
31/// Per-operator metrics collected during a profile run.
32#[derive(Debug, Clone, Default)]
33pub struct OperatorProfile {
34    pub rows: u64,
35    pub elapsed_ns: u64,
36    pub next_calls: u64,
37}
38
39/// Shared accumulator. Pass an `Arc` clone to each `MeteredSource` so
40/// every operator writes into the same map.
41#[derive(Debug, Default)]
42pub struct MetricsCollector {
43    inner: Mutex<BTreeMap<PhysicalNodeId, OperatorProfile>>,
44}
45
46impl MetricsCollector {
47    pub fn new() -> Self {
48        Self::default()
49    }
50
51    pub fn record(&self, op: PhysicalNodeId, elapsed_ns: u64, produced_row: bool) {
52        if let Ok(mut map) = self.inner.lock() {
53            let entry = map.entry(op).or_default();
54            entry.next_calls += 1;
55            entry.elapsed_ns = entry.elapsed_ns.saturating_add(elapsed_ns);
56            if produced_row {
57                entry.rows += 1;
58            }
59        }
60    }
61
62    pub fn snapshot(&self) -> BTreeMap<PhysicalNodeId, OperatorProfile> {
63        self.inner.lock().map(|m| m.clone()).unwrap_or_default()
64    }
65}
66
67thread_local! {
68    static CURRENT: RefCell<Option<Arc<MetricsCollector>>> = const { RefCell::new(None) };
69}
70
71/// RAII guard. While alive, the thread-local current collector is set
72/// to the supplied collector; on drop it is cleared.
73pub struct CollectorGuard {
74    _private: (),
75}
76
77impl CollectorGuard {
78    pub fn install(collector: Arc<MetricsCollector>) -> Self {
79        CURRENT.with(|cell| {
80            *cell.borrow_mut() = Some(collector);
81        });
82        Self { _private: () }
83    }
84}
85
86impl Drop for CollectorGuard {
87    fn drop(&mut self) {
88        CURRENT.with(|cell| {
89            *cell.borrow_mut() = None;
90        });
91    }
92}
93
94/// Wrap `inner` with timing instrumentation if a collector is currently
95/// installed; otherwise return `inner` unchanged.
96pub(crate) fn wrap_metered<'a>(
97    op_id: PhysicalNodeId,
98    inner: Box<dyn RowSource + 'a>,
99) -> Box<dyn RowSource + 'a> {
100    let collector = CURRENT.with(|cell| cell.borrow().clone());
101    match collector {
102        Some(c) => Box::new(MeteredSource {
103            inner,
104            op_id,
105            collector: c,
106        }),
107        None => inner,
108    }
109}
110
111struct MeteredSource<'a> {
112    inner: Box<dyn RowSource + 'a>,
113    op_id: PhysicalNodeId,
114    collector: Arc<MetricsCollector>,
115}
116
117impl<'a> RowSource for MeteredSource<'a> {
118    fn next_row(&mut self) -> ExecResult<Option<Row>> {
119        let t0 = Instant::now();
120        let result = self.inner.next_row();
121        let elapsed_ns = t0.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64;
122        let produced = matches!(&result, Ok(Some(_)));
123        self.collector.record(self.op_id, elapsed_ns, produced);
124        result
125    }
126}