Skip to main content

vti_common/telemetry/
ring.rs

1//! In-memory bounded ring-buffer telemetry sink.
2
3use std::collections::VecDeque;
4
5use async_trait::async_trait;
6use tokio::sync::RwLock;
7
8use super::{TelemetryError, TelemetryEvent, TelemetryFilter, TelemetrySink};
9
10const DEFAULT_CAPACITY: usize = 10_000;
11
12pub struct RingBufferTelemetry {
13    buf: RwLock<VecDeque<TelemetryEvent>>,
14    capacity: usize,
15}
16
17impl RingBufferTelemetry {
18    pub fn new() -> Self {
19        Self::with_capacity(DEFAULT_CAPACITY)
20    }
21
22    pub fn with_capacity(capacity: usize) -> Self {
23        assert!(capacity > 0, "RingBufferTelemetry capacity must be > 0");
24        Self {
25            buf: RwLock::new(VecDeque::with_capacity(capacity)),
26            capacity,
27        }
28    }
29
30    pub fn capacity(&self) -> usize {
31        self.capacity
32    }
33}
34
35impl Default for RingBufferTelemetry {
36    fn default() -> Self {
37        Self::new()
38    }
39}
40
41#[async_trait]
42impl TelemetrySink for RingBufferTelemetry {
43    async fn record(&self, event: TelemetryEvent) -> Result<(), TelemetryError> {
44        let mut buf = self.buf.write().await;
45        if buf.len() == self.capacity {
46            buf.pop_front();
47        }
48        buf.push_back(event);
49        Ok(())
50    }
51
52    async fn query(&self, filter: &TelemetryFilter) -> Result<Vec<TelemetryEvent>, TelemetryError> {
53        let buf = self.buf.read().await;
54        Ok(buf
55            .iter()
56            .rev()
57            .filter(|e| filter.matches(e))
58            .cloned()
59            .collect())
60    }
61}
62
63#[cfg(test)]
64mod tests {
65    use super::*;
66    use crate::telemetry::TelemetryKind;
67    use chrono::{Duration, Utc};
68
69    fn ev(kind: TelemetryKind) -> TelemetryEvent {
70        TelemetryEvent::new(kind)
71    }
72
73    #[tokio::test]
74    async fn round_trip_record_and_query() {
75        let sink = RingBufferTelemetry::with_capacity(8);
76        sink.record(ev(TelemetryKind::DidcommInbound).with_mediator("did:test:A"))
77            .await
78            .unwrap();
79        let out = sink.query(&TelemetryFilter::new()).await.unwrap();
80        assert_eq!(out.len(), 1);
81        assert_eq!(out[0].kind, TelemetryKind::DidcommInbound);
82        assert_eq!(out[0].mediator_did.as_deref(), Some("did:test:A"));
83    }
84
85    #[tokio::test]
86    async fn capacity_overflow_drops_oldest() {
87        let sink = RingBufferTelemetry::with_capacity(3);
88        for i in 0..5 {
89            sink.record(
90                ev(TelemetryKind::DidcommInbound).with_field("i", serde_json::Value::from(i)),
91            )
92            .await
93            .unwrap();
94        }
95        let out = sink.query(&TelemetryFilter::new()).await.unwrap();
96        assert_eq!(out.len(), 3, "buffer cap respected");
97        let ids: Vec<i64> = out
98            .iter()
99            .map(|e| e.fields["i"].as_i64().unwrap())
100            .collect();
101        assert_eq!(ids, vec![4, 3, 2], "newest-first; 0 and 1 dropped");
102    }
103
104    #[tokio::test]
105    async fn time_range_filter() {
106        let sink = RingBufferTelemetry::with_capacity(16);
107        let t0 = Utc::now();
108        for offset_secs in [0, 60, 120] {
109            let mut e = ev(TelemetryKind::DidcommInbound);
110            e.at = t0 + Duration::seconds(offset_secs);
111            sink.record(e).await.unwrap();
112        }
113        let out = sink
114            .query(
115                &TelemetryFilter::new()
116                    .since(t0 + Duration::seconds(30))
117                    .until(t0 + Duration::seconds(90)),
118            )
119            .await
120            .unwrap();
121        assert_eq!(out.len(), 1);
122    }
123
124    #[tokio::test]
125    async fn kind_filter() {
126        let sink = RingBufferTelemetry::with_capacity(16);
127        sink.record(ev(TelemetryKind::DidcommInbound))
128            .await
129            .unwrap();
130        sink.record(ev(TelemetryKind::MediatorHandshakeOk))
131            .await
132            .unwrap();
133        sink.record(ev(TelemetryKind::MediatorDrainExpire))
134            .await
135            .unwrap();
136        let out = sink
137            .query(&TelemetryFilter::new().kind(TelemetryKind::MediatorHandshakeOk))
138            .await
139            .unwrap();
140        assert_eq!(out.len(), 1);
141        assert_eq!(out[0].kind, TelemetryKind::MediatorHandshakeOk);
142    }
143
144    #[tokio::test]
145    async fn mediator_filter() {
146        let sink = RingBufferTelemetry::with_capacity(16);
147        sink.record(ev(TelemetryKind::DidcommInbound).with_mediator("did:test:A"))
148            .await
149            .unwrap();
150        sink.record(ev(TelemetryKind::DidcommInbound).with_mediator("did:test:B"))
151            .await
152            .unwrap();
153        let out = sink
154            .query(&TelemetryFilter::new().mediator("did:test:A"))
155            .await
156            .unwrap();
157        assert_eq!(out.len(), 1);
158        assert_eq!(out[0].mediator_did.as_deref(), Some("did:test:A"));
159    }
160
161    #[tokio::test]
162    async fn sender_filter() {
163        let sink = RingBufferTelemetry::with_capacity(16);
164        sink.record(ev(TelemetryKind::DidcommInbound).with_sender("did:peer:alice"))
165            .await
166            .unwrap();
167        sink.record(ev(TelemetryKind::DidcommInbound).with_sender("did:peer:bob"))
168            .await
169            .unwrap();
170        let out = sink
171            .query(&TelemetryFilter::new().sender("did:peer:alice"))
172            .await
173            .unwrap();
174        assert_eq!(out.len(), 1);
175        assert_eq!(out[0].sender_did.as_deref(), Some("did:peer:alice"));
176    }
177
178    #[tokio::test]
179    async fn newest_first_ordering() {
180        let sink = RingBufferTelemetry::with_capacity(16);
181        for i in 0..5 {
182            sink.record(
183                ev(TelemetryKind::DidcommInbound).with_field("seq", serde_json::Value::from(i)),
184            )
185            .await
186            .unwrap();
187        }
188        let out = sink.query(&TelemetryFilter::new()).await.unwrap();
189        let seqs: Vec<i64> = out
190            .iter()
191            .map(|e| e.fields["seq"].as_i64().unwrap())
192            .collect();
193        assert_eq!(seqs, vec![4, 3, 2, 1, 0]);
194    }
195
196    #[tokio::test]
197    async fn combined_filters() {
198        let sink = RingBufferTelemetry::with_capacity(16);
199        sink.record(
200            ev(TelemetryKind::DidcommInbound)
201                .with_mediator("did:test:A")
202                .with_sender("did:peer:alice"),
203        )
204        .await
205        .unwrap();
206        sink.record(
207            ev(TelemetryKind::DidcommInbound)
208                .with_mediator("did:test:A")
209                .with_sender("did:peer:bob"),
210        )
211        .await
212        .unwrap();
213        sink.record(ev(TelemetryKind::MediatorHandshakeOk).with_mediator("did:test:A"))
214            .await
215            .unwrap();
216
217        let out = sink
218            .query(
219                &TelemetryFilter::new()
220                    .kind(TelemetryKind::DidcommInbound)
221                    .mediator("did:test:A")
222                    .sender("did:peer:alice"),
223            )
224            .await
225            .unwrap();
226        assert_eq!(out.len(), 1);
227        assert_eq!(out[0].sender_did.as_deref(), Some("did:peer:alice"));
228    }
229
230    #[tokio::test]
231    async fn arc_dyn_dispatch_works() {
232        // Proves the trait is dyn-compatible — caller can swap impls behind
233        // an Arc<dyn TelemetrySink>. Foreshadows P5.3 (criterion #17).
234        let sink: super::super::SharedTelemetrySink =
235            std::sync::Arc::new(RingBufferTelemetry::with_capacity(4));
236        sink.record(ev(TelemetryKind::DidcommInbound))
237            .await
238            .unwrap();
239        let out = sink.query(&TelemetryFilter::new()).await.unwrap();
240        assert_eq!(out.len(), 1);
241    }
242
243    #[test]
244    #[should_panic(expected = "capacity must be > 0")]
245    fn zero_capacity_panics() {
246        let _ = RingBufferTelemetry::with_capacity(0);
247    }
248}