1use std::fmt;
15use std::sync::atomic::{
16 AtomicBool, AtomicU64,
17 Ordering::{Acquire, Relaxed, Release},
18};
19use std::time::{Duration, Instant};
20
21use serde::{Deserialize, Serialize};
22use serde_json::json;
23
24pub const LATENCY_BUCKETS_MS: [u64; 17] = [
29 50, 100, 150, 200, 300, 400, 500, 650, 800, 1000, 1300, 1600, 2000, 2500, 3000, 4000, 5000,
30];
31
32pub const LATENCY_RECENT_WINDOW: usize = 256;
38
39const BUCKET_COUNT: usize = LATENCY_BUCKETS_MS.len() + 1;
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
43pub struct LatencyBucket {
44 pub upper_ms: Option<u64>,
46 pub count: u64,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
55pub struct LatencyStats {
56 pub count: u64,
58 pub last_ms: u64,
60 pub min_ms: u64,
62 pub max_ms: u64,
64 pub mean_ms: u64,
66 pub p50_ms: u64,
68 pub p90_ms: u64,
70 pub p99_ms: u64,
72 pub histogram: Vec<LatencyBucket>,
75}
76
77impl fmt::Display for LatencyStats {
78 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
80 if self.count == 0 {
81 return write!(f, "turns=0 (no response measured yet)");
82 }
83 write!(
84 f,
85 "turns={} last={}ms p50={}ms p90={}ms p99={}ms min={}ms max={}ms",
86 self.count,
87 self.last_ms,
88 self.p50_ms,
89 self.p90_ms,
90 self.p99_ms,
91 self.min_ms,
92 self.max_ms
93 )
94 }
95}
96
97struct LatencyRecorder {
100 last_ns: AtomicU64,
101 sum_ns: AtomicU64,
102 count: AtomicU64,
103 min_ns: AtomicU64,
104 max_ns: AtomicU64,
105 buckets: [AtomicU64; BUCKET_COUNT],
106 recent: [AtomicU64; LATENCY_RECENT_WINDOW],
107 recent_next: AtomicU64,
108}
109
110impl LatencyRecorder {
111 fn new() -> Self {
112 Self {
113 last_ns: AtomicU64::new(0),
114 sum_ns: AtomicU64::new(0),
115 count: AtomicU64::new(0),
116 min_ns: AtomicU64::new(u64::MAX),
117 max_ns: AtomicU64::new(0),
118 buckets: std::array::from_fn(|_| AtomicU64::new(0)),
119 recent: std::array::from_fn(|_| AtomicU64::new(0)),
120 recent_next: AtomicU64::new(0),
121 }
122 }
123
124 #[inline]
133 fn record(&self, latency_ns: u64) {
134 self.last_ns.store(latency_ns, Relaxed);
135 self.sum_ns.fetch_add(latency_ns, Relaxed);
136 self.min_ns.fetch_min(latency_ns, Relaxed);
137 self.max_ns.fetch_max(latency_ns, Relaxed);
138
139 let ms = latency_ns / 1_000_000;
140 let bucket = LATENCY_BUCKETS_MS
141 .iter()
142 .position(|&upper| ms <= upper)
143 .unwrap_or(BUCKET_COUNT - 1);
144 self.buckets[bucket].fetch_add(1, Relaxed);
145
146 let slot = self.recent_next.fetch_add(1, Relaxed) as usize % LATENCY_RECENT_WINDOW;
148 self.recent[slot].store(latency_ns, Relaxed);
149
150 self.count.fetch_add(1, Release);
151 }
152
153 fn stats(&self) -> LatencyStats {
154 let count = self.count.load(Acquire);
159 let histogram = self
160 .buckets
161 .iter()
162 .enumerate()
163 .map(|(i, b)| LatencyBucket {
164 upper_ms: LATENCY_BUCKETS_MS.get(i).copied(),
165 count: b.load(Relaxed),
166 })
167 .collect();
168 if count == 0 {
169 return LatencyStats {
170 histogram,
171 ..LatencyStats::default()
172 };
173 }
174
175 let filled = (count as usize).min(LATENCY_RECENT_WINDOW);
182 let mut recent: Vec<u64> = self.recent[..filled]
183 .iter()
184 .map(|s| s.load(Relaxed))
185 .collect();
186 recent.sort_unstable();
187 let pct = |p: usize| -> u64 {
189 let rank = (p * recent.len()).div_ceil(100).max(1);
190 recent.get(rank - 1).copied().unwrap_or(0) / 1_000_000
191 };
192
193 LatencyStats {
194 count,
195 last_ms: self.last_ns.load(Relaxed) / 1_000_000,
196 min_ms: self.min_ns.load(Relaxed) / 1_000_000,
197 max_ms: self.max_ns.load(Relaxed) / 1_000_000,
198 mean_ms: self.sum_ns.load(Relaxed) / count / 1_000_000,
199 p50_ms: pct(50),
200 p90_ms: pct(90),
201 p99_ms: pct(99),
202 histogram,
203 }
204 }
205}
206
207pub struct SessionTelemetry {
216 start: Instant,
217
218 audio_chunks_out: AtomicU64,
220 audio_bytes_out: AtomicU64,
221
222 interruptions: AtomicU64,
224
225 vad_end_ns: AtomicU64,
228 awaiting_response: AtomicBool,
229 text_send_ns: AtomicU64,
231 awaiting_text_response: AtomicBool,
232 latency: LatencyRecorder,
235
236 turn_complete_count: AtomicU64,
238 last_turn_start_ns: AtomicU64,
239 turn_duration_sum_ns: AtomicU64,
240 turn_duration_count: AtomicU64,
241
242 total_token_count: AtomicU64,
245 prompt_token_count: AtomicU64,
247 response_token_count: AtomicU64,
249 cached_content_token_count: AtomicU64,
251 thoughts_token_count: AtomicU64,
253}
254
255impl SessionTelemetry {
256 pub fn new() -> Self {
258 Self {
259 start: Instant::now(),
260 audio_chunks_out: AtomicU64::new(0),
261 audio_bytes_out: AtomicU64::new(0),
262 interruptions: AtomicU64::new(0),
263 vad_end_ns: AtomicU64::new(0),
264 awaiting_response: AtomicBool::new(false),
265 text_send_ns: AtomicU64::new(0),
266 awaiting_text_response: AtomicBool::new(false),
267 latency: LatencyRecorder::new(),
268 turn_complete_count: AtomicU64::new(0),
269 last_turn_start_ns: AtomicU64::new(0),
270 turn_duration_sum_ns: AtomicU64::new(0),
271 turn_duration_count: AtomicU64::new(0),
272 total_token_count: AtomicU64::new(0),
273 prompt_token_count: AtomicU64::new(0),
274 response_token_count: AtomicU64::new(0),
275 cached_content_token_count: AtomicU64::new(0),
276 thoughts_token_count: AtomicU64::new(0),
277 }
278 }
279
280 #[inline]
288 pub fn record_audio_out(&self, byte_len: usize) -> Option<Duration> {
289 self.audio_chunks_out.fetch_add(1, Relaxed);
290 self.audio_bytes_out.fetch_add(byte_len as u64, Relaxed);
291
292 let text = self.record_text_response_latency();
294
295 if self
298 .awaiting_response
299 .compare_exchange(true, false, Relaxed, Relaxed)
300 .is_ok()
301 {
302 let now_ns = self.elapsed_ns();
303 let vad_end = self.vad_end_ns.load(Relaxed);
304 if now_ns > vad_end && vad_end > 0 {
305 let latency = now_ns - vad_end;
306 self.latency.record(latency);
307 return Some(Duration::from_nanos(latency));
308 }
309 }
310 text
311 }
312
313 #[inline]
315 pub fn record_vad_end(&self) {
316 self.vad_end_ns.store(self.elapsed_ns(), Relaxed);
317 self.awaiting_response.store(true, Relaxed);
318 }
319
320 #[inline]
322 pub fn record_text_send(&self) {
323 self.text_send_ns.store(self.elapsed_ns(), Relaxed);
324 self.awaiting_text_response.store(true, Relaxed);
325 }
326
327 #[inline]
330 fn record_text_response_latency(&self) -> Option<Duration> {
331 if self
332 .awaiting_text_response
333 .compare_exchange(true, false, Relaxed, Relaxed)
334 .is_ok()
335 {
336 let now_ns = self.elapsed_ns();
337 let send_ns = self.text_send_ns.load(Relaxed);
338 if now_ns > send_ns && send_ns > 0 {
339 let latency = now_ns - send_ns;
340 self.latency.record(latency);
341 return Some(Duration::from_nanos(latency));
342 }
343 }
344 None
345 }
346
347 #[inline]
352 pub fn record_text_out(&self) -> Option<Duration> {
353 self.record_text_response_latency()
354 }
355
356 #[inline]
358 pub fn record_interruption(&self) {
359 self.interruptions.fetch_add(1, Relaxed);
360 }
361
362 #[inline]
364 pub fn record_turn_complete(&self) {
365 self.turn_complete_count.fetch_add(1, Relaxed);
366 let now = self.elapsed_ns();
367 let turn_start = self.last_turn_start_ns.swap(now, Relaxed);
368 if turn_start > 0 {
369 let duration = now.saturating_sub(turn_start);
370 self.turn_duration_sum_ns.fetch_add(duration, Relaxed);
371 self.turn_duration_count.fetch_add(1, Relaxed);
372 }
373 }
374
375 #[inline]
377 pub fn record_usage(
378 &self,
379 total: Option<u32>,
380 prompt: Option<u32>,
381 response: Option<u32>,
382 cached: Option<u32>,
383 thoughts: Option<u32>,
384 ) {
385 if let Some(v) = total {
386 self.total_token_count.store(v as u64, Relaxed);
387 }
388 if let Some(v) = prompt {
389 self.prompt_token_count.store(v as u64, Relaxed);
390 }
391 if let Some(v) = response {
392 self.response_token_count.store(v as u64, Relaxed);
393 }
394 if let Some(v) = cached {
395 self.cached_content_token_count.store(v as u64, Relaxed);
396 }
397 if let Some(v) = thoughts {
398 self.thoughts_token_count.store(v as u64, Relaxed);
399 }
400 }
401
402 #[inline]
404 pub fn mark_turn_start(&self) {
405 let now = self.elapsed_ns();
406 self.last_turn_start_ns
408 .compare_exchange(0, now, Relaxed, Relaxed)
409 .ok();
410 }
411
412 pub fn latency(&self) -> LatencyStats {
420 self.latency.stats()
421 }
422
423 pub fn snapshot(&self) -> serde_json::Value {
429 let elapsed = self.start.elapsed();
430 let elapsed_secs = elapsed.as_secs_f64();
431
432 let chunks = self.audio_chunks_out.load(Relaxed);
433 let bytes = self.audio_bytes_out.load(Relaxed);
434 let latency = self.latency.stats();
435
436 let turn_count = self.turn_duration_count.load(Relaxed);
437 let turn_complete_count = self.turn_complete_count.load(Relaxed);
438 let avg_turn_ms = if turn_count > 0 {
439 self.turn_duration_sum_ns.load(Relaxed) / turn_count / 1_000_000
440 } else {
441 0
442 };
443
444 let throughput_kbps = if elapsed_secs > 0.0 {
446 (bytes as f64 / 1024.0) / elapsed_secs
447 } else {
448 0.0
449 };
450
451 let total_tokens = self.total_token_count.load(Relaxed);
452 let prompt_tokens = self.prompt_token_count.load(Relaxed);
453 let response_tokens = self.response_token_count.load(Relaxed);
454 let cached_tokens = self.cached_content_token_count.load(Relaxed);
455 let thoughts_tokens = self.thoughts_token_count.load(Relaxed);
456
457 json!({
458 "uptime_secs": elapsed.as_secs(),
459 "audio_chunks_out": chunks,
460 "audio_kbytes_out": bytes / 1024,
461 "audio_throughput_kbps": (throughput_kbps * 10.0).round() / 10.0,
462 "interruptions": self.interruptions.load(Relaxed),
463 "last_response_latency_ms": latency.last_ms,
464 "avg_response_latency_ms": latency.mean_ms,
465 "min_response_latency_ms": latency.min_ms,
466 "max_response_latency_ms": latency.max_ms,
467 "response_count": latency.count,
468 "response_latency": latency,
469 "turn_count": turn_complete_count,
470 "avg_turn_duration_ms": avg_turn_ms,
471 "total_token_count": total_tokens,
472 "prompt_token_count": prompt_tokens,
473 "response_token_count": response_tokens,
474 "cached_content_token_count": cached_tokens,
475 "thoughts_token_count": thoughts_tokens,
476 })
477 }
478
479 #[inline]
480 fn elapsed_ns(&self) -> u64 {
481 self.start.elapsed().as_nanos() as u64
482 }
483}
484
485impl Default for SessionTelemetry {
486 fn default() -> Self {
487 Self::new()
488 }
489}
490
491#[cfg(test)]
492mod tests {
493 use super::*;
494
495 #[test]
496 fn new_snapshot_is_zeroed() {
497 let t = SessionTelemetry::new();
498 let snap = t.snapshot();
499 assert_eq!(snap["audio_chunks_out"], 0);
500 assert_eq!(snap["interruptions"], 0);
501 assert_eq!(snap["last_response_latency_ms"], 0);
502 assert_eq!(snap["response_count"], 0);
503 assert_eq!(snap["turn_count"], 0);
504 assert_eq!(snap["response_latency"]["count"], 0);
505 assert_eq!(
506 t.latency(),
507 LatencyStats {
508 histogram: t.latency().histogram.clone(),
509 ..LatencyStats::default()
510 }
511 );
512 assert_eq!(
513 t.latency().to_string(),
514 "turns=0 (no response measured yet)"
515 );
516 }
517
518 #[test]
519 fn audio_counters_accumulate() {
520 let t = SessionTelemetry::new();
521 t.record_audio_out(480);
522 t.record_audio_out(480);
523 t.record_audio_out(480);
524 let snap = t.snapshot();
525 assert_eq!(snap["audio_chunks_out"], 3);
526 }
527
528 #[test]
529 fn interruption_counter() {
530 let t = SessionTelemetry::new();
531 t.record_interruption();
532 t.record_interruption();
533 assert_eq!(t.snapshot()["interruptions"], 2);
534 }
535
536 #[test]
537 fn turn_complete_counter_is_independent_of_latency() {
538 let t = SessionTelemetry::new();
539 t.record_turn_complete();
540 t.record_turn_complete();
541
542 let snap = t.snapshot();
543 assert_eq!(snap["turn_count"], 2);
544 assert_eq!(snap["response_count"], 0);
545 }
546
547 #[test]
548 fn latency_tracking() {
549 let t = SessionTelemetry::new();
550 t.record_vad_end();
552 std::thread::sleep(std::time::Duration::from_millis(10));
553 let first = t.record_audio_out(480);
554 assert!(
555 first.is_some(),
556 "first chunk after VAD end reports the latency"
557 );
558 assert!(t.record_audio_out(480).is_none());
560 assert!(t.record_audio_out(480).is_none());
561
562 let snap = t.snapshot();
563 assert_eq!(snap["response_count"], 1);
564 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
566 assert_eq!(snap["response_latency"]["count"], 1);
567 }
568
569 #[test]
570 fn multiple_turns_average_latency() {
571 let t = SessionTelemetry::new();
572
573 t.record_vad_end();
575 std::thread::sleep(std::time::Duration::from_millis(10));
576 t.record_audio_out(480);
577
578 t.record_vad_end();
580 std::thread::sleep(std::time::Duration::from_millis(10));
581 t.record_audio_out(480);
582
583 let snap = t.snapshot();
584 assert_eq!(snap["response_count"], 2);
585 assert!(snap["avg_response_latency_ms"].as_u64().unwrap() >= 5);
586 }
587
588 #[test]
589 fn text_input_latency_via_text_out() {
590 let t = SessionTelemetry::new();
591 t.record_text_send();
593 std::thread::sleep(std::time::Duration::from_millis(10));
594 assert!(t.record_text_out().is_some());
595 assert!(t.record_text_out().is_none());
597
598 let snap = t.snapshot();
599 assert_eq!(snap["response_count"], 1);
600 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
601 }
602
603 #[test]
604 fn text_input_latency_via_audio_out() {
605 let t = SessionTelemetry::new();
606 t.record_text_send();
608 std::thread::sleep(std::time::Duration::from_millis(10));
609 assert!(t.record_audio_out(480).is_some());
610
611 let snap = t.snapshot();
612 assert_eq!(snap["response_count"], 1);
614 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
615 }
616
617 #[test]
618 fn mixed_voice_and_text_turns() {
619 let t = SessionTelemetry::new();
620
621 t.record_vad_end();
623 std::thread::sleep(std::time::Duration::from_millis(10));
624 t.record_audio_out(480);
625
626 t.record_text_send();
628 std::thread::sleep(std::time::Duration::from_millis(10));
629 t.record_text_out();
630
631 let snap = t.snapshot();
632 assert_eq!(snap["response_count"], 2);
633 }
634
635 fn ms(n: u64) -> u64 {
638 n * 1_000_000
639 }
640
641 #[test]
642 fn recorder_scalars_and_percentiles() {
643 let r = LatencyRecorder::new();
644 for v in [300, 100, 500, 200, 400] {
645 r.record(ms(v));
646 }
647 let s = r.stats();
648 assert_eq!(s.count, 5);
649 assert_eq!(s.last_ms, 400);
650 assert_eq!(s.min_ms, 100);
651 assert_eq!(s.max_ms, 500);
652 assert_eq!(s.mean_ms, 300);
653 assert_eq!(s.p50_ms, 300);
656 assert_eq!(s.p90_ms, 500);
657 assert_eq!(s.p99_ms, 500);
658 assert_eq!(
659 s.to_string(),
660 "turns=5 last=400ms p50=300ms p90=500ms p99=500ms min=100ms max=500ms"
661 );
662 }
663
664 #[test]
665 fn recorder_first_sample_is_visible_to_percentiles() {
666 let r = LatencyRecorder::new();
669 r.record(ms(420));
670 let s = r.stats();
671 assert_eq!(s.count, 1);
672 assert_eq!((s.p50_ms, s.p90_ms, s.p99_ms), (420, 420, 420));
673 assert_eq!(s.mean_ms, 420);
674 }
675
676 #[test]
677 fn recorder_stats_survives_a_count_without_its_sample() {
678 let r = LatencyRecorder::new();
684 r.count.store(1, Release); let s = r.stats();
686 assert_eq!(s.count, 1);
687 assert_eq!((s.p50_ms, s.p90_ms, s.p99_ms), (0, 0, 0));
688 }
689
690 #[test]
691 fn recorder_snapshots_stay_consistent_under_concurrent_records() {
692 use std::sync::Arc;
696 use std::sync::atomic::AtomicBool;
697
698 let r = Arc::new(LatencyRecorder::new());
699 let done = Arc::new(AtomicBool::new(false));
700
701 let reader = {
702 let r = Arc::clone(&r);
703 let done = Arc::clone(&done);
704 std::thread::spawn(move || {
705 while !done.load(Relaxed) {
706 let s = r.stats();
707 if s.count > 0 {
708 assert!(s.p99_ms >= s.p50_ms);
709 assert!(s.max_ms >= s.p99_ms);
710 assert!(s.min_ms <= s.max_ms);
711 }
712 }
713 })
714 };
715
716 for _ in 0..2_000 {
717 r.record(ms(100));
718 }
719 done.store(true, Relaxed);
720 reader.join().expect("reader saw an inconsistent snapshot");
721 assert_eq!(r.stats().count, 2_000);
722 }
723
724 #[test]
725 fn recorder_histogram_buckets_by_upper_bound() {
726 let r = LatencyRecorder::new();
727 r.record(ms(50)); r.record(ms(51)); r.record(ms(999)); r.record(ms(9_000)); let h = r.stats().histogram;
732 assert_eq!(h.len(), LATENCY_BUCKETS_MS.len() + 1);
733 assert_eq!(
734 h[0],
735 LatencyBucket {
736 upper_ms: Some(50),
737 count: 1
738 }
739 );
740 assert_eq!(
741 h[1],
742 LatencyBucket {
743 upper_ms: Some(100),
744 count: 1
745 }
746 );
747 assert_eq!(
748 h[9],
749 LatencyBucket {
750 upper_ms: Some(1000),
751 count: 1
752 }
753 );
754 assert_eq!(
755 h[h.len() - 1],
756 LatencyBucket {
757 upper_ms: None,
758 count: 1
759 }
760 );
761 assert_eq!(h.iter().map(|b| b.count).sum::<u64>(), 4);
762 }
763
764 #[test]
765 fn recorder_percentiles_use_recent_window_only() {
766 let r = LatencyRecorder::new();
767 for _ in 0..LATENCY_RECENT_WINDOW {
771 r.record(ms(2_000));
772 }
773 for _ in 0..LATENCY_RECENT_WINDOW {
774 r.record(ms(200));
775 }
776 let s = r.stats();
777 assert_eq!(s.count, 2 * LATENCY_RECENT_WINDOW as u64);
778 assert_eq!(s.p50_ms, 200);
779 assert_eq!(s.p99_ms, 200);
780 assert_eq!(s.max_ms, 2_000);
781 assert_eq!(s.min_ms, 200);
782 let slow: u64 = s
783 .histogram
784 .iter()
785 .filter(|b| b.upper_ms == Some(2000))
786 .map(|b| b.count)
787 .sum();
788 assert_eq!(slow, LATENCY_RECENT_WINDOW as u64);
789 }
790
791 #[test]
792 fn stats_round_trip_through_json() {
793 let r = LatencyRecorder::new();
794 r.record(ms(420));
795 let s = r.stats();
796 let v = serde_json::to_value(&s).unwrap();
797 assert_eq!(v["p50_ms"], 420);
798 let back: LatencyStats = serde_json::from_value(v).unwrap();
799 assert_eq!(back, s);
800 }
801}