1use 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#[derive(Debug, Clone, Default)]
33pub struct OperatorProfile {
34 pub rows: u64,
35 pub elapsed_ns: u64,
36 pub next_calls: u64,
37}
38
39#[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
71pub 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
94pub(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}