use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use chrono::Utc;
use super::BucketProjection;
use crate::domain::entity::Bucket;
use crate::domain::events::{
BucketCreatedEvent, BucketUpdatedEvent, BucketDeletedEvent,
};
#[async_trait]
pub trait BucketEventHandler: Send + Sync {
async fn on_created(&self, event: BucketCreatedEvent, sequence: i64) -> Result<()>;
async fn on_updated(&self, event: BucketUpdatedEvent, sequence: i64) -> Result<()>;
async fn on_deleted(&self, event: BucketDeletedEvent, sequence: i64) -> Result<()>;
}
#[async_trait]
pub trait BucketProjectionRepository: Send + Sync {
async fn save(&self, projection: &BucketProjection) -> Result<()>;
async fn find_by_id(&self, id: uuid::Uuid) -> Result<Option<BucketProjection>>;
async fn delete(&self, id: uuid::Uuid) -> Result<()>;
async fn rebuild_all(&self) -> Result<u64>;
}
pub struct BucketProjector<R: BucketProjectionRepository> {
repository: Arc<R>,
}
impl<R: BucketProjectionRepository> BucketProjector<R> {
pub fn new(repository: Arc<R>) -> Self {
Self { repository }
}
pub async fn rebuild(&self) -> Result<u64> {
self.repository.rebuild_all().await
}
pub async fn process_stored_event(&self, event: &crate::infrastructure::event_store::StoredEvent) -> Result<bool>
where
Self: BucketEventHandler,
{
let sequence = event.sequence;
match event.event_type.as_str() {
"BucketCreatedEvent" => {
let created: BucketCreatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_created(created, sequence).await?;
Ok(true)
}
"BucketUpdatedEvent" => {
let updated: BucketUpdatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_updated(updated, sequence).await?;
Ok(true)
}
"BucketDeletedEvent" => {
let deleted: BucketDeletedEvent = serde_json::from_value(event.payload.clone())?;
self.on_deleted(deleted, sequence).await?;
Ok(true)
}
_ => Ok(false), }
}
}
#[async_trait]
impl<R: BucketProjectionRepository> BucketEventHandler for BucketProjector<R> {
async fn on_created(&self, event: BucketCreatedEvent, sequence: i64) -> Result<()> {
let mut projection = BucketProjection::new(
event.id,
event.name,
event.slug,
event.description,
event.owner_id,
event.bucket_type,
event.status,
event.storage_backend,
event.root_path,
event.file_count,
event.total_size_bytes,
event.max_file_size,
event.allowed_mime_types,
event.auto_delete_after_days,
event.enable_cdn,
event.enable_versioning,
event.enable_deduplication,
event.metadata
);
projection.apply_event(sequence);
self.repository.save(&projection).await
}
async fn on_updated(&self, event: BucketUpdatedEvent, sequence: i64) -> Result<()> {
if let Some(mut projection) = self.repository.find_by_id(event.id).await? {
projection.name = event.name;
projection.slug = event.slug;
projection.description = event.description;
projection.owner_id = event.owner_id;
projection.bucket_type = event.bucket_type;
projection.status = event.status;
projection.storage_backend = event.storage_backend;
projection.root_path = event.root_path;
projection.file_count = event.file_count;
projection.total_size_bytes = event.total_size_bytes;
projection.max_file_size = event.max_file_size;
projection.allowed_mime_types = event.allowed_mime_types;
projection.auto_delete_after_days = event.auto_delete_after_days;
projection.enable_cdn = event.enable_cdn;
projection.enable_versioning = event.enable_versioning;
projection.enable_deduplication = event.enable_deduplication;
projection.metadata = event.metadata;
projection.apply_event(sequence);
self.repository.save(&projection).await?;
}
Ok(())
}
async fn on_deleted(&self, event: BucketDeletedEvent, _sequence: i64) -> Result<()> {
self.repository.delete(event.id).await
}
}