pub mod compose;
pub mod secrets;
pub use compose::{detect, merge, ComposeError, ComposeFile, ProjectOverrides};
pub use secrets::{check, SecretsBundle, SecretsError};
use std::path::{Path, PathBuf};
use std::time::Duration;
use nyx_agent_core::project::Project;
use serde::Deserialize;
use thiserror::Error;
use tokio::io::AsyncReadExt;
use tokio::process::Command;
use tokio::time::timeout;
const DEFAULT_DOCKER_TIMEOUT: Duration = Duration::from_secs(180);
const DOWN_TIMEOUT: Duration = Duration::from_secs(60);
pub const MIN_COMPOSE_VERSION: &str = "2.20.0";
#[derive(Debug, Clone)]
pub struct ComposeVersion {
pub raw: String,
pub parsed: semver::Version,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum PullPolicy {
#[default]
Missing,
Always,
Never,
}
impl PullPolicy {
pub fn as_flag(&self) -> &'static str {
match self {
PullPolicy::Missing => "missing",
PullPolicy::Always => "always",
PullPolicy::Never => "never",
}
}
}
#[derive(Debug, Error)]
pub enum EnvError {
#[error(transparent)]
Compose(#[from] ComposeError),
#[error(transparent)]
Secrets(#[from] SecretsError),
#[error("`docker` binary not found on PATH; install Docker before running env-builder")]
DockerMissing,
#[error("`docker compose` is unsupported: {reason}")]
ComposeUnsupported { reason: String },
#[error("`docker compose up` failed (exit {code:?}); stderr: {stderr}")]
UpFailed { code: Option<i32>, stderr: String },
#[error("docker compose project `{project}` is already in use; pick a different project name")]
ProjectInUse { project: String },
#[error("`docker compose ls` failed: {reason}")]
LsFailed { reason: String },
#[error("`docker compose down` failed (exit {code:?}); stderr: {stderr}")]
DownFailed { code: Option<i32>, stderr: String },
#[error("`docker compose ps` failed (exit {code:?}); stderr: {stderr}")]
PsFailed { code: Option<i32>, stderr: String },
#[error("`docker compose logs` failed (exit {code:?}); stderr: {stderr}")]
LogsFailed { code: Option<i32>, stderr: String },
#[error("`docker compose ps` returned malformed JSON: {0}")]
MalformedPs(#[source] serde_json::Error),
#[error("docker subcommand timed out after {0:?}")]
Timeout(Duration),
#[error("io error invoking docker: {0}")]
Io(#[from] std::io::Error),
}
#[derive(Debug, Clone)]
pub struct RepoInput {
pub name: String,
pub root: PathBuf,
}
#[derive(Debug, Clone, Deserialize)]
pub struct ServiceHealth {
#[serde(rename = "Name", default)]
pub container_name: String,
#[serde(rename = "Service", default)]
pub service: String,
#[serde(rename = "State", default)]
pub state: String,
#[serde(rename = "Health", default)]
pub health: String,
#[serde(rename = "Status", default)]
pub status: String,
}
#[derive(Debug, Clone)]
pub struct EnvBuilder {
pub docker_binary: PathBuf,
pub workspace: PathBuf,
pub state_root: PathBuf,
pub project_name: String,
pub target_base_url: Option<String>,
pub env_config: Option<serde_json::Value>,
pub repos: Vec<RepoInput>,
pub command_timeout: Duration,
pub pull_policy: PullPolicy,
}
impl EnvBuilder {
pub fn discover(
workspace: PathBuf,
state_root: PathBuf,
project: &Project,
repos: Vec<RepoInput>,
) -> Result<Self, EnvError> {
let docker = which_on_path("docker").ok_or(EnvError::DockerMissing)?;
probe_compose_version(&docker)?;
Ok(Self {
docker_binary: docker,
workspace,
state_root,
project_name: project.name.clone(),
target_base_url: project.target_base_url.clone(),
env_config: project.env_config.clone(),
repos,
command_timeout: DEFAULT_DOCKER_TIMEOUT,
pull_policy: PullPolicy::default(),
})
}
pub fn super_compose_filename(&self) -> String {
format!("nyx-super-compose-{}.yml", sanitise_filename(&self.project_name))
}
pub async fn up(&self) -> Result<RunningEnv, EnvError> {
let secrets_bundle = check(&self.state_root)?;
let compose_files = self.detect_compose_files();
let super_compose = self.workspace.join(self.super_compose_filename());
let overrides = ProjectOverrides {
target_base_url: self.target_base_url.as_deref(),
env_config: self.env_config.as_ref(),
};
let services = merge(&compose_files, &super_compose, &overrides)?;
self.refuse_if_project_in_use().await?;
let mut cmd = self.compose_command(&super_compose, &secrets_bundle.path);
cmd.arg("up").arg("-d").arg("--build").arg("--pull").arg(self.pull_policy.as_flag());
let outcome = run_command(cmd, self.command_timeout).await?;
if !outcome.status_ok() {
return Err(EnvError::UpFailed {
code: outcome.exit_code,
stderr: outcome.stderr_string(),
});
}
let env = RunningEnv {
docker_binary: self.docker_binary.clone(),
super_compose,
secrets_path: secrets_bundle.path.clone(),
project_name: self.project_name.clone(),
services,
command_timeout: self.command_timeout,
running: true,
};
Ok(env)
}
fn detect_compose_files(&self) -> Vec<ComposeFile> {
let mut out = Vec::new();
for repo in &self.repos {
if let Some(f) = detect(&repo.root, &repo.name) {
out.push(f);
}
}
out
}
pub async fn refuse_if_project_in_use(&self) -> Result<(), EnvError> {
let mut cmd = Command::new(&self.docker_binary);
cmd.arg("compose").arg("ls").arg("--all").arg("--format").arg("json");
let outcome = run_command(cmd, self.command_timeout).await?;
if !outcome.status_ok() {
return Err(EnvError::LsFailed {
reason: format!(
"`docker compose ls --all --format json` exited {code:?}: {stderr}",
code = outcome.exit_code,
stderr = outcome.stderr_string(),
),
});
}
let names = parse_compose_ls_names(&outcome.stdout)?;
if names.iter().any(|n| n == &self.project_name) {
return Err(EnvError::ProjectInUse { project: self.project_name.clone() });
}
Ok(())
}
fn compose_command(&self, super_compose: &Path, env_file: &Path) -> Command {
let mut cmd = Command::new(&self.docker_binary);
cmd.arg("compose")
.arg("--project-name")
.arg(&self.project_name)
.arg("-f")
.arg(super_compose)
.arg("--env-file")
.arg(env_file);
cmd
}
}
#[derive(Debug)]
pub struct RunningEnv {
docker_binary: PathBuf,
super_compose: PathBuf,
secrets_path: PathBuf,
project_name: String,
services: Vec<String>,
command_timeout: Duration,
running: bool,
}
impl RunningEnv {
pub fn project_name(&self) -> &str {
&self.project_name
}
pub fn services(&self) -> &[String] {
&self.services
}
pub fn super_compose_path(&self) -> &Path {
&self.super_compose
}
pub async fn services_health(&self) -> Result<Vec<ServiceHealth>, EnvError> {
let mut cmd = self.compose_command();
cmd.arg("ps").arg("--format").arg("json");
let outcome = run_command(cmd, self.command_timeout).await?;
if !outcome.status_ok() {
return Err(EnvError::PsFailed {
code: outcome.exit_code,
stderr: outcome.stderr_string(),
});
}
parse_ps_json(&outcome.stdout)
}
pub async fn service_logs(&self, service: &str) -> Result<Vec<u8>, EnvError> {
let mut cmd = self.compose_command();
cmd.arg("logs").arg("--no-color").arg("--timestamps").arg("--no-log-prefix").arg(service);
let outcome = run_command(cmd, self.command_timeout).await?;
if !outcome.status_ok() {
return Err(EnvError::LogsFailed {
code: outcome.exit_code,
stderr: outcome.stderr_string(),
});
}
Ok(outcome.stdout)
}
pub async fn down(mut self) -> Result<(), EnvError> {
let result = self.down_inner().await;
self.running = false;
result
}
pub async fn reset(&mut self) -> Result<(), EnvError> {
self.down_inner().await?;
let mut cmd = self.compose_command();
cmd.arg("up").arg("-d").arg("--build");
let outcome = run_command(cmd, self.command_timeout).await?;
if !outcome.status_ok() {
self.running = false;
return Err(EnvError::UpFailed {
code: outcome.exit_code,
stderr: outcome.stderr_string(),
});
}
self.running = true;
Ok(())
}
async fn down_inner(&self) -> Result<(), EnvError> {
let mut cmd = self.compose_command();
cmd.arg("down").arg("--volumes").arg("--remove-orphans");
let outcome = run_command(cmd, DOWN_TIMEOUT).await?;
if !outcome.status_ok() {
return Err(EnvError::DownFailed {
code: outcome.exit_code,
stderr: outcome.stderr_string(),
});
}
Ok(())
}
fn compose_command(&self) -> Command {
let mut cmd = Command::new(&self.docker_binary);
cmd.arg("compose")
.arg("--project-name")
.arg(&self.project_name)
.arg("-f")
.arg(&self.super_compose)
.arg("--env-file")
.arg(&self.secrets_path);
cmd
}
}
impl Drop for RunningEnv {
fn drop(&mut self) {
if !self.running {
return;
}
let docker_binary = self.docker_binary.clone();
let super_compose = self.super_compose.clone();
let secrets_path = self.secrets_path.clone();
let project_name = self.project_name.clone();
let _ = std::thread::Builder::new()
.name(format!("nyx-env-down-{project_name}"))
.spawn(move || {
let mut cmd = std::process::Command::new(&docker_binary);
cmd.arg("compose")
.arg("--project-name")
.arg(&project_name)
.arg("-f")
.arg(&super_compose)
.arg("--env-file")
.arg(&secrets_path)
.arg("down")
.arg("--volumes")
.arg("--remove-orphans")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) => {
tracing::warn!(
error = %e,
project = %project_name,
"docker compose down failed to spawn on RunningEnv drop; containers may leak"
);
return;
}
};
let deadline = std::time::Instant::now() + DOWN_TIMEOUT;
loop {
match child.try_wait() {
Ok(Some(_)) => return,
Ok(None) => {
if std::time::Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
tracing::warn!(
project = %project_name,
timeout_secs = DOWN_TIMEOUT.as_secs(),
"docker compose down timed out on RunningEnv drop; killed subprocess"
);
return;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
Err(e) => {
tracing::warn!(
error = %e,
project = %project_name,
"docker compose down wait errored on RunningEnv drop"
);
return;
}
}
}
});
}
}
fn parse_ps_json(raw: &[u8]) -> Result<Vec<ServiceHealth>, EnvError> {
let text = std::str::from_utf8(raw).unwrap_or("").trim();
if text.is_empty() {
return Ok(Vec::new());
}
if text.starts_with('[') {
return serde_json::from_str::<Vec<ServiceHealth>>(text).map_err(EnvError::MalformedPs);
}
let mut out = Vec::new();
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let row: ServiceHealth = serde_json::from_str(line).map_err(EnvError::MalformedPs)?;
out.push(row);
}
Ok(out)
}
fn parse_compose_ls_names(raw: &[u8]) -> Result<Vec<String>, EnvError> {
#[derive(Deserialize)]
struct ComposeLsRow {
#[serde(rename = "Name", default)]
name: String,
}
let text = std::str::from_utf8(raw).unwrap_or("").trim();
if text.is_empty() {
return Ok(Vec::new());
}
let rows: Vec<ComposeLsRow> = if text.starts_with('[') {
serde_json::from_str(text).map_err(|e| EnvError::LsFailed {
reason: format!("malformed `docker compose ls` array: {e}"),
})?
} else {
let mut out = Vec::new();
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let row: ComposeLsRow = serde_json::from_str(line).map_err(|e| EnvError::LsFailed {
reason: format!("malformed `docker compose ls` ndjson row: {e}"),
})?;
out.push(row);
}
out
};
Ok(rows.into_iter().map(|r| r.name).filter(|n| !n.is_empty()).collect())
}
#[derive(Debug)]
struct CommandOutcome {
exit_code: Option<i32>,
stdout: Vec<u8>,
stderr: Vec<u8>,
}
impl CommandOutcome {
fn status_ok(&self) -> bool {
matches!(self.exit_code, Some(0))
}
fn stderr_string(&self) -> String {
String::from_utf8_lossy(&self.stderr).trim().to_string()
}
}
async fn run_command(mut cmd: Command, cap: Duration) -> Result<CommandOutcome, EnvError> {
cmd.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
let mut child = cmd.spawn()?;
let mut stdout = child.stdout.take().expect("stdout piped");
let mut stderr = child.stderr.take().expect("stderr piped");
let drain = async {
let mut out = Vec::new();
let mut err = Vec::new();
let (a, b) = tokio::join!(stdout.read_to_end(&mut out), stderr.read_to_end(&mut err));
a?;
b?;
Ok::<_, std::io::Error>((out, err))
};
let waited = timeout(cap, async {
let (out, err) = drain.await?;
let status = child.wait().await?;
Ok::<_, std::io::Error>((status, out, err))
})
.await;
match waited {
Err(_) => {
let _ = child.start_kill();
let _ = child.wait().await;
Err(EnvError::Timeout(cap))
}
Ok(Err(io)) => Err(EnvError::Io(io)),
Ok(Ok((status, stdout, stderr))) => {
Ok(CommandOutcome { exit_code: status.code(), stdout, stderr })
}
}
}
fn sanitise_filename(name: &str) -> String {
let mut out = String::with_capacity(name.len());
for c in name.chars() {
if c.is_ascii_alphanumeric() {
out.push(c.to_ascii_lowercase());
} else {
out.push('_');
}
}
if out.is_empty() || out.chars().all(|c| c == '_') {
"project".to_string()
} else {
out
}
}
fn which_on_path(bin: &str) -> Option<PathBuf> {
let path = std::env::var_os("PATH")?;
for dir in std::env::split_paths(&path) {
let candidate = dir.join(bin);
if candidate.is_file() {
return Some(candidate);
}
}
None
}
pub fn probe_compose_version(docker: &Path) -> Result<ComposeVersion, EnvError> {
let output = std::process::Command::new(docker)
.arg("compose")
.arg("version")
.arg("--short")
.stdin(std::process::Stdio::null())
.output()
.map_err(EnvError::Io)?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(EnvError::ComposeUnsupported {
reason: format!(
"`docker compose version --short` exited non-zero (status {status:?}); stderr: {stderr}",
status = output.status.code(),
),
});
}
let raw = String::from_utf8_lossy(&output.stdout).trim().to_string();
parse_and_gate_compose_version(&raw)
}
fn parse_and_gate_compose_version(raw: &str) -> Result<ComposeVersion, EnvError> {
let trimmed = raw.trim();
let candidate = trimmed.strip_prefix('v').unwrap_or(trimmed);
let parsed = match semver::Version::parse(candidate) {
Ok(v) => v,
Err(e) => {
return Err(EnvError::ComposeUnsupported {
reason: format!("could not parse compose version `{trimmed}`: {e}"),
});
}
};
let min = semver::Version::parse(MIN_COMPOSE_VERSION)
.expect("MIN_COMPOSE_VERSION is a valid semver triple");
if parsed < min {
return Err(EnvError::ComposeUnsupported {
reason: format!(
"docker compose {parsed} is below the {min} floor required by the env-builder"
),
});
}
Ok(ComposeVersion { raw: trimmed.to_string(), parsed })
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_ps_json_ndjson() {
let raw = b"{\"Name\":\"a_1\",\"Service\":\"a\",\"State\":\"running\",\"Health\":\"\",\"Status\":\"Up\"}\n{\"Name\":\"b_1\",\"Service\":\"b\",\"State\":\"running\",\"Health\":\"healthy\",\"Status\":\"Up\"}\n";
let rows = parse_ps_json(raw).expect("parse");
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].service, "a");
assert_eq!(rows[1].health, "healthy");
}
#[test]
fn parse_ps_json_array() {
let raw = b"[{\"Name\":\"a_1\",\"Service\":\"a\",\"State\":\"running\",\"Health\":\"\",\"Status\":\"Up\"}]";
let rows = parse_ps_json(raw).expect("parse");
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].service, "a");
}
#[test]
fn parse_ps_json_empty() {
let rows = parse_ps_json(b"").expect("parse empty");
assert!(rows.is_empty());
}
#[test]
fn parse_and_gate_compose_version_accepts_floor() {
let v = parse_and_gate_compose_version("2.20.0").expect("floor must be accepted");
assert_eq!(v.raw, "2.20.0");
assert_eq!(v.parsed.major, 2);
assert_eq!(v.parsed.minor, 20);
}
#[test]
fn parse_and_gate_compose_version_accepts_newer() {
let v = parse_and_gate_compose_version("2.27.1").expect("newer must be accepted");
assert_eq!(v.parsed.minor, 27);
}
#[test]
fn parse_and_gate_compose_version_strips_v_prefix() {
let v = parse_and_gate_compose_version("v2.27.1").expect("`v` prefix must be tolerated");
assert_eq!(v.parsed.minor, 27);
}
#[test]
fn parse_and_gate_compose_version_refuses_old_v2() {
let err = parse_and_gate_compose_version("2.19.0").expect_err("below-floor must refuse");
match err {
EnvError::ComposeUnsupported { reason } => {
assert!(reason.contains("2.19"), "reason must name the seen version: {reason}");
assert!(reason.contains("2.20"), "reason must name the floor: {reason}");
}
other => panic!("expected ComposeUnsupported, got {other:?}"),
}
}
#[test]
fn parse_and_gate_compose_version_refuses_v1_legacy_shape() {
let err = parse_and_gate_compose_version("1.29.2").expect_err("v1 must refuse");
assert!(matches!(err, EnvError::ComposeUnsupported { .. }));
}
#[test]
fn parse_compose_ls_names_array() {
let raw = br#"[{"Name":"alpha","Status":"running(1)","ConfigFiles":"/a.yml"},{"Name":"beta","Status":"exited","ConfigFiles":"/b.yml"}]"#;
let names = parse_compose_ls_names(raw).expect("parse");
assert_eq!(names, vec!["alpha".to_string(), "beta".to_string()]);
}
#[test]
fn parse_compose_ls_names_ndjson() {
let raw = b"{\"Name\":\"alpha\",\"Status\":\"running(1)\"}\n{\"Name\":\"beta\",\"Status\":\"exited\"}\n";
let names = parse_compose_ls_names(raw).expect("parse");
assert_eq!(names, vec!["alpha".to_string(), "beta".to_string()]);
}
#[test]
fn parse_compose_ls_names_empty() {
let names = parse_compose_ls_names(b"").expect("parse empty");
assert!(names.is_empty());
let names = parse_compose_ls_names(b"[]").expect("parse empty array");
assert!(names.is_empty());
}
#[test]
fn parse_compose_ls_names_drops_blank_rows() {
let raw = br#"[{"Name":""},{"Name":"real"}]"#;
let names = parse_compose_ls_names(raw).expect("parse");
assert_eq!(names, vec!["real".to_string()]);
}
#[test]
fn pull_policy_default_matches_docker_default() {
assert_eq!(PullPolicy::default(), PullPolicy::Missing);
}
#[test]
fn pull_policy_flag_strings_match_docker_compose_grammar() {
assert_eq!(PullPolicy::Missing.as_flag(), "missing");
assert_eq!(PullPolicy::Always.as_flag(), "always");
assert_eq!(PullPolicy::Never.as_flag(), "never");
}
#[test]
fn parse_and_gate_compose_version_refuses_garbage() {
let err = parse_and_gate_compose_version("not-a-version")
.expect_err("non-semver string must refuse");
match err {
EnvError::ComposeUnsupported { reason } => {
assert!(reason.contains("not-a-version"), "reason must echo the input: {reason}");
}
other => panic!("expected ComposeUnsupported, got {other:?}"),
}
}
}