use std::sync::Arc;
use sherpack_core::{LoadedPack, ReleaseInfo, TemplateContext, Values};
use sherpack_engine::Engine;
use sherpack_engine::cluster_reader::ClusterReader;
use crate::actions::{InstallOptions, RollbackOptions, UninstallOptions, UpgradeOptions};
use crate::diff::{DiffEngine, DiffResult};
use crate::error::{KubeError, Result};
use crate::health::{HealthCheckConfig, HealthChecker, HealthStatus};
use crate::hooks::{HookExecutor, HookPhase, parse_hooks_from_manifest};
use crate::lookup::KubeClusterReader;
use crate::release::{ReleaseState, StoredRelease};
use crate::resources::ResourceManager;
use crate::storage::StorageDriver;
pub struct KubeClient<S: StorageDriver> {
client: kube::Client,
storage: S,
diff_engine: DiffEngine,
}
impl<S: StorageDriver> KubeClient<S> {
pub async fn new(storage: S) -> Result<Self> {
let client = kube::Client::try_default().await?;
let diff_engine = DiffEngine::new();
Ok(Self {
client,
storage,
diff_engine,
})
}
pub fn with_client(client: kube::Client, storage: S) -> Self {
let diff_engine = DiffEngine::new();
Self {
client,
storage,
diff_engine,
}
}
pub fn kube_client(&self) -> &kube::Client {
&self.client
}
pub fn storage(&self) -> &S {
&self.storage
}
async fn engine_with_lookup(&self) -> Engine {
match KubeClusterReader::new(self.client.clone()).await {
Ok(mut reader) => {
if let Ok(secs) = std::env::var("SHERPACK_LOOKUP_TIMEOUT_SECS")
&& let Ok(parsed) = secs.parse::<u64>()
&& parsed > 0
{
reader = reader.with_timeout(std::time::Duration::from_secs(parsed));
}
let arc: Arc<dyn ClusterReader> = Arc::new(reader);
Engine::builder()
.strict(true)
.with_cluster_reader(arc)
.build()
}
Err(e) => {
tracing::warn!(
"Cluster discovery for lookup() failed; lookup() will return empty: {}",
e
);
Engine::builder().strict(true).build()
}
}
}
fn surface_lookup_warnings(engine: &Engine) {
if let Some(state) = engine.lookup_state() {
for w in state.take_warnings() {
tracing::warn!("{}", w);
}
}
}
pub async fn install(
&self,
pack: &LoadedPack,
values: Values,
options: &InstallOptions,
) -> Result<StoredRelease> {
if self
.storage
.exists(&options.namespace, &options.name)
.await?
{
return Err(KubeError::ReleaseAlreadyExists {
name: options.name.clone(),
namespace: options.namespace.clone(),
});
}
let release_info = ReleaseInfo::for_install(&options.name, &options.namespace);
let context = TemplateContext::new(values.clone(), release_info, &pack.pack.metadata);
let engine = self.engine_with_lookup().await;
let render_result = engine
.render_pack(pack, &context)
.map_err(|e| KubeError::Template(e.to_string()))?;
Self::surface_lookup_warnings(&engine);
let mut release = StoredRelease::for_install(
options.name.clone(),
options.namespace.clone(),
pack.pack.metadata.clone(),
values,
render_result
.manifests
.values()
.cloned()
.collect::<Vec<_>>()
.join("\n---\n"),
);
release.notes = render_result.notes;
release.hooks = parse_hooks_from_manifest(&release.manifest);
for (k, v) in &options.labels {
release.labels.insert(k.clone(), v.clone());
}
if options.dry_run {
return Ok(release);
}
if options.show_diff {
println!("Resources to be created:");
for manifest in render_result.manifests.keys() {
println!(" + {}", manifest);
}
}
self.storage.create(&release).await?;
let mut hook_executor = HookExecutor::new();
if let Err(e) = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PreInstall,
&release.name,
release.version,
&self.client,
)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
return Err(e);
}
if let Err(e) = self
.apply_manifest(&release.namespace, &release.manifest)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
if options.atomic {
let _ = self.cleanup_release(&release).await;
}
return Err(e);
}
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::DuringInstall,
&release.name,
release.version,
&self.client,
)
.await;
if options.wait {
let _timeout = options.timeout.unwrap_or(chrono::Duration::minutes(5));
let health_config = options.health_check.clone().unwrap_or_default();
let checker = HealthChecker::new(health_config);
let status = checker.check(&release, &self.client).await?;
if !status.healthy {
let err_msg = status.summary();
release.mark_failed(err_msg.clone(), true);
self.storage.update(&release).await?;
if options.atomic {
let _ = self.cleanup_release(&release).await;
}
return Err(KubeError::HealthCheckFailed {
name: release.name.clone(),
message: err_msg,
});
}
}
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PostInstall,
&release.name,
release.version,
&self.client,
)
.await;
release.mark_deployed();
self.storage.update(&release).await?;
Ok(release)
}
pub async fn upgrade(
&self,
pack: &LoadedPack,
values: Values,
options: &UpgradeOptions,
) -> Result<StoredRelease> {
let existing = match self
.storage
.get_latest(&options.namespace, &options.name)
.await
{
Ok(r) => Some(r),
Err(KubeError::ReleaseNotFound { .. }) if options.install => None,
Err(e) => return Err(e),
};
if existing.is_none() {
let install_opts = InstallOptions {
name: options.name.clone(),
namespace: options.namespace.clone(),
wait: options.wait,
timeout: options.timeout,
health_check: options.health_check.clone(),
atomic: options.atomic,
dry_run: options.dry_run,
show_diff: options.show_diff,
labels: options.labels.clone(),
description: options.description.clone(),
..Default::default()
};
return self.install(pack, values, &install_opts).await;
}
let existing = existing.unwrap();
if existing.state.is_pending() {
if existing.is_stuck() {
return Err(KubeError::StuckRelease {
name: existing.name.clone(),
status: existing.state.status_name().to_string(),
elapsed: existing
.state
.elapsed()
.map(|d| format!("{} seconds", d.num_seconds()))
.unwrap_or_else(|| "unknown".to_string()),
});
} else {
return Err(KubeError::OperationInProgress {
name: existing.name.clone(),
status: existing.state.to_string(),
});
}
}
let final_values = if options.reset_values {
values
} else if options.reuse_values {
let mut merged = existing.values.clone();
merged.merge(&values);
merged
} else {
values
};
let release_info =
ReleaseInfo::for_upgrade(&options.name, &options.namespace, existing.version + 1);
let context = TemplateContext::new(final_values.clone(), release_info, &pack.pack.metadata);
let engine = self.engine_with_lookup().await;
let render_result = engine
.render_pack(pack, &context)
.map_err(|e| KubeError::Template(e.to_string()))?;
Self::surface_lookup_warnings(&engine);
let manifest = render_result
.manifests
.values()
.cloned()
.collect::<Vec<_>>()
.join("\n---\n");
let mut release = StoredRelease::for_upgrade(&existing, final_values, manifest);
release.notes = render_result.notes;
release.hooks = parse_hooks_from_manifest(&release.manifest);
for (k, v) in &options.labels {
release.labels.insert(k.clone(), v.clone());
}
if options.show_diff {
let diff = self.diff_engine.diff_releases(&existing, &release);
println!("Changes: {}", self.diff_engine.summary(&diff));
}
if options.dry_run {
return Ok(release);
}
self.storage.create(&release).await?;
let mut prev = existing;
prev.mark_superseded();
self.storage.update(&prev).await?;
let mut hook_executor = HookExecutor::new();
if !options.no_hooks
&& let Err(e) = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PreUpgrade,
&release.name,
release.version,
&self.client,
)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
if options.atomic {
return self.rollback_to(&release, prev.version).await;
}
return Err(e);
}
if let Err(e) = self
.apply_manifest(&release.namespace, &release.manifest)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
if options.atomic {
return self.rollback_to(&release, prev.version).await;
}
return Err(e);
}
if !options.no_hooks {
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::DuringUpgrade,
&release.name,
release.version,
&self.client,
)
.await;
}
if options.wait {
let health_config = options.health_check.clone().unwrap_or_default();
let checker = HealthChecker::new(health_config);
let status = checker.check(&release, &self.client).await?;
if !status.healthy {
let err_msg = status.summary();
release.mark_failed(err_msg.clone(), true);
self.storage.update(&release).await?;
if options.atomic {
return self.rollback_to(&release, prev.version).await;
}
return Err(KubeError::HealthCheckFailed {
name: release.name.clone(),
message: err_msg,
});
}
}
if !options.no_hooks {
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PostUpgrade,
&release.name,
release.version,
&self.client,
)
.await;
}
release.mark_deployed();
self.storage.update(&release).await?;
if let Some(max_history) = options.max_history {
self.cleanup_history(&release.namespace, &release.name, max_history)
.await?;
}
Ok(release)
}
pub async fn uninstall(&self, options: &UninstallOptions) -> Result<StoredRelease> {
let mut release = self
.storage
.get_latest(&options.namespace, &options.name)
.await?;
release.state = ReleaseState::PendingUninstall {
started_at: chrono::Utc::now(),
timeout: options.timeout.unwrap_or(chrono::Duration::minutes(5)),
};
self.storage.update(&release).await?;
if options.dry_run {
return Ok(release);
}
let mut hook_executor = HookExecutor::new();
if !options.no_hooks {
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PreDelete,
&release.name,
release.version,
&self.client,
)
.await;
}
if let Err(e) = self
.delete_manifest(&release.namespace, &release.manifest)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
return Err(e);
}
if !options.no_hooks {
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PostDelete,
&release.name,
release.version,
&self.client,
)
.await;
}
release.mark_uninstalled();
self.storage.update(&release).await?;
if !options.keep_history {
self.storage
.delete_all(&options.namespace, &options.name)
.await?;
}
Ok(release)
}
pub async fn rollback(&self, options: &RollbackOptions) -> Result<StoredRelease> {
let history = self
.storage
.history(&options.namespace, &options.name)
.await?;
if history.is_empty() {
return Err(KubeError::ReleaseNotFound {
name: options.name.clone(),
namespace: options.namespace.clone(),
});
}
let target_version = if options.revision == 0 {
if history.len() < 2 {
return Err(KubeError::RollbackNotPossible {
name: options.name.clone(),
reason: "no previous revision available".to_string(),
});
}
history[1].version
} else {
options.revision
};
let target = history
.iter()
.find(|r| r.version == target_version)
.ok_or_else(|| KubeError::RollbackNotPossible {
name: options.name.clone(),
reason: format!("revision {} not found", target_version),
})?;
let current = &history[0];
if options.show_diff {
let diff = self.diff_engine.diff_releases(current, target);
println!("Rollback changes: {}", self.diff_engine.summary(&diff));
}
if options.dry_run {
return Ok(target.clone());
}
let mut release =
StoredRelease::for_upgrade(current, target.values.clone(), target.manifest.clone());
release.state = ReleaseState::PendingRollback {
started_at: chrono::Utc::now(),
timeout: options.timeout.unwrap_or(chrono::Duration::minutes(5)),
target_version,
};
self.storage.create(&release).await?;
let mut prev = current.clone();
prev.mark_superseded();
self.storage.update(&prev).await?;
let mut hook_executor = HookExecutor::new();
if !options.no_hooks
&& let Err(e) = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PreRollback,
&release.name,
release.version,
&self.client,
)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
return Err(e);
}
if let Err(e) = self
.apply_manifest(&release.namespace, &release.manifest)
.await
{
release.mark_failed(e.to_string(), true);
self.storage.update(&release).await?;
return Err(e);
}
if options.wait {
let health_config = options.health_check.clone().unwrap_or_default();
let checker = HealthChecker::new(health_config);
let status = checker.check(&release, &self.client).await?;
if !status.healthy {
let err_msg = status.summary();
release.mark_failed(err_msg.clone(), true);
self.storage.update(&release).await?;
return Err(KubeError::HealthCheckFailed {
name: release.name.clone(),
message: err_msg,
});
}
}
if !options.no_hooks {
let _ = hook_executor
.execute_phase(
&release.hooks,
HookPhase::PostRollback,
&release.name,
release.version,
&self.client,
)
.await;
}
release.mark_deployed();
self.storage.update(&release).await?;
if let Some(max_history) = options.max_history {
self.cleanup_history(&release.namespace, &release.name, max_history)
.await?;
}
Ok(release)
}
pub async fn list(
&self,
namespace: Option<&str>,
all_namespaces: bool,
) -> Result<Vec<StoredRelease>> {
let ns = if all_namespaces { None } else { namespace };
self.storage.list(ns, None, false).await
}
pub async fn history(&self, namespace: &str, name: &str) -> Result<Vec<StoredRelease>> {
self.storage.history(namespace, name).await
}
pub async fn status(&self, namespace: &str, name: &str) -> Result<StoredRelease> {
self.storage.get_latest(namespace, name).await
}
pub async fn health(
&self,
namespace: &str,
name: &str,
config: Option<HealthCheckConfig>,
) -> Result<HealthStatus> {
let release = self.storage.get_latest(namespace, name).await?;
let checker = HealthChecker::new(config.unwrap_or_default());
checker.check_once(&release, &self.client).await
}
pub async fn diff(
&self,
namespace: &str,
name: &str,
revision1: u32,
revision2: u32,
) -> Result<DiffResult> {
let r1 = self.storage.get(namespace, name, revision1).await?;
let r2 = self.storage.get(namespace, name, revision2).await?;
Ok(self.diff_engine.diff_releases(&r1, &r2))
}
pub async fn recover(&self, namespace: &str, name: &str) -> Result<StoredRelease> {
let mut release = self.storage.get_latest(namespace, name).await?;
if !release.state.is_pending() {
return Err(KubeError::InvalidConfig(format!(
"release '{}' is not in a pending state",
name
)));
}
release.mark_failed("Manually recovered from stuck state".to_string(), true);
self.storage.update(&release).await?;
Ok(release)
}
async fn resource_manager(&self) -> Result<ResourceManager> {
ResourceManager::new(self.client.clone()).await
}
async fn apply_manifest(&self, namespace: &str, manifest: &str) -> Result<()> {
let manager = self.resource_manager().await?;
let summary = manager.apply_manifest(namespace, manifest, false).await?;
if !summary.is_success() {
let errors: Vec<String> = summary
.failed
.iter()
.map(|(name, err)| format!("{}: {}", name, err))
.collect();
return Err(KubeError::InvalidConfig(format!(
"Failed to apply resources: {}",
errors.join("; ")
)));
}
Ok(())
}
#[allow(dead_code)]
async fn apply_manifest_dry_run(
&self,
namespace: &str,
manifest: &str,
) -> Result<crate::resources::OperationSummary> {
let manager = self.resource_manager().await?;
manager.apply_manifest(namespace, manifest, true).await
}
async fn delete_manifest(&self, namespace: &str, manifest: &str) -> Result<()> {
let manager = self.resource_manager().await?;
let summary = manager.delete_manifest(namespace, manifest, false).await?;
if !summary.is_success() {
let errors: Vec<String> = summary
.failed
.iter()
.map(|(name, err)| format!("{}: {}", name, err))
.collect();
return Err(KubeError::InvalidConfig(format!(
"Failed to delete resources: {}",
errors.join("; ")
)));
}
Ok(())
}
async fn cleanup_release(&self, release: &StoredRelease) -> Result<()> {
self.delete_manifest(&release.namespace, &release.manifest)
.await
}
async fn rollback_to(
&self,
current: &StoredRelease,
target_version: u32,
) -> Result<StoredRelease> {
let _target = self
.storage
.get(¤t.namespace, ¤t.name, target_version)
.await?;
let options = RollbackOptions {
name: current.name.clone(),
namespace: current.namespace.clone(),
revision: target_version,
wait: true,
timeout: Some(chrono::Duration::minutes(5)),
..Default::default()
};
self.rollback(&options).await
}
async fn cleanup_history(&self, namespace: &str, name: &str, max_history: u32) -> Result<()> {
let history = self.storage.history(namespace, name).await?;
if history.len() as u32 <= max_history {
return Ok(());
}
for release in history.iter().skip(max_history as usize) {
self.storage
.delete(namespace, name, release.version)
.await?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
}