mod reservoir;
mod stats;
use reservoir::ValueReservoir;
use std::collections::BTreeMap;
use std::time::Duration;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Kind {
Duration,
Bytes,
Count,
}
#[derive(Clone, Copy, Debug, PartialEq)]
pub struct Metric {
pub name: &'static str,
pub kind: Kind,
}
impl Metric {
pub const fn duration(name: &'static str) -> Self {
Self {
name,
kind: Kind::Duration,
}
}
pub const fn bytes(name: &'static str) -> Self {
Self {
name,
kind: Kind::Bytes,
}
}
pub const fn count(name: &'static str) -> Self {
Self {
name,
kind: Kind::Count,
}
}
}
pub mod metrics {
use super::Metric;
pub const ENDPOINT_TOTAL: Metric = Metric::duration("endpoint.total");
pub const RPC_TOTAL: Metric = Metric::duration("rpc.total");
pub const REQUEST_BYTES: Metric = Metric::bytes("request.bytes");
pub const RESPONSE_BYTES: Metric = Metric::bytes("response.bytes");
pub const ENDPOINT_DECODE: Metric = Metric::duration("endpoint.decode");
pub const ENDPOINT_USER: Metric = Metric::duration("endpoint.user");
pub const ENDPOINT_ENCODE: Metric = Metric::duration("endpoint.encode");
pub const ENDPOINT_QUEUE: Metric = Metric::duration("endpoint.queue");
pub const PREDICT_IN_FLIGHT: Metric = Metric::count("predict.in_flight");
pub const PREDICT_ADAPTER: Metric = Metric::duration("predict.adapter");
pub const HELD_EPISODES: Metric = Metric::count("held.episodes");
pub const HELD_BYTES: Metric = Metric::bytes("held.bytes");
pub const GROUP_SIZE: Metric = Metric::count("group.size");
pub const LANE_SKEW: Metric = Metric::duration("lane.skew");
pub const HISTORY_ROWS: Metric = Metric::count("history.rows");
pub const ALL: &[Metric] = &[
ENDPOINT_TOTAL,
ENDPOINT_DECODE,
ENDPOINT_USER,
ENDPOINT_ENCODE,
ENDPOINT_QUEUE,
PREDICT_IN_FLIGHT,
PREDICT_ADAPTER,
HELD_EPISODES,
HELD_BYTES,
RPC_TOTAL,
REQUEST_BYTES,
RESPONSE_BYTES,
GROUP_SIZE,
LANE_SKEW,
HISTORY_ROWS,
];
}
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct Source {
pub op: &'static str,
pub component: &'static str,
}
#[derive(Clone, Copy, Debug)]
pub struct Sample {
pub source: Source,
pub metric: Metric,
pub value: f64,
}
impl Sample {
pub fn dur(source: Source, metric: Metric, d: Duration) -> Self {
Self {
source,
metric,
value: d.as_secs_f64() * 1e3,
}
}
pub fn bytes(source: Source, metric: Metric, n: u64) -> Self {
Self {
source,
metric,
value: n as f64,
}
}
pub fn count(source: Source, metric: Metric, n: u64) -> Self {
Self {
source,
metric,
value: n as f64,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Horizon {
Window,
Session,
}
#[derive(Clone, Debug, PartialEq)]
pub struct Row {
pub source: Source,
pub metric: Metric,
pub count: u64,
pub avg: f64,
pub p50: f64,
pub p95: f64,
pub p99: f64,
}
#[derive(Clone, Debug, PartialEq)]
pub struct Snapshot {
pub horizon: Horizon,
pub rows: Vec<Row>,
}
struct Series {
source: Source,
metric: Metric,
window: ValueReservoir,
session: ValueReservoir,
}
impl Series {
fn new(source: Source, metric: Metric) -> Self {
Self {
source,
metric,
window: ValueReservoir::default(),
session: ValueReservoir::default(),
}
}
fn summarize(&self, horizon: Horizon) -> Option<Row> {
let reservoir = match horizon {
Horizon::Window => &self.window,
Horizon::Session => &self.session,
};
let stats = stats::summary(reservoir.samples())?;
Some(Row {
source: self.source,
metric: self.metric,
count: reservoir.seen(),
avg: stats.avg,
p50: stats.p50,
p95: stats.p95,
p99: stats.p99,
})
}
}
#[derive(Default)]
pub struct Aggregator {
series: BTreeMap<(&'static str, &'static str, &'static str), Series>,
}
impl Aggregator {
pub fn record(&mut self, sample: Sample) {
debug_assert!(
metrics::ALL.iter().any(|m| m.name == sample.metric.name),
"unregistered metric `{}` (add it to telemetry::metrics)",
sample.metric.name,
);
let key = (
sample.source.op,
sample.source.component,
sample.metric.name,
);
let series = self
.series
.entry(key)
.or_insert_with(|| Series::new(sample.source, sample.metric));
series.window.push(sample.value);
series.session.push(sample.value);
}
pub fn snapshot(&self, horizon: Horizon) -> Snapshot {
Snapshot {
horizon,
rows: self
.series
.values()
.filter_map(|s| s.summarize(horizon))
.collect(),
}
}
pub fn flush_window(&mut self) {
for s in self.series.values_mut() {
s.window.clear();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn src() -> Source {
Source {
op: "model.predict",
component: "model",
}
}
fn row<'a>(snap: &'a Snapshot, name: &str) -> &'a Row {
snap.rows
.iter()
.find(|r| r.metric.name == name)
.expect("row present")
}
#[test]
fn aggregates_percentiles_and_flushes_window() {
let mut agg = Aggregator::default();
for ms in [10u64, 20, 30, 40] {
agg.record(Sample::dur(
src(),
metrics::RPC_TOTAL,
Duration::from_millis(ms),
));
}
let w = agg.snapshot(Horizon::Window);
let r = row(&w, "rpc.total");
assert_eq!(r.count, 4);
assert!((r.avg - 25.0).abs() < 1e-9);
assert!(r.p50 <= r.p95 && r.p95 <= r.p99);
assert!((r.p99 - 40.0).abs() < 1e-9);
agg.flush_window();
assert!(agg.snapshot(Horizon::Window).rows.is_empty());
assert_eq!(row(&agg.snapshot(Horizon::Session), "rpc.total").count, 4);
}
}