use std::collections::BTreeMap;
use std::env::var;
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::account::ClaudeSubscriptionKind;
use crate::auth_proxy::{
ADMIN_KEY_ENV, AuthProxyClient, AuthProxyError, IssuedToken, POD_BASE_URL_ENV, POD_TOKEN_ENV,
TokenRequest, resolve_credential,
};
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 crate::providers::claude::rate_limit_event;
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"];
const PROXY_FORBIDDEN_ENV: [&str; 5] = [
"CLAUDE_CODE_OAUTH_TOKEN",
"ANTHROPIC_API_KEY",
POD_TOKEN_ENV,
POD_BASE_URL_ENV,
CREDENTIALS_ENV_VAR,
];
const NONESSENTIAL_TRAFFIC_ENV: &str = "CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC";
fn proxy_forbidden(name: &str) -> AgentError {
AgentError::ProcessFailed {
exit_code: -1,
stderr: format!(
"auth_proxy is set: the pod must not receive {name}; the proxy injects the credential"
),
}
}
fn auth_proxy_error(e: AuthProxyError) -> AgentError {
AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("auth proxy: {e}"),
}
}
async fn abort_launch_configmap(configmaps: &Api<ConfigMap>, name: Option<&str>) {
if let Some(name) = name
&& let Err(e) = configmaps.delete(name, &DeleteParams::default()).await
{
warn!(
configmap = %name,
error = %e,
"failed to delete the prompt ConfigMap of an aborted pod"
);
}
}
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,
runtime_class: Option<String>,
auth_proxy_url: Option<String>,
auth_proxy_admin_key: Option<String>,
}
fn apply_runtime_class(pod: &mut Pod, runtime_class: Option<&str>) {
if let (Some(name), Some(spec)) = (runtime_class, pod.spec.as_mut()) {
spec.runtime_class_name = Some(name.to_string());
}
}
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),
runtime_class: None,
auth_proxy_url: None,
auth_proxy_admin_key: None,
}
}
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 auth_proxy(mut self, url: &str) -> Self {
assert!(
url.starts_with("http://") || url.starts_with("https://"),
"auth_proxy url must start with http:// or https://"
);
self.auth_proxy_url = Some(url.trim_end_matches('/').to_string());
self
}
pub fn auth_proxy_admin_key(mut self, key: &str) -> Self {
self.auth_proxy_admin_key = Some(key.to_string());
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 runtime_class(mut self, name: &str) -> Self {
assert!(!name.trim().is_empty(), "runtime class must not be empty");
self.runtime_class = Some(name.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>>,
proxy_token_id: Option<String>,
}
#[derive(Debug)]
struct MergedPodInputs<'a> {
volumes: &'a [(String, String)],
pvc_volumes: &'a [(String, String)],
secret_env: Vec<SecretEnvVar>,
service_account: Option<String>,
read_only_volumes: Vec<ReadOnlyVolume>,
managed_settings_configmap: Option<String>,
labels: BTreeMap<String, String>,
runtime_class: Option<String>,
}
impl K8sEphemeralProvider {
fn merged_pod_inputs<'a>(
&'a self,
config: &AgentConfig,
) -> Result<MergedPodInputs<'a>, AgentError> {
if self.auth_proxy_url.is_some() {
self.check_no_proxy_credential()?;
}
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());
}
if self.auth_proxy_url.is_some()
&& let Some(entry) = secret_env
.iter()
.find(|s| PROXY_FORBIDDEN_ENV.contains(&s.name.as_str()))
{
return Err(proxy_forbidden(&entry.name));
}
let service_account = config
.pod
.service_account
.clone()
.or_else(|| self.service_account.clone());
let runtime_class = config
.pod
.runtime_class
.clone()
.or_else(|| self.runtime_class.clone());
if runtime_class
.as_deref()
.is_some_and(|c| c.trim().is_empty())
{
return Err(AgentError::ProcessFailed {
exit_code: -1,
stderr: "runtime class must not be empty".to_string(),
});
}
let (volumes, pvc_volumes): (&[_], &[_]) = if config.pod.without_provider_volumes {
(&[], &[])
} else {
(&self.volumes, &self.pvc_volumes)
};
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 {
volumes,
pvc_volumes,
secret_env,
service_account,
read_only_volumes,
managed_settings_configmap,
labels,
runtime_class,
})
}
fn check_no_proxy_credential(&self) -> Result<(), AgentError> {
if self.oauth_credentials.is_some() {
return Err(proxy_forbidden("oauth_credentials"));
}
if self.oauth_credentials_secret.is_some() {
return Err(proxy_forbidden(CREDENTIALS_ENV_VAR));
}
let forbidden = self.env_vars.iter().find(|(key, value)| {
PROXY_FORBIDDEN_ENV.contains(&key.as_str()) || value.starts_with("sk-ant")
});
match forbidden {
Some((key, _)) => Err(proxy_forbidden(key)),
None => Ok(()),
}
}
fn pod_env_vars(&self, proxy_token: Option<&str>) -> Vec<(String, String)> {
let mut env = self.env_vars.clone();
if let (Some(url), Some(token)) = (&self.auth_proxy_url, proxy_token) {
env.push((POD_BASE_URL_ENV.to_string(), url.clone()));
env.push((POD_TOKEN_ENV.to_string(), token.to_string()));
env.push((NONESSENTIAL_TRAFFIC_ENV.to_string(), "1".to_string()));
}
env
}
fn auth_proxy_client(&self) -> Result<Option<AuthProxyClient>, AgentError> {
let Some(url) = &self.auth_proxy_url else {
return Ok(None);
};
let key = self
.auth_proxy_admin_key
.clone()
.or_else(|| var(ADMIN_KEY_ENV).ok().filter(|key| !key.is_empty()))
.ok_or_else(|| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("auth_proxy requires {ADMIN_KEY_ENV}"),
})?;
Ok(Some(AuthProxyClient::new(url, &key)))
}
async fn revoke_proxy_token(&self, id: Option<&str>) {
let Some(id) = id else {
return;
};
let short = id.get(..12).unwrap_or(id);
let result = match self.auth_proxy_client() {
Ok(Some(client)) => client.revoke(id).await.map_err(|e| e.to_string()),
Ok(None) => Ok(()),
Err(e) => Err(e.to_string()),
};
match result {
Ok(()) => debug!(token = %short, "auth proxy token revoked"),
Err(e) => warn!(
token = %short,
error = %e,
"auth proxy token revocation failed; it expires at expires-at"
),
}
}
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 issued = match self
.issue_proxy_token(config, run_id, step, &pod_name, expires_at)
.await
{
Ok(issued) => issued,
Err(e) => {
abort_launch_configmap(&configmaps, prompt_configmap_name.as_deref()).await;
return Err(e);
}
};
let env_vars = self.pod_env_vars(issued.as_ref().map(|t| t.token.as_str()));
let proxy_token_id = issued.map(|t| t.id);
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: &env_vars,
image_pull_secrets: &self.image_pull_secrets,
extra_labels: &merged.labels,
node_selector: &self.node_selector,
tolerations: &self.tolerations,
volumes: merged.volumes,
pvc_volumes: merged.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,
step_pvc_volumes: &config.pod.pvc_volumes,
managed_settings_configmap: merged.managed_settings_configmap.as_deref(),
claude_profiles: &self.claude_profiles,
annotations: Some(&annotations),
},
});
let created = match pod_spec {
Ok(mut pod_spec) => {
apply_active_deadline_seconds(&mut pod_spec, self.effective_deadline());
apply_runtime_class(&mut pod_spec, merged.runtime_class.as_deref());
pods.create(&PostParams::default(), &pod_spec)
.await
.map_err(|e| AgentError::ProcessFailed {
exit_code: -1,
stderr: format!("failed to create K8s pod: {e}"),
})
}
Err(e) => Err(e),
};
if let Err(e) = created {
self.revoke_proxy_token(proxy_token_id.as_deref()).await;
abort_launch_configmap(&configmaps, prompt_configmap_name.as_deref()).await;
return Err(e);
}
Ok(CreatedPod {
pods,
pod_name,
start,
prompt_configmap: prompt_configmap_name,
configmaps: Some(configmaps),
proxy_token_id,
})
}
async fn issue_proxy_token(
&self,
config: &AgentConfig,
run_id: Option<&String>,
step: Option<&String>,
pod_name: &str,
expires_at: u64,
) -> Result<Option<IssuedToken>, AgentError> {
let Some(client) = self.auth_proxy_client()? else {
return Ok(None);
};
let credential = resolve_credential(config.account.as_ref(), |k| var(k).ok());
let credential = credential.map_err(auth_proxy_error)?;
let request = TokenRequest {
run_id: run_id.map_or_else(|| pod_name.to_string(), String::clone),
step: step.map_or_else(|| "agent".to_string(), String::clone),
expires_at,
credential,
};
let issued = client.issue(&request).await.map_err(auth_proxy_error)?;
info!(
token = %issued.short_id(),
pod = %pod_name,
run_id = %request.run_id,
step = %request.step,
"auth proxy token issued"
);
Ok(Some(issued))
}
fn finalize_pod(
&self,
logs: &str,
pod_phase: &str,
timed_out: bool,
pod_name: &str,
config: &AgentConfig,
start: Instant,
) -> Result<AgentOutput, AgentError> {
rate_limit_event::record_rate_limits(config, logs);
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 K8sEphemeralProvider {
async fn run_created(
&self,
config: &AgentConfig,
created: &CreatedPod,
) -> Result<AgentOutput, AgentError> {
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)
}
async fn run_created_with_logs(
&self,
config: &AgentConfig,
created: &CreatedPod,
log_sink: Arc<dyn LogSink>,
) -> Result<AgentOutput, AgentError> {
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,
)
}
}
impl AgentProvider for K8sEphemeralProvider {
fn release_run<'a>(&'a self, run_id: &'a str) -> ReleaseFuture<'a> {
let limit = self.previous_attempt_timeout;
Box::pin(async move {
let cleanup = release_run(&self.cluster_config, &self.namespace, run_id, limit).await;
if let Ok(Some(client)) = self.auth_proxy_client()
&& let Err(e) = client.revoke_run(run_id).await
{
warn!(
run_id = %run_id,
error = %e,
"auth proxy token revocation of the run failed; they expire on their own"
);
}
cleanup
})
}
fn account_kind(&self) -> Option<&'static str> {
self.auth_proxy_url
.as_ref()
.map(|_| ClaudeSubscriptionKind::ID)
}
fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
Box::pin(async move {
let created = self.create_pod(config).await?;
let result = self.run_created(config, &created).await;
self.revoke_proxy_token(created.proxy_token_id.as_deref())
.await;
result
})
}
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 result = self.run_created_with_logs(config, &created, log_sink).await;
self.revoke_proxy_token(created.proxy_token_id.as_deref())
.await;
result
})
}
}
#[cfg(test)]
mod label_tests;
#[cfg(test)]
mod tests {
use serde_json::{Value, to_value};
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);
}
#[test]
fn k8s_runtime_class_default_is_none_and_builder_stores_value() {
assert!(K8sEphemeralProvider::new("img:v1").runtime_class.is_none());
let provider = K8sEphemeralProvider::new("img:v1").runtime_class("gvisor");
assert_eq!(provider.runtime_class.as_deref(), Some("gvisor"));
}
#[test]
#[should_panic(expected = "runtime class must not be empty")]
fn k8s_runtime_class_blank_builder_panics() {
let _ = K8sEphemeralProvider::new("img:v1").runtime_class(" ");
}
#[test]
fn k8s_merged_step_runtime_class_overrides_provider() {
let provider = K8sEphemeralProvider::new("img:v1").runtime_class("gvisor");
let step = AgentConfig::new("hi").runtime_class("kata");
assert_eq!(
merge(&provider, &step).runtime_class.as_deref(),
Some("kata")
);
assert_eq!(
merge(&provider, &AgentConfig::new("hi"))
.runtime_class
.as_deref(),
Some("gvisor")
);
let plain = K8sEphemeralProvider::new("img:v1");
assert_eq!(merge(&plain, &AgentConfig::new("hi")).runtime_class, None);
}
#[test]
fn k8s_merged_blank_step_runtime_class_is_an_error() {
let provider = K8sEphemeralProvider::new("img:v1");
let config = AgentConfig::new("hi").runtime_class(" ");
let err = err_text(provider.merged_pod_inputs(&config));
assert!(err.contains("runtime class must not be empty"), "{err}");
}
#[test]
fn k8s_apply_runtime_class_sets_runtime_class_name() {
let mut pod: Pod =
from_value(json!({"spec": {"containers": []}})).expect("valid minimal pod");
apply_runtime_class(&mut pod, Some("gvisor"));
assert_eq!(
pod.spec.unwrap().runtime_class_name.as_deref(),
Some("gvisor")
);
}
#[test]
fn k8s_apply_runtime_class_none_leaves_field_absent() {
let mut pod: Pod =
from_value(json!({"spec": {"containers": []}})).expect("valid minimal pod");
apply_runtime_class(&mut pod, None);
assert_eq!(pod.spec.unwrap().runtime_class_name, None);
}
fn err_text(result: Result<MergedPodInputs<'_>, AgentError>) -> String {
result.unwrap_err().to_string()
}
fn merge<'a>(provider: &'a K8sEphemeralProvider, config: &AgentConfig) -> MergedPodInputs<'a> {
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"));
}
fn provider_with_volumes() -> K8sEphemeralProvider {
K8sEphemeralProvider::sandboxed("img:v1")
.volume("/srv/work", "/data/work")
.pvc_volume("repos", "/data/repos")
}
#[test]
fn merged_provider_volumes_kept_without_flag() {
let provider = provider_with_volumes();
let config = AgentConfig::new("hi");
let merged = merge(&provider, &config);
assert_eq!(merged.volumes.len(), 1);
assert_eq!(merged.pvc_volumes.len(), 1);
}
#[test]
fn merged_without_provider_volumes_drops_them() {
let provider = provider_with_volumes();
let config = AgentConfig::new("hi").without_provider_volumes();
let merged = merge(&provider, &config);
assert!(merged.volumes.is_empty());
assert!(merged.pvc_volumes.is_empty());
}
#[test]
fn merged_step_pvc_volumes_keep_provider_ones() {
let provider = provider_with_volumes();
let config = AgentConfig::new("hi").pvc_volume("scratch", "/scratch", Some("a"), false);
let merged = merge(&provider, &config);
assert_eq!(
merged.pvc_volumes,
[("repos".to_string(), "/data/repos".to_string())]
);
assert_eq!(merged.volumes.len(), 1);
}
#[test]
fn pod_spec_without_provider_volumes_keeps_home_and_tmp() {
let provider = provider_with_volumes().read_only_pvc("ro-claim", "/data/ro");
let config = AgentConfig::new("hi")
.without_provider_volumes()
.pvc_volume("scratch", "/data/work", None, false);
let merged = merge(&provider, &config);
let pod = to_value(
build_pod_spec(&PodConfig {
name: "test-pod",
image: "img:v1",
command: vec!["sh".to_string()],
namespace: "default",
resources: &K8sResources::default(),
service_account: None,
restart_policy: "Never",
image_pull_policy: &ImagePullPolicy::default(),
env_vars: &[],
image_pull_secrets: &[],
extra_labels: &BTreeMap::new(),
node_selector: &BTreeMap::new(),
tolerations: &[],
volumes: merged.volumes,
pvc_volumes: merged.pvc_volumes,
inputs: &[],
input_init_image: DEFAULT_INPUT_INIT_IMAGE,
prompt_configmap: None,
prompt_mount_path: "",
hardening: PodHardening {
sandbox: provider.sandbox.as_ref(),
read_only_volumes: &merged.read_only_volumes,
step_pvc_volumes: &config.pod.pvc_volumes,
..PodHardening::default()
},
})
.unwrap(),
)
.unwrap();
let names: Vec<&str> = pod["spec"]["volumes"]
.as_array()
.unwrap()
.iter()
.map(|v| v["name"].as_str().unwrap())
.collect();
assert!(!names.contains(&"pvc-0"));
assert!(!names.contains(&"vol-0"));
assert!(names.contains(&"step-pvc-0"));
assert!(names.contains(&"ro-0"));
assert!(names.contains(&"ironflow-home"));
assert!(names.contains(&"ironflow-tmp"));
}
#[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));
}
const PROXY_URL: &str = "http://ironflow-auth-proxy.ironflow-system";
fn proxied() -> K8sEphemeralProvider {
K8sEphemeralProvider::sandboxed("img:v1").auth_proxy(PROXY_URL)
}
fn pod_env(
provider: &K8sEphemeralProvider,
config: &AgentConfig,
token: Option<&str>,
) -> Vec<Value> {
let merged = merge(provider, config);
let env_vars = provider.pod_env_vars(token);
let pod = build_pod_spec(&PodConfig {
name: "claude-code-test",
image: &provider.image,
command: vec!["sh".to_string()],
namespace: &provider.namespace,
resources: &provider.resources,
service_account: merged.service_account.as_deref(),
restart_policy: "Never",
image_pull_policy: &provider.image_pull_policy,
env_vars: &env_vars,
image_pull_secrets: &provider.image_pull_secrets,
extra_labels: &merged.labels,
node_selector: &provider.node_selector,
tolerations: &provider.tolerations,
volumes: &provider.volumes,
pvc_volumes: &provider.pvc_volumes,
inputs: &config.inputs,
input_init_image: DEFAULT_INPUT_INIT_IMAGE,
prompt_configmap: None,
prompt_mount_path: PROMPT_MOUNT_PATH,
hardening: PodHardening {
sandbox: provider.sandbox.as_ref(),
secret_env: &merged.secret_env,
read_only_volumes: &merged.read_only_volumes,
step_pvc_volumes: &config.pod.pvc_volumes,
managed_settings_configmap: None,
claude_profiles: &provider.claude_profiles,
annotations: None,
},
})
.unwrap();
let pod = to_value(&pod).unwrap();
pod["spec"]["containers"][0]["env"]
.as_array()
.cloned()
.unwrap_or_default()
}
fn env_value<'a>(env: &'a [Value], name: &str) -> Option<&'a str> {
env.iter()
.find(|entry| entry["name"] == name)
.and_then(|entry| entry["value"].as_str())
}
#[test]
fn auth_proxy_builder_sets_url_and_trims_slash() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.auth_proxy("http://ironflow-auth-proxy.ironflow-system/")
.auth_proxy_admin_key("0123456789abcdef0123456789abcdef");
assert_eq!(provider.auth_proxy_url.as_deref(), Some(PROXY_URL));
assert_eq!(
provider.auth_proxy_admin_key.as_deref(),
Some("0123456789abcdef0123456789abcdef")
);
assert!(provider.auth_proxy_client().unwrap().is_some());
let https = K8sEphemeralProvider::new("img:v1").auth_proxy("https://proxy");
assert_eq!(https.auth_proxy_url.as_deref(), Some("https://proxy"));
}
#[test]
#[should_panic(expected = "auth_proxy url must start with http:// or https://")]
fn auth_proxy_invalid_url_panics() {
let _ = K8sEphemeralProvider::sandboxed("img:v1").auth_proxy("ironflow-auth-proxy:80");
}
#[test]
fn auth_proxy_client_is_none_without_proxy() {
let provider = K8sEphemeralProvider::sandboxed("img:v1").auth_proxy_admin_key("k");
assert!(provider.auth_proxy_client().unwrap().is_none());
}
#[test]
fn auth_proxy_rejects_oauth_token_from_secret() {
let provider = proxied().oauth_token_from_secret("claude-oauth", "token");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("auth_proxy is set"), "{err}");
assert!(err.contains("CLAUDE_CODE_OAUTH_TOKEN"), "{err}");
}
#[test]
fn auth_proxy_rejects_step_secret_env_oauth_token() {
let config = AgentConfig::new("hi").env_from_secret("CLAUDE_CODE_OAUTH_TOKEN", "o", "k");
let err = err_text(proxied().merged_pod_inputs(&config));
assert!(err.contains("CLAUDE_CODE_OAUTH_TOKEN"), "{err}");
let config = AgentConfig::new("hi").env_from_secret("ANTHROPIC_API_KEY", "anthropic", "k");
let err = err_text(proxied().merged_pod_inputs(&config));
assert!(err.contains("ANTHROPIC_API_KEY"), "{err}");
}
#[test]
fn auth_proxy_rejects_inline_oauth_credentials() {
let json = r#"{"claudeAiOauth":{"accessToken":"sk-ant-oat01-x"}}"#;
let provider = proxied().oauth_credentials(json);
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("auth_proxy is set"), "{err}");
assert!(!err.contains("sk-ant"), "{err}");
let provider = proxied().oauth_credentials_from_secret("claude-credentials", "creds");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains(CREDENTIALS_ENV_VAR), "{err}");
let provider = K8sEphemeralProvider::new("img:v1")
.auth_proxy(PROXY_URL)
.oauth_credentials("{}");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("auth_proxy is set"), "{err}");
}
#[test]
fn auth_proxy_rejects_sk_ant_plain_value() {
let provider = K8sEphemeralProvider::new("img:v1")
.auth_proxy(PROXY_URL)
.env("SOME_VAR", "sk-ant-oat01-leak");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains("SOME_VAR"), "{err}");
assert!(!err.contains("sk-ant-oat01-leak"), "{err}");
for name in ["ANTHROPIC_BASE_URL", "ANTHROPIC_AUTH_TOKEN"] {
let provider = proxied().env(name, "x");
let err = err_text(provider.merged_pod_inputs(&AgentConfig::new("hi")));
assert!(err.contains(name), "{err}");
}
}
#[test]
fn auth_proxy_pod_env_has_base_url_and_opaque_token_only() {
let provider = proxied().env("TEAM", "infra");
let config = AgentConfig::new("hi").run_scope("run-1", "review");
let env = pod_env(&provider, &config, Some("ifap_x"));
assert_eq!(env_value(&env, "ANTHROPIC_BASE_URL"), Some(PROXY_URL));
assert_eq!(env_value(&env, "ANTHROPIC_AUTH_TOKEN"), Some("ifap_x"));
assert_eq!(
env_value(&env, "CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC"),
Some("1")
);
assert_eq!(env_value(&env, "TEAM"), Some("infra"));
assert!(
env.iter().all(|entry| entry.get("valueFrom").is_none()),
"no secretKeyRef expected: {env:?}"
);
for name in ["CLAUDE_CODE_OAUTH_TOKEN", "ANTHROPIC_API_KEY"] {
assert!(
env_value(&env, name).is_none_or(str::is_empty),
"{name} must not carry a value: {env:?}"
);
}
let rendered = Value::Array(env).to_string();
assert!(!rendered.contains("sk-ant"), "{rendered}");
}
#[test]
fn auth_proxy_account_kind_is_subscription_only_with_proxy() {
assert_eq!(proxied().account_kind(), Some(ClaudeSubscriptionKind::ID));
assert_eq!(
K8sEphemeralProvider::sandboxed("img:v1").account_kind(),
None
);
assert_eq!(K8sEphemeralProvider::new("img:v1").account_kind(), None);
}
#[test]
fn non_proxy_pod_env_unchanged() {
let provider = K8sEphemeralProvider::sandboxed("img:v1")
.env("TEAM", "infra")
.oauth_token_from_secret("claude-oauth", "token");
assert_eq!(
provider.pod_env_vars(Some("ifap_x")),
vec![("TEAM".to_string(), "infra".to_string())]
);
assert_eq!(proxied().pod_env_vars(None), Vec::new());
let env = pod_env(&provider, &AgentConfig::new("hi"), None);
assert!(env_value(&env, "ANTHROPIC_BASE_URL").is_none());
assert!(env_value(&env, "ANTHROPIC_AUTH_TOKEN").is_none());
let token = env
.iter()
.find(|entry| entry["name"] == "CLAUDE_CODE_OAUTH_TOKEN")
.expect("token entry");
assert_eq!(
token["valueFrom"]["secretKeyRef"]["name"],
Value::from("claude-oauth")
);
}
}