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 pub fn can_retry(&self) -> bool {
284 self.status == JobStatus::Failed && self.retry_count < self.max_retries
285 }
286
287 pub fn increment_retry(&mut self) {
289 self.retry_count += 1;
290 self.metadata.touch();
291 }
292
293 pub fn mark_started(&mut self) -> Result<(), JobError> {
295 self.status = JobStatus::Running;
296 self.started_at = Some(Utc::now());
297 self.metadata.touch();
298 Ok(())
299 }
300
301 pub fn mark_completed(&mut self, result: serde_json::Value) -> Result<(), JobError> {
303 self.status = JobStatus::Completed;
304 self.result_data = Some(result);
305 self.completed_at = Some(Utc::now());
306 self.metadata.touch();
307 Ok(())
308 }
309
310 pub fn mark_failed(&mut self, error: String) -> Result<(), JobError> {
312 self.status = JobStatus::Failed;
313 self.error_message = Some(error);
314 self.metadata.touch();
315 Ok(())
316 }
317
318 pub fn cancel(&mut self) -> Result<(), JobError> {
320 self.status = JobStatus::Cancelled;
321 self.metadata.touch();
322 Ok(())
323 }
324
325 pub fn duration(&self) -> Option<Duration> {
327 match (self.started_at, self.completed_at) {
328 (Some(start), Some(end)) => Some(end - start),
329 _ => None,
330 }
331 }
332
333 pub fn check_invariants(&self) -> Result<(), Vec<&'static str>> {
335 let mut errors = Vec::new();
336 if self.retry_count > self.max_retries {
337 errors.push("retry_count must not exceed max_retries");
338 }
339 if self.status == JobStatus::Running && self.started_at.is_none() {
340 errors.push("running job must have started_at");
341 }
342 if self.status == JobStatus::Completed && self.completed_at.is_none() {
343 errors.push("completed job must have completed_at");
344 }
345 if errors.is_empty() { Ok(()) } else { Err(errors) }
346 }
347 }
349
350impl super::Entity for ProcessingJob {
351 type Id = Uuid;
352
353 fn entity_id(&self) -> &Self::Id {
354 &self.id
355 }
356
357 fn entity_type() -> &'static str {
358 "ProcessingJob"
359 }
360}
361
362impl backbone_core::PersistentEntity for ProcessingJob {
363 fn entity_id(&self) -> String {
364 self.id.to_string()
365 }
366 fn set_entity_id(&mut self, id: String) {
367 if let Ok(uuid) = uuid::Uuid::parse_str(&id) {
368 self.id = uuid;
369 }
370 }
371 fn created_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
372 self.metadata.created_at
373 }
374 fn set_created_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
375 self.metadata.created_at = Some(ts);
376 }
377 fn updated_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
378 self.metadata.updated_at
379 }
380 fn set_updated_at(&mut self, ts: chrono::DateTime<chrono::Utc>) {
381 self.metadata.updated_at = Some(ts);
382 }
383 fn deleted_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
384 self.metadata.deleted_at
385 }
386 fn set_deleted_at(&mut self, ts: Option<chrono::DateTime<chrono::Utc>>) {
387 self.metadata.deleted_at = ts;
388 }
389}
390
391impl backbone_orm::EntityRepoMeta for ProcessingJob {
392 fn column_types() -> std::collections::HashMap<String, String> {
393 let mut m = std::collections::HashMap::new();
394 m.insert("id".to_string(), "uuid".to_string());
395 m.insert("file_id".to_string(), "uuid".to_string());
396 m.insert("job_type".to_string(), "processing_job_type".to_string());
397 m.insert("status".to_string(), "job_status".to_string());
398 m.insert("started_at".to_string(), "timestamptz".to_string());
399 m.insert("completed_at".to_string(), "timestamptz".to_string());
400 m
401 }
402 fn search_fields() -> &'static [&'static str] {
403 &[]
404 }
405}
406
407#[derive(Debug, Clone, Default)]
412pub struct ProcessingJobBuilder {
413 file_id: Option<Uuid>,
414 job_type: Option<ProcessingJobType>,
415 status: Option<JobStatus>,
416 priority: Option<i32>,
417 input_data: Option<serde_json::Value>,
418 result_data: Option<serde_json::Value>,
419 error_message: Option<String>,
420 started_at: Option<DateTime<Utc>>,
421 completed_at: Option<DateTime<Utc>>,
422 retry_count: Option<i32>,
423 max_retries: Option<i32>,
424}
425
426impl ProcessingJobBuilder {
427 pub fn file_id(mut self, value: Uuid) -> Self {
429 self.file_id = Some(value);
430 self
431 }
432
433 pub fn job_type(mut self, value: ProcessingJobType) -> Self {
435 self.job_type = Some(value);
436 self
437 }
438
439 pub fn status(mut self, value: JobStatus) -> Self {
441 self.status = Some(value);
442 self
443 }
444
445 pub fn priority(mut self, value: i32) -> Self {
447 self.priority = Some(value);
448 self
449 }
450
451 pub fn input_data(mut self, value: serde_json::Value) -> Self {
453 self.input_data = Some(value);
454 self
455 }
456
457 pub fn result_data(mut self, value: serde_json::Value) -> Self {
459 self.result_data = Some(value);
460 self
461 }
462
463 pub fn error_message(mut self, value: String) -> Self {
465 self.error_message = Some(value);
466 self
467 }
468
469 pub fn started_at(mut self, value: DateTime<Utc>) -> Self {
471 self.started_at = Some(value);
472 self
473 }
474
475 pub fn completed_at(mut self, value: DateTime<Utc>) -> Self {
477 self.completed_at = Some(value);
478 self
479 }
480
481 pub fn retry_count(mut self, value: i32) -> Self {
483 self.retry_count = Some(value);
484 self
485 }
486
487 pub fn max_retries(mut self, value: i32) -> Self {
489 self.max_retries = Some(value);
490 self
491 }
492
493 pub fn build(self) -> Result<ProcessingJob, String> {
497 let file_id = self.file_id.ok_or_else(|| "file_id is required".to_string())?;
498 let job_type = self.job_type.ok_or_else(|| "job_type is required".to_string())?;
499
500 Ok(ProcessingJob {
501 id: Uuid::new_v4(),
502 file_id,
503 job_type,
504 status: self.status.unwrap_or_default(),
505 priority: self.priority.unwrap_or(0),
506 input_data: self.input_data,
507 result_data: self.result_data,
508 error_message: self.error_message,
509 started_at: self.started_at,
510 completed_at: self.completed_at,
511 retry_count: self.retry_count.unwrap_or(0),
512 max_retries: self.max_retries.unwrap_or(3),
513 metadata: AuditMetadata::default(),
514 })
515 }
516}