use async_trait::async_trait;
use anyhow::Result;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::domain::entity::ProcessingJob;
use crate::infrastructure::persistence::ProcessingJobRepository;
use super::CommandHandler;
#[deprecated(since = "0.2.0", note = "Use CommandHandler<C> trait instead")]
#[async_trait]
pub trait ProcessingJobCommand: Send + Sync {
type Output;
async fn execute(&self) -> Result<Self::Output>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CreateProcessingJobCommand {
pub file_id: Uuid,
pub job_type: ProcessingJobType,
pub status: JobStatus,
pub priority: i32,
pub input_data: Option<serde_json::Value>,
pub result_data: Option<serde_json::Value>,
pub error_message: Option<String>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
pub retry_count: i32,
pub max_retries: i32,
pub metadata: serde_json::Value,
}
pub struct CreateProcessingJobHandler<R: ProcessingJobRepository> {
repository: std::sync::Arc<R>,
}
impl<R: ProcessingJobRepository> CreateProcessingJobHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: ProcessingJobRepository + 'static> CommandHandler<CreateProcessingJobCommand> for CreateProcessingJobHandler<R> {
type Output = ProcessingJob;
async fn handle(&self, cmd: CreateProcessingJobCommand) -> Result<Self::Output> {
let entity = ProcessingJob::builder()
.id(Uuid::new_v4().to_string())
.file_id(cmd.file_id)
.job_type(cmd.job_type)
.status(cmd.status)
.priority(cmd.priority)
.input_data(cmd.input_data)
.result_data(cmd.result_data)
.error_message(cmd.error_message)
.started_at(cmd.started_at)
.completed_at(cmd.completed_at)
.retry_count(cmd.retry_count)
.max_retries(cmd.max_retries)
.metadata(cmd.metadata)
.build()?;
self.repository.save(&entity).await
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct UpdateProcessingJobCommand {
pub id: String,
pub file_id: Option<Uuid>,
pub job_type: Option<ProcessingJobType>,
pub status: Option<JobStatus>,
pub priority: Option<i32>,
pub input_data: Option<serde_json::Value>,
pub result_data: Option<serde_json::Value>,
pub error_message: Option<String>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
pub retry_count: Option<i32>,
pub max_retries: Option<i32>,
pub metadata: Option<serde_json::Value>,
}
pub struct UpdateProcessingJobHandler<R: ProcessingJobRepository> {
repository: std::sync::Arc<R>,
}
impl<R: ProcessingJobRepository> UpdateProcessingJobHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: ProcessingJobRepository + 'static> CommandHandler<UpdateProcessingJobCommand> for UpdateProcessingJobHandler<R> {
type Output = Option<ProcessingJob>;
async fn handle(&self, cmd: UpdateProcessingJobCommand) -> 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.file_id {
entity.file_id = value;
}
if let Some(value) = cmd.job_type {
entity.job_type = value;
}
if let Some(value) = cmd.status {
entity.status = value;
}
if let Some(value) = cmd.priority {
entity.priority = value;
}
entity.input_data = cmd.input_data;
entity.result_data = cmd.result_data;
entity.error_message = cmd.error_message;
entity.started_at = cmd.started_at;
entity.completed_at = cmd.completed_at;
if let Some(value) = cmd.retry_count {
entity.retry_count = value;
}
if let Some(value) = cmd.max_retries {
entity.max_retries = value;
}
if let Some(value) = cmd.metadata {
entity.metadata = value;
}
self.repository.update(&cmd.id, &entity).await
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeleteProcessingJobCommand {
pub id: String,
pub hard_delete: bool,
}
pub struct DeleteProcessingJobHandler<R: ProcessingJobRepository> {
repository: std::sync::Arc<R>,
}
impl<R: ProcessingJobRepository> DeleteProcessingJobHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: ProcessingJobRepository + 'static> CommandHandler<DeleteProcessingJobCommand> for DeleteProcessingJobHandler<R> {
type Output = bool;
async fn handle(&self, cmd: DeleteProcessingJobCommand) -> Result<Self::Output> {
if cmd.hard_delete {
self.repository.delete(&cmd.id).await
} else {
self.repository.soft_delete(&cmd.id).await
}
}
}