use std::collections::BTreeMap;
use kube::CustomResource;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
#[derive(CustomResource, Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
#[kube(
group = "polychrome.dev",
version = "v1alpha1",
kind = "Workflow",
namespaced,
status = "WorkflowStatus",
shortname = "wf",
category = "polychrome",
derive = "PartialEq",
printcolumn = r#"{"name":"Ready","type":"boolean","jsonPath":".status.ready"}"#,
printcolumn = r#"{"name":"Phase","type":"string","jsonPath":".status.phase"}"#,
printcolumn = r#"{"name":"Stages","type":"integer","jsonPath":".status.stageCount"}"#,
printcolumn = r#"{"name":"Age","type":"date","jsonPath":".metadata.creationTimestamp"}"#
)]
#[serde(rename_all = "camelCase")]
pub struct WorkflowSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub owner: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub stages: Vec<WorkflowStage>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct WorkflowStage {
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub kind: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub template: Option<String>,
pub image: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub port: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub replicas: Option<i32>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub env: BTreeMap<String, String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, Default, PartialEq, Eq, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct WorkflowStatus {
#[serde(default)]
pub ready: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub phase: Option<String>,
#[serde(default)]
pub stage_count: i32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub message: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub stages: Vec<StageStatus>,
}
#[derive(Serialize, Deserialize, Clone, Debug, Default, PartialEq, Eq, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct StageStatus {
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub kind: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub template: Option<String>,
pub phase: String,
#[serde(default)]
pub ready: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub service_definition: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use super::*;
use kube::CustomResourceExt;
use serde_json::{Value, json};
#[test]
fn crd_identity_is_polychrome_workflow() {
let crd = Workflow::crd();
assert_eq!(crd.spec.group, "polychrome.dev");
assert_eq!(crd.spec.names.kind, "Workflow");
assert_eq!(crd.spec.names.plural, "workflows");
}
#[test]
fn stage_round_trips_camel_case() {
let stage = WorkflowStage {
name: "ingest".to_owned(),
kind: Some("dataset".to_owned()),
depends_on: vec!["source".to_owned()],
description: Some("ingest the source feed".to_owned()),
template: Some("rust-dataset".to_owned()),
image: "example.test/dataset:1".to_owned(),
port: Some(8080),
replicas: Some(1),
env: BTreeMap::from([("FEED".to_owned(), "main".to_owned())]),
};
let v = serde_json::to_value(&stage).unwrap();
assert!(v.get("dependsOn").is_some(), "depends_on → dependsOn");
assert_eq!(v["dependsOn"][0], "source");
assert_eq!(serde_json::from_value::<WorkflowStage>(v).unwrap(), stage);
}
#[test]
fn minimal_stage_needs_only_name_and_image() {
let stage: WorkflowStage =
serde_json::from_value(json!({ "name": "api", "image": "img:1" })).unwrap();
assert!(stage.depends_on.is_empty());
assert!(stage.kind.is_none());
let out = serde_json::to_value(&stage).unwrap();
assert!(out.get("dependsOn").is_none());
assert!(out.get("env").is_none());
}
#[test]
fn status_stages_round_trip() {
let status = WorkflowStatus {
ready: false,
phase: Some("Progressing".to_owned()),
stage_count: 2,
message: None,
stages: vec![
StageStatus {
name: "data".to_owned(),
kind: Some("dataset".to_owned()),
template: Some("rust-dataset".to_owned()),
phase: "Ready".to_owned(),
ready: true,
service_definition: Some("pipe-data".to_owned()),
depends_on: vec![],
},
StageStatus {
name: "api".to_owned(),
kind: Some("service".to_owned()),
template: None,
phase: "Pending".to_owned(),
ready: false,
service_definition: Some("pipe-api".to_owned()),
depends_on: vec!["data".to_owned()],
},
],
};
let v = serde_json::to_value(&status).unwrap();
assert_eq!(v["stageCount"], 2);
assert_eq!(v["stages"][1]["serviceDefinition"], "pipe-api");
assert_eq!(serde_json::from_value::<WorkflowStatus>(v).unwrap(), status);
}
#[test]
fn minimal_workflow_spec_is_valid() {
let spec: WorkflowSpec = serde_json::from_value(json!({
"stages": [{ "name": "only", "image": "img:1" }]
}))
.unwrap();
assert_eq!(spec.stages.len(), 1);
let _: Value = serde_json::to_value(spec).unwrap();
}
}