use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
pub trait IntegrationEvent: Clone + Send + Sync + Serialize + for<'de> Deserialize<'de> + 'static {
fn event_type(&self) -> &'static str;
fn source_context(&self) -> &'static str;
fn aggregate_id(&self) -> &str;
fn occurred_at(&self) -> DateTime<Utc>;
fn version(&self) -> u32 {
1
}
fn correlation_id(&self) -> Option<&str> {
None
}
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct IntegrationEventEnvelope {
pub id: String,
pub event_type: String,
pub source_context: String,
pub aggregate_id: String,
pub occurred_at: DateTime<Utc>,
pub published_at: DateTime<Utc>,
pub version: u32,
pub correlation_id: Option<String>,
pub causation_id: Option<String>,
pub payload: serde_json::Value,
}
impl IntegrationEventEnvelope {
pub fn from_event<E: IntegrationEvent>(event: &E) -> Result<Self, serde_json::Error> {
Ok(Self {
id: Uuid::new_v4().to_string(),
event_type: event.event_type().to_string(),
source_context: event.source_context().to_string(),
aggregate_id: event.aggregate_id().to_string(),
occurred_at: event.occurred_at(),
published_at: Utc::now(),
version: event.version(),
correlation_id: event.correlation_id().map(String::from),
causation_id: None,
payload: serde_json::to_value(event)?,
})
}
pub fn deserialize<E: IntegrationEvent>(&self) -> Result<E, serde_json::Error> {
serde_json::from_value(self.payload.clone())
}
pub fn with_causation_id(mut self, causation_id: impl Into<String>) -> Self {
self.causation_id = Some(causation_id.into());
self
}
pub fn with_correlation_id(mut self, correlation_id: impl Into<String>) -> Self {
self.correlation_id = Some(correlation_id.into());
self
}
pub fn matches_pattern(&self, pattern: &str) -> bool {
if pattern == "*" {
return true;
}
if let Some(prefix) = pattern.strip_suffix(".*") {
return self.event_type.starts_with(prefix);
}
pattern == self.event_type
}
}
#[cfg(test)]
mod tests {
use super::*;
#[derive(Clone, Debug, Serialize, Deserialize)]
struct TestIntegrationEvent {
user_id: String,
email: String,
occurred_at: DateTime<Utc>,
}
impl IntegrationEvent for TestIntegrationEvent {
fn event_type(&self) -> &'static str {
"test.user.created"
}
fn source_context(&self) -> &'static str {
"test"
}
fn aggregate_id(&self) -> &str {
&self.user_id
}
fn occurred_at(&self) -> DateTime<Utc> {
self.occurred_at
}
}
#[test]
fn test_envelope_from_event() {
let event = TestIntegrationEvent {
user_id: "user-123".to_string(),
email: "test@example.com".to_string(),
occurred_at: Utc::now(),
};
let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
assert_eq!(envelope.event_type, "test.user.created");
assert_eq!(envelope.source_context, "test");
assert_eq!(envelope.aggregate_id, "user-123");
assert_eq!(envelope.version, 1);
assert!(!envelope.id.is_empty());
}
#[test]
fn test_envelope_deserialize() {
let event = TestIntegrationEvent {
user_id: "user-123".to_string(),
email: "test@example.com".to_string(),
occurred_at: Utc::now(),
};
let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
let restored: TestIntegrationEvent = envelope.deserialize().unwrap();
assert_eq!(restored.user_id, "user-123");
assert_eq!(restored.email, "test@example.com");
}
#[test]
fn test_pattern_matching() {
let event = TestIntegrationEvent {
user_id: "user-123".to_string(),
email: "test@example.com".to_string(),
occurred_at: Utc::now(),
};
let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
assert!(envelope.matches_pattern("test.user.created"));
assert!(!envelope.matches_pattern("test.user.deleted"));
assert!(envelope.matches_pattern("test.user.*"));
assert!(envelope.matches_pattern("test.*"));
assert!(!envelope.matches_pattern("other.*"));
assert!(envelope.matches_pattern("*"));
}
#[test]
fn test_envelope_with_causation() {
let event = TestIntegrationEvent {
user_id: "user-123".to_string(),
email: "test@example.com".to_string(),
occurred_at: Utc::now(),
};
let envelope = IntegrationEventEnvelope::from_event(&event)
.unwrap()
.with_causation_id("parent-event-id")
.with_correlation_id("trace-123");
assert_eq!(envelope.causation_id, Some("parent-event-id".to_string()));
assert_eq!(envelope.correlation_id, Some("trace-123".to_string()));
}
}