1use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use uuid::Uuid;
7
8use crate::JobId;
9
10#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
11#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
12#[serde(rename_all = "lowercase")]
13pub enum JobPriority {
14 Low,
15 #[default]
16 Normal,
17 High,
18 Critical,
19}
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
22#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
23#[serde(rename_all = "lowercase")]
24pub enum JobStatus {
25 Pending,
26 Running,
27 Completed,
28 Failed,
29 Cancelled,
30}
31
32#[derive(Debug, Clone, Serialize, Deserialize)]
33#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
34#[serde(rename_all = "camelCase", deny_unknown_fields)]
35pub struct SubmitJobRequest {
36 pub name: String,
37 #[serde(default, skip_serializing_if = "Option::is_none")]
38 pub payload: Option<Value>,
39 #[serde(default, skip_serializing_if = "Option::is_none")]
40 pub priority: Option<JobPriority>,
41 #[serde(default, skip_serializing_if = "Option::is_none")]
42 pub max_retries: Option<u32>,
43 #[serde(default, skip_serializing_if = "Option::is_none")]
44 pub timeout_ms: Option<u64>,
45}
46
47impl SubmitJobRequest {
48 pub(crate) fn validate(&self) -> Result<(), &'static str> {
49 if self.name.trim().is_empty() {
50 return Err("job name must not be empty");
51 }
52 Ok(())
53 }
54}
55
56#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
57#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
58#[serde(rename_all = "camelCase", deny_unknown_fields)]
59pub struct JobSubmitResponse {
60 pub job_id: JobId,
61 pub status: JobStatus,
62}
63
64#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
65#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
66#[serde(rename_all = "camelCase", deny_unknown_fields)]
67pub struct JobResult {
68 pub job_id: JobId,
69 pub success: bool,
70 pub output: Value,
71 pub error: Option<String>,
72 pub completed_at: DateTime<Utc>,
73}
74
75#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
76#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
77#[serde(rename_all = "camelCase", deny_unknown_fields)]
78pub struct JobStatusResponse {
79 pub id: JobId,
80 pub name: String,
81 pub priority: JobPriority,
82 pub status: JobStatus,
83 pub payload: Value,
84 pub created_at: DateTime<Utc>,
85 pub started_at: Option<DateTime<Utc>>,
86 pub completed_at: Option<DateTime<Utc>>,
87 pub retry_count: u32,
88 pub max_retries: u32,
89 pub result: Option<JobResult>,
90 pub error: Option<String>,
91 pub timeout_ms: Option<u64>,
92 pub correlation_id: Option<Uuid>,
93}
94
95impl JobStatusResponse {
96 pub(crate) fn validate(&self) -> Result<(), &'static str> {
97 if self.name.trim().is_empty() || self.retry_count > self.max_retries {
98 return Err("job status response violates bounded fields");
99 }
100 if self
101 .result
102 .as_ref()
103 .is_some_and(|result| result.job_id != self.id)
104 {
105 return Err("job result identifier does not match its status");
106 }
107
108 let valid = match self.status {
109 JobStatus::Pending => {
110 self.started_at.is_none()
111 && self.completed_at.is_none()
112 && self.result.is_none()
113 && self.error.is_none()
114 }
115 JobStatus::Running => {
116 self.started_at.is_some()
117 && self.completed_at.is_none()
118 && self.result.is_none()
119 && self.error.is_none()
120 }
121 JobStatus::Completed => {
122 self.started_at.is_some()
123 && self.completed_at.is_some()
124 && self
125 .result
126 .as_ref()
127 .is_some_and(|result| result.success && result.error.is_none())
128 && self.error.is_none()
129 }
130 JobStatus::Failed => {
131 self.started_at.is_some()
132 && self.completed_at.is_some()
133 && self.result.is_none()
134 && self.error.as_ref().is_some_and(|error| !error.is_empty())
135 }
136 JobStatus::Cancelled => {
137 self.completed_at.is_some() && self.result.is_none() && self.error.is_none()
138 }
139 };
140 valid
141 .then_some(())
142 .ok_or("job status response violates lifecycle invariants")
143 }
144}
145
146#[cfg(test)]
147mod tests {
148 use super::*;
149 use serde_json::json;
150
151 #[test]
152 fn job_identifiers_require_the_runtime_prefix_and_uuid() {
153 assert!(serde_json::from_value::<JobId>(json!(format!("job_{}", Uuid::new_v4()))).is_ok());
154 assert!(serde_json::from_value::<JobId>(json!(Uuid::new_v4().to_string())).is_err());
155 assert!(serde_json::from_value::<JobId>(json!("job_not-a-uuid")).is_err());
156 }
157
158 #[test]
159 fn job_status_validation_enforces_lifecycle_invariants() {
160 let id = JobId::new();
161 let mut response: JobStatusResponse = serde_json::from_value(json!({
162 "id": id,
163 "name": "fetch-report",
164 "priority": "normal",
165 "status": "running",
166 "payload": null,
167 "createdAt": "2026-09-19T00:00:00Z",
168 "startedAt": "2026-09-19T00:00:01Z",
169 "completedAt": null,
170 "retryCount": 0,
171 "maxRetries": 3,
172 "result": null,
173 "error": null,
174 "timeoutMs": null,
175 "correlationId": null,
176 }))
177 .unwrap();
178 assert!(response.validate().is_ok());
179 response.started_at = None;
180 assert!(response.validate().is_err());
181 }
182}