gdo: the code that opens the iris
A durable consumer on portal's WORMHOLE stream, each message handed to the host's SMTP as text with an HTML alternative; acked on acceptance, retried after a minute otherwise, left for a newer gdo when the contract version is unknown. `gdo probe <to>` sends one through and waits for the consumer to drain. GDO_TYPST_THEME attaches a Typst PDF. A systemd unit and an installer for one gdo per host. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01L4jrCgLiKKHAEFuZUJjckH
This commit is contained in:
co-authored by
Claude Fable 5.1
commit
983491295e
+372
@@ -0,0 +1,372 @@
|
||||
//! gdo - the transmitter that sends the code through the gate so the
|
||||
//! iris opens. Portal (uhhm/portal) decides when a person hears
|
||||
//! something and renders the mail; gdo holds a durable consumer on
|
||||
//! the `WORMHOLE` stream and hands each message to the host's SMTP.
|
||||
//! No site knowledge, no state of its own: if it dies, the stream
|
||||
//! keeps the mail until it is back.
|
||||
//!
|
||||
//! `gdo` run the consumer
|
||||
//! `gdo probe <to>` send one message through and report when it
|
||||
//! was delivered
|
||||
//! `gdo --version`
|
||||
//!
|
||||
//! Env: NATS_URL (nats://127.0.0.1:4222), SMTP_HOST (127.0.0.1),
|
||||
//! SMTP_PORT (25), GDO_CONSUMER (gdo), GDO_TYPST_THEME (a .typ file;
|
||||
//! set, every mail also carries a PDF rendered from it).
|
||||
|
||||
use anyhow::{anyhow, Context};
|
||||
use futures::StreamExt;
|
||||
use lettre::{message::{header::ContentType, Attachment, MultiPart, SinglePart}, AsyncSmtpTransport, AsyncTransport, Message, Tokio1Executor};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::time::Duration;
|
||||
|
||||
/// Portal's side of this contract is `portal::mail`. Same names, same
|
||||
/// stream config, so whichever boots first creates it.
|
||||
const STREAM: &str = "WORMHOLE";
|
||||
const SUBJECTS: &str = "portal.mail.>";
|
||||
const SUBJECT: &str = "portal.mail.send";
|
||||
const CONTRACT: u8 = 1;
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
|
||||
struct Outgoing {
|
||||
v: u8,
|
||||
id: String,
|
||||
site: String,
|
||||
from: String,
|
||||
#[serde(default)]
|
||||
reply_to: Option<String>,
|
||||
to: String,
|
||||
subject: String,
|
||||
text: String,
|
||||
html: String,
|
||||
}
|
||||
|
||||
struct Config {
|
||||
nats_url: String,
|
||||
smtp_host: String,
|
||||
smtp_port: u16,
|
||||
/// What gdo says in EHLO. A strict server (stalwart) refuses a bare
|
||||
/// hostname, so this is the host's fully qualified name, or
|
||||
/// `SMTP_HELO` when the host does not know its own.
|
||||
smtp_helo: String,
|
||||
consumer: String,
|
||||
typst_theme: Option<String>,
|
||||
}
|
||||
|
||||
impl Config {
|
||||
fn from_env() -> Self {
|
||||
let var = |k: &str, d: &str| std::env::var(k).ok().filter(|v| !v.is_empty()).unwrap_or_else(|| d.to_string());
|
||||
Config {
|
||||
nats_url: var("NATS_URL", "nats://127.0.0.1:4222"),
|
||||
smtp_host: var("SMTP_HOST", "127.0.0.1"),
|
||||
smtp_port: var("SMTP_PORT", "25").parse().unwrap_or(25),
|
||||
smtp_helo: var("SMTP_HELO", &fqdn()),
|
||||
consumer: var("GDO_CONSUMER", "gdo"),
|
||||
typst_theme: std::env::var("GDO_TYPST_THEME").ok().filter(|v| !v.is_empty()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn connect(url: &str) -> anyhow::Result<async_nats::Client> {
|
||||
let parsed = url::parse(url)?;
|
||||
let mut opts = async_nats::ConnectOptions::new();
|
||||
if let Some((user, pass)) = parsed {
|
||||
opts = opts.user_and_password(user, pass);
|
||||
}
|
||||
Ok(opts.connect(url).await.with_context(|| format!("connecting to {url}"))?)
|
||||
}
|
||||
|
||||
/// The user and password of a `nats://user:pass@host` URL, without a
|
||||
/// URL crate: the one shape portal's own env uses.
|
||||
mod url {
|
||||
pub fn parse(url: &str) -> anyhow::Result<Option<(String, String)>> {
|
||||
let rest = url.split_once("://").map(|(_, r)| r).unwrap_or(url);
|
||||
let Some((cred, _)) = rest.rsplit_once('@') else { return Ok(None) };
|
||||
let (user, pass) = cred.split_once(':').unwrap_or((cred, ""));
|
||||
Ok(Some((user.to_string(), pass.to_string())))
|
||||
}
|
||||
}
|
||||
|
||||
async fn ensure_stream(js: &async_nats::jetstream::Context) -> anyhow::Result<async_nats::jetstream::stream::Stream> {
|
||||
Ok(js
|
||||
.get_or_create_stream(async_nats::jetstream::stream::Config {
|
||||
name: STREAM.to_string(),
|
||||
subjects: vec![SUBJECTS.to_string()],
|
||||
max_age: Duration::from_secs(7 * 24 * 3600),
|
||||
duplicate_window: Duration::from_secs(24 * 3600),
|
||||
..Default::default()
|
||||
})
|
||||
.await?)
|
||||
}
|
||||
|
||||
/// The mail as SMTP sees it: text and HTML alternatives, and the
|
||||
/// themed PDF when a theme is set.
|
||||
fn build_message(mail: &Outgoing, pdf: Option<Vec<u8>>) -> anyhow::Result<Message> {
|
||||
let mut builder = Message::builder()
|
||||
.from(mail.from.parse().with_context(|| format!("from address {:?}", mail.from))?)
|
||||
.to(mail.to.parse().with_context(|| format!("to address {:?}", mail.to))?)
|
||||
.subject(&mail.subject);
|
||||
if let Some(reply_to) = &mail.reply_to {
|
||||
builder = builder.reply_to(reply_to.parse().with_context(|| format!("reply-to address {reply_to:?}"))?);
|
||||
}
|
||||
let alternative = MultiPart::alternative()
|
||||
.singlepart(SinglePart::builder().header(ContentType::TEXT_PLAIN).body(mail.text.clone()))
|
||||
.singlepart(SinglePart::builder().header(ContentType::TEXT_HTML).body(mail.html.clone()));
|
||||
let body = match pdf {
|
||||
Some(bytes) => MultiPart::mixed()
|
||||
.multipart(alternative)
|
||||
.singlepart(Attachment::new(format!("{}.pdf", mail.site)).body(bytes, ContentType::parse("application/pdf")?)),
|
||||
None => alternative,
|
||||
};
|
||||
Ok(builder.multipart(body)?)
|
||||
}
|
||||
|
||||
/// A Typst theme: `typst compile` on the theme with the mail's fields
|
||||
/// as `--input`s (`site`, `subject`, `text`, `to`), giving a PDF. A
|
||||
/// theme that fails to render costs the attachment, never the mail.
|
||||
async fn render_theme(theme: &str, mail: &Outgoing) -> Option<Vec<u8>> {
|
||||
let out = std::env::temp_dir().join(format!("gdo-{}.pdf", mail.id.replace(['/', ':'], "_")));
|
||||
let status = tokio::process::Command::new("typst")
|
||||
.arg("compile")
|
||||
.arg("--input").arg(format!("site={}", mail.site))
|
||||
.arg("--input").arg(format!("subject={}", mail.subject))
|
||||
.arg("--input").arg(format!("text={}", mail.text))
|
||||
.arg("--input").arg(format!("to={}", mail.to))
|
||||
.arg(theme)
|
||||
.arg(&out)
|
||||
.stdout(std::process::Stdio::null())
|
||||
.stderr(std::process::Stdio::inherit())
|
||||
.status()
|
||||
.await;
|
||||
match status {
|
||||
Ok(s) if s.success() => {
|
||||
let bytes = tokio::fs::read(&out).await.ok();
|
||||
let _ = tokio::fs::remove_file(&out).await;
|
||||
bytes
|
||||
}
|
||||
Ok(s) => {
|
||||
tracing::warn!(theme, id = %mail.id, "typst exited with {s}; sending without the PDF");
|
||||
None
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(theme, "typst could not be run: {e}; sending without the PDF");
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn serve(cfg: Config) -> anyhow::Result<()> {
|
||||
let nats = connect(&cfg.nats_url).await?;
|
||||
let js = async_nats::jetstream::new(nats);
|
||||
let stream = ensure_stream(&js).await?;
|
||||
let transport = AsyncSmtpTransport::<Tokio1Executor>::builder_dangerous(&cfg.smtp_host)
|
||||
.port(cfg.smtp_port)
|
||||
.hello_name(lettre::transport::smtp::extension::ClientId::Domain(cfg.smtp_helo.clone()))
|
||||
.build();
|
||||
match transport.test_connection().await {
|
||||
Ok(true) => tracing::info!(host = %cfg.smtp_host, port = cfg.smtp_port, helo = %cfg.smtp_helo, "chevron seven, locked"),
|
||||
Ok(false) => return Err(anyhow!("SMTP at {}:{} did not answer", cfg.smtp_host, cfg.smtp_port)),
|
||||
Err(e) => return Err(anyhow!("SMTP at {}:{} refused the greeting as {:?}: {e} (set SMTP_HELO to this host's full name)", cfg.smtp_host, cfg.smtp_port, cfg.smtp_helo)),
|
||||
}
|
||||
let consumer = stream
|
||||
.get_or_create_consumer(
|
||||
&cfg.consumer,
|
||||
async_nats::jetstream::consumer::pull::Config {
|
||||
durable_name: Some(cfg.consumer.clone()),
|
||||
filter_subject: SUBJECT.to_string(),
|
||||
ack_policy: async_nats::jetstream::consumer::AckPolicy::Explicit,
|
||||
ack_wait: Duration::from_secs(120),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
let mut messages = consumer.messages().await?;
|
||||
tracing::info!(consumer = %cfg.consumer, "listening on {STREAM}");
|
||||
while let Some(next) = messages.next().await {
|
||||
let msg = match next {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
tracing::warn!("consumer stream error: {e}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let mail: Outgoing = match serde_json::from_slice(&msg.payload) {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
// Not ours to understand; acked so it does not block the
|
||||
// rest, logged so someone sees it.
|
||||
tracing::error!("unreadable message on {SUBJECT}: {e}");
|
||||
let _ = msg.ack().await;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
if mail.v != CONTRACT {
|
||||
tracing::error!(id = %mail.id, v = mail.v, "contract version {CONTRACT} expected; leaving it for a gdo that knows it");
|
||||
let _ = msg.ack_with(async_nats::jetstream::AckKind::Nak(Some(Duration::from_secs(3600)))).await;
|
||||
continue;
|
||||
}
|
||||
let pdf = match &cfg.typst_theme {
|
||||
Some(theme) => render_theme(theme, &mail).await,
|
||||
None => None,
|
||||
};
|
||||
let message = match build_message(&mail, pdf) {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
tracing::error!(id = %mail.id, to = %mail.to, "cannot build the message: {e}; dropped");
|
||||
let _ = msg.ack().await;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
match transport.send(message).await {
|
||||
Ok(_) => {
|
||||
tracing::info!(id = %mail.id, to = %mail.to, site = %mail.site, "sent");
|
||||
if let Err(e) = msg.double_ack().await {
|
||||
tracing::warn!(id = %mail.id, "sent, but the ack failed: {e}; it may go twice");
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(id = %mail.id, to = %mail.to, "SMTP refused: {e}; retrying in a minute");
|
||||
let _ = msg.ack_with(async_nats::jetstream::AckKind::Nak(Some(Duration::from_secs(60)))).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Publish one message and watch the consumer's pending count go back
|
||||
/// to zero: delivered means the running gdo handed it to SMTP and
|
||||
/// acked, which is what delivered means from here.
|
||||
async fn probe(cfg: Config, to: &str, from: &str) -> anyhow::Result<()> {
|
||||
let nats = connect(&cfg.nats_url).await?;
|
||||
let js = async_nats::jetstream::new(nats);
|
||||
let stream = ensure_stream(&js).await?;
|
||||
let id = format!("probe:{}", std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis());
|
||||
let mail = Outgoing {
|
||||
v: CONTRACT,
|
||||
id: id.clone(),
|
||||
site: "gdo".into(),
|
||||
from: from.into(),
|
||||
reply_to: None,
|
||||
to: to.into(),
|
||||
subject: "gdo probe: chevron seven".into(),
|
||||
text: format!("A probe through the gate, id {id}. If you read this, gdo delivers on this host."),
|
||||
html: format!("<!doctype html><html><body><p>A probe through the gate, id {id}. If you read this, gdo delivers on this host.</p></body></html>"),
|
||||
};
|
||||
let mut headers = async_nats::HeaderMap::new();
|
||||
headers.insert("Nats-Msg-Id", id.as_str());
|
||||
js.publish_with_headers(SUBJECT, headers, serde_json::to_vec(&mail)?.into()).await?.await?;
|
||||
println!("sent {id} into {STREAM}");
|
||||
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
|
||||
loop {
|
||||
let Ok(mut consumer) = stream.get_consumer::<async_nats::jetstream::consumer::pull::Config>(&cfg.consumer).await else {
|
||||
println!("no consumer {:?} on {STREAM}: gdo has never run against this NATS; the probe waits in the stream", cfg.consumer);
|
||||
return Ok(());
|
||||
};
|
||||
let info = consumer.info().await?;
|
||||
if info.num_pending == 0 && info.num_ack_pending == 0 {
|
||||
println!("delivered: gdo acked everything up to and including the probe");
|
||||
return Ok(());
|
||||
}
|
||||
if tokio::time::Instant::now() > deadline {
|
||||
println!("not delivered within a minute: {} pending, {} awaiting ack - is gdo running, and does SMTP accept?", info.num_pending, info.num_ack_pending);
|
||||
return Ok(());
|
||||
}
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> anyhow::Result<()> {
|
||||
tracing_subscriber::fmt().with_env_filter(tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| "info".into())).init();
|
||||
let args: Vec<String> = std::env::args().skip(1).collect();
|
||||
let cfg = Config::from_env();
|
||||
match args.first().map(String::as_str) {
|
||||
None => serve(cfg).await,
|
||||
Some("probe") => {
|
||||
let to = args.get(1).ok_or_else(|| anyhow!("usage: gdo probe <to> [--from <address>]"))?;
|
||||
let from = args.iter().position(|a| a == "--from").and_then(|i| args.get(i + 1)).cloned().unwrap_or_else(|| format!("gdo@{}", hostname()));
|
||||
probe(cfg, to, &from).await
|
||||
}
|
||||
Some("--version") | Some("-V") => {
|
||||
println!("gdo {} (contract v{CONTRACT})", env!("CARGO_PKG_VERSION"));
|
||||
Ok(())
|
||||
}
|
||||
Some(other) => Err(anyhow!("unknown argument {other:?}; usage: gdo | gdo probe <to> [--from <address>] | gdo --version")),
|
||||
}
|
||||
}
|
||||
|
||||
fn hostname() -> String {
|
||||
std::fs::read_to_string("/etc/hostname").ok().map(|s| s.trim().to_string()).filter(|s| !s.is_empty()).unwrap_or_else(|| "localhost".into())
|
||||
}
|
||||
|
||||
/// The host's full name for EHLO: /etc/hostname when it has a dot,
|
||||
/// else the first dotted name /etc/hosts gives it, else the bare name
|
||||
/// with `.localdomain`, which strict servers at least parse.
|
||||
fn fqdn() -> String {
|
||||
let short = hostname();
|
||||
if short.contains('.') {
|
||||
return short;
|
||||
}
|
||||
if let Ok(hosts) = std::fs::read_to_string("/etc/hosts") {
|
||||
for line in hosts.lines() {
|
||||
let names: Vec<&str> = line.split('#').next().unwrap_or("").split_whitespace().skip(1).collect();
|
||||
if names.contains(&short.as_str()) {
|
||||
if let Some(full) = names.iter().find(|n| n.contains('.')) {
|
||||
return full.to_string();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
format!("{short}.localdomain")
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn mail() -> Outgoing {
|
||||
Outgoing {
|
||||
v: CONTRACT,
|
||||
id: "reports:abc:fixed".into(),
|
||||
site: "Tomter Vel".into(),
|
||||
from: "Tomter Vel <vel@example.no>".into(),
|
||||
reply_to: Some("post@example.no".into()),
|
||||
to: "Kari <kari@example.no>".into(),
|
||||
subject: "Ordnet".into(),
|
||||
text: "Det er ordnet.".into(),
|
||||
html: "<p>Det er ordnet.</p>".into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_message_carries_both_alternatives_and_its_headers() {
|
||||
let m = build_message(&mail(), None).unwrap();
|
||||
let raw = String::from_utf8(m.formatted()).unwrap();
|
||||
assert!(raw.contains("From: \"Tomter Vel\" <vel@example.no>"));
|
||||
assert!(raw.contains("Reply-To: post@example.no"));
|
||||
assert!(raw.contains("Subject: Ordnet"));
|
||||
assert!(raw.contains("multipart/alternative"));
|
||||
assert!(raw.contains("text/plain") && raw.contains("text/html"));
|
||||
assert!(!raw.contains("application/pdf"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_theme_makes_it_mixed_with_a_pdf() {
|
||||
let m = build_message(&mail(), Some(b"%PDF-1.7".to_vec())).unwrap();
|
||||
let raw = String::from_utf8(m.formatted()).unwrap();
|
||||
assert!(raw.contains("multipart/mixed") && raw.contains("application/pdf") && raw.contains("Tomter Vel.pdf"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_bad_address_is_refused_before_smtp() {
|
||||
let mut bad = mail();
|
||||
bad.to = "not an address".into();
|
||||
assert!(build_message(&bad, None).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_nats_url_credentials_are_read() {
|
||||
assert_eq!(url::parse("nats://infra:secret@127.0.0.1:4222").unwrap(), Some(("infra".into(), "secret".into())));
|
||||
assert_eq!(url::parse("nats://127.0.0.1:4222").unwrap(), None);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user