Skip to main content

rlmesh_runtime/telemetry/
mod.rs

1//! Telemetry — windowed/session aggregation of per-op metrics.
2//!
3//! One ingest path ([`Aggregator::record`]), one shape out ([`Snapshot`]). A
4//! metric is identified by its `&'static str` name in the [`metrics`] catalog;
5//! window and session are retention horizons over the same Vitter reservoirs.
6//! The per-step hot scalar `endpoint_total_ns` is the *wire* contract, not part
7//! of this system — it is recorded here like any other Duration sample
8//! (`metrics::ENDPOINT_TOTAL`).
9
10mod reservoir;
11mod stats;
12
13use reservoir::ValueReservoir;
14use std::collections::BTreeMap;
15use std::time::Duration;
16
17/// What a metric measures — fixes its unit (ms for `Duration`, bytes for `Bytes`).
18#[derive(Clone, Copy, PartialEq, Eq, Debug)]
19pub enum Kind {
20    Duration,
21    Bytes,
22}
23
24/// A metric's identity (`name`) + `kind`. Defined only in the [`metrics`] catalog.
25#[derive(Clone, Copy, Debug, PartialEq)]
26pub struct Metric {
27    pub name: &'static str,
28    pub kind: Kind,
29}
30
31impl Metric {
32    pub const fn duration(name: &'static str) -> Self {
33        Self {
34            name,
35            kind: Kind::Duration,
36        }
37    }
38    pub const fn bytes(name: &'static str) -> Self {
39        Self {
40            name,
41            kind: Kind::Bytes,
42        }
43    }
44}
45
46/// The metric catalog — adding a metric is one const here (plus `ALL`).
47pub mod metrics {
48    use super::Metric;
49
50    pub const ENDPOINT_TOTAL: Metric = Metric::duration("endpoint.total");
51    pub const RPC_TOTAL: Metric = Metric::duration("rpc.total");
52    pub const REQUEST_BYTES: Metric = Metric::bytes("request.bytes");
53    pub const RESPONSE_BYTES: Metric = Metric::bytes("response.bytes");
54
55    /// The cardinality allowlist — derived from the catalog, not a second table.
56    pub const ALL: &[Metric] = &[ENDPOINT_TOTAL, RPC_TOTAL, REQUEST_BYTES, RESPONSE_BYTES];
57}
58
59/// Who produced a sample: an `operation` on a `component`.
60#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
61pub struct Source {
62    pub op: &'static str,
63    pub component: &'static str,
64}
65
66/// One observation — the only thing [`Aggregator::record`] ingests.
67#[derive(Clone, Copy, Debug)]
68pub struct Sample {
69    pub source: Source,
70    pub metric: Metric,
71    /// ms for a `Duration` metric; raw byte count for a `Bytes` metric.
72    pub value: f64,
73}
74
75impl Sample {
76    pub fn dur(source: Source, metric: Metric, d: Duration) -> Self {
77        Self {
78            source,
79            metric,
80            value: d.as_secs_f64() * 1e3,
81        }
82    }
83    pub fn bytes(source: Source, metric: Metric, n: u64) -> Self {
84        Self {
85            source,
86            metric,
87            value: n as f64,
88        }
89    }
90}
91
92/// Retention horizon over the same recorded samples.
93#[derive(Clone, Copy, PartialEq, Eq, Debug)]
94pub enum Horizon {
95    /// Cleared each `flush_window`.
96    Window,
97    /// Whole session.
98    Session,
99}
100
101/// One aggregated metric row. `avg`/`p*` are in the metric's unit (ms or bytes).
102#[derive(Clone, Debug, PartialEq)]
103pub struct Row {
104    pub source: Source,
105    pub metric: Metric,
106    pub count: u64,
107    pub avg: f64,
108    pub p50: f64,
109    pub p95: f64,
110    pub p99: f64,
111}
112
113/// A point-in-time aggregate for one horizon — the one shape every consumer sees.
114#[derive(Clone, Debug, PartialEq)]
115pub struct Snapshot {
116    pub horizon: Horizon,
117    pub rows: Vec<Row>,
118}
119
120struct Series {
121    source: Source,
122    metric: Metric,
123    window: ValueReservoir,
124    session: ValueReservoir,
125}
126
127impl Series {
128    fn new(source: Source, metric: Metric) -> Self {
129        Self {
130            source,
131            metric,
132            window: ValueReservoir::default(),
133            session: ValueReservoir::default(),
134        }
135    }
136
137    fn summarize(&self, horizon: Horizon) -> Option<Row> {
138        let reservoir = match horizon {
139            Horizon::Window => &self.window,
140            Horizon::Session => &self.session,
141        };
142        let stats = stats::summary(reservoir.samples())?;
143        Some(Row {
144            source: self.source,
145            metric: self.metric,
146            count: reservoir.seen(),
147            avg: stats.avg,
148            p50: stats.p50,
149            p95: stats.p95,
150            p99: stats.p99,
151        })
152    }
153}
154
155/// Windowed / session metric aggregator. One `record` in, `snapshot` out.
156#[derive(Default)]
157pub struct Aggregator {
158    series: BTreeMap<(&'static str, &'static str, &'static str), Series>,
159}
160
161impl Aggregator {
162    /// The one ingest path. Every measurement in the system goes through here.
163    pub fn record(&mut self, sample: Sample) {
164        debug_assert!(
165            metrics::ALL.iter().any(|m| m.name == sample.metric.name),
166            "unregistered metric `{}` (add it to telemetry::metrics)",
167            sample.metric.name,
168        );
169        let key = (
170            sample.source.op,
171            sample.source.component,
172            sample.metric.name,
173        );
174        let series = self
175            .series
176            .entry(key)
177            .or_insert_with(|| Series::new(sample.source, sample.metric));
178        series.window.push(sample.value);
179        series.session.push(sample.value);
180    }
181
182    /// Aggregate one horizon into a snapshot. Series with no samples this horizon
183    /// are omitted (e.g. every series right after `flush_window` for `Window`).
184    pub fn snapshot(&self, horizon: Horizon) -> Snapshot {
185        Snapshot {
186            horizon,
187            rows: self
188                .series
189                .values()
190                .filter_map(|s| s.summarize(horizon))
191                .collect(),
192        }
193    }
194
195    /// Clear the window horizon — call after emitting a window snapshot. Reuses
196    /// each reservoir's backing capacity (no per-flush realloc).
197    pub fn flush_window(&mut self) {
198        for s in self.series.values_mut() {
199            s.window.clear();
200        }
201    }
202}
203
204#[cfg(test)]
205mod tests {
206    use super::*;
207
208    fn src() -> Source {
209        Source {
210            op: "model.predict",
211            component: "model",
212        }
213    }
214
215    fn row<'a>(snap: &'a Snapshot, name: &str) -> &'a Row {
216        snap.rows
217            .iter()
218            .find(|r| r.metric.name == name)
219            .expect("row present")
220    }
221
222    #[test]
223    fn aggregates_percentiles_and_flushes_window() {
224        let mut agg = Aggregator::default();
225        for ms in [10u64, 20, 30, 40] {
226            agg.record(Sample::dur(
227                src(),
228                metrics::RPC_TOTAL,
229                Duration::from_millis(ms),
230            ));
231        }
232
233        // Window snapshot: count + ordered percentiles, durations reported in ms.
234        let w = agg.snapshot(Horizon::Window);
235        let r = row(&w, "rpc.total");
236        assert_eq!(r.count, 4);
237        assert!((r.avg - 25.0).abs() < 1e-9);
238        assert!(r.p50 <= r.p95 && r.p95 <= r.p99);
239        assert!((r.p99 - 40.0).abs() < 1e-9);
240
241        // flush clears the window, but session retains.
242        agg.flush_window();
243        assert!(agg.snapshot(Horizon::Window).rows.is_empty());
244        assert_eq!(row(&agg.snapshot(Horizon::Session), "rpc.total").count, 4);
245    }
246}