use async_trait::async_trait;
use anyhow::Result;
use serde::{Deserialize, Serialize};
use crate::domain::entity::ProcessingJob;
use crate::infrastructure::persistence::{ProcessingJobRepository, PaginationParams, PaginatedResult, ProcessingJobFilter};
use super::QueryHandler;
#[deprecated(since = "0.2.0", note = "Use QueryHandler<Q> trait instead")]
#[async_trait]
pub trait ProcessingJobQuery: Send + Sync {
type Output;
async fn execute(&self) -> Result<Self::Output>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GetProcessingJobByIdQuery {
pub id: String,
}
pub struct GetProcessingJobByIdHandler<R: ProcessingJobRepository> {
repository: std::sync::Arc<R>,
}
impl<R: ProcessingJobRepository> GetProcessingJobByIdHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: ProcessingJobRepository + 'static> QueryHandler<GetProcessingJobByIdQuery> for GetProcessingJobByIdHandler<R> {
type Output = Option<ProcessingJob>;
async fn handle(&self, query: GetProcessingJobByIdQuery) -> Result<Self::Output> {
self.repository.find_by_id(&query.id).await
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ListProcessingJobQuery {
pub page: u32,
pub per_page: u32,
pub filter_file_id: Option<Uuid>,
pub filter_job_type: Option<ProcessingJobType>,
pub filter_status: Option<JobStatus>,
pub filter_error_message: Option<String>,
}
impl Default for ListProcessingJobQuery {
fn default() -> Self {
Self {
page: 1,
per_page: 20,
filter_file_id: None,
filter_job_type: None,
filter_status: None,
filter_error_message: None,
}
}
}
pub struct ListProcessingJobHandler<R: ProcessingJobRepository> {
repository: std::sync::Arc<R>,
}
impl<R: ProcessingJobRepository> ListProcessingJobHandler<R> {
pub fn new(repository: std::sync::Arc<R>) -> Self {
Self { repository }
}
}
#[async_trait]
impl<R: ProcessingJobRepository + 'static> QueryHandler<ListProcessingJobQuery> for ListProcessingJobHandler<R> {
type Output = PaginatedResult<ProcessingJob>;
async fn handle(&self, query: ListProcessingJobQuery) -> Result<Self::Output> {
let params = PaginationParams::new(query.page, query.per_page);
let filters = ProcessingJobFilter {
file_id: query.filter_file_id.clone(),
job_type: query.filter_job_type.clone(),
status: query.filter_status.clone(),
error_message: query.filter_error_message.clone(),
..Default::default()
};
if filters.has_filters() {
self.repository.list_with_filters(params, filters).await
} else {
self.repository.list(params).await
}
}
}