lean_ctx/core/a2a/
telemetry.rs1use serde::{Deserialize, Serialize};
4use std::sync::atomic::{AtomicU64, Ordering};
5
6static DELIVERIES_ATTEMPTED: AtomicU64 = AtomicU64::new(0);
7static DELIVERIES_SUCCEEDED: AtomicU64 = AtomicU64::new(0);
8static DELIVERIES_FAILED: AtomicU64 = AtomicU64::new(0);
9static DELIVERIES_DLQ: AtomicU64 = AtomicU64::new(0);
10static TOTAL_PAYLOAD_BYTES: AtomicU64 = AtomicU64::new(0);
11static TOTAL_LATENCY_US: AtomicU64 = AtomicU64::new(0);
12static RELAY_HOPS_TOTAL: AtomicU64 = AtomicU64::new(0);
13
14pub fn record_delivery(success: bool, payload_bytes: u64, latency_us: u64) {
16 DELIVERIES_ATTEMPTED.fetch_add(1, Ordering::Relaxed);
17 TOTAL_PAYLOAD_BYTES.fetch_add(payload_bytes, Ordering::Relaxed);
18 TOTAL_LATENCY_US.fetch_add(latency_us, Ordering::Relaxed);
19 if success {
20 DELIVERIES_SUCCEEDED.fetch_add(1, Ordering::Relaxed);
21 } else {
22 DELIVERIES_FAILED.fetch_add(1, Ordering::Relaxed);
23 }
24}
25
26pub fn record_dlq() {
28 DELIVERIES_DLQ.fetch_add(1, Ordering::Relaxed);
29}
30
31pub fn record_relay_hop() {
33 RELAY_HOPS_TOTAL.fetch_add(1, Ordering::Relaxed);
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize)]
38pub struct TransportSnapshot {
39 pub deliveries_attempted: u64,
40 pub deliveries_succeeded: u64,
41 pub deliveries_failed: u64,
42 pub deliveries_dlq: u64,
43 pub total_payload_bytes: u64,
44 pub total_latency_us: u64,
45 pub relay_hops_total: u64,
46 pub success_rate: f64,
47 pub avg_latency_us: u64,
48 pub avg_payload_bytes: u64,
49}
50
51pub fn snapshot() -> TransportSnapshot {
53 let deliveries_attempted = DELIVERIES_ATTEMPTED.load(Ordering::Relaxed);
54 let deliveries_succeeded = DELIVERIES_SUCCEEDED.load(Ordering::Relaxed);
55 let deliveries_failed = DELIVERIES_FAILED.load(Ordering::Relaxed);
56 let total_payload_bytes = TOTAL_PAYLOAD_BYTES.load(Ordering::Relaxed);
57 let total_latency_us = TOTAL_LATENCY_US.load(Ordering::Relaxed);
58
59 TransportSnapshot {
60 deliveries_attempted,
61 deliveries_succeeded,
62 deliveries_failed,
63 deliveries_dlq: DELIVERIES_DLQ.load(Ordering::Relaxed),
64 total_payload_bytes,
65 total_latency_us,
66 relay_hops_total: RELAY_HOPS_TOTAL.load(Ordering::Relaxed),
67 success_rate: if deliveries_attempted > 0 {
68 deliveries_succeeded as f64 / deliveries_attempted as f64
69 } else {
70 0.0
71 },
72 avg_latency_us: total_latency_us
73 .checked_div(deliveries_attempted)
74 .unwrap_or(0),
75 avg_payload_bytes: total_payload_bytes
76 .checked_div(deliveries_attempted)
77 .unwrap_or(0),
78 }
79}
80
81#[cfg(test)]
82mod tests {
83 use super::*;
84 use std::sync::Mutex;
85
86 static TEST_LOCK: Mutex<()> = Mutex::new(());
87
88 #[test]
89 fn records_delivery_counters() {
90 let _guard = TEST_LOCK.lock().unwrap();
91 let before = snapshot();
92
93 record_delivery(true, 128, 40);
94 record_delivery(false, 256, 80);
95
96 let after = snapshot();
97 assert_eq!(after.deliveries_attempted, before.deliveries_attempted + 2);
98 assert_eq!(after.deliveries_succeeded, before.deliveries_succeeded + 1);
99 assert_eq!(after.deliveries_failed, before.deliveries_failed + 1);
100 assert_eq!(after.total_payload_bytes, before.total_payload_bytes + 384);
101 assert_eq!(after.total_latency_us, before.total_latency_us + 120);
102 }
103
104 #[test]
105 fn snapshot_includes_dlq_relay_and_averages() {
106 let _guard = TEST_LOCK.lock().unwrap();
107 let before = snapshot();
108
109 record_delivery(true, 300, 90);
110 record_dlq();
111 record_relay_hop();
112
113 let after = snapshot();
114 assert_eq!(after.deliveries_dlq, before.deliveries_dlq + 1);
115 assert_eq!(after.relay_hops_total, before.relay_hops_total + 1);
116 assert_eq!(
117 after.avg_latency_us,
118 after.total_latency_us / after.deliveries_attempted
119 );
120 assert_eq!(
121 after.avg_payload_bytes,
122 after.total_payload_bytes / after.deliveries_attempted
123 );
124 }
125
126 #[test]
127 fn snapshot_reports_success_rate() {
128 let _guard = TEST_LOCK.lock().unwrap();
129 record_delivery(true, 1, 1);
130 record_delivery(true, 1, 1);
131 record_delivery(false, 1, 1);
132
133 let current = snapshot();
134 let expected = current.deliveries_succeeded as f64 / current.deliveries_attempted as f64;
135 assert!((current.success_rate - expected).abs() < f64::EPSILON);
136 assert!((0.0..=1.0).contains(¤t.success_rate));
137 }
138}