1use serde::{Deserialize, Serialize};
8use std::collections::VecDeque;
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex};
11use std::time::{Duration, Instant};
12
13pub const METRICS_SCHEMA_VERSION: u32 = 2;
15pub const OP_EVENT_SCHEMA_VERSION: u32 = 1;
17pub const DEFAULT_EVENT_QUEUE_CAP: usize = 256;
19
20#[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#[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#[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#[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#[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 #[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 pub fn to_json(&self) -> Result<String, serde_json::Error> {
208 serde_json::to_string(self)
209 }
210}
211
212pub trait EventSink: Send + Sync {
218 fn emit(&self, event: OpEvent);
219}
220
221#[derive(Debug, Default, Clone, Copy)]
223pub struct NoopEventSink;
224
225impl EventSink for NoopEventSink {
226 fn emit(&self, _event: OpEvent) {}
227}
228
229#[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
286pub 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 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 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 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 }
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#[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#[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 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
676pub 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 if !self.finished {
741 self.finish(TerminalCategory::Failed, false);
742 }
743 }
744}
745
746#[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
770pub fn process_metrics() -> Arc<Metrics> {
772 process_metrics_arc()
773}
774
775pub 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
785pub 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); }
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 struct ReentrantSink {
865 metrics: Arc<Metrics>,
866 hit: Arc<AtomicBool>,
867 _guard: StdMutex<()>,
869 }
870
871 impl EventSink for ReentrantSink {
872 fn emit(&self, _event: OpEvent) {
873 self.hit.store(true, AtomicOrdering::SeqCst);
874 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 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}