1use std::collections::VecDeque;
7use std::num::NonZeroUsize;
8
9use lru::LruCache;
10
11use crate::normalize::NormalizedEvent;
12
13#[derive(Debug, Clone)]
15pub struct WindowConfig {
16 pub max_events_per_trace: usize,
18 pub trace_ttl_ms: u64,
20 pub max_active_traces: NonZeroUsize,
22}
23
24const DEFAULT_MAX_ACTIVE_TRACES: NonZeroUsize =
26 NonZeroUsize::new(10_000).expect("non-zero literal");
27
28impl Default for WindowConfig {
29 fn default() -> Self {
30 Self {
31 max_events_per_trace: 1000,
32 trace_ttl_ms: 30_000,
33 max_active_traces: DEFAULT_MAX_ACTIVE_TRACES,
34 }
35 }
36}
37
38struct TraceBuffer {
40 events: VecDeque<NormalizedEvent>,
41 last_seen_ms: u64,
44}
45
46pub struct TraceWindow {
50 config: WindowConfig,
51 traces: LruCache<String, TraceBuffer>,
52}
53
54impl TraceWindow {
55 #[must_use]
56 pub fn new(config: WindowConfig) -> Self {
57 let cap = config.max_active_traces;
58 Self {
59 config,
60 traces: LruCache::new(cap),
61 }
62 }
63
64 pub fn push(
69 &mut self,
70 event: NormalizedEvent,
71 now_ms: u64,
72 ) -> Option<(String, Vec<NormalizedEvent>)> {
73 if let Some(buf) = self.traces.get_mut(event.event.trace_id.as_str()) {
75 buf.last_seen_ms = now_ms;
76 buf.events.push_back(event);
77 if buf.events.len() > self.config.max_events_per_trace {
79 buf.events.pop_front();
80 }
81 return None;
82 }
83
84 let trace_id = event.event.trace_id.clone();
86 let mut events = VecDeque::with_capacity(8);
87 events.push_back(event);
88
89 self.traces
90 .push(
91 trace_id,
92 TraceBuffer {
93 events,
94 last_seen_ms: now_ms,
95 },
96 )
97 .map(|(id, buf)| (id, Vec::from(buf.events)))
98 }
99
100 pub fn evict(&mut self, now_ms: u64) {
112 for key in self.collect_expired_keys(now_ms) {
113 self.traces.pop(&key);
114 }
115 }
116
117 pub fn evict_expired(&mut self, now_ms: u64) -> Vec<(String, Vec<NormalizedEvent>)> {
123 let expired_keys = self.collect_expired_keys(now_ms);
124 let mut expired = Vec::with_capacity(expired_keys.len());
125 for key in expired_keys {
126 if let Some((_id, buf)) = self.traces.pop_entry(&key) {
127 expired.push((key, Vec::from(buf.events)));
128 }
129 }
130 expired
131 }
132
133 fn collect_expired_keys(&self, now_ms: u64) -> Vec<String> {
136 let ttl = self.config.trace_ttl_ms;
137 self.traces
138 .iter()
139 .filter(|(_, buf)| now_ms.saturating_sub(buf.last_seen_ms) > ttl)
140 .map(|(id, _)| id.clone())
141 .collect()
142 }
143
144 pub fn drain_all(&mut self) -> Vec<(String, Vec<NormalizedEvent>)> {
146 let mut result = Vec::with_capacity(self.traces.len());
147 while let Some((id, buf)) = self.traces.pop_lru() {
148 result.push((id, Vec::from(buf.events)));
149 }
150 result
151 }
152
153 #[must_use]
155 pub fn active_traces(&self) -> usize {
156 self.traces.len()
157 }
158
159 #[must_use]
162 pub fn peek_clone(&self, trace_id: &str) -> Option<Vec<NormalizedEvent>> {
163 self.traces
164 .peek(trace_id)
165 .map(|buf| buf.events.iter().cloned().collect())
166 }
167}
168
169#[cfg(test)]
170mod tests {
171 use std::sync::Arc;
172
173 use super::*;
174 use crate::event::{EventSource, EventType, SpanEvent};
175 use crate::normalize;
176
177 fn make_event(trace_id: &str, target: &str) -> NormalizedEvent {
178 let event = SpanEvent {
179 timestamp: "2025-07-10T14:32:01.123Z".to_string(),
180 trace_id: trace_id.to_string(),
181 span_id: "span-1".to_string(),
182 parent_span_id: None,
183 service: Arc::from("test"),
184 cloud_region: None,
185 event_type: EventType::Sql,
186 operation: "SELECT".to_string(),
187 target: target.to_string(),
188 duration_us: 100,
189 source: EventSource {
190 endpoint: "GET /test".to_string(),
191 method: "Test::test".to_string(),
192 },
193 status_code: None,
194 response_size_bytes: None,
195 code_function: None,
196 code_filepath: None,
197 code_lineno: None,
198 code_namespace: None,
199 instrumentation_scopes: Vec::new(),
200 };
201 normalize::normalize(event)
202 }
203
204 #[test]
205 fn accumulates_events_by_trace() {
206 let mut w = TraceWindow::new(WindowConfig::default());
207 w.push(make_event("t1", "SELECT 1"), 0);
208 w.push(make_event("t1", "SELECT 2"), 10);
209 w.push(make_event("t2", "SELECT 3"), 20);
210
211 assert_eq!(w.active_traces(), 2);
212 let drained = w.drain_all();
213 let t1 = drained.iter().find(|(id, _)| id == "t1").unwrap();
214 assert_eq!(t1.1.len(), 2);
215 }
216
217 #[test]
218 fn ring_buffer_overflow() {
219 let config = WindowConfig {
220 max_events_per_trace: 3,
221 ..Default::default()
222 };
223 let mut w = TraceWindow::new(config);
224 for i in 0..5 {
225 w.push(
226 make_event("t1", &format!("SELECT {i}")),
227 u64::try_from(i).unwrap(),
228 );
229 }
230
231 let drained = w.drain_all();
232 let t1 = drained.iter().find(|(id, _)| id == "t1").unwrap();
233 assert_eq!(t1.1.len(), 3);
234 assert_eq!(t1.1[0].event.target, "SELECT 2");
236 assert_eq!(t1.1[2].event.target, "SELECT 4");
237 }
238
239 #[test]
240 fn ttl_eviction() {
241 let config = WindowConfig {
242 trace_ttl_ms: 100,
243 ..Default::default()
244 };
245 let mut w = TraceWindow::new(config);
246 w.push(make_event("t1", "SELECT 1"), 0);
247 w.push(make_event("t2", "SELECT 2"), 50);
248
249 w.evict(150);
250 assert_eq!(w.active_traces(), 1);
253 let drained = w.drain_all();
254 assert_eq!(drained[0].0, "t2");
255 }
256
257 #[test]
258 fn lru_eviction() {
259 let config = WindowConfig {
260 max_active_traces: NonZeroUsize::new(2).unwrap(),
261 ..Default::default()
262 };
263 let mut w = TraceWindow::new(config);
264 w.push(make_event("t1", "SELECT 1"), 0);
265 w.push(make_event("t2", "SELECT 2"), 10);
266 let evicted = w.push(make_event("t3", "SELECT 3"), 20);
268
269 assert!(evicted.is_some());
270 assert_eq!(evicted.unwrap().0, "t1");
271 assert_eq!(w.active_traces(), 2);
272 assert!(w.traces.peek(&"t2".to_string()).is_some());
273 assert!(w.traces.peek(&"t3".to_string()).is_some());
274 assert!(w.traces.peek(&"t1".to_string()).is_none());
275 }
276
277 #[test]
278 fn drain_empties_window() {
279 let mut w = TraceWindow::new(WindowConfig::default());
280 w.push(make_event("t1", "SELECT 1"), 0);
281 let drained = w.drain_all();
282 assert_eq!(drained.len(), 1);
283 assert_eq!(w.active_traces(), 0);
284 }
285
286 #[test]
287 fn lru_touch_prevents_eviction() {
288 let config = WindowConfig {
289 max_active_traces: NonZeroUsize::new(2).unwrap(),
290 ..Default::default()
291 };
292 let mut w = TraceWindow::new(config);
293 w.push(make_event("t1", "SELECT 1"), 0);
294 w.push(make_event("t2", "SELECT 2"), 10);
295 w.push(make_event("t1", "SELECT 1b"), 20);
297 let evicted = w.push(make_event("t3", "SELECT 3"), 30);
299
300 assert!(evicted.is_some());
301 assert_eq!(evicted.unwrap().0, "t2");
302 assert_eq!(w.active_traces(), 2);
303 assert!(w.traces.peek(&"t1".to_string()).is_some());
304 assert!(w.traces.peek(&"t3".to_string()).is_some());
305 assert!(w.traces.peek(&"t2".to_string()).is_none());
306 }
307
308 #[test]
309 fn evict_on_empty_window() {
310 let mut w = TraceWindow::new(WindowConfig::default());
311 w.evict(1000);
312 assert_eq!(w.active_traces(), 0);
313 }
314
315 #[test]
316 fn ttl_evicts_all_expired() {
317 let config = WindowConfig {
318 trace_ttl_ms: 50,
319 ..Default::default()
320 };
321 let mut w = TraceWindow::new(config);
322 w.push(make_event("t1", "SELECT 1"), 0);
323 w.push(make_event("t2", "SELECT 2"), 10);
324 w.evict(200);
326 assert_eq!(w.active_traces(), 0);
327 }
328
329 #[test]
330 fn drain_empty_window() {
331 let mut w = TraceWindow::new(WindowConfig::default());
332 let drained = w.drain_all();
333 assert!(drained.is_empty());
334 }
335
336 #[test]
337 fn lru_eviction_chain() {
338 let config = WindowConfig {
339 max_active_traces: NonZeroUsize::new(1).unwrap(),
340 ..Default::default()
341 };
342 let mut w = TraceWindow::new(config);
343
344 let evicted1 = w.push(make_event("t1", "SELECT 1"), 0);
345 assert!(evicted1.is_none()); let evicted2 = w.push(make_event("t2", "SELECT 2"), 10);
348 assert!(evicted2.is_some());
350 assert_eq!(evicted2.unwrap().0, "t1");
351 assert_eq!(w.active_traces(), 1);
352 assert!(w.traces.peek(&"t2".to_string()).is_some());
353
354 let evicted3 = w.push(make_event("t3", "SELECT 3"), 20);
355 assert!(evicted3.is_some());
357 assert_eq!(evicted3.unwrap().0, "t2");
358 assert_eq!(w.active_traces(), 1);
359 assert!(w.traces.peek(&"t3".to_string()).is_some());
360 }
361
362 #[test]
363 fn evict_expired_returns_traces() {
364 let config = WindowConfig {
365 trace_ttl_ms: 100,
366 ..Default::default()
367 };
368 let mut w = TraceWindow::new(config);
369 w.push(make_event("t1", "SELECT 1"), 0);
370 w.push(make_event("t2", "SELECT 2"), 50);
371
372 let expired = w.evict_expired(50);
374 assert!(expired.is_empty());
375 assert_eq!(w.active_traces(), 2);
376
377 let expired = w.evict_expired(150);
379 assert_eq!(expired.len(), 1);
380 assert_eq!(expired[0].0, "t1");
381 assert_eq!(w.active_traces(), 1);
382 }
383
384 #[test]
385 fn push_returns_evicted_events() {
386 let config = WindowConfig {
387 max_active_traces: NonZeroUsize::new(1).unwrap(),
388 ..Default::default()
389 };
390 let mut w = TraceWindow::new(config);
391 w.push(make_event("t1", "SELECT 1"), 0);
392 w.push(make_event("t1", "SELECT 2"), 5);
393
394 let evicted = w.push(make_event("t2", "SELECT 3"), 10);
395 assert!(evicted.is_some());
396 let (id, events) = evicted.unwrap();
397 assert_eq!(id, "t1");
398 assert_eq!(events.len(), 2); }
400}