use async_trait::async_trait;
use chrono::{DateTime, Utc};
use std::sync::Arc;
use crate::error::EventError;
use crate::event::DomainEvent;
#[derive(Debug, Clone)]
pub struct EventMetadata {
pub event_id: String,
pub timestamp: DateTime<Utc>,
pub correlation_id: Option<String>,
pub entity_type: &'static str,
pub aggregate_id: String,
}
impl EventMetadata {
pub fn new(aggregate_id: impl Into<String>, entity_type: &'static str) -> Self {
Self {
event_id: uuid::Uuid::new_v4().to_string(),
timestamp: Utc::now(),
correlation_id: None,
entity_type,
aggregate_id: aggregate_id.into(),
}
}
pub fn with_correlation(mut self, id: impl Into<String>) -> Self {
self.correlation_id = Some(id.into());
self
}
}
#[derive(Debug, Clone)]
pub enum CrudEvent<E: Clone + Send + Sync + 'static> {
Created { entity: E, metadata: EventMetadata },
Updated { before: E, after: E, metadata: EventMetadata },
Patched { before: E, after: E, metadata: EventMetadata },
SoftDeleted { entity: E, metadata: EventMetadata },
Restored { entity: E, metadata: EventMetadata },
HardDeleted { entity_id: String, metadata: EventMetadata },
BulkCreated { entities: Vec<E>, metadata: EventMetadata },
}
impl<E: Clone + Send + Sync + 'static> CrudEvent<E> {
pub fn metadata(&self) -> &EventMetadata {
match self {
CrudEvent::Created { metadata, .. } => metadata,
CrudEvent::Updated { metadata, .. } => metadata,
CrudEvent::Patched { metadata, .. } => metadata,
CrudEvent::SoftDeleted { metadata, .. } => metadata,
CrudEvent::Restored { metadata, .. } => metadata,
CrudEvent::HardDeleted { metadata, .. } => metadata,
CrudEvent::BulkCreated { metadata, .. } => metadata,
}
}
}
impl<E: Clone + Send + Sync + 'static> DomainEvent for CrudEvent<E> {
fn event_type(&self) -> &'static str {
match self {
CrudEvent::Created { .. } => "entity.created",
CrudEvent::Updated { .. } => "entity.updated",
CrudEvent::Patched { .. } => "entity.patched",
CrudEvent::SoftDeleted { .. } => "entity.soft_deleted",
CrudEvent::Restored { .. } => "entity.restored",
CrudEvent::HardDeleted { .. } => "entity.hard_deleted",
CrudEvent::BulkCreated { .. } => "entity.bulk_created",
}
}
fn aggregate_id(&self) -> &str {
&self.metadata().aggregate_id
}
fn occurred_at(&self) -> DateTime<Utc> {
self.metadata().timestamp
}
fn aggregate_type(&self) -> &'static str {
self.metadata().entity_type
}
}
#[async_trait]
pub trait CrudEventPublisher<E: Clone + Send + Sync + 'static>: Send + Sync {
async fn publish(&self, event: CrudEvent<E>) -> Result<(), EventError>;
async fn publish_many(&self, events: Vec<CrudEvent<E>>) -> Result<(), EventError> {
for event in events {
self.publish(event).await?;
}
Ok(())
}
async fn publish_created(&self, entity: E, _user_id: Option<String>) -> Result<(), EventError> {
let meta = EventMetadata::new("", std::any::type_name::<E>());
self.publish(CrudEvent::Created { entity, metadata: meta }).await
}
async fn publish_updated(&self, entity: E, _user_id: Option<String>) -> Result<(), EventError> {
let meta = EventMetadata::new("", std::any::type_name::<E>());
self.publish(CrudEvent::Updated { before: entity.clone(), after: entity, metadata: meta }).await
}
async fn publish_deleted(&self, entity_id: String, _user_id: Option<String>) -> Result<(), EventError> {
let meta = EventMetadata::new(entity_id.clone(), std::any::type_name::<E>());
self.publish(CrudEvent::HardDeleted { entity_id, metadata: meta }).await
}
}
pub struct NoOpCrudEventPublisher<E: Clone + Send + Sync + 'static> {
_phantom: std::marker::PhantomData<E>,
}
impl<E: Clone + Send + Sync + 'static> NoOpCrudEventPublisher<E> {
pub fn new() -> Self {
Self {
_phantom: std::marker::PhantomData,
}
}
pub fn arc() -> Arc<dyn CrudEventPublisher<E>> {
Arc::new(Self::new())
}
}
impl<E: Clone + Send + Sync + 'static> Default for NoOpCrudEventPublisher<E> {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl<E: Clone + Send + Sync + 'static> CrudEventPublisher<E> for NoOpCrudEventPublisher<E> {
async fn publish(&self, _event: CrudEvent<E>) -> Result<(), EventError> {
Ok(())
}
async fn publish_many(&self, _events: Vec<CrudEvent<E>>) -> Result<(), EventError> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[derive(Debug, Clone)]
struct FakeEntity {
id: String,
}
fn meta(id: &str) -> EventMetadata {
EventMetadata::new(id.to_string(), "fake_entity")
}
#[tokio::test]
async fn noop_publisher_returns_ok_for_all_variants() {
let publisher = NoOpCrudEventPublisher::<FakeEntity>::new();
let events = vec![
CrudEvent::Created {
entity: FakeEntity { id: "1".into() },
metadata: meta("1"),
},
CrudEvent::SoftDeleted {
entity: FakeEntity { id: "1".into() },
metadata: meta("1"),
},
CrudEvent::HardDeleted {
entity_id: "1".into(),
metadata: meta("1"),
},
];
for event in events {
assert!(publisher.publish(event).await.is_ok());
}
}
#[test]
fn event_type_names_are_stable() {
let created: CrudEvent<FakeEntity> = CrudEvent::Created {
entity: FakeEntity { id: "1".into() },
metadata: meta("1"),
};
assert_eq!(created.event_type(), "entity.created");
let deleted: CrudEvent<FakeEntity> = CrudEvent::HardDeleted {
entity_id: "1".into(),
metadata: meta("1"),
};
assert_eq!(deleted.event_type(), "entity.hard_deleted");
}
#[test]
fn aggregate_id_roundtrips() {
let event: CrudEvent<FakeEntity> = CrudEvent::Created {
entity: FakeEntity { id: "abc".into() },
metadata: meta("abc"),
};
assert_eq!(event.aggregate_id(), "abc");
}
#[test]
fn metadata_correlation_id() {
let m = EventMetadata::new("e1", "entity").with_correlation("corr-42");
assert_eq!(m.correlation_id.as_deref(), Some("corr-42"));
assert_eq!(m.aggregate_id, "e1");
assert_eq!(m.entity_type, "entity");
}
}