1use std::collections::HashSet;
6
7use serde::{Deserialize, Serialize};
8
9use crate::{ConditionStatus, ModelError, ModelResult, TaskCondition, TaskPhase};
10
11#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
13#[serde(rename_all = "camelCase", try_from = "raw::TaskStatusRaw")]
14pub struct TaskStatus {
15 pub(crate) observed_generation: u64,
16 pub(crate) phase: TaskPhase,
17 pub(crate) attempt: u32,
18 #[serde(skip_serializing_if = "Option::is_none")]
19 pub(crate) exit_code: Option<i32>,
20 #[serde(skip_serializing_if = "Option::is_none")]
21 pub(crate) error: Option<String>,
22 pub(crate) conditions: Vec<TaskCondition>,
23}
24
25#[cfg(feature = "schema")]
26impl schemars::JsonSchema for TaskStatus {
27 fn schema_name() -> std::borrow::Cow<'static, str> {
28 "TaskStatus".into()
29 }
30
31 fn json_schema(generator: &mut schemars::SchemaGenerator) -> schemars::Schema {
32 crate::schema::task_status(generator)
33 }
34}
35
36impl TaskStatus {
37 pub fn pending(desired_generation: u64) -> ModelResult<Self> {
46 if desired_generation == 0 {
47 return Err(ModelError::Invalid(
48 "desired generation must be greater than zero".into(),
49 ));
50 }
51 Ok(Self {
52 observed_generation: 0,
53 phase: TaskPhase::Pending,
54 exit_code: None,
55 error: None,
56 attempt: 0,
57 conditions: vec![TaskCondition::reconciled_unknown(desired_generation)],
58 })
59 }
60
61 pub fn from_parts(
68 observed_generation: u64,
69 phase: TaskPhase,
70 attempt: u32,
71 exit_code: Option<i32>,
72 error: Option<String>,
73 conditions: Vec<TaskCondition>,
74 ) -> ModelResult<Self> {
75 let status = Self {
76 observed_generation,
77 phase,
78 attempt,
79 exit_code,
80 error,
81 conditions,
82 };
83 status.validate()?;
84 Ok(status)
85 }
86
87 pub(crate) fn pending_after(&self, desired_generation: u64) -> Self {
88 let mut pending = Self {
89 observed_generation: self.observed_generation,
90 phase: TaskPhase::Pending,
91 exit_code: None,
92 error: None,
93 attempt: 0,
94 conditions: self.conditions.clone(),
95 };
96 pending.mark_reconciliation_pending(desired_generation);
97 pending
98 }
99
100 pub fn observed_generation(&self) -> u64 {
102 self.observed_generation
103 }
104
105 pub fn phase(&self) -> TaskPhase {
107 self.phase
108 }
109
110 pub fn attempt(&self) -> u32 {
114 self.attempt
115 }
116
117 pub fn exit_code(&self) -> Option<i32> {
119 self.exit_code
120 }
121
122 pub fn error(&self) -> Option<&str> {
124 self.error.as_deref()
125 }
126
127 pub fn conditions(&self) -> &[TaskCondition] {
129 &self.conditions
130 }
131
132 pub fn into_parts(
134 self,
135 ) -> (
136 u64,
137 TaskPhase,
138 u32,
139 Option<i32>,
140 Option<String>,
141 Vec<TaskCondition>,
142 ) {
143 (
144 self.observed_generation,
145 self.phase,
146 self.attempt,
147 self.exit_code,
148 self.error,
149 self.conditions,
150 )
151 }
152
153 pub fn condition(&self, condition_type: &crate::TaskConditionType) -> Option<&TaskCondition> {
155 self.conditions
156 .iter()
157 .find(|condition| condition.condition_type() == condition_type)
158 }
159
160 pub fn reconciled(&self) -> &TaskCondition {
162 self.conditions
163 .iter()
164 .find(|condition| condition.condition_type().is_reconciled())
165 .expect("validated TaskStatus has a Reconciled condition")
166 }
167
168 pub fn reconciliation_failed(&self) -> bool {
170 self.reconciled().status() == ConditionStatus::False
171 }
172
173 pub(crate) fn validate(&self) -> ModelResult<()> {
174 let mut condition_types = HashSet::with_capacity(self.conditions.len());
175 let mut reconciled_count = 0;
176 for condition in &self.conditions {
177 condition.validate()?;
178 if !condition_types.insert(condition.condition_type().as_str()) {
179 return Err(ModelError::Invalid(
180 format!(
181 "status.conditions contains duplicate type `{}`",
182 condition.condition_type()
183 )
184 .into(),
185 ));
186 }
187 if condition.condition_type().is_reconciled() {
188 reconciled_count += 1;
189 }
190 }
191 if reconciled_count != 1 {
192 return Err(ModelError::Invalid(
193 "status.conditions must contain one Reconciled condition".into(),
194 ));
195 }
196 let reconciled = self.reconciled();
197 if reconciled.observed_generation() == 0 {
198 return Err(ModelError::Invalid(
199 "status.conditions[type=Reconciled].observedGeneration must be greater than zero"
200 .into(),
201 ));
202 }
203 if reconciled.observed_generation() < self.observed_generation {
204 return Err(ModelError::Invalid(
205 "status.conditions[type=Reconciled].observedGeneration cannot be less than status.observedGeneration"
206 .into(),
207 ));
208 }
209 if reconciled.status() != ConditionStatus::Unknown
210 && reconciled.observed_generation() != self.observed_generation
211 {
212 return Err(ModelError::Invalid(
213 "status.conditions[type=Reconciled].observedGeneration must equal status.observedGeneration when Reconciled is True or False"
214 .into(),
215 ));
216 }
217 if self.phase != TaskPhase::Pending && reconciled.status() != ConditionStatus::True {
218 return Err(ModelError::Invalid(
219 "status.phase requires a Reconciled=True condition unless phase is pending".into(),
220 ));
221 }
222 Self::validate_execution_fields(
223 self.phase,
224 self.attempt,
225 self.exit_code,
226 self.error.as_deref(),
227 )?;
228 Ok(())
229 }
230
231 fn validate_execution_fields(
232 phase: TaskPhase,
233 attempt: u32,
234 exit_code: Option<i32>,
235 error: Option<&str>,
236 ) -> ModelResult<()> {
237 match phase {
238 TaskPhase::Pending => {
239 if attempt != 0 {
240 return Err(ModelError::Invalid(
241 "status.attempt must be zero while status.phase is pending".into(),
242 ));
243 }
244 if exit_code.is_some() || error.is_some() {
245 return Err(ModelError::Invalid(
246 "status.exitCode and status.error must be absent while status.phase is pending"
247 .into(),
248 ));
249 }
250 }
251 TaskPhase::Running => {
252 if attempt == 0 {
253 return Err(ModelError::Invalid(
254 "status.attempt must be greater than zero while status.phase is running"
255 .into(),
256 ));
257 }
258 if exit_code.is_some() || error.is_some() {
259 return Err(ModelError::Invalid(
260 "status.exitCode and status.error must be absent while status.phase is running"
261 .into(),
262 ));
263 }
264 }
265 TaskPhase::Succeeded
266 | TaskPhase::Failed
267 | TaskPhase::Timeout
268 | TaskPhase::Canceled
269 | TaskPhase::Exhausted => {}
270 }
271 Ok(())
272 }
273
274 pub(crate) fn reconciled_required(&self) -> &TaskCondition {
275 self.reconciled()
276 }
277
278 pub(crate) fn mark_reconciliation_pending(&mut self, generation: u64) -> bool {
279 self.reconciled_mut().transition(
280 ConditionStatus::Unknown,
281 generation,
282 "ReconciliationScheduled",
283 "runtime reconciliation is scheduled",
284 )
285 }
286
287 pub(crate) fn mark_reconciled(&mut self, generation: u64) -> bool {
288 let changed = self.reconciled_mut().transition(
289 ConditionStatus::True,
290 generation,
291 "RuntimeAccepted",
292 "runtime accepted the desired state",
293 );
294 self.observed_generation = generation;
295 changed
296 }
297
298 pub(crate) fn mark_reconciliation_failed(
299 &mut self,
300 generation: u64,
301 reason: impl Into<String>,
302 message: impl Into<String>,
303 ) -> bool {
304 let changed =
305 self.reconciled_mut()
306 .transition(ConditionStatus::False, generation, reason, message);
307 self.observed_generation = generation;
308 changed
309 }
310
311 fn reconciled_mut(&mut self) -> &mut TaskCondition {
312 self.conditions
313 .iter_mut()
314 .find(|condition| condition.condition_type().is_reconciled())
315 .expect("validated TaskStatus has a Reconciled condition")
316 }
317}
318
319mod raw {
320 use super::*;
321
322 #[derive(Deserialize)]
323 #[serde(rename_all = "camelCase", deny_unknown_fields)]
324 pub(super) struct TaskStatusRaw {
325 observed_generation: u64,
326 phase: TaskPhase,
327 attempt: u32,
328 #[serde(default)]
329 exit_code: Option<i32>,
330 #[serde(default)]
331 error: Option<String>,
332 conditions: Vec<TaskCondition>,
333 }
334
335 impl TryFrom<TaskStatusRaw> for TaskStatus {
336 type Error = ModelError;
337
338 fn try_from(raw: TaskStatusRaw) -> Result<Self, Self::Error> {
339 TaskStatus::from_parts(
340 raw.observed_generation,
341 raw.phase,
342 raw.attempt,
343 raw.exit_code,
344 raw.error,
345 raw.conditions,
346 )
347 }
348 }
349}
350
351#[cfg(test)]
352mod tests {
353 use super::*;
354 use crate::TaskConditionType;
355 use std::time::SystemTime;
356
357 fn condition(condition_type: TaskConditionType, status: ConditionStatus) -> TaskCondition {
358 TaskCondition::new(
359 condition_type,
360 status,
361 1,
362 SystemTime::UNIX_EPOCH,
363 "Observed",
364 "observed state",
365 )
366 .unwrap()
367 }
368
369 #[test]
370 fn pending_generation_is_explicit() {
371 let status = TaskStatus::pending(3).unwrap();
372 assert_eq!(status.phase(), TaskPhase::Pending);
373 assert_eq!(status.observed_generation(), 0);
374 assert_eq!(status.attempt(), 0);
375 assert!(status.error().is_none());
376 assert_eq!(status.reconciled().status(), ConditionStatus::Unknown);
377 assert_eq!(status.reconciled().observed_generation(), 3);
378 assert!(TaskStatus::pending(0).is_err());
379 }
380
381 #[test]
382 fn standalone_status_rejects_missing_reconciled_condition() {
383 let json = serde_json::json!({
384 "observedGeneration": 0,
385 "phase": "pending",
386 "attempt": 0,
387 "conditions": []
388 });
389 assert!(serde_json::from_value::<TaskStatus>(json).is_err());
390 }
391
392 #[test]
393 fn status_accepts_one_reconciled_and_extensible_conditions() {
394 let reconciled = condition(TaskConditionType::reconciled(), ConditionStatus::True);
395 let available_type = TaskConditionType::new("Available").unwrap();
396 let available = condition(available_type.clone(), ConditionStatus::False);
397
398 let status = TaskStatus::from_parts(
399 1,
400 TaskPhase::Running,
401 1,
402 None,
403 None,
404 vec![reconciled, available],
405 )
406 .unwrap();
407
408 assert_eq!(status.conditions().len(), 2);
409 assert_eq!(
410 status.condition(&available_type).unwrap().status(),
411 ConditionStatus::False
412 );
413 let back: TaskStatus =
414 serde_json::from_value(serde_json::to_value(&status).unwrap()).unwrap();
415 assert_eq!(back, status);
416 }
417
418 #[test]
419 fn status_rejects_duplicate_condition_types() {
420 let reconciled = condition(TaskConditionType::reconciled(), ConditionStatus::True);
421 let duplicate = reconciled.clone();
422
423 assert!(
424 TaskStatus::from_parts(
425 1,
426 TaskPhase::Running,
427 1,
428 None,
429 None,
430 vec![reconciled, duplicate],
431 )
432 .is_err()
433 );
434 }
435
436 #[test]
437 fn status_rejects_inconsistent_phase_fields() {
438 let cases = [
439 (TaskPhase::Pending, 1, None, None),
440 (TaskPhase::Pending, 0, Some(0), None),
441 (TaskPhase::Pending, 0, None, Some("error".into())),
442 (TaskPhase::Running, 0, None, None),
443 (TaskPhase::Running, 1, Some(0), None),
444 (TaskPhase::Running, 1, None, Some("error".into())),
445 ];
446 for (phase, attempt, exit_code, error) in cases {
447 assert!(
448 TaskStatus::from_parts(
449 1,
450 phase,
451 attempt,
452 exit_code,
453 error,
454 vec![condition(
455 TaskConditionType::reconciled(),
456 ConditionStatus::True
457 )],
458 )
459 .is_err()
460 );
461 }
462 }
463
464 #[test]
465 fn status_enforces_reconciled_generation_contract() {
466 let reconciled = |status, observed_generation| {
467 TaskCondition::new(
468 TaskConditionType::reconciled(),
469 status,
470 observed_generation,
471 SystemTime::UNIX_EPOCH,
472 "Observed",
473 "observed state",
474 )
475 .unwrap()
476 };
477
478 assert!(
479 TaskStatus::from_parts(
480 1,
481 TaskPhase::Pending,
482 0,
483 None,
484 None,
485 vec![reconciled(ConditionStatus::Unknown, 2)],
486 )
487 .is_ok()
488 );
489 for (observed_generation, phase, condition_status, condition_generation) in [
490 (0, TaskPhase::Pending, ConditionStatus::Unknown, 0),
491 (2, TaskPhase::Pending, ConditionStatus::Unknown, 1),
492 (1, TaskPhase::Pending, ConditionStatus::True, 2),
493 (1, TaskPhase::Running, ConditionStatus::Unknown, 1),
494 (1, TaskPhase::Failed, ConditionStatus::False, 1),
495 ] {
496 assert!(
497 TaskStatus::from_parts(
498 observed_generation,
499 phase,
500 if phase == TaskPhase::Running { 1 } else { 0 },
501 None,
502 None,
503 vec![reconciled(condition_status, condition_generation)],
504 )
505 .is_err()
506 );
507 }
508 }
509
510 #[test]
511 fn terminal_status_allows_an_unknown_attempt() {
512 let status = TaskStatus::from_parts(
513 1,
514 TaskPhase::Failed,
515 0,
516 None,
517 Some("submission failed before an attempt started".into()),
518 vec![condition(
519 TaskConditionType::reconciled(),
520 ConditionStatus::True,
521 )],
522 )
523 .unwrap();
524
525 assert_eq!(status.attempt(), 0);
526 }
527
528 #[test]
529 fn status_and_conditions_reject_unknown_fields() {
530 let mut status = serde_json::to_value(TaskStatus::pending(1).unwrap()).unwrap();
531 status["unexpected"] = serde_json::json!(true);
532 assert!(serde_json::from_value::<TaskStatus>(status).is_err());
533
534 let mut status = serde_json::to_value(TaskStatus::pending(1).unwrap()).unwrap();
535 status["conditions"][0]["unexpected"] = serde_json::json!(true);
536 assert!(serde_json::from_value::<TaskStatus>(status).is_err());
537 }
538}