Skip to main content

solti_model/resource/
status.rs

1//! # Task status
2//!
3//! [`TaskStatus`] is observed reconciliation and execution state.
4
5use std::collections::HashSet;
6
7use serde::{Deserialize, Serialize};
8
9use crate::{ConditionStatus, ModelError, ModelResult, TaskCondition, TaskPhase};
10
11/// Observed runtime state of a task.
12#[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    /// Creates an unobserved pending status for a desired generation.
38    ///
39    /// `observedGeneration` starts at zero.
40    /// The `Reconciled` condition refers to `desired_generation`.
41    ///
42    /// # Errors
43    ///
44    /// Returns [`ModelError::Invalid`] when `desired_generation` is zero.
45    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    /// Reconstructs status from serialized fields.
62    ///
63    /// # Errors
64    ///
65    /// Returns [`ModelError::Invalid`] when status fields are inconsistent.
66    /// This includes lifecycle fields, conditions, and generations.
67    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    /// Latest generation processed by the controller.
101    pub fn observed_generation(&self) -> u64 {
102        self.observed_generation
103    }
104
105    /// Current lifecycle phase.
106    pub fn phase(&self) -> TaskPhase {
107        self.phase
108    }
109
110    /// Current attempt number.
111    ///
112    /// Zero means no attempt has started.
113    pub fn attempt(&self) -> u32 {
114        self.attempt
115    }
116
117    /// Process exit code, when available.
118    pub fn exit_code(&self) -> Option<i32> {
119        self.exit_code
120    }
121
122    /// Current lifecycle diagnostic, when available.
123    pub fn error(&self) -> Option<&str> {
124        self.error.as_deref()
125    }
126
127    /// All status conditions.
128    pub fn conditions(&self) -> &[TaskCondition] {
129        &self.conditions
130    }
131
132    /// Returns the serialized status fields.
133    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    /// Returns a condition by type.
154    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    /// Returns the required `Reconciled` condition.
161    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    /// Returns whether reconciliation failed.
169    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}