use async_trait::async_trait;
use k8s_openapi::api::core::v1::ConfigMap;
use k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta;
use kube::Client;
use kube::api::{Api, DeleteParams, ListParams, PostParams};
use std::collections::BTreeMap;
use super::{
CompressionMethod, LargeReleaseStrategy, MAX_RESOURCE_SIZE, StorageConfig, StorageDriver,
decode_from_storage, encode_for_storage, storage_labels,
};
use crate::error::{KubeError, Result};
use crate::release::StoredRelease;
pub struct ConfigMapDriver {
client: Client,
config: StorageConfig,
}
impl ConfigMapDriver {
pub async fn new(config: StorageConfig) -> Result<Self> {
let client = Client::try_default().await?;
Ok(Self { client, config })
}
pub fn with_client(client: Client, config: StorageConfig) -> Self {
Self { client, config }
}
fn configmaps_api(&self, namespace: &str) -> Api<ConfigMap> {
Api::namespaced(self.client.clone(), namespace)
}
fn build_configmap(&self, release: &StoredRelease) -> Result<ConfigMap> {
let encoded = encode_for_storage(release, &self.config)?;
if encoded.len() > MAX_RESOURCE_SIZE {
match &self.config.large_release_strategy {
LargeReleaseStrategy::Fail => {
return Err(KubeError::ReleaseTooLarge {
size: encoded.len(),
max: MAX_RESOURCE_SIZE,
});
}
_ => {
return Err(KubeError::Storage(
"Large release strategies not yet implemented for ConfigMaps".to_string(),
));
}
}
}
let mut labels = storage_labels(release);
labels.insert(
"sherpack.io/storage-driver".to_string(),
"configmap".to_string(),
);
let compression_type = match self.config.compression {
CompressionMethod::None => "none",
CompressionMethod::Gzip { .. } => "gzip",
CompressionMethod::Zstd { .. } => "zstd",
};
labels.insert(
"sherpack.io/compression".to_string(),
compression_type.to_string(),
);
let mut data = BTreeMap::new();
data.insert("release".to_string(), encoded);
Ok(ConfigMap {
metadata: ObjectMeta {
name: Some(release.storage_key()),
namespace: Some(release.namespace.clone()),
labels: Some(labels),
..Default::default()
},
data: Some(data),
..Default::default()
})
}
fn parse_configmap(&self, cm: &ConfigMap) -> Result<StoredRelease> {
let encoded = cm
.data
.as_ref()
.and_then(|d| d.get("release"))
.ok_or_else(|| KubeError::Storage("ConfigMap missing 'release' data".to_string()))?;
let compression = cm
.metadata
.labels
.as_ref()
.and_then(|l| l.get("sherpack.io/compression"))
.map(|c| match c.as_str() {
"none" => CompressionMethod::None,
"gzip" => CompressionMethod::Gzip { level: 6 },
"zstd" => CompressionMethod::Zstd { level: 3 },
_ => self.config.compression,
})
.unwrap_or(self.config.compression);
decode_from_storage(encoded, compression)
}
}
#[async_trait]
impl StorageDriver for ConfigMapDriver {
async fn get(&self, namespace: &str, name: &str, version: u32) -> Result<StoredRelease> {
let api = self.configmaps_api(namespace);
let key = format!("sh.sherpack.release.v1.{}.v{}", name, version);
match api.get(&key).await {
Ok(cm) => self.parse_configmap(&cm),
Err(kube::Error::Api(e)) if e.code == 404 => Err(KubeError::ReleaseNotFound {
name: name.to_string(),
namespace: namespace.to_string(),
}),
Err(e) => Err(e.into()),
}
}
async fn get_latest(&self, namespace: &str, name: &str) -> Result<StoredRelease> {
let history = self.history(namespace, name).await?;
history
.into_iter()
.next()
.ok_or_else(|| KubeError::ReleaseNotFound {
name: name.to_string(),
namespace: namespace.to_string(),
})
}
async fn list(
&self,
namespace: Option<&str>,
name: Option<&str>,
include_superseded: bool,
) -> Result<Vec<StoredRelease>> {
let mut label_selector = "app.kubernetes.io/managed-by=sherpack".to_string();
if let Some(n) = name {
label_selector.push_str(&format!(",sherpack.io/release-name={}", n));
}
let lp = ListParams::default().labels(&label_selector);
let configmaps = if let Some(ns) = namespace {
self.configmaps_api(ns).list(&lp).await?
} else {
let api: Api<ConfigMap> = Api::all(self.client.clone());
api.list(&lp).await?
};
let mut releases: Vec<StoredRelease> = configmaps
.items
.iter()
.filter_map(|cm| self.parse_configmap(cm).ok())
.collect();
releases.sort_by_key(|r| std::cmp::Reverse(r.version));
if !include_superseded {
let mut seen = std::collections::HashSet::new();
releases.retain(|r| {
let key = format!("{}/{}", r.namespace, r.name);
seen.insert(key)
});
}
Ok(releases)
}
async fn history(&self, namespace: &str, name: &str) -> Result<Vec<StoredRelease>> {
let label_selector = format!(
"app.kubernetes.io/managed-by=sherpack,sherpack.io/release-name={}",
name
);
let lp = ListParams::default().labels(&label_selector);
let configmaps = self.configmaps_api(namespace).list(&lp).await?;
let mut releases: Vec<StoredRelease> = configmaps
.items
.iter()
.filter_map(|cm| self.parse_configmap(cm).ok())
.collect();
releases.sort_by_key(|r| std::cmp::Reverse(r.version));
if releases.is_empty() {
return Err(KubeError::ReleaseNotFound {
name: name.to_string(),
namespace: namespace.to_string(),
});
}
Ok(releases)
}
async fn create(&self, release: &StoredRelease) -> Result<()> {
let api = self.configmaps_api(&release.namespace);
let cm = self.build_configmap(release)?;
match api.get(&release.storage_key()).await {
Ok(_) => {
return Err(KubeError::ReleaseAlreadyExists {
name: release.name.clone(),
namespace: release.namespace.clone(),
});
}
Err(kube::Error::Api(e)) if e.code == 404 => {}
Err(e) => return Err(e.into()),
}
api.create(&PostParams::default(), &cm).await?;
Ok(())
}
async fn update(&self, release: &StoredRelease) -> Result<()> {
let api = self.configmaps_api(&release.namespace);
let cm = self.build_configmap(release)?;
api.replace(&release.storage_key(), &PostParams::default(), &cm)
.await?;
Ok(())
}
async fn delete(&self, namespace: &str, name: &str, version: u32) -> Result<StoredRelease> {
let release = self.get(namespace, name, version).await?;
let api = self.configmaps_api(namespace);
let key = format!("sh.sherpack.release.v1.{}.v{}", name, version);
api.delete(&key, &DeleteParams::default()).await?;
Ok(release)
}
async fn delete_all(&self, namespace: &str, name: &str) -> Result<Vec<StoredRelease>> {
let releases = self.history(namespace, name).await?;
let api = self.configmaps_api(namespace);
for release in &releases {
let _ = api
.delete(&release.storage_key(), &DeleteParams::default())
.await;
}
Ok(releases)
}
}