Skip to main content

el_telemetry/
lib.rs

1//! `el-telemetry` — a one-way, downstream subscriber that folds content-free
2//! [`el_core::DomainEvent`]s into performance snapshots (ADR-007).
3//!
4//! It depends on `el-core` and **nothing depends on it**; it has no network
5//! channel. Because it can only read the numeric fields of already-content-free
6//! events, "no user content in telemetry" is structural.
7
8#![forbid(unsafe_code)]
9
10use el_core::{DomainEvent, EventEnvelope};
11
12/// A content-free sample of session performance.
13#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
14pub struct TelemetrySnapshot {
15    pub prefill_tps: u32,
16    pub decode_tps: u32,
17    pub ttft_ms: u32,
18    pub peak_bytes: u64,
19    pub tokens_generated: u32,
20    pub compressions: u32,
21    pub safety_violations: u32,
22}
23
24/// Subscribes to the domain-event stream and maintains a running snapshot.
25#[derive(Debug, Default)]
26pub struct MetricsCollector {
27    snap: TelemetrySnapshot,
28}
29
30impl MetricsCollector {
31    pub fn new() -> Self {
32        Self::default()
33    }
34
35    pub fn snapshot(&self) -> TelemetrySnapshot {
36        self.snap
37    }
38
39    /// Fold one event into the snapshot. Reads only numeric/enum fields.
40    pub fn observe(&mut self, env: &EventEnvelope) {
41        match env.event {
42            DomainEvent::PrefillCompleted { prefill_tps, .. } => {
43                self.snap.prefill_tps = prefill_tps;
44            }
45            DomainEvent::TokenCommitted { .. } => {
46                self.snap.tokens_generated += 1;
47            }
48            DomainEvent::PromptCompressed { .. } => {
49                self.snap.compressions += 1;
50            }
51            DomainEvent::SafetyViolationDetected { .. } => {
52                self.snap.safety_violations += 1;
53            }
54            DomainEvent::MemoryPlanCreated { total_bytes, .. } => {
55                self.snap.peak_bytes = self.snap.peak_bytes.max(total_bytes);
56            }
57            DomainEvent::MetricsSampled {
58                decode_tps,
59                ttft_ms,
60                peak_bytes,
61            } => {
62                self.snap.decode_tps = decode_tps;
63                self.snap.ttft_ms = ttft_ms;
64                self.snap.peak_bytes = self.snap.peak_bytes.max(peak_bytes);
65            }
66            _ => {}
67        }
68    }
69}
70
71#[cfg(test)]
72mod tests {
73    use super::*;
74    use el_core::SessionId;
75
76    fn env(step: u32, event: DomainEvent) -> EventEnvelope {
77        EventEnvelope::new(SessionId(1), step, event)
78    }
79
80    #[test]
81    fn folds_events_into_counters() {
82        let mut c = MetricsCollector::new();
83        c.observe(&env(
84            0,
85            DomainEvent::PrefillCompleted {
86                prompt_tokens: 100,
87                kv_len: 100,
88                prefill_tps: 480,
89            },
90        ));
91        c.observe(&env(1, DomainEvent::TokenCommitted { kv_len: 101 }));
92        c.observe(&env(2, DomainEvent::TokenCommitted { kv_len: 102 }));
93        c.observe(&env(
94            3,
95            DomainEvent::MetricsSampled {
96                decode_tps: 55,
97                ttft_ms: 180,
98                peak_bytes: 600_000_000,
99            },
100        ));
101
102        let s = c.snapshot();
103        assert_eq!(s.prefill_tps, 480);
104        assert_eq!(s.tokens_generated, 2);
105        assert_eq!(s.decode_tps, 55);
106        assert_eq!(s.peak_bytes, 600_000_000);
107    }
108
109    #[test]
110    fn peak_bytes_is_monotonic() {
111        let mut c = MetricsCollector::new();
112        c.observe(&env(
113            0,
114            DomainEvent::MemoryPlanCreated {
115                total_bytes: 500,
116                sram_bytes: 100,
117                dram_bytes: 400,
118            },
119        ));
120        c.observe(&env(
121            1,
122            DomainEvent::MetricsSampled {
123                decode_tps: 1,
124                ttft_ms: 1,
125                peak_bytes: 200,
126            },
127        ));
128        assert_eq!(c.snapshot().peak_bytes, 500, "high-water mark only rises");
129    }
130}