use chrono::{DateTime, Utc, Duration};
use serde::{Deserialize, Serialize};
use sqlx::FromRow;
use uuid::Uuid;
use super::ProcessingJobType;
use super::JobStatus;
use super::AuditMetadata;
use super::*;
use crate::domain::state_machine::{ProcessingJobStateMachine, ProcessingJobState, StateMachineError};
use thiserror::Error;
#[derive(Debug, Clone, Error)]
pub enum JobError {
#[error("{0}")]
Message(String),
#[error("Not found: {0}")]
NotFound(String),
#[error("Validation failed: {0}")]
ValidationFailed(String),
#[error("Conflict: {0}")]
Conflict(String),
}
impl From<String> for JobError {
fn from(msg: String) -> Self { Self::Message(msg) }
}
impl From<&str> for JobError {
fn from(msg: &str) -> Self { Self::Message(msg.to_string()) }
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct ProcessingJobId(pub Uuid);
impl ProcessingJobId {
pub fn new(id: Uuid) -> Self { Self(id) }
pub fn generate() -> Self { Self(Uuid::new_v4()) }
pub fn into_inner(self) -> Uuid { self.0 }
}
impl std::fmt::Display for ProcessingJobId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
impl std::str::FromStr for ProcessingJobId {
type Err = uuid::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Ok(Self(Uuid::parse_str(s)?))
}
}
impl From<Uuid> for ProcessingJobId {
fn from(id: Uuid) -> Self { Self(id) }
}
impl From<ProcessingJobId> for Uuid {
fn from(id: ProcessingJobId) -> Self { id.0 }
}
impl AsRef<Uuid> for ProcessingJobId {
fn as_ref(&self) -> &Uuid { &self.0 }
}
impl std::ops::Deref for ProcessingJobId {
type Target = Uuid;
fn deref(&self) -> &Self::Target { &self.0 }
}
#[derive(Debug, Clone, Serialize, Deserialize, FromRow)]
pub struct ProcessingJob {
pub id: Uuid,
pub file_id: Uuid,
pub job_type: ProcessingJobType,
pub(crate) 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,
#[serde(default)]
#[sqlx(json)]
pub metadata: AuditMetadata,
}
impl ProcessingJob {
pub fn builder() -> ProcessingJobBuilder {
<ProcessingJobBuilder as Default>::default()
}
pub fn new(file_id: Uuid, job_type: ProcessingJobType, status: JobStatus, priority: i32, retry_count: i32, max_retries: i32) -> Self {
Self {
id: Uuid::new_v4(),
file_id,
job_type,
status,
priority,
input_data: None,
result_data: None,
error_message: None,
started_at: None,
completed_at: None,
retry_count,
max_retries,
metadata: AuditMetadata::default(),
}
}
pub fn id(&self) -> &Uuid {
&self.id
}
pub fn typed_id(&self) -> ProcessingJobId {
ProcessingJobId(self.id)
}
pub fn created_at(&self) -> Option<&DateTime<Utc>> {
self.metadata.created_at.as_ref()
}
pub fn updated_at(&self) -> Option<&DateTime<Utc>> {
self.metadata.updated_at.as_ref()
}
pub fn is_deleted(&self) -> bool {
self.metadata.deleted_at.is_some()
}
pub fn is_active(&self) -> bool {
self.metadata.deleted_at.is_none()
}
pub fn deleted_at(&self) -> Option<&DateTime<Utc>> {
self.metadata.deleted_at.as_ref()
}
pub fn created_by(&self) -> Option<&Uuid> {
self.metadata.created_by.as_ref()
}
pub fn updated_by(&self) -> Option<&Uuid> {
self.metadata.updated_by.as_ref()
}
pub fn deleted_by(&self) -> Option<&Uuid> {
self.metadata.deleted_by.as_ref()
}
pub fn status(&self) -> &JobStatus {
&self.status
}
pub fn with_input_data(mut self, value: serde_json::Value) -> Self {
self.input_data = Some(value);
self
}
pub fn with_result_data(mut self, value: serde_json::Value) -> Self {
self.result_data = Some(value);
self
}
pub fn with_error_message(mut self, value: String) -> Self {
self.error_message = Some(value);
self
}
pub fn with_started_at(mut self, value: DateTime<Utc>) -> Self {
self.started_at = Some(value);
self
}
pub fn with_completed_at(mut self, value: DateTime<Utc>) -> Self {
self.completed_at = Some(value);
self
}
pub fn transition_to(&mut self, new_state: ProcessingJobState) -> Result<(), StateMachineError> {
let current = self.status.to_string().parse::<ProcessingJobState>()?;
let mut sm = ProcessingJobStateMachine::from_state(current);
sm.transition_to_state(new_state)?;
self.status = new_state.to_string().parse::<JobStatus>()
.map_err(|e| StateMachineError::InvalidState(e.to_string()))?;
Ok(())
}
pub fn apply_patch(&mut self, fields: std::collections::HashMap<String, serde_json::Value>) {
for (key, value) in fields {
match key.as_str() {
"file_id" => {
if let Ok(v) = serde_json::from_value(value) { self.file_id = v; }
}
"job_type" => {
if let Ok(v) = serde_json::from_value(value) { self.job_type = v; }
}
"priority" => {
if let Ok(v) = serde_json::from_value(value) { self.priority = v; }
}
"input_data" => {
if let Ok(v) = serde_json::from_value(value) { self.input_data = v; }
}
"result_data" => {
if let Ok(v) = serde_json::from_value(value) { self.result_data = v; }
}
"error_message" => {
if let Ok(v) = serde_json::from_value(value) { self.error_message = v; }
}
"started_at" => {
if let Ok(v) = serde_json::from_value(value) { self.started_at = v; }
}
"completed_at" => {
if let Ok(v) = serde_json::from_value(value) { self.completed_at = v; }
}
"retry_count" => {
if let Ok(v) = serde_json::from_value(value) { self.retry_count = v; }
}
"max_retries" => {
if let Ok(v) = serde_json::from_value(value) { self.max_retries = v; }
}
_ => {} }
}
}
}
impl super::Entity for ProcessingJob {
type Id = Uuid;
fn entity_id(&self) -> &Self::Id {
&self.id
}
fn entity_type() -> &'static str {
"ProcessingJob"
}
}
impl backbone_core::PersistentEntity for ProcessingJob {
fn entity_id(&self) -> String {
self.id.to_string()
}
fn set_entity_id(&mut self, id: String) {
if let Ok(uuid) = uuid::Uuid::parse_str(&id) {
self.id = uuid;
}
}
fn created_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
self.metadata.created_at
}
fn set_created_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
self.metadata.created_at = Some(ts);
}
fn updated_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
self.metadata.updated_at
}
fn set_updated_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
self.metadata.updated_at = Some(ts);
}
fn deleted_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
self.metadata.deleted_at
}
fn set_deleted_at(&mut self, ts: Option<chrono::DateTime<chrono::Utc>>) {
self.metadata.deleted_at = ts;
}
}
impl backbone_orm::EntityRepoMeta for ProcessingJob {
fn column_types() -> std::collections::HashMap<String, String> {
let mut m = std::collections::HashMap::new();
m.insert("id".to_string(), "uuid".to_string());
m.insert("file_id".to_string(), "uuid".to_string());
m.insert("job_type".to_string(), "processing_job_type".to_string());
m.insert("status".to_string(), "job_status".to_string());
m.insert("started_at".to_string(), "timestamptz".to_string());
m.insert("completed_at".to_string(), "timestamptz".to_string());
m
}
fn search_fields() -> &'static [&'static str] {
&[]
}
}
#[derive(Debug, Clone, Default)]
pub struct ProcessingJobBuilder {
file_id: Option<Uuid>,
job_type: Option<ProcessingJobType>,
status: Option<JobStatus>,
priority: Option<i32>,
input_data: Option<serde_json::Value>,
result_data: Option<serde_json::Value>,
error_message: Option<String>,
started_at: Option<DateTime<Utc>>,
completed_at: Option<DateTime<Utc>>,
retry_count: Option<i32>,
max_retries: Option<i32>,
}
impl ProcessingJobBuilder {
pub fn file_id(mut self, value: Uuid) -> Self {
self.file_id = Some(value);
self
}
pub fn job_type(mut self, value: ProcessingJobType) -> Self {
self.job_type = Some(value);
self
}
pub fn status(mut self, value: JobStatus) -> Self {
self.status = Some(value);
self
}
pub fn priority(mut self, value: i32) -> Self {
self.priority = Some(value);
self
}
pub fn input_data(mut self, value: serde_json::Value) -> Self {
self.input_data = Some(value);
self
}
pub fn result_data(mut self, value: serde_json::Value) -> Self {
self.result_data = Some(value);
self
}
pub fn error_message(mut self, value: String) -> Self {
self.error_message = Some(value);
self
}
pub fn started_at(mut self, value: DateTime<Utc>) -> Self {
self.started_at = Some(value);
self
}
pub fn completed_at(mut self, value: DateTime<Utc>) -> Self {
self.completed_at = Some(value);
self
}
pub fn retry_count(mut self, value: i32) -> Self {
self.retry_count = Some(value);
self
}
pub fn max_retries(mut self, value: i32) -> Self {
self.max_retries = Some(value);
self
}
pub fn build(self) -> Result<ProcessingJob, String> {
let file_id = self.file_id.ok_or_else(|| "file_id is required".to_string())?;
let job_type = self.job_type.ok_or_else(|| "job_type is required".to_string())?;
Ok(ProcessingJob {
id: Uuid::new_v4(),
file_id,
job_type,
status: self.status.unwrap_or_default(),
priority: self.priority.unwrap_or(0),
input_data: self.input_data,
result_data: self.result_data,
error_message: self.error_message,
started_at: self.started_at,
completed_at: self.completed_at,
retry_count: self.retry_count.unwrap_or(0),
max_retries: self.max_retries.unwrap_or(3),
metadata: AuditMetadata::default(),
})
}
}
#[path = "processing_job.ext.rs"]
mod processing_job_ext;
pub use processing_job_ext::*;