use std::collections::BTreeMap;
use std::future::Future;
use std::pin::Pin;
use std::time::Duration;
use anyhow::{bail, Context, Result};
use async_trait::async_trait;
use local_driver::{canonical_name, ContainerRunSpec, ContainerState, LocalRuntime};
use serde::Deserialize;
use super::{ReconcileCtx, Reconciler, RunningWorkload};
use crate::config::{CloudConfig, Provider};
use crate::local_container_spec_from_provider;
use crate::MirrorShape;
const SLOT: &str = "compute";
#[derive(Debug, Clone, Default)]
pub struct ContainerOptions {
pub adopt_only: bool,
}
#[derive(Debug, Default)]
pub struct ContainerReconciler {
opts: ContainerOptions,
}
impl ContainerReconciler {
pub fn new() -> Self {
Self::default()
}
pub fn with_options(mut self, opts: ContainerOptions) -> Self {
self.opts = opts;
self
}
}
#[derive(Debug, Default, Deserialize)]
struct ContainerComponent {
#[serde(default)]
build: BuildSpec,
#[serde(default)]
run: RunSpec,
}
#[derive(Debug, Deserialize)]
struct BuildSpec {
#[serde(default = "default_dockerfile")]
dockerfile: String,
#[serde(default)]
context: Option<String>,
#[serde(default)]
image: Option<String>,
}
impl Default for BuildSpec {
fn default() -> Self {
Self {
dockerfile: default_dockerfile(),
context: None,
image: None,
}
}
}
fn default_dockerfile() -> String {
"Dockerfile".to_string()
}
#[derive(Debug, Default, Deserialize)]
struct RunSpec {
#[serde(default)]
port: Option<u16>,
#[serde(default)]
host_port: Option<u16>,
#[serde(default)]
env: BTreeMap<String, String>,
#[serde(default)]
mounts: Vec<MountSpec>,
}
#[derive(Debug, Deserialize)]
struct MountSpec {
host: String,
container: String,
#[serde(default = "default_read_only")]
read_only: bool,
}
fn default_read_only() -> bool {
true
}
#[async_trait]
impl Reconciler for ContainerReconciler {
fn kind(&self) -> &'static str {
"container"
}
async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload> {
ctx.materialize().await?;
if !matches!(ctx.mirror.shape, MirrorShape::Local) {
bail!(
"component {}: kind \"container\" has only a local reconciler — mirror \
shape is {:?}, not `local`. Deploy the cloud tier via \
`yah cloud workload deploy` against a yubaba machine.",
ctx.component.id,
ctx.mirror.shape,
);
}
let spec = load_container_component(&ctx)?;
let container_port = spec.run.port.with_context(|| {
format!(
"component {}: workload.toml is missing [run].port — the container reconciler \
needs the port the process listens on",
ctx.component.id,
)
})?;
let host_port = spec.run.host_port.unwrap_or(container_port);
let runtime = detect_local_runtime(&ctx)
.await
.context("detecting local container runtime (orbstack/colima/docker)")?;
let name = canonical_name(&ctx.service.name, ctx.env, &ctx.component.id);
if self.opts.adopt_only {
return match runtime.container_state(&name).await? {
Some(ContainerState::Running) => {
let hp = runtime
.container_host_port(&name, container_port)
.await
.unwrap_or(host_port);
Ok(RunningWorkload::adopted(
"container",
SLOT,
Some(format!("http://127.0.0.1:{hp}")),
)
.with_teardown(teardown_for(&ctx, name.clone())))
}
other => bail!(
"adopt_only: no running container named {name} for component {} \
(state: {other:?}) — nothing to adopt",
ctx.component.id,
),
};
}
let image = spec
.build
.image
.clone()
.unwrap_or_else(|| default_image_tag(&ctx.service.name, &ctx.component.id));
let dockerfile = ctx.workload_dir().join(&spec.build.dockerfile);
let context = match &spec.build.context {
Some(rel) => ctx.workspace_root.join(rel),
None => ctx.workload_dir(),
};
runtime
.build_image(&image, &dockerfile, &context)
.await
.with_context(|| {
format!(
"building image {image} for component {} (dockerfile {}, context {})",
ctx.component.id,
dockerfile.display(),
context.display(),
)
})?;
let mut run_spec =
ContainerRunSpec::new(&ctx.service.name, ctx.env, &ctx.component.id, image);
run_spec.ports = vec![(host_port, container_port)];
run_spec.env = spec.run.env.clone();
run_spec.volumes = resolve_mounts(&spec.run.mounts, ctx.workspace_root)?;
runtime
.run(&run_spec)
.await
.with_context(|| format!("running container for component {}", ctx.component.id))?;
let actual = runtime
.container_host_port(&name, container_port)
.await
.unwrap_or(host_port);
Ok(RunningWorkload::adopted(
"container",
SLOT,
Some(format!("http://127.0.0.1:{actual}")),
)
.with_teardown(teardown_for(&ctx, name.clone())))
}
}
const TEARDOWN_GRACE: Duration = Duration::from_secs(5);
fn teardown_for(
ctx: &ReconcileCtx<'_>,
name: String,
) -> impl FnOnce() -> Pin<Box<dyn Future<Output = Result<()>> + Send>> + Send + Sync + 'static {
let workspace_root = ctx.workspace_root.to_path_buf();
move || {
Box::pin(async move {
let runtime = detect_local_runtime_at(&workspace_root)
.await
.context("detecting local container runtime for teardown")?;
runtime
.stop_and_remove(&name, TEARDOWN_GRACE)
.await
.with_context(|| format!("stopping container {name}"))
})
}
}
fn load_container_component(ctx: &ReconcileCtx<'_>) -> Result<ContainerComponent> {
let path = ctx.workload_dir().join("workload.toml");
let src =
std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))
}
fn default_image_tag(service: &str, component: &str) -> String {
format!("yah-local/{service}-{component}:dev")
}
fn resolve_mounts(
mounts: &[MountSpec],
workspace_root: &std::path::Path,
) -> Result<Vec<(std::path::PathBuf, String)>> {
mounts
.iter()
.map(|m| {
let host = std::path::Path::new(&m.host);
let host = if host.is_absolute() {
host.to_path_buf()
} else {
workspace_root.join(host)
};
if !host.exists() {
bail!(
"mount source {} does not exist (declared as `{}` in workload.toml \
[[run.mounts]]) — the container would silently get an empty directory",
host.display(),
m.host,
);
}
let target = if m.read_only {
format!("{}:ro", m.container)
} else {
m.container.clone()
};
Ok((host, target))
})
.collect()
}
async fn detect_local_runtime(ctx: &ReconcileCtx<'_>) -> Result<LocalRuntime> {
detect_local_runtime_at(ctx.workspace_root).await
}
async fn detect_local_runtime_at(workspace_root: &std::path::Path) -> Result<LocalRuntime> {
let cfg = CloudConfig::load(workspace_root)
.context("loading CloudConfig for local-container provider lookup")?;
let provider = cfg
.providers
.iter()
.find(|p| matches!(p.kind, Provider::LocalContainer))
.with_context(|| {
format!(
"no `kind = \"local-container\"` provider declared in {}/.yah/infra/providers/ — \
the container reconciler needs orbstack.toml or equivalent",
workspace_root.display(),
)
})?;
let local_spec = local_container_spec_from_provider(provider)?;
LocalRuntime::detect(&local_spec).await
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_build_and_run_sections() {
let src = r#"
schema_version = 1
name = "yah-cloud-admin"
kind = "container"
[build]
dockerfile = "Dockerfile"
context = "."
image = "yah-local/yah-cloud-admin:dev"
[run]
port = 4325
host_port = 4325
[run.env]
YAH_CLOUD_ADMIN_ADDR = "0.0.0.0:4325"
YAH_CLOUD_ADMIN_DEV_ANON = "1"
"#;
let c: ContainerComponent = toml::from_str(src).unwrap();
assert_eq!(c.build.dockerfile, "Dockerfile");
assert_eq!(c.build.context.as_deref(), Some("."));
assert_eq!(
c.build.image.as_deref(),
Some("yah-local/yah-cloud-admin:dev")
);
assert_eq!(c.run.port, Some(4325));
assert_eq!(c.run.host_port, Some(4325));
assert_eq!(
c.run.env.get("YAH_CLOUD_ADMIN_ADDR").map(String::as_str),
Some("0.0.0.0:4325")
);
assert_eq!(
c.run
.env
.get("YAH_CLOUD_ADMIN_DEV_ANON")
.map(String::as_str),
Some("1")
);
}
#[test]
fn mounts_resolve_against_the_workspace_root_and_default_read_only() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(tmp.path().join(".yah/infra")).unwrap();
let src = r#"
[run]
port = 4325
[[run.mounts]]
host = ".yah/infra"
container = "/workspace/.yah/infra"
"#;
let c: ContainerComponent = toml::from_str(src).unwrap();
let out = resolve_mounts(&c.run.mounts, tmp.path()).unwrap();
assert_eq!(out.len(), 1);
assert_eq!(out[0].0, tmp.path().join(".yah/infra"));
assert_eq!(out[0].1, "/workspace/.yah/infra:ro");
}
#[test]
fn a_writable_mount_must_be_asked_for() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(tmp.path().join("state")).unwrap();
let src = r#"
[run]
port = 1
[[run.mounts]]
host = "state"
container = "/var/lib/state"
read_only = false
"#;
let c: ContainerComponent = toml::from_str(src).unwrap();
let out = resolve_mounts(&c.run.mounts, tmp.path()).unwrap();
assert_eq!(out[0].1, "/var/lib/state");
}
#[test]
fn a_missing_mount_source_fails_the_reconcile() {
let tmp = tempfile::tempdir().unwrap();
let src = r#"
[run]
port = 1
[[run.mounts]]
host = "nope"
container = "/nope"
"#;
let c: ContainerComponent = toml::from_str(src).unwrap();
let err = resolve_mounts(&c.run.mounts, tmp.path()).unwrap_err();
assert!(err.to_string().contains("does not exist"), "{err}");
}
#[test]
fn no_mounts_declared_is_no_volumes() {
let c: ContainerComponent = toml::from_str("[run]\nport = 1\n").unwrap();
assert!(c.run.mounts.is_empty());
assert!(resolve_mounts(&c.run.mounts, std::path::Path::new("/"))
.unwrap()
.is_empty());
}
#[test]
fn build_defaults_when_section_absent() {
let src = r#"
kind = "container"
[run]
port = 8080
"#;
let c: ContainerComponent = toml::from_str(src).unwrap();
assert_eq!(c.build.dockerfile, "Dockerfile");
assert!(c.build.context.is_none());
assert!(c.build.image.is_none());
assert_eq!(c.run.port, Some(8080));
assert!(c.run.host_port.is_none());
assert!(c.run.env.is_empty());
}
#[test]
fn the_teardown_name_matches_the_name_run_created() {
let (service, env, component) = ("yah-cloud-admin", "pond", "cloud-admin");
let run_spec = ContainerRunSpec::new(service, env, component, "img:dev");
let teardown_target = canonical_name(service, env, component);
assert_eq!(
run_spec.name, teardown_target,
"teardown would docker-stop a name that was never created"
);
}
#[test]
fn default_image_tag_derives_from_service_and_component() {
assert_eq!(
default_image_tag("yah-cloud-admin", "cloud-admin"),
"yah-local/yah-cloud-admin-cloud-admin:dev"
);
}
#[test]
fn run_spec_publishes_declared_ports_and_env() {
let mut run_spec =
ContainerRunSpec::new("yah-cloud-admin", "dev", "cloud-admin", "img:dev");
run_spec.ports = vec![(4325, 4325)];
run_spec.env.insert("K".into(), "V".into());
let args = run_spec.docker_run_args();
let joined = args.join(" ");
assert!(joined.contains("-p 4325:4325"), "args: {joined}");
assert!(joined.contains("-e K=V"), "args: {joined}");
assert!(
joined.ends_with("img:dev"),
"image is the final arg: {joined}"
);
}
}