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
18/// `Bytes`, a dimensionless count for `Count`).
19#[derive(Clone, Copy, PartialEq, Eq, Debug)]
20pub enum Kind {
21    Duration,
22    Bytes,
23    Count,
24}
25
26/// A metric's identity (`name`) + `kind`. Defined only in the [`metrics`] catalog.
27#[derive(Clone, Copy, Debug, PartialEq)]
28pub struct Metric {
29    pub name: &'static str,
30    pub kind: Kind,
31}
32
33impl Metric {
34    pub const fn duration(name: &'static str) -> Self {
35        Self {
36            name,
37            kind: Kind::Duration,
38        }
39    }
40    pub const fn bytes(name: &'static str) -> Self {
41        Self {
42            name,
43            kind: Kind::Bytes,
44        }
45    }
46    pub const fn count(name: &'static str) -> Self {
47        Self {
48            name,
49            kind: Kind::Count,
50        }
51    }
52}
53
54/// The metric catalog — adding a metric is one const here (plus `ALL`).
55pub mod metrics {
56    use super::Metric;
57
58    pub const ENDPOINT_TOTAL: Metric = Metric::duration("endpoint.total");
59    pub const RPC_TOTAL: Metric = Metric::duration("rpc.total");
60    pub const REQUEST_BYTES: Metric = Metric::bytes("request.bytes");
61    pub const RESPONSE_BYTES: Metric = Metric::bytes("response.bytes");
62    /// The split of [`ENDPOINT_TOTAL`] a peer reports: wire decode of the
63    /// request, the env/model implementation's own work, wire encode of the
64    /// response. Recorded only for a peer that stamps them, so their absence
65    /// means an older peer, not a zero cost.
66    pub const ENDPOINT_DECODE: Metric = Metric::duration("endpoint.decode");
67    pub const ENDPOINT_USER: Metric = Metric::duration("endpoint.user");
68    pub const ENDPOINT_ENCODE: Metric = Metric::duration("endpoint.encode");
69    /// Wait at the endpoint before the op's handler ran — at a model the
70    /// concurrency permit, route gate, and handler lock; at an env the env
71    /// lock — the queueing cost a client cannot otherwise tell from a slow op.
72    pub const ENDPOINT_QUEUE: Metric = Metric::duration("endpoint.queue");
73    /// Requests holding a slot at the model endpoint when this predict was
74    /// dispatched to its handler (>= 1): slot occupancy, not parallelism — the
75    /// handler lock runs them one at a time.
76    pub const PREDICT_IN_FLIGHT: Metric = Metric::count("predict.in_flight");
77    /// Adapter work inside the predict handler's own time (obs assembly +
78    /// action apply) — a sub-span of `endpoint.user`, so the model's own
79    /// forward is `endpoint.user` minus this. Recorded only for a handler
80    /// that measures it (the adapter engine); zero for a spec-less route.
81    pub const PREDICT_ADAPTER: Metric = Metric::duration("predict.adapter");
82    /// Episodes whose frame-stack windows the model endpoint's adapter engine
83    /// held when the predict was stamped, across all its routes — the state
84    /// that grows with concurrent episodes and shrinks on episode-end GC. A
85    /// measured zero is recorded (that is the eviction working); a handler
86    /// that keeps no such accounting records nothing.
87    pub const HELD_EPISODES: Metric = Metric::count("held.episodes");
88    /// Bytes those held frame-stack windows occupy.
89    pub const HELD_BYTES: Metric = Metric::bytes("held.bytes");
90    /// Lanes fused into the model forward this predict rode in (1 = unfused;
91    /// recorded only when the transport reports it — grouped/coalesced predicts).
92    pub const GROUP_SIZE: Metric = Metric::count("group.size");
93    /// How much longer a vector env's slowest lane took than its median lane on
94    /// this op — the straggler cost the whole batch pays, which the `env.step`
95    /// aggregate only blurs. One dispersion sample per op, never a series per
96    /// lane; recorded only for an env that times its own lanes, including its
97    /// zeros, so the percentiles say how often a straggler appears.
98    pub const LANE_SKEW: Metric = Metric::duration("lane.skew");
99    /// Replayed env steps delivered as observation-history rows on a predict
100    /// (a route that negotiated history at resolve): one sample per predict
101    /// that carried any, so the series says how much history each re-plan
102    /// hauls across the wire.
103    pub const HISTORY_ROWS: Metric = Metric::count("history.rows");
104
105    /// The cardinality allowlist — derived from the catalog, not a second table.
106    pub const ALL: &[Metric] = &[
107        ENDPOINT_TOTAL,
108        ENDPOINT_DECODE,
109        ENDPOINT_USER,
110        ENDPOINT_ENCODE,
111        ENDPOINT_QUEUE,
112        PREDICT_IN_FLIGHT,
113        PREDICT_ADAPTER,
114        HELD_EPISODES,
115        HELD_BYTES,
116        RPC_TOTAL,
117        REQUEST_BYTES,
118        RESPONSE_BYTES,
119        GROUP_SIZE,
120        LANE_SKEW,
121        HISTORY_ROWS,
122    ];
123}
124
125/// Who produced a sample: an `operation` on a `component`.
126#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
127pub struct Source {
128    pub op: &'static str,
129    pub component: &'static str,
130}
131
132/// One observation — the only thing [`Aggregator::record`] ingests.
133#[derive(Clone, Copy, Debug)]
134pub struct Sample {
135    pub source: Source,
136    pub metric: Metric,
137    /// ms for a `Duration` metric; raw byte count for a `Bytes` metric.
138    pub value: f64,
139}
140
141impl Sample {
142    pub fn dur(source: Source, metric: Metric, d: Duration) -> Self {
143        Self {
144            source,
145            metric,
146            value: d.as_secs_f64() * 1e3,
147        }
148    }
149    pub fn bytes(source: Source, metric: Metric, n: u64) -> Self {
150        Self {
151            source,
152            metric,
153            value: n as f64,
154        }
155    }
156    pub fn count(source: Source, metric: Metric, n: u64) -> Self {
157        Self {
158            source,
159            metric,
160            value: n as f64,
161        }
162    }
163}
164
165/// Retention horizon over the same recorded samples.
166#[derive(Clone, Copy, PartialEq, Eq, Debug)]
167pub enum Horizon {
168    /// Cleared each `flush_window`.
169    Window,
170    /// Whole session.
171    Session,
172}
173
174/// One aggregated metric row. `avg`/`p*` are in the metric's unit (ms or bytes).
175#[derive(Clone, Debug, PartialEq)]
176pub struct Row {
177    pub source: Source,
178    pub metric: Metric,
179    pub count: u64,
180    pub avg: f64,
181    pub p50: f64,
182    pub p95: f64,
183    pub p99: f64,
184}
185
186/// A point-in-time aggregate for one horizon — the one shape every consumer sees.
187#[derive(Clone, Debug, PartialEq)]
188pub struct Snapshot {
189    pub horizon: Horizon,
190    pub rows: Vec<Row>,
191}
192
193struct Series {
194    source: Source,
195    metric: Metric,
196    window: ValueReservoir,
197    session: ValueReservoir,
198}
199
200impl Series {
201    fn new(source: Source, metric: Metric) -> Self {
202        Self {
203            source,
204            metric,
205            window: ValueReservoir::default(),
206            session: ValueReservoir::default(),
207        }
208    }
209
210    fn summarize(&self, horizon: Horizon) -> Option<Row> {
211        let reservoir = match horizon {
212            Horizon::Window => &self.window,
213            Horizon::Session => &self.session,
214        };
215        let stats = stats::summary(reservoir.samples())?;
216        Some(Row {
217            source: self.source,
218            metric: self.metric,
219            count: reservoir.seen(),
220            avg: stats.avg,
221            p50: stats.p50,
222            p95: stats.p95,
223            p99: stats.p99,
224        })
225    }
226}
227
228/// Windowed / session metric aggregator. One `record` in, `snapshot` out.
229#[derive(Default)]
230pub struct Aggregator {
231    series: BTreeMap<(&'static str, &'static str, &'static str), Series>,
232}
233
234impl Aggregator {
235    /// The one ingest path. Every measurement in the system goes through here.
236    pub fn record(&mut self, sample: Sample) {
237        debug_assert!(
238            metrics::ALL.iter().any(|m| m.name == sample.metric.name),
239            "unregistered metric `{}` (add it to telemetry::metrics)",
240            sample.metric.name,
241        );
242        let key = (
243            sample.source.op,
244            sample.source.component,
245            sample.metric.name,
246        );
247        let series = self
248            .series
249            .entry(key)
250            .or_insert_with(|| Series::new(sample.source, sample.metric));
251        series.window.push(sample.value);
252        series.session.push(sample.value);
253    }
254
255    /// Aggregate one horizon into a snapshot. Series with no samples this horizon
256    /// are omitted (e.g. every series right after `flush_window` for `Window`).
257    pub fn snapshot(&self, horizon: Horizon) -> Snapshot {
258        Snapshot {
259            horizon,
260            rows: self
261                .series
262                .values()
263                .filter_map(|s| s.summarize(horizon))
264                .collect(),
265        }
266    }
267
268    /// Clear the window horizon — call after emitting a window snapshot. Reuses
269    /// each reservoir's backing capacity (no per-flush realloc).
270    pub fn flush_window(&mut self) {
271        for s in self.series.values_mut() {
272            s.window.clear();
273        }
274    }
275}
276
277#[cfg(test)]
278mod tests {
279    use super::*;
280
281    fn src() -> Source {
282        Source {
283            op: "model.predict",
284            component: "model",
285        }
286    }
287
288    fn row<'a>(snap: &'a Snapshot, name: &str) -> &'a Row {
289        snap.rows
290            .iter()
291            .find(|r| r.metric.name == name)
292            .expect("row present")
293    }
294
295    #[test]
296    fn aggregates_percentiles_and_flushes_window() {
297        let mut agg = Aggregator::default();
298        for ms in [10u64, 20, 30, 40] {
299            agg.record(Sample::dur(
300                src(),
301                metrics::RPC_TOTAL,
302                Duration::from_millis(ms),
303            ));
304        }
305
306        // Window snapshot: count + ordered percentiles, durations reported in ms.
307        let w = agg.snapshot(Horizon::Window);
308        let r = row(&w, "rpc.total");
309        assert_eq!(r.count, 4);
310        assert!((r.avg - 25.0).abs() < 1e-9);
311        assert!(r.p50 <= r.p95 && r.p95 <= r.p99);
312        assert!((r.p99 - 40.0).abs() < 1e-9);
313
314        // flush clears the window, but session retains.
315        agg.flush_window();
316        assert!(agg.snapshot(Horizon::Window).rows.is_empty());
317        assert_eq!(row(&agg.snapshot(Horizon::Session), "rpc.total").count, 4);
318    }
319}