vti_common/telemetry/
ring.rs1use 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 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}