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 tokens_by_modality: parking_lot::Mutex<std::collections::BTreeMap<String, [u64; 2]>>,
256}
257
258impl SessionTelemetry {
259 pub fn new() -> Self {
261 Self {
262 start: Instant::now(),
263 audio_chunks_out: AtomicU64::new(0),
264 audio_bytes_out: AtomicU64::new(0),
265 interruptions: AtomicU64::new(0),
266 vad_end_ns: AtomicU64::new(0),
267 awaiting_response: AtomicBool::new(false),
268 text_send_ns: AtomicU64::new(0),
269 awaiting_text_response: AtomicBool::new(false),
270 latency: LatencyRecorder::new(),
271 turn_complete_count: AtomicU64::new(0),
272 last_turn_start_ns: AtomicU64::new(0),
273 turn_duration_sum_ns: AtomicU64::new(0),
274 turn_duration_count: AtomicU64::new(0),
275 total_token_count: AtomicU64::new(0),
276 prompt_token_count: AtomicU64::new(0),
277 response_token_count: AtomicU64::new(0),
278 cached_content_token_count: AtomicU64::new(0),
279 thoughts_token_count: AtomicU64::new(0),
280 tokens_by_modality: parking_lot::Mutex::new(std::collections::BTreeMap::new()),
281 }
282 }
283
284 #[inline]
292 pub fn record_audio_out(&self, byte_len: usize) -> Option<Duration> {
293 self.audio_chunks_out.fetch_add(1, Relaxed);
294 self.audio_bytes_out.fetch_add(byte_len as u64, Relaxed);
295
296 let text = self.record_text_response_latency();
298
299 if self
302 .awaiting_response
303 .compare_exchange(true, false, Relaxed, Relaxed)
304 .is_ok()
305 {
306 let now_ns = self.elapsed_ns();
307 let vad_end = self.vad_end_ns.load(Relaxed);
308 if now_ns > vad_end && vad_end > 0 {
309 let latency = now_ns - vad_end;
310 self.latency.record(latency);
311 gemini_genai_rs::telemetry::metrics::record_response_latency(latency as f64 / 1e6);
312 return Some(Duration::from_nanos(latency));
313 }
314 }
315 text
316 }
317
318 #[inline]
320 pub fn record_vad_end(&self) {
321 self.vad_end_ns.store(self.elapsed_ns(), Relaxed);
322 self.awaiting_response.store(true, Relaxed);
323 }
324
325 #[inline]
327 pub fn record_text_send(&self) {
328 self.text_send_ns.store(self.elapsed_ns(), Relaxed);
329 self.awaiting_text_response.store(true, Relaxed);
330 }
331
332 #[inline]
335 fn record_text_response_latency(&self) -> Option<Duration> {
336 if self
337 .awaiting_text_response
338 .compare_exchange(true, false, Relaxed, Relaxed)
339 .is_ok()
340 {
341 let now_ns = self.elapsed_ns();
342 let send_ns = self.text_send_ns.load(Relaxed);
343 if now_ns > send_ns && send_ns > 0 {
344 let latency = now_ns - send_ns;
345 self.latency.record(latency);
346 gemini_genai_rs::telemetry::metrics::record_response_latency(latency as f64 / 1e6);
347 return Some(Duration::from_nanos(latency));
348 }
349 }
350 None
351 }
352
353 #[inline]
358 pub fn record_text_out(&self) -> Option<Duration> {
359 self.record_text_response_latency()
360 }
361
362 #[inline]
364 pub fn record_interruption(&self) {
365 self.interruptions.fetch_add(1, Relaxed);
366 }
367
368 #[inline]
370 pub fn record_turn_complete(&self) {
371 self.turn_complete_count.fetch_add(1, Relaxed);
372 let now = self.elapsed_ns();
373 let turn_start = self.last_turn_start_ns.swap(now, Relaxed);
374 if turn_start > 0 {
375 let duration = now.saturating_sub(turn_start);
376 self.turn_duration_sum_ns.fetch_add(duration, Relaxed);
377 self.turn_duration_count.fetch_add(1, Relaxed);
378 }
379 }
380
381 #[inline]
383 pub fn record_usage(
384 &self,
385 total: Option<u32>,
386 prompt: Option<u32>,
387 response: Option<u32>,
388 cached: Option<u32>,
389 thoughts: Option<u32>,
390 ) {
391 if let Some(v) = total {
392 self.total_token_count.store(v as u64, Relaxed);
393 }
394 if let Some(v) = prompt {
395 self.prompt_token_count.store(v as u64, Relaxed);
396 }
397 if let Some(v) = response {
398 self.response_token_count.store(v as u64, Relaxed);
399 }
400 if let Some(v) = cached {
401 self.cached_content_token_count.store(v as u64, Relaxed);
402 }
403 if let Some(v) = thoughts {
404 self.thoughts_token_count.store(v as u64, Relaxed);
405 }
406 }
407
408 pub fn record_turn_usage(&self, usage: &gemini_genai_rs::prelude::UsageMetadata) {
412 let mut totals = self.tokens_by_modality.lock();
413 for (index, direction, details) in [
414 (0, "prompt", &usage.prompt_tokens_details),
415 (1, "response", &usage.response_tokens_details),
416 ] {
417 for detail in details {
418 let (Some(modality), Some(count)) = (&detail.modality, detail.token_count) else {
419 continue;
420 };
421 totals.entry(modality.clone()).or_default()[index] += u64::from(count);
422 gemini_genai_rs::telemetry::metrics::record_tokens(
423 direction,
424 modality,
425 count.into(),
426 );
427 }
428 }
429 }
430
431 #[inline]
433 pub fn mark_turn_start(&self) {
434 let now = self.elapsed_ns();
435 self.last_turn_start_ns
437 .compare_exchange(0, now, Relaxed, Relaxed)
438 .ok();
439 }
440
441 pub fn latency(&self) -> LatencyStats {
449 self.latency.stats()
450 }
451
452 pub fn snapshot(&self) -> serde_json::Value {
458 let elapsed = self.start.elapsed();
459 let elapsed_secs = elapsed.as_secs_f64();
460
461 let chunks = self.audio_chunks_out.load(Relaxed);
462 let bytes = self.audio_bytes_out.load(Relaxed);
463 let latency = self.latency.stats();
464
465 let turn_count = self.turn_duration_count.load(Relaxed);
466 let turn_complete_count = self.turn_complete_count.load(Relaxed);
467 let avg_turn_ms = if turn_count > 0 {
468 self.turn_duration_sum_ns.load(Relaxed) / turn_count / 1_000_000
469 } else {
470 0
471 };
472
473 let throughput_kbps = if elapsed_secs > 0.0 {
475 (bytes as f64 / 1024.0) / elapsed_secs
476 } else {
477 0.0
478 };
479
480 let total_tokens = self.total_token_count.load(Relaxed);
481 let prompt_tokens = self.prompt_token_count.load(Relaxed);
482 let response_tokens = self.response_token_count.load(Relaxed);
483 let cached_tokens = self.cached_content_token_count.load(Relaxed);
484 let thoughts_tokens = self.thoughts_token_count.load(Relaxed);
485
486 json!({
487 "uptime_secs": elapsed.as_secs(),
488 "audio_chunks_out": chunks,
489 "audio_kbytes_out": bytes / 1024,
490 "audio_throughput_kbps": (throughput_kbps * 10.0).round() / 10.0,
491 "interruptions": self.interruptions.load(Relaxed),
492 "last_response_latency_ms": latency.last_ms,
493 "avg_response_latency_ms": latency.mean_ms,
494 "min_response_latency_ms": latency.min_ms,
495 "max_response_latency_ms": latency.max_ms,
496 "response_count": latency.count,
497 "response_latency": latency,
498 "turn_count": turn_complete_count,
499 "avg_turn_duration_ms": avg_turn_ms,
500 "total_token_count": total_tokens,
501 "prompt_token_count": prompt_tokens,
502 "response_token_count": response_tokens,
503 "cached_content_token_count": cached_tokens,
504 "thoughts_token_count": thoughts_tokens,
505 "tokens_by_modality": self
506 .tokens_by_modality
507 .lock()
508 .iter()
509 .map(|(modality, [prompt, response])| {
510 (modality.clone(), json!({ "prompt": prompt, "response": response }))
511 })
512 .collect::<serde_json::Map<_, _>>(),
513 })
514 }
515
516 #[inline]
517 fn elapsed_ns(&self) -> u64 {
518 self.start.elapsed().as_nanos() as u64
519 }
520}
521
522impl Default for SessionTelemetry {
523 fn default() -> Self {
524 Self::new()
525 }
526}
527
528#[cfg(test)]
529mod tests {
530 use super::*;
531
532 #[test]
533 fn turn_usage_is_summed_by_modality() {
534 let t = SessionTelemetry::new();
535 let usage: gemini_genai_rs::prelude::UsageMetadata = serde_json::from_value(json!({
536 "promptTokenCount": 30,
537 "promptTokensDetails": [
538 { "modality": "AUDIO", "tokenCount": 20 },
539 { "modality": "TEXT", "tokenCount": 10 }
540 ],
541 "responseTokensDetails": [{ "modality": "AUDIO", "tokenCount": 40 }]
542 }))
543 .unwrap();
544 t.record_turn_usage(&usage);
545 t.record_turn_usage(&usage);
546 let snap = t.snapshot();
547 assert_eq!(
548 snap["tokens_by_modality"]["AUDIO"],
549 json!({ "prompt": 40, "response": 80 })
550 );
551 assert_eq!(
552 snap["tokens_by_modality"]["TEXT"],
553 json!({ "prompt": 20, "response": 0 })
554 );
555 }
556
557 #[test]
558 fn new_snapshot_is_zeroed() {
559 let t = SessionTelemetry::new();
560 let snap = t.snapshot();
561 assert_eq!(snap["audio_chunks_out"], 0);
562 assert_eq!(snap["interruptions"], 0);
563 assert_eq!(snap["last_response_latency_ms"], 0);
564 assert_eq!(snap["response_count"], 0);
565 assert_eq!(snap["turn_count"], 0);
566 assert_eq!(snap["response_latency"]["count"], 0);
567 assert_eq!(
568 t.latency(),
569 LatencyStats {
570 histogram: t.latency().histogram.clone(),
571 ..LatencyStats::default()
572 }
573 );
574 assert_eq!(
575 t.latency().to_string(),
576 "turns=0 (no response measured yet)"
577 );
578 }
579
580 #[test]
581 fn audio_counters_accumulate() {
582 let t = SessionTelemetry::new();
583 t.record_audio_out(480);
584 t.record_audio_out(480);
585 t.record_audio_out(480);
586 let snap = t.snapshot();
587 assert_eq!(snap["audio_chunks_out"], 3);
588 }
589
590 #[test]
591 fn interruption_counter() {
592 let t = SessionTelemetry::new();
593 t.record_interruption();
594 t.record_interruption();
595 assert_eq!(t.snapshot()["interruptions"], 2);
596 }
597
598 #[test]
599 fn turn_complete_counter_is_independent_of_latency() {
600 let t = SessionTelemetry::new();
601 t.record_turn_complete();
602 t.record_turn_complete();
603
604 let snap = t.snapshot();
605 assert_eq!(snap["turn_count"], 2);
606 assert_eq!(snap["response_count"], 0);
607 }
608
609 #[test]
610 fn latency_tracking() {
611 let t = SessionTelemetry::new();
612 t.record_vad_end();
614 std::thread::sleep(std::time::Duration::from_millis(10));
615 let first = t.record_audio_out(480);
616 assert!(
617 first.is_some(),
618 "first chunk after VAD end reports the latency"
619 );
620 assert!(t.record_audio_out(480).is_none());
622 assert!(t.record_audio_out(480).is_none());
623
624 let snap = t.snapshot();
625 assert_eq!(snap["response_count"], 1);
626 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
628 assert_eq!(snap["response_latency"]["count"], 1);
629 }
630
631 #[test]
632 fn multiple_turns_average_latency() {
633 let t = SessionTelemetry::new();
634
635 t.record_vad_end();
637 std::thread::sleep(std::time::Duration::from_millis(10));
638 t.record_audio_out(480);
639
640 t.record_vad_end();
642 std::thread::sleep(std::time::Duration::from_millis(10));
643 t.record_audio_out(480);
644
645 let snap = t.snapshot();
646 assert_eq!(snap["response_count"], 2);
647 assert!(snap["avg_response_latency_ms"].as_u64().unwrap() >= 5);
648 }
649
650 #[test]
651 fn text_input_latency_via_text_out() {
652 let t = SessionTelemetry::new();
653 t.record_text_send();
655 std::thread::sleep(std::time::Duration::from_millis(10));
656 assert!(t.record_text_out().is_some());
657 assert!(t.record_text_out().is_none());
659
660 let snap = t.snapshot();
661 assert_eq!(snap["response_count"], 1);
662 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
663 }
664
665 #[test]
666 fn text_input_latency_via_audio_out() {
667 let t = SessionTelemetry::new();
668 t.record_text_send();
670 std::thread::sleep(std::time::Duration::from_millis(10));
671 assert!(t.record_audio_out(480).is_some());
672
673 let snap = t.snapshot();
674 assert_eq!(snap["response_count"], 1);
676 assert!(snap["last_response_latency_ms"].as_u64().unwrap() >= 5);
677 }
678
679 #[test]
680 fn mixed_voice_and_text_turns() {
681 let t = SessionTelemetry::new();
682
683 t.record_vad_end();
685 std::thread::sleep(std::time::Duration::from_millis(10));
686 t.record_audio_out(480);
687
688 t.record_text_send();
690 std::thread::sleep(std::time::Duration::from_millis(10));
691 t.record_text_out();
692
693 let snap = t.snapshot();
694 assert_eq!(snap["response_count"], 2);
695 }
696
697 fn ms(n: u64) -> u64 {
700 n * 1_000_000
701 }
702
703 #[test]
704 fn recorder_scalars_and_percentiles() {
705 let r = LatencyRecorder::new();
706 for v in [300, 100, 500, 200, 400] {
707 r.record(ms(v));
708 }
709 let s = r.stats();
710 assert_eq!(s.count, 5);
711 assert_eq!(s.last_ms, 400);
712 assert_eq!(s.min_ms, 100);
713 assert_eq!(s.max_ms, 500);
714 assert_eq!(s.mean_ms, 300);
715 assert_eq!(s.p50_ms, 300);
718 assert_eq!(s.p90_ms, 500);
719 assert_eq!(s.p99_ms, 500);
720 assert_eq!(
721 s.to_string(),
722 "turns=5 last=400ms p50=300ms p90=500ms p99=500ms min=100ms max=500ms"
723 );
724 }
725
726 #[test]
727 fn recorder_first_sample_is_visible_to_percentiles() {
728 let r = LatencyRecorder::new();
731 r.record(ms(420));
732 let s = r.stats();
733 assert_eq!(s.count, 1);
734 assert_eq!((s.p50_ms, s.p90_ms, s.p99_ms), (420, 420, 420));
735 assert_eq!(s.mean_ms, 420);
736 }
737
738 #[test]
739 fn recorder_stats_survives_a_count_without_its_sample() {
740 let r = LatencyRecorder::new();
746 r.count.store(1, Release); let s = r.stats();
748 assert_eq!(s.count, 1);
749 assert_eq!((s.p50_ms, s.p90_ms, s.p99_ms), (0, 0, 0));
750 }
751
752 #[test]
753 fn recorder_snapshots_stay_consistent_under_concurrent_records() {
754 use std::sync::Arc;
758 use std::sync::atomic::AtomicBool;
759
760 let r = Arc::new(LatencyRecorder::new());
761 let done = Arc::new(AtomicBool::new(false));
762
763 let reader = {
764 let r = Arc::clone(&r);
765 let done = Arc::clone(&done);
766 std::thread::spawn(move || {
767 while !done.load(Relaxed) {
768 let s = r.stats();
769 if s.count > 0 {
770 assert!(s.p99_ms >= s.p50_ms);
771 assert!(s.max_ms >= s.p99_ms);
772 assert!(s.min_ms <= s.max_ms);
773 }
774 }
775 })
776 };
777
778 for _ in 0..2_000 {
779 r.record(ms(100));
780 }
781 done.store(true, Relaxed);
782 reader.join().expect("reader saw an inconsistent snapshot");
783 assert_eq!(r.stats().count, 2_000);
784 }
785
786 #[test]
787 fn recorder_histogram_buckets_by_upper_bound() {
788 let r = LatencyRecorder::new();
789 r.record(ms(50)); r.record(ms(51)); r.record(ms(999)); r.record(ms(9_000)); let h = r.stats().histogram;
794 assert_eq!(h.len(), LATENCY_BUCKETS_MS.len() + 1);
795 assert_eq!(
796 h[0],
797 LatencyBucket {
798 upper_ms: Some(50),
799 count: 1
800 }
801 );
802 assert_eq!(
803 h[1],
804 LatencyBucket {
805 upper_ms: Some(100),
806 count: 1
807 }
808 );
809 assert_eq!(
810 h[9],
811 LatencyBucket {
812 upper_ms: Some(1000),
813 count: 1
814 }
815 );
816 assert_eq!(
817 h[h.len() - 1],
818 LatencyBucket {
819 upper_ms: None,
820 count: 1
821 }
822 );
823 assert_eq!(h.iter().map(|b| b.count).sum::<u64>(), 4);
824 }
825
826 #[test]
827 fn recorder_percentiles_use_recent_window_only() {
828 let r = LatencyRecorder::new();
829 for _ in 0..LATENCY_RECENT_WINDOW {
833 r.record(ms(2_000));
834 }
835 for _ in 0..LATENCY_RECENT_WINDOW {
836 r.record(ms(200));
837 }
838 let s = r.stats();
839 assert_eq!(s.count, 2 * LATENCY_RECENT_WINDOW as u64);
840 assert_eq!(s.p50_ms, 200);
841 assert_eq!(s.p99_ms, 200);
842 assert_eq!(s.max_ms, 2_000);
843 assert_eq!(s.min_ms, 200);
844 let slow: u64 = s
845 .histogram
846 .iter()
847 .filter(|b| b.upper_ms == Some(2000))
848 .map(|b| b.count)
849 .sum();
850 assert_eq!(slow, LATENCY_RECENT_WINDOW as u64);
851 }
852
853 #[test]
854 fn stats_round_trip_through_json() {
855 let r = LatencyRecorder::new();
856 r.record(ms(420));
857 let s = r.stats();
858 let v = serde_json::to_value(&s).unwrap();
859 assert_eq!(v["p50_ms"], 420);
860 let back: LatencyStats = serde_json::from_value(v).unwrap();
861 assert_eq!(back, s);
862 }
863}