use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use chrono::Utc;
use super::AccessLogProjection;
use crate::domain::entity::AccessLog;
use crate::domain::events::{
AccessLogCreatedEvent, AccessLogUpdatedEvent, AccessLogDeletedEvent,
};
#[async_trait]
pub trait AccessLogEventHandler: Send + Sync {
async fn on_created(&self, event: AccessLogCreatedEvent, sequence: i64) -> Result<()>;
async fn on_updated(&self, event: AccessLogUpdatedEvent, sequence: i64) -> Result<()>;
async fn on_deleted(&self, event: AccessLogDeletedEvent, sequence: i64) -> Result<()>;
}
#[async_trait]
pub trait AccessLogProjectionRepository: Send + Sync {
async fn save(&self, projection: &AccessLogProjection) -> Result<()>;
async fn find_by_id(&self, id: uuid::Uuid) -> Result<Option<AccessLogProjection>>;
async fn delete(&self, id: uuid::Uuid) -> Result<()>;
async fn rebuild_all(&self) -> Result<u64>;
}
pub struct AccessLogProjector<R: AccessLogProjectionRepository> {
repository: Arc<R>,
}
impl<R: AccessLogProjectionRepository> AccessLogProjector<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: AccessLogEventHandler,
{
let sequence = event.sequence;
match event.event_type.as_str() {
"AccessLogCreatedEvent" => {
let created: AccessLogCreatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_created(created, sequence).await?;
Ok(true)
}
"AccessLogUpdatedEvent" => {
let updated: AccessLogUpdatedEvent = serde_json::from_value(event.payload.clone())?;
self.on_updated(updated, sequence).await?;
Ok(true)
}
"AccessLogDeletedEvent" => {
let deleted: AccessLogDeletedEvent = serde_json::from_value(event.payload.clone())?;
self.on_deleted(deleted, sequence).await?;
Ok(true)
}
_ => Ok(false), }
}
}
#[async_trait]
impl<R: AccessLogProjectionRepository> AccessLogEventHandler for AccessLogProjector<R> {
async fn on_created(&self, event: AccessLogCreatedEvent, sequence: i64) -> Result<()> {
let mut projection = AccessLogProjection::new(
event.id,
event.file_id,
event.bucket_id,
event.action,
event.user_id,
event.share_id,
event.is_owner,
event.is_shared,
event.is_public,
event.ip_address,
event.user_agent,
event.referer,
event.country_code,
event.city,
event.bytes_transferred,
event.duration_ms,
event.success,
event.error_message,
event.accessed_at,
event.metadata
);
projection.apply_event(sequence);
self.repository.save(&projection).await
}
async fn on_updated(&self, event: AccessLogUpdatedEvent, sequence: i64) -> Result<()> {
if let Some(mut projection) = self.repository.find_by_id(event.id).await? {
projection.file_id = event.file_id;
projection.bucket_id = event.bucket_id;
projection.action = event.action;
projection.user_id = event.user_id;
projection.share_id = event.share_id;
projection.is_owner = event.is_owner;
projection.is_shared = event.is_shared;
projection.is_public = event.is_public;
projection.ip_address = event.ip_address;
projection.user_agent = event.user_agent;
projection.referer = event.referer;
projection.country_code = event.country_code;
projection.city = event.city;
projection.bytes_transferred = event.bytes_transferred;
projection.duration_ms = event.duration_ms;
projection.success = event.success;
projection.error_message = event.error_message;
projection.accessed_at = event.accessed_at;
projection.metadata = event.metadata;
projection.apply_event(sequence);
self.repository.save(&projection).await?;
}
Ok(())
}
async fn on_deleted(&self, event: AccessLogDeletedEvent, _sequence: i64) -> Result<()> {
self.repository.delete(event.id).await
}
}