use std::{collections::BTreeMap, sync::Arc, time::Duration};
use futures::StreamExt;
use k8s_openapi::{
api::{
apps::v1::{Deployment, DeploymentSpec},
core::v1::{
Capabilities, Container, ContainerPort, EnvVar, PodSpec, PodTemplateSpec, Probe,
ResourceRequirements, SeccompProfile, SecurityContext, Service, ServicePort,
ServiceSpec, TCPSocketAction,
},
},
apimachinery::pkg::{
api::resource::Quantity, apis::meta::v1::LabelSelector, util::intstr::IntOrString,
},
};
use kube::{
Api, Client, Resource, ResourceExt,
api::{ObjectMeta, Patch, PatchParams},
runtime::{
controller::{Action, Controller},
watcher,
},
};
use serde_json::json;
use crate::{
fanout::{self},
servicedefinition::ServiceDefinition,
};
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("kube api: {0}")]
Kube(#[from] kube::Error),
#[error("service definition has no namespace")]
NoNamespace,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ServiceDefinitionAction {
Apply,
Noop,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct WorkloadReadiness {
pub available_replicas: i32,
pub ready: bool,
}
#[must_use]
pub fn plan(svc: &ServiceDefinition) -> ServiceDefinitionAction {
if svc.meta().deletion_timestamp.is_some() {
ServiceDefinitionAction::Noop
} else {
ServiceDefinitionAction::Apply
}
}
#[must_use]
pub fn workload_readiness(deployment: Option<&Deployment>, desired: i32) -> WorkloadReadiness {
let Some(deployment) = deployment else {
return WorkloadReadiness::default();
};
let generation = deployment.metadata.generation.unwrap_or(0);
let status = deployment.status.as_ref();
let observed_generation = status.and_then(|s| s.observed_generation).unwrap_or(0);
let available = status.and_then(|s| s.available_replicas).unwrap_or(0);
let updated = status.and_then(|s| s.updated_replicas).unwrap_or(0);
let total = status.and_then(|s| s.replicas).unwrap_or(0);
WorkloadReadiness {
available_replicas: available,
ready: desired > 0
&& observed_generation >= generation
&& updated >= desired
&& available >= desired
&& total == updated,
}
}
pub struct Context {
pub client: Client,
}
#[tracing::instrument(skip_all, fields(servicedefinition = %svc.name_any()))]
pub async fn reconcile(svc: Arc<ServiceDefinition>, ctx: Arc<Context>) -> Result<Action, Error> {
let ns = svc.namespace().ok_or(Error::NoNamespace)?;
let name = svc.name_any();
match plan(&svc) {
ServiceDefinitionAction::Noop => Ok(Action::await_change()),
ServiceDefinitionAction::Apply => {
let pp = PatchParams::apply("polychrome.dev/controller").force();
let deployments: Api<Deployment> = Api::namespaced(ctx.client.clone(), &ns);
let services: Api<Service> = Api::namespaced(ctx.client.clone(), &ns);
let observed = deployments
.patch(&name, &pp, &Patch::Apply(&build_deployment(&svc, &ns)))
.await?;
services
.patch(&name, &pp, &Patch::Apply(&build_service(&svc, &ns)))
.await?;
let readiness = workload_readiness(Some(&observed), svc.spec.replicas);
let current = svc.status.clone().unwrap_or_default();
if current.ready != readiness.ready
|| current.available_replicas != readiness.available_replicas
{
let phase = if readiness.ready { "Ready" } else { "Pending" };
let status = json!({ "status": {
"ready": readiness.ready,
"availableReplicas": readiness.available_replicas,
"phase": phase,
"message": null,
} });
let svcdefs: Api<ServiceDefinition> = Api::namespaced(ctx.client.clone(), &ns);
svcdefs
.patch_status(&name, &PatchParams::default(), &Patch::Merge(&status))
.await?;
tracing::info!(
ready = readiness.ready,
available = readiness.available_replicas,
"synced service definition status from deployment"
);
}
if readiness.ready {
Ok(Action::requeue(Duration::from_mins(5)))
} else {
Ok(Action::requeue(Duration::from_secs(15)))
}
}
}
}
fn labels_for(svc: &ServiceDefinition) -> BTreeMap<String, String> {
crate::fanout::labels(&svc.name_any(), "polychrome.dev/service-definition")
}
#[must_use]
pub fn build_deployment(svc: &ServiceDefinition, ns: &str) -> Deployment {
let name = svc.name_any();
let labels = labels_for(svc);
let mut env: Vec<EnvVar> = svc
.spec
.env
.iter()
.map(|(k, v)| EnvVar {
name: k.clone(),
value: Some(v.clone()),
value_from: None,
})
.collect();
if !svc.spec.env.contains_key("PORT") {
env.push(EnvVar {
name: "PORT".to_owned(),
value: Some(svc.spec.port.to_string()),
value_from: None,
});
}
Deployment {
metadata: ObjectMeta {
name: Some(name.clone()),
namespace: Some(ns.to_owned()),
labels: Some(labels.clone()),
owner_references: svc.controller_owner_ref(&()).map(|r| vec![r]),
..ObjectMeta::default()
},
spec: Some(DeploymentSpec {
replicas: Some(svc.spec.replicas),
selector: LabelSelector {
match_labels: Some(BTreeMap::from([(
"app.kubernetes.io/name".to_owned(),
name,
)])),
..LabelSelector::default()
},
template: PodTemplateSpec {
metadata: Some(ObjectMeta {
labels: Some(labels),
..ObjectMeta::default()
}),
spec: Some(PodSpec {
automount_service_account_token: Some(false),
security_context: Some(k8s_openapi::api::core::v1::PodSecurityContext {
run_as_non_root: Some(true),
run_as_user: Some(65532),
run_as_group: Some(65532),
seccomp_profile: Some(SeccompProfile {
type_: "RuntimeDefault".to_owned(),
localhost_profile: None,
}),
..Default::default()
}),
containers: vec![Container {
name: "app".to_owned(),
image: Some(svc.spec.image.clone()),
ports: Some(vec![ContainerPort {
name: Some("http".to_owned()),
container_port: svc.spec.port,
..ContainerPort::default()
}]),
env: (!env.is_empty()).then_some(env),
security_context: Some(SecurityContext {
allow_privilege_escalation: Some(false),
read_only_root_filesystem: Some(true),
capabilities: Some(Capabilities {
drop: Some(vec!["ALL".to_owned()]),
add: None,
}),
..Default::default()
}),
readiness_probe: Some(Probe {
tcp_socket: Some(TCPSocketAction {
port: IntOrString::String("http".to_owned()),
host: None,
}),
initial_delay_seconds: Some(2),
period_seconds: Some(10),
..Probe::default()
}),
resources: Some(ResourceRequirements {
requests: Some(BTreeMap::from([
("cpu".to_owned(), Quantity("50m".to_owned())),
("memory".to_owned(), Quantity("64Mi".to_owned())),
])),
limits: Some(BTreeMap::from([
("cpu".to_owned(), Quantity("500m".to_owned())),
("memory".to_owned(), Quantity("256Mi".to_owned())),
])),
..ResourceRequirements::default()
}),
..Container::default()
}],
..PodSpec::default()
}),
},
..DeploymentSpec::default()
}),
status: None,
}
}
#[must_use]
pub fn build_service(svc: &ServiceDefinition, ns: &str) -> Service {
let name = svc.name_any();
Service {
metadata: ObjectMeta {
name: Some(name.clone()),
namespace: Some(ns.to_owned()),
labels: Some(labels_for(svc)),
owner_references: svc.controller_owner_ref(&()).map(|r| vec![r]),
..ObjectMeta::default()
},
spec: Some(ServiceSpec {
selector: Some(BTreeMap::from([(
"app.kubernetes.io/name".to_owned(),
name,
)])),
ports: Some(vec![ServicePort {
name: Some("http".to_owned()),
port: 80,
target_port: Some(IntOrString::String("http".to_owned())),
..ServicePort::default()
}]),
..ServiceSpec::default()
}),
status: None,
}
}
#[must_use]
pub fn error_policy(_svc: Arc<ServiceDefinition>, err: &Error, _ctx: Arc<Context>) -> Action {
tracing::warn!(error = %err, "service definition reconcile failed; requeuing");
Action::requeue(Duration::from_secs(10))
}
pub async fn run_servicedefinition(
client: Client,
watch_client: Client,
namespace: &str,
) -> Result<(), Error> {
let svcdefs: Api<ServiceDefinition> = Api::namespaced(watch_client.clone(), namespace);
fanout::await_watchable(&svcdefs, namespace, "ServiceDefinition").await;
let deployments: Api<Deployment> = Api::namespaced(watch_client.clone(), namespace);
let services: Api<Service> = Api::namespaced(watch_client, namespace);
let ctx = Arc::new(Context { client });
Controller::new(svcdefs, watcher::Config::default())
.owns(deployments, watcher::Config::default())
.owns(services, watcher::Config::default())
.run(reconcile, error_policy, ctx)
.for_each(|res| async move {
if let Err(e) = res {
tracing::warn!(error = %e, "service definition reconcile stream item errored");
}
})
.await;
Ok(())
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use k8s_openapi::api::apps::v1::DeploymentStatus;
use super::*;
use crate::servicedefinition::ServiceDefinitionSpec;
fn svcdef(name: &str) -> ServiceDefinition {
let mut svc = ServiceDefinition::new(
name,
ServiceDefinitionSpec {
description: Some("test".to_owned()),
owner: Some("user:slack/U1".to_owned()),
template: Some("rust-hello-world".to_owned()),
image: "example.test/hello:1".to_owned(),
port: 8080,
replicas: 1,
env: BTreeMap::new(),
},
);
svc.metadata.namespace = Some("polychrome-apps".to_owned());
svc.metadata.uid = Some("uid-1".to_owned());
svc
}
#[test]
fn live_object_plans_apply() {
assert_eq!(plan(&svcdef("hello")), ServiceDefinitionAction::Apply);
}
#[test]
fn deleting_object_plans_noop() {
let mut svc = svcdef("hello");
svc.metadata.deletion_timestamp =
Some(k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(
"2026-06-10T00:00:00Z".parse().unwrap(),
));
assert_eq!(plan(&svc), ServiceDefinitionAction::Noop);
}
fn rolled_out(generation: i64, count: i32) -> Deployment {
Deployment {
metadata: ObjectMeta {
generation: Some(generation),
..ObjectMeta::default()
},
status: Some(DeploymentStatus {
observed_generation: Some(generation),
replicas: Some(count),
updated_replicas: Some(count),
available_replicas: Some(count),
..Default::default()
}),
..Default::default()
}
}
#[test]
fn readiness_unobserved_is_not_ready() {
assert_eq!(
workload_readiness(None, 1),
WorkloadReadiness {
available_replicas: 0,
ready: false,
}
);
}
#[test]
fn readiness_covers_desired_replicas() {
let d = rolled_out(1, 2);
assert_eq!(
workload_readiness(Some(&d), 2),
WorkloadReadiness {
available_replicas: 2,
ready: true,
}
);
assert_eq!(
workload_readiness(Some(&d), 3),
WorkloadReadiness {
available_replicas: 2,
ready: false,
}
);
}
#[test]
fn rolling_update_is_not_ready_until_new_generation_is_serving() {
let mid_rollout = Deployment {
metadata: ObjectMeta {
generation: Some(2),
..ObjectMeta::default()
},
status: Some(DeploymentStatus {
observed_generation: Some(1),
replicas: Some(3),
updated_replicas: Some(1),
available_replicas: Some(2),
..Default::default()
}),
..Default::default()
};
assert!(!workload_readiness(Some(&mid_rollout), 2).ready);
assert!(workload_readiness(Some(&rolled_out(2, 2)), 2).ready);
}
#[test]
fn surge_replicas_block_ready_until_old_generation_drains() {
let surging = Deployment {
metadata: ObjectMeta {
generation: Some(2),
..ObjectMeta::default()
},
status: Some(DeploymentStatus {
observed_generation: Some(2),
replicas: Some(3),
updated_replicas: Some(2),
available_replicas: Some(2),
..Default::default()
}),
..Default::default()
};
assert!(!workload_readiness(Some(&surging), 2).ready);
}
#[test]
fn zero_desired_replicas_is_never_ready() {
let d = Deployment::default();
assert!(!workload_readiness(Some(&d), 0).ready);
}
#[test]
fn deployment_is_owned_and_hardened() {
let svc = svcdef("hello");
let d = build_deployment(&svc, "polychrome-apps");
let owners = d.metadata.owner_references.expect("owner ref set");
assert_eq!(owners[0].kind, "ServiceDefinition");
let pod = d.spec.unwrap().template.spec.unwrap();
assert_eq!(pod.automount_service_account_token, Some(false));
let container = &pod.containers[0];
assert_eq!(container.image.as_deref(), Some("example.test/hello:1"));
let sc = container.security_context.as_ref().unwrap();
assert_eq!(sc.read_only_root_filesystem, Some(true));
assert_eq!(sc.allow_privilege_escalation, Some(false));
assert!(container.resources.as_ref().unwrap().limits.is_some());
}
#[test]
fn service_targets_declared_port_by_name() {
let svc = svcdef("hello");
let s = build_service(&svc, "polychrome-apps");
let ports = s.spec.unwrap().ports.unwrap();
assert_eq!(ports[0].port, 80);
assert_eq!(
ports[0].target_port,
Some(IntOrString::String("http".to_owned()))
);
}
}