use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use futures_util::{AsyncBufReadExt, TryStreamExt};
use k8s_openapi::api::core::v1::{ConfigMap, Pod};
use kube::Client;
use kube::api::{Api, DeleteParams, LogParams, PostParams};
use kube::runtime::wait::await_condition;
use serde_json::{from_value, json};
use tokio::spawn;
use tokio::task::JoinHandle;
use tokio::time;
use tracing::{debug, info, warn};
use crate::error::AgentError;
use crate::provider::{
AgentConfig, AgentOutput, AgentProvider, InvokeFuture, LABEL_EGRESS_PROFILE, LABEL_EXPIRES_AT,
LABEL_ROOT_RUN_ID, LABEL_RUN_ID, LABEL_STEP, LogSink, PodVolumeSource, ReadOnlyVolume,
ReleaseFuture, SecretEnvVar, assert_pod_label_allowed, is_reserved_pod_label,
upsert_secret_env,
};
use crate::providers::claude::common as claude_common;
use crate::providers::claude::common::DEFAULT_TIMEOUT;
use super::cleanup::{delete_and_wait, release_run, step_selection};
use super::common::{
DEFAULT_DEADLINE_MARGIN, DEFAULT_INPUT_INIT_IMAGE, ImagePullPolicy, K8sClusterConfig,
K8sResources, PodConfig, PodHardening, SandboxSettings, build_credentials_from_env_prefix,
build_credentials_prefix, build_pod_spec, create_client, generate_pod_name,
};
use super::profile::{ClaudeProfile, build_profile_copy_prefix};
use super::reaper::{ReapReport, reap_orphans};
use super::toleration::K8sToleration;
const CREDENTIALS_ENV_VAR: &str = "IRONFLOW_CLAUDE_CREDENTIALS";
const PLAIN_TEXT_SECRETS: [&str; 2] = ["CLAUDE_CODE_OAUTH_TOKEN", "ANTHROPIC_API_KEY"];
pub(super) const MANAGED_SELECTOR: &str = "app.kubernetes.io/managed-by=ironflow";
pub(super) const RUNNER_SELECTOR: &str =
"app.kubernetes.io/managed-by=ironflow,app.kubernetes.io/component=claude-runner";
pub(super) const PROMPT_SELECTOR: &str =
"app.kubernetes.io/managed-by=ironflow,app.kubernetes.io/component=prompt-data";
pub(super) fn now_unix() -> Result<u64, AgentError> {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.map_err(|e| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("system clock is before the unix epoch: {e}"),
})
}
fn is_terminal_phase(phase: &str) -> bool {
phase == "Succeeded" || phase == "Failed"
}
fn is_pod_completed() -> impl kube::runtime::wait::Condition<Pod> {
|obj: Option<&Pod>| {
obj.and_then(|pod| pod.status.as_ref())
.and_then(|status| status.phase.as_deref())
.is_some_and(is_terminal_phase)
}
}
fn is_pod_running_or_terminal() -> impl kube::runtime::wait::Condition<Pod> {
|obj: Option<&Pod>| {
obj.and_then(|pod| pod.status.as_ref())
.and_then(|status| status.phase.as_deref())
.is_some_and(|phase| phase == "Running" || is_terminal_phase(phase))
}
}
#[derive(Clone)]
pub struct K8sEphemeralProvider {
image: String,
namespace: String,
claude_path: String,
working_dir: Option<String>,
resources: K8sResources,
service_account: Option<String>,
image_pull_policy: ImagePullPolicy,
env_vars: Vec<(String, String)>,
image_pull_secrets: Vec<String>,
oauth_credentials: Option<String>,
cluster_config: K8sClusterConfig,
timeout: Duration,
pod_labels: BTreeMap<String, String>,
volumes: Vec<(String, String)>,
pvc_volumes: Vec<(String, String)>,
input_init_image: String,
node_selector: BTreeMap<String, String>,
tolerations: Vec<K8sToleration>,
active_deadline_seconds: Option<Duration>,
sandbox: Option<SandboxSettings>,
secret_env: Vec<SecretEnvVar>,
oauth_credentials_secret: Option<(String, String)>,
read_only_volumes: Vec<ReadOnlyVolume>,
managed_settings_presets: BTreeMap<String, String>,
default_managed_settings: Option<String>,
pub(super) claude_profiles: Vec<ClaudeProfile>,
egress_profile: Option<String>,
previous_attempt_timeout: Duration,
}
fn apply_active_deadline_seconds(pod: &mut Pod, deadline: Option<Duration>) {
if let (Some(d), Some(spec)) = (deadline, pod.spec.as_mut()) {
spec.active_deadline_seconds = Some(i64::try_from(d.as_secs()).unwrap_or(i64::MAX));
}
}
impl K8sEphemeralProvider {
pub fn new(image: &str) -> Self {
Self {
image: image.to_string(),
namespace: "default".to_string(),
claude_path: "claude".to_string(),
working_dir: None,
resources: K8sResources::default(),
service_account: None,
image_pull_policy: ImagePullPolicy::default(),
env_vars: Vec::new(),
image_pull_secrets: Vec::new(),
oauth_credentials: None,
cluster_config: K8sClusterConfig::default(),
timeout: DEFAULT_TIMEOUT,
pod_labels: BTreeMap::new(),
volumes: Vec::new(),
pvc_volumes: Vec::new(),
input_init_image: DEFAULT_INPUT_INIT_IMAGE.to_string(),
node_selector: BTreeMap::new(),
tolerations: Vec::new(),
active_deadline_seconds: None,
sandbox: None,
secret_env: Vec::new(),
oauth_credentials_secret: None,
read_only_volumes: Vec::new(),
managed_settings_presets: BTreeMap::new(),
default_managed_settings: None,
claude_profiles: Vec::new(),
egress_profile: None,
previous_attempt_timeout: Duration::from_secs(60),
}
}
pub fn sandboxed(image: &str) -> Self {
Self {
sandbox: Some(SandboxSettings::default()),
..Self::new(image)
}
}
fn sandbox_mut(&mut self, method: &str) -> &mut SandboxSettings {
self.sandbox.as_mut().unwrap_or_else(|| {
panic!("{method} requires a provider built with K8sEphemeralProvider::sandboxed")
})
}
pub fn allow_writable_root(mut self) -> Self {
self.sandbox_mut("allow_writable_root").writable_root = true;
self
}
pub fn home_size_limit(mut self, limit: &str) -> Self {
self.sandbox_mut("home_size_limit").home_size_limit = limit.to_string();
self
}
pub fn tmp_size_limit(mut self, limit: &str) -> Self {
self.sandbox_mut("tmp_size_limit").tmp_size_limit = limit.to_string();
self
}
pub fn deadline_margin(mut self, margin: Duration) -> Self {
self.sandbox_mut("deadline_margin").deadline_margin = margin;
self
}
pub fn run_as_user(mut self, uid: i64) -> Self {
assert!(uid > 0, "run_as_user must be greater than 0");
self.sandbox_mut("run_as_user").run_as_user = uid;
self
}
pub fn env_from_secret(mut self, var: &str, secret: &str, key: &str) -> Self {
let entry = SecretEnvVar {
name: var.to_string(),
secret: secret.to_string(),
key: key.to_string(),
};
upsert_secret_env(&mut self.secret_env, entry);
self
}
pub fn oauth_credentials_from_secret(mut self, secret: &str, key: &str) -> Self {
self.oauth_credentials_secret = Some((secret.to_string(), key.to_string()));
self
}
pub fn oauth_token_from_secret(self, secret: &str, key: &str) -> Self {
self.env_from_secret("CLAUDE_CODE_OAUTH_TOKEN", secret, key)
}
pub fn read_only_volume(mut self, volume: ReadOnlyVolume) -> Self {
self.read_only_volumes.push(volume);
self
}
pub fn read_only_pvc(self, claim: &str, mount_path: &str) -> Self {
self.read_only_volume(ReadOnlyVolume {
source: PodVolumeSource::PersistentVolumeClaim {
claim_name: claim.to_string(),
},
mount_path: mount_path.to_string(),
sub_path: None,
})
}
pub fn managed_settings_preset(mut self, name: &str, configmap: &str) -> Self {
self.managed_settings_presets
.insert(name.to_string(), configmap.to_string());
self
}
pub fn default_managed_settings(mut self, name: &str) -> Self {
self.default_managed_settings = Some(name.to_string());
self
}
pub fn egress_profile(mut self, name: &str) -> Self {
self.egress_profile = Some(name.to_string());
self
}
pub fn previous_attempt_timeout(mut self, timeout: Duration) -> Self {
self.previous_attempt_timeout = timeout;
self
}
pub fn namespace(mut self, ns: &str) -> Self {
self.namespace = ns.to_string();
self
}
pub fn claude_path(mut self, path: &str) -> Self {
self.claude_path = path.to_string();
self
}
pub fn working_dir(mut self, dir: &str) -> Self {
self.working_dir = Some(dir.to_string());
self
}
pub fn resources(mut self, resources: K8sResources) -> Self {
self.resources = resources;
self
}
pub fn service_account(mut self, sa: &str) -> Self {
self.service_account = Some(sa.to_string());
self
}
pub fn image_pull_policy(mut self, policy: ImagePullPolicy) -> Self {
self.image_pull_policy = policy;
self
}
pub fn oauth_credentials(mut self, json: &str) -> Self {
self.oauth_credentials = Some(json.to_string());
self
}
pub fn image_pull_secret(mut self, secret_name: &str) -> Self {
self.image_pull_secrets.push(secret_name.to_string());
self
}
pub fn env(mut self, key: &str, value: &str) -> Self {
self.env_vars.push((key.to_string(), value.to_string()));
self
}
pub fn cluster_config(mut self, config: K8sClusterConfig) -> Self {
self.cluster_config = config;
self
}
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
pub fn active_deadline_seconds(mut self, duration: Duration) -> Self {
self.active_deadline_seconds = Some(duration);
self
}
pub fn pod_label(mut self, key: &str, value: &str) -> Self {
assert_pod_label_allowed(key);
self.pod_labels.insert(key.to_string(), value.to_string());
self
}
pub fn pod_labels(mut self, labels: BTreeMap<String, String>) -> Self {
labels.keys().for_each(|key| assert_pod_label_allowed(key));
self.pod_labels = labels;
self
}
pub fn volume(mut self, host_path: &str, container_path: &str) -> Self {
self.volumes
.push((host_path.to_string(), container_path.to_string()));
self
}
pub fn pvc_volume(mut self, claim_name: &str, mount_path: &str) -> Self {
self.pvc_volumes
.push((claim_name.to_string(), mount_path.to_string()));
self
}
pub fn input_init_image(mut self, image: &str) -> Self {
self.input_init_image = image.to_string();
self
}
pub fn node_selector(mut self, key: &str, value: &str) -> Self {
self.node_selector
.insert(key.to_string(), value.to_string());
self
}
pub fn toleration(mut self, toleration: K8sToleration) -> Self {
self.tolerations.push(toleration);
self
}
}
const PROMPT_MOUNT_PATH: &str = "/mnt/ironflow-prompt";
const PROMPT_CM_KEY: &str = "prompt";
struct CreatedPod {
pods: Api<Pod>,
pod_name: String,
start: Instant,
prompt_configmap: Option<String>,
configmaps: Option<Api<ConfigMap>>,
}
#[derive(Debug)]
struct MergedPodInputs {
secret_env: Vec<SecretEnvVar>,
service_account: Option<String>,
read_only_volumes: Vec<ReadOnlyVolume>,
managed_settings_configmap: Option<String>,
labels: BTreeMap<String, String>,
}
impl K8sEphemeralProvider {
fn merged_pod_inputs(&self, config: &AgentConfig) -> Result<MergedPodInputs, AgentError> {
if self.sandbox.is_some() {
if self.oauth_credentials.is_some() {
return Err(AgentError::ProcessFailed {
exit_code: -1,
stderr: "sandboxed provider refuses inline oauth_credentials; use oauth_credentials_from_secret".to_string(),
});
}
if let Some((key, _)) = self
.env_vars
.iter()
.find(|(k, _)| PLAIN_TEXT_SECRETS.contains(&k.as_str()))
{
return Err(AgentError::ProcessFailed {
exit_code: -1,
stderr: format!(
"sandboxed provider refuses {key} as a plain env var; use env_from_secret"
),
});
}
}
let mut secret_env = self.secret_env.clone();
if let Some((secret, key)) = &self.oauth_credentials_secret {
secret_env.push(SecretEnvVar {
name: CREDENTIALS_ENV_VAR.to_string(),
secret: secret.clone(),
key: key.clone(),
});
}
for entry in &config.pod.secret_env {
upsert_secret_env(&mut secret_env, entry.clone());
}
let service_account = config
.pod
.service_account
.clone()
.or_else(|| self.service_account.clone());
let mut read_only_volumes = self.read_only_volumes.clone();
read_only_volumes.extend(config.pod.read_only_volumes.iter().cloned());
let preset = config
.pod
.managed_settings
.as_ref()
.or(self.default_managed_settings.as_ref());
let managed_settings_configmap = preset
.map(|name| {
let configmap = self.managed_settings_presets.get(name).cloned();
configmap.ok_or_else(|| {
let known: Vec<&str> = self
.managed_settings_presets
.keys()
.map(String::as_str)
.collect();
AgentError::ProcessFailed {
exit_code: -1,
stderr: format!(
"unknown managed settings preset '{name}', known presets: [{}]",
known.join(", ")
),
}
})
})
.transpose()?;
if let Some(key) = config.pod_labels.keys().find(|k| is_reserved_pod_label(k)) {
return Err(AgentError::ProcessFailed {
exit_code: -1,
stderr: format!(
"pod label '{key}' is reserved: ironflow sets it on every object it creates"
),
});
}
let mut labels = self.pod_labels.clone();
if let Some(profile) = &self.egress_profile {
labels.insert(LABEL_EGRESS_PROFILE.to_string(), profile.clone());
}
labels.extend(config.pod_labels.clone());
Ok(MergedPodInputs {
secret_env,
service_account,
read_only_volumes,
managed_settings_configmap,
labels,
})
}
pub(super) fn home_setup_prefix(&self) -> String {
let credentials = if self.oauth_credentials_secret.is_some() {
build_credentials_from_env_prefix(CREDENTIALS_ENV_VAR)
} else {
build_credentials_prefix(self.oauth_credentials.as_deref())
};
let profiles = build_profile_copy_prefix(&self.claude_profiles);
format!("{profiles}{credentials}")
}
fn deadline_margin_or_default(&self) -> Duration {
self.sandbox
.as_ref()
.map_or(DEFAULT_DEADLINE_MARGIN, |s| s.deadline_margin)
}
fn effective_deadline(&self) -> Option<Duration> {
self.active_deadline_seconds.or_else(|| {
self.sandbox
.as_ref()
.map(|s| self.timeout + s.deadline_margin)
})
}
async fn delete_previous_attempt(
&self,
client: &Client,
run_id: &str,
step: &str,
) -> Result<usize, AgentError> {
let selection = step_selection(run_id, step);
let limit = self.previous_attempt_timeout;
let deleted = delete_and_wait(client, &self.namespace, &selection, limit).await?;
if deleted > 0 {
info!(
run_id = %run_id,
step = %step,
pods_deleted = deleted,
"deleted pods of a previous attempt before retrying"
);
}
Ok(deleted)
}
pub async fn reap_orphans(&self) -> Result<ReapReport, AgentError> {
reap_orphans(&self.cluster_config, &self.namespace).await
}
pub fn spawn_orphan_reaper(&self, interval: Duration) -> JoinHandle<()> {
assert!(
!interval.is_zero(),
"orphan reaper interval must be greater than zero"
);
let cluster_config = self.cluster_config.clone();
let namespace = self.namespace.clone();
spawn(async move {
let mut ticker = time::interval(interval);
loop {
ticker.tick().await;
if let Err(e) = reap_orphans(&cluster_config, &namespace).await {
warn!(error = %e, "orphan reaping pass failed");
}
}
})
}
async fn create_pod(&self, config: &AgentConfig) -> Result<CreatedPod, AgentError> {
claude_common::validate_prompt_size(config)?;
let built = claude_common::build_command(config)?;
let merged = self.merged_pod_inputs(config)?;
let pod_name = generate_pod_name("claude-code");
let creds_prefix = self.home_setup_prefix();
let start = Instant::now();
let client = create_client(&self.cluster_config).await?;
let pods: Api<Pod> = Api::namespaced(client.clone(), &self.namespace);
let mut prompt_configmap_name: Option<String> = None;
let configmaps: Api<ConfigMap> = Api::namespaced(client.clone(), &self.namespace);
let run_id = merged.labels.get(LABEL_RUN_ID);
let step = merged.labels.get(LABEL_STEP);
if let (Some(run_id), Some(step)) = (run_id, step) {
self.delete_previous_attempt(&client, run_id, step).await?;
}
let lifetime = self.timeout + self.deadline_margin_or_default();
let expires_at = now_unix()? + lifetime.as_secs();
let mut annotations = BTreeMap::new();
annotations.insert(LABEL_EXPIRES_AT.to_string(), expires_at.to_string());
let trace_prefix = config
.trace_context
.as_ref()
.map(|ctx| format!("export TRACEPARENT='{}'; ", ctx.to_traceparent()))
.unwrap_or_default();
let full_cmd = if let Some(ref prompt) = built.stdin_prompt {
let cm_name = format!("{pod_name}-prompt");
let mut cm_labels = BTreeMap::new();
cm_labels.insert("app.kubernetes.io/managed-by", "ironflow");
cm_labels.insert("app.kubernetes.io/component", "prompt-data");
if let Some(root) = merged.labels.get(LABEL_ROOT_RUN_ID) {
cm_labels.insert(LABEL_ROOT_RUN_ID, root.as_str());
}
if let (Some(run_id), Some(step)) = (run_id, step) {
cm_labels.insert(LABEL_RUN_ID, run_id.as_str());
cm_labels.insert(LABEL_STEP, step.as_str());
}
let cm: ConfigMap = from_value(json!({
"apiVersion": "v1",
"kind": "ConfigMap",
"metadata": {
"name": &cm_name,
"namespace": &self.namespace,
"labels": cm_labels,
"annotations": &annotations
},
"data": {
PROMPT_CM_KEY: prompt
}
}))
.map_err(|e| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed to build prompt ConfigMap: {e}"),
})?;
configmaps
.create(&PostParams::default(), &cm)
.await
.map_err(|e| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed to create prompt ConfigMap: {e}"),
})?;
prompt_configmap_name = Some(cm_name);
info!(
pod = %pod_name,
prompt_bytes = prompt.len(),
"prompt too large for CLI args, using ConfigMap + stdin pipe"
);
let prompt_file = format!("{PROMPT_MOUNT_PATH}/{PROMPT_CM_KEY}");
let claude_cmd = claude_common::build_shell_command(&self.claude_path, &built.args);
let pipe_prefix = format!(
"cat {} | ",
claude_common::build_shell_command(&prompt_file, &[])
);
match (&self.working_dir, &config.working_dir) {
(_, Some(dir)) | (Some(dir), None) => {
format!(
"{trace_prefix}{creds_prefix}cd {} && {pipe_prefix}{claude_cmd}",
claude_common::build_shell_command(dir, &[]),
)
}
(None, None) => format!("{trace_prefix}{creds_prefix}{pipe_prefix}{claude_cmd}"),
}
} else {
let claude_cmd = claude_common::build_shell_command(&self.claude_path, &built.args);
match (&self.working_dir, &config.working_dir) {
(_, Some(dir)) | (Some(dir), None) => {
format!(
"{trace_prefix}{creds_prefix}cd {} && {}",
claude_common::build_shell_command(dir, &[]),
claude_cmd
)
}
(None, None) => format!("{trace_prefix}{creds_prefix}{claude_cmd}"),
}
};
debug!(
pod_name = %pod_name,
namespace = %self.namespace,
image = %self.image,
model = %config.model,
prompt_via_configmap = prompt_configmap_name.is_some(),
"creating ephemeral K8s pod"
);
let pod_spec = build_pod_spec(&PodConfig {
name: &pod_name,
image: &self.image,
command: vec!["sh".to_string(), "-c".to_string(), full_cmd],
namespace: &self.namespace,
resources: &self.resources,
service_account: merged.service_account.as_deref(),
restart_policy: "Never",
image_pull_policy: &self.image_pull_policy,
env_vars: &self.env_vars,
image_pull_secrets: &self.image_pull_secrets,
extra_labels: &merged.labels,
node_selector: &self.node_selector,
tolerations: &self.tolerations,
volumes: &self.volumes,
pvc_volumes: &self.pvc_volumes,
inputs: &config.inputs,
input_init_image: &self.input_init_image,
prompt_configmap: prompt_configmap_name.as_deref(),
prompt_mount_path: PROMPT_MOUNT_PATH,
hardening: PodHardening {
sandbox: self.sandbox.as_ref(),
secret_env: &merged.secret_env,
read_only_volumes: &merged.read_only_volumes,
managed_settings_configmap: merged.managed_settings_configmap.as_deref(),
claude_profiles: &self.claude_profiles,
annotations: Some(&annotations),
},
})?;
let mut pod_spec = pod_spec;
apply_active_deadline_seconds(&mut pod_spec, self.effective_deadline());
pods.create(&PostParams::default(), &pod_spec)
.await
.map_err(|e| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed to create K8s pod: {e}"),
})?;
Ok(CreatedPod {
pods,
pod_name,
start,
prompt_configmap: prompt_configmap_name,
configmaps: Some(configmaps),
})
}
fn finalize_pod(
&self,
logs: &str,
pod_phase: &str,
timed_out: bool,
pod_name: &str,
config: &AgentConfig,
start: Instant,
) -> Result<AgentOutput, AgentError> {
if timed_out {
warn!(timeout = ?self.timeout, pod = %pod_name, "K8s pod timed out");
return Err(AgentError::Timeout {
limit: self.timeout,
});
}
let duration_ms = start.elapsed().as_millis() as u64;
let exit_code = if pod_phase == "Succeeded" { 0 } else { 1 };
if exit_code != 0 {
return claude_common::handle_nonzero_exit(
exit_code,
logs,
"",
config,
duration_ms,
"ephemeral k8s",
);
}
debug!(stdout_len = logs.len(), "ephemeral claude pod completed");
claude_common::parse_output(logs, config, duration_ms)
}
}
impl AgentProvider for K8sEphemeralProvider {
fn release_run<'a>(&'a self, run_id: &'a str) -> ReleaseFuture<'a> {
let limit = self.previous_attempt_timeout;
Box::pin(release_run(
&self.cluster_config,
&self.namespace,
run_id,
limit,
))
}
fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
Box::pin(async move {
let created = self.create_pod(config).await?;
let CreatedPod {
pods,
pod_name,
start,
prompt_configmap,
configmaps,
} = &created;
let wait_result = time::timeout(
self.timeout,
await_condition(pods.clone(), pod_name, is_pod_completed()),
)
.await;
let timed_out = wait_result.is_err();
let pod_phase = if timed_out {
"TimedOut".to_string()
} else {
let condition_result =
wait_result.expect("timeout already handled").map_err(|e| {
AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed waiting for pod completion: {e}"),
}
})?;
condition_result
.and_then(|p| p.status)
.and_then(|s| s.phase)
.unwrap_or_else(|| "Unknown".to_string())
};
let logs = pods
.logs(pod_name, &LogParams::default())
.await
.unwrap_or_default();
let _ = pods.delete(pod_name, &DeleteParams::default()).await;
if let (Some(cm_name), Some(cm_api)) = (prompt_configmap, configmaps) {
let _ = cm_api.delete(cm_name, &DeleteParams::default()).await;
}
self.finalize_pod(&logs, &pod_phase, timed_out, pod_name, config, *start)
})
}
fn invoke_with_logs<'a>(
&'a self,
config: &'a AgentConfig,
log_sink: Arc<dyn LogSink>,
) -> InvokeFuture<'a> {
Box::pin(async move {
let effective_config;
let config = if !config.verbose {
debug!("forcing verbose=true for log streaming (stream-json output required)");
effective_config = config.clone().verbose(true);
&effective_config
} else {
config
};
let created = self.create_pod(config).await?;
let CreatedPod {
pods,
pod_name,
start,
prompt_configmap,
configmaps,
} = &created;
let ready_result = time::timeout(
self.timeout,
await_condition(pods.clone(), pod_name, is_pod_running_or_terminal()),
)
.await;
if let Err(_elapsed) = ready_result {
let _ = pods.delete(pod_name, &DeleteParams::default()).await;
if let (Some(cm_name), Some(cm_api)) = (prompt_configmap, configmaps) {
let _ = cm_api.delete(cm_name, &DeleteParams::default()).await;
}
warn!(timeout = ?self.timeout, pod = %pod_name, "K8s pod timed out waiting for Running");
return Err(AgentError::Timeout {
limit: self.timeout,
});
}
if let Err(e) = ready_result.expect("timeout already handled") {
let _ = pods.delete(pod_name, &DeleteParams::default()).await;
if let (Some(cm_name), Some(cm_api)) = (prompt_configmap, configmaps) {
let _ = cm_api.delete(cm_name, &DeleteParams::default()).await;
}
return Err(AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed waiting for pod to start: {e}"),
});
}
let phase_after_ready = pods
.get(pod_name)
.await
.ok()
.and_then(|p| p.status)
.and_then(|s| s.phase);
let already_terminal = phase_after_ready.as_deref().is_some_and(is_terminal_phase);
let mut accumulated = String::new();
let mut timed_out = false;
if already_terminal {
debug!(pod = %pod_name, phase = ?phase_after_ready, "pod already terminal, skipping log stream");
accumulated = pods
.logs(pod_name, &LogParams::default())
.await
.unwrap_or_default();
for line in accumulated.lines() {
log_sink.log("stdout", line);
}
} else {
let log_params = LogParams {
follow: true,
..Default::default()
};
const MAX_ACCUMULATED_BYTES: usize = 50 * 1024 * 1024;
let mut truncated = false;
let completion_notify = Arc::new(tokio::sync::Notify::new());
let watcher_handle = {
let pods = pods.clone();
let pod_name = pod_name.to_string();
let notify = completion_notify.clone();
tokio::spawn(async move {
let _ = await_condition(pods, &pod_name, is_pod_completed()).await;
notify.notify_waiters();
})
};
let stream_result = time::timeout(
self.timeout,
async {
match pods.log_stream(pod_name, &log_params).await {
Ok(stream) => {
let mut lines = stream.lines();
loop {
tokio::select! {
line_result = lines.try_next() => {
match line_result {
Ok(Some(line)) => {
log_sink.log("stdout", &line);
if !truncated {
if accumulated.len() + line.len() + 1
> MAX_ACCUMULATED_BYTES
{
truncated = true;
warn!(pod = %pod_name, "log accumulation cap reached, further output will only be streamed");
} else {
accumulated.push_str(&line);
accumulated.push('\n');
}
}
}
Ok(None) => break,
Err(_) => break,
}
}
_ = completion_notify.notified() => {
debug!(pod = %pod_name, "pod completed, draining remaining log lines");
while let Ok(Some(line)) = time::timeout(
Duration::from_secs(2),
lines.try_next(),
).await.unwrap_or(Ok(None)) {
log_sink.log("stdout", &line);
if !truncated {
if accumulated.len() + line.len() + 1
> MAX_ACCUMULATED_BYTES
{
truncated = true;
} else {
accumulated.push_str(&line);
accumulated.push('\n');
}
}
}
break;
}
}
}
Ok(())
}
Err(e) => Err(e),
}
},
)
.await;
watcher_handle.abort();
timed_out = stream_result.is_err();
if let Ok(Err(e)) = stream_result {
warn!(pod = %pod_name, error = %e, "failed to open log stream, falling back to batch read");
let _ = time::timeout(
self.timeout,
await_condition(pods.clone(), pod_name, is_pod_completed()),
)
.await;
accumulated = pods
.logs(pod_name, &LogParams::default())
.await
.unwrap_or_default();
for line in accumulated.lines() {
log_sink.log("stdout", line);
}
}
}
let pod_phase = if already_terminal {
phase_after_ready.unwrap_or_else(|| "Unknown".to_string())
} else {
match pods.get(pod_name).await {
Ok(pod) => pod
.status
.and_then(|s| s.phase)
.unwrap_or_else(|| "Unknown".to_string()),
Err(_) => "Unknown".to_string(),
}
};
let _ = pods.delete(pod_name, &DeleteParams::default()).await;
if let (Some(cm_name), Some(cm_api)) = (prompt_configmap, configmaps) {
let _ = cm_api.delete(cm_name, &DeleteParams::default()).await;
}
self.finalize_pod(
&accumulated,
&pod_phase,
timed_out,
pod_name,
config,
*start,
)
})
}
}
#[cfg(test)]
mod label_tests;
#[cfg(test)]
mod tests {
use super::super::toleration::{TolerationEffect, TolerationOperator};
use super::*;
#[test]
fn ephemeral_provider_defaults() {
let provider = K8sEphemeralProvider::new("my-image:v1");
assert_eq!(provider.image, "my-image:v1");
assert_eq!(provider.namespace, "default");
assert_eq!(provider.claude_path, "claude");
assert!(provider.working_dir.is_none());
assert!(provider.service_account.is_none());
assert_eq!(provider.timeout, DEFAULT_TIMEOUT);
}
#[test]
fn ephemeral_provider_builder_chain() {
let provider = K8sEphemeralProvider::new("img:v2")
.namespace("ci")
.claude_path("/usr/bin/claude")
.working_dir("/workspace")
.service_account("claude-sa")
.resources(K8sResources {
cpu_limit: Some("1".to_string()),
memory_limit: Some("2Gi".to_string()),
})
.timeout(Duration::from_secs(600));
assert_eq!(provider.namespace, "ci");
assert_eq!(provider.claude_path, "/usr/bin/claude");
assert_eq!(provider.working_dir, Some("/workspace".to_string()));
assert_eq!(provider.service_account, Some("claude-sa".to_string()));
assert_eq!(provider.resources.cpu_limit, Some("1".to_string()));
assert_eq!(provider.resources.memory_limit, Some("2Gi".to_string()));
assert_eq!(provider.timeout, Duration::from_secs(600));
}
#[test]
fn ephemeral_provider_image_pull_secrets() {
let provider = K8sEphemeralProvider::new("registry.gitlab.com/org/img:v1")
.image_pull_secret("gitlab-registry")
.image_pull_secret("dockerhub");
assert_eq!(provider.image_pull_secrets.len(), 2);
assert_eq!(provider.image_pull_secrets[0], "gitlab-registry");
assert_eq!(provider.image_pull_secrets[1], "dockerhub");
}
#[test]
fn ephemeral_provider_clone() {
let provider = K8sEphemeralProvider::new("img")
.namespace("ns")
.timeout(Duration::from_secs(42));
let cloned = provider.clone();
assert_eq!(cloned.namespace, "ns");
assert_eq!(cloned.timeout, Duration::from_secs(42));
}
#[test]
fn ephemeral_provider_pod_labels_default_empty() {
let provider = K8sEphemeralProvider::new("img:v1");
assert!(provider.pod_labels.is_empty());
}
#[test]
fn ephemeral_provider_pod_labels_builder() {
let mut labels = BTreeMap::new();
labels.insert("env".to_string(), "staging".to_string());
labels.insert("team".to_string(), "platform".to_string());
let provider = K8sEphemeralProvider::new("img:v1").pod_labels(labels);
assert_eq!(provider.pod_labels.len(), 2);
assert_eq!(provider.pod_labels["env"], "staging");
assert_eq!(provider.pod_labels["team"], "platform");
}
#[test]
fn ephemeral_provider_pod_label_builder() {
let provider = K8sEphemeralProvider::new("img:v1")
.pod_label("env", "prod")
.pod_label("team", "infra");
assert_eq!(provider.pod_labels.len(), 2);
assert_eq!(provider.pod_labels["env"], "prod");
assert_eq!(provider.pod_labels["team"], "infra");
}
#[test]
fn ephemeral_provider_volume_builder() {
let provider = K8sEphemeralProvider::new("img:v1")
.volume("/tmp/worktrees", "/data/worktrees")
.volume("/tmp/repos", "/data/repos");
assert_eq!(provider.volumes.len(), 2);
assert_eq!(
provider.volumes[0],
("/tmp/worktrees".to_string(), "/data/worktrees".to_string())
);
assert_eq!(
provider.volumes[1],
("/tmp/repos".to_string(), "/data/repos".to_string())
);
}
#[test]
fn ephemeral_provider_volumes_default_empty() {
let provider = K8sEphemeralProvider::new("img:v1");
assert!(provider.volumes.is_empty());
}
#[test]
fn ephemeral_provider_pvc_volume_builder() {
let provider = K8sEphemeralProvider::new("img:v1")
.pvc_volume("jarvis-repos", "/data/repos")
.pvc_volume("jarvis-worktrees", "/data/worktrees");
assert_eq!(provider.pvc_volumes.len(), 2);
assert_eq!(
provider.pvc_volumes[0],
("jarvis-repos".to_string(), "/data/repos".to_string())
);
assert_eq!(
provider.pvc_volumes[1],
(
"jarvis-worktrees".to_string(),
"/data/worktrees".to_string()
)
);
}
#[test]
fn ephemeral_provider_pvc_volumes_default_empty() {
let provider = K8sEphemeralProvider::new("img:v1");
assert!(provider.pvc_volumes.is_empty());
}
#[test]
fn ephemeral_provider_node_selector_default_empty() {
let provider = K8sEphemeralProvider::new("img:v1");
assert!(provider.node_selector.is_empty());
}
#[test]
fn ephemeral_provider_node_selector_accumulates() {
let provider = K8sEphemeralProvider::new("img:v1")
.node_selector("kubernetes.io/hostname", "ryzen1")
.node_selector("workload", "agent");
assert_eq!(provider.node_selector.len(), 2);
assert_eq!(provider.node_selector["kubernetes.io/hostname"], "ryzen1");
assert_eq!(provider.node_selector["workload"], "agent");
}
#[test]
fn ephemeral_provider_toleration_default_empty() {
let provider = K8sEphemeralProvider::new("img:v1");
assert!(provider.tolerations.is_empty());
}
#[test]
fn ephemeral_provider_toleration_accumulates() {
let provider = K8sEphemeralProvider::new("img:v1")
.toleration(K8sToleration {
key: "dedicated".to_string(),
operator: TolerationOperator::Equal,
value: Some("worker".to_string()),
effect: TolerationEffect::NoSchedule,
toleration_seconds: None,
})
.toleration(K8sToleration {
key: "gpu".to_string(),
operator: TolerationOperator::Exists,
value: None,
effect: TolerationEffect::NoExecute,
toleration_seconds: None,
});
assert_eq!(provider.tolerations.len(), 2);
assert_eq!(provider.tolerations[0].key, "dedicated");
assert_eq!(provider.tolerations[1].operator, TolerationOperator::Exists);
}
#[test]
fn ephemeral_provider_active_deadline_default_and_builder() {
let default = K8sEphemeralProvider::new("img:v1");
assert!(default.active_deadline_seconds.is_none());
let provider =
K8sEphemeralProvider::new("img:v1").active_deadline_seconds(Duration::from_secs(900));
assert_eq!(
provider.active_deadline_seconds,
Some(Duration::from_secs(900))
);
}
#[test]
fn apply_active_deadline_sets_field_in_seconds() {
let mut pod: Pod =
from_value(json!({"spec": {"containers": []}})).expect("valid minimal pod");
apply_active_deadline_seconds(&mut pod, Some(Duration::from_secs(600)));
assert_eq!(pod.spec.unwrap().active_deadline_seconds, Some(600));
}
#[test]
fn apply_active_deadline_none_leaves_field_absent() {
let mut pod: Pod =
from_value(json!({"spec": {"containers": []}})).expect("valid minimal pod");
apply_active_deadline_seconds(&mut pod, None);
assert_eq!(pod.spec.unwrap().active_deadline_seconds, None);
}
fn err_text(result: Result<MergedPodInputs, AgentError>) -> String {
result.unwrap_err().to_string()
}
fn merge(provider: &K8sEphemeralProvider, config: &AgentConfig) -> MergedPodInputs {
provider.merged_pod_inputs(config).unwrap()
}
#[test]
fn sandboxed_defaults() {
let provider = K8sEphemeralProvider::sandboxed("img:v1");
let sandbox = provider.sandbox.as_ref().expect("sandbox set");
assert_eq!(sandbox, &SandboxSettings::default());
assert_eq!(sandbox.deadline_margin, Duration::from_secs(60));
assert_eq!(provider.image, "img:v1");
assert_eq!(provider.previous_attempt_timeout, Duration::from_secs(60));
assert!(K8sEphemeralProvider::new("img:v1").sandbox.is_none());
}
#[test]
fn sandboxed_relaxation_builders() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.allow_writable_root()
.home_size_limit("4Gi")
.tmp_size_limit("2Gi")
.deadline_margin(Duration::from_secs(120))
.run_as_user(20000);
let sandbox = provider.sandbox.unwrap();
assert!(sandbox.writable_root);
assert_eq!(sandbox.home_size_limit, "4Gi");
assert_eq!(sandbox.tmp_size_limit, "2Gi");
assert_eq!(sandbox.deadline_margin, Duration::from_secs(120));
assert_eq!(sandbox.run_as_user, 20000);
}
#[test]
#[should_panic(expected = "allow_writable_root requires a provider built with")]
fn allow_writable_root_on_non_sandboxed_panics() {
let _ = K8sEphemeralProvider::new("img:v1").allow_writable_root();
}
#[test]
#[should_panic(expected = "home_size_limit requires a provider built with")]
fn home_size_limit_on_non_sandboxed_panics() {
let _ = K8sEphemeralProvider::new("img:v1").home_size_limit("1Gi");
}
#[test]
#[should_panic(expected = "run_as_user must be greater than 0")]
fn run_as_user_zero_panics() {
let _ = K8sEphemeralProvider::sandboxed("img:v1").run_as_user(0);
}
#[test]
fn oauth_token_from_secret_builder() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.oauth_token_from_secret("claude-oauth", "token");
assert_eq!(provider.secret_env.len(), 1);
assert_eq!(provider.secret_env[0].name, "CLAUDE_CODE_OAUTH_TOKEN");
assert_eq!(provider.secret_env[0].secret, "claude-oauth");
assert_eq!(provider.secret_env[0].key, "token");
}
#[test]
fn oauth_credentials_from_secret_adds_credentials_env() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.oauth_credentials_from_secret("claude-credentials", "credentials.json");
let (secret, key) = provider.oauth_credentials_secret.clone().unwrap();
assert_eq!(secret, "claude-credentials");
assert_eq!(key, "credentials.json");
let merged = merge(&provider, &AgentConfig::new("hi"));
assert_eq!(merged.secret_env.len(), 1);
assert_eq!(merged.secret_env[0].name, CREDENTIALS_ENV_VAR);
assert_eq!(merged.secret_env[0].secret, "claude-credentials");
}
#[test]
fn env_from_secret_replaces_same_var() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.env_from_secret("TOKEN", "a", "k")
.env_from_secret("TOKEN", "b", "k");
assert_eq!(provider.secret_env.len(), 1);
assert_eq!(provider.secret_env[0].secret, "b");
}
#[test]
fn managed_settings_preset_map() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.managed_settings_preset("locked", "claude-managed-locked")
.managed_settings_preset("readonly", "claude-managed-readonly")
.default_managed_settings("locked");
assert_eq!(provider.managed_settings_presets.len(), 2);
assert_eq!(
provider.managed_settings_presets["readonly"],
"claude-managed-readonly"
);
assert_eq!(provider.default_managed_settings.as_deref(), Some("locked"));
}
#[test]
fn other_sandbox_builders() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.read_only_pvc("repos", "/data/repos")
.claude_profile_configmap("claude-profile")
.egress_profile("anthropic-only")
.previous_attempt_timeout(Duration::from_secs(5));
assert_eq!(provider.read_only_volumes.len(), 1);
assert_eq!(provider.read_only_volumes[0].mount_path, "/data/repos");
assert_eq!(provider.claude_profiles[0].configmap, "claude-profile");
assert_eq!(provider.egress_profile.as_deref(), Some("anthropic-only"));
assert_eq!(provider.previous_attempt_timeout, Duration::from_secs(5));
}
#[test]
fn merged_step_secret_overrides_provider_secret() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.env_from_secret("TOKEN", "provider-secret", "k")
.env_from_secret("OTHER", "other", "k");
let config = AgentConfig::new("hi")
.env_from_secret("TOKEN", "step-secret", "k2")
.env_from_secret("STEP_ONLY", "step", "k");
let merged = merge(&provider, &config);
assert_eq!(merged.secret_env.len(), 3);
let token = &merged.secret_env[0];
assert_eq!(token.name, "TOKEN");
assert_eq!(token.secret, "step-secret");
assert_eq!(token.key, "k2");
}
#[test]
fn merged_step_service_account_overrides_provider() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").service_account("provider-sa");
let merged = merge(&provider, &AgentConfig::new("hi"));
assert_eq!(merged.service_account.as_deref(), Some("provider-sa"));
let config = AgentConfig::new("hi").service_account("step-sa");
let merged = merge(&provider, &config);
assert_eq!(merged.service_account.as_deref(), Some("step-sa"));
}
#[test]
fn merged_read_only_volumes_provider_then_step() {
let provider =
K8sEphemeralProvider::sandboxed("img:v1").read_only_pvc("repos", "/data/repos");
let config = AgentConfig::new("hi").read_only_config_map("cm", "/data/cm");
let merged = merge(&provider, &config);
let paths: Vec<&str> = merged
.read_only_volumes
.iter()
.map(|v| v.mount_path.as_str())
.collect();
assert_eq!(paths, vec!["/data/repos", "/data/cm"]);
}
#[test]
fn merged_unknown_preset_is_an_error() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.managed_settings_preset("locked", "claude-managed-locked");
let config = AgentConfig::new("hi").managed_settings("typo");
let err = err_text(provider.merged_pod_inputs(&config));
assert!(
err.contains("unknown managed settings preset 'typo'"),
"{err}"
);
assert!(err.contains("locked"), "{err}");
}
#[test]
fn merged_unknown_default_preset_is_an_error() {
let provider =
K8sEphemeralProvider::sandboxed("img:v1").default_managed_settings("missing");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(
err.contains("unknown managed settings preset 'missing'"),
"{err}"
);
}
#[test]
fn merged_default_preset_used_when_step_has_none() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.managed_settings_preset("locked", "claude-managed-locked")
.managed_settings_preset("readonly", "claude-managed-readonly")
.default_managed_settings("locked");
let merged = merge(&provider, &AgentConfig::new("hi"));
assert_eq!(
merged.managed_settings_configmap.as_deref(),
Some("claude-managed-locked")
);
let config = AgentConfig::new("hi").managed_settings("readonly");
let merged = merge(&provider, &config);
assert_eq!(
merged.managed_settings_configmap.as_deref(),
Some("claude-managed-readonly")
);
}
#[test]
fn merged_no_preset_means_no_managed_settings() {
let provider = K8sEphemeralProvider::new("img:v1");
let merged = merge(&provider, &AgentConfig::new("hi"));
assert!(merged.managed_settings_configmap.is_none());
}
#[test]
fn merged_step_egress_label_overrides_provider() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").egress_profile("anthropic-only");
let merged = merge(&provider, &AgentConfig::new("hi"));
assert_eq!(merged.labels[LABEL_EGRESS_PROFILE], "anthropic-only");
let config = AgentConfig::new("hi").egress_profile("gitlab");
let merged = merge(&provider, &config);
assert_eq!(merged.labels[LABEL_EGRESS_PROFILE], "gitlab");
}
#[test]
fn merged_labels_carry_run_scope() {
let provider = K8sEphemeralProvider::new("img:v1").pod_label("team", "infra");
let config = AgentConfig::new("hi").run_scope("run-1", "investigate");
let merged = merge(&provider, &config);
assert_eq!(merged.labels["team"], "infra");
assert_eq!(merged.labels[LABEL_RUN_ID], "run-1");
assert_eq!(merged.labels[LABEL_STEP], "investigate");
}
#[test]
fn sandboxed_refuses_inline_oauth_credentials() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").oauth_credentials("{}");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("refuses inline oauth_credentials"), "{err}");
}
#[test]
fn sandboxed_refuses_plain_api_key() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").env("ANTHROPIC_API_KEY", "sk");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("ANTHROPIC_API_KEY"), "{err}");
assert!(err.contains("env_from_secret"), "{err}");
}
#[test]
fn sandboxed_refuses_plain_oauth_token() {
let provider =
K8sEphemeralProvider::sandboxed("img:v1").env("CLAUDE_CODE_OAUTH_TOKEN", "tok");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("CLAUDE_CODE_OAUTH_TOKEN"), "{err}");
}
#[test]
fn non_sandboxed_keeps_inline_credentials() {
let provider = K8sEphemeralProvider::new("img:v1")
.oauth_credentials("{}")
.env("ANTHROPIC_API_KEY", "sk");
assert!(provider.merged_pod_inputs(&AgentConfig::new("hi")).is_ok());
}
#[test]
fn effective_deadline_explicit_value_wins() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.timeout(Duration::from_secs(600))
.active_deadline_seconds(Duration::from_secs(30));
assert_eq!(provider.effective_deadline(), Some(Duration::from_secs(30)));
}
#[test]
fn effective_deadline_sandboxed_is_timeout_plus_margin() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").timeout(Duration::from_secs(600));
assert_eq!(
provider.effective_deadline(),
Some(Duration::from_secs(660))
);
}
#[test]
fn effective_deadline_non_sandboxed_defaults_to_none() {
let provider = K8sEphemeralProvider::new("img:v1").timeout(Duration::from_secs(600));
assert_eq!(provider.effective_deadline(), None);
}
#[test]
fn deadline_margin_defaults_to_sixty_seconds_when_not_sandboxed() {
let provider = K8sEphemeralProvider::new("img:v1");
assert_eq!(
provider.deadline_margin_or_default(),
DEFAULT_DEADLINE_MARGIN
);
}
#[test]
#[should_panic(expected = "orphan reaper interval must be greater than zero")]
fn spawn_orphan_reaper_zero_interval_panics() {
drop(K8sEphemeralProvider::sandboxed("img:v1").spawn_orphan_reaper(Duration::ZERO));
}
}