rlmesh_runtime/telemetry/
mod.rs1mod reservoir;
11mod stats;
12
13use reservoir::ValueReservoir;
14use std::collections::BTreeMap;
15use std::time::Duration;
16
17#[derive(Clone, Copy, PartialEq, Eq, Debug)]
20pub enum Kind {
21 Duration,
22 Bytes,
23 Count,
24}
25
26#[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
54pub 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 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 pub const ENDPOINT_QUEUE: Metric = Metric::duration("endpoint.queue");
73 pub const PREDICT_IN_FLIGHT: Metric = Metric::count("predict.in_flight");
77 pub const PREDICT_ADAPTER: Metric = Metric::duration("predict.adapter");
82 pub const HELD_EPISODES: Metric = Metric::count("held.episodes");
88 pub const HELD_BYTES: Metric = Metric::bytes("held.bytes");
90 pub const GROUP_SIZE: Metric = Metric::count("group.size");
93 pub const LANE_SKEW: Metric = Metric::duration("lane.skew");
99 pub const HISTORY_ROWS: Metric = Metric::count("history.rows");
104
105 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#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
127pub struct Source {
128 pub op: &'static str,
129 pub component: &'static str,
130}
131
132#[derive(Clone, Copy, Debug)]
134pub struct Sample {
135 pub source: Source,
136 pub metric: Metric,
137 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#[derive(Clone, Copy, PartialEq, Eq, Debug)]
167pub enum Horizon {
168 Window,
170 Session,
172}
173
174#[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#[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#[derive(Default)]
230pub struct Aggregator {
231 series: BTreeMap<(&'static str, &'static str, &'static str), Series>,
232}
233
234impl Aggregator {
235 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 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 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 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 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}