use anyhow::{Context, Result};
use colored::*;
use indicatif::{ProgressBar, ProgressStyle};
use jsonnet::JsonnetVm;
use k8s_openapi::api::core::v1::ConfigMap;
use kube::{
api::{Api, DynamicObject, GroupVersionKind, Patch, PatchParams, PostParams},
discovery::{self, Scope},
Client,
};
use serde_json::Value;
use std::{collections::BTreeMap, path::PathBuf, time::Duration};
use tabled::Tabled;
#[cfg(test)]
use mockall::automock;
pub const SET_LABEL_KEY: &str = "nethermit.set";
pub const SETS_CONFIGMAP_NAME: &str = "nethermit-sets";
#[derive(Debug, serde::Serialize, serde::Deserialize, Clone)]
pub struct ResourceInfo {
pub api_version: String,
pub kind: String,
pub name: String,
pub namespace: Option<String>,
}
#[derive(Tabled)]
pub struct SetInfo {
#[tabled(rename = "Set Name")]
pub name: String,
#[tabled(rename = "Resources")]
pub resource_count: usize,
}
#[cfg_attr(test, automock)]
#[async_trait::async_trait]
pub trait KubeClient {
async fn get_configmap(&self, name: &str, namespace: &str) -> Result<Option<ConfigMap>>;
async fn create_configmap(&self, cm: &ConfigMap, namespace: &str) -> Result<()>;
async fn patch_configmap<'a>(
&'a self,
name: &'a str,
namespace: &'a str,
patch: &'a Patch<&'a Value>,
) -> Result<()>;
async fn delete_configmap(&self, name: &str, namespace: &str) -> Result<()>;
async fn create_pod(&self, pod: &Value, namespace: &str) -> Result<()>;
async fn delete_pod(&self, name: &str, namespace: &str) -> Result<()>;
}
pub struct RealKubeClient {
client: Client,
}
#[async_trait::async_trait]
impl KubeClient for RealKubeClient {
async fn get_configmap(&self, name: &str, namespace: &str) -> Result<Option<ConfigMap>> {
let api: Api<ConfigMap> = Api::namespaced(self.client.clone(), namespace);
match api.get(name).await {
Ok(cm) => Ok(Some(cm)),
Err(kube::Error::Api(err)) if err.code == 404 => Ok(None),
Err(e) => Err(e.into()),
}
}
async fn create_configmap(&self, cm: &ConfigMap, namespace: &str) -> Result<()> {
let api: Api<ConfigMap> = Api::namespaced(self.client.clone(), namespace);
api.create(&PostParams::default(), cm).await?;
Ok(())
}
async fn patch_configmap<'a>(
&'a self,
name: &'a str,
namespace: &'a str,
patch: &'a Patch<&'a Value>,
) -> Result<()> {
let api: Api<ConfigMap> = Api::namespaced(self.client.clone(), namespace);
api.patch(name, &PatchParams::apply("nethermit"), patch)
.await?;
Ok(())
}
async fn delete_configmap(&self, name: &str, namespace: &str) -> Result<()> {
let api: Api<ConfigMap> = Api::namespaced(self.client.clone(), namespace);
api.delete(name, &Default::default()).await?;
Ok(())
}
async fn create_pod(&self, pod: &Value, namespace: &str) -> Result<()> {
let api: Api<DynamicObject> = Api::namespaced_with(
self.client.clone(),
namespace,
&kube::core::ApiResource::erase::<k8s_openapi::api::core::v1::Pod>(&()),
);
api.create(
&PostParams::default(),
&serde_json::from_value(pod.clone())?,
)
.await?;
Ok(())
}
async fn delete_pod(&self, name: &str, namespace: &str) -> Result<()> {
let api: Api<DynamicObject> = Api::namespaced_with(
self.client.clone(),
namespace,
&kube::core::ApiResource::erase::<k8s_openapi::api::core::v1::Pod>(&()),
);
api.delete(name, &Default::default()).await?;
Ok(())
}
}
impl RealKubeClient {
pub async fn new() -> Result<Self> {
Ok(Self {
client: Client::try_default().await?,
})
}
}
pub fn create_spinner(msg: &str) -> ProgressBar {
let pb = ProgressBar::new_spinner();
pb.enable_steady_tick(Duration::from_millis(100));
pb.set_style(
ProgressStyle::default_spinner()
.tick_chars("⠁⠂⠄⡀⢀⠠⠐⠈")
.template("{spinner:.blue} {msg}")
.unwrap(),
);
pb.set_message(msg.to_string());
pb
}
pub fn format_success(msg: &str) -> String {
format!("✓ {msg}").green().to_string()
}
pub fn format_error(msg: &str) -> String {
format!("✗ {msg}").red().to_string()
}
pub fn format_info(msg: &str) -> String {
format!("ℹ {msg}").blue().to_string()
}
pub async fn set_exists(client: &impl KubeClient, name: &str) -> Result<bool> {
match client.get_configmap(SETS_CONFIGMAP_NAME, "default").await? {
Some(cm) => Ok(cm
.data
.as_ref()
.map(|data| data.contains_key(&format!("{name}.resources")))
.unwrap_or(false)),
None => Ok(false),
}
}
pub async fn list_sets_info(client: &impl KubeClient) -> Result<Vec<SetInfo>> {
match client.get_configmap(SETS_CONFIGMAP_NAME, "default").await? {
Some(cm) => {
let data = cm.data.unwrap_or_default();
Ok(data
.iter()
.filter(|(k, _)| k.ends_with(".resources"))
.map(|(k, v)| {
let name = k.trim_end_matches(".resources").to_string();
let resources: Vec<ResourceInfo> = serde_json::from_str(v).unwrap_or_default();
SetInfo {
name,
resource_count: resources.len(),
}
})
.collect())
}
None => Ok(Vec::new()),
}
}
pub async fn add_set(
client: &impl KubeClient,
name: &str,
resources: Vec<ResourceInfo>,
) -> Result<()> {
if set_exists(client, name).await? {
return Err(anyhow::anyhow!("Set '{}' already exists", name));
}
let mut data = BTreeMap::new();
data.insert(
format!("{name}.resources"),
serde_json::to_string(&resources)?,
);
let patch_json = serde_json::json!({
"apiVersion": "v1",
"kind": "ConfigMap",
"metadata": {
"name": SETS_CONFIGMAP_NAME,
},
"data": data
});
let patch = Patch::Apply(&patch_json);
client
.patch_configmap(SETS_CONFIGMAP_NAME, "default", &patch)
.await?;
Ok(())
}
pub async fn remove_set(client: &impl KubeClient, name: &str) -> Result<()> {
if !set_exists(client, name).await? {
return Err(anyhow::anyhow!("Set '{}' does not exist", name));
}
match client.get_configmap(SETS_CONFIGMAP_NAME, "default").await? {
Some(cm) => {
let mut data = cm.data.unwrap_or_default();
data.remove(&format!("{name}.resources"));
let patch_json = serde_json::json!({
"apiVersion": "v1",
"kind": "ConfigMap",
"metadata": {
"name": SETS_CONFIGMAP_NAME,
},
"data": data
});
let patch = Patch::Apply(&patch_json);
client
.patch_configmap(SETS_CONFIGMAP_NAME, "default", &patch)
.await?;
Ok(())
}
None => Err(anyhow::anyhow!("ConfigMap not found")),
}
}
pub fn jsonnet_to_yaml(content: &str, import_paths: &[PathBuf]) -> Result<String> {
let mut vm = JsonnetVm::new();
vm.max_stack(100);
vm.max_trace(Some(100));
for path in import_paths {
vm.jpath_add(path.to_string_lossy().as_ref());
}
let json = vm
.evaluate_snippet("snippet", content)
.map_err(|e| anyhow::anyhow!("Failed to evaluate Jsonnet: {}", e))?;
let value: Value = serde_json::from_str(&json).context("Failed to parse JSON")?;
let yaml = match value {
Value::Array(resources) => {
resources
.into_iter()
.map(|r| serde_yaml::to_string(&r).context("Failed to convert to YAML"))
.collect::<Result<Vec<_>>>()?
.join("\n---\n")
}
_ => serde_yaml::to_string(&value).context("Failed to convert to YAML")?,
};
Ok(yaml)
}
pub async fn apply_kubernetes(yaml: &str, set_name: Option<&str>) -> Result<()> {
let client = RealKubeClient::new().await?;
let kube_client = client.client.clone();
let documents: Vec<&str> = yaml
.split("\n---\n")
.filter(|s| !s.trim().is_empty())
.collect();
let pb = create_spinner("Applying resources");
let mut resources = Vec::new();
for (i, doc) in documents.iter().enumerate() {
let mut obj: DynamicObject =
serde_yaml::from_str(doc).context("Failed to parse YAML document")?;
let type_meta = obj.types.clone().context("Missing type metadata")?;
if let Some(name) = set_name {
let labels = obj.metadata.labels.get_or_insert_with(BTreeMap::new);
labels.insert(SET_LABEL_KEY.to_string(), name.to_string());
}
let gvk = GroupVersionKind::try_from(&type_meta)?;
let (api_resource, caps) = discovery::pinned_kind(&kube_client, &gvk).await?;
let api: Api<DynamicObject> = match caps.scope {
Scope::Cluster => Api::all_with(kube_client.clone(), &api_resource),
Scope::Namespaced => {
let namespace = obj.metadata.namespace.as_deref().unwrap_or("default");
Api::namespaced_with(kube_client.clone(), namespace, &api_resource)
}
};
let name = obj
.metadata
.name
.as_deref()
.context("Resource must have a name")?;
pb.set_message(format!(
"Applying {} '{}' ({}/{})",
type_meta.kind,
name,
i + 1,
documents.len()
));
api.patch(name, &PatchParams::apply("nethermit"), &Patch::Apply(&obj))
.await
.with_context(|| format!("Failed to apply resource {name}"))?;
if let Some(_set_name) = set_name {
resources.push(ResourceInfo {
api_version: type_meta.api_version,
kind: type_meta.kind,
name: name.to_string(),
namespace: obj.metadata.namespace,
});
}
}
pb.finish_with_message(format_success(&format!(
"Applied {} resources",
documents.len()
)));
if let Some(name) = set_name {
add_set(&client, name, resources).await?;
}
Ok(())
}
pub async fn delete_set_resources(client: &impl KubeClient, set_name: &str) -> Result<()> {
let cm = match client.get_configmap(SETS_CONFIGMAP_NAME, "default").await? {
Some(cm) => cm,
None => return Err(anyhow::anyhow!("ConfigMap not found")),
};
let resources: Vec<ResourceInfo> = if let Some(data) = cm.data {
if let Some(resources_str) = data.get(&format!("{set_name}.resources")) {
serde_json::from_str(resources_str)?
} else {
return Err(anyhow::anyhow!("Set '{}' not found", set_name));
}
} else {
return Err(anyhow::anyhow!("No sets found"));
};
let pb = create_spinner("Deleting resources");
for (i, resource) in resources.iter().enumerate() {
pb.set_message(format!(
"Deleting {} '{}' ({}/{})",
resource.kind,
resource.name,
i + 1,
resources.len()
));
if resource.api_version == "v1" && resource.kind == "Pod" {
if let Some(ref ns) = resource.namespace {
client.delete_pod(&resource.name, ns).await?;
}
} else {
return Err(anyhow::anyhow!(
"Unsupported resource type: {}",
resource.kind
));
}
}
pb.finish_with_message(format_success(&format!(
"Deleted {} resources",
resources.len()
)));
remove_set(client, set_name).await?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use mockall::predicate::*;
#[tokio::test]
async fn test_set_management() -> Result<()> {
let mut mock_client = MockKubeClient::new();
let set_name = "test-set";
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(|_, _| Ok(None));
mock_client
.expect_patch_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"), always())
.times(1)
.returning(|_, _, _| Ok(()));
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{name}.resources", name = set_name),
"[]".to_string(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{name}.resources", name = set_name),
"[]".to_string(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{name}.resources", name = set_name),
"[]".to_string(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{name}.resources", name = set_name),
"[]".to_string(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_patch_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"), always())
.times(1)
.returning(|_, _, _| Ok(()));
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(|_, _| {
Ok(Some(ConfigMap {
data: Some(BTreeMap::new()),
..Default::default()
}))
});
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(|_, _| {
Ok(Some(ConfigMap {
data: Some(BTreeMap::new()),
..Default::default()
}))
});
add_set(&mock_client, set_name, Vec::new()).await?;
assert!(set_exists(&mock_client, set_name).await?);
let sets = list_sets_info(&mock_client).await?;
assert!(sets.iter().any(|s| s.name == set_name));
remove_set(&mock_client, set_name).await?;
assert!(!set_exists(&mock_client, set_name).await?);
let sets = list_sets_info(&mock_client).await?;
assert!(!sets.iter().any(|s| s.name == set_name));
Ok(())
}
#[tokio::test]
async fn test_resource_labels() -> Result<()> {
let set_name = "test-label-set";
let mut mock_client = MockKubeClient::new();
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.returning(|_, _| Ok(None));
mock_client
.expect_patch_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"), always())
.returning(|_, _, _| Ok(()));
let pod = serde_json::json!({
"apiVersion": "v1",
"kind": "Pod",
"metadata": {
"name": "test-pod",
"namespace": "default",
"labels": {
SET_LABEL_KEY: set_name
}
},
"spec": {
"containers": [{
"name": "test",
"image": "nginx"
}]
}
});
mock_client
.expect_create_pod()
.with(always(), eq("default"))
.returning(|_, _| Ok(()));
mock_client.create_pod(&pod, "default").await?;
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.returning(|_, _| Ok(None));
let resources = vec![ResourceInfo {
api_version: "v1".to_string(),
kind: "Pod".to_string(),
name: "test-pod".to_string(),
namespace: Some("default".to_string()),
}];
add_set(&mock_client, set_name, resources).await?;
Ok(())
}
#[tokio::test]
async fn test_delete_set_resources() -> Result<()> {
let mut mock_client = MockKubeClient::new();
let set_name = "test-delete-set";
let pod_name = "test-delete-pod";
let pod = serde_json::json!({
"apiVersion": "v1",
"kind": "Pod",
"metadata": {
"name": pod_name,
"namespace": "default"
},
"spec": {
"containers": [{
"name": "test",
"image": "nginx"
}]
}
});
let resources = vec![ResourceInfo {
api_version: "v1".to_string(),
kind: "Pod".to_string(),
name: pod_name.to_string(),
namespace: Some("default".to_string()),
}];
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(|_, _| Ok(None));
mock_client
.expect_patch_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"), always())
.times(1)
.returning(|_, _, _| Ok(()));
mock_client
.expect_create_pod()
.with(always(), eq("default"))
.times(1)
.returning(|_, _| Ok(()));
let resources_clone = resources.clone();
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{set_name}.resources"),
serde_json::to_string(&resources_clone).unwrap(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_delete_pod()
.with(eq(pod_name), eq("default"))
.times(1)
.returning(|_, _| Ok(()));
let resources_clone = resources.clone();
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{set_name}.resources"),
serde_json::to_string(&resources_clone).unwrap(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
let resources_clone = resources.clone();
mock_client
.expect_get_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"))
.times(1)
.returning(move |_, _| {
let mut data = BTreeMap::new();
data.insert(
format!("{set_name}.resources"),
serde_json::to_string(&resources_clone).unwrap(),
);
Ok(Some(ConfigMap {
data: Some(data),
..Default::default()
}))
});
mock_client
.expect_patch_configmap()
.with(eq(SETS_CONFIGMAP_NAME), eq("default"), always())
.times(1)
.returning(|_, _, _| Ok(()));
mock_client.create_pod(&pod, "default").await?;
add_set(&mock_client, set_name, resources).await?;
delete_set_resources(&mock_client, set_name).await?;
Ok(())
}
}