use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use super::handlers::*;
#[derive(Debug, Clone)]
pub struct Subscription {
pub source_module: String,
pub event_type: String,
pub handler_name: String,
}
#[async_trait]
pub trait EventBus: Send + Sync {
async fn subscribe(&self, subscription: Subscription) -> Result<()>;
async fn unsubscribe(&self, subscription: &Subscription) -> Result<()>;
}
pub struct SubscriptionRegistry {
subscriptions: Vec<Subscription>,
}
impl SubscriptionRegistry {
pub fn new() -> Self {
Self {
subscriptions: Self::default_subscriptions(),
}
}
fn default_subscriptions() -> Vec<Subscription> {
vec![
]
}
pub async fn register_all<B: EventBus>(&self, bus: &B) -> Result<()> {
for subscription in &self.subscriptions {
tracing::info!(
"Registering subscription: {}::{}",
subscription.source_module,
subscription.event_type
);
bus.subscribe(subscription.clone()).await?;
}
Ok(())
}
pub async fn unregister_all<B: EventBus>(&self, bus: &B) -> Result<()> {
for subscription in &self.subscriptions {
bus.unsubscribe(subscription).await?;
}
Ok(())
}
pub fn subscriptions(&self) -> &[Subscription] {
&self.subscriptions
}
pub fn add(&mut self, subscription: Subscription) {
self.subscriptions.push(subscription);
}
}
impl Default for SubscriptionRegistry {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone)]
pub struct FailedEvent {
pub event_id: uuid::Uuid,
pub source_module: String,
pub event_type: String,
pub payload: serde_json::Value,
pub error: String,
pub attempts: u32,
pub failed_at: chrono::DateTime<chrono::Utc>,
}
#[async_trait]
pub trait DeadLetterQueue: Send + Sync {
async fn send(&self, event: FailedEvent) -> Result<()>;
async fn retrieve(&self, limit: usize) -> Result<Vec<FailedEvent>>;
async fn retry(&self, event_id: uuid::Uuid) -> Result<()>;
async fn discard(&self, event_id: uuid::Uuid) -> Result<()>;
}