Skip to main content

solti_model/resource/
task.rs

1//! # Task resource
2//!
3//! [`TaskManifest`] is caller-owned desired state.
4//! [`Task`] is a stored resource with server metadata and status.
5//!
6//! ## Apply
7//!
8//! ```text
9//! identical desired state ──▶ DesiredChange::None
10//! labels or annotations   ──▶ DesiredChange::Metadata
11//! spec changed            ──▶ DesiredChange::Spec
12//!                              └─ generation increments
13//! ```
14//!
15//! ## Status Flow
16//!
17//! ```text
18//! Reconciled: Unknown        ── accepted     ────▶ True
19//!             False          ── manual retry ────▶ Unknown
20//!             Unknown | True ── failure      ────▶ False
21//!
22//! Pending ── attempt starts ──▶ Running ── attempt ends ──▶ terminal phase
23//! ```
24//!
25//! Generation is checked before attempt transitions.
26//! Stale generation updates are ignored.
27//! Repeating an identical update is a no-op.
28
29use serde::{Deserialize, Serialize};
30
31use crate::{
32    Annotations, ConditionStatus, Labels, ModelError, ModelResult, ObjectMeta, Slot, TaskCondition,
33    TaskId, TaskPhase, TaskSpec, TaskStatus, Uid,
34};
35
36macro_rules! task_api_major {
37    () => {
38        1
39    };
40}
41
42/// Major version of the built-in Task resource API.
43pub const TASK_API_VERSION_MAJOR: u32 = task_api_major!();
44
45/// API group and version of the built-in Task resource.
46pub const TASK_API_VERSION: &str = concat!("solti.io/v", task_api_major!());
47
48/// Kind of the built-in Task resource.
49pub const TASK_KIND: &str = "Task";
50
51/// Classification of an apply operation.
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub enum DesiredChange {
54    /// Desired state and user-owned metadata were already identical.
55    None,
56    /// Only labels and/or annotations changed.
57    Metadata,
58    /// Spec changed, optionally together with labels or annotations.
59    Spec,
60}
61
62impl DesiredChange {
63    /// Returns whether apply changed the resource.
64    #[inline]
65    pub fn is_changed(self) -> bool {
66        !matches!(self, Self::None)
67    }
68}
69
70/// Group/version and kind of resource schema.
71#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
72#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
73#[serde(rename_all = "camelCase", deny_unknown_fields)]
74pub struct TypeMeta {
75    #[cfg_attr(
76        feature = "schema",
77        schemars(schema_with = "crate::schema::task_api_version")
78    )]
79    api_version: String,
80    #[cfg_attr(feature = "schema", schemars(schema_with = "crate::schema::task_kind"))]
81    kind: String,
82}
83
84/// User-owned metadata accepted in a [`TaskManifest`].
85///
86/// Runtime identity, resource version, generation, creation time and status are deliberately absent because the state store owns them.
87#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
88#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
89#[serde(rename_all = "camelCase", deny_unknown_fields)]
90pub struct TaskManifestMeta {
91    name: TaskId,
92    #[serde(default, skip_serializing_if = "Labels::is_empty")]
93    labels: Labels,
94    #[serde(default, skip_serializing_if = "Annotations::is_empty")]
95    annotations: Annotations,
96}
97
98impl TaskManifestMeta {
99    /// Creates user-owned metadata.
100    ///
101    /// # Errors
102    ///
103    /// Returns [`ModelError::Invalid`] when `name` is invalid.
104    pub fn new(name: impl AsRef<str>) -> ModelResult<Self> {
105        let metadata = Self {
106            name: TaskId::new(name)?,
107            labels: Labels::new(),
108            annotations: Annotations::new(),
109        };
110        metadata.name.validate_format()?;
111        Ok(metadata)
112    }
113
114    /// Stable resource name.
115    #[inline]
116    pub fn name(&self) -> &TaskId {
117        &self.name
118    }
119
120    /// Selector metadata.
121    #[inline]
122    pub fn labels(&self) -> &Labels {
123        &self.labels
124    }
125
126    /// Free-form metadata.
127    #[inline]
128    pub fn annotations(&self) -> &Annotations {
129        &self.annotations
130    }
131
132    fn with_labels(mut self, labels: Labels) -> Self {
133        self.labels = labels;
134        self
135    }
136
137    fn with_annotations(mut self, annotations: Annotations) -> Self {
138        self.annotations = annotations;
139        self
140    }
141
142    fn validate(&self) -> ModelResult<()> {
143        self.name.validate_format()?;
144        self.labels.validate()?;
145        self.annotations.validate()
146    }
147}
148
149/// Caller-owned desired state for create and apply.
150///
151/// The serialized shape is `apiVersion`, `kind`, `metadata`, and `spec`.
152/// A stored [`Task`] adds server metadata and `status`.
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
155#[cfg_attr(feature = "schema", schemars(!try_from, deny_unknown_fields))]
156#[serde(rename_all = "camelCase")]
157#[serde(try_from = "raw::TaskManifestRaw")]
158pub struct TaskManifest {
159    #[serde(flatten)]
160    type_meta: TypeMeta,
161    metadata: TaskManifestMeta,
162    spec: TaskSpec,
163}
164
165impl TaskManifest {
166    /// Creates a Task manifest.
167    ///
168    /// # Errors
169    ///
170    /// Returns [`ModelError::Invalid`] when the name or spec is invalid.
171    pub fn new(name: impl AsRef<str>, spec: TaskSpec) -> ModelResult<Self> {
172        Self::from_parts(TypeMeta::task(), TaskManifestMeta::new(name)?, spec)
173    }
174
175    /// Reconstructs a manifest from serialized fields.
176    ///
177    /// # Errors
178    ///
179    /// Returns [`ModelError::Invalid`] when GVK, metadata, or spec is invalid.
180    pub fn from_parts(
181        type_meta: TypeMeta,
182        metadata: TaskManifestMeta,
183        spec: TaskSpec,
184    ) -> ModelResult<Self> {
185        let manifest = Self {
186            type_meta,
187            metadata,
188            spec,
189        };
190        manifest.validate()?;
191        Ok(manifest)
192    }
193
194    /// Sets manifest labels.
195    ///
196    /// # Errors
197    ///
198    /// Returns [`ModelError::Invalid`] when a label is invalid.
199    pub fn with_labels(mut self, labels: Labels) -> ModelResult<Self> {
200        labels.validate()?;
201        self.metadata = self.metadata.with_labels(labels);
202        Ok(self)
203    }
204
205    /// Sets manifest annotations.
206    ///
207    /// # Errors
208    ///
209    /// Returns [`ModelError::Invalid`] when an annotation is invalid.
210    pub fn with_annotations(mut self, annotations: Annotations) -> ModelResult<Self> {
211        annotations.validate()?;
212        self.metadata = self.metadata.with_annotations(annotations);
213        Ok(self)
214    }
215
216    /// Validates the manifest.
217    ///
218    /// # Errors
219    ///
220    /// Returns [`ModelError::Invalid`] when GVK, metadata, or spec is invalid.
221    pub fn validate(&self) -> ModelResult<()> {
222        self.type_meta.validate_task()?;
223        self.metadata.validate()?;
224        self.spec.validate()
225    }
226
227    /// Resource type metadata.
228    #[inline]
229    pub fn type_meta(&self) -> &TypeMeta {
230        &self.type_meta
231    }
232
233    /// User-owned resource metadata.
234    #[inline]
235    pub fn metadata(&self) -> &TaskManifestMeta {
236        &self.metadata
237    }
238
239    /// Stable resource name.
240    #[inline]
241    pub fn name(&self) -> &TaskId {
242        self.metadata.name()
243    }
244
245    /// Desired state.
246    #[inline]
247    pub fn spec(&self) -> &TaskSpec {
248        &self.spec
249    }
250
251    /// Logical concurrency slot.
252    #[inline]
253    pub fn slot(&self) -> &Slot {
254        self.spec.slot()
255    }
256
257    /// Returns the serialized manifest fields.
258    pub fn into_parts(self) -> (TypeMeta, TaskManifestMeta, TaskSpec) {
259        (self.type_meta, self.metadata, self.spec)
260    }
261}
262
263impl TypeMeta {
264    /// Type metadata for the built-in Task resource.
265    pub fn task() -> Self {
266        Self {
267            api_version: TASK_API_VERSION.to_owned(),
268            kind: TASK_KIND.to_owned(),
269        }
270    }
271
272    /// Resource API group and version.
273    #[inline]
274    pub fn api_version(&self) -> &str {
275        &self.api_version
276    }
277
278    /// Resource kind.
279    #[inline]
280    pub fn kind(&self) -> &str {
281        &self.kind
282    }
283
284    fn validate_task(&self) -> ModelResult<()> {
285        if self.api_version != TASK_API_VERSION {
286            return Err(ModelError::Invalid(
287                format!(
288                    "Task apiVersion must be `{TASK_API_VERSION}`, got `{}`",
289                    self.api_version
290                )
291                .into(),
292            ));
293        }
294        if self.kind != TASK_KIND {
295            return Err(ModelError::Invalid(
296                format!("Task kind must be `{TASK_KIND}`, got `{}`", self.kind).into(),
297            ));
298        }
299        Ok(())
300    }
301}
302
303/// Stored Task resource.
304///
305/// The serialized shape is `apiVersion`, `kind`, `metadata`, `spec`, `status`.
306/// Name, UID, and creation time are preserved by [`Self::apply_desired`].
307#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
308#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
309#[cfg_attr(feature = "schema", schemars(!try_from, deny_unknown_fields))]
310#[serde(rename_all = "camelCase")]
311#[serde(try_from = "raw::TaskRaw")]
312pub struct Task {
313    #[serde(flatten)]
314    type_meta: TypeMeta,
315    metadata: ObjectMeta,
316    spec: TaskSpec,
317    status: TaskStatus,
318}
319
320impl Task {
321    /// Creates a stored Task with server-owned defaults.
322    ///
323    /// The UID and creation timestamp are generated, generation starts at `1`, and status starts as unobserved `Pending`.
324    /// The state store assigns the initial resource version separately.
325    ///
326    /// # Errors
327    ///
328    /// Returns [`ModelError::Invalid`] when the name or spec is invalid, or the entropy source is unavailable.
329    pub fn new(name: impl AsRef<str>, spec: TaskSpec) -> ModelResult<Self> {
330        Self::from_manifest(TaskManifest::new(name, spec)?)
331    }
332
333    /// Creates a stored resource from a manifest.
334    ///
335    /// # Errors
336    ///
337    /// Returns [`ModelError::Invalid`] when the manifest is invalid or the entropy source is unavailable.
338    pub fn from_manifest(manifest: TaskManifest) -> ModelResult<Self> {
339        manifest.validate()?;
340        let (_, metadata, spec) = manifest.into_parts();
341        let mut object_meta = ObjectMeta::new(metadata.name.clone())?;
342        object_meta.apply_metadata(metadata.labels, metadata.annotations);
343        let task = Self {
344            type_meta: TypeMeta::task(),
345            metadata: object_meta,
346            spec,
347            status: TaskStatus::pending(1)?,
348        };
349        task.validate()?;
350        Ok(task)
351    }
352
353    /// Reconstructs a resource from persisted fields.
354    ///
355    /// # Errors
356    ///
357    /// Returns [`ModelError::Invalid`] when a resource invariant is violated.
358    pub fn from_parts(
359        type_meta: TypeMeta,
360        metadata: ObjectMeta,
361        spec: TaskSpec,
362        status: TaskStatus,
363    ) -> ModelResult<Self> {
364        let task = Self {
365            type_meta,
366            metadata,
367            spec,
368            status,
369        };
370        task.validate()?;
371        Ok(task)
372    }
373
374    /// Validates the complete resource.
375    ///
376    /// # Errors
377    ///
378    /// Returns [`ModelError::Invalid`] when GVK, metadata, spec, status, generation, or conditions are inconsistent.
379    pub fn validate(&self) -> ModelResult<()> {
380        self.type_meta.validate_task()?;
381        self.metadata.name().validate_format()?;
382        self.metadata.labels().validate()?;
383        self.metadata.annotations().validate()?;
384        if self.metadata.generation() == 0 {
385            return Err(ModelError::Invalid(
386                "metadata.generation must be greater than zero".into(),
387            ));
388        }
389        if self.status.observed_generation > self.metadata.generation() {
390            return Err(ModelError::Invalid(
391                "status.observedGeneration cannot exceed metadata.generation".into(),
392            ));
393        }
394        self.status.validate()?;
395        if self
396            .status
397            .conditions()
398            .iter()
399            .any(|condition| condition.observed_generation() > self.metadata.generation())
400        {
401            return Err(ModelError::Invalid(
402                "status.conditions[].observedGeneration cannot exceed metadata.generation".into(),
403            ));
404        }
405        self.spec.validate()
406    }
407
408    /// Resource type metadata.
409    #[inline]
410    pub fn type_meta(&self) -> &TypeMeta {
411        &self.type_meta
412    }
413
414    /// Resource metadata.
415    #[inline]
416    pub fn metadata(&self) -> &ObjectMeta {
417        &self.metadata
418    }
419
420    /// Desired state.
421    #[inline]
422    pub fn spec(&self) -> &TaskSpec {
423        &self.spec
424    }
425
426    /// Observed state.
427    #[inline]
428    pub fn status(&self) -> &TaskStatus {
429        &self.status
430    }
431
432    /// Returns the serialized resource fields.
433    pub fn into_parts(self) -> (TypeMeta, ObjectMeta, TaskSpec, TaskStatus) {
434        (self.type_meta, self.metadata, self.spec, self.status)
435    }
436
437    /// Stable resource address (`metadata.name`).
438    #[inline]
439    pub fn name(&self) -> &TaskId {
440        self.metadata.name()
441    }
442
443    /// Identity of this resource incarnation.
444    #[inline]
445    pub fn uid(&self) -> &Uid {
446        self.metadata.uid()
447    }
448
449    /// Logical concurrency slot.
450    #[inline]
451    pub fn slot(&self) -> &Slot {
452        self.spec.slot()
453    }
454
455    /// Resource labels.
456    #[inline]
457    pub fn labels(&self) -> &Labels {
458        self.metadata.labels()
459    }
460
461    /// Current lifecycle phase.
462    #[inline]
463    pub fn phase(&self) -> &TaskPhase {
464        &self.status.phase
465    }
466
467    /// Assigns a state-store resource version.
468    ///
469    /// # Errors
470    ///
471    /// Returns [`ModelError::Invalid`] when the value is empty.
472    pub fn set_resource_version(&mut self, resource_version: impl Into<String>) -> ModelResult<()> {
473        self.metadata.set_resource_version(resource_version)
474    }
475
476    /// Applies caller-owned metadata and desired state.
477    ///
478    /// UID and creation time are preserved.
479    /// Metadata-only changes preserve generation and status.
480    /// Spec changes increment generation and reset phase and attempt.
481    /// The previous `observedGeneration` is retained.
482    ///
483    /// Identical desired state returns [`DesiredChange::None`].
484    /// In that case, `resource_version` is not assigned.
485    ///
486    /// # Errors
487    ///
488    /// Returns [`ModelError::Invalid`] when metadata, spec, or a changed `resource_version` is invalid.
489    pub fn apply_desired(
490        &mut self,
491        labels: Labels,
492        annotations: Annotations,
493        spec: TaskSpec,
494        resource_version: impl Into<String>,
495    ) -> ModelResult<DesiredChange> {
496        spec.validate()?;
497        labels.validate()?;
498        annotations.validate()?;
499        let metadata_changed =
500            self.metadata.labels() != &labels || self.metadata.annotations() != &annotations;
501        let spec_changed = self.spec != spec;
502        if !metadata_changed && !spec_changed {
503            return Ok(DesiredChange::None);
504        }
505
506        self.metadata.set_resource_version(resource_version)?;
507        if metadata_changed {
508            self.metadata.apply_metadata(labels, annotations);
509        }
510        if spec_changed {
511            self.spec = spec;
512            self.metadata.bump_generation();
513            self.status = self.status.pending_after(self.metadata.generation());
514        }
515        Ok(if spec_changed {
516            DesiredChange::Spec
517        } else {
518            DesiredChange::Metadata
519        })
520    }
521
522    /// Marks the current generation as reconciled.
523    ///
524    /// Returns `true` when status changed.
525    /// Returns `false` when the same generation was already reconciled.
526    ///
527    /// # Errors
528    ///
529    /// Returns [`ModelError::Invalid`] when a changed `resource_version` is empty.
530    pub fn mark_observed(&mut self, resource_version: impl Into<String>) -> ModelResult<bool> {
531        let generation = self.metadata.generation();
532        let condition = self.status.reconciled_required();
533        let changed = self.status.observed_generation != generation
534            || condition.status() != ConditionStatus::True
535            || condition.observed_generation() != generation
536            || condition.reason() != "RuntimeAccepted"
537            || condition.message() != "runtime accepted the desired state";
538        if !changed {
539            return Ok(false);
540        }
541        self.metadata.set_resource_version(resource_version)?;
542        self.status.mark_reconciled(generation);
543        Ok(true)
544    }
545
546    /// Reschedules reconciliation after a recorded failure.
547    ///
548    /// Returns `false` when reconciliation is not failed.
549    /// A change resets phase, attempt, exit code, and lifecycle error.
550    ///
551    /// # Errors
552    ///
553    /// Returns [`ModelError::Invalid`] when a changed `resource_version` is empty.
554    pub fn mark_reconciliation_pending(
555        &mut self,
556        resource_version: impl Into<String>,
557    ) -> ModelResult<bool> {
558        if !self.status.reconciliation_failed() {
559            return Ok(false);
560        }
561        let generation = self.metadata.generation();
562        self.metadata.set_resource_version(resource_version)?;
563        self.status.mark_reconciliation_pending(generation);
564        self.status.phase = TaskPhase::Pending;
565        self.status.attempt = 0;
566        self.status.exit_code = None;
567        self.status.error = None;
568        Ok(true)
569    }
570
571    /// Records a reconciliation failure.
572    ///
573    /// Desired state is retained.
574    /// Execution phase and diagnostics are reset.
575    ///
576    /// Returns `true` when status changed.
577    ///
578    /// # Errors
579    ///
580    /// Returns [`ModelError::Invalid`] when reason, message, or a changed `resource_version` is invalid.
581    pub fn mark_reconciliation_failed(
582        &mut self,
583        reason: impl Into<String>,
584        message: impl Into<String>,
585        resource_version: impl Into<String>,
586    ) -> ModelResult<bool> {
587        let reason = reason.into();
588        let message = message.into();
589        TaskCondition::validate_reason_message(&reason, &message)?;
590        let generation = self.metadata.generation();
591        let condition = self.status.reconciled_required();
592        let changed = condition.status() != ConditionStatus::False
593            || condition.observed_generation() != generation
594            || condition.reason() != reason
595            || condition.message() != message
596            || self.status.observed_generation != generation
597            || self.status.phase != TaskPhase::Pending
598            || self.status.attempt != 0
599            || self.status.exit_code.is_some()
600            || self.status.error.is_some();
601        if !changed {
602            return Ok(false);
603        }
604        self.metadata.set_resource_version(resource_version)?;
605        self.status
606            .mark_reconciliation_failed(generation, reason, message);
607        self.status.phase = TaskPhase::Pending;
608        self.status.attempt = 0;
609        self.status.exit_code = None;
610        self.status.error = None;
611        Ok(true)
612    }
613
614    /// Records an authoritative attempt start.
615    ///
616    /// A stale generation returns `false`.
617    /// An identical transition also returns `false`.
618    /// Attempt numbers come from the execution source of truth.
619    ///
620    /// # Errors
621    ///
622    /// Returns [`ModelError::Invalid`] when the current generation has attempt zero or a changed `resource_version` is empty.
623    pub fn transition_starting(
624        &mut self,
625        generation: u64,
626        attempt: u32,
627        resource_version: impl Into<String>,
628    ) -> ModelResult<bool> {
629        if generation != self.metadata.generation() {
630            return Ok(false);
631        }
632        if attempt == 0 {
633            return Err(ModelError::Invalid(
634                "attempt must be greater than zero".into(),
635            ));
636        }
637        let changed = self.status.observed_generation != generation
638            || self.status.reconciled_required().status() != ConditionStatus::True
639            || self.status.reconciled_required().observed_generation() != generation
640            || self.status.phase != TaskPhase::Running
641            || self.status.attempt != attempt
642            || self.status.exit_code.is_some()
643            || self.status.error.is_some();
644        if !changed {
645            return Ok(false);
646        }
647        self.metadata.set_resource_version(resource_version)?;
648        self.status.mark_reconciled(generation);
649        self.status.phase = TaskPhase::Running;
650        self.status.attempt = attempt;
651        self.status.exit_code = None;
652        self.status.error = None;
653        Ok(true)
654    }
655
656    /// Records a terminal attempt phase.
657    ///
658    /// A stale generation or older attempt returns `false`.
659    /// Terminal phases are sticky.
660    /// `Failed` may be refined to `Exhausted` or `Timeout`.
661    ///
662    /// # Errors
663    ///
664    /// Returns [`ModelError::Invalid`] when the current generation has attempt zero, `phase` is not terminal, or a changed `resource_version` is empty.
665    pub fn transition_finished(
666        &mut self,
667        generation: u64,
668        attempt: u32,
669        phase: TaskPhase,
670        error: Option<String>,
671        exit_code: Option<i32>,
672        resource_version: impl Into<String>,
673    ) -> ModelResult<bool> {
674        if generation != self.metadata.generation() {
675            return Ok(false);
676        }
677        if attempt == 0 {
678            return Err(ModelError::Invalid(
679                "attempt must be greater than zero".into(),
680            ));
681        }
682        if !phase.is_terminal() {
683            return Err(ModelError::Invalid(
684                format!("transition_finished requires a terminal phase, got {phase}").into(),
685            ));
686        }
687        if attempt < self.status.attempt {
688            return Ok(false);
689        }
690        let same_attempt = attempt == self.status.attempt;
691        let reconciled = self.status.reconciled_required().status() == ConditionStatus::True
692            && self.status.reconciled_required().observed_generation() == generation;
693        if same_attempt && self.status.phase == phase && reconciled {
694            return Ok(false);
695        }
696        if same_attempt && self.status.phase.is_terminal() && reconciled {
697            let refines_failed = self.status.phase == TaskPhase::Failed
698                && matches!(phase, TaskPhase::Exhausted | TaskPhase::Timeout);
699            if !refines_failed {
700                return Ok(false);
701            }
702        }
703        self.metadata.set_resource_version(resource_version)?;
704        self.status.attempt = attempt;
705        self.set_terminal(generation, phase, error, exit_code);
706        Ok(true)
707    }
708
709    /// Records an authoritative final lifecycle outcome.
710    ///
711    /// Unlike [`Self::transition_finished`], this may replace a conflicting terminal attempt phase.
712    /// A stale generation or identical outcome returns `false`.
713    ///
714    /// # Errors
715    ///
716    /// Returns [`ModelError::Invalid`] when `phase` is not terminal for the current generation or a changed `resource_version` is empty.
717    pub fn reconcile_finished(
718        &mut self,
719        generation: u64,
720        phase: TaskPhase,
721        error: Option<String>,
722        exit_code: Option<i32>,
723        resource_version: impl Into<String>,
724    ) -> ModelResult<bool> {
725        if generation != self.metadata.generation() {
726            return Ok(false);
727        }
728        if !phase.is_terminal() {
729            return Err(ModelError::Invalid(
730                format!("reconcile_finished requires a terminal phase, got {phase}").into(),
731            ));
732        }
733        let changed = self.status.observed_generation != generation
734            || self.status.reconciled_required().status() != ConditionStatus::True
735            || self.status.reconciled_required().observed_generation() != generation
736            || self.status.phase != phase
737            || self.status.error != error
738            || self.status.exit_code != exit_code;
739        if !changed {
740            return Ok(false);
741        }
742        self.metadata.set_resource_version(resource_version)?;
743        self.set_terminal(generation, phase, error, exit_code);
744        Ok(true)
745    }
746
747    fn set_terminal(
748        &mut self,
749        generation: u64,
750        phase: TaskPhase,
751        error: Option<String>,
752        exit_code: Option<i32>,
753    ) {
754        self.status.mark_reconciled(generation);
755        self.status.phase = phase;
756        self.status.error = error;
757        self.status.exit_code = exit_code;
758    }
759}
760
761impl From<&Task> for TaskManifest {
762    fn from(task: &Task) -> Self {
763        Self {
764            type_meta: task.type_meta.clone(),
765            metadata: TaskManifestMeta {
766                name: task.name().clone(),
767                labels: task.metadata.labels().clone(),
768                annotations: task.metadata.annotations().clone(),
769            },
770            spec: task.spec.clone(),
771        }
772    }
773}
774
775impl From<Task> for TaskManifest {
776    fn from(task: Task) -> Self {
777        Self::from(&task)
778    }
779}
780
781mod raw {
782    use super::*;
783
784    #[derive(Deserialize)]
785    #[serde(rename_all = "camelCase", deny_unknown_fields)]
786    pub(super) struct TaskRaw {
787        api_version: String,
788        kind: String,
789        metadata: ObjectMeta,
790        spec: TaskSpec,
791        status: TaskStatus,
792    }
793
794    #[derive(Deserialize)]
795    #[serde(rename_all = "camelCase", deny_unknown_fields)]
796    pub(super) struct TaskManifestRaw {
797        api_version: String,
798        kind: String,
799        metadata: TaskManifestMeta,
800        spec: TaskSpec,
801    }
802
803    impl TryFrom<TaskRaw> for Task {
804        type Error = ModelError;
805
806        fn try_from(raw: TaskRaw) -> Result<Self, Self::Error> {
807            Task::from_parts(
808                TypeMeta {
809                    api_version: raw.api_version,
810                    kind: raw.kind,
811                },
812                raw.metadata,
813                raw.spec,
814                raw.status,
815            )
816        }
817    }
818
819    impl TryFrom<TaskManifestRaw> for TaskManifest {
820        type Error = ModelError;
821
822        fn try_from(raw: TaskManifestRaw) -> Result<Self, Self::Error> {
823            TaskManifest::from_parts(
824                TypeMeta {
825                    api_version: raw.api_version,
826                    kind: raw.kind,
827                },
828                raw.metadata,
829                raw.spec,
830            )
831        }
832    }
833}
834
835#[cfg(test)]
836mod tests {
837    use super::*;
838    use crate::{EmbeddedSpec, TaskWorkload};
839
840    fn spec(slot: &str) -> TaskSpec {
841        TaskSpec::builder(
842            slot,
843            TaskWorkload::Embedded(EmbeddedSpec::new("test-v1").unwrap()),
844            5_000_u64,
845        )
846        .build()
847        .unwrap()
848    }
849
850    fn task() -> Task {
851        let mut task = Task::new("task-a", spec("slot-a")).unwrap();
852        task.set_resource_version("1").unwrap();
853        task
854    }
855
856    #[test]
857    fn new_creates_valid_unobserved_resource() {
858        let task = Task::new("task-a", spec("slot-a")).unwrap();
859
860        assert_eq!(task.type_meta().api_version(), TASK_API_VERSION);
861        assert_eq!(task.type_meta().kind(), TASK_KIND);
862        assert_eq!(task.name(), "task-a");
863        assert!(!task.uid().as_str().is_empty());
864        assert_eq!(task.metadata().generation(), 1);
865        assert_eq!(task.status().observed_generation(), 0);
866        assert_eq!(*task.phase(), TaskPhase::Pending);
867        assert_eq!(
868            task.status().reconciled().status(),
869            ConditionStatus::Unknown
870        );
871        assert_eq!(task.status().reconciled().observed_generation(), 1);
872    }
873
874    #[test]
875    fn serde_shape_and_roundtrip_are_crd_shaped() {
876        let task = task();
877        let json = serde_json::to_value(&task).unwrap();
878
879        assert_eq!(json["apiVersion"], TASK_API_VERSION);
880        assert_eq!(json["kind"], TASK_KIND);
881        assert!(json.get("metadata").is_some());
882        assert!(json.get("spec").is_some());
883        assert!(json.get("status").is_some());
884
885        let back: Task = serde_json::from_value(json).unwrap();
886        assert_eq!(back, task);
887    }
888
889    #[test]
890    fn manifest_serde_contains_only_user_owned_resource_fields() {
891        let mut labels = Labels::new();
892        labels.insert("tier", "worker");
893        let manifest = TaskManifest::new("task-a", spec("slot-a"))
894            .unwrap()
895            .with_labels(labels)
896            .unwrap();
897
898        let json = serde_json::to_value(&manifest).unwrap();
899        assert_eq!(json["apiVersion"], TASK_API_VERSION);
900        assert_eq!(json["kind"], TASK_KIND);
901        assert_eq!(json["metadata"]["name"], "task-a");
902        assert!(json.get("spec").is_some());
903        assert!(json.get("status").is_none());
904        for server_owned in ["uid", "resourceVersion", "generation", "creationTimestamp"] {
905            assert!(json["metadata"].get(server_owned).is_none());
906        }
907
908        let back: TaskManifest = serde_json::from_value(json).unwrap();
909        assert_eq!(back, manifest);
910    }
911
912    #[test]
913    fn stored_task_roundtrips_through_its_desired_manifest() {
914        let stored = task();
915        let manifest = TaskManifest::from(&stored);
916        let rematerialized = Task::from_manifest(manifest).unwrap();
917
918        assert_eq!(rematerialized.name(), stored.name());
919        assert_eq!(rematerialized.spec(), stored.spec());
920        assert_eq!(
921            rematerialized.metadata().labels(),
922            stored.metadata().labels()
923        );
924        assert_ne!(rematerialized.uid(), stored.uid());
925        assert_eq!(rematerialized.status().phase(), TaskPhase::Pending);
926    }
927
928    #[test]
929    fn manifest_deserialization_rejects_wrong_resource_gvk() {
930        let manifest = TaskManifest::new("task-a", spec("slot-a")).unwrap();
931        let mut json = serde_json::to_value(manifest).unwrap();
932        json["kind"] = serde_json::json!("Other");
933
934        let error = serde_json::from_value::<TaskManifest>(json).unwrap_err();
935        assert!(error.to_string().contains("Task kind"));
936    }
937
938    #[test]
939    fn manifest_deserialization_rejects_invalid_metadata() {
940        let manifest = TaskManifest::new("task-a", spec("slot-a")).unwrap();
941        let mut json = serde_json::to_value(manifest).unwrap();
942        json["metadata"]["labels"] = serde_json::json!({ "example.io/bad key": "value" });
943
944        let error = serde_json::from_value::<TaskManifest>(json).unwrap_err();
945        assert!(error.to_string().contains("label key"));
946    }
947
948    #[test]
949    fn manifest_rejects_unknown_fields_at_resource_and_metadata_levels() {
950        let manifest = TaskManifest::new("task-a", spec("slot-a")).unwrap();
951        let mut resource = serde_json::to_value(&manifest).unwrap();
952        resource["unexpected"] = serde_json::json!(true);
953        assert!(serde_json::from_value::<TaskManifest>(resource).is_err());
954
955        let mut metadata = serde_json::to_value(manifest).unwrap();
956        metadata["metadata"]["unexpected"] = serde_json::json!(true);
957        assert!(serde_json::from_value::<TaskManifest>(metadata).is_err());
958    }
959
960    #[test]
961    fn stored_task_rejects_unknown_status_and_server_metadata_fields() {
962        let stored = task();
963        let mut metadata = serde_json::to_value(&stored).unwrap();
964        metadata["metadata"]["unexpected"] = serde_json::json!(true);
965        assert!(serde_json::from_value::<Task>(metadata).is_err());
966
967        let mut status = serde_json::to_value(stored).unwrap();
968        status["status"]["unexpected"] = serde_json::json!(true);
969        assert!(serde_json::from_value::<Task>(status).is_err());
970    }
971
972    #[test]
973    fn serde_rejects_wrong_resource_gvk() {
974        let mut json = serde_json::to_value(task()).unwrap();
975        json["kind"] = serde_json::json!("Other");
976
977        let error = serde_json::from_value::<Task>(json).unwrap_err();
978        assert!(error.to_string().contains("Task kind"));
979    }
980
981    #[test]
982    fn metadata_only_apply_preserves_generation_and_status() {
983        let mut task = task();
984        let uid = task.uid().clone();
985        let creation = task.metadata().creation_timestamp();
986        let mut labels = Labels::new();
987        labels.insert("tier", "prod");
988
989        let changed = task
990            .apply_desired(labels, Annotations::new(), spec("slot-a"), "2")
991            .unwrap();
992
993        assert_eq!(changed, DesiredChange::Metadata);
994        assert_eq!(task.metadata().generation(), 1);
995        assert_eq!(task.uid(), &uid);
996        assert_eq!(task.metadata().creation_timestamp(), creation);
997        assert_eq!(task.metadata().resource_version(), "2");
998        assert_eq!(*task.phase(), TaskPhase::Pending);
999    }
1000
1001    #[test]
1002    fn spec_apply_increments_generation_and_resets_execution_state() {
1003        let mut task = task();
1004        task.transition_starting(1, 4, "2").unwrap();
1005
1006        let change = task
1007            .apply_desired(Labels::new(), Annotations::new(), spec("slot-b"), "3")
1008            .unwrap();
1009
1010        assert_eq!(change, DesiredChange::Spec);
1011        assert_eq!(task.metadata().generation(), 2);
1012        assert_eq!(task.status().observed_generation(), 1);
1013        assert_eq!(task.status().attempt(), 0);
1014        assert_eq!(*task.phase(), TaskPhase::Pending);
1015        assert_eq!(
1016            task.status().reconciled().status(),
1017            ConditionStatus::Unknown
1018        );
1019        assert_eq!(task.status().reconciled().observed_generation(), 2);
1020        assert_eq!(task.slot(), "slot-b");
1021    }
1022
1023    #[test]
1024    fn embedded_revision_change_is_a_spec_change() {
1025        let mut task = task();
1026        let transition_time = task.status().reconciled().last_transition_time();
1027        let changed_spec = TaskSpec::builder(
1028            "slot-a",
1029            TaskWorkload::Embedded(EmbeddedSpec::new("test-v2").unwrap()),
1030            5_000_u64,
1031        )
1032        .build()
1033        .unwrap();
1034
1035        let change = task
1036            .apply_desired(Labels::new(), Annotations::new(), changed_spec, "2")
1037            .unwrap();
1038
1039        assert_eq!(change, DesiredChange::Spec);
1040        assert_eq!(task.metadata().generation(), 2);
1041        let TaskWorkload::Embedded(embedded) = task.spec().workload() else {
1042            panic!("workload must remain Embedded");
1043        };
1044        assert_eq!(embedded.revision(), "test-v2");
1045        assert_eq!(
1046            task.status().reconciled().last_transition_time(),
1047            transition_time,
1048            "lastTransitionTime changes only when condition status changes"
1049        );
1050    }
1051
1052    #[test]
1053    fn starting_uses_authoritative_attempt_and_observes_generation() {
1054        let mut task = task();
1055
1056        assert!(task.transition_starting(1, 7, "2").unwrap());
1057        assert_eq!(task.status().attempt(), 7);
1058        assert_eq!(task.status().observed_generation(), 1);
1059        assert_eq!(*task.phase(), TaskPhase::Running);
1060    }
1061
1062    #[test]
1063    fn terminal_event_sets_attempt_when_start_event_was_not_observed() {
1064        let mut task = task();
1065
1066        task.transition_finished(1, 3, TaskPhase::Succeeded, None, Some(0), "2")
1067            .unwrap();
1068
1069        assert_eq!(task.status().attempt(), 3);
1070        assert_eq!(task.status().observed_generation(), 1);
1071        assert_eq!(*task.phase(), TaskPhase::Succeeded);
1072    }
1073
1074    #[test]
1075    fn no_op_apply_does_not_consume_resource_version() {
1076        let mut task = task();
1077
1078        let change = task
1079            .apply_desired(Labels::new(), Annotations::new(), spec("slot-a"), "2")
1080            .unwrap();
1081
1082        assert_eq!(change, DesiredChange::None);
1083        assert_eq!(task.metadata().resource_version(), "1");
1084    }
1085
1086    #[test]
1087    fn invalid_metadata_apply_is_rejected_without_mutation() {
1088        let mut task = task();
1089        let before = task.clone();
1090        let mut labels = Labels::new();
1091        labels.insert("bad key", "value");
1092
1093        assert!(
1094            task.apply_desired(labels, Annotations::new(), spec("slot-b"), "2")
1095                .is_err()
1096        );
1097        assert_eq!(task, before);
1098    }
1099
1100    #[test]
1101    fn stale_generation_status_is_ignored_without_bumping_version() {
1102        let mut task = task();
1103
1104        assert!(!task.transition_starting(2, 1, "2").unwrap());
1105        assert_eq!(task.metadata().resource_version(), "1");
1106        assert_eq!(*task.phase(), TaskPhase::Pending);
1107    }
1108
1109    #[test]
1110    fn reconciliation_failure_observes_generation_and_retains_spec() {
1111        let mut task = task();
1112
1113        assert!(
1114            task.mark_reconciliation_failed("RunnerBuildFailed", "no runner", "2")
1115                .unwrap()
1116        );
1117        assert_eq!(task.status().observed_generation(), 1);
1118        assert_eq!(*task.phase(), TaskPhase::Pending);
1119        assert_eq!(task.status().attempt(), 0);
1120        assert!(task.status().error().is_none());
1121        assert_eq!(task.status().reconciled().status(), ConditionStatus::False);
1122        assert_eq!(task.status().reconciled().reason(), "RunnerBuildFailed");
1123        assert_eq!(task.status().reconciled().message(), "no runner");
1124        assert_eq!(task.slot(), "slot-a");
1125    }
1126
1127    #[test]
1128    fn failed_reconciliation_can_be_rescheduled_without_changing_generation() {
1129        let mut task = task();
1130        task.mark_reconciliation_failed("RunnerBuildFailed", "no runner", "2")
1131            .unwrap();
1132
1133        assert!(task.mark_reconciliation_pending("3").unwrap());
1134        assert_eq!(task.metadata().generation(), 1);
1135        assert_eq!(task.metadata().resource_version(), "3");
1136        assert_eq!(
1137            task.status().reconciled().status(),
1138            ConditionStatus::Unknown
1139        );
1140        assert_eq!(task.status().reconciled().observed_generation(), 1);
1141    }
1142
1143    #[test]
1144    fn sticky_terminal_can_only_refine_failed() {
1145        let mut task = task();
1146        task.transition_starting(1, 1, "2").unwrap();
1147        task.transition_finished(1, 1, TaskPhase::Failed, Some("attempt".into()), None, "3")
1148            .unwrap();
1149
1150        assert!(
1151            !task
1152                .transition_finished(1, 1, TaskPhase::Succeeded, None, Some(0), "4")
1153                .unwrap()
1154        );
1155        assert!(
1156            task.transition_finished(1, 1, TaskPhase::Exhausted, Some("budget".into()), None, "4",)
1157                .unwrap()
1158        );
1159        assert_eq!(*task.phase(), TaskPhase::Exhausted);
1160    }
1161}