aurum_core/
observability.rs1use serde::{Deserialize, Serialize};
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::Arc;
8use std::time::{Duration, Instant};
9
10pub const METRICS_SCHEMA_VERSION: u32 = 1;
12
13#[derive(Debug, Default)]
15pub struct Metrics {
16 pub ops_started: AtomicU64,
17 pub ops_completed: AtomicU64,
18 pub ops_failed: AtomicU64,
19 pub ops_cancelled: AtomicU64,
20 pub ops_deadline: AtomicU64,
21 pub queue_wait_ms_total: AtomicU64,
22 pub inference_ms_total: AtomicU64,
23 pub model_loads: AtomicU64,
24 pub remote_errors: AtomicU64,
25 pub busy_rejections: AtomicU64,
26 pub overload_rejections: AtomicU64,
27}
28
29impl Metrics {
30 pub fn new() -> Self {
31 Self::default()
32 }
33
34 pub fn shared() -> Arc<Self> {
36 process_metrics_arc()
37 }
38
39 pub fn record_start(&self) {
40 self.ops_started.fetch_add(1, Ordering::Relaxed);
41 }
42
43 pub fn record_complete(&self, inference: Duration) {
44 self.ops_completed.fetch_add(1, Ordering::Relaxed);
45 self.inference_ms_total
46 .fetch_add(inference.as_millis() as u64, Ordering::Relaxed);
47 }
48
49 pub fn record_failed(&self) {
50 self.ops_failed.fetch_add(1, Ordering::Relaxed);
51 }
52
53 pub fn record_cancelled(&self) {
54 self.ops_cancelled.fetch_add(1, Ordering::Relaxed);
55 }
56
57 pub fn record_deadline(&self) {
58 self.ops_deadline.fetch_add(1, Ordering::Relaxed);
59 }
60
61 pub fn record_model_load(&self) {
62 self.model_loads.fetch_add(1, Ordering::Relaxed);
63 }
64
65 pub fn record_busy(&self) {
66 self.busy_rejections.fetch_add(1, Ordering::Relaxed);
67 }
68
69 pub fn record_overload(&self) {
70 self.overload_rejections.fetch_add(1, Ordering::Relaxed);
71 }
72
73 pub fn record_queue_wait(&self, d: Duration) {
74 self.queue_wait_ms_total
75 .fetch_add(d.as_millis() as u64, Ordering::Relaxed);
76 }
77
78 pub fn snapshot(&self) -> MetricsSnapshot {
79 MetricsSnapshot {
80 schema_version: METRICS_SCHEMA_VERSION,
81 ops_started: self.ops_started.load(Ordering::Relaxed),
82 ops_completed: self.ops_completed.load(Ordering::Relaxed),
83 ops_failed: self.ops_failed.load(Ordering::Relaxed),
84 ops_cancelled: self.ops_cancelled.load(Ordering::Relaxed),
85 ops_deadline: self.ops_deadline.load(Ordering::Relaxed),
86 queue_wait_ms_total: self.queue_wait_ms_total.load(Ordering::Relaxed),
87 inference_ms_total: self.inference_ms_total.load(Ordering::Relaxed),
88 model_loads: self.model_loads.load(Ordering::Relaxed),
89 remote_errors: self.remote_errors.load(Ordering::Relaxed),
90 busy_rejections: self.busy_rejections.load(Ordering::Relaxed),
91 overload_rejections: self.overload_rejections.load(Ordering::Relaxed),
92 }
93 }
94}
95
96#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
98pub struct MetricsSnapshot {
99 pub schema_version: u32,
100 pub ops_started: u64,
101 pub ops_completed: u64,
102 pub ops_failed: u64,
103 pub ops_cancelled: u64,
104 pub ops_deadline: u64,
105 pub queue_wait_ms_total: u64,
106 pub inference_ms_total: u64,
107 pub model_loads: u64,
108 pub remote_errors: u64,
109 pub busy_rejections: u64,
110 pub overload_rejections: u64,
111}
112
113#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
115pub struct DiagnosticBundle {
116 pub schema_version: u32,
117 pub request_id: Option<String>,
118 pub operation: Option<String>,
119 pub metrics: MetricsSnapshot,
120 pub notes: Vec<String>,
122}
123
124impl DiagnosticBundle {
125 pub fn from_metrics(metrics: &Metrics) -> Self {
126 Self {
127 schema_version: METRICS_SCHEMA_VERSION,
128 request_id: None,
129 operation: None,
130 metrics: metrics.snapshot(),
131 notes: vec![
132 "payloads (PCM/text/API keys) are never included".into(),
133 "counters are process-local best-effort".into(),
134 ],
135 }
136 }
137}
138
139#[derive(Debug)]
141pub struct SpanTimer {
142 start: Instant,
143}
144
145impl SpanTimer {
146 pub fn start() -> Self {
147 Self {
148 start: Instant::now(),
149 }
150 }
151
152 pub fn elapsed(&self) -> Duration {
153 self.start.elapsed()
154 }
155}
156
157fn process_metrics_arc() -> Arc<Metrics> {
158 use once_cell::sync::Lazy;
159 static SHARED: Lazy<Arc<Metrics>> = Lazy::new(|| Arc::new(Metrics::new()));
160 Arc::clone(&SHARED)
161}
162
163pub fn process_metrics() -> Arc<Metrics> {
165 process_metrics_arc()
166}
167
168#[cfg(test)]
169mod tests {
170 use super::*;
171
172 #[test]
173 fn metrics_roundtrip() {
174 let m = Metrics::new();
175 m.record_start();
176 m.record_complete(Duration::from_millis(12));
177 let s = m.snapshot();
178 assert_eq!(s.ops_started, 1);
179 assert_eq!(s.ops_completed, 1);
180 assert!(s.inference_ms_total >= 12);
181 let bundle = DiagnosticBundle::from_metrics(&m);
182 assert!(serde_json::to_string(&bundle)
183 .unwrap()
184 .contains("ops_started"));
185 }
186}