Skip to main content

aurum_core/
observability.rs

1//! Privacy-safe per-operation metrics, events, and diagnostics (JOE-1627 / JOE-2222).
2//!
3//! Counters are process-local (or engine-local via [`Arc`]) and never store
4//! payload text, PCM, API keys, or absolute private paths. Hosts may attach a
5//! bounded event sink; the default is a no-op with minimal overhead.
6
7use serde::{Deserialize, Serialize};
8use std::collections::VecDeque;
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex};
11use std::time::{Duration, Instant};
12
13/// Schema version for metrics / diagnostic / event JSON.
14pub const METRICS_SCHEMA_VERSION: u32 = 2;
15/// Operation event schema version.
16pub const OP_EVENT_SCHEMA_VERSION: u32 = 1;
17/// Default max buffered events per sink (overflow drops, never blocks).
18pub const DEFAULT_EVENT_QUEUE_CAP: usize = 256;
19
20// ---------------------------------------------------------------------------
21// Controlled enums (no free-form user labels)
22// ---------------------------------------------------------------------------
23
24/// High-level operation class.
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
26#[serde(rename_all = "snake_case")]
27pub enum OpKind {
28    Stt,
29    Tts,
30    Cleanup,
31    ModelLoad,
32    Download,
33    BatchItem,
34}
35
36impl OpKind {
37    pub fn as_str(self) -> &'static str {
38        match self {
39            Self::Stt => "stt",
40            Self::Tts => "tts",
41            Self::Cleanup => "cleanup",
42            Self::ModelLoad => "model_load",
43            Self::Download => "download",
44            Self::BatchItem => "batch_item",
45        }
46    }
47}
48
49/// Controlled stage labels (not free-form).
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
51#[serde(rename_all = "snake_case")]
52pub enum OpStage {
53    Start,
54    QueueWait,
55    Admitted,
56    ModelLoad,
57    Encode,
58    NetworkSend,
59    NetworkBody,
60    Inference,
61    Normalize,
62    Cleanup,
63    Chunk,
64    Stitch,
65    Output,
66    Terminal,
67}
68
69/// Terminal outcome category (exactly one per operation).
70#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
71#[serde(rename_all = "snake_case")]
72pub enum TerminalCategory {
73    Completed,
74    Failed,
75    Cancelled,
76    Deadline,
77    Overload,
78    Busy,
79}
80
81impl TerminalCategory {
82    pub fn as_str(self) -> &'static str {
83        match self {
84            Self::Completed => "completed",
85            Self::Failed => "failed",
86            Self::Cancelled => "cancelled",
87            Self::Deadline => "deadline",
88            Self::Overload => "overload",
89            Self::Busy => "busy",
90        }
91    }
92}
93
94/// Scope for metrics isolation.
95#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
96#[serde(rename_all = "snake_case")]
97pub enum MetricsScope {
98    EngineLocal,
99    #[default]
100    ProcessGlobal,
101}
102
103// ---------------------------------------------------------------------------
104// Events
105// ---------------------------------------------------------------------------
106
107/// Versioned privacy-safe operation event (no payloads/secrets/paths).
108#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
109pub struct OpEvent {
110    pub schema_version: u32,
111    pub request_id: u64,
112    pub operation: OpKind,
113    pub stage: OpStage,
114    #[serde(default, skip_serializing_if = "Option::is_none")]
115    pub provider_id: Option<String>,
116    #[serde(default, skip_serializing_if = "Option::is_none")]
117    pub backend_class: Option<String>,
118    #[serde(default, skip_serializing_if = "Option::is_none")]
119    pub model_id: Option<String>,
120    pub scope: MetricsScope,
121    /// Monotonic elapsed ms since operation start (when known).
122    #[serde(default, skip_serializing_if = "Option::is_none")]
123    pub elapsed_ms: Option<u64>,
124    #[serde(default, skip_serializing_if = "Option::is_none")]
125    pub queue_ms: Option<u64>,
126    #[serde(default, skip_serializing_if = "Option::is_none")]
127    pub encoded_bytes: Option<u64>,
128    #[serde(default, skip_serializing_if = "Option::is_none")]
129    pub decoded_bytes: Option<u64>,
130    #[serde(default, skip_serializing_if = "Option::is_none")]
131    pub chunk_index: Option<u32>,
132    #[serde(default, skip_serializing_if = "Option::is_none")]
133    pub chunk_count: Option<u32>,
134    #[serde(default, skip_serializing_if = "Option::is_none")]
135    pub cache_state: Option<String>,
136    #[serde(default, skip_serializing_if = "Option::is_none")]
137    pub terminal: Option<TerminalCategory>,
138    #[serde(default, skip_serializing_if = "Option::is_none")]
139    pub retryable: Option<bool>,
140    #[serde(default, skip_serializing_if = "Option::is_none")]
141    pub error_category: Option<String>,
142}
143
144impl OpEvent {
145    pub fn stage(request_id: u64, operation: OpKind, stage: OpStage, scope: MetricsScope) -> Self {
146        Self {
147            schema_version: OP_EVENT_SCHEMA_VERSION,
148            request_id,
149            operation,
150            stage,
151            provider_id: None,
152            backend_class: None,
153            model_id: None,
154            scope,
155            elapsed_ms: None,
156            queue_ms: None,
157            encoded_bytes: None,
158            decoded_bytes: None,
159            chunk_index: None,
160            chunk_count: None,
161            cache_state: None,
162            terminal: None,
163            retryable: None,
164            error_category: None,
165        }
166    }
167
168    pub fn with_provider(mut self, id: impl Into<String>) -> Self {
169        self.provider_id = Some(id.into());
170        self
171    }
172
173    pub fn with_model(mut self, id: impl Into<String>) -> Self {
174        self.model_id = Some(id.into());
175        self
176    }
177
178    pub fn with_elapsed_ms(mut self, ms: u64) -> Self {
179        self.elapsed_ms = Some(ms);
180        self
181    }
182
183    pub fn with_terminal(mut self, cat: TerminalCategory, retryable: bool) -> Self {
184        self.stage = OpStage::Terminal;
185        self.terminal = Some(cat);
186        self.retryable = Some(retryable);
187        self
188    }
189
190    /// Scan serialized form for forbidden markers (privacy canary helper).
191    pub fn to_json(&self) -> Result<String, serde_json::Error> {
192        serde_json::to_string(self)
193    }
194}
195
196// ---------------------------------------------------------------------------
197// Event sink
198// ---------------------------------------------------------------------------
199
200/// Host-facing event sink. Must not block core locks or retain payloads.
201pub trait EventSink: Send + Sync {
202    fn emit(&self, event: OpEvent);
203}
204
205/// No-op sink (default).
206#[derive(Debug, Default, Clone, Copy)]
207pub struct NoopEventSink;
208
209impl EventSink for NoopEventSink {
210    fn emit(&self, _event: OpEvent) {}
211}
212
213/// Bounded in-memory queue; overflow increments `dropped` and never blocks.
214#[derive(Debug)]
215pub struct BoundedEventSink {
216    cap: usize,
217    inner: Mutex<VecDeque<OpEvent>>,
218    dropped: AtomicU64,
219}
220
221impl BoundedEventSink {
222    pub fn new(cap: usize) -> Self {
223        Self {
224            cap: cap.max(1),
225            inner: Mutex::new(VecDeque::with_capacity(cap.min(64))),
226            dropped: AtomicU64::new(0),
227        }
228    }
229
230    pub fn dropped(&self) -> u64 {
231        self.dropped.load(Ordering::Relaxed)
232    }
233
234    pub fn drain(&self) -> Vec<OpEvent> {
235        self.inner
236            .lock()
237            .map(|mut q| q.drain(..).collect())
238            .unwrap_or_default()
239    }
240
241    pub fn len(&self) -> usize {
242        self.inner.lock().map(|q| q.len()).unwrap_or(0)
243    }
244
245    pub fn is_empty(&self) -> bool {
246        self.len() == 0
247    }
248}
249
250impl Default for BoundedEventSink {
251    fn default() -> Self {
252        Self::new(DEFAULT_EVENT_QUEUE_CAP)
253    }
254}
255
256impl EventSink for BoundedEventSink {
257    fn emit(&self, event: OpEvent) {
258        if let Ok(mut q) = self.inner.lock() {
259            if q.len() >= self.cap {
260                self.dropped.fetch_add(1, Ordering::Relaxed);
261                return;
262            }
263            q.push_back(event);
264        } else {
265            self.dropped.fetch_add(1, Ordering::Relaxed);
266        }
267    }
268}
269
270// ---------------------------------------------------------------------------
271// Metrics
272// ---------------------------------------------------------------------------
273
274/// Process-wide or engine-local metrics sink (also attachable via [`Arc`]).
275pub struct Metrics {
276    pub ops_started: AtomicU64,
277    pub ops_completed: AtomicU64,
278    pub ops_failed: AtomicU64,
279    pub ops_cancelled: AtomicU64,
280    pub ops_deadline: AtomicU64,
281    pub queue_wait_ms_total: AtomicU64,
282    pub inference_ms_total: AtomicU64,
283    pub encode_ms_total: AtomicU64,
284    pub network_ms_total: AtomicU64,
285    pub model_loads: AtomicU64,
286    pub cache_hits: AtomicU64,
287    pub cache_misses: AtomicU64,
288    pub cache_evictions: AtomicU64,
289    pub remote_errors: AtomicU64,
290    pub remote_auth_errors: AtomicU64,
291    pub remote_rate_limits: AtomicU64,
292    pub remote_quota_errors: AtomicU64,
293    pub remote_network_errors: AtomicU64,
294    pub remote_invalid_payload: AtomicU64,
295    pub upload_bytes_total: AtomicU64,
296    pub response_bytes_total: AtomicU64,
297    pub decoded_bytes_total: AtomicU64,
298    pub long_form_chunks_attempted: AtomicU64,
299    pub long_form_chunks_completed: AtomicU64,
300    pub long_form_chunks_failed: AtomicU64,
301    pub tts_chunks_total: AtomicU64,
302    pub tts_chars_total: AtomicU64,
303    pub batch_items_succeeded: AtomicU64,
304    pub batch_items_failed: AtomicU64,
305    pub batch_stale_reprocess: AtomicU64,
306    pub output_tx_success: AtomicU64,
307    pub output_tx_failure: AtomicU64,
308    pub busy_rejections: AtomicU64,
309    pub overload_rejections: AtomicU64,
310    pub events_dropped: AtomicU64,
311    /// Optional host sink (outside hot locks when possible).
312    event_sink: Mutex<Option<Arc<dyn EventSink>>>,
313    scope: MetricsScope,
314}
315
316impl std::fmt::Debug for Metrics {
317    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
318        f.debug_struct("Metrics")
319            .field("scope", &self.scope)
320            .field("snapshot", &self.snapshot())
321            .finish()
322    }
323}
324
325impl Default for Metrics {
326    fn default() -> Self {
327        Self::new()
328    }
329}
330
331impl Metrics {
332    pub fn new() -> Self {
333        Self {
334            ops_started: AtomicU64::new(0),
335            ops_completed: AtomicU64::new(0),
336            ops_failed: AtomicU64::new(0),
337            ops_cancelled: AtomicU64::new(0),
338            ops_deadline: AtomicU64::new(0),
339            queue_wait_ms_total: AtomicU64::new(0),
340            inference_ms_total: AtomicU64::new(0),
341            encode_ms_total: AtomicU64::new(0),
342            network_ms_total: AtomicU64::new(0),
343            model_loads: AtomicU64::new(0),
344            cache_hits: AtomicU64::new(0),
345            cache_misses: AtomicU64::new(0),
346            cache_evictions: AtomicU64::new(0),
347            remote_errors: AtomicU64::new(0),
348            remote_auth_errors: AtomicU64::new(0),
349            remote_rate_limits: AtomicU64::new(0),
350            remote_quota_errors: AtomicU64::new(0),
351            remote_network_errors: AtomicU64::new(0),
352            remote_invalid_payload: AtomicU64::new(0),
353            upload_bytes_total: AtomicU64::new(0),
354            response_bytes_total: AtomicU64::new(0),
355            decoded_bytes_total: AtomicU64::new(0),
356            long_form_chunks_attempted: AtomicU64::new(0),
357            long_form_chunks_completed: AtomicU64::new(0),
358            long_form_chunks_failed: AtomicU64::new(0),
359            tts_chunks_total: AtomicU64::new(0),
360            tts_chars_total: AtomicU64::new(0),
361            batch_items_succeeded: AtomicU64::new(0),
362            batch_items_failed: AtomicU64::new(0),
363            batch_stale_reprocess: AtomicU64::new(0),
364            output_tx_success: AtomicU64::new(0),
365            output_tx_failure: AtomicU64::new(0),
366            busy_rejections: AtomicU64::new(0),
367            overload_rejections: AtomicU64::new(0),
368            events_dropped: AtomicU64::new(0),
369            event_sink: Mutex::new(None),
370            scope: MetricsScope::ProcessGlobal,
371        }
372    }
373
374    pub fn engine_local() -> Self {
375        let mut m = Self::new();
376        m.scope = MetricsScope::EngineLocal;
377        m
378    }
379
380    pub fn scope(&self) -> MetricsScope {
381        self.scope
382    }
383
384    /// Process-global shared metrics sink (same instance as [`process_metrics`]).
385    pub fn shared() -> Arc<Self> {
386        process_metrics_arc()
387    }
388
389    pub fn set_event_sink(&self, sink: Option<Arc<dyn EventSink>>) {
390        if let Ok(mut g) = self.event_sink.lock() {
391            *g = sink;
392        }
393    }
394
395    pub fn emit(&self, event: OpEvent) {
396        if let Ok(g) = self.event_sink.lock() {
397            if let Some(sink) = g.as_ref() {
398                sink.emit(event);
399            }
400        }
401    }
402
403    pub fn record_start(&self) {
404        self.ops_started.fetch_add(1, Ordering::Relaxed);
405    }
406
407    pub fn record_complete(&self, inference: Duration) {
408        self.ops_completed.fetch_add(1, Ordering::Relaxed);
409        self.inference_ms_total
410            .fetch_add(inference.as_millis() as u64, Ordering::Relaxed);
411    }
412
413    pub fn record_failed(&self) {
414        self.ops_failed.fetch_add(1, Ordering::Relaxed);
415    }
416
417    pub fn record_cancelled(&self) {
418        self.ops_cancelled.fetch_add(1, Ordering::Relaxed);
419    }
420
421    pub fn record_deadline(&self) {
422        self.ops_deadline.fetch_add(1, Ordering::Relaxed);
423    }
424
425    pub fn record_model_load(&self) {
426        self.model_loads.fetch_add(1, Ordering::Relaxed);
427    }
428
429    pub fn record_cache_hit(&self) {
430        self.cache_hits.fetch_add(1, Ordering::Relaxed);
431    }
432
433    pub fn record_cache_miss(&self) {
434        self.cache_misses.fetch_add(1, Ordering::Relaxed);
435    }
436
437    pub fn record_busy(&self) {
438        self.busy_rejections.fetch_add(1, Ordering::Relaxed);
439    }
440
441    pub fn record_overload(&self) {
442        self.overload_rejections.fetch_add(1, Ordering::Relaxed);
443    }
444
445    pub fn record_queue_wait(&self, d: Duration) {
446        self.queue_wait_ms_total
447            .fetch_add(d.as_millis() as u64, Ordering::Relaxed);
448    }
449
450    pub fn record_remote_error(&self) {
451        self.remote_errors.fetch_add(1, Ordering::Relaxed);
452    }
453
454    pub fn record_terminal(&self, cat: TerminalCategory) {
455        match cat {
456            TerminalCategory::Completed => {
457                // completed also recorded via record_complete in many paths
458            }
459            TerminalCategory::Failed => self.record_failed(),
460            TerminalCategory::Cancelled => self.record_cancelled(),
461            TerminalCategory::Deadline => self.record_deadline(),
462            TerminalCategory::Overload => self.record_overload(),
463            TerminalCategory::Busy => self.record_busy(),
464        }
465    }
466
467    pub fn snapshot(&self) -> MetricsSnapshot {
468        MetricsSnapshot {
469            schema_version: METRICS_SCHEMA_VERSION,
470            scope: self.scope,
471            ops_started: self.ops_started.load(Ordering::Relaxed),
472            ops_completed: self.ops_completed.load(Ordering::Relaxed),
473            ops_failed: self.ops_failed.load(Ordering::Relaxed),
474            ops_cancelled: self.ops_cancelled.load(Ordering::Relaxed),
475            ops_deadline: self.ops_deadline.load(Ordering::Relaxed),
476            queue_wait_ms_total: self.queue_wait_ms_total.load(Ordering::Relaxed),
477            inference_ms_total: self.inference_ms_total.load(Ordering::Relaxed),
478            encode_ms_total: self.encode_ms_total.load(Ordering::Relaxed),
479            network_ms_total: self.network_ms_total.load(Ordering::Relaxed),
480            model_loads: self.model_loads.load(Ordering::Relaxed),
481            cache_hits: self.cache_hits.load(Ordering::Relaxed),
482            cache_misses: self.cache_misses.load(Ordering::Relaxed),
483            cache_evictions: self.cache_evictions.load(Ordering::Relaxed),
484            remote_errors: self.remote_errors.load(Ordering::Relaxed),
485            remote_auth_errors: self.remote_auth_errors.load(Ordering::Relaxed),
486            remote_rate_limits: self.remote_rate_limits.load(Ordering::Relaxed),
487            remote_quota_errors: self.remote_quota_errors.load(Ordering::Relaxed),
488            remote_network_errors: self.remote_network_errors.load(Ordering::Relaxed),
489            remote_invalid_payload: self.remote_invalid_payload.load(Ordering::Relaxed),
490            upload_bytes_total: self.upload_bytes_total.load(Ordering::Relaxed),
491            response_bytes_total: self.response_bytes_total.load(Ordering::Relaxed),
492            decoded_bytes_total: self.decoded_bytes_total.load(Ordering::Relaxed),
493            long_form_chunks_attempted: self.long_form_chunks_attempted.load(Ordering::Relaxed),
494            long_form_chunks_completed: self.long_form_chunks_completed.load(Ordering::Relaxed),
495            long_form_chunks_failed: self.long_form_chunks_failed.load(Ordering::Relaxed),
496            tts_chunks_total: self.tts_chunks_total.load(Ordering::Relaxed),
497            tts_chars_total: self.tts_chars_total.load(Ordering::Relaxed),
498            batch_items_succeeded: self.batch_items_succeeded.load(Ordering::Relaxed),
499            batch_items_failed: self.batch_items_failed.load(Ordering::Relaxed),
500            batch_stale_reprocess: self.batch_stale_reprocess.load(Ordering::Relaxed),
501            output_tx_success: self.output_tx_success.load(Ordering::Relaxed),
502            output_tx_failure: self.output_tx_failure.load(Ordering::Relaxed),
503            busy_rejections: self.busy_rejections.load(Ordering::Relaxed),
504            overload_rejections: self.overload_rejections.load(Ordering::Relaxed),
505            events_dropped: self.events_dropped.load(Ordering::Relaxed),
506        }
507    }
508}
509
510/// Serializable metrics snapshot (no averages labelled as p95).
511#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
512pub struct MetricsSnapshot {
513    pub schema_version: u32,
514    #[serde(default = "default_process_global")]
515    pub scope: MetricsScope,
516    pub ops_started: u64,
517    pub ops_completed: u64,
518    pub ops_failed: u64,
519    pub ops_cancelled: u64,
520    pub ops_deadline: u64,
521    pub queue_wait_ms_total: u64,
522    pub inference_ms_total: u64,
523    #[serde(default)]
524    pub encode_ms_total: u64,
525    #[serde(default)]
526    pub network_ms_total: u64,
527    pub model_loads: u64,
528    #[serde(default)]
529    pub cache_hits: u64,
530    #[serde(default)]
531    pub cache_misses: u64,
532    #[serde(default)]
533    pub cache_evictions: u64,
534    pub remote_errors: u64,
535    #[serde(default)]
536    pub remote_auth_errors: u64,
537    #[serde(default)]
538    pub remote_rate_limits: u64,
539    #[serde(default)]
540    pub remote_quota_errors: u64,
541    #[serde(default)]
542    pub remote_network_errors: u64,
543    #[serde(default)]
544    pub remote_invalid_payload: u64,
545    #[serde(default)]
546    pub upload_bytes_total: u64,
547    #[serde(default)]
548    pub response_bytes_total: u64,
549    #[serde(default)]
550    pub decoded_bytes_total: u64,
551    #[serde(default)]
552    pub long_form_chunks_attempted: u64,
553    #[serde(default)]
554    pub long_form_chunks_completed: u64,
555    #[serde(default)]
556    pub long_form_chunks_failed: u64,
557    #[serde(default)]
558    pub tts_chunks_total: u64,
559    #[serde(default)]
560    pub tts_chars_total: u64,
561    #[serde(default)]
562    pub batch_items_succeeded: u64,
563    #[serde(default)]
564    pub batch_items_failed: u64,
565    #[serde(default)]
566    pub batch_stale_reprocess: u64,
567    #[serde(default)]
568    pub output_tx_success: u64,
569    #[serde(default)]
570    pub output_tx_failure: u64,
571    pub busy_rejections: u64,
572    pub overload_rejections: u64,
573    #[serde(default)]
574    pub events_dropped: u64,
575}
576
577fn default_process_global() -> MetricsScope {
578    MetricsScope::ProcessGlobal
579}
580
581/// Redacted diagnostic bundle for support / `aurum doctor --json` enrichment.
582#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
583pub struct DiagnosticBundle {
584    pub schema_version: u32,
585    pub request_id: Option<String>,
586    pub operation: Option<String>,
587    pub metrics: MetricsSnapshot,
588    /// Never contains API keys or raw audio.
589    pub notes: Vec<String>,
590}
591
592impl DiagnosticBundle {
593    pub fn from_metrics(metrics: &Metrics) -> Self {
594        Self {
595            schema_version: METRICS_SCHEMA_VERSION,
596            request_id: None,
597            operation: None,
598            metrics: metrics.snapshot(),
599            notes: vec![
600                "payloads (PCM/text/API keys) are never included".into(),
601                "counters are process-local or engine-local best-effort".into(),
602                "totals are sums; averages are not p95".into(),
603            ],
604        }
605    }
606
607    pub fn with_request_id(mut self, id: u64) -> Self {
608        self.request_id = Some(id.to_string());
609        self
610    }
611}
612
613// ---------------------------------------------------------------------------
614// Terminal guard (exactly one outcome)
615// ---------------------------------------------------------------------------
616
617/// Ensures a terminal outcome is recorded at most once (including drop).
618pub struct TerminalGuard {
619    metrics: Arc<Metrics>,
620    request_id: u64,
621    operation: OpKind,
622    scope: MetricsScope,
623    finished: bool,
624    start: Instant,
625}
626
627impl TerminalGuard {
628    pub fn start(metrics: Arc<Metrics>, request_id: u64, operation: OpKind) -> Self {
629        metrics.record_start();
630        let scope = metrics.scope();
631        metrics.emit(OpEvent::stage(request_id, operation, OpStage::Start, scope));
632        Self {
633            metrics,
634            request_id,
635            operation,
636            scope,
637            finished: false,
638            start: Instant::now(),
639        }
640    }
641
642    pub fn finish(&mut self, cat: TerminalCategory, retryable: bool) {
643        if self.finished {
644            return;
645        }
646        self.finished = true;
647        match cat {
648            TerminalCategory::Completed => {
649                self.metrics.record_complete(self.start.elapsed());
650            }
651            other => self.metrics.record_terminal(other),
652        }
653        self.metrics.emit(
654            OpEvent::stage(
655                self.request_id,
656                self.operation,
657                OpStage::Terminal,
658                self.scope,
659            )
660            .with_elapsed_ms(self.start.elapsed().as_millis() as u64)
661            .with_terminal(cat, retryable),
662        );
663    }
664
665    pub fn request_id(&self) -> u64 {
666        self.request_id
667    }
668
669    pub fn elapsed(&self) -> Duration {
670        self.start.elapsed()
671    }
672}
673
674impl Drop for TerminalGuard {
675    fn drop(&mut self) {
676        // Panic / early return: count as failed once.
677        if !self.finished {
678            self.finish(TerminalCategory::Failed, false);
679        }
680    }
681}
682
683/// Simple wall-clock timer for inference spans.
684#[derive(Debug)]
685pub struct SpanTimer {
686    start: Instant,
687}
688
689impl SpanTimer {
690    pub fn start() -> Self {
691        Self {
692            start: Instant::now(),
693        }
694    }
695
696    pub fn elapsed(&self) -> Duration {
697        self.start.elapsed()
698    }
699}
700
701fn process_metrics_arc() -> Arc<Metrics> {
702    use once_cell::sync::Lazy;
703    static SHARED: Lazy<Arc<Metrics>> = Lazy::new(|| Arc::new(Metrics::new()));
704    Arc::clone(&SHARED)
705}
706
707/// Process-global metrics (shared by CLI/FFI).
708pub fn process_metrics() -> Arc<Metrics> {
709    process_metrics_arc()
710}
711
712/// Forbidden substrings that must never appear in public observability output.
713pub const PRIVACY_CANARY_MARKERS: &[&str] = &[
714    "sk-test-secret-key",
715    "BEGIN_PRIVATE_AUDIO",
716    "USER_TRANSCRIPT_PAYLOAD",
717    "SYNTHESIS_TEXT_SECRET",
718    "/Users/private/home/",
719    "Authorization: Bearer",
720];
721
722/// Scan JSON/text for privacy canary markers.
723pub fn privacy_scan(text: &str) -> Result<(), String> {
724    for m in PRIVACY_CANARY_MARKERS {
725        if text.contains(m) {
726            return Err(format!("privacy canary hit: marker present: {m}"));
727        }
728    }
729    Ok(())
730}
731
732#[cfg(test)]
733mod tests {
734    use super::*;
735
736    #[test]
737    fn metrics_roundtrip() {
738        let m = Metrics::new();
739        m.record_start();
740        m.record_complete(Duration::from_millis(12));
741        let s = m.snapshot();
742        assert_eq!(s.ops_started, 1);
743        assert_eq!(s.ops_completed, 1);
744        assert!(s.inference_ms_total >= 12);
745        assert_eq!(s.schema_version, METRICS_SCHEMA_VERSION);
746        let bundle = DiagnosticBundle::from_metrics(&m);
747        assert!(serde_json::to_string(&bundle)
748            .unwrap()
749            .contains("ops_started"));
750    }
751
752    #[test]
753    fn terminal_guard_once() {
754        let m = Arc::new(Metrics::engine_local());
755        let sink = Arc::new(BoundedEventSink::new(32));
756        m.set_event_sink(Some(sink.clone()));
757        {
758            let mut g = TerminalGuard::start(m.clone(), 42, OpKind::Stt);
759            g.finish(TerminalCategory::Completed, false);
760            g.finish(TerminalCategory::Failed, false); // no-op
761        }
762        assert_eq!(m.snapshot().ops_started, 1);
763        assert_eq!(m.snapshot().ops_completed, 1);
764        assert_eq!(m.snapshot().ops_failed, 0);
765        let events = sink.drain();
766        assert!(events.iter().any(|e| e.stage == OpStage::Start));
767        assert_eq!(events.iter().filter(|e| e.terminal.is_some()).count(), 1);
768        assert!(events.iter().all(|e| e.request_id == 42));
769    }
770
771    #[test]
772    fn terminal_guard_drop_counts_failed() {
773        let m = Arc::new(Metrics::new());
774        {
775            let _g = TerminalGuard::start(m.clone(), 1, OpKind::Tts);
776        }
777        assert_eq!(m.snapshot().ops_failed, 1);
778    }
779
780    #[test]
781    fn bounded_sink_drops() {
782        let sink = BoundedEventSink::new(2);
783        for i in 0..5 {
784            sink.emit(OpEvent::stage(
785                i,
786                OpKind::Stt,
787                OpStage::Start,
788                MetricsScope::ProcessGlobal,
789            ));
790        }
791        assert_eq!(sink.len(), 2);
792        assert_eq!(sink.dropped(), 3);
793    }
794
795    #[test]
796    fn distinct_outcomes() {
797        let m = Metrics::new();
798        m.record_cancelled();
799        m.record_deadline();
800        m.record_overload();
801        m.record_busy();
802        m.record_failed();
803        let s = m.snapshot();
804        assert_eq!(s.ops_cancelled, 1);
805        assert_eq!(s.ops_deadline, 1);
806        assert_eq!(s.overload_rejections, 1);
807        assert_eq!(s.busy_rejections, 1);
808        assert_eq!(s.ops_failed, 1);
809    }
810
811    #[test]
812    fn engine_vs_process_scope() {
813        let eng = Metrics::engine_local();
814        let proc = Metrics::new();
815        assert_eq!(eng.scope(), MetricsScope::EngineLocal);
816        assert_eq!(proc.scope(), MetricsScope::ProcessGlobal);
817        eng.record_start();
818        assert_eq!(eng.snapshot().ops_started, 1);
819        assert_eq!(proc.snapshot().ops_started, 0);
820    }
821
822    #[test]
823    fn privacy_canary_on_event_and_snapshot() {
824        let m = Metrics::new();
825        m.record_start();
826        let event = OpEvent::stage(
827            7,
828            OpKind::Stt,
829            OpStage::Inference,
830            MetricsScope::EngineLocal,
831        )
832        .with_provider("local")
833        .with_model("base");
834        let json = event.to_json().unwrap();
835        privacy_scan(&json).unwrap();
836        privacy_scan(&serde_json::to_string(&m.snapshot()).unwrap()).unwrap();
837        // Inject markers into notes would fail — bundle notes are fixed.
838        let mut bad = json;
839        bad.push_str("sk-test-secret-key");
840        assert!(privacy_scan(&bad).is_err());
841    }
842
843    #[test]
844    fn event_schema_stable() {
845        let e = OpEvent::stage(
846            1,
847            OpKind::Cleanup,
848            OpStage::Normalize,
849            MetricsScope::ProcessGlobal,
850        );
851        let j1 = e.to_json().unwrap();
852        let j2 = e.to_json().unwrap();
853        assert_eq!(j1, j2);
854        assert!(j1.contains("schema_version"));
855    }
856}