use futures::TryStreamExt;
use kube_client::{Api, Resource};
use serde::de::DeserializeOwned;
use snafu::{futures::TryStreamExt as _, Snafu};
use std::fmt::Debug;
use crate::watcher::{self, watch_object};
#[derive(Debug, Snafu)]
pub enum Error {
#[snafu(display("failed to probe for whether the condition is fulfilled yet: {}", source))]
ProbeFailed {
#[snafu(backtrace)]
source: watcher::Error,
},
}
pub async fn await_condition<K>(api: Api<K>, name: &str, cond: impl Condition<K>) -> Result<(), Error>
where
K: Clone + Debug + Send + DeserializeOwned + Resource + 'static,
{
watch_object(api, name)
.context(ProbeFailed)
.try_take_while(|obj| {
let result = !cond.matches_object(obj.as_ref());
async move { Ok(result) }
})
.try_for_each(|_| async { Ok(()) })
.await
}
pub trait Condition<K> {
fn matches_object(&self, obj: Option<&K>) -> bool;
}
impl<K, F: Fn(Option<&K>) -> bool> Condition<K> for F {
fn matches_object(&self, obj: Option<&K>) -> bool {
(self)(obj)
}
}
pub mod conditions {
pub use super::Condition;
use k8s_openapi::apiextensions_apiserver::pkg::apis::apiextensions::v1::CustomResourceDefinition;
use kube_client::Resource;
#[must_use]
pub fn is_deleted<K: Resource>(uid: &str) -> impl Condition<K> + '_ {
move |obj: Option<&K>| {
obj.map_or(
true,
|obj| obj.meta().uid.as_deref() != Some(uid),
)
}
}
#[must_use]
pub fn is_crd_established() -> impl Condition<CustomResourceDefinition> {
|obj: Option<&CustomResourceDefinition>| {
if let Some(o) = obj {
if let Some(s) = &o.status {
if let Some(conds) = &s.conditions {
if let Some(pcond) = conds.iter().find(|c| c.type_ == "Established") {
return pcond.status == "True";
}
}
}
}
false
}
}
}
pub mod delete {
use super::{await_condition, conditions};
use kube_client::{api::DeleteParams, Api, Resource};
use serde::de::DeserializeOwned;
use snafu::{OptionExt, ResultExt, Snafu};
use std::fmt::Debug;
#[derive(Snafu, Debug)]
pub enum Error {
#[snafu(display("deleted object has no UID to wait for"))]
NoUid,
#[snafu(display("failed to delete object: {}", source))]
Delete { source: kube_client::Error },
#[snafu(display("failed to wait for object to be deleted: {}", source))]
Await { source: super::Error },
}
#[allow(clippy::module_name_repetitions)]
pub async fn delete_and_finalize<K: Clone + Debug + Send + DeserializeOwned + Resource + 'static>(
api: Api<K>,
name: &str,
delete_params: &DeleteParams,
) -> Result<(), Error> {
let deleted_obj_uid = api
.delete(name, delete_params)
.await
.context(Delete)?
.either(
|mut obj| obj.meta_mut().uid.take(),
|status| status.details.map(|details| details.uid),
)
.context(NoUid)?;
await_condition(api, name, conditions::is_deleted(&deleted_obj_uid))
.await
.context(Await)
}
}