agentsight_capture/analyzers/
timestamp_normalizer.rs1use super::Analyzer;
10use crate::event::Event;
11use crate::time::boot_ns_to_epoch_ms;
12use async_trait::async_trait;
13use futures::stream::{Stream, StreamExt};
14use std::pin::Pin;
15
16type EventStream = Pin<Box<dyn Stream<Item = Event> + Send>>;
17
18#[derive(Debug)]
19pub struct TimestampNormalizer {}
20
21impl TimestampNormalizer {
22 pub fn new() -> Self {
23 Self {}
24 }
25}
26
27impl Default for TimestampNormalizer {
28 fn default() -> Self {
29 Self::new()
30 }
31}
32
33#[async_trait]
34impl Analyzer for TimestampNormalizer {
35 async fn process(
36 &mut self,
37 stream: EventStream,
38 ) -> Result<EventStream, Box<dyn std::error::Error + Send + Sync>> {
39 let normalized_stream = stream.map(|mut event| {
40 let timestamp_ms = boot_ns_to_epoch_ms(event.timestamp);
42 event.timestamp = timestamp_ms;
43 event
44 });
45
46 Ok(Box::pin(normalized_stream))
47 }
48}
49
50#[cfg(test)]
51mod tests {
52 use super::*;
53 use futures::stream;
54 use serde_json::json;
55
56 #[tokio::test]
57 async fn test_timestamp_normalizer() {
58 let mut normalizer = TimestampNormalizer::new();
59
60 let test_event = Event::new_with_timestamp(
62 1_000_000_000, "test".to_string(),
64 1234,
65 "test_comm".to_string(),
66 json!({"test": "data"}),
67 );
68
69 let input_stream = stream::iter(vec![test_event]);
70 let output_stream = normalizer.process(Box::pin(input_stream)).await.unwrap();
71
72 let results: Vec<Event> = output_stream.collect().await;
73 assert_eq!(results.len(), 1);
74
75 assert!(results[0].timestamp > 1_000_000_000_000); }
79
80 #[tokio::test]
81 async fn test_timestamp_normalizer_multiple_events() {
82 let mut normalizer = TimestampNormalizer::new();
83
84 let events = vec![
85 Event::new_with_timestamp(
86 1_000_000_000, "test".to_string(),
88 1234,
89 "test1".to_string(),
90 json!({"id": 1}),
91 ),
92 Event::new_with_timestamp(
93 2_000_000_000, "test".to_string(),
95 1234,
96 "test2".to_string(),
97 json!({"id": 2}),
98 ),
99 Event::new_with_timestamp(
100 3_000_000_000, "test".to_string(),
102 1234,
103 "test3".to_string(),
104 json!({"id": 3}),
105 ),
106 ];
107
108 let input_stream = stream::iter(events);
109 let output_stream = normalizer.process(Box::pin(input_stream)).await.unwrap();
110
111 let results: Vec<Event> = output_stream.collect().await;
112 assert_eq!(results.len(), 3);
113
114 assert!(results[0].timestamp < results[1].timestamp);
116 assert!(results[1].timestamp < results[2].timestamp);
117
118 for result in &results {
120 assert!(result.timestamp > 1_000_000_000_000); }
122 }
123}