Skip to main content

aurum_core/
observability.rs

1//! Structured operation metrics and privacy-safe diagnostics (JOE-1627).
2//!
3//! Counters are process-local, lock-free, and never store payload text/PCM.
4
5use serde::{Deserialize, Serialize};
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::Arc;
8use std::time::{Duration, Instant};
9
10/// Schema version for metrics / diagnostic JSON.
11pub const METRICS_SCHEMA_VERSION: u32 = 1;
12
13/// Process-wide metrics sink (also attachable to an engine via [`Arc`]).
14#[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    /// Process-global shared metrics sink (same instance as [`process_metrics`]).
35    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/// Serializable metrics snapshot.
97#[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/// Redacted diagnostic bundle for support / `aurum doctor --json` enrichment.
114#[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    /// Never contains API keys or raw audio.
121    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/// Simple wall-clock timer for inference spans.
140#[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
163/// Process-global metrics (shared by CLI/FFI).
164pub 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}