1use 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
42pub const TASK_API_VERSION_MAJOR: u32 = task_api_major!();
44
45pub const TASK_API_VERSION: &str = concat!("solti.io/v", task_api_major!());
47
48pub const TASK_KIND: &str = "Task";
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub enum DesiredChange {
54 None,
56 Metadata,
58 Spec,
60}
61
62impl DesiredChange {
63 #[inline]
65 pub fn is_changed(self) -> bool {
66 !matches!(self, Self::None)
67 }
68}
69
70#[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#[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 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 #[inline]
116 pub fn name(&self) -> &TaskId {
117 &self.name
118 }
119
120 #[inline]
122 pub fn labels(&self) -> &Labels {
123 &self.labels
124 }
125
126 #[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#[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 pub fn new(name: impl AsRef<str>, spec: TaskSpec) -> ModelResult<Self> {
172 Self::from_parts(TypeMeta::task(), TaskManifestMeta::new(name)?, spec)
173 }
174
175 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 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 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 pub fn validate(&self) -> ModelResult<()> {
222 self.type_meta.validate_task()?;
223 self.metadata.validate()?;
224 self.spec.validate()
225 }
226
227 #[inline]
229 pub fn type_meta(&self) -> &TypeMeta {
230 &self.type_meta
231 }
232
233 #[inline]
235 pub fn metadata(&self) -> &TaskManifestMeta {
236 &self.metadata
237 }
238
239 #[inline]
241 pub fn name(&self) -> &TaskId {
242 self.metadata.name()
243 }
244
245 #[inline]
247 pub fn spec(&self) -> &TaskSpec {
248 &self.spec
249 }
250
251 #[inline]
253 pub fn slot(&self) -> &Slot {
254 self.spec.slot()
255 }
256
257 pub fn into_parts(self) -> (TypeMeta, TaskManifestMeta, TaskSpec) {
259 (self.type_meta, self.metadata, self.spec)
260 }
261}
262
263impl TypeMeta {
264 pub fn task() -> Self {
266 Self {
267 api_version: TASK_API_VERSION.to_owned(),
268 kind: TASK_KIND.to_owned(),
269 }
270 }
271
272 #[inline]
274 pub fn api_version(&self) -> &str {
275 &self.api_version
276 }
277
278 #[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#[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 pub fn new(name: impl AsRef<str>, spec: TaskSpec) -> ModelResult<Self> {
330 Self::from_manifest(TaskManifest::new(name, spec)?)
331 }
332
333 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 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 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 #[inline]
410 pub fn type_meta(&self) -> &TypeMeta {
411 &self.type_meta
412 }
413
414 #[inline]
416 pub fn metadata(&self) -> &ObjectMeta {
417 &self.metadata
418 }
419
420 #[inline]
422 pub fn spec(&self) -> &TaskSpec {
423 &self.spec
424 }
425
426 #[inline]
428 pub fn status(&self) -> &TaskStatus {
429 &self.status
430 }
431
432 pub fn into_parts(self) -> (TypeMeta, ObjectMeta, TaskSpec, TaskStatus) {
434 (self.type_meta, self.metadata, self.spec, self.status)
435 }
436
437 #[inline]
439 pub fn name(&self) -> &TaskId {
440 self.metadata.name()
441 }
442
443 #[inline]
445 pub fn uid(&self) -> &Uid {
446 self.metadata.uid()
447 }
448
449 #[inline]
451 pub fn slot(&self) -> &Slot {
452 self.spec.slot()
453 }
454
455 #[inline]
457 pub fn labels(&self) -> &Labels {
458 self.metadata.labels()
459 }
460
461 #[inline]
463 pub fn phase(&self) -> &TaskPhase {
464 &self.status.phase
465 }
466
467 pub fn set_resource_version(&mut self, resource_version: impl Into<String>) -> ModelResult<()> {
473 self.metadata.set_resource_version(resource_version)
474 }
475
476 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 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 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 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 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 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 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}