1#![forbid(unsafe_code)]
9
10use el_core::{DomainEvent, EventEnvelope};
11
12#[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#[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 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}