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_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 pub fn to_json(&self) -> Result<String, serde_json::Error> {
192 serde_json::to_string(self)
193 }
194}
195
196pub trait EventSink: Send + Sync {
202 fn emit(&self, event: OpEvent);
203}
204
205#[derive(Debug, Default, Clone, Copy)]
207pub struct NoopEventSink;
208
209impl EventSink for NoopEventSink {
210 fn emit(&self, _event: OpEvent) {}
211}
212
213#[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
270pub 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 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 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 }
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#[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#[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 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
613pub 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 if !self.finished {
678 self.finish(TerminalCategory::Failed, false);
679 }
680 }
681}
682
683#[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
707pub fn process_metrics() -> Arc<Metrics> {
709 process_metrics_arc()
710}
711
712pub 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
722pub 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); }
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 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}