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_decoded_bytes(mut self, n: u64) -> Self {
184        self.decoded_bytes = Some(n);
185        self
186    }
187
188    pub fn with_encoded_bytes(mut self, n: u64) -> Self {
189        self.encoded_bytes = Some(n);
190        self
191    }
192
193    pub fn with_chunk(mut self, index: u32, count: u32) -> Self {
194        self.chunk_index = Some(index);
195        self.chunk_count = Some(count);
196        self
197    }
198
199    pub fn with_terminal(mut self, cat: TerminalCategory, retryable: bool) -> Self {
200        self.stage = OpStage::Terminal;
201        self.terminal = Some(cat);
202        self.retryable = Some(retryable);
203        self
204    }
205
206    /// Scan serialized form for forbidden markers (privacy canary helper).
207    pub fn to_json(&self) -> Result<String, serde_json::Error> {
208        serde_json::to_string(self)
209    }
210}
211
212// ---------------------------------------------------------------------------
213// Event sink
214// ---------------------------------------------------------------------------
215
216/// Host-facing event sink. Must not block core locks or retain payloads.
217pub trait EventSink: Send + Sync {
218    fn emit(&self, event: OpEvent);
219}
220
221/// No-op sink (default).
222#[derive(Debug, Default, Clone, Copy)]
223pub struct NoopEventSink;
224
225impl EventSink for NoopEventSink {
226    fn emit(&self, _event: OpEvent) {}
227}
228
229/// Bounded in-memory queue; overflow increments `dropped` and never blocks.
230#[derive(Debug)]
231pub struct BoundedEventSink {
232    cap: usize,
233    inner: Mutex<VecDeque<OpEvent>>,
234    dropped: AtomicU64,
235}
236
237impl BoundedEventSink {
238    pub fn new(cap: usize) -> Self {
239        Self {
240            cap: cap.max(1),
241            inner: Mutex::new(VecDeque::with_capacity(cap.min(64))),
242            dropped: AtomicU64::new(0),
243        }
244    }
245
246    pub fn dropped(&self) -> u64 {
247        self.dropped.load(Ordering::Relaxed)
248    }
249
250    pub fn drain(&self) -> Vec<OpEvent> {
251        self.inner
252            .lock()
253            .map(|mut q| q.drain(..).collect())
254            .unwrap_or_default()
255    }
256
257    pub fn len(&self) -> usize {
258        self.inner.lock().map(|q| q.len()).unwrap_or(0)
259    }
260
261    pub fn is_empty(&self) -> bool {
262        self.len() == 0
263    }
264}
265
266impl Default for BoundedEventSink {
267    fn default() -> Self {
268        Self::new(DEFAULT_EVENT_QUEUE_CAP)
269    }
270}
271
272impl EventSink for BoundedEventSink {
273    fn emit(&self, event: OpEvent) {
274        if let Ok(mut q) = self.inner.lock() {
275            if q.len() >= self.cap {
276                self.dropped.fetch_add(1, Ordering::Relaxed);
277                return;
278            }
279            q.push_back(event);
280        } else {
281            self.dropped.fetch_add(1, Ordering::Relaxed);
282        }
283    }
284}
285
286// ---------------------------------------------------------------------------
287// Metrics
288// ---------------------------------------------------------------------------
289
290/// Process-wide or engine-local metrics sink (also attachable via [`Arc`]).
291pub struct Metrics {
292    pub ops_started: AtomicU64,
293    pub ops_completed: AtomicU64,
294    pub ops_failed: AtomicU64,
295    pub ops_cancelled: AtomicU64,
296    pub ops_deadline: AtomicU64,
297    pub queue_wait_ms_total: AtomicU64,
298    pub inference_ms_total: AtomicU64,
299    pub encode_ms_total: AtomicU64,
300    pub network_ms_total: AtomicU64,
301    pub model_loads: AtomicU64,
302    pub cache_hits: AtomicU64,
303    pub cache_misses: AtomicU64,
304    pub cache_evictions: AtomicU64,
305    pub remote_errors: AtomicU64,
306    pub remote_auth_errors: AtomicU64,
307    pub remote_rate_limits: AtomicU64,
308    pub remote_quota_errors: AtomicU64,
309    pub remote_network_errors: AtomicU64,
310    pub remote_invalid_payload: AtomicU64,
311    pub upload_bytes_total: AtomicU64,
312    pub response_bytes_total: AtomicU64,
313    pub decoded_bytes_total: AtomicU64,
314    pub long_form_chunks_attempted: AtomicU64,
315    pub long_form_chunks_completed: AtomicU64,
316    pub long_form_chunks_failed: AtomicU64,
317    pub tts_chunks_total: AtomicU64,
318    pub tts_chars_total: AtomicU64,
319    pub batch_items_succeeded: AtomicU64,
320    pub batch_items_failed: AtomicU64,
321    pub batch_stale_reprocess: AtomicU64,
322    pub output_tx_success: AtomicU64,
323    pub output_tx_failure: AtomicU64,
324    pub busy_rejections: AtomicU64,
325    pub overload_rejections: AtomicU64,
326    pub events_dropped: AtomicU64,
327    /// Optional host sink (outside hot locks when possible).
328    event_sink: Mutex<Option<Arc<dyn EventSink>>>,
329    scope: MetricsScope,
330}
331
332impl std::fmt::Debug for Metrics {
333    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
334        f.debug_struct("Metrics")
335            .field("scope", &self.scope)
336            .field("snapshot", &self.snapshot())
337            .finish()
338    }
339}
340
341impl Default for Metrics {
342    fn default() -> Self {
343        Self::new()
344    }
345}
346
347impl Metrics {
348    pub fn new() -> Self {
349        Self {
350            ops_started: AtomicU64::new(0),
351            ops_completed: AtomicU64::new(0),
352            ops_failed: AtomicU64::new(0),
353            ops_cancelled: AtomicU64::new(0),
354            ops_deadline: AtomicU64::new(0),
355            queue_wait_ms_total: AtomicU64::new(0),
356            inference_ms_total: AtomicU64::new(0),
357            encode_ms_total: AtomicU64::new(0),
358            network_ms_total: AtomicU64::new(0),
359            model_loads: AtomicU64::new(0),
360            cache_hits: AtomicU64::new(0),
361            cache_misses: AtomicU64::new(0),
362            cache_evictions: AtomicU64::new(0),
363            remote_errors: AtomicU64::new(0),
364            remote_auth_errors: AtomicU64::new(0),
365            remote_rate_limits: AtomicU64::new(0),
366            remote_quota_errors: AtomicU64::new(0),
367            remote_network_errors: AtomicU64::new(0),
368            remote_invalid_payload: AtomicU64::new(0),
369            upload_bytes_total: AtomicU64::new(0),
370            response_bytes_total: AtomicU64::new(0),
371            decoded_bytes_total: AtomicU64::new(0),
372            long_form_chunks_attempted: AtomicU64::new(0),
373            long_form_chunks_completed: AtomicU64::new(0),
374            long_form_chunks_failed: AtomicU64::new(0),
375            tts_chunks_total: AtomicU64::new(0),
376            tts_chars_total: AtomicU64::new(0),
377            batch_items_succeeded: AtomicU64::new(0),
378            batch_items_failed: AtomicU64::new(0),
379            batch_stale_reprocess: AtomicU64::new(0),
380            output_tx_success: AtomicU64::new(0),
381            output_tx_failure: AtomicU64::new(0),
382            busy_rejections: AtomicU64::new(0),
383            overload_rejections: AtomicU64::new(0),
384            events_dropped: AtomicU64::new(0),
385            event_sink: Mutex::new(None),
386            scope: MetricsScope::ProcessGlobal,
387        }
388    }
389
390    pub fn engine_local() -> Self {
391        let mut m = Self::new();
392        m.scope = MetricsScope::EngineLocal;
393        m
394    }
395
396    pub fn scope(&self) -> MetricsScope {
397        self.scope
398    }
399
400    /// Process-global shared metrics sink (same instance as [`process_metrics`]).
401    pub fn shared() -> Arc<Self> {
402        process_metrics_arc()
403    }
404
405    pub fn set_event_sink(&self, sink: Option<Arc<dyn EventSink>>) {
406        if let Ok(mut g) = self.event_sink.lock() {
407            *g = sink;
408        }
409    }
410
411    pub fn emit(&self, event: OpEvent) {
412        // Clone the Arc under the lock, then invoke the host callback *outside*
413        // so a slow/re-entrant sink cannot hold the metrics mutex (v0.0.23).
414        let sink = self
415            .event_sink
416            .lock()
417            .ok()
418            .and_then(|g| g.as_ref().map(Arc::clone));
419        if let Some(sink) = sink {
420            sink.emit(event);
421        }
422    }
423
424    pub fn record_start(&self) {
425        self.ops_started.fetch_add(1, Ordering::Relaxed);
426    }
427
428    pub fn record_complete(&self, inference: Duration) {
429        self.ops_completed.fetch_add(1, Ordering::Relaxed);
430        self.inference_ms_total
431            .fetch_add(inference.as_millis() as u64, Ordering::Relaxed);
432    }
433
434    pub fn record_failed(&self) {
435        self.ops_failed.fetch_add(1, Ordering::Relaxed);
436    }
437
438    pub fn record_cancelled(&self) {
439        self.ops_cancelled.fetch_add(1, Ordering::Relaxed);
440    }
441
442    pub fn record_deadline(&self) {
443        self.ops_deadline.fetch_add(1, Ordering::Relaxed);
444    }
445
446    pub fn record_model_load(&self) {
447        self.model_loads.fetch_add(1, Ordering::Relaxed);
448    }
449
450    pub fn record_cache_hit(&self) {
451        self.cache_hits.fetch_add(1, Ordering::Relaxed);
452    }
453
454    pub fn record_cache_miss(&self) {
455        self.cache_misses.fetch_add(1, Ordering::Relaxed);
456    }
457
458    pub fn record_busy(&self) {
459        self.busy_rejections.fetch_add(1, Ordering::Relaxed);
460    }
461
462    pub fn record_overload(&self) {
463        self.overload_rejections.fetch_add(1, Ordering::Relaxed);
464    }
465
466    pub fn record_queue_wait(&self, d: Duration) {
467        self.queue_wait_ms_total
468            .fetch_add(d.as_millis() as u64, Ordering::Relaxed);
469    }
470
471    pub fn record_remote_error(&self) {
472        self.remote_errors.fetch_add(1, Ordering::Relaxed);
473    }
474
475    pub fn record_decoded_bytes(&self, n: u64) {
476        self.decoded_bytes_total.fetch_add(n, Ordering::Relaxed);
477    }
478
479    pub fn record_upload_bytes(&self, n: u64) {
480        self.upload_bytes_total.fetch_add(n, Ordering::Relaxed);
481    }
482
483    pub fn record_response_bytes(&self, n: u64) {
484        self.response_bytes_total.fetch_add(n, Ordering::Relaxed);
485    }
486
487    pub fn record_tts_chars(&self, n: u64) {
488        self.tts_chars_total.fetch_add(n, Ordering::Relaxed);
489    }
490
491    pub fn record_tts_chunks(&self, n: u64) {
492        self.tts_chunks_total.fetch_add(n, Ordering::Relaxed);
493    }
494
495    pub fn record_long_form_chunk_attempted(&self) {
496        self.long_form_chunks_attempted
497            .fetch_add(1, Ordering::Relaxed);
498    }
499
500    pub fn record_long_form_chunk_completed(&self) {
501        self.long_form_chunks_completed
502            .fetch_add(1, Ordering::Relaxed);
503    }
504
505    pub fn record_long_form_chunk_failed(&self) {
506        self.long_form_chunks_failed.fetch_add(1, Ordering::Relaxed);
507    }
508
509    pub fn record_batch_item_succeeded(&self) {
510        self.batch_items_succeeded.fetch_add(1, Ordering::Relaxed);
511    }
512
513    pub fn record_batch_item_failed(&self) {
514        self.batch_items_failed.fetch_add(1, Ordering::Relaxed);
515    }
516
517    pub fn record_terminal(&self, cat: TerminalCategory) {
518        match cat {
519            TerminalCategory::Completed => {
520                // completed also recorded via record_complete in many paths
521            }
522            TerminalCategory::Failed => self.record_failed(),
523            TerminalCategory::Cancelled => self.record_cancelled(),
524            TerminalCategory::Deadline => self.record_deadline(),
525            TerminalCategory::Overload => self.record_overload(),
526            TerminalCategory::Busy => self.record_busy(),
527        }
528    }
529
530    pub fn snapshot(&self) -> MetricsSnapshot {
531        MetricsSnapshot {
532            schema_version: METRICS_SCHEMA_VERSION,
533            scope: self.scope,
534            ops_started: self.ops_started.load(Ordering::Relaxed),
535            ops_completed: self.ops_completed.load(Ordering::Relaxed),
536            ops_failed: self.ops_failed.load(Ordering::Relaxed),
537            ops_cancelled: self.ops_cancelled.load(Ordering::Relaxed),
538            ops_deadline: self.ops_deadline.load(Ordering::Relaxed),
539            queue_wait_ms_total: self.queue_wait_ms_total.load(Ordering::Relaxed),
540            inference_ms_total: self.inference_ms_total.load(Ordering::Relaxed),
541            encode_ms_total: self.encode_ms_total.load(Ordering::Relaxed),
542            network_ms_total: self.network_ms_total.load(Ordering::Relaxed),
543            model_loads: self.model_loads.load(Ordering::Relaxed),
544            cache_hits: self.cache_hits.load(Ordering::Relaxed),
545            cache_misses: self.cache_misses.load(Ordering::Relaxed),
546            cache_evictions: self.cache_evictions.load(Ordering::Relaxed),
547            remote_errors: self.remote_errors.load(Ordering::Relaxed),
548            remote_auth_errors: self.remote_auth_errors.load(Ordering::Relaxed),
549            remote_rate_limits: self.remote_rate_limits.load(Ordering::Relaxed),
550            remote_quota_errors: self.remote_quota_errors.load(Ordering::Relaxed),
551            remote_network_errors: self.remote_network_errors.load(Ordering::Relaxed),
552            remote_invalid_payload: self.remote_invalid_payload.load(Ordering::Relaxed),
553            upload_bytes_total: self.upload_bytes_total.load(Ordering::Relaxed),
554            response_bytes_total: self.response_bytes_total.load(Ordering::Relaxed),
555            decoded_bytes_total: self.decoded_bytes_total.load(Ordering::Relaxed),
556            long_form_chunks_attempted: self.long_form_chunks_attempted.load(Ordering::Relaxed),
557            long_form_chunks_completed: self.long_form_chunks_completed.load(Ordering::Relaxed),
558            long_form_chunks_failed: self.long_form_chunks_failed.load(Ordering::Relaxed),
559            tts_chunks_total: self.tts_chunks_total.load(Ordering::Relaxed),
560            tts_chars_total: self.tts_chars_total.load(Ordering::Relaxed),
561            batch_items_succeeded: self.batch_items_succeeded.load(Ordering::Relaxed),
562            batch_items_failed: self.batch_items_failed.load(Ordering::Relaxed),
563            batch_stale_reprocess: self.batch_stale_reprocess.load(Ordering::Relaxed),
564            output_tx_success: self.output_tx_success.load(Ordering::Relaxed),
565            output_tx_failure: self.output_tx_failure.load(Ordering::Relaxed),
566            busy_rejections: self.busy_rejections.load(Ordering::Relaxed),
567            overload_rejections: self.overload_rejections.load(Ordering::Relaxed),
568            events_dropped: self.events_dropped.load(Ordering::Relaxed),
569        }
570    }
571}
572
573/// Serializable metrics snapshot (no averages labelled as p95).
574#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
575pub struct MetricsSnapshot {
576    pub schema_version: u32,
577    #[serde(default = "default_process_global")]
578    pub scope: MetricsScope,
579    pub ops_started: u64,
580    pub ops_completed: u64,
581    pub ops_failed: u64,
582    pub ops_cancelled: u64,
583    pub ops_deadline: u64,
584    pub queue_wait_ms_total: u64,
585    pub inference_ms_total: u64,
586    #[serde(default)]
587    pub encode_ms_total: u64,
588    #[serde(default)]
589    pub network_ms_total: u64,
590    pub model_loads: u64,
591    #[serde(default)]
592    pub cache_hits: u64,
593    #[serde(default)]
594    pub cache_misses: u64,
595    #[serde(default)]
596    pub cache_evictions: u64,
597    pub remote_errors: u64,
598    #[serde(default)]
599    pub remote_auth_errors: u64,
600    #[serde(default)]
601    pub remote_rate_limits: u64,
602    #[serde(default)]
603    pub remote_quota_errors: u64,
604    #[serde(default)]
605    pub remote_network_errors: u64,
606    #[serde(default)]
607    pub remote_invalid_payload: u64,
608    #[serde(default)]
609    pub upload_bytes_total: u64,
610    #[serde(default)]
611    pub response_bytes_total: u64,
612    #[serde(default)]
613    pub decoded_bytes_total: u64,
614    #[serde(default)]
615    pub long_form_chunks_attempted: u64,
616    #[serde(default)]
617    pub long_form_chunks_completed: u64,
618    #[serde(default)]
619    pub long_form_chunks_failed: u64,
620    #[serde(default)]
621    pub tts_chunks_total: u64,
622    #[serde(default)]
623    pub tts_chars_total: u64,
624    #[serde(default)]
625    pub batch_items_succeeded: u64,
626    #[serde(default)]
627    pub batch_items_failed: u64,
628    #[serde(default)]
629    pub batch_stale_reprocess: u64,
630    #[serde(default)]
631    pub output_tx_success: u64,
632    #[serde(default)]
633    pub output_tx_failure: u64,
634    pub busy_rejections: u64,
635    pub overload_rejections: u64,
636    #[serde(default)]
637    pub events_dropped: u64,
638}
639
640fn default_process_global() -> MetricsScope {
641    MetricsScope::ProcessGlobal
642}
643
644/// Redacted diagnostic bundle for support / `aurum doctor --json` enrichment.
645#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
646pub struct DiagnosticBundle {
647    pub schema_version: u32,
648    pub request_id: Option<String>,
649    pub operation: Option<String>,
650    pub metrics: MetricsSnapshot,
651    /// Never contains API keys or raw audio.
652    pub notes: Vec<String>,
653}
654
655impl DiagnosticBundle {
656    pub fn from_metrics(metrics: &Metrics) -> Self {
657        Self {
658            schema_version: METRICS_SCHEMA_VERSION,
659            request_id: None,
660            operation: None,
661            metrics: metrics.snapshot(),
662            notes: vec![
663                "payloads (PCM/text/API keys) are never included".into(),
664                "counters are process-local or engine-local best-effort".into(),
665                "totals are sums; averages are not p95".into(),
666            ],
667        }
668    }
669
670    pub fn with_request_id(mut self, id: u64) -> Self {
671        self.request_id = Some(id.to_string());
672        self
673    }
674}
675
676// ---------------------------------------------------------------------------
677// Terminal guard (exactly one outcome)
678// ---------------------------------------------------------------------------
679
680/// Ensures a terminal outcome is recorded at most once (including drop).
681pub struct TerminalGuard {
682    metrics: Arc<Metrics>,
683    request_id: u64,
684    operation: OpKind,
685    scope: MetricsScope,
686    finished: bool,
687    start: Instant,
688}
689
690impl TerminalGuard {
691    pub fn start(metrics: Arc<Metrics>, request_id: u64, operation: OpKind) -> Self {
692        metrics.record_start();
693        let scope = metrics.scope();
694        metrics.emit(OpEvent::stage(request_id, operation, OpStage::Start, scope));
695        Self {
696            metrics,
697            request_id,
698            operation,
699            scope,
700            finished: false,
701            start: Instant::now(),
702        }
703    }
704
705    pub fn finish(&mut self, cat: TerminalCategory, retryable: bool) {
706        if self.finished {
707            return;
708        }
709        self.finished = true;
710        match cat {
711            TerminalCategory::Completed => {
712                self.metrics.record_complete(self.start.elapsed());
713            }
714            other => self.metrics.record_terminal(other),
715        }
716        self.metrics.emit(
717            OpEvent::stage(
718                self.request_id,
719                self.operation,
720                OpStage::Terminal,
721                self.scope,
722            )
723            .with_elapsed_ms(self.start.elapsed().as_millis() as u64)
724            .with_terminal(cat, retryable),
725        );
726    }
727
728    pub fn request_id(&self) -> u64 {
729        self.request_id
730    }
731
732    pub fn elapsed(&self) -> Duration {
733        self.start.elapsed()
734    }
735}
736
737impl Drop for TerminalGuard {
738    fn drop(&mut self) {
739        // Panic / early return: count as failed once.
740        if !self.finished {
741            self.finish(TerminalCategory::Failed, false);
742        }
743    }
744}
745
746/// Simple wall-clock timer for inference spans.
747#[derive(Debug)]
748pub struct SpanTimer {
749    start: Instant,
750}
751
752impl SpanTimer {
753    pub fn start() -> Self {
754        Self {
755            start: Instant::now(),
756        }
757    }
758
759    pub fn elapsed(&self) -> Duration {
760        self.start.elapsed()
761    }
762}
763
764fn process_metrics_arc() -> Arc<Metrics> {
765    use once_cell::sync::Lazy;
766    static SHARED: Lazy<Arc<Metrics>> = Lazy::new(|| Arc::new(Metrics::new()));
767    Arc::clone(&SHARED)
768}
769
770/// Process-global metrics (shared by CLI/FFI).
771pub fn process_metrics() -> Arc<Metrics> {
772    process_metrics_arc()
773}
774
775/// Forbidden substrings that must never appear in public observability output.
776pub const PRIVACY_CANARY_MARKERS: &[&str] = &[
777    "sk-test-secret-key",
778    "BEGIN_PRIVATE_AUDIO",
779    "USER_TRANSCRIPT_PAYLOAD",
780    "SYNTHESIS_TEXT_SECRET",
781    "/Users/private/home/",
782    "Authorization: Bearer",
783];
784
785/// Scan JSON/text for privacy canary markers.
786pub fn privacy_scan(text: &str) -> Result<(), String> {
787    for m in PRIVACY_CANARY_MARKERS {
788        if text.contains(m) {
789            return Err(format!("privacy canary hit: marker present: {m}"));
790        }
791    }
792    Ok(())
793}
794
795#[cfg(test)]
796mod tests {
797    use super::*;
798
799    #[test]
800    fn metrics_roundtrip() {
801        let m = Metrics::new();
802        m.record_start();
803        m.record_complete(Duration::from_millis(12));
804        let s = m.snapshot();
805        assert_eq!(s.ops_started, 1);
806        assert_eq!(s.ops_completed, 1);
807        assert!(s.inference_ms_total >= 12);
808        assert_eq!(s.schema_version, METRICS_SCHEMA_VERSION);
809        let bundle = DiagnosticBundle::from_metrics(&m);
810        assert!(serde_json::to_string(&bundle)
811            .unwrap()
812            .contains("ops_started"));
813    }
814
815    #[test]
816    fn terminal_guard_once() {
817        let m = Arc::new(Metrics::engine_local());
818        let sink = Arc::new(BoundedEventSink::new(32));
819        m.set_event_sink(Some(sink.clone()));
820        {
821            let mut g = TerminalGuard::start(m.clone(), 42, OpKind::Stt);
822            g.finish(TerminalCategory::Completed, false);
823            g.finish(TerminalCategory::Failed, false); // no-op
824        }
825        assert_eq!(m.snapshot().ops_started, 1);
826        assert_eq!(m.snapshot().ops_completed, 1);
827        assert_eq!(m.snapshot().ops_failed, 0);
828        let events = sink.drain();
829        assert!(events.iter().any(|e| e.stage == OpStage::Start));
830        assert_eq!(events.iter().filter(|e| e.terminal.is_some()).count(), 1);
831        assert!(events.iter().all(|e| e.request_id == 42));
832    }
833
834    #[test]
835    fn terminal_guard_drop_counts_failed() {
836        let m = Arc::new(Metrics::new());
837        {
838            let _g = TerminalGuard::start(m.clone(), 1, OpKind::Tts);
839        }
840        assert_eq!(m.snapshot().ops_failed, 1);
841    }
842
843    #[test]
844    fn bounded_sink_drops() {
845        let sink = BoundedEventSink::new(2);
846        for i in 0..5 {
847            sink.emit(OpEvent::stage(
848                i,
849                OpKind::Stt,
850                OpStage::Start,
851                MetricsScope::ProcessGlobal,
852            ));
853        }
854        assert_eq!(sink.len(), 2);
855        assert_eq!(sink.dropped(), 3);
856    }
857
858    #[test]
859    fn emit_does_not_hold_metrics_mutex_during_sink_callback() {
860        use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
861        use std::sync::Mutex as StdMutex;
862
863        /// Sink that re-enters Metrics::set_event_sink while handling emit.
864        struct ReentrantSink {
865            metrics: Arc<Metrics>,
866            hit: Arc<AtomicBool>,
867            // Keep a strong ref so Drop of prior sink does not race.
868            _guard: StdMutex<()>,
869        }
870
871        impl EventSink for ReentrantSink {
872            fn emit(&self, _event: OpEvent) {
873                self.hit.store(true, AtomicOrdering::SeqCst);
874                // Would deadlock if Metrics::emit held event_sink while calling us.
875                self.metrics.set_event_sink(None);
876            }
877        }
878
879        let m = Arc::new(Metrics::new());
880        let hit = Arc::new(AtomicBool::new(false));
881        let sink: Arc<dyn EventSink> = Arc::new(ReentrantSink {
882            metrics: Arc::clone(&m),
883            hit: Arc::clone(&hit),
884            _guard: StdMutex::new(()),
885        });
886        m.set_event_sink(Some(sink));
887        m.emit(OpEvent::stage(
888            1,
889            OpKind::Stt,
890            OpStage::Start,
891            MetricsScope::EngineLocal,
892        ));
893        assert!(hit.load(AtomicOrdering::SeqCst));
894    }
895
896    #[test]
897    fn distinct_outcomes() {
898        let m = Metrics::new();
899        m.record_cancelled();
900        m.record_deadline();
901        m.record_overload();
902        m.record_busy();
903        m.record_failed();
904        let s = m.snapshot();
905        assert_eq!(s.ops_cancelled, 1);
906        assert_eq!(s.ops_deadline, 1);
907        assert_eq!(s.overload_rejections, 1);
908        assert_eq!(s.busy_rejections, 1);
909        assert_eq!(s.ops_failed, 1);
910    }
911
912    #[test]
913    fn engine_vs_process_scope() {
914        let eng = Metrics::engine_local();
915        let proc = Metrics::new();
916        assert_eq!(eng.scope(), MetricsScope::EngineLocal);
917        assert_eq!(proc.scope(), MetricsScope::ProcessGlobal);
918        eng.record_start();
919        assert_eq!(eng.snapshot().ops_started, 1);
920        assert_eq!(proc.snapshot().ops_started, 0);
921    }
922
923    #[test]
924    fn privacy_canary_on_event_and_snapshot() {
925        let m = Metrics::new();
926        m.record_start();
927        let event = OpEvent::stage(
928            7,
929            OpKind::Stt,
930            OpStage::Inference,
931            MetricsScope::EngineLocal,
932        )
933        .with_provider("local")
934        .with_model("base");
935        let json = event.to_json().unwrap();
936        privacy_scan(&json).unwrap();
937        privacy_scan(&serde_json::to_string(&m.snapshot()).unwrap()).unwrap();
938        // Inject markers into notes would fail — bundle notes are fixed.
939        let mut bad = json;
940        bad.push_str("sk-test-secret-key");
941        assert!(privacy_scan(&bad).is_err());
942    }
943
944    #[test]
945    fn event_schema_stable() {
946        let e = OpEvent::stage(
947            1,
948            OpKind::Cleanup,
949            OpStage::Normalize,
950            MetricsScope::ProcessGlobal,
951        );
952        let j1 = e.to_json().unwrap();
953        let j2 = e.to_json().unwrap();
954        assert_eq!(j1, j2);
955        assert!(j1.contains("schema_version"));
956    }
957}