use std::collections::BTreeMap;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{bail, Context, Result};
use async_trait::async_trait;
use kamaji::native::NativeRuntime;
use kamaji::{Kamaji, MeshAssignment, MeshIdent};
use serde::Deserialize;
use tokio::sync::oneshot;
use tracing::{info, warn};
use workload_spec::EnvVar;
use super::native_support::{
capture_paths, native_spec, sanitize_ident, spawn_native_log_supervisor,
};
use super::{into_running, wait_for_port, LogBuffer, ReconcileCtx, Reconciler, RunningWorkload};
use crate::{MirrorProviderSlot, MirrorShape, Provider};
const SLOT: &str = "compute";
const READY_TIMEOUT: Duration = Duration::from_secs(20);
pub fn slot_declared(mirror: &crate::MirrorConfig) -> bool {
matches!(
mirror.providers.get(SLOT),
Some(MirrorProviderSlot::Inline {
kind: Provider::LocalProcess,
..
})
)
}
#[derive(Debug, Default)]
pub struct LocalProcessReconciler;
impl LocalProcessReconciler {
pub fn new() -> Self {
Self::default()
}
}
#[derive(Debug, Default, Deserialize)]
struct ProcessComponent {
#[serde(default)]
process: Option<ProcessSpec>,
}
#[derive(Debug, Deserialize)]
struct ProcessSpec {
#[serde(default)]
cargo_package: Option<String>,
#[serde(default)]
bin: Option<String>,
#[serde(default)]
args: Vec<String>,
port: u16,
#[serde(default)]
env: BTreeMap<String, String>,
#[serde(default = "default_profile")]
profile: String,
}
fn default_profile() -> String {
"debug".to_string()
}
#[async_trait]
impl Reconciler for LocalProcessReconciler {
fn kind(&self) -> &'static str {
"local-process"
}
async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload> {
ctx.materialize().await?;
if !matches!(ctx.mirror.shape, MirrorShape::Local) {
bail!(
"component {}: `local-process` is a dev-tier compute slot — mirror shape is \
{:?}, not `local`. Deploy the cloud tier via `yah cloud workload deploy`.",
ctx.component.id,
ctx.mirror.shape,
);
}
let spec = load_process_spec(&ctx)?;
if let Some(pkg) = &spec.cargo_package {
cargo_build(ctx.workspace_root, pkg, &spec.profile).await?;
}
let bin = resolve_binary(&spec, ctx.workspace_root)?;
let state_dir = ctx.workspace_root.join(".yah/jit/native");
let ident_str = sanitize_ident(&format!(
"local-process-{}-{}-{}",
ctx.service.name, ctx.env, ctx.component.id
));
let ident = MeshIdent(ident_str.clone());
let owner_path = state_dir.join(&ident_str).join("owner.json");
reap_predecessor(&owner_path).await;
let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), spec.port);
if tokio::net::TcpStream::connect(addr).await.is_ok() {
bail!(
"component {}: {addr} is already held by a process this reconciler did not \
start. Stop it, or give [process] a different `port`.",
ctx.component.id,
);
}
let mut argv = vec![bin.display().to_string()];
argv.extend(spec.args.iter().cloned());
let env: Vec<EnvVar> = spec
.env
.iter()
.map(|(name, value)| EnvVar {
name: name.clone(),
value: workload_spec::EnvValue::Literal {
value: value.clone(),
},
})
.collect();
let workload = native_spec(&ident_str, argv, env);
let runtime = Arc::new(NativeRuntime::new(&state_dir));
let mesh = MeshAssignment::inlined(Ipv4Addr::LOCALHOST);
info!(
binary = %bin.display(),
port = spec.port,
ident = %ident_str,
"spawning local-process component (kamaji native backend)",
);
let deployed = runtime
.deploy_workload(&workload, &mesh)
.await
.with_context(|| {
format!(
"deploying component {} via kamaji native backend (binary {})",
ctx.component.id,
bin.display(),
)
})?;
let (stdout_path, stderr_path) = capture_paths(&state_dir, &ident_str);
write_owner(&owner_path, deployed.task_pid as i32, spec.port).await;
if !wait_for_port(addr, READY_TIMEOUT).await {
warn!(addr = %addr, "local-process did not bind within timeout; tearing down");
let tail = read_capture_tail(&stdout_path, &stderr_path).await;
runtime.teardown_workload(&ident).await.ok();
bail!(
"component {} did not bind {addr} within {READY_TIMEOUT:?}{tail}",
ctx.component.id,
);
}
let dev_url = format!("http://{addr}");
info!(dev_url = %dev_url, pid = deployed.task_pid, "local-process ready");
let log_buf = LogBuffer::new();
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
let supervisor = spawn_native_log_supervisor(
runtime,
ident,
log_buf.clone(),
stdout_path,
stderr_path,
shutdown_rx,
);
Ok(into_running(
"local-process",
SLOT,
Some(dev_url),
None,
Some(log_buf),
shutdown_tx,
supervisor,
))
}
}
fn load_process_spec(ctx: &ReconcileCtx<'_>) -> Result<ProcessSpec> {
let path = ctx.workload_dir().join("workload.toml");
let src =
std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
let parsed: ProcessComponent =
toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
parsed.process.with_context(|| {
format!(
"component {} is bound to a `local-process` compute slot but {} declares no \
[process] section — add one with at least `port`",
ctx.component.id,
path.display(),
)
})
}
fn resolve_binary(spec: &ProcessSpec, workspace_root: &std::path::Path) -> Result<PathBuf> {
let rel = match (&spec.bin, &spec.cargo_package) {
(Some(bin), _) => PathBuf::from(bin),
(None, Some(pkg)) => PathBuf::from(format!("target/{}/{pkg}", spec.profile)),
(None, None) => bail!(
"[process] declares neither `bin` nor `cargo_package` — one is needed to know \
what to run"
),
};
let path = if rel.is_absolute() {
rel
} else {
workspace_root.join(rel)
};
if !path.exists() {
bail!(
"[process] binary {} does not exist — set `cargo_package` to have it built, or \
point `bin` at an existing file",
path.display(),
);
}
Ok(path)
}
async fn cargo_build(workspace_root: &std::path::Path, pkg: &str, profile: &str) -> Result<()> {
let mut cmd = tokio::process::Command::new("cargo");
cmd.arg("build").arg("-p").arg(pkg);
if profile == "release" {
cmd.arg("--release");
} else if profile != "debug" {
cmd.arg("--profile").arg(profile);
}
cmd.current_dir(workspace_root);
info!(
package = pkg,
profile, "cargo build for local-process component"
);
let out = cmd
.output()
.await
.with_context(|| format!("spawning `cargo build -p {pkg}`"))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
let tail: Vec<&str> = stderr.lines().rev().take(30).collect();
let tail: Vec<&str> = tail.into_iter().rev().collect();
bail!("cargo build -p {pkg} failed:\n{}", tail.join("\n"));
}
Ok(())
}
#[derive(Debug, serde::Serialize, Deserialize)]
struct OwnerRecord {
pid: i32,
port: u16,
}
async fn write_owner(path: &std::path::Path, pid: i32, port: u16) {
if let Some(dir) = path.parent() {
if tokio::fs::create_dir_all(dir).await.is_err() {
return;
}
}
if let Ok(json) = serde_json::to_vec(&OwnerRecord { pid, port }) {
if let Err(e) = tokio::fs::write(path, json).await {
warn!(path = %path.display(), error = %e, "could not record local-process owner");
}
}
}
async fn reap_predecessor(owner_path: &std::path::Path) {
let Ok(bytes) = tokio::fs::read(owner_path).await else {
return;
};
let _ = tokio::fs::remove_file(owner_path).await;
let Ok(owner) = serde_json::from_slice::<OwnerRecord>(&bytes) else {
return;
};
if !pid_alive(owner.pid) {
return;
}
info!(
pid = owner.pid,
port = owner.port,
"replacing predecessor local-process"
);
signal_pid(owner.pid, libc::SIGTERM);
for _ in 0..50 {
if !pid_alive(owner.pid) {
return;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
warn!(
pid = owner.pid,
"predecessor ignored SIGTERM; sending SIGKILL"
);
signal_pid(owner.pid, libc::SIGKILL);
tokio::time::sleep(Duration::from_millis(200)).await;
}
fn pid_alive(pid: i32) -> bool {
unsafe { libc::kill(pid, 0) == 0 }
}
fn signal_pid(pid: i32, sig: i32) {
unsafe {
libc::kill(pid, sig);
}
}
async fn read_capture_tail(stdout_path: &std::path::Path, stderr_path: &std::path::Path) -> String {
let mut lines: Vec<String> = Vec::new();
for path in [stderr_path, stdout_path] {
if let Ok(s) = tokio::fs::read_to_string(path).await {
lines.extend(s.lines().rev().take(10).map(str::to_string));
}
}
if lines.is_empty() {
return String::new();
}
lines.reverse();
format!("\nlast output:\n{}", lines.join("\n"))
}
#[cfg(test)]
mod tests {
use super::*;
fn spec(bin: Option<&str>, pkg: Option<&str>) -> ProcessSpec {
ProcessSpec {
cargo_package: pkg.map(str::to_string),
bin: bin.map(str::to_string),
args: vec![],
port: 4325,
env: BTreeMap::new(),
profile: "debug".to_string(),
}
}
#[test]
fn parses_a_process_section_alongside_container_sections() {
let src = r#"
schema_version = 1
kind = "container"
[build]
image = "yah-local/x:dev"
[run]
port = 4325
[process]
cargo_package = "yah-cloud-admin"
port = 4325
args = ["--verbose"]
[process.env]
YAH_CLOUD_ADMIN_ADDR = "127.0.0.1:4325"
"#;
let c: ProcessComponent = toml::from_str(src).unwrap();
let p = c.process.unwrap();
assert_eq!(p.cargo_package.as_deref(), Some("yah-cloud-admin"));
assert_eq!(p.port, 4325);
assert_eq!(p.args, vec!["--verbose".to_string()]);
assert_eq!(p.profile, "debug");
assert_eq!(
p.env.get("YAH_CLOUD_ADMIN_ADDR").map(String::as_str),
Some("127.0.0.1:4325")
);
}
#[test]
fn a_workload_without_a_process_section_parses_as_none() {
let c: ProcessComponent = toml::from_str("[run]\nport = 1\n").unwrap();
assert!(c.process.is_none());
}
#[test]
fn binary_defaults_to_the_cargo_convention() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(tmp.path().join("target/debug")).unwrap();
std::fs::write(tmp.path().join("target/debug/yah-cloud-admin"), b"").unwrap();
let got = resolve_binary(&spec(None, Some("yah-cloud-admin")), tmp.path()).unwrap();
assert_eq!(got, tmp.path().join("target/debug/yah-cloud-admin"));
}
#[test]
fn explicit_bin_wins_over_the_cargo_convention() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(tmp.path().join("bin")).unwrap();
std::fs::write(tmp.path().join("bin/custom"), b"").unwrap();
let got = resolve_binary(&spec(Some("bin/custom"), Some("pkg")), tmp.path()).unwrap();
assert_eq!(got, tmp.path().join("bin/custom"));
}
#[test]
fn a_missing_binary_is_named_before_we_try_to_exec_it() {
let tmp = tempfile::tempdir().unwrap();
let err = resolve_binary(&spec(None, Some("nope")), tmp.path()).unwrap_err();
assert!(err.to_string().contains("does not exist"), "{err}");
}
#[test]
fn neither_bin_nor_package_is_an_error() {
let tmp = tempfile::tempdir().unwrap();
let err = resolve_binary(&spec(None, None), tmp.path()).unwrap_err();
assert!(err.to_string().contains("neither"), "{err}");
}
#[tokio::test]
async fn reap_stops_a_live_predecessor_and_clears_the_record() {
let tmp = tempfile::tempdir().unwrap();
let owner_path = tmp.path().join("owner.json");
let mut child = tokio::process::Command::new("sleep")
.arg("120")
.spawn()
.unwrap();
let pid = child.id().unwrap() as i32;
assert!(pid_alive(pid));
write_owner(&owner_path, pid, 4325).await;
reap_predecessor(&owner_path).await;
let status = tokio::time::timeout(Duration::from_secs(5), child.wait())
.await
.expect("a signalled `sleep 120` must have exited well inside 5s")
.unwrap();
assert!(!status.success(), "predecessor should have been signalled");
assert!(!owner_path.exists(), "sidecar should be cleared");
}
#[tokio::test]
async fn reap_clears_a_stale_record_without_signalling_anything() {
let tmp = tempfile::tempdir().unwrap();
let owner_path = tmp.path().join("owner.json");
let mut child = tokio::process::Command::new("true").spawn().unwrap();
let pid = child.id().unwrap() as i32;
child.wait().await.unwrap();
write_owner(&owner_path, pid, 4325).await;
reap_predecessor(&owner_path).await;
assert!(!owner_path.exists());
}
#[tokio::test]
async fn reap_with_no_record_is_a_no_op() {
let tmp = tempfile::tempdir().unwrap();
reap_predecessor(&tmp.path().join("owner.json")).await;
}
#[test]
fn slot_declared_only_matches_the_local_process_compute_slot() {
let mut m = crate::MirrorConfig {
schema_version: 1,
shape: MirrorShape::Local,
ingress: Default::default(),
providers: Default::default(),
drivers: Default::default(),
asset_aliases: Default::default(),
};
assert!(!slot_declared(&m));
m.providers.insert(
SLOT.to_string(),
MirrorProviderSlot::Inline {
kind: Provider::LocalContainer,
fields: Default::default(),
},
);
assert!(!slot_declared(&m), "a container slot is not a process slot");
m.providers.insert(
SLOT.to_string(),
MirrorProviderSlot::Inline {
kind: Provider::LocalProcess,
fields: Default::default(),
},
);
assert!(slot_declared(&m));
}
}