1use std::collections::BTreeMap;
36use std::net::Ipv4Addr;
37use std::path::{Path, PathBuf};
38use std::sync::Arc;
39use std::time::{Duration, Instant};
40
41use anyhow::{Context, Result};
42use kamaji::native::NativeRuntime;
43use kamaji::{Kamaji, MeshAssignment, MeshIdent};
44use tracing::{info, warn};
45use workload_spec::MeshPort;
46
47use super::native_support::{native_spec, sanitize_ident};
48use crate::capability::Capability;
49use crate::config::{MirrorConfig, Provider, ServiceWithMirrors};
50
51pub const SMTP_DRIVER_IDENT: &str = "yah-smtp-dev";
55
56pub const SMTP_DEV_BIN_ENV: &str = "YAH_SMTP_DEV_BIN";
59
60pub const PORT_NAME_SMTP: &str = "smtp";
65pub const PORT_NAME_HTTP: &str = "http";
66
67const CATCHER_ENVS: &[&str] = &["dev", "pond"];
74
75#[derive(Debug, Clone, Default)]
77pub struct SmtpDriverOptions {
78 pub binary: Option<PathBuf>,
81 pub ready_timeout: Option<Duration>,
85}
86
87impl SmtpDriverOptions {
88 fn resolved_binary(&self) -> PathBuf {
89 if let Some(ref p) = self.binary {
90 return p.clone();
91 }
92 if let Some(p) = std::env::var_os(SMTP_DEV_BIN_ENV) {
93 return PathBuf::from(p);
94 }
95 PathBuf::from("yah-smtp-dev")
96 }
97
98 fn ready_timeout(&self) -> Duration {
99 self.ready_timeout.unwrap_or(Duration::from_secs(120))
100 }
101}
102
103pub struct RunningSmtpDriver {
105 pub smtp_port: u16,
107 pub http_port: u16,
109 pub inbox_url: String,
111 runtime: Arc<NativeRuntime>,
112 ident: MeshIdent,
113}
114
115impl RunningSmtpDriver {
116 pub async fn teardown(&self) {
123 self.runtime.teardown_workload(&self.ident).await.ok();
124 }
125}
126
127pub fn camp_binds_smtp_driver(services: &BTreeMap<String, ServiceWithMirrors>) -> bool {
136 services.values().any(|svc| {
137 CATCHER_ENVS.iter().any(|env| {
138 svc.mirrors
139 .get(*env)
140 .is_some_and(|mirror| binds_local_mailcrab(mirror))
141 })
142 })
143}
144
145fn binds_local_mailcrab(mirror: &MirrorConfig) -> bool {
147 mirror
148 .driver(Capability::Smtp)
149 .and_then(|slot| slot.inline_kind())
150 == Some(Provider::LocalMailcrab)
151}
152
153pub fn coords_path(workspace_root: &Path) -> PathBuf {
155 workspace_root.join(".yah/infra/state/dev/smtp/coords.json")
156}
157
158pub async fn up_smtp_driver(
161 workspace_root: &Path,
162 opts: &SmtpDriverOptions,
163) -> Result<RunningSmtpDriver> {
164 let binary = opts.resolved_binary();
165 let ident_str = sanitize_ident(SMTP_DRIVER_IDENT);
166 let ident = MeshIdent(ident_str.clone());
167
168 let argv: Vec<String> = vec![
169 binary.display().to_string(),
170 "serve".to_string(),
171 "--workspace".to_string(),
172 workspace_root.display().to_string(),
173 ];
174
175 let coords = coords_path(workspace_root);
179 let _ = std::fs::remove_file(&coords);
180
181 let mut spec = native_spec(&ident_str, argv, Vec::new());
182 spec.expose.mesh.ports = vec![
185 MeshPort::named(PORT_NAME_SMTP),
186 MeshPort::named(PORT_NAME_HTTP),
187 ];
188
189 let state_dir = workspace_root.join(".yah/jit/native");
190 let runtime = Arc::new(NativeRuntime::new(&state_dir));
191 let mesh = MeshAssignment::inlined(Ipv4Addr::LOCALHOST);
192
193 info!(
194 binary = %binary.display(),
195 ident = %ident_str,
196 "spawning yah-smtp-dev (kamaji native backend)",
197 );
198
199 runtime
200 .deploy_workload(&spec, &mesh)
201 .await
202 .with_context(|| {
203 format!(
204 "deploying the dev-tier smtp driver via kamaji — install it with \
205 `cargo install --path crates/yah/smtp-dev` or point {SMTP_DEV_BIN_ENV} \
206 at the binary ({})",
207 binary.display(),
208 )
209 })?;
210
211 let timeout = opts.ready_timeout();
212 let Some(ready) = wait_for_coords(&coords, timeout).await else {
213 warn!(timeout = ?timeout, "yah-smtp-dev did not publish coords; tearing down");
214 runtime.teardown_workload(&ident).await.ok();
215 let (_out, err) = super::native_support::capture_paths(&state_dir, &ident_str);
216 anyhow::bail!(
217 "the dev-tier smtp driver did not become ready within {timeout:?} — \
218 check {} for why",
219 err.display(),
220 );
221 };
222
223 info!(
224 smtp_port = ready.smtp_port,
225 http_port = ready.http_port,
226 inbox = %ready.inbox_url,
227 "dev-tier smtp driver ready",
228 );
229 Ok(RunningSmtpDriver {
230 smtp_port: ready.smtp_port,
231 http_port: ready.http_port,
232 inbox_url: ready.inbox_url,
233 runtime,
234 ident,
235 })
236}
237
238#[derive(Debug, Clone, PartialEq, Eq)]
243struct ReadyCoords {
244 smtp_port: u16,
245 http_port: u16,
246 inbox_url: String,
247}
248
249async fn wait_for_coords(path: &Path, timeout: Duration) -> Option<ReadyCoords> {
252 let deadline = Instant::now() + timeout;
253 while Instant::now() < deadline {
254 if let Some(coords) = read_coords(path) {
255 return Some(coords);
256 }
257 tokio::time::sleep(Duration::from_millis(100)).await;
258 }
259 None
260}
261
262fn read_coords(path: &Path) -> Option<ReadyCoords> {
269 let bytes = std::fs::read(path).ok()?;
270 let v: serde_json::Value = serde_json::from_slice(&bytes).ok()?;
271 let port = |key: &str| -> Option<u16> {
272 let n = u16::try_from(v.get(key)?.as_u64()?).ok()?;
273 (n != 0).then_some(n)
274 };
275 let smtp_port = port("smtp_port")?;
276 let http_port = port("http_port")?;
277 let inbox_url = v
278 .get("inbox_url")
279 .and_then(|u| u.as_str())
280 .map(str::to_string)
281 .unwrap_or_else(|| format!("http://127.0.0.1:{http_port}/"));
282 Some(ReadyCoords {
283 smtp_port,
284 http_port,
285 inbox_url,
286 })
287}
288
289#[cfg(test)]
290mod tests {
291 use super::*;
292 use crate::config::ServiceConfig;
293
294 fn service(mirrors: &[(&str, &str)]) -> ServiceWithMirrors {
295 let service: ServiceConfig =
296 toml::from_str("schema_version = 1\nname = \"svc\"\ndomain = \"svc.example\"\n")
297 .expect("parse service");
298 ServiceWithMirrors {
299 service,
300 mirrors: mirrors
301 .iter()
302 .map(|(env, src)| {
303 (
304 (*env).to_string(),
305 toml::from_str::<MirrorConfig>(src).expect("parse mirror"),
306 )
307 })
308 .collect(),
309 component_transform_recipes: BTreeMap::new(),
310 passway_machines: BTreeMap::new(),
311 }
312 }
313
314 const BINDS_SMTP: &str = r#"
315schema_version = 1
316shape = "local"
317[drivers.smtp]
318kind = "local-mailcrab"
319"#;
320
321 const BINDS_PG: &str = r#"
322schema_version = 1
323shape = "local"
324[drivers.pg]
325kind = "local-pg-dev"
326"#;
327
328 fn services(entries: Vec<(&str, ServiceWithMirrors)>) -> BTreeMap<String, ServiceWithMirrors> {
329 entries
330 .into_iter()
331 .map(|(n, s)| (n.to_string(), s))
332 .collect()
333 }
334
335 #[test]
336 fn one_mirror_binding_smtp_activates_the_camps_driver() {
337 let svcs = services(vec![
338 ("quiet", service(&[("dev", BINDS_PG)])),
339 ("mailer", service(&[("dev", BINDS_SMTP)])),
340 ]);
341 assert!(camp_binds_smtp_driver(&svcs));
342 }
343
344 #[test]
345 fn a_camp_that_binds_no_smtp_driver_spawns_nothing() {
346 let svcs = services(vec![("quiet", service(&[("dev", BINDS_PG)]))]);
347 assert!(!camp_binds_smtp_driver(&svcs));
348 }
349
350 #[test]
352 fn pond_counts_and_cloud_does_not() {
353 assert!(camp_binds_smtp_driver(&services(vec![(
354 "mailer",
355 service(&[("pond", BINDS_SMTP)])
356 )])));
357 assert!(!camp_binds_smtp_driver(&services(vec![(
358 "mailer",
359 service(&[("prod", BINDS_SMTP)])
360 )])));
361 }
362
363 #[test]
367 fn the_spec_declares_both_listeners_by_name() {
368 let mut spec = native_spec("yah-smtp-dev", vec!["yah-smtp-dev".to_string()], Vec::new());
369 spec.expose.mesh.ports = vec![
370 MeshPort::named(PORT_NAME_SMTP),
371 MeshPort::named(PORT_NAME_HTTP),
372 ];
373 let names = spec.expose.mesh.names();
374 assert!(names.contains(&"smtp"), "missing smtp: {names:?}");
375 assert!(names.contains(&"http"), "missing http: {names:?}");
376 assert!(spec.expose.mesh.ports.iter().all(|p| p.number.is_none()));
379 workload_spec::validate::shape(&spec).expect("spec must validate");
380 }
381
382 #[tokio::test]
383 async fn coords_are_incomplete_until_both_listeners_report() {
384 let tmp = tempfile::tempdir().unwrap();
385 let path = tmp.path().join("coords.json");
386 let brief = Duration::from_millis(150);
387
388 assert_eq!(wait_for_coords(&path, brief).await, None);
389 std::fs::write(&path, br#"{"smtp_port":1025,"http_port":0}"#).unwrap();
391 assert_eq!(wait_for_coords(&path, brief).await, None);
392 std::fs::write(&path, br#"{"smtp_port":102"#).unwrap();
394 assert_eq!(wait_for_coords(&path, brief).await, None);
395 }
396
397 #[tokio::test]
398 async fn a_complete_coords_file_yields_both_ports_and_the_inbox_url() {
399 let tmp = tempfile::tempdir().unwrap();
400 let path = tmp.path().join("coords.json");
401 std::fs::write(
402 &path,
403 br#"{"smtp_port":51001,"http_port":51002,"inbox_url":"http://127.0.0.1:51002/"}"#,
404 )
405 .unwrap();
406 assert_eq!(
407 wait_for_coords(&path, Duration::from_secs(1)).await,
408 Some(ReadyCoords {
409 smtp_port: 51001,
410 http_port: 51002,
411 inbox_url: "http://127.0.0.1:51002/".to_string(),
412 })
413 );
414 }
415
416 #[test]
419 fn a_missing_inbox_url_is_derived_from_the_http_port() {
420 let tmp = tempfile::tempdir().unwrap();
421 let path = tmp.path().join("coords.json");
422 std::fs::write(&path, br#"{"smtp_port":1025,"http_port":1080}"#).unwrap();
423 assert_eq!(
424 read_coords(&path).unwrap().inbox_url,
425 "http://127.0.0.1:1080/"
426 );
427 }
428
429 #[test]
430 fn binary_resolution_prefers_explicit_over_env_over_path() {
431 let explicit = SmtpDriverOptions {
432 binary: Some(PathBuf::from("/opt/yah-smtp-dev")),
433 ..Default::default()
434 };
435 assert_eq!(
436 explicit.resolved_binary(),
437 PathBuf::from("/opt/yah-smtp-dev")
438 );
439 if std::env::var_os(SMTP_DEV_BIN_ENV).is_none() {
440 assert_eq!(
441 SmtpDriverOptions::default().resolved_binary(),
442 PathBuf::from("yah-smtp-dev")
443 );
444 }
445 }
446}