WP-12: Kubernetes overview and workload actions
kube-rs gateway (nodes, namespaces, deployments/statefulsets/daemonsets with images, PVCs; rollout restart, scale, set image via patches) using the microk8s kubeconfig by default, fake gateway for dev, /api/cluster routes, Kubernetes page with node cards, workload table with admin actions and PVC list. Also implements the streaming command runner that WP-11 relied on. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@ -43,11 +43,52 @@ impl CommandRunner for SystemCommandRunner {
|
||||
|
||||
async fn run_streaming(
|
||||
&self,
|
||||
_program: &str,
|
||||
_args: &[&str],
|
||||
_out: &dyn LineSink,
|
||||
program: &str,
|
||||
args: &[&str],
|
||||
out: &dyn LineSink,
|
||||
) -> Result<bool, DomainError> {
|
||||
todo!()
|
||||
use tokio::io::{AsyncBufReadExt, BufReader};
|
||||
let mut child = tokio::process::Command::new(program)
|
||||
.args(args)
|
||||
.env("DEBIAN_FRONTEND", "noninteractive")
|
||||
.env("LC_ALL", "C")
|
||||
.stdin(std::process::Stdio::null())
|
||||
.stdout(std::process::Stdio::piped())
|
||||
.stderr(std::process::Stdio::piped())
|
||||
.spawn()
|
||||
.map_err(|e| DomainError::Unavailable(format!("{program}: {e}")))?;
|
||||
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<String>();
|
||||
let mut readers = Vec::new();
|
||||
if let Some(stdout) = child.stdout.take() {
|
||||
let tx = tx.clone();
|
||||
readers.push(tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stdout).lines();
|
||||
while let Ok(Some(l)) = lines.next_line().await {
|
||||
let _ = tx.send(l);
|
||||
}
|
||||
}));
|
||||
}
|
||||
if let Some(stderr) = child.stderr.take() {
|
||||
let tx = tx.clone();
|
||||
readers.push(tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stderr).lines();
|
||||
while let Ok(Some(l)) = lines.next_line().await {
|
||||
let _ = tx.send(l);
|
||||
}
|
||||
}));
|
||||
}
|
||||
drop(tx);
|
||||
while let Some(line) = rx.recv().await {
|
||||
out.line(&line);
|
||||
}
|
||||
for r in readers {
|
||||
let _ = r.await;
|
||||
}
|
||||
let status = child
|
||||
.wait()
|
||||
.await
|
||||
.map_err(|e| DomainError::Unavailable(format!("{program}: {e}")))?;
|
||||
Ok(status.success())
|
||||
}
|
||||
|
||||
async fn read_file(&self, path: &str) -> Result<Option<String>, DomainError> {
|
||||
|
||||
383
backend/crates/infrastructure/src/k8s.rs
Normal file
383
backend/crates/infrastructure/src/k8s.rs
Normal file
@ -0,0 +1,383 @@
|
||||
//! Kubernetes gateway on top of kube-rs, plus a fake for development.
|
||||
use async_trait::async_trait;
|
||||
use chrono::Utc;
|
||||
use domain::cluster::{
|
||||
ClusterOverview, Container, NodeInfo, VolumeClaim, Workload, WorkloadKind, WorkloadRef,
|
||||
};
|
||||
use domain::ports::ClusterGateway;
|
||||
use domain::DomainError;
|
||||
use k8s_openapi::api::apps::v1::{DaemonSet, Deployment, StatefulSet};
|
||||
use k8s_openapi::api::core::v1::{Namespace, Node, PersistentVolumeClaim, PodTemplateSpec};
|
||||
use kube::api::{Patch, PatchParams};
|
||||
use kube::config::{KubeConfigOptions, Kubeconfig};
|
||||
use kube::{Api, Client, Config};
|
||||
use serde_json::json;
|
||||
use tokio::sync::OnceCell;
|
||||
|
||||
pub struct KubeGateway {
|
||||
kubeconfig: Option<String>,
|
||||
client: OnceCell<Client>,
|
||||
}
|
||||
|
||||
impl KubeGateway {
|
||||
/// `kubeconfig`: path to a kubeconfig file; `None` infers (in-cluster or ~/.kube/config).
|
||||
pub fn new(kubeconfig: Option<String>) -> Self {
|
||||
Self {
|
||||
kubeconfig,
|
||||
client: OnceCell::new(),
|
||||
}
|
||||
}
|
||||
|
||||
async fn client(&self) -> Result<Client, DomainError> {
|
||||
let unavailable = |e: String| DomainError::Unavailable(format!("kubernetes: {e}"));
|
||||
self.client
|
||||
.get_or_try_init(|| async {
|
||||
let config = match &self.kubeconfig {
|
||||
Some(path) => {
|
||||
let kc =
|
||||
Kubeconfig::read_from(path).map_err(|e| unavailable(e.to_string()))?;
|
||||
Config::from_custom_kubeconfig(kc, &KubeConfigOptions::default())
|
||||
.await
|
||||
.map_err(|e| unavailable(e.to_string()))?
|
||||
}
|
||||
None => Config::infer()
|
||||
.await
|
||||
.map_err(|e| unavailable(e.to_string()))?,
|
||||
};
|
||||
Client::try_from(config).map_err(|e| unavailable(e.to_string()))
|
||||
})
|
||||
.await
|
||||
.cloned()
|
||||
}
|
||||
}
|
||||
|
||||
fn map_err(e: kube::Error) -> DomainError {
|
||||
match e {
|
||||
kube::Error::Api(ref r) if r.code == 404 => DomainError::NotFound,
|
||||
kube::Error::Api(ref r) if r.code == 422 => DomainError::Validation(r.message.clone()),
|
||||
e => DomainError::Unavailable(format!("kubernetes: {e}")),
|
||||
}
|
||||
}
|
||||
|
||||
fn containers(t: &Option<PodTemplateSpec>) -> Vec<Container> {
|
||||
t.as_ref()
|
||||
.and_then(|t| t.spec.as_ref())
|
||||
.map(|s| {
|
||||
s.containers
|
||||
.iter()
|
||||
.map(|c| Container {
|
||||
name: c.name.clone(),
|
||||
image: c.image.clone().unwrap_or_default(),
|
||||
})
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
fn restart_patch() -> Patch<serde_json::Value> {
|
||||
Patch::Merge(
|
||||
json!({"spec": {"template": {"metadata": {"annotations": {"kubectl.kubernetes.io/restartedAt": Utc::now().to_rfc3339()}}}}}),
|
||||
)
|
||||
}
|
||||
|
||||
fn image_patch(container: &str, image: &str) -> Patch<serde_json::Value> {
|
||||
Patch::Strategic(
|
||||
json!({"spec": {"template": {"spec": {"containers": [{"name": container, "image": image}]}}}}),
|
||||
)
|
||||
}
|
||||
|
||||
macro_rules! patch_kind {
|
||||
($client:expr, $w:expr, $patch:expr) => {{
|
||||
let pp = PatchParams::apply("softvisor-monitoring").force();
|
||||
let pp = PatchParams {
|
||||
field_manager: pp.field_manager,
|
||||
..PatchParams::default()
|
||||
};
|
||||
match $w.kind {
|
||||
WorkloadKind::Deployment => Api::<Deployment>::namespaced($client, &$w.namespace)
|
||||
.patch(&$w.name, &pp, &$patch)
|
||||
.await
|
||||
.map(|_| ()),
|
||||
WorkloadKind::StatefulSet => Api::<StatefulSet>::namespaced($client, &$w.namespace)
|
||||
.patch(&$w.name, &pp, &$patch)
|
||||
.await
|
||||
.map(|_| ()),
|
||||
WorkloadKind::DaemonSet => Api::<DaemonSet>::namespaced($client, &$w.namespace)
|
||||
.patch(&$w.name, &pp, &$patch)
|
||||
.await
|
||||
.map(|_| ()),
|
||||
}
|
||||
.map_err(map_err)
|
||||
}};
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ClusterGateway for KubeGateway {
|
||||
async fn overview(&self) -> Result<ClusterOverview, DomainError> {
|
||||
let client = self.client().await?;
|
||||
let lp = Default::default();
|
||||
let nodes = Api::<Node>::all(client.clone())
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
let namespaces = Api::<Namespace>::all(client.clone())
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
let deployments = Api::<Deployment>::all(client.clone())
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
let statefulsets = Api::<StatefulSet>::all(client.clone())
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
let daemonsets = Api::<DaemonSet>::all(client.clone())
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
let pvcs = Api::<PersistentVolumeClaim>::all(client)
|
||||
.list(&lp)
|
||||
.await
|
||||
.map_err(map_err)?;
|
||||
|
||||
let mut workloads = Vec::new();
|
||||
for d in deployments {
|
||||
workloads.push(Workload {
|
||||
namespace: d.metadata.namespace.clone().unwrap_or_default(),
|
||||
kind: WorkloadKind::Deployment,
|
||||
name: d.metadata.name.clone().unwrap_or_default(),
|
||||
ready: d
|
||||
.status
|
||||
.as_ref()
|
||||
.and_then(|s| s.ready_replicas)
|
||||
.unwrap_or(0),
|
||||
desired: d.spec.as_ref().and_then(|s| s.replicas).unwrap_or(0),
|
||||
containers: containers(&d.spec.as_ref().map(|s| s.template.clone())),
|
||||
});
|
||||
}
|
||||
for s in statefulsets {
|
||||
workloads.push(Workload {
|
||||
namespace: s.metadata.namespace.clone().unwrap_or_default(),
|
||||
kind: WorkloadKind::StatefulSet,
|
||||
name: s.metadata.name.clone().unwrap_or_default(),
|
||||
ready: s
|
||||
.status
|
||||
.as_ref()
|
||||
.and_then(|s| s.ready_replicas)
|
||||
.unwrap_or(0),
|
||||
desired: s.spec.as_ref().and_then(|s| s.replicas).unwrap_or(0),
|
||||
containers: containers(&s.spec.as_ref().map(|s| s.template.clone())),
|
||||
});
|
||||
}
|
||||
for d in daemonsets {
|
||||
workloads.push(Workload {
|
||||
namespace: d.metadata.namespace.clone().unwrap_or_default(),
|
||||
kind: WorkloadKind::DaemonSet,
|
||||
name: d.metadata.name.clone().unwrap_or_default(),
|
||||
ready: d.status.as_ref().map(|s| s.number_ready).unwrap_or(0),
|
||||
desired: d
|
||||
.status
|
||||
.as_ref()
|
||||
.map(|s| s.desired_number_scheduled)
|
||||
.unwrap_or(0),
|
||||
containers: containers(&d.spec.as_ref().map(|s| s.template.clone())),
|
||||
});
|
||||
}
|
||||
workloads.sort_by(|a, b| (&a.namespace, &a.name).cmp(&(&b.namespace, &b.name)));
|
||||
|
||||
Ok(ClusterOverview {
|
||||
nodes: nodes
|
||||
.iter()
|
||||
.map(|n| {
|
||||
let info = n.status.as_ref().and_then(|s| s.node_info.as_ref());
|
||||
NodeInfo {
|
||||
name: n.metadata.name.clone().unwrap_or_default(),
|
||||
version: info.map(|i| i.kubelet_version.clone()).unwrap_or_default(),
|
||||
ready: n
|
||||
.status
|
||||
.as_ref()
|
||||
.and_then(|s| s.conditions.as_ref())
|
||||
.is_some_and(|c| {
|
||||
c.iter().any(|c| c.type_ == "Ready" && c.status == "True")
|
||||
}),
|
||||
os_image: info.map(|i| i.os_image.clone()).unwrap_or_default(),
|
||||
kernel: info.map(|i| i.kernel_version.clone()).unwrap_or_default(),
|
||||
container_runtime: info
|
||||
.map(|i| i.container_runtime_version.clone())
|
||||
.unwrap_or_default(),
|
||||
}
|
||||
})
|
||||
.collect(),
|
||||
namespaces: namespaces
|
||||
.iter()
|
||||
.filter_map(|n| n.metadata.name.clone())
|
||||
.collect(),
|
||||
workloads,
|
||||
volume_claims: pvcs
|
||||
.iter()
|
||||
.map(|p| VolumeClaim {
|
||||
namespace: p.metadata.namespace.clone().unwrap_or_default(),
|
||||
name: p.metadata.name.clone().unwrap_or_default(),
|
||||
capacity: p
|
||||
.status
|
||||
.as_ref()
|
||||
.and_then(|s| s.capacity.as_ref())
|
||||
.and_then(|c| c.get("storage"))
|
||||
.map(|q| q.0.clone())
|
||||
.unwrap_or_default(),
|
||||
storage_class: p
|
||||
.spec
|
||||
.as_ref()
|
||||
.and_then(|s| s.storage_class_name.clone())
|
||||
.unwrap_or_default(),
|
||||
status: p
|
||||
.status
|
||||
.as_ref()
|
||||
.and_then(|s| s.phase.clone())
|
||||
.unwrap_or_default(),
|
||||
})
|
||||
.collect(),
|
||||
fetched_at: Utc::now(),
|
||||
})
|
||||
}
|
||||
|
||||
async fn restart(&self, w: &WorkloadRef) -> Result<(), DomainError> {
|
||||
let client = self.client().await?;
|
||||
patch_kind!(client.clone(), w, restart_patch())
|
||||
}
|
||||
|
||||
async fn scale(&self, w: &WorkloadRef, replicas: i32) -> Result<(), DomainError> {
|
||||
if w.kind == WorkloadKind::DaemonSet {
|
||||
return Err(DomainError::Validation(
|
||||
"daemonsets cannot be scaled".into(),
|
||||
));
|
||||
}
|
||||
let client = self.client().await?;
|
||||
patch_kind!(
|
||||
client.clone(),
|
||||
w,
|
||||
Patch::<serde_json::Value>::Merge(json!({"spec": {"replicas": replicas}}))
|
||||
)
|
||||
}
|
||||
|
||||
async fn set_image(
|
||||
&self,
|
||||
w: &WorkloadRef,
|
||||
container: &str,
|
||||
image: &str,
|
||||
) -> Result<(), DomainError> {
|
||||
let client = self.client().await?;
|
||||
patch_kind!(client.clone(), w, image_patch(container, image))
|
||||
}
|
||||
}
|
||||
|
||||
/// Sample cluster for development machines (FAKE_HOST=true).
|
||||
pub struct FakeClusterGateway;
|
||||
|
||||
#[async_trait]
|
||||
impl ClusterGateway for FakeClusterGateway {
|
||||
async fn overview(&self) -> Result<ClusterOverview, DomainError> {
|
||||
let c = |name: &str, image: &str| {
|
||||
vec![Container {
|
||||
name: name.into(),
|
||||
image: image.into(),
|
||||
}]
|
||||
};
|
||||
let w = |ns: &str, kind, name: &str, ready, desired, containers| Workload {
|
||||
namespace: ns.into(),
|
||||
kind,
|
||||
name: name.into(),
|
||||
ready,
|
||||
desired,
|
||||
containers,
|
||||
};
|
||||
Ok(ClusterOverview {
|
||||
nodes: vec![NodeInfo {
|
||||
name: "fake-node".into(),
|
||||
version: "v1.32.13".into(),
|
||||
ready: true,
|
||||
os_image: "Debian GNU/Linux 12 (bookworm)".into(),
|
||||
kernel: "6.1.0-42-amd64".into(),
|
||||
container_runtime: "containerd://1.6.36".into(),
|
||||
}],
|
||||
namespaces: vec![
|
||||
"default".into(),
|
||||
"gitea".into(),
|
||||
"cert-manager".into(),
|
||||
"kube-system".into(),
|
||||
],
|
||||
workloads: vec![
|
||||
w(
|
||||
"cert-manager",
|
||||
WorkloadKind::Deployment,
|
||||
"cert-manager",
|
||||
1,
|
||||
1,
|
||||
c(
|
||||
"cert-manager",
|
||||
"quay.io/jetstack/cert-manager-controller:v1.16.1",
|
||||
),
|
||||
),
|
||||
w(
|
||||
"gitea",
|
||||
WorkloadKind::Deployment,
|
||||
"gitea",
|
||||
1,
|
||||
1,
|
||||
c("gitea", "gitea/gitea:1.22.3"),
|
||||
),
|
||||
w(
|
||||
"gitea",
|
||||
WorkloadKind::StatefulSet,
|
||||
"gitea-postgresql",
|
||||
1,
|
||||
1,
|
||||
c("postgresql", "bitnami/postgresql:16.4.0"),
|
||||
),
|
||||
w(
|
||||
"gitea",
|
||||
WorkloadKind::StatefulSet,
|
||||
"gitea-valkey-primary",
|
||||
1,
|
||||
1,
|
||||
c("valkey", "bitnami/valkey:8.0.1"),
|
||||
),
|
||||
w(
|
||||
"kube-system",
|
||||
WorkloadKind::DaemonSet,
|
||||
"calico-node",
|
||||
1,
|
||||
1,
|
||||
c("calico-node", "docker.io/calico/node:v3.28.1"),
|
||||
),
|
||||
],
|
||||
volume_claims: vec![
|
||||
VolumeClaim {
|
||||
namespace: "gitea".into(),
|
||||
name: "gitea-shared-storage".into(),
|
||||
capacity: "10Gi".into(),
|
||||
storage_class: "microk8s-hostpath".into(),
|
||||
status: "Bound".into(),
|
||||
},
|
||||
VolumeClaim {
|
||||
namespace: "gitea".into(),
|
||||
name: "data-gitea-postgresql-0".into(),
|
||||
capacity: "10Gi".into(),
|
||||
storage_class: "microk8s-hostpath".into(),
|
||||
status: "Bound".into(),
|
||||
},
|
||||
],
|
||||
fetched_at: Utc::now(),
|
||||
})
|
||||
}
|
||||
async fn restart(&self, _: &WorkloadRef) -> Result<(), DomainError> {
|
||||
Ok(())
|
||||
}
|
||||
async fn scale(&self, _: &WorkloadRef, _: i32) -> Result<(), DomainError> {
|
||||
Ok(())
|
||||
}
|
||||
async fn set_image(&self, _: &WorkloadRef, _: &str, _: &str) -> Result<(), DomainError> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@ -2,6 +2,7 @@
|
||||
pub mod cipher;
|
||||
pub mod db;
|
||||
pub mod host;
|
||||
pub mod k8s;
|
||||
pub mod mail;
|
||||
pub mod password;
|
||||
pub mod sqlite;
|
||||
@ -12,6 +13,7 @@ pub use db::{connect, DbPool};
|
||||
pub use host::{
|
||||
DebianInspector, DebianUpdater, FakeHostInspector, FakeHostUpdater, SystemCommandRunner,
|
||||
};
|
||||
pub use k8s::{FakeClusterGateway, KubeGateway};
|
||||
pub use mail::LettreMailer;
|
||||
pub use password::Argon2Hasher;
|
||||
pub use sqlite::SqliteInventory;
|
||||
|
||||
Reference in New Issue
Block a user