use async_trait::async_trait;
use anyhow::Result;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::domain::entity::UploadSession;
use crate::infrastructure::persistence::UploadSessionRepository;
use super::CommandHandler;
#[deprecated(since = "0.2.0", note = "Use CommandHandler<C> trait instead")]
#[async_trait]
pub trait UploadSessionCommand: Send + Sync {
type Output;
async fn execute(&self) -> Result<Self::Output>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CreateUploadSessionCommand {
pub bucket_id: Uuid,
pub user_id: Uuid,
pub path: String,
pub filename: String,
pub mime_type: Option<String>,
pub file_size: i64,
pub chunk_size: i32,
pub total_chunks: i32,
pub uploaded_chunks: i32,
pub status: UploadStatus,
pub storage_backend: StorageBackend,
pub completed_parts: Vec<i32>,
pub part_etags: Option<serde_json::Value>,
pub expires_at: DateTime<Utc>,
pub metadata: serde_json::Value,
}
pub struct CreateUploadSessionHandler<R: UploadSessionRepository> {
repository: std::sync::Arc<R>,
}
impl<R: UploadSessionRepository> CreateUploadSessionHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: UploadSessionRepository + 'static> CommandHandler<CreateUploadSessionCommand> for CreateUploadSessionHandler<R> {
type Output = UploadSession;
async fn handle(&self, cmd: CreateUploadSessionCommand) -> Result<Self::Output> {
let entity = UploadSession::builder()
.id(Uuid::new_v4().to_string())
.bucket_id(cmd.bucket_id)
.user_id(cmd.user_id)
.path(cmd.path)
.filename(cmd.filename)
.mime_type(cmd.mime_type)
.file_size(cmd.file_size)
.chunk_size(cmd.chunk_size)
.total_chunks(cmd.total_chunks)
.uploaded_chunks(cmd.uploaded_chunks)
.status(cmd.status)
.storage_backend(cmd.storage_backend)
.completed_parts(cmd.completed_parts)
.part_etags(cmd.part_etags)
.expires_at(cmd.expires_at)
.metadata(cmd.metadata)
.build()?;
self.repository.save(&entity).await
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct UpdateUploadSessionCommand {
pub id: String,
pub bucket_id: Option<Uuid>,
pub user_id: Option<Uuid>,
pub path: Option<String>,
pub filename: Option<String>,
pub mime_type: Option<String>,
pub file_size: Option<i64>,
pub chunk_size: Option<i32>,
pub total_chunks: Option<i32>,
pub uploaded_chunks: Option<i32>,
pub status: Option<UploadStatus>,
pub storage_backend: Option<StorageBackend>,
pub completed_parts: Option<Vec<i32>>,
pub part_etags: Option<serde_json::Value>,
pub expires_at: Option<DateTime<Utc>>,
pub metadata: Option<serde_json::Value>,
}
pub struct UpdateUploadSessionHandler<R: UploadSessionRepository> {
repository: std::sync::Arc<R>,
}
impl<R: UploadSessionRepository> UpdateUploadSessionHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: UploadSessionRepository + 'static> CommandHandler<UpdateUploadSessionCommand> for UpdateUploadSessionHandler<R> {
type Output = Option<UploadSession>;
async fn handle(&self, cmd: UpdateUploadSessionCommand) -> Result<Self::Output> {
let existing = self.repository.find_by_id(&cmd.id).await?;
let Some(mut entity) = existing else {
return Ok(None);
};
if let Some(value) = cmd.bucket_id {
entity.bucket_id = value;
}
if let Some(value) = cmd.user_id {
entity.user_id = value;
}
if let Some(value) = cmd.path {
entity.path = value;
}
if let Some(value) = cmd.filename {
entity.filename = value;
}
entity.mime_type = cmd.mime_type;
if let Some(value) = cmd.file_size {
entity.file_size = value;
}
if let Some(value) = cmd.chunk_size {
entity.chunk_size = value;
}
if let Some(value) = cmd.total_chunks {
entity.total_chunks = value;
}
if let Some(value) = cmd.uploaded_chunks {
entity.uploaded_chunks = value;
}
if let Some(value) = cmd.status {
entity.status = value;
}
if let Some(value) = cmd.storage_backend {
entity.storage_backend = value;
}
if let Some(value) = cmd.completed_parts {
entity.completed_parts = value;
}
entity.part_etags = cmd.part_etags;
if let Some(value) = cmd.expires_at {
entity.expires_at = value;
}
if let Some(value) = cmd.metadata {
entity.metadata = value;
}
self.repository.update(&cmd.id, &entity).await
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeleteUploadSessionCommand {
pub id: String,
pub hard_delete: bool,
}
pub struct DeleteUploadSessionHandler<R: UploadSessionRepository> {
repository: std::sync::Arc<R>,
}
impl<R: UploadSessionRepository> DeleteUploadSessionHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: UploadSessionRepository + 'static> CommandHandler<DeleteUploadSessionCommand> for DeleteUploadSessionHandler<R> {
type Output = bool;
async fn handle(&self, cmd: DeleteUploadSessionCommand) -> Result<Self::Output> {
if cmd.hard_delete {
self.repository.delete(&cmd.id).await
} else {
self.repository.soft_delete(&cmd.id).await
}
}
}