Skip to main content

lean_ctx/core/a2a/
telemetry.rs

1//! Process-global telemetry for A2A transport operations.
2
3use 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
14/// Record one bounded A2A transport delivery attempt.
15pub 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
26/// Record one delivery routed to the dead-letter queue.
27pub fn record_dlq() {
28    DELIVERIES_DLQ.fetch_add(1, Ordering::Relaxed);
29}
30
31/// Record one relay hop traversed by an A2A delivery.
32pub fn record_relay_hop() {
33    RELAY_HOPS_TOTAL.fetch_add(1, Ordering::Relaxed);
34}
35
36/// Point-in-time snapshot of A2A transport telemetry.
37#[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
51/// Capture current process-global A2A transport telemetry.
52pub 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(&current.success_rate));
137    }
138}