use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use chrono::Utc;
use super::StoredFileProjection;
use crate::domain::entity::StoredFile;
use crate::domain::events::{
StoredFileCreatedEvent, StoredFileUpdatedEvent, StoredFileDeletedEvent,
};
#[async_trait]
pub trait StoredFileEventHandler: Send + Sync {
async fn on_created(&self, event: StoredFileCreatedEvent, sequence: i64) -> Result<()>;
async fn on_updated(&self, event: StoredFileUpdatedEvent, sequence: i64) -> Result<()>;
async fn on_deleted(&self, event: StoredFileDeletedEvent, sequence: i64) -> Result<()>;
}
#[async_trait]
pub trait StoredFileProjectionRepository: Send + Sync {
async fn save(&self, projection: &StoredFileProjection) -> Result<()>;
async fn find_by_id(&self, id: uuid::Uuid) -> Result<Option<StoredFileProjection>>;
async fn delete(&self, id: uuid::Uuid) -> Result<()>;
async fn rebuild_all(&self) -> Result<u64>;
}
pub struct StoredFileProjector<R: StoredFileProjectionRepository> {
repository: Arc<R>,
}
impl<R: StoredFileProjectionRepository> StoredFileProjector<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: StoredFileEventHandler,
{
let sequence = event.sequence;
match event.event_type.as_str() {
"StoredFileCreatedEvent" => {
let created: StoredFileCreatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_created(created, sequence).await?;
Ok(true)
}
"StoredFileUpdatedEvent" => {
let updated: StoredFileUpdatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_updated(updated, sequence).await?;
Ok(true)
}
"StoredFileDeletedEvent" => {
let deleted: StoredFileDeletedEvent = serde_json::from_value(event.payload.clone())?;
self.on_deleted(deleted, sequence).await?;
Ok(true)
}
_ => Ok(false), }
}
}
#[async_trait]
impl<R: StoredFileProjectionRepository> StoredFileEventHandler for StoredFileProjector<R> {
async fn on_created(&self, event: StoredFileCreatedEvent, sequence: i64) -> Result<()> {
let mut projection = StoredFileProjection::new(
event.id,
event.bucket_id,
event.owner_id,
event.path,
event.original_name,
event.size_bytes,
event.mime_type,
event.checksum,
event.is_compressed,
event.original_size,
event.compression_algorithm,
event.is_scanned,
event.scan_result,
event.threat_level,
event.has_thumbnail,
event.thumbnail_path,
event.has_video_thumbnail,
event.has_document_preview,
event.processing_status,
event.content_hash_id,
event.cdn_url,
event.cdn_url_expires_at,
event.owner_module,
event.owner_entity,
event.owner_entity_id,
event.field_name,
event.sort_order,
event.status,
event.storage_key,
event.version,
event.previous_version_id,
event.download_count,
event.last_accessed_at,
event.metadata
);
projection.apply_event(sequence);
self.repository.save(&projection).await
}
async fn on_updated(&self, event: StoredFileUpdatedEvent, sequence: i64) -> Result<()> {
if let Some(mut projection) = self.repository.find_by_id(event.id).await? {
projection.bucket_id = event.bucket_id;
projection.owner_id = event.owner_id;
projection.path = event.path;
projection.original_name = event.original_name;
projection.size_bytes = event.size_bytes;
projection.mime_type = event.mime_type;
projection.checksum = event.checksum;
projection.is_compressed = event.is_compressed;
projection.original_size = event.original_size;
projection.compression_algorithm = event.compression_algorithm;
projection.is_scanned = event.is_scanned;
projection.scan_result = event.scan_result;
projection.threat_level = event.threat_level;
projection.has_thumbnail = event.has_thumbnail;
projection.thumbnail_path = event.thumbnail_path;
projection.has_video_thumbnail = event.has_video_thumbnail;
projection.has_document_preview = event.has_document_preview;
projection.processing_status = event.processing_status;
projection.content_hash_id = event.content_hash_id;
projection.cdn_url = event.cdn_url;
projection.cdn_url_expires_at = event.cdn_url_expires_at;
projection.owner_module = event.owner_module;
projection.owner_entity = event.owner_entity;
projection.owner_entity_id = event.owner_entity_id;
projection.field_name = event.field_name;
projection.sort_order = event.sort_order;
projection.status = event.status;
projection.storage_key = event.storage_key;
projection.version = event.version;
projection.previous_version_id = event.previous_version_id;
projection.download_count = event.download_count;
projection.last_accessed_at = event.last_accessed_at;
projection.metadata = event.metadata;
projection.apply_event(sequence);
self.repository.save(&projection).await?;
}
Ok(())
}
async fn on_deleted(&self, event: StoredFileDeletedEvent, _sequence: i64) -> Result<()> {
self.repository.delete(event.id).await
}
}