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)]
19pub enum Kind {
20 Duration,
21 Bytes,
22}
23
24#[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
46pub 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 pub const ALL: &[Metric] = &[ENDPOINT_TOTAL, RPC_TOTAL, REQUEST_BYTES, RESPONSE_BYTES];
57}
58
59#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
61pub struct Source {
62 pub op: &'static str,
63 pub component: &'static str,
64}
65
66#[derive(Clone, Copy, Debug)]
68pub struct Sample {
69 pub source: Source,
70 pub metric: Metric,
71 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#[derive(Clone, Copy, PartialEq, Eq, Debug)]
94pub enum Horizon {
95 Window,
97 Session,
99}
100
101#[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#[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#[derive(Default)]
157pub struct Aggregator {
158 series: BTreeMap<(&'static str, &'static str, &'static str), Series>,
159}
160
161impl Aggregator {
162 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 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 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 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 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}