use std::collections::BTreeMap;
use std::time::Duration;
use futures_util::future::join_all;
use k8s_openapi::api::batch::v1::Job;
use k8s_openapi::api::core::v1::{ConfigMap, Pod};
use kube::api::{Api, DeleteParams, ListParams};
use kube::runtime::wait::{await_condition, conditions};
use kube::{Client, Error as KubeError};
use tokio::time;
use tracing::{info, warn};
use crate::error::AgentError;
use crate::provider::{LABEL_ROOT_RUN_ID, LABEL_RUN_ID, LABEL_STEP, sanitize_label_value};
use super::common::{K8sClusterConfig, create_client};
use super::ephemeral::{MANAGED_SELECTOR, PROMPT_SELECTOR, RUNNER_SELECTOR};
pub(super) struct Selection {
pods: Vec<String>,
jobs: Vec<String>,
configmaps: Vec<String>,
what: String,
}
pub(super) fn step_selection(run_id: &str, step: &str) -> Selection {
let scope = format!("{LABEL_RUN_ID}={run_id},{LABEL_STEP}={step}");
Selection {
pods: vec![format!("{RUNNER_SELECTOR},{scope}")],
jobs: Vec::new(),
configmaps: vec![format!("{PROMPT_SELECTOR},{scope}")],
what: format!("the previous attempt of step {step} of run {run_id}"),
}
}
fn run_selection(run_id: &str) -> Selection {
let id = sanitize_label_value(run_id);
let scopes = [
format!("{LABEL_RUN_ID}={id}"),
format!("{LABEL_ROOT_RUN_ID}={id}"),
];
let managed = scopes.iter().map(|s| format!("{MANAGED_SELECTOR},{s}"));
Selection {
pods: managed.clone().collect(),
jobs: managed.collect(),
configmaps: scopes
.iter()
.map(|s| format!("{PROMPT_SELECTOR},{s}"))
.collect(),
what: format!("run {run_id}"),
}
}
pub(super) async fn release_run(
cluster_config: &K8sClusterConfig,
namespace: &str,
run_id: &str,
limit: Duration,
) -> Result<(), AgentError> {
let client = create_client(cluster_config).await?;
let deleted = delete_and_wait(&client, namespace, &run_selection(run_id), limit).await?;
if deleted > 0 {
info!(run_id = %run_id, pods_deleted = deleted, "released the pods of a previous execution");
}
Ok(())
}
fn failure(stderr: String) -> AgentError {
AgentError::ProcessFailed {
exit_code: -1,
stderr,
}
}
pub(super) async fn delete_and_wait(
client: &Client,
namespace: &str,
selection: &Selection,
limit: Duration,
) -> Result<usize, AgentError> {
let what = &selection.what;
delete_jobs(&Api::namespaced(client.clone(), namespace), selection).await?;
let pods: Api<Pod> = Api::namespaced(client.clone(), namespace);
let mut listed: BTreeMap<String, String> = BTreeMap::new();
for selector in &selection.pods {
let list = pods.list(&ListParams::default().labels(selector)).await;
let list = list.map_err(|e| failure(format!("failed to list pods of {what}: {e}")))?;
for pod in list.items {
if let (Some(name), Some(uid)) = (pod.metadata.name, pod.metadata.uid) {
listed.insert(uid, name);
}
}
}
let mut deleted: Vec<(String, String)> = Vec::new();
for (uid, name) in listed {
match pods.delete(&name, &DeleteParams::default()).await {
Ok(_) => deleted.push((name, uid)),
Err(KubeError::Api(e)) if e.code == 404 => {}
Err(e) => {
return Err(failure(format!(
"failed to delete pod '{name}' of {what}: {e}"
)));
}
}
}
let waits = deleted
.iter()
.map(|(name, uid)| await_condition(pods.clone(), name, conditions::is_deleted(uid)));
let waited = time::timeout(limit, join_all(waits)).await;
let results = waited.map_err(|e| {
failure(format!(
"pods of {what} still terminating after {limit:?}: {e}"
))
})?;
for result in results {
result.map_err(|e| failure(format!("failed waiting for the pods of {what}: {e}")))?;
}
delete_configmaps(&Api::namespaced(client.clone(), namespace), selection).await;
Ok(deleted.len())
}
async fn delete_jobs(jobs: &Api<Job>, selection: &Selection) -> Result<(), AgentError> {
let what = &selection.what;
for selector in &selection.jobs {
let list = match jobs.list(&ListParams::default().labels(selector)).await {
Ok(list) => list,
Err(KubeError::Api(e)) if e.code == 403 => {
warn!(error = %e, "cannot list Jobs: JobRun Jobs of {what} are not deleted; grant list and delete on jobs");
continue;
}
Err(e) => return Err(failure(format!("failed to list Jobs of {what}: {e}"))),
};
for name in list.items.into_iter().filter_map(|job| job.metadata.name) {
match jobs.delete(&name, &DeleteParams::background()).await {
Ok(_) => {}
Err(KubeError::Api(e)) if e.code == 404 => {}
Err(e) => {
return Err(failure(format!(
"failed to delete Job '{name}' of {what}: {e}"
)));
}
}
}
}
Ok(())
}
async fn delete_configmaps(configmaps: &Api<ConfigMap>, selection: &Selection) {
let mut names: Vec<String> = Vec::new();
for selector in &selection.configmaps {
match configmaps
.list(&ListParams::default().labels(selector))
.await
{
Ok(list) => names.extend(list.items.into_iter().filter_map(|cm| cm.metadata.name)),
Err(e) => warn!(error = %e, "failed to list prompt ConfigMaps of {}", selection.what),
}
}
names.sort();
names.dedup();
for name in names {
if let Err(e) = configmaps.delete(&name, &DeleteParams::default()).await {
warn!(configmap = %name, error = %e, "failed to delete prompt ConfigMap");
}
}
}