use crate::error::{EngineError, Result};
use crate::sandbox_container::{self, ContainerRuntime};
use crate::workspace_contract::{PortPolicy, PreviewSpec, WorkspaceContract};
use crate::workspace_gate::{
CommandOutcome, GatePhase, BOOTSTRAP_SUMMARY_PREFIX, READINESS_SUMMARY_PREFIX,
};
use crate::workspace_provider::{
report_gate_outcomes, PreviewPlaceholder, ProgressSink, ProvisionSpec, ReadinessOutcome,
TeardownMode, WorkspaceHandle, WorkspaceProvider, WorkspaceProviderKind,
};
use serde_json::json;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
const PROJECT_PREFIX: &str = "kranz-ws-";
const WORKSPACE_SERVICE: &str = "workspace";
const DYNAMIC_CONTAINER_PORT: u16 = 8080;
const COMPOSE_FILE_NAME: &str = "compose.json";
const WORKSPACE_DIR: &str = "workspace";
const COMPOSE_PROBE_TIMEOUT: Duration = Duration::from_secs(30);
const COMPOSE_UP_TIMEOUT: Duration = Duration::from_secs(600);
const EXEC_TIMEOUT: Duration = Duration::from_secs(600);
const HEALTH_CHECK_TIMEOUT: Duration = Duration::from_secs(120);
const HEALTH_POLL_INTERVAL: Duration = Duration::from_secs(1);
const HEALTH_EXEC_TIMEOUT: Duration = Duration::from_secs(30);
const TEARDOWN_TIMEOUT: Duration = Duration::from_secs(300);
const OUTPUT_TAIL: usize = 1500;
#[derive(Debug, Clone)]
pub struct ContainerWorkspace {
pub runtime: ContainerRuntime,
pub project: String,
pub compose_file: PathBuf,
pub assigned_ports: Vec<(String, u16)>,
}
type DetectHook = Arc<dyn Fn() -> Option<ContainerRuntime> + Send + Sync>;
type RunHook = Arc<dyn Fn(&[String], Duration) -> (Option<i32>, String) + Send + Sync>;
#[derive(Clone)]
pub(crate) struct RuntimeHooks {
detect: DetectHook,
run: RunHook,
}
impl RuntimeHooks {
fn host() -> Self {
Self {
detect: Arc::new(sandbox_container::detect),
run: Arc::new(spawn_bounded),
}
}
}
fn spawn_bounded(argv: &[String], timeout: Duration) -> (Option<i32>, String) {
let Some((program, args)) = argv.split_first() else {
return (None, "empty argv".to_string());
};
match crate::command_exec::run_with_timeout(Path::new(program), args, timeout) {
Some(output) => {
let mut text = String::from_utf8_lossy(&output.stdout).into_owned();
let stderr = String::from_utf8_lossy(&output.stderr);
if !stderr.is_empty() {
if !text.is_empty() {
text.push('\n');
}
text.push_str(&stderr);
}
(
output.status.code(),
crate::command_exec::last_chars_local(&text, OUTPUT_TAIL),
)
}
None => (
None,
format!(
"no exit code (spawn failure or timeout after {}s)",
timeout.as_secs()
),
),
}
}
pub struct LocalContainerProvider {
hooks: RuntimeHooks,
}
impl LocalContainerProvider {
pub fn new() -> Self {
Self {
hooks: RuntimeHooks::host(),
}
}
#[cfg(test)]
pub(crate) fn with_hooks(hooks: RuntimeHooks) -> Self {
Self { hooks }
}
async fn run_argv(&self, argv: &[String], timeout: Duration) -> (Option<i32>, String) {
let run = self.hooks.run.clone();
let argv = argv.to_vec();
tokio::task::spawn_blocking(move || run(&argv, timeout))
.await
.unwrap_or_else(|e| (None, format!("runtime invocation failed to complete: {e}")))
}
async fn cleanup_project(&self, runtime: ContainerRuntime, project: &str, compose_file: &Path) {
let argv = compose_argv(runtime, project, compose_file, &["down", "-v"]);
let _ = self.run_argv(&argv, TEARDOWN_TIMEOUT).await;
}
}
impl Default for LocalContainerProvider {
fn default() -> Self {
Self::new()
}
}
pub(crate) fn compose_project_name(mission_id: &str) -> String {
let sanitized: String = mission_id
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() {
c.to_ascii_lowercase()
} else if c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
if sanitized.is_empty() {
format!("{PROJECT_PREFIX}mission")
} else {
format!("{PROJECT_PREFIX}{sanitized}")
}
}
fn compose_argv(
runtime: ContainerRuntime,
project: &str,
compose_file: &Path,
args: &[&str],
) -> Vec<String> {
let mut argv = vec![
runtime.binary().to_string(),
"compose".to_string(),
"-p".to_string(),
project.to_string(),
"-f".to_string(),
compose_file.display().to_string(),
];
argv.extend(args.iter().map(|arg| (*arg).to_string()));
argv
}
fn exec_argv(
runtime: ContainerRuntime,
project: &str,
compose_file: &Path,
service: &str,
command: &str,
) -> Vec<String> {
compose_argv(
runtime,
project,
compose_file,
&["exec", "-T", service, "sh", "-c", command],
)
}
fn code_phrase(code: Option<i32>) -> String {
match code {
Some(code) => format!("exit code {code}"),
None => "no exit code (spawn failure or timeout)".to_string(),
}
}
fn check_fixed_port_collisions(contract: &WorkspaceContract) -> Result<()> {
for service in &contract.services {
if let PortPolicy::Fixed(port) = service.port.policy {
if std::net::TcpListener::bind((std::net::Ipv4Addr::UNSPECIFIED, port)).is_err() {
return Err(EngineError::InvalidState(format!(
"workspace.provider \"container\": fixed port {port} for service {:?} is \
already bound on the host; refusing rather than silently rebinding \
(owner: repo-setup — free the port or pick another fixed port)",
service.name
)));
}
}
}
Ok(())
}
fn bind_mount(source: &str, target: &str) -> serde_json::Value {
json!({ "type": "bind", "source": source, "target": target })
}
pub(crate) fn render_compose_file(
project: &str,
contract: &WorkspaceContract,
worktree: &Path,
env: &HashMap<String, String>,
) -> serde_json::Value {
let worktree_abs = crate::sandbox::absolutize(worktree).display().to_string();
let mut volumes = vec![bind_mount(&worktree_abs, &worktree_abs)];
for mount in &contract.mounts {
volumes.push(bind_mount(mount, mount));
}
let mut workspace = json!({
"image": sandbox_container::DEFAULT_IMAGE,
"command": ["sleep", "infinity"],
"working_dir": worktree_abs,
"volumes": volumes,
});
if !env.is_empty() {
workspace["environment"] = json!(env);
}
let mut services = serde_json::Map::new();
services.insert(WORKSPACE_SERVICE.to_string(), workspace);
for service in &contract.services {
let mut def = json!({
"image": sandbox_container::DEFAULT_IMAGE,
"command": ["sh", "-c", service.start],
});
let ports = match &service.port.policy {
PortPolicy::Dynamic => vec![format!("0:{DYNAMIC_CONTAINER_PORT}")],
PortPolicy::Fixed(port) => vec![format!("{port}:{port}")],
};
def["ports"] = json!(ports);
if let Some(health_check) = &service.health_check {
def["healthcheck"] = json!({
"test": ["CMD-SHELL", health_check],
"interval": "2s",
"timeout": "10s",
"retries": 30,
"start_period": "5s",
});
}
services.insert(service.name.clone(), def);
}
json!({
"name": project,
"services": services,
"x-kranz": {
"workspaceService": WORKSPACE_SERVICE,
"bootstrap": contract.bootstrap,
"readiness": contract.readiness,
},
})
}
fn parse_compose_port(output: &str) -> Option<u16> {
output
.lines()
.filter_map(|line| line.trim().rsplit(':').next())
.find_map(|segment| segment.parse::<u16>().ok())
}
fn fill_previews(
previews: &[PreviewSpec],
assigned_ports: &[(String, u16)],
) -> Vec<PreviewPlaceholder> {
previews
.iter()
.map(|preview| {
let url_template = match assigned_ports {
[(_, port)] => preview.url_template.replace("{port}", &port.to_string()),
_ => preview.url_template.clone(),
};
PreviewPlaceholder {
name: preview.name.clone(),
url_template,
}
})
.collect()
}
#[async_trait::async_trait]
impl WorkspaceProvider for LocalContainerProvider {
fn kind(&self) -> WorkspaceProviderKind {
WorkspaceProviderKind::Container
}
async fn provision(&self, spec: &ProvisionSpec) -> Result<WorkspaceHandle> {
if !spec.repo_root.is_dir() {
return Err(EngineError::InvalidState(format!(
"container provision: execution cwd {} does not exist",
spec.repo_root.display()
)));
}
let env = crate::runner::contract_env(spec.base_sha.as_deref());
let Some(contract) = &spec.contract else {
return Ok(WorkspaceHandle {
cwd: spec.repo_root.clone(),
env,
previews: Vec::new(),
contract: None,
detail: None,
container: None,
remote: None,
gate_env: spec.gate_env.clone(),
});
};
check_fixed_port_collisions(contract)?;
let runtime = (self.hooks.detect)().ok_or_else(|| {
EngineError::Config(
"workspace.provider \"container\" needs a container runtime on PATH (docker, \
podman, or nerdctl — none found; owner: operator — install a runtime or choose \
another workspace.provider); refusing rather than sharing the host port namespace"
.to_string(),
)
})?;
if runtime == ContainerRuntime::AppleContainer {
return Err(EngineError::Config(
"workspace.provider \"container\": the `container` runtime (Apple Container) has \
no compose subcommand; install docker, podman, or nerdctl (owner: operator) or \
choose another workspace.provider"
.to_string(),
));
}
let probe = self
.run_argv(
&[
runtime.binary().to_string(),
"compose".to_string(),
"version".to_string(),
],
COMPOSE_PROBE_TIMEOUT,
)
.await;
if probe.0 != Some(0) {
return Err(EngineError::Config(format!(
"workspace.provider \"container\": `{} compose` is unavailable ({}); the \
local-container provider needs a compose subcommand — install the compose \
plugin (owner: operator) or choose another workspace.provider",
runtime.binary(),
crate::scrub::scrub(&probe.1)
)));
}
let project = compose_project_name(&spec.mission_id);
let workspace_dir = spec.runtime_dir.join(WORKSPACE_DIR);
std::fs::create_dir_all(&workspace_dir)?;
let compose_file = workspace_dir.join(COMPOSE_FILE_NAME);
let doc = render_compose_file(&project, contract, &spec.repo_root, &env);
std::fs::write(&compose_file, serde_json::to_vec_pretty(&doc)?)?;
let up = self
.run_argv(
&compose_argv(runtime, &project, &compose_file, &["up", "-d"]),
COMPOSE_UP_TIMEOUT,
)
.await;
if up.0 != Some(0) {
self.cleanup_project(runtime, &project, &compose_file).await;
return Err(EngineError::InvalidState(format!(
"container provision: `{} compose up -d` failed ({}): {}",
runtime.binary(),
code_phrase(up.0),
crate::scrub::scrub(&up.1)
)));
}
for service in contract
.services
.iter()
.filter(|service| service.health_check.is_some())
{
let health_check = service.health_check.as_deref().expect("filtered on Some");
let deadline = std::time::Instant::now() + HEALTH_CHECK_TIMEOUT;
loop {
let argv = exec_argv(
runtime,
&project,
&compose_file,
&service.name,
health_check,
);
let (code, _tail) = self.run_argv(&argv, HEALTH_EXEC_TIMEOUT).await;
if code == Some(0) {
break;
}
if std::time::Instant::now() >= deadline {
self.cleanup_project(runtime, &project, &compose_file).await;
return Err(EngineError::InvalidState(format!(
"container provision: health check for service {:?} did not pass within \
{}s: `{health_check}`",
service.name,
HEALTH_CHECK_TIMEOUT.as_secs()
)));
}
tokio::time::sleep(HEALTH_POLL_INTERVAL).await;
}
}
let mut assigned_ports = Vec::new();
for service in contract
.services
.iter()
.filter(|service| matches!(service.port.policy, PortPolicy::Dynamic))
{
let argv = compose_argv(
runtime,
&project,
&compose_file,
&["port", &service.name, &DYNAMIC_CONTAINER_PORT.to_string()],
);
let (code, output) = self.run_argv(&argv, COMPOSE_PROBE_TIMEOUT).await;
let port = if code == Some(0) {
parse_compose_port(&output)
} else {
None
};
let Some(port) = port else {
self.cleanup_project(runtime, &project, &compose_file).await;
return Err(EngineError::InvalidState(format!(
"container provision: could not read back the OS-assigned host port for \
dynamic service {:?} (`compose port` {}): {}",
service.name,
code_phrase(code),
crate::scrub::scrub(&output)
)));
};
assigned_ports.push((service.name.clone(), port));
}
Ok(WorkspaceHandle {
cwd: spec.repo_root.clone(),
env,
previews: fill_previews(&contract.previews, &assigned_ports),
contract: Some(contract.clone()),
detail: Some(format!("compose project {project}")),
container: Some(ContainerWorkspace {
runtime,
project,
compose_file,
assigned_ports,
}),
remote: None,
gate_env: spec.gate_env.clone(),
})
}
async fn readiness(
&self,
handle: &WorkspaceHandle,
progress: &mut ProgressSink<'_>,
) -> Result<ReadinessOutcome> {
let Some(contract) = &handle.contract else {
return Ok(ReadinessOutcome::Ready);
};
let Some(workspace) = &handle.container else {
return Err(EngineError::InvalidState(
"container readiness: the handle carries a contract but no compose project — \
provision did not complete"
.to_string(),
));
};
if let Some(data) = &contract.data {
for (hook, command) in [
(crate::workspace_data::DataHookKind::Clone, &data.clone),
(crate::workspace_data::DataHookKind::Migrate, &data.migrate),
] {
if let Some(command) = command {
if let Some(failed) =
self.run_data_hook(handle, hook, command, progress).await?
{
return Ok(ReadinessOutcome::Failed {
kind: hook.gate_kind(),
failed,
});
}
}
}
}
let phase = GatePhase {
kind: "bootstrap command",
unit: "command",
plural: "commands",
prefix: BOOTSTRAP_SUMMARY_PREFIX,
commands: &contract.bootstrap,
stop_at_first_failure: true,
};
if let Some(failed) = self.run_exec_phase(workspace, &phase, progress).await? {
return Ok(ReadinessOutcome::Failed {
kind: "bootstrap command",
failed,
});
}
let phase = GatePhase {
kind: "readiness check",
unit: "check",
plural: "checks",
prefix: READINESS_SUMMARY_PREFIX,
commands: &contract.readiness,
stop_at_first_failure: false,
};
if let Some(failed) = self.run_exec_phase(workspace, &phase, progress).await? {
return Ok(ReadinessOutcome::Failed {
kind: "readiness check",
failed,
});
}
if let Some(command) = contract
.data
.as_ref()
.and_then(|data| data.skew_check.as_ref())
{
if let Some(failed) = self
.run_data_hook(
handle,
crate::workspace_data::DataHookKind::SkewCheck,
command,
progress,
)
.await?
{
return Ok(ReadinessOutcome::DataSkew { failed });
}
}
Ok(ReadinessOutcome::Ready)
}
async fn run_data_hook(
&self,
handle: &WorkspaceHandle,
hook: crate::workspace_data::DataHookKind,
command: &str,
progress: &mut ProgressSink<'_>,
) -> Result<Option<CommandOutcome>> {
let Some(workspace) = &handle.container else {
return Err(EngineError::InvalidState(
"container data hook: the handle carries no compose project — provision did \
not complete"
.to_string(),
));
};
let argv = exec_argv(
workspace.runtime,
&workspace.project,
&workspace.compose_file,
WORKSPACE_SERVICE,
command,
);
let (code, output_tail) = self.run_argv(&argv, EXEC_TIMEOUT).await;
crate::workspace_data::hook_outcome(hook, command, code, output_tail, progress)
}
async fn teardown(&self, handle: WorkspaceHandle, mode: TeardownMode) -> Result<()> {
let Some(workspace) = &handle.container else {
return Ok(()); };
match mode {
TeardownMode::Keep => Ok(()),
TeardownMode::Hibernate => {
let (code, tail) = self
.run_argv(
&compose_argv(
workspace.runtime,
&workspace.project,
&workspace.compose_file,
&["stop"],
),
TEARDOWN_TIMEOUT,
)
.await;
if code == Some(0) {
Ok(())
} else {
Err(EngineError::InvalidState(format!(
"container teardown (hibernate): `compose stop` failed ({}): {}",
code_phrase(code),
crate::scrub::scrub(&tail)
)))
}
}
TeardownMode::Destroy => {
let (code, tail) = self
.run_argv(
&compose_argv(
workspace.runtime,
&workspace.project,
&workspace.compose_file,
&["down", "-v"],
),
TEARDOWN_TIMEOUT,
)
.await;
if code != Some(0) {
return Err(EngineError::InvalidState(format!(
"container teardown (destroy): `compose down -v` failed ({}): {}",
code_phrase(code),
crate::scrub::scrub(&tail)
)));
}
if let Some(prune) = handle
.contract
.as_ref()
.and_then(|contract| contract.disk.as_ref())
.and_then(|disk| disk.prune.as_deref())
{
let cwd = workspace
.compose_file
.parent()
.unwrap_or(handle.cwd.as_path());
let env = crate::workspace_gate::gate_command_env(
&handle.gate_env,
&handle.env,
handle.contract.as_ref(),
);
let (code, tail) =
crate::command_exec::run_shell_command_with_code_cleared(cwd, prune, &env)
.await;
if code != Some(0) {
return Err(EngineError::InvalidState(format!(
"container teardown (destroy): disk prune command `{prune}` failed \
({}): {}",
code_phrase(code),
crate::scrub::scrub(&tail)
)));
}
}
Ok(())
}
}
}
}
impl LocalContainerProvider {
async fn run_exec_phase(
&self,
workspace: &ContainerWorkspace,
phase: &GatePhase<'_>,
progress: &mut ProgressSink<'_>,
) -> Result<Option<CommandOutcome>> {
progress(
&format!(
"{} running {} {}",
phase.prefix,
phase.commands.len(),
phase.plural
),
None,
)?;
let total = phase.commands.len();
let mut outcomes = Vec::with_capacity(total);
for (i, command) in phase.commands.iter().enumerate() {
let argv = exec_argv(
workspace.runtime,
&workspace.project,
&workspace.compose_file,
WORKSPACE_SERVICE,
command,
);
let (code, output_tail) = self.run_argv(&argv, EXEC_TIMEOUT).await;
let outcome = CommandOutcome {
ordinal: i + 1,
total,
command: command.clone(),
code,
output_tail,
};
let failed = !outcome.ok();
outcomes.push(outcome);
if failed && phase.stop_at_first_failure {
break;
}
}
report_gate_outcomes(phase, outcomes, progress)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::workspace_contract::parse_workspace_contract;
use std::sync::Mutex;
fn spec(root: &Path, mission_id: &str, contract: Option<WorkspaceContract>) -> ProvisionSpec {
let runtime_dir = root.join(".kranz").join("missions").join(mission_id);
ProvisionSpec {
mission_id: mission_id.to_string(),
repo_root: root.to_path_buf(),
gate_env: crate::workspace_provider::GateEnvPolicy::for_mission(&runtime_dir, &[]),
runtime_dir,
base_sha: Some("deadbeefcafe".to_string()),
contract,
}
}
fn contract(json: &[u8]) -> WorkspaceContract {
parse_workspace_contract(json).expect("valid contract")
}
fn minimal_contract() -> WorkspaceContract {
contract(br#"{"schemaVersion": 1, "readiness": ["true"]}"#)
}
fn fake_contract() -> WorkspaceContract {
contract(
br#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"services": [
{
"name": "api",
"start": "sleep infinity",
"healthCheck": "true",
"port": { "policy": "dynamic" }
}
],
"readiness": ["test -f .boot-marker"],
"previews": [{ "name": "app", "urlTemplate": "http://localhost:{port}/" }],
"disk": { "prune": "echo pruned > prune-marker.txt" }
}"#,
)
}
#[derive(Default)]
struct Progress(Vec<(String, Option<String>)>);
impl Progress {
fn sink(&mut self) -> impl FnMut(&str, Option<String>) -> Result<()> + Send + use<'_> {
|summary, detail| {
self.0.push((summary.to_string(), detail));
Ok(())
}
}
fn summaries(&self) -> Vec<&str> {
self.0.iter().map(|(s, _)| s.as_str()).collect()
}
}
#[derive(Default)]
struct FakeRuntime {
calls: Mutex<Vec<Vec<String>>>,
fail_exec_containing: Option<String>,
}
impl FakeRuntime {
fn hooks(self: &Arc<Self>) -> RuntimeHooks {
let this = Arc::clone(self);
RuntimeHooks {
detect: Arc::new(|| Some(ContainerRuntime::Docker)),
run: Arc::new(move |argv, _timeout| this.answer(argv)),
}
}
fn answer(&self, argv: &[String]) -> (Option<i32>, String) {
self.calls.lock().unwrap().push(argv.to_vec());
let args: Vec<&str> = argv.iter().map(String::as_str).collect();
if args == ["docker", "compose", "version"] {
return (Some(0), "v2.27.0".to_string());
}
let rest = &args[6..];
match rest[0] {
"up" | "stop" | "down" | "ps" => (Some(0), String::new()),
"port" => (Some(0), "0.0.0.0:32768\n".to_string()),
"exec" => {
let command = rest[5];
if let Some(needle) = &self.fail_exec_containing {
if command.contains(needle.as_str()) {
return (Some(3), format!("boom running `{command}`"));
}
}
(Some(0), String::new())
}
other => (
Some(1),
format!("fake runtime: unexpected subcommand {other:?}"),
),
}
}
fn calls(&self) -> Vec<Vec<String>> {
self.calls.lock().unwrap().clone()
}
fn called_with(&self, suffix: &[&str]) -> bool {
self.calls().iter().any(|argv| {
let args: Vec<&str> = argv.iter().map(String::as_str).collect();
args.ends_with(suffix)
})
}
fn execed(&self, service: &str, command: &str) -> bool {
self.calls().iter().any(|argv| {
let args: Vec<&str> = argv.iter().map(String::as_str).collect();
args.windows(6)
.any(|w| w == ["exec", "-T", service, "sh", "-c", command])
})
}
}
#[test]
fn compose_project_name_is_sanitized_and_mission_owned() {
assert_eq!(compose_project_name("m-abc123"), "kranz-ws-m-abc123");
assert_eq!(
compose_project_name("m-Foo_Bar/baz.qux"),
"kranz-ws-m-foo_bar-baz-qux"
);
assert_eq!(compose_project_name(""), "kranz-ws-mission");
assert_ne!(
compose_project_name("m-a"),
compose_project_name("m-b"),
"parallel missions never share a project"
);
}
#[test]
fn render_compose_file_maps_contract_to_services_ports_and_volumes() {
let dir = tempfile::tempdir().expect("tempdir");
let contract = contract(
br#"{
"schemaVersion": 1,
"bootstrap": ["cargo fetch"],
"services": [
{
"name": "api",
"start": "./run-api",
"healthCheck": "curl -sf localhost:8080/health",
"port": { "policy": "dynamic" }
},
{
"name": "db",
"start": "./run-db",
"port": { "policy": { "fixed": 5432 } }
}
],
"readiness": ["curl -sf localhost:8080/health"],
"mounts": ["/var/cache/cargo"],
"previews": [{ "name": "app", "urlTemplate": "http://localhost:{port}/" }]
}"#,
);
let env = crate::runner::contract_env(Some("deadbeefcafe"));
let doc = render_compose_file("kranz-ws-m-render1", &contract, dir.path(), &env);
assert_eq!(doc["name"], "kranz-ws-m-render1");
let abs = crate::sandbox::absolutize(dir.path()).display().to_string();
let workspace = &doc["services"][WORKSPACE_SERVICE];
assert_eq!(workspace["image"], sandbox_container::DEFAULT_IMAGE);
assert_eq!(workspace["command"], json!(["sleep", "infinity"]));
assert_eq!(workspace["working_dir"], json!(abs));
assert_eq!(workspace["environment"]["KRANZ_BASE_SHA"], "deadbeefcafe");
let volumes = workspace["volumes"].as_array().expect("volumes array");
assert!(
volumes.contains(&bind_mount(&abs, &abs)),
"worktree bind mount: {volumes:?}"
);
assert!(
volumes.contains(&bind_mount("/var/cache/cargo", "/var/cache/cargo")),
"declared mounts[] become volumes: {volumes:?}"
);
assert_eq!(
doc["services"]["api"]["ports"],
json!([format!("0:{DYNAMIC_CONTAINER_PORT}")])
);
assert_eq!(doc["services"]["db"]["ports"], json!(["5432:5432"]));
assert_eq!(
doc["services"]["api"]["healthcheck"]["test"],
json!(["CMD-SHELL", "curl -sf localhost:8080/health"])
);
assert!(doc["services"]["db"].get("healthcheck").is_none());
assert_eq!(
doc["services"]["api"]["command"],
json!(["sh", "-c", "./run-api"])
);
assert_eq!(doc["x-kranz"]["bootstrap"], json!(["cargo fetch"]));
assert_eq!(
doc["x-kranz"]["readiness"],
json!(["curl -sf localhost:8080/health"])
);
}
#[tokio::test]
async fn fixed_port_collision_refuses_at_provision_naming_service_and_port() {
let probe = std::net::TcpListener::bind((std::net::Ipv4Addr::UNSPECIFIED, 0))
.expect("bind probe socket");
let port = probe.local_addr().expect("local addr").port();
let dir = tempfile::tempdir().expect("tempdir");
let contract_json = format!(
r#"{{"schemaVersion": 1, "services": [
{{"name": "db", "start": "./run-db", "port": {{"policy": {{"fixed": {port}}}}}}}
]}}"#
);
let err = LocalContainerProvider::new()
.provision(&spec(
dir.path(),
"m-collision",
Some(contract(contract_json.as_bytes())),
))
.await
.expect_err("a host-bound fixed port refuses provision");
let msg = err.to_string();
assert!(msg.contains("\"db\""), "{msg}");
assert!(msg.contains(&port.to_string()), "{msg}");
assert!(msg.contains("refusing"), "{msg}");
assert!(
!dir.path()
.join(".kranz/missions/m-collision/workspace/compose.json")
.exists(),
"the refusal precedes any runtime work: no compose file written"
);
}
#[tokio::test]
async fn provision_fails_closed_when_no_runtime_is_detected() {
let dir = tempfile::tempdir().expect("tempdir");
let provider = LocalContainerProvider::with_hooks(RuntimeHooks {
detect: Arc::new(|| None),
run: Arc::new(|_argv, _timeout| {
panic!("no runtime invocation may happen without a detected runtime")
}),
});
let err = provider
.provision(&spec(dir.path(), "m-noruntime", Some(minimal_contract())))
.await
.expect_err("no runtime ⇒ fail closed");
let msg = err.to_string();
assert!(msg.contains("workspace.provider"), "{msg}");
assert!(msg.contains("docker"), "{msg}");
assert!(msg.contains("podman"), "{msg}");
assert!(msg.contains("nerdctl"), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
}
#[tokio::test]
async fn provision_fails_closed_on_a_runtimeless_host() {
if sandbox_container::detect().is_some() {
eprintln!(
"host has a container runtime; skipping the real-detection refusal \
(covered deterministically by the injected-detection test)"
);
return;
}
let dir = tempfile::tempdir().expect("tempdir");
let err = LocalContainerProvider::new()
.provision(&spec(
dir.path(),
"m-noruntime-host",
Some(minimal_contract()),
))
.await
.expect_err("a runtime-less host fails closed at provision");
let msg = err.to_string();
assert!(msg.contains("workspace.provider"), "{msg}");
assert!(msg.contains("docker"), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
}
#[tokio::test]
async fn provision_writes_compose_goes_up_and_reads_back_dynamic_ports() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-Fake1", Some(fake_contract())))
.await
.expect("provision");
assert_eq!(provider.kind(), WorkspaceProviderKind::Container);
assert_eq!(handle.cwd, dir.path());
assert_eq!(
handle.env.get("KRANZ_BASE_SHA").map(String::as_str),
Some("deadbeefcafe")
);
assert_eq!(
handle.detail.as_deref(),
Some("compose project kranz-ws-m-fake1"),
"the provisioned event's detail carries the compose project"
);
let workspace = handle.container.as_ref().expect("container state");
assert_eq!(workspace.project, "kranz-ws-m-fake1");
assert_eq!(
workspace.assigned_ports,
vec![("api".to_string(), 32768)],
"the OS-assigned port is read back from the runtime"
);
assert_eq!(
handle.previews,
vec![PreviewPlaceholder {
name: "app".to_string(),
url_template: "http://localhost:32768/".to_string(),
}]
);
let compose_file = dir
.path()
.join(".kranz/missions/m-Fake1/workspace/compose.json");
let doc: serde_json::Value =
serde_json::from_slice(&std::fs::read(&compose_file).expect("compose file written"))
.expect("compose file is JSON");
assert_eq!(doc["name"], "kranz-ws-m-fake1");
assert_eq!(workspace.compose_file, compose_file);
assert!(fake.called_with(&["up", "-d"]), "{:?}", fake.calls());
assert!(fake.execed("api", "true"), "{:?}", fake.calls());
}
#[tokio::test]
async fn readiness_execs_bootstrap_then_readiness_inside_the_container_network() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fake2", Some(fake_contract())))
.await
.expect("provision");
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(matches!(outcome, ReadinessOutcome::Ready), "{outcome:?}");
assert!(
fake.execed(WORKSPACE_SERVICE, "echo boot > .boot-marker"),
"bootstrap exec'd inside the workspace container: {:?}",
fake.calls()
);
assert!(
fake.execed(WORKSPACE_SERVICE, "test -f .boot-marker"),
"readiness exec'd inside the workspace container: {:?}",
fake.calls()
);
assert_eq!(
progress.summaries(),
vec![
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok",
"workspace readiness: running 1 checks",
"workspace readiness: 1/1 checks ok",
],
"the gate's decision lines are byte-identical to the host path"
);
}
#[tokio::test]
async fn readiness_failure_reports_the_gate_outcome() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime {
fail_exec_containing: Some("test -f".to_string()),
..FakeRuntime::default()
});
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fake3", Some(fake_contract())))
.await
.expect("provision (health checks do not match the failure needle)");
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
let ReadinessOutcome::Failed { kind, failed } = outcome else {
panic!("readiness failure must be Failed, got {outcome:?}");
};
assert_eq!(kind, "readiness check");
assert_eq!(failed.code, Some(3));
assert_eq!(
progress.summaries(),
vec![
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok",
"workspace readiness: running 1 checks",
"workspace readiness: FAILED at check 1/1 — blocking mission (owner: repo-setup)",
]
);
}
fn fake_data_contract() -> WorkspaceContract {
contract(
br#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .boot-marker"],
"data": {
"clone": "clone-golden",
"migrate": "migrate-golden",
"reset": "reseed-golden",
"skewCheck": "check-skew",
"resetBetweenRounds": true
}
}"#,
)
}
fn exec_index(fake: &FakeRuntime, service: &str, command: &str) -> usize {
fake.calls()
.iter()
.position(|argv| {
let args: Vec<&str> = argv.iter().map(String::as_str).collect();
args.windows(6)
.any(|w| w == ["exec", "-T", service, "sh", "-c", command])
})
.unwrap_or_else(|| panic!("no exec of {command:?} in {service}: {:?}", fake.calls()))
}
#[tokio::test]
async fn data_hooks_exec_inside_the_container_network_in_lifecycle_order() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fakedata", Some(fake_data_contract())))
.await
.expect("provision");
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(matches!(outcome, ReadinessOutcome::Ready), "{outcome:?}");
let clone = exec_index(&fake, WORKSPACE_SERVICE, "clone-golden");
let migrate = exec_index(&fake, WORKSPACE_SERVICE, "migrate-golden");
let bootstrap = exec_index(&fake, WORKSPACE_SERVICE, "echo boot > .boot-marker");
let readiness = exec_index(&fake, WORKSPACE_SERVICE, "test -f .boot-marker");
let skew = exec_index(&fake, WORKSPACE_SERVICE, "check-skew");
assert!(
clone < migrate && migrate < bootstrap && bootstrap < readiness && readiness < skew,
"lifecycle order (clone<{clone} migrate<{migrate} bootstrap<{bootstrap} readiness<{readiness} skew<{skew})"
);
assert_eq!(
progress.summaries(),
vec![
"workspace data: clone `clone-golden` → ok (exit code 0)",
"workspace data: migrate `migrate-golden` → ok (exit code 0)",
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok",
"workspace readiness: running 1 checks",
"workspace readiness: 1/1 checks ok",
"workspace data: skewCheck `check-skew` → ok (exit code 0)",
],
"the data decision lines are byte-identical to the host path"
);
let mut progress = Progress::default();
let failed = provider
.run_data_hook(
&handle,
crate::workspace_data::DataHookKind::Reset,
"reseed-golden",
&mut progress.sink(),
)
.await
.expect("run_data_hook");
assert!(failed.is_none(), "{failed:?}");
assert!(
fake.execed(WORKSPACE_SERVICE, "reseed-golden"),
"reset exec'd inside the workspace container: {:?}",
fake.calls()
);
assert_eq!(
progress.summaries(),
vec!["workspace data: reset `reseed-golden` → ok (exit code 0)"]
);
}
#[tokio::test]
async fn container_skew_failure_is_the_distinct_skew_outcome() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime {
fail_exec_containing: Some("check-skew".to_string()),
..FakeRuntime::default()
});
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fakeskew", Some(fake_data_contract())))
.await
.expect("provision");
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
let ReadinessOutcome::DataSkew { failed } = outcome else {
panic!("a skewCheck failure must be DataSkew, got {outcome:?}");
};
assert_eq!(failed.code, Some(3));
let summaries = progress.summaries();
let last = summaries.last().expect("a skew decision line");
assert_eq!(
*last,
"workspace data: skewCheck `check-skew` → FAILED (exit code 3) — blocking mission (owner: repo-setup)"
);
}
#[tokio::test]
async fn teardown_keep_hibernate_and_destroy_semantics() {
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fake4", Some(fake_contract())))
.await
.expect("provision");
let calls_before = fake.calls().len();
provider
.teardown(handle, TeardownMode::Keep)
.await
.expect("keep is a no-op");
assert_eq!(
fake.calls().len(),
calls_before,
"Keep leaves the project running: no runtime calls"
);
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fake5", Some(fake_contract())))
.await
.expect("provision");
provider
.teardown(handle, TeardownMode::Hibernate)
.await
.expect("hibernate");
assert!(fake.called_with(&["stop"]), "{:?}", fake.calls());
assert!(!fake.called_with(&["down", "-v"]), "{:?}", fake.calls());
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let handle = provider
.provision(&spec(dir.path(), "m-fake6", Some(fake_contract())))
.await
.expect("provision");
provider
.teardown(handle, TeardownMode::Destroy)
.await
.expect("destroy");
assert!(fake.called_with(&["down", "-v"]), "{:?}", fake.calls());
assert!(
dir.path()
.join(".kranz/missions/m-fake6/workspace/prune-marker.txt")
.exists(),
"the disk prune hint ran with the compose dir as cwd"
);
}
#[cfg(unix)]
#[tokio::test]
async fn disk_prune_runs_with_the_same_cleared_gate_env_as_the_other_lanes() {
let _guard = crate::agent_env::EnvTestGuard::engage(&[(
"KRANZ_SECRET_TEST",
"must-not-reach-disk-prune",
)]);
let dir = tempfile::tempdir().expect("tempdir");
let fake = Arc::new(FakeRuntime::default());
let provider = LocalContainerProvider::with_hooks(fake.hooks());
let contract = contract(
br#"{
"schemaVersion": 1,
"readiness": ["true"],
"services": [
{
"name": "api",
"start": "sleep infinity",
"healthCheck": "true",
"port": { "policy": "dynamic" }
}
],
"disk": {
"prune": "printf '%s|%s' \"$KRANZ_SECRET_TEST\" \"$HOME\" > prune-env.txt"
}
}"#,
);
let handle = provider
.provision(&spec(dir.path(), "m-prune", Some(contract)))
.await
.expect("provision");
let gate_home = handle.gate_env.home.clone();
provider
.teardown(handle, TeardownMode::Destroy)
.await
.expect("destroy");
let recorded = std::fs::read_to_string(
dir.path()
.join(".kranz/missions/m-prune/workspace/prune-env.txt"),
)
.expect("the prune command ran");
let (secret, home) = recorded.split_once('|').expect("both values recorded");
assert!(
secret.is_empty(),
"an undeclared ambient credential reached disk.prune: {recorded:?}"
);
assert_eq!(
home,
gate_home.to_str().expect("utf-8 gate home"),
"disk.prune must run with the mission's shared gate HOME, not the operator's"
);
}
#[tokio::test]
async fn provision_without_a_contract_starts_no_containers_and_needs_no_runtime() {
let dir = tempfile::tempdir().expect("tempdir");
let provider = LocalContainerProvider::new();
let handle = provider
.provision(&spec(dir.path(), "m-nocontract", None))
.await
.expect("provision");
assert!(handle.container.is_none());
assert!(handle.detail.is_none());
assert!(handle.previews.is_empty());
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(matches!(outcome, ReadinessOutcome::Ready));
assert!(progress.0.is_empty(), "no contract ⇒ no gate lines");
provider
.teardown(handle, TeardownMode::Destroy)
.await
.expect("teardown with no containers is a no-op");
}
#[test]
fn parse_compose_port_reads_the_assigned_host_port() {
assert_eq!(parse_compose_port("0.0.0.0:32768\n"), Some(32768));
assert_eq!(
parse_compose_port("[::]:49153\n0.0.0.0:49153\n"),
Some(49153)
);
assert_eq!(parse_compose_port(""), None);
assert_eq!(parse_compose_port("Error: no such service\n"), None);
}
#[test]
fn previews_substitute_only_a_single_actually_assigned_dynamic_port() {
let previews = vec![PreviewSpec {
name: "app".to_string(),
url_template: "http://localhost:{port}/".to_string(),
}];
assert_eq!(
fill_previews(&previews, &[])[0].url_template,
"http://localhost:{port}/"
);
assert_eq!(
fill_previews(&previews, &[("api".to_string(), 32768)])[0].url_template,
"http://localhost:32768/"
);
assert_eq!(
fill_previews(
&previews,
&[("api".to_string(), 32768), ("web".to_string(), 32769)]
)[0]
.url_template,
"http://localhost:{port}/"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn container_workspace_smoke_provisions_isolates_and_destroys() {
let _env = crate::agent_env::EnvTestGuard::engage(&[]);
if !sandbox_container::host_supports_container_contract() {
crate::test_capability::skip(
crate::test_capability::capability::CONTAINER,
"live container contract is supported only on Linux",
);
return;
}
let Some(runtime) = sandbox_container::detect() else {
eprintln!(
"no container runtime (docker/podman/nerdctl/container) on PATH; \
skipping container workspace smoke test"
);
return;
};
if runtime == ContainerRuntime::AppleContainer {
eprintln!("the `container` runtime has no compose subcommand; skipping smoke test");
return;
}
let compose_ok = std::process::Command::new(runtime.binary())
.args(["compose", "version"])
.stdin(std::process::Stdio::null())
.output()
.map(|output| output.status.success())
.unwrap_or(false);
if !compose_ok {
eprintln!(
"`{} compose` unavailable; skipping smoke test",
runtime.binary()
);
return;
}
let parent = std::env::current_dir().unwrap();
let dir_a = tempfile::tempdir_in(&parent).expect("tempdir a");
let dir_b = tempfile::tempdir_in(&parent).expect("tempdir b");
let provider = LocalContainerProvider::new();
let handle_a = provider
.provision(&spec(dir_a.path(), "m-smoke-a", Some(fake_contract())))
.await
.expect("provision a");
let handle_b = provider
.provision(&spec(dir_b.path(), "m-smoke-b", Some(fake_contract())))
.await
.expect("provision b (parallel)");
let project_a = handle_a.container.as_ref().unwrap().project.clone();
let project_b = handle_b.container.as_ref().unwrap().project.clone();
assert_ne!(
project_a, project_b,
"parallel missions get distinct projects (distinct networks)"
);
for handle in [&handle_a, &handle_b] {
let mut progress = Progress::default();
let outcome = provider
.readiness(handle, &mut progress.sink())
.await
.expect("readiness");
assert!(
matches!(outcome, ReadinessOutcome::Ready),
"readiness passes inside the container network: {outcome:?}"
);
}
assert!(
dir_a.path().join(".boot-marker").exists(),
"bootstrap wrote through the worktree mount to the host"
);
let port_a = handle_a.container.as_ref().unwrap().assigned_ports[0].1;
let port_b = handle_b.container.as_ref().unwrap().assigned_ports[0].1;
assert!(port_a > 0 && port_b > 0, "OS-assigned ports read back");
assert_ne!(port_a, port_b, "no shared port namespace");
assert!(
handle_a.previews[0]
.url_template
.contains(&port_a.to_string()),
"the preview got the actually-assigned port"
);
provider
.teardown(handle_a, TeardownMode::Destroy)
.await
.expect("destroy a");
provider
.teardown(handle_b, TeardownMode::Destroy)
.await
.expect("destroy b");
for project in [&project_a, &project_b] {
let output = std::process::Command::new(runtime.binary())
.args([
"ps",
"-aq",
"--filter",
&format!("label=com.docker.compose.project={project}"),
])
.stdin(std::process::Stdio::null())
.output()
.expect("spawn compose ps");
assert!(
output.status.success()
&& String::from_utf8_lossy(&output.stdout).trim().is_empty(),
"Destroy removes the project {project}: {}",
String::from_utf8_lossy(&output.stderr)
);
}
}
}