backbone_bucket/domain/entity/
processing_job.rs1use chrono::{DateTime, Utc, Duration};
2use serde::{Deserialize, Serialize};
3use sqlx::FromRow;
4use uuid::Uuid;
5
6use super::ProcessingJobType;
7use super::JobStatus;
8use super::AuditMetadata;
9
10use super::*;
11
12use crate::domain::state_machine::{ProcessingJobStateMachine, ProcessingJobState, StateMachineError};
13
14use thiserror::Error;
15
16#[derive(Debug, Clone, Error)]
18pub enum JobError {
19 #[error("{0}")]
20 Message(String),
21
22 #[error("Not found: {0}")]
23 NotFound(String),
24
25 #[error("Validation failed: {0}")]
26 ValidationFailed(String),
27
28 #[error("Conflict: {0}")]
29 Conflict(String),
30}
31
32impl From<String> for JobError {
33 fn from(msg: String) -> Self { Self::Message(msg) }
34}
35
36impl From<&str> for JobError {
37 fn from(msg: &str) -> Self { Self::Message(msg.to_string()) }
38}
39
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
43#[serde(transparent)]
44pub struct ProcessingJobId(pub Uuid);
45
46impl ProcessingJobId {
47 pub fn new(id: Uuid) -> Self { Self(id) }
48 pub fn generate() -> Self { Self(Uuid::new_v4()) }
49 pub fn into_inner(self) -> Uuid { self.0 }
50}
51
52impl std::fmt::Display for ProcessingJobId {
53 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
54 write!(f, "{}", self.0)
55 }
56}
57
58impl std::str::FromStr for ProcessingJobId {
59 type Err = uuid::Error;
60 fn from_str(s: &str) -> Result<Self, Self::Err> {
61 Ok(Self(Uuid::parse_str(s)?))
62 }
63}
64
65impl From<Uuid> for ProcessingJobId {
66 fn from(id: Uuid) -> Self { Self(id) }
67}
68
69impl From<ProcessingJobId> for Uuid {
70 fn from(id: ProcessingJobId) -> Self { id.0 }
71}
72
73impl AsRef<Uuid> for ProcessingJobId {
74 fn as_ref(&self) -> &Uuid { &self.0 }
75}
76
77impl std::ops::Deref for ProcessingJobId {
78 type Target = Uuid;
79 fn deref(&self) -> &Self::Target { &self.0 }
80}
81
82#[derive(Debug, Clone, Serialize, Deserialize, FromRow)]
83pub struct ProcessingJob {
84 pub id: Uuid,
85 pub file_id: Uuid,
86 pub job_type: ProcessingJobType,
87 pub(crate) status: JobStatus,
88 pub priority: i32,
89 pub input_data: Option<serde_json::Value>,
90 pub result_data: Option<serde_json::Value>,
91 pub error_message: Option<String>,
92 pub started_at: Option<DateTime<Utc>>,
93 pub completed_at: Option<DateTime<Utc>>,
94 pub retry_count: i32,
95 pub max_retries: i32,
96 #[serde(default)]
97 #[sqlx(json)]
98 pub metadata: AuditMetadata,
99}
100
101impl ProcessingJob {
102 pub fn builder() -> ProcessingJobBuilder {
104 <ProcessingJobBuilder as Default>::default()
105 }
106
107 pub fn new(file_id: Uuid, job_type: ProcessingJobType, status: JobStatus, priority: i32, retry_count: i32, max_retries: i32) -> Self {
109 Self {
110 id: Uuid::new_v4(),
111 file_id,
112 job_type,
113 status,
114 priority,
115 input_data: None,
116 result_data: None,
117 error_message: None,
118 started_at: None,
119 completed_at: None,
120 retry_count,
121 max_retries,
122 metadata: AuditMetadata::default(),
123 }
124 }
125
126 pub fn id(&self) -> &Uuid {
128 &self.id
129 }
130
131 pub fn typed_id(&self) -> ProcessingJobId {
133 ProcessingJobId(self.id)
134 }
135
136 pub fn created_at(&self) -> Option<&DateTime<Utc>> {
138 self.metadata.created_at.as_ref()
139 }
140
141 pub fn updated_at(&self) -> Option<&DateTime<Utc>> {
143 self.metadata.updated_at.as_ref()
144 }
145
146 pub fn is_deleted(&self) -> bool {
148 self.metadata.deleted_at.is_some()
149 }
150
151 pub fn is_active(&self) -> bool {
153 self.metadata.deleted_at.is_none()
154 }
155
156 pub fn deleted_at(&self) -> Option<&DateTime<Utc>> {
158 self.metadata.deleted_at.as_ref()
159 }
160
161 pub fn created_by(&self) -> Option<&Uuid> {
163 self.metadata.created_by.as_ref()
164 }
165
166 pub fn updated_by(&self) -> Option<&Uuid> {
168 self.metadata.updated_by.as_ref()
169 }
170
171 pub fn deleted_by(&self) -> Option<&Uuid> {
173 self.metadata.deleted_by.as_ref()
174 }
175
176 pub fn status(&self) -> &JobStatus {
178 &self.status
179 }
180
181
182 pub fn with_input_data(mut self, value: serde_json::Value) -> Self {
188 self.input_data = Some(value);
189 self
190 }
191
192 pub fn with_result_data(mut self, value: serde_json::Value) -> Self {
194 self.result_data = Some(value);
195 self
196 }
197
198 pub fn with_error_message(mut self, value: String) -> Self {
200 self.error_message = Some(value);
201 self
202 }
203
204 pub fn with_started_at(mut self, value: DateTime<Utc>) -> Self {
206 self.started_at = Some(value);
207 self
208 }
209
210 pub fn with_completed_at(mut self, value: DateTime<Utc>) -> Self {
212 self.completed_at = Some(value);
213 self
214 }
215
216 pub fn transition_to(&mut self, new_state: ProcessingJobState) -> Result<(), StateMachineError> {
225 let current = self.status.to_string().parse::<ProcessingJobState>()?;
226 let mut sm = ProcessingJobStateMachine::from_state(current);
227 sm.transition_to_state(new_state)?;
228 self.status = new_state.to_string().parse::<JobStatus>()
229 .map_err(|e| StateMachineError::InvalidState(e.to_string()))?;
230 Ok(())
231 }
232
233 pub fn apply_patch(&mut self, fields: std::collections::HashMap<String, serde_json::Value>) {
239 for (key, value) in fields {
240 match key.as_str() {
241 "file_id" => {
242 if let Ok(v) = serde_json::from_value(value) { self.file_id = v; }
243 }
244 "job_type" => {
245 if let Ok(v) = serde_json::from_value(value) { self.job_type = v; }
246 }
247 "priority" => {
248 if let Ok(v) = serde_json::from_value(value) { self.priority = v; }
249 }
250 "input_data" => {
251 if let Ok(v) = serde_json::from_value(value) { self.input_data = v; }
252 }
253 "result_data" => {
254 if let Ok(v) = serde_json::from_value(value) { self.result_data = v; }
255 }
256 "error_message" => {
257 if let Ok(v) = serde_json::from_value(value) { self.error_message = v; }
258 }
259 "started_at" => {
260 if let Ok(v) = serde_json::from_value(value) { self.started_at = v; }
261 }
262 "completed_at" => {
263 if let Ok(v) = serde_json::from_value(value) { self.completed_at = v; }
264 }
265 "retry_count" => {
266 if let Ok(v) = serde_json::from_value(value) { self.retry_count = v; }
267 }
268 "max_retries" => {
269 if let Ok(v) = serde_json::from_value(value) { self.max_retries = v; }
270 }
271 _ => {} }
273 }
274 }
275
276}
277
278impl super::Entity for ProcessingJob {
279 type Id = Uuid;
280
281 fn entity_id(&self) -> &Self::Id {
282 &self.id
283 }
284
285 fn entity_type() -> &'static str {
286 "ProcessingJob"
287 }
288}
289
290impl backbone_core::PersistentEntity for ProcessingJob {
291 fn entity_id(&self) -> String {
292 self.id.to_string()
293 }
294 fn set_entity_id(&mut self, id: String) {
295 if let Ok(uuid) = uuid::Uuid::parse_str(&id) {
296 self.id = uuid;
297 }
298 }
299 fn created_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
300 self.metadata.created_at
301 }
302 fn set_created_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
303 self.metadata.created_at = Some(ts);
304 }
305 fn updated_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
306 self.metadata.updated_at
307 }
308 fn set_updated_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
309 self.metadata.updated_at = Some(ts);
310 }
311 fn deleted_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
312 self.metadata.deleted_at
313 }
314 fn set_deleted_at(&mut self, ts: Option<chrono::DateTime<chrono::Utc>>) {
315 self.metadata.deleted_at = ts;
316 }
317}
318
319impl backbone_orm::EntityRepoMeta for ProcessingJob {
320 fn column_types() -> std::collections::HashMap<String, String> {
321 let mut m = std::collections::HashMap::new();
322 m.insert("id".to_string(), "uuid".to_string());
323 m.insert("file_id".to_string(), "uuid".to_string());
324 m.insert("job_type".to_string(), "processing_job_type".to_string());
325 m.insert("status".to_string(), "job_status".to_string());
326 m.insert("started_at".to_string(), "timestamptz".to_string());
327 m.insert("completed_at".to_string(), "timestamptz".to_string());
328 m
329 }
330 fn search_fields() -> &'static [&'static str] {
331 &[]
332 }
333}
334
335#[derive(Debug, Clone, Default)]
340pub struct ProcessingJobBuilder {
341 file_id: Option<Uuid>,
342 job_type: Option<ProcessingJobType>,
343 status: Option<JobStatus>,
344 priority: Option<i32>,
345 input_data: Option<serde_json::Value>,
346 result_data: Option<serde_json::Value>,
347 error_message: Option<String>,
348 started_at: Option<DateTime<Utc>>,
349 completed_at: Option<DateTime<Utc>>,
350 retry_count: Option<i32>,
351 max_retries: Option<i32>,
352}
353
354impl ProcessingJobBuilder {
355 pub fn file_id(mut self, value: Uuid) -> Self {
357 self.file_id = Some(value);
358 self
359 }
360
361 pub fn job_type(mut self, value: ProcessingJobType) -> Self {
363 self.job_type = Some(value);
364 self
365 }
366
367 pub fn status(mut self, value: JobStatus) -> Self {
369 self.status = Some(value);
370 self
371 }
372
373 pub fn priority(mut self, value: i32) -> Self {
375 self.priority = Some(value);
376 self
377 }
378
379 pub fn input_data(mut self, value: serde_json::Value) -> Self {
381 self.input_data = Some(value);
382 self
383 }
384
385 pub fn result_data(mut self, value: serde_json::Value) -> Self {
387 self.result_data = Some(value);
388 self
389 }
390
391 pub fn error_message(mut self, value: String) -> Self {
393 self.error_message = Some(value);
394 self
395 }
396
397 pub fn started_at(mut self, value: DateTime<Utc>) -> Self {
399 self.started_at = Some(value);
400 self
401 }
402
403 pub fn completed_at(mut self, value: DateTime<Utc>) -> Self {
405 self.completed_at = Some(value);
406 self
407 }
408
409 pub fn retry_count(mut self, value: i32) -> Self {
411 self.retry_count = Some(value);
412 self
413 }
414
415 pub fn max_retries(mut self, value: i32) -> Self {
417 self.max_retries = Some(value);
418 self
419 }
420
421 pub fn build(self) -> Result<ProcessingJob, String> {
425 let file_id = self.file_id.ok_or_else(|| "file_id is required".to_string())?;
426 let job_type = self.job_type.ok_or_else(|| "job_type is required".to_string())?;
427
428 Ok(ProcessingJob {
429 id: Uuid::new_v4(),
430 file_id,
431 job_type,
432 status: self.status.unwrap_or_default(),
433 priority: self.priority.unwrap_or(0),
434 input_data: self.input_data,
435 result_data: self.result_data,
436 error_message: self.error_message,
437 started_at: self.started_at,
438 completed_at: self.completed_at,
439 retry_count: self.retry_count.unwrap_or(0),
440 max_retries: self.max_retries.unwrap_or(3),
441 metadata: AuditMetadata::default(),
442 })
443 }
444}
445
446#[path = "processing_job.ext.rs"]
448mod processing_job_ext;
449pub use processing_job_ext::*;