use std::collections::BTreeMap;
use std::net::Ipv4Addr;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use kamaji::native::NativeRuntime;
use kamaji::{Kamaji, MeshAssignment, MeshIdent};
use tracing::{info, warn};
use workload_spec::MeshPort;
use super::native_support::{native_spec, sanitize_ident};
use crate::capability::Capability;
use crate::config::{MirrorConfig, Provider, ServiceWithMirrors};
pub const SMTP_DRIVER_IDENT: &str = "yah-smtp-dev";
pub const SMTP_DEV_BIN_ENV: &str = "YAH_SMTP_DEV_BIN";
pub const PORT_NAME_SMTP: &str = "smtp";
pub const PORT_NAME_HTTP: &str = "http";
const CATCHER_ENVS: &[&str] = &["dev", "pond"];
#[derive(Debug, Clone, Default)]
pub struct SmtpDriverOptions {
pub binary: Option<PathBuf>,
pub ready_timeout: Option<Duration>,
}
impl SmtpDriverOptions {
fn resolved_binary(&self) -> PathBuf {
if let Some(ref p) = self.binary {
return p.clone();
}
if let Some(p) = std::env::var_os(SMTP_DEV_BIN_ENV) {
return PathBuf::from(p);
}
PathBuf::from("yah-smtp-dev")
}
fn ready_timeout(&self) -> Duration {
self.ready_timeout.unwrap_or(Duration::from_secs(120))
}
}
pub struct RunningSmtpDriver {
pub smtp_port: u16,
pub http_port: u16,
pub inbox_url: String,
runtime: Arc<NativeRuntime>,
ident: MeshIdent,
}
impl RunningSmtpDriver {
pub async fn teardown(&self) {
self.runtime.teardown_workload(&self.ident).await.ok();
}
}
pub fn camp_binds_smtp_driver(services: &BTreeMap<String, ServiceWithMirrors>) -> bool {
services.values().any(|svc| {
CATCHER_ENVS.iter().any(|env| {
svc.mirrors
.get(*env)
.is_some_and(|mirror| binds_local_mailcrab(mirror))
})
})
}
fn binds_local_mailcrab(mirror: &MirrorConfig) -> bool {
mirror
.driver(Capability::Smtp)
.and_then(|slot| slot.inline_kind())
== Some(Provider::LocalMailcrab)
}
pub fn coords_path(workspace_root: &Path) -> PathBuf {
workspace_root.join(".yah/infra/state/dev/smtp/coords.json")
}
pub async fn up_smtp_driver(
workspace_root: &Path,
opts: &SmtpDriverOptions,
) -> Result<RunningSmtpDriver> {
let binary = opts.resolved_binary();
let ident_str = sanitize_ident(SMTP_DRIVER_IDENT);
let ident = MeshIdent(ident_str.clone());
let argv: Vec<String> = vec![
binary.display().to_string(),
"serve".to_string(),
"--workspace".to_string(),
workspace_root.display().to_string(),
];
let coords = coords_path(workspace_root);
let _ = std::fs::remove_file(&coords);
let mut spec = native_spec(&ident_str, argv, Vec::new());
spec.expose.mesh.ports = vec![
MeshPort::named(PORT_NAME_SMTP),
MeshPort::named(PORT_NAME_HTTP),
];
let state_dir = workspace_root.join(".yah/jit/native");
let runtime = Arc::new(NativeRuntime::new(&state_dir));
let mesh = MeshAssignment::inlined(Ipv4Addr::LOCALHOST);
info!(
binary = %binary.display(),
ident = %ident_str,
"spawning yah-smtp-dev (kamaji native backend)",
);
runtime
.deploy_workload(&spec, &mesh)
.await
.with_context(|| {
format!(
"deploying the dev-tier smtp driver via kamaji — install it with \
`cargo install --path crates/yah/smtp-dev` or point {SMTP_DEV_BIN_ENV} \
at the binary ({})",
binary.display(),
)
})?;
let timeout = opts.ready_timeout();
let Some(ready) = wait_for_coords(&coords, timeout).await else {
warn!(timeout = ?timeout, "yah-smtp-dev did not publish coords; tearing down");
runtime.teardown_workload(&ident).await.ok();
anyhow::bail!(super::native_support::ready_timeout_message(
"smtp",
timeout,
&state_dir,
&ident_str,
));
};
info!(
smtp_port = ready.smtp_port,
http_port = ready.http_port,
inbox = %ready.inbox_url,
"dev-tier smtp driver ready",
);
Ok(RunningSmtpDriver {
smtp_port: ready.smtp_port,
http_port: ready.http_port,
inbox_url: ready.inbox_url,
runtime,
ident,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ReadyCoords {
smtp_port: u16,
http_port: u16,
inbox_url: String,
}
async fn wait_for_coords(path: &Path, timeout: Duration) -> Option<ReadyCoords> {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if let Some(coords) = read_coords(path) {
return Some(coords);
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
None
}
fn read_coords(path: &Path) -> Option<ReadyCoords> {
let bytes = std::fs::read(path).ok()?;
let v: serde_json::Value = serde_json::from_slice(&bytes).ok()?;
let port = |key: &str| -> Option<u16> {
let n = u16::try_from(v.get(key)?.as_u64()?).ok()?;
(n != 0).then_some(n)
};
let smtp_port = port("smtp_port")?;
let http_port = port("http_port")?;
let inbox_url = v
.get("inbox_url")
.and_then(|u| u.as_str())
.map(str::to_string)
.unwrap_or_else(|| format!("http://127.0.0.1:{http_port}/"));
Some(ReadyCoords {
smtp_port,
http_port,
inbox_url,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::ServiceConfig;
fn service(mirrors: &[(&str, &str)]) -> ServiceWithMirrors {
let service: ServiceConfig =
toml::from_str("schema_version = 1\nname = \"svc\"\n[address]\nkind = \"front-door\"\ndomain = \"svc.example\"\n")
.expect("parse service");
ServiceWithMirrors {
service,
mirrors: mirrors
.iter()
.map(|(env, src)| {
(
(*env).to_string(),
toml::from_str::<MirrorConfig>(src).expect("parse mirror"),
)
})
.collect(),
component_transform_recipes: BTreeMap::new(),
passway_machines: BTreeMap::new(),
}
}
const BINDS_SMTP: &str = r#"
schema_version = 1
shape = "local"
[drivers.smtp]
kind = "local-mailcrab"
"#;
const BINDS_PG: &str = r#"
schema_version = 1
shape = "local"
[drivers.pg]
kind = "local-pg-dev"
"#;
fn services(entries: Vec<(&str, ServiceWithMirrors)>) -> BTreeMap<String, ServiceWithMirrors> {
entries
.into_iter()
.map(|(n, s)| (n.to_string(), s))
.collect()
}
#[test]
fn one_mirror_binding_smtp_activates_the_camps_driver() {
let svcs = services(vec![
("quiet", service(&[("dev", BINDS_PG)])),
("mailer", service(&[("dev", BINDS_SMTP)])),
]);
assert!(camp_binds_smtp_driver(&svcs));
}
#[test]
fn a_camp_that_binds_no_smtp_driver_spawns_nothing() {
let svcs = services(vec![("quiet", service(&[("dev", BINDS_PG)]))]);
assert!(!camp_binds_smtp_driver(&svcs));
}
#[test]
fn pond_counts_and_cloud_does_not() {
assert!(camp_binds_smtp_driver(&services(vec![(
"mailer",
service(&[("pond", BINDS_SMTP)])
)])));
assert!(!camp_binds_smtp_driver(&services(vec![(
"mailer",
service(&[("prod", BINDS_SMTP)])
)])));
}
#[test]
fn the_spec_declares_both_listeners_by_name() {
let mut spec = native_spec("yah-smtp-dev", vec!["yah-smtp-dev".to_string()], Vec::new());
spec.expose.mesh.ports = vec![
MeshPort::named(PORT_NAME_SMTP),
MeshPort::named(PORT_NAME_HTTP),
];
let names = spec.expose.mesh.names();
assert!(names.contains(&"smtp"), "missing smtp: {names:?}");
assert!(names.contains(&"http"), "missing http: {names:?}");
assert!(spec.expose.mesh.ports.iter().all(|p| p.number.is_none()));
workload_spec::validate::shape(&spec).expect("spec must validate");
}
#[tokio::test]
async fn coords_are_incomplete_until_both_listeners_report() {
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("coords.json");
let brief = Duration::from_millis(150);
assert_eq!(wait_for_coords(&path, brief).await, None);
std::fs::write(&path, br#"{"smtp_port":1025,"http_port":0}"#).unwrap();
assert_eq!(wait_for_coords(&path, brief).await, None);
std::fs::write(&path, br#"{"smtp_port":102"#).unwrap();
assert_eq!(wait_for_coords(&path, brief).await, None);
}
#[tokio::test]
async fn a_complete_coords_file_yields_both_ports_and_the_inbox_url() {
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("coords.json");
std::fs::write(
&path,
br#"{"smtp_port":51001,"http_port":51002,"inbox_url":"http://127.0.0.1:51002/"}"#,
)
.unwrap();
assert_eq!(
wait_for_coords(&path, Duration::from_secs(1)).await,
Some(ReadyCoords {
smtp_port: 51001,
http_port: 51002,
inbox_url: "http://127.0.0.1:51002/".to_string(),
})
);
}
#[test]
fn a_missing_inbox_url_is_derived_from_the_http_port() {
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("coords.json");
std::fs::write(&path, br#"{"smtp_port":1025,"http_port":1080}"#).unwrap();
assert_eq!(
read_coords(&path).unwrap().inbox_url,
"http://127.0.0.1:1080/"
);
}
#[test]
fn binary_resolution_prefers_explicit_over_env_over_path() {
let explicit = SmtpDriverOptions {
binary: Some(PathBuf::from("/opt/yah-smtp-dev")),
..Default::default()
};
assert_eq!(
explicit.resolved_binary(),
PathBuf::from("/opt/yah-smtp-dev")
);
if std::env::var_os(SMTP_DEV_BIN_ENV).is_none() {
assert_eq!(
SmtpDriverOptions::default().resolved_binary(),
PathBuf::from("yah-smtp-dev")
);
}
}
}