use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use chrono::Utc;
use super::ThumbnailProjection;
use crate::domain::entity::Thumbnail;
use crate::domain::events::{
ThumbnailCreatedEvent, ThumbnailUpdatedEvent, ThumbnailDeletedEvent,
};
#[async_trait]
pub trait ThumbnailEventHandler: Send + Sync {
async fn on_created(&self, event: ThumbnailCreatedEvent, sequence: i64) -> Result<()>;
async fn on_updated(&self, event: ThumbnailUpdatedEvent, sequence: i64) -> Result<()>;
async fn on_deleted(&self, event: ThumbnailDeletedEvent, sequence: i64) -> Result<()>;
}
#[async_trait]
pub trait ThumbnailProjectionRepository: Send + Sync {
async fn save(&self, projection: &ThumbnailProjection) -> Result<()>;
async fn find_by_id(&self, id: uuid::Uuid) -> Result<Option<ThumbnailProjection>>;
async fn delete(&self, id: uuid::Uuid) -> Result<()>;
async fn rebuild_all(&self) -> Result<u64>;
}
pub struct ThumbnailProjector<R: ThumbnailProjectionRepository> {
repository: Arc<R>,
}
impl<R: ThumbnailProjectionRepository> ThumbnailProjector<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: ThumbnailEventHandler,
{
let sequence = event.sequence;
match event.event_type.as_str() {
"ThumbnailCreatedEvent" => {
let created: ThumbnailCreatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_created(created, sequence).await?;
Ok(true)
}
"ThumbnailUpdatedEvent" => {
let updated: ThumbnailUpdatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_updated(updated, sequence).await?;
Ok(true)
}
"ThumbnailDeletedEvent" => {
let deleted: ThumbnailDeletedEvent = serde_json::from_value(event.payload.clone())?;
self.on_deleted(deleted, sequence).await?;
Ok(true)
}
_ => Ok(false), }
}
}
#[async_trait]
impl<R: ThumbnailProjectionRepository> ThumbnailEventHandler for ThumbnailProjector<R> {
async fn on_created(&self, event: ThumbnailCreatedEvent, sequence: i64) -> Result<()> {
let mut projection = ThumbnailProjection::new(
event.id,
event.file_id,
event.size,
event.width,
event.height,
event.storage_key,
event.storage_backend,
event.mime_type,
event.format,
event.quality,
event.size_bytes,
event.generated_at,
event.generation_time_ms,
event.source_version,
event.cdn_url,
event.cache_expires_at,
event.is_stale,
event.metadata
);
projection.apply_event(sequence);
self.repository.save(&projection).await
}
async fn on_updated(&self, event: ThumbnailUpdatedEvent, sequence: i64) -> Result<()> {
if let Some(mut projection) = self.repository.find_by_id(event.id).await? {
projection.file_id = event.file_id;
projection.size = event.size;
projection.width = event.width;
projection.height = event.height;
projection.storage_key = event.storage_key;
projection.storage_backend = event.storage_backend;
projection.mime_type = event.mime_type;
projection.format = event.format;
projection.quality = event.quality;
projection.size_bytes = event.size_bytes;
projection.generated_at = event.generated_at;
projection.generation_time_ms = event.generation_time_ms;
projection.source_version = event.source_version;
projection.cdn_url = event.cdn_url;
projection.cache_expires_at = event.cache_expires_at;
projection.is_stale = event.is_stale;
projection.metadata = event.metadata;
projection.apply_event(sequence);
self.repository.save(&projection).await?;
}
Ok(())
}
async fn on_deleted(&self, event: ThumbnailDeletedEvent, _sequence: i64) -> Result<()> {
self.repository.delete(event.id).await
}
}