use std::{
collections::{BTreeMap, BTreeSet},
sync::Arc,
time::Duration,
};
use futures::{StreamExt, future::try_join_all};
use kube::{
Api, Client, Resource, ResourceExt,
api::{Patch, PatchParams},
runtime::{
controller::{Action, Controller},
watcher,
},
};
use serde_json::json;
use crate::{
fanout::{self, MANAGED_BY_KEY, MANAGED_BY_VALUE},
servicedefinition::{ServiceDefinition, ServiceDefinitionSpec, default_port, default_replicas},
workflow::{StageStatus, Workflow, WorkflowSpec, WorkflowStage, WorkflowStatus},
};
const MAX_LABEL_LEN: usize = 63;
const COLLISION_RETRY_INTERVAL: Duration = Duration::from_secs(30);
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("kube api: {0}")]
Kube(#[from] kube::Error),
#[error("workflow has no namespace")]
NoNamespace,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum WorkflowAction {
Apply(Vec<usize>),
Invalid(String),
Noop,
}
#[must_use]
pub fn plan(wf: &Workflow) -> WorkflowAction {
if wf.meta().deletion_timestamp.is_some() {
return WorkflowAction::Noop;
}
match topo_sort(&wf.spec) {
Err(reason) => WorkflowAction::Invalid(reason),
Ok(order) => match check_stage_names(&wf.name_any(), &wf.spec) {
Err(reason) => WorkflowAction::Invalid(reason),
Ok(()) => WorkflowAction::Apply(order),
},
}
}
pub fn topo_sort(spec: &WorkflowSpec) -> Result<Vec<usize>, String> {
let stages = &spec.stages;
if stages.is_empty() {
return Err("workflow has no stages".to_owned());
}
let mut index: BTreeMap<&str, usize> = BTreeMap::new();
for (i, s) in stages.iter().enumerate() {
if index.insert(s.name.as_str(), i).is_some() {
return Err(format!("duplicate stage name `{}`", s.name));
}
}
let mut in_degree = vec![0usize; stages.len()];
let mut dependents: Vec<Vec<usize>> = vec![Vec::new(); stages.len()];
for (i, s) in stages.iter().enumerate() {
for dep in &s.depends_on {
if dep == &s.name {
return Err(format!("stage `{}` depends on itself", s.name));
}
let Some(&j) = index.get(dep.as_str()) else {
return Err(format!(
"stage `{}` depends on `{dep}`, which is not a stage in this workflow",
s.name
));
};
dependents[j].push(i);
in_degree[i] += 1;
}
}
let mut queue: Vec<usize> = (0..stages.len()).filter(|&i| in_degree[i] == 0).collect();
let mut order = Vec::with_capacity(stages.len());
let mut head = 0;
while head < queue.len() {
let i = queue[head];
head += 1;
order.push(i);
for &k in &dependents[i] {
in_degree[k] -= 1;
if in_degree[k] == 0 {
queue.push(k);
}
}
}
if order.len() != stages.len() {
let mut stuck: Vec<&str> = stages
.iter()
.enumerate()
.filter(|(i, _)| !order.contains(i))
.map(|(_, s)| s.name.as_str())
.collect();
stuck.sort_unstable();
return Err(format!(
"dependsOn graph has a cycle through: {}",
stuck.join(", ")
));
}
Ok(order)
}
pub fn validate(wf: &Workflow) -> Result<(), String> {
topo_sort(&wf.spec)?;
check_stage_names(&wf.name_any(), &wf.spec)
}
fn check_stage_names(workflow: &str, spec: &WorkflowSpec) -> Result<(), String> {
if !is_label_fragment(workflow) {
return Err(format!(
"workflow name `{workflow}` must be a DNS-1123 label (lowercase alphanumeric and \
'-', starting with a letter)"
));
}
for s in &spec.stages {
if !is_label_fragment(&s.name) {
return Err(format!(
"stage name `{}` must be a DNS-1123 label (lowercase alphanumeric and '-', \
starting with a letter)",
s.name
));
}
let combined = stage_sd_name(workflow, &s.name);
if combined.len() > MAX_LABEL_LEN {
return Err(format!(
"stage `{}` yields name `{combined}` ({} chars), over the {MAX_LABEL_LEN}-char limit",
s.name,
combined.len()
));
}
}
Ok(())
}
fn is_label_fragment(name: &str) -> bool {
!name.is_empty()
&& name.as_bytes()[0].is_ascii_lowercase()
&& name.ends_with(|c: char| c.is_ascii_lowercase() || c.is_ascii_digit())
&& name
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-')
}
#[must_use]
pub fn stage_sd_name(workflow: &str, stage: &str) -> String {
format!("{workflow}-{stage}")
}
#[must_use]
pub fn unblocked(spec: &WorkflowSpec, ready: &BTreeSet<String>) -> BTreeSet<String> {
spec.stages
.iter()
.filter(|s| s.depends_on.iter().all(|d| ready.contains(d)))
.map(|s| s.name.clone())
.collect()
}
fn sd_ready(svc: Option<&ServiceDefinition>) -> bool {
svc.and_then(|s| s.status.as_ref())
.is_some_and(|st| st.ready)
}
fn is_owned_by(sd: &ServiceDefinition, wf: &Workflow) -> bool {
let wf_uid = wf.uid();
let wf_name = wf.name_any();
sd.metadata.owner_references.iter().flatten().any(|r| {
r.controller == Some(true)
&& r.kind == "Workflow"
&& wf_uid
.as_ref()
.map_or_else(|| r.name == wf_name, |uid| &r.uid == uid)
})
}
pub struct Context {
pub client: Client,
}
#[allow(clippy::too_many_lines)]
#[tracing::instrument(skip_all, fields(workflow = %wf.name_any()))]
pub async fn reconcile(wf: Arc<Workflow>, ctx: Arc<Context>) -> Result<Action, Error> {
let ns = wf.namespace().ok_or(Error::NoNamespace)?;
let name = wf.name_any();
let svcdefs: Api<ServiceDefinition> = Api::namespaced(ctx.client.clone(), &ns);
let order = match plan(&wf) {
WorkflowAction::Noop => return Ok(Action::await_change()),
WorkflowAction::Invalid(reason) => {
sync_degraded(&ctx.client, &ns, &wf, &reason).await?;
return Ok(Action::await_change());
}
WorkflowAction::Apply(order) => order,
};
let sd_names: Vec<String> = wf
.spec
.stages
.iter()
.map(|s| stage_sd_name(&name, &s.name))
.collect();
let fetched = try_join_all(sd_names.iter().map(|n| svcdefs.get_opt(n))).await?;
let by_name: BTreeMap<String, ServiceDefinition> = sd_names
.iter()
.cloned()
.zip(fetched)
.filter_map(|(n, sd)| sd.map(|sd| (n, sd)))
.collect();
let observed = |s: &WorkflowStage| by_name.get(&stage_sd_name(&name, &s.name));
for s in &wf.spec.stages {
if let Some(sd) = observed(s)
&& !is_owned_by(sd, &wf)
{
let reason = format!(
"stage `{}` collides with existing ServiceDefinition `{}` not owned by this \
workflow; refusing to adopt it",
s.name,
stage_sd_name(&name, &s.name)
);
sync_degraded(&ctx.client, &ns, &wf, &reason).await?;
return Ok(Action::requeue(COLLISION_RETRY_INTERVAL));
}
}
let ready: BTreeSet<String> = wf
.spec
.stages
.iter()
.filter(|s| sd_ready(observed(s)))
.map(|s| s.name.clone())
.collect();
let unblocked = unblocked(&wf.spec, &ready);
let pp = PatchParams::apply("polychrome.dev/workflow");
let mut apply_error: Option<String> = None;
let mut materialized: BTreeSet<String> = wf
.spec
.stages
.iter()
.filter(|s| observed(s).is_some())
.map(|s| s.name.clone())
.collect();
for s in &wf.spec.stages {
if !materialized.contains(&s.name) && !unblocked.contains(&s.name) {
continue;
}
let desired = build_stage_service_definition(&wf, s, &ns);
if observed(s).is_none_or(|c| c.spec != desired.spec) {
match svcdefs
.patch(&stage_sd_name(&name, &s.name), &pp, &Patch::Apply(&desired))
.await
{
Ok(_) => {
materialized.insert(s.name.clone());
tracing::info!(stage = %s.name, "applied workflow stage");
}
Err(e) => {
apply_error = Some(format!("stage `{}`: {e}", s.name));
break;
}
}
}
}
let mut stage_statuses = Vec::with_capacity(order.len());
for &i in &order {
let s = &wf.spec.stages[i];
let is_ready = ready.contains(&s.name);
let has_sd = materialized.contains(&s.name);
let sd_name = stage_sd_name(&name, &s.name);
let (phase, svc_def) = if is_ready {
("Ready", Some(sd_name))
} else if has_sd {
("Pending", Some(sd_name))
} else {
("Blocked", None)
};
stage_statuses.push(StageStatus {
name: s.name.clone(),
kind: s.kind.clone(),
template: s.template.clone(),
phase: phase.to_owned(),
ready: is_ready,
service_definition: svc_def,
depends_on: s.depends_on.clone(),
});
}
let all_ready = !stage_statuses.is_empty() && stage_statuses.iter().all(|s| s.ready);
let any_ready = stage_statuses.iter().any(|s| s.ready);
let reconciled = all_ready && apply_error.is_none();
let phase = if reconciled {
"Ready"
} else if apply_error.is_some() || any_ready {
"Progressing"
} else {
"Pending"
};
let ready_count = stage_statuses.iter().filter(|s| s.ready).count();
let message = (!reconciled).then(|| {
apply_error
.clone()
.unwrap_or_else(|| format!("{ready_count}/{} stages ready", stage_statuses.len()))
});
sync_status(
&ctx.client,
&ns,
&name,
&wf,
reconciled,
phase,
message.as_deref(),
&stage_statuses,
)
.await?;
if reconciled {
Ok(Action::requeue(Duration::from_mins(5)))
} else {
Ok(Action::requeue(Duration::from_secs(15)))
}
}
async fn sync_degraded(
client: &Client,
ns: &str,
wf: &Workflow,
reason: &str,
) -> Result<(), Error> {
sync_status(
client,
ns,
&wf.name_any(),
wf,
false,
"Degraded",
Some(reason),
&[],
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn sync_status(
client: &Client,
ns: &str,
name: &str,
wf: &Workflow,
ready: bool,
phase: &str,
message: Option<&str>,
stages: &[StageStatus],
) -> Result<(), Error> {
let stage_count = i32::try_from(wf.spec.stages.len()).unwrap_or(i32::MAX);
let default = WorkflowStatus::default();
let current = wf.status.as_ref().unwrap_or(&default);
let unchanged = current.ready == ready
&& current.phase.as_deref() == Some(phase)
&& current.message.as_deref() == message
&& current.stages == stages
&& current.stage_count == stage_count;
if unchanged {
return Ok(());
}
let workflows: Api<Workflow> = Api::namespaced(client.clone(), ns);
let status = json!({ "status": {
"ready": ready,
"phase": phase,
"stageCount": stage_count,
"message": message,
"stages": stages,
} });
workflows
.patch_status(name, &PatchParams::default(), &Patch::Merge(&status))
.await?;
tracing::info!(%phase, ready, "synced workflow status");
Ok(())
}
#[must_use]
pub fn build_stage_service_definition(
wf: &Workflow,
stage: &WorkflowStage,
ns: &str,
) -> ServiceDefinition {
let workflow_name = wf.name_any();
let sd_name = stage_sd_name(&workflow_name, &stage.name);
let mut env = stage.env.clone();
env.entry("SERVICE_NAME".to_owned())
.or_insert_with(|| sd_name.clone());
let mut svc = ServiceDefinition::new(
&sd_name,
ServiceDefinitionSpec {
description: Some(stage.description.clone().unwrap_or_else(|| {
format!("Stage `{}` of workflow `{workflow_name}`", stage.name)
})),
owner: wf.spec.owner.clone(),
template: stage.template.clone(),
image: stage.image.clone(),
port: stage.port.unwrap_or_else(default_port),
replicas: stage.replicas.unwrap_or_else(default_replicas),
env,
},
);
svc.metadata.namespace = Some(ns.to_owned());
svc.metadata.owner_references = wf.controller_owner_ref(&()).map(|r| vec![r]);
svc.metadata.labels = Some(BTreeMap::from([
(MANAGED_BY_KEY.to_owned(), MANAGED_BY_VALUE.to_owned()),
("polychrome.dev/workflow".to_owned(), workflow_name),
(
"polychrome.dev/workflow-stage".to_owned(),
stage.name.clone(),
),
]));
svc
}
#[must_use]
pub fn error_policy(_wf: Arc<Workflow>, err: &Error, _ctx: Arc<Context>) -> Action {
tracing::warn!(error = %err, "workflow reconcile failed; requeuing");
Action::requeue(Duration::from_secs(10))
}
pub async fn run_workflow(
client: Client,
watch_client: Client,
namespace: &str,
) -> Result<(), Error> {
let workflows: Api<Workflow> = Api::namespaced(watch_client.clone(), namespace);
fanout::await_watchable(&workflows, namespace, "Workflow").await;
let svcdefs: Api<ServiceDefinition> = Api::namespaced(watch_client, namespace);
let ctx = Arc::new(Context { client });
Controller::new(workflows, watcher::Config::default())
.owns(svcdefs, watcher::Config::default())
.run(reconcile, error_policy, ctx)
.for_each(|res| async move {
if let Err(e) = res {
tracing::warn!(error = %e, "workflow reconcile stream item errored");
}
})
.await;
Ok(())
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use super::*;
use crate::workflow::WorkflowSpec;
fn stage(name: &str, deps: &[&str]) -> WorkflowStage {
WorkflowStage {
name: name.to_owned(),
kind: None,
depends_on: deps.iter().map(|d| (*d).to_owned()).collect(),
description: None,
template: None,
image: "example.test/img:1".to_owned(),
port: None,
replicas: None,
env: BTreeMap::new(),
}
}
fn workflow(name: &str, stages: Vec<WorkflowStage>) -> Workflow {
let mut wf = Workflow::new(
name,
WorkflowSpec {
description: None,
owner: Some("user:slack/U1".to_owned()),
stages,
},
);
wf.metadata.namespace = Some("polychrome-apps".to_owned());
wf.metadata.uid = Some("uid-1".to_owned());
wf
}
#[test]
fn topo_sort_orders_upstreams_first() {
let spec = workflow(
"pipe",
vec![
stage("api", &["data"]),
stage("data", &["src"]),
stage("src", &[]),
],
)
.spec;
let order = topo_sort(&spec).expect("valid DAG");
let names: Vec<&str> = order
.iter()
.map(|&i| spec.stages[i].name.as_str())
.collect();
assert_eq!(names, vec!["src", "data", "api"]);
}
#[test]
fn topo_sort_rejects_cycles() {
let spec = workflow("pipe", vec![stage("a", &["b"]), stage("b", &["a"])]).spec;
let err = topo_sort(&spec).unwrap_err();
assert!(err.contains("cycle"), "{err}");
assert!(err.contains('a') && err.contains('b'));
}
#[test]
fn topo_sort_rejects_dangling_dependency() {
let spec = workflow("pipe", vec![stage("a", &["ghost"])]).spec;
let err = topo_sort(&spec).unwrap_err();
assert!(err.contains("ghost"), "{err}");
}
#[test]
fn topo_sort_rejects_self_dependency() {
let spec = workflow("pipe", vec![stage("a", &["a"])]).spec;
assert!(topo_sort(&spec).unwrap_err().contains("itself"));
}
#[test]
fn topo_sort_rejects_duplicate_and_empty() {
let dup = workflow("pipe", vec![stage("a", &[]), stage("a", &[])]).spec;
assert!(topo_sort(&dup).unwrap_err().contains("duplicate"));
let empty = workflow("pipe", vec![]).spec;
assert!(topo_sort(&empty).unwrap_err().contains("no stages"));
}
#[test]
fn plan_rejects_overlong_fanned_out_name() {
let long = "a".repeat(60);
let wf = workflow("pipeline", vec![stage(&long, &[])]);
match plan(&wf) {
WorkflowAction::Invalid(r) => assert!(r.contains("over the"), "{r}"),
other => panic!("expected Invalid, got {other:?}"),
}
}
#[test]
fn plan_rejects_bad_stage_name() {
let wf = workflow("pipe", vec![stage("Bad_Name", &[])]);
match plan(&wf) {
WorkflowAction::Invalid(r) => assert!(r.contains("DNS-1123"), "{r}"),
other => panic!("expected Invalid, got {other:?}"),
}
}
#[test]
fn plan_rejects_bad_workflow_name() {
for bad in ["my.pipe", "1pipe", "Pipe"] {
let wf = workflow(bad, vec![stage("data", &[])]);
match plan(&wf) {
WorkflowAction::Invalid(r) => {
assert!(r.contains("workflow name") && r.contains("DNS-1123"), "{r}");
}
other => panic!("expected Invalid for `{bad}`, got {other:?}"),
}
}
}
#[test]
fn deleting_workflow_plans_noop() {
let mut wf = workflow("pipe", vec![stage("a", &[])]);
wf.metadata.deletion_timestamp =
Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
"2026-06-10T00:00:00Z".parse().unwrap(),
));
assert_eq!(plan(&wf), WorkflowAction::Noop);
}
#[test]
fn unblocked_gates_on_upstream_readiness() {
let spec = workflow("pipe", vec![stage("data", &[]), stage("api", &["data"])]).spec;
let none = unblocked(&spec, &BTreeSet::new());
assert!(none.contains("data") && !none.contains("api"));
let data_ready: BTreeSet<String> = ["data".to_owned()].into_iter().collect();
let after = unblocked(&spec, &data_ready);
assert!(after.contains("data") && after.contains("api"));
}
#[test]
fn stage_service_definition_is_owned_and_named() {
let wf = workflow("pipe", vec![stage("data", &[])]);
let s = &wf.spec.stages[0];
let sd = build_stage_service_definition(&wf, s, "polychrome-apps");
assert_eq!(sd.metadata.name.as_deref(), Some("pipe-data"));
let owners = sd.metadata.owner_references.expect("owner ref set");
assert_eq!(owners[0].kind, "Workflow");
assert_eq!(
sd.spec.env.get("SERVICE_NAME").map(String::as_str),
Some("pipe-data")
);
assert_eq!(sd.spec.owner.as_deref(), Some("user:slack/U1"));
assert_eq!(sd.metadata.namespace.as_deref(), Some("polychrome-apps"));
}
#[test]
fn stage_name_length_check_uses_combined_name() {
let wf_name = "wf";
let ok = workflow(wf_name, vec![stage(&format!("d{}", "a".repeat(59)), &[])]);
assert!(matches!(plan(&ok), WorkflowAction::Apply(_)));
let over = workflow(wf_name, vec![stage(&format!("d{}", "a".repeat(60)), &[])]);
assert!(matches!(plan(&over), WorkflowAction::Invalid(_)));
}
#[test]
fn validate_matches_plan_validation() {
let ok = workflow("pipe", vec![stage("data", &[]), stage("api", &["data"])]);
assert!(validate(&ok).is_ok());
let cyclic = workflow("pipe", vec![stage("a", &["b"]), stage("b", &["a"])]);
assert!(validate(&cyclic).unwrap_err().contains("cycle"));
let over = workflow("wf", vec![stage(&format!("d{}", "a".repeat(60)), &[])]);
assert!(validate(&over).unwrap_err().contains("over the"));
}
#[test]
fn is_owned_by_distinguishes_own_foreign_and_unowned() {
let wf = workflow("pipe", vec![stage("data", &[])]);
let owned = build_stage_service_definition(&wf, &wf.spec.stages[0], "polychrome-apps");
assert!(is_owned_by(&owned, &wf));
let mut foreign = ServiceDefinition::new(
"pipe-data",
crate::servicedefinition::ServiceDefinitionSpec {
description: None,
owner: None,
template: None,
image: "img:1".to_owned(),
port: 8080,
replicas: 1,
env: BTreeMap::new(),
},
);
assert!(!is_owned_by(&foreign, &wf));
let mut other = workflow("pipe", vec![stage("data", &[])]);
other.metadata.uid = Some("uid-2".to_owned());
foreign.metadata.owner_references = other.controller_owner_ref(&()).map(|r| vec![r]);
assert!(!is_owned_by(&foreign, &wf));
foreign.metadata.owner_references = Some(vec![
k8s_openapi::apimachinery::pkg::apis::meta::v1::OwnerReference {
api_version: "polychrome.dev/v1alpha1".to_owned(),
kind: "Workflow".to_owned(),
name: "pipe".to_owned(),
uid: String::new(),
controller: Some(true),
block_owner_deletion: None,
},
]);
assert!(!is_owned_by(&foreign, &wf));
}
}