use std::collections::BTreeMap;
use k8s_openapi::api::apps::v1::{
Deployment, DeploymentSpec, RollingUpdateStatefulSetStrategy, StatefulSet,
StatefulSetPersistentVolumeClaimRetentionPolicy, StatefulSetSpec, StatefulSetUpdateStrategy,
};
use k8s_openapi::api::autoscaling::v2::{
CrossVersionObjectReference, HorizontalPodAutoscaler, HorizontalPodAutoscalerSpec, MetricSpec,
MetricTarget, ResourceMetricSource,
};
use k8s_openapi::api::core::v1::{
ConfigMap, ConfigMapVolumeSource, Container, ContainerPort, EnvVar, EnvVarSource,
HTTPGetAction, ObjectFieldSelector, PersistentVolumeClaim, PersistentVolumeClaimSpec, PodSpec,
PodTemplateSpec, Probe, Service, ServicePort, ServiceSpec, TCPSocketAction, Volume,
VolumeMount, VolumeResourceRequirements,
};
use k8s_openapi::api::policy::v1::{PodDisruptionBudget, PodDisruptionBudgetSpec};
use k8s_openapi::apimachinery::pkg::api::resource::Quantity;
use k8s_openapi::apimachinery::pkg::apis::meta::v1::{LabelSelector, ObjectMeta};
use k8s_openapi::apimachinery::pkg::util::intstr::IntOrString;
use kube::Resource;
use super::crd::{BoatRampCluster, ClusterMode};
const PORT: i32 = 8080;
const MESH_PORT: i32 = 7000;
const CONFIG_MOUNT: &str = "/etc/boatramp";
const DATA_MOUNT: &str = "/data";
const DEFAULT_IMAGE: &str = "ghcr.io/boatramp/boatramp:latest";
fn labels(name: &str) -> BTreeMap<String, String> {
[
("app.kubernetes.io/name".to_string(), "boatramp".to_string()),
("app.kubernetes.io/instance".to_string(), name.to_string()),
(
"app.kubernetes.io/managed-by".to_string(),
"boatramp-operator".to_string(),
),
]
.into()
}
fn child_meta(brc: &BoatRampCluster, name: String) -> ObjectMeta {
ObjectMeta {
name: Some(name),
namespace: brc.metadata.namespace.clone(),
labels: Some(labels(&instance(brc))),
owner_references: brc.controller_owner_ref(&()).map(|r| vec![r]),
..Default::default()
}
}
pub fn instance(brc: &BoatRampCluster) -> String {
brc.metadata
.name
.clone()
.unwrap_or_else(|| "boatramp".to_string())
}
fn image(brc: &BoatRampCluster) -> String {
brc.spec
.image
.clone()
.unwrap_or_else(|| DEFAULT_IMAGE.to_string())
}
fn headless_name(brc: &BoatRampCluster) -> String {
format!("{}-headless", instance(brc))
}
fn config_ron(brc: &BoatRampCluster) -> String {
let posture = brc.spec.posture.as_deref().unwrap_or("multi-tenant");
let auth = match brc.spec.root_pubkey.as_deref() {
Some(pubkey) => format!(" auth_root_public_key: \"{pubkey}\",\n"),
None => String::new(),
};
let cluster = if brc.spec.mode == ClusterMode::Cluster {
format!(
" cluster: (\n \
listen: \"0.0.0.0:{MESH_PORT}\",\n \
store_dir: \"{DATA_MOUNT}/raft\",\n \
),\n"
)
} else {
String::new()
};
format!(
"(\n \
serve: (\n \
addr: \"0.0.0.0:{PORT}\",\n \
data_dir: \"{DATA_MOUNT}\",\n\
{auth} \
),\n\
{cluster} \
security: ( profile: \"{posture}\" ),\n\
)\n"
)
}
pub fn config_map(brc: &BoatRampCluster) -> ConfigMap {
ConfigMap {
metadata: child_meta(brc, format!("{}-config", instance(brc))),
data: Some([("boatramp.cfg".to_string(), config_ron(brc))].into()),
..Default::default()
}
}
pub fn headless_service(brc: &BoatRampCluster) -> Service {
Service {
metadata: child_meta(brc, headless_name(brc)),
spec: Some(ServiceSpec {
cluster_ip: Some("None".to_string()),
selector: Some(labels(&instance(brc))),
ports: Some(vec![port("http", PORT), port("mesh", MESH_PORT)]),
publish_not_ready_addresses: Some(true),
..Default::default()
}),
..Default::default()
}
}
pub fn client_service(brc: &BoatRampCluster) -> Service {
Service {
metadata: child_meta(brc, instance(brc)),
spec: Some(ServiceSpec {
selector: Some(labels(&instance(brc))),
ports: Some(vec![port("http", PORT)]),
..Default::default()
}),
..Default::default()
}
}
fn downward_env(name: &str, field_path: &str) -> EnvVar {
EnvVar {
name: name.to_string(),
value_from: Some(EnvVarSource {
field_ref: Some(ObjectFieldSelector {
field_path: field_path.to_string(),
..Default::default()
}),
..Default::default()
}),
..Default::default()
}
}
fn secret_env(name: &str, secret: &str, key: &str) -> EnvVar {
EnvVar {
name: name.to_string(),
value_from: Some(EnvVarSource {
secret_key_ref: Some(k8s_openapi::api::core::v1::SecretKeySelector {
name: secret.to_string(),
key: key.to_string(),
optional: Some(true),
}),
..Default::default()
}),
..Default::default()
}
}
fn port(name: &str, p: i32) -> ServicePort {
ServicePort {
name: Some(name.to_string()),
port: p,
target_port: Some(IntOrString::Int(p)),
..Default::default()
}
}
fn container(brc: &BoatRampCluster) -> Container {
Container {
name: "boatramp".to_string(),
image: Some(image(brc)),
image_pull_policy: Some("IfNotPresent".to_string()),
command: Some(vec!["boatramp".to_string()]),
args: Some({
let mut args = vec![
"serve".to_string(),
"--config".to_string(),
format!("{CONFIG_MOUNT}/boatramp.cfg"),
];
if brc.spec.mode == ClusterMode::Cluster {
args.push("--tls".to_string());
args.push("rpk".to_string());
}
args
}),
ports: Some({
let mut ports = vec![ContainerPort {
name: Some("http".to_string()),
container_port: PORT,
..Default::default()
}];
if brc.spec.mode == ClusterMode::Cluster {
ports.push(ContainerPort {
name: Some("mesh".to_string()),
container_port: MESH_PORT,
..Default::default()
});
}
ports
}),
env: Some(
vec![
downward_env("BOATRAMP_POD_NAME", "metadata.name"),
{
let (secret, key) = super::executor::join_env_source(brc);
EnvVar {
name: "BOATRAMP_CLUSTER_JOIN".to_string(),
value_from: Some(EnvVarSource {
secret_key_ref: Some(k8s_openapi::api::core::v1::SecretKeySelector {
name: secret,
key: key.to_string(),
optional: Some(true),
}),
..Default::default()
}),
..Default::default()
}
},
]
.into_iter()
.chain(brc.spec.auth_secret.as_deref().into_iter().flat_map(|s| {
[
secret_env("BOATRAMP_AUTH_ROOT_PRIVATE_KEY", s, "root-private-key"),
secret_env("BOATRAMP_BOOTSTRAP_SECRET", s, "bootstrap-secret"),
]
}))
.chain(
(brc.spec.mode == ClusterMode::Cluster)
.then(|| {
vec![
downward_env("POD_NAMESPACE", "metadata.namespace"),
EnvVar {
name: "BOATRAMP_CLUSTER_ADVERTISE_ADDR".to_string(),
value: Some(format!(
"https://$(BOATRAMP_POD_NAME).{}.$(POD_NAMESPACE).svc:{MESH_PORT}",
headless_name(brc)
)),
..Default::default()
},
]
})
.into_iter()
.flatten(),
)
.collect(),
),
liveness_probe: Some(if brc.spec.mode == ClusterMode::Cluster {
tcp_probe()
} else {
http_probe("/healthz")
}),
readiness_probe: Some(if brc.spec.mode == ClusterMode::Cluster {
tcp_probe()
} else {
http_probe("/readyz")
}),
volume_mounts: Some(vec![
VolumeMount {
name: "config".to_string(),
mount_path: CONFIG_MOUNT.to_string(),
read_only: Some(true),
..Default::default()
},
VolumeMount {
name: "data".to_string(),
mount_path: DATA_MOUNT.to_string(),
..Default::default()
},
]),
..Default::default()
}
}
fn http_probe(path: &str) -> Probe {
Probe {
http_get: Some(HTTPGetAction {
path: Some(path.to_string()),
port: IntOrString::Int(PORT),
..Default::default()
}),
period_seconds: Some(10),
..Default::default()
}
}
fn tcp_probe() -> Probe {
Probe {
tcp_socket: Some(TCPSocketAction {
port: IntOrString::Int(PORT),
..Default::default()
}),
period_seconds: Some(10),
..Default::default()
}
}
fn pod_template(brc: &BoatRampCluster, data_volume: Option<Volume>) -> PodTemplateSpec {
let mut volumes = vec![Volume {
name: "config".to_string(),
config_map: Some(ConfigMapVolumeSource {
name: format!("{}-config", instance(brc)),
..Default::default()
}),
..Default::default()
}];
volumes.extend(data_volume);
PodTemplateSpec {
metadata: Some(ObjectMeta {
labels: Some(labels(&instance(brc))),
..Default::default()
}),
spec: Some(PodSpec {
containers: vec![container(brc)],
volumes: Some(volumes),
..Default::default()
}),
}
}
fn selector(brc: &BoatRampCluster) -> LabelSelector {
LabelSelector {
match_labels: Some(labels(&instance(brc))),
..Default::default()
}
}
pub fn stateful_set(brc: &BoatRampCluster, roll_partition: i32) -> StatefulSet {
let storage = brc
.spec
.storage
.clone()
.unwrap_or_else(|| "10Gi".to_string());
let pvc = PersistentVolumeClaim {
metadata: ObjectMeta {
name: Some("data".to_string()),
..Default::default()
},
spec: Some(PersistentVolumeClaimSpec {
access_modes: Some(vec!["ReadWriteOnce".to_string()]),
resources: Some(VolumeResourceRequirements {
requests: Some([("storage".to_string(), Quantity(storage))].into()),
..Default::default()
}),
..Default::default()
}),
..Default::default()
};
StatefulSet {
metadata: child_meta(brc, instance(brc)),
spec: Some(StatefulSetSpec {
replicas: Some(brc.spec.replicas as i32),
service_name: Some(headless_name(brc)),
selector: selector(brc),
template: pod_template(brc, None),
volume_claim_templates: Some(vec![pvc]),
update_strategy: Some(StatefulSetUpdateStrategy {
type_: Some("RollingUpdate".to_string()),
rolling_update: Some(RollingUpdateStatefulSetStrategy {
partition: Some(roll_partition),
..Default::default()
}),
}),
persistent_volume_claim_retention_policy: Some(
StatefulSetPersistentVolumeClaimRetentionPolicy {
when_deleted: Some("Retain".to_string()),
when_scaled: Some("Retain".to_string()),
},
),
..Default::default()
}),
..Default::default()
}
}
pub fn deployment(brc: &BoatRampCluster) -> Deployment {
let data = Volume {
name: "data".to_string(),
empty_dir: Some(Default::default()),
..Default::default()
};
Deployment {
metadata: child_meta(brc, instance(brc)),
spec: Some(DeploymentSpec {
replicas: Some(brc.spec.replicas as i32),
selector: selector(brc),
template: pod_template(brc, Some(data)),
..Default::default()
}),
..Default::default()
}
}
pub fn pod_disruption_budget(brc: &BoatRampCluster) -> PodDisruptionBudget {
let min_available = (brc.spec.replicas / 2) + 1;
PodDisruptionBudget {
metadata: child_meta(brc, instance(brc)),
spec: Some(PodDisruptionBudgetSpec {
min_available: Some(IntOrString::Int(min_available as i32)),
selector: Some(selector(brc)),
..Default::default()
}),
..Default::default()
}
}
pub fn hpa(brc: &BoatRampCluster) -> HorizontalPodAutoscaler {
HorizontalPodAutoscaler {
metadata: child_meta(brc, instance(brc)),
spec: Some(HorizontalPodAutoscalerSpec {
scale_target_ref: CrossVersionObjectReference {
api_version: Some("apps/v1".to_string()),
kind: "Deployment".to_string(),
name: instance(brc),
},
min_replicas: Some(brc.spec.replicas.max(1) as i32),
max_replicas: (brc.spec.replicas.max(1) * 4) as i32,
metrics: Some(vec![MetricSpec {
type_: "Resource".to_string(),
resource: Some(ResourceMetricSource {
name: "cpu".to_string(),
target: MetricTarget {
type_: "Utilization".to_string(),
average_utilization: Some(75),
..Default::default()
},
}),
..Default::default()
}]),
..Default::default()
}),
..Default::default()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::operator::crd::{BoatRampClusterSpec, ClusterMode};
fn cluster(name: &str, mode: ClusterMode, replicas: u32) -> BoatRampCluster {
let mut brc = BoatRampCluster::new(
name,
BoatRampClusterSpec {
mode,
replicas,
image: None,
storage: None,
posture: None,
admin_token_secret: None,
root_pubkey: None,
auth_secret: None,
},
);
brc.metadata.namespace = Some("tenant-a".to_string());
brc.metadata.uid = Some("uid-123".to_string());
brc
}
#[test]
fn statefulset_has_pvc_stable_dns_probes_and_owner_ref() {
let brc = cluster("db", ClusterMode::Cluster, 3);
let sts = stateful_set(&brc, 3);
let spec = sts.spec.unwrap();
assert_eq!(spec.replicas, Some(3));
assert_eq!(spec.service_name.as_deref(), Some("db-headless"));
assert_eq!(spec.volume_claim_templates.as_ref().unwrap().len(), 1);
assert_eq!(
spec.update_strategy
.as_ref()
.and_then(|u| u.rolling_update.as_ref())
.and_then(|r| r.partition),
Some(3)
);
let retain = spec
.persistent_volume_claim_retention_policy
.as_ref()
.unwrap();
assert_eq!(retain.when_deleted.as_deref(), Some("Retain"));
assert_eq!(retain.when_scaled.as_deref(), Some("Retain"));
let c = &spec.template.spec.unwrap().containers[0];
assert!(c.liveness_probe.as_ref().unwrap().tcp_socket.is_some());
assert!(c.readiness_probe.as_ref().unwrap().tcp_socket.is_some());
let owners = sts.metadata.owner_references.unwrap();
assert_eq!(owners[0].kind, "BoatRampCluster");
assert_eq!(owners[0].controller, Some(true));
}
#[test]
fn cluster_pods_serve_rpk_tls_and_probe_by_tcp() {
let brc = cluster("db", ClusterMode::Cluster, 3);
let c = stateful_set(&brc, 0)
.spec
.unwrap()
.template
.spec
.unwrap()
.containers
.remove(0);
let args = c.args.unwrap();
assert_eq!(
args,
vec![
"serve".to_string(),
"--config".to_string(),
format!("{CONFIG_MOUNT}/boatramp.cfg"),
"--tls".to_string(),
"rpk".to_string(),
]
);
assert!(c.readiness_probe.as_ref().unwrap().tcp_socket.is_some());
assert!(c.readiness_probe.as_ref().unwrap().http_get.is_none());
assert!(c.liveness_probe.as_ref().unwrap().tcp_socket.is_some());
let web = cluster("web", ClusterMode::Stateless, 2);
let wc = deployment(&web)
.spec
.unwrap()
.template
.spec
.unwrap()
.containers
.remove(0);
assert!(!wc.args.unwrap().contains(&"rpk".to_string()));
let rp = wc.readiness_probe.unwrap();
assert!(rp.tcp_socket.is_none());
assert_eq!(rp.http_get.unwrap().path.as_deref(), Some("/readyz"));
}
#[test]
fn headless_service_is_headless_and_publishes_not_ready() {
let svc = headless_service(&cluster("db", ClusterMode::Cluster, 3));
let spec = svc.spec.unwrap();
assert_eq!(spec.cluster_ip.as_deref(), Some("None"));
assert_eq!(spec.publish_not_ready_addresses, Some(true));
}
#[test]
fn pdb_keeps_a_quorum_majority() {
for (n, want) in [(1, 1), (3, 2), (5, 3)] {
let pdb = pod_disruption_budget(&cluster("db", ClusterMode::Cluster, n));
let min = pdb.spec.unwrap().min_available.unwrap();
assert_eq!(min, IntOrString::Int(want), "n={n}");
}
}
#[test]
fn stateless_mode_is_a_deployment_with_hpa_and_ephemeral_data() {
let brc = cluster("web", ClusterMode::Stateless, 2);
let dep = deployment(&brc);
let tspec = dep.spec.unwrap().template.spec.unwrap();
let data_vol = tspec
.volumes
.unwrap()
.into_iter()
.find(|v| v.name == "data")
.unwrap();
assert!(data_vol.empty_dir.is_some());
let hpa = hpa(&brc);
let hspec = hpa.spec.unwrap();
assert_eq!(hspec.scale_target_ref.kind, "Deployment");
assert_eq!(hspec.min_replicas, Some(2));
assert_eq!(hspec.max_replicas, 8);
}
#[test]
fn config_map_carries_posture_and_bind() {
let mut brc = cluster("db", ClusterMode::Cluster, 3);
brc.spec.posture = Some("single-tenant".to_string());
let cm = config_map(&brc);
let cfg = &cm.data.unwrap()["boatramp.cfg"];
assert!(cfg.contains("profile: \"single-tenant\""));
assert!(cfg.contains("0.0.0.0:8080"));
}
#[test]
fn cluster_config_and_auth_are_wired_and_parse() {
let mut brc = cluster("db", ClusterMode::Cluster, 3);
brc.spec.root_pubkey = Some("es256:03a1".to_string());
brc.spec.auth_secret = Some("db-auth".to_string());
let cfg = &config_map(&brc).data.unwrap()["boatramp.cfg"];
let parsed = crate::config::ServerConfig::parse(cfg).expect("valid boatramp.cfg");
assert!(parsed.cluster.is_some(), "cluster mode config present");
assert_eq!(
parsed.serve.unwrap().auth_root_public_key.as_deref(),
Some("es256:03a1")
);
let sts = stateful_set(&brc, 0);
let env = sts
.spec
.unwrap()
.template
.spec
.unwrap()
.containers
.remove(0)
.env
.unwrap();
let names: Vec<&str> = env.iter().map(|e| e.name.as_str()).collect();
assert!(names.contains(&"BOATRAMP_AUTH_ROOT_PRIVATE_KEY"));
assert!(names.contains(&"BOATRAMP_BOOTSTRAP_SECRET"));
let advertise = env
.iter()
.find(|e| e.name == "BOATRAMP_CLUSTER_ADVERTISE_ADDR")
.and_then(|e| e.value.as_deref())
.unwrap();
assert_eq!(
advertise,
"https://$(BOATRAMP_POD_NAME).db-headless.$(POD_NAMESPACE).svc:7000"
);
}
}