use crate::event::types::EventId;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CorrelationPayload {
pub entry_time_ns: u64,
pub entry_event_id: EventId,
#[serde(
default,
skip_serializing_if = "Option::is_none",
deserialize_with = "crate::serde_support::present_json"
)]
pub metadata: Option<serde_json::Value>,
}
impl CorrelationPayload {
pub fn new(entry_event_id: EventId) -> Self {
Self {
entry_time_ns: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos() as u64,
entry_event_id,
metadata: None,
}
}
pub fn calculate_latency(&self) -> std::time::Duration {
let now_ns = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos() as u64;
let latency_ns = now_ns.saturating_sub(self.entry_time_ns);
std::time::Duration::from_nanos(latency_ns)
}
pub fn entry_time(&self) -> std::time::SystemTime {
std::time::UNIX_EPOCH + std::time::Duration::from_nanos(self.entry_time_ns)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_correlation_payload_creation() {
let event_id = EventId::new();
let payload = CorrelationPayload::new(event_id);
assert_eq!(payload.entry_event_id, event_id);
assert!(payload.entry_time_ns > 0);
assert!(payload.metadata.is_none());
}
#[test]
fn test_latency_calculation() {
let event_id = EventId::new();
let payload = CorrelationPayload::new(event_id);
std::thread::sleep(std::time::Duration::from_millis(10));
let latency = payload.calculate_latency();
assert!(latency.as_millis() >= 10);
}
}