//! Minimal SMTP/LMTP-ish receiver for inbound user mail. //! //! Supported commands: EHLO, HELO, AUTH PLAIN, MAIL FROM, RCPT TO, DATA, //! RSET, QUIT, NOOP. No TLS (plaintext on a trusted network); AUTH PLAIN //! (RFC 4954) is required before MAIL FROM. The authenticated identity //! authorizes the envelope sender (RFC 4954 §4); RCPT TO must be under the //! daemon domain. use std::net::SocketAddr; use std::sync::Arc; use anyhow::Context; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{TcpListener, TcpStream}; use crate::config::{DaemonConfig, SmtpConfig}; use crate::store::Store; /// A condition that will never succeed on retry (bad recipient, unknown /// provider, unknown sender). Maps to SMTP 550, not 451 — tempfail tells the /// client to retry forever. Transient failures (DB down, etc.) stay as /// `anyhow::Error` and yield 451. #[derive(Debug, thiserror::Error)] #[error("{0}")] struct Permanent(String); impl Permanent { fn new(msg: impl Into) -> Self { Self(msg.into()) } } #[derive(Debug, Clone)] pub struct SmtpServer { cfg: SmtpConfig, daemon: DaemonConfig, store: Arc, } impl SmtpServer { pub fn new(cfg: SmtpConfig, daemon: DaemonConfig, store: Arc) -> Self { Self { cfg, daemon, store } } pub async fn listen(&self) -> anyhow::Result { let listener = TcpListener::bind((self.cfg.host, self.cfg.port)) .await .with_context(|| format!("bind SMTP {}:{}", self.cfg.host, self.cfg.port))?; let addr = listener.local_addr()?; tracing::info!(%addr, "SMTP server listening"); let store = self.store.clone(); let daemon = self.daemon.clone(); let max_bytes = self.cfg.max_message_bytes; tokio::spawn(async move { loop { let (stream, peer) = match listener.accept().await { Ok(v) => v, Err(e) => { tracing::warn!(error = %e, "SMTP accept failed"); continue; } }; let store = store.clone(); let daemon = daemon.clone(); tokio::spawn(async move { tracing::debug!(%peer, "SMTP client connected"); if let Err(e) = handle_client(stream, peer, store, daemon, max_bytes).await { tracing::warn!(error = %e, peer = %peer, "SMTP session ended with error"); } }); } }); Ok(addr) } } #[derive(Debug, Default)] struct SessionState { mail_from: Option, rcpt_to: Vec, /// Set once AUTH PLAIN succeeds. MAIL/RCPT/DATA require it, matching /// IMAP. Empty until authenticated. authed_user: Option, } async fn handle_client( stream: TcpStream, peer: SocketAddr, store: Arc, daemon: DaemonConfig, max_message_bytes: usize, ) -> anyhow::Result<()> { let (reader, mut writer) = stream.into_split(); let mut reader = BufReader::new(reader); let mut line: Vec = Vec::new(); send(&mut writer, "220 posthorn ESMTP ready\r\n").await?; let mut state = SessionState::default(); loop { line.clear(); let n = reader.read_until(b'\n', &mut line).await?; if n == 0 { tracing::debug!(peer = %peer, "SMTP client disconnected"); return Ok(()); } // Lossy-decode for parsing: addresses and commands are ASCII, so any // non-UTF-8 bytes here (raw 8-bit in EHLO/MAIL FROM) degrade to U+FFFD // without killing the session the way read_line on a String would. let line_str = String::from_utf8_lossy(&line); let cmd = line_str.trim(); if cmd.is_empty() { continue; } tracing::debug!(peer = %peer, cmd = %cmd, "SMTP command"); let upper = cmd.to_ascii_uppercase(); if upper.starts_with("QUIT ") || upper == "QUIT" { send(&mut writer, "221 bye\r\n").await?; return Ok(()); } if upper.starts_with("NOOP") { send(&mut writer, "250 OK\r\n").await?; continue; } if upper.starts_with("RSET") { state = SessionState::default(); send(&mut writer, "250 OK\r\n").await?; continue; } if upper.starts_with("EHLO ") || upper.starts_with("HELO ") { state.mail_from = None; state.rcpt_to.clear(); if upper.starts_with("EHLO ") { // We are now genuinely 8-bit-clean end to end (read_data reads // bytes); advertise it, the SIZE limit we enforce, and AUTH // PLAIN so real MUAs (which require SMTP auth) will proceed. // No STARTTLS: v1 is plaintext on a trusted network. let resp = format!( "250-posthorn\r\n250-8BITMIME\r\n250-PIPELINING\r\n\ 250-AUTH PLAIN\r\n250 SIZE {}\r\n", max_message_bytes ); send(&mut writer, &resp).await?; } else { send(&mut writer, "250 posthorn\r\n").await?; } continue; } if upper.starts_with("AUTH ") { // RFC 4954. We support AUTH PLAIN, either inline // (`AUTH PLAIN `) or as a two-step challenge // (`AUTH PLAIN` -> `334 ` -> client sends b64). let rest = &cmd[5..]; if !rest.to_ascii_uppercase().starts_with("PLAIN") { send(&mut writer, "504 unsupported mechanism\r\n").await?; continue; } let inline = rest["PLAIN".len()..].trim(); let blob = if inline.is_empty() { // Two-step: issue an empty challenge. The client replies with // the base64 SASL response, or `*` to cancel. send(&mut writer, "334 \r\n").await?; line.clear(); if reader.read_until(b'\n', &mut line).await? == 0 { return Ok(()); } String::from_utf8_lossy(&line).trim().to_string() } else { inline.to_string() }; if blob == "*" { send(&mut writer, "501 auth cancelled\r\n").await?; continue; } match decode_auth_plain(&blob) { Some((user, pass)) => { if let Some(rec) = store.find_user_by_email(&user.to_lowercase()).await { if rec.password == pass { state.authed_user = Some(rec.email); send(&mut writer, "235 OK authenticated\r\n").await?; continue; } } tracing::warn!(peer = %peer, "SMTP AUTH failed"); send(&mut writer, "535 auth failed\r\n").await?; continue; } None => { send(&mut writer, "501 bad auth response\r\n").await?; continue; } } } if upper.starts_with("MAIL FROM:") { // Require AUTH before accepting mail, matching IMAP. The old // "look up MAIL FROM as a known user" gate was implicit auth and // left MUAs that wait for an advertised AUTH extension stranded. if state.authed_user.is_none() { send(&mut writer, "530 auth required\r\n").await?; continue; } let from = extract_addr(cmd); if from.is_empty() { send(&mut writer, "501 bad address\r\n").await?; continue; } // The sender must be the authenticated user (RFC 4954 §4: the // SASL identity authorizes the envelope). `from` is already // lowercased by extract_addr; authed_user is the stored email. if Some(&from) != state.authed_user.as_ref() { tracing::warn!( from = %from, authed = ?state.authed_user, peer = %peer, "rejecting sender mismatch", ); send(&mut writer, &format!("550 not authorized as {from}\r\n")).await?; continue; } state.mail_from = Some(from); state.rcpt_to.clear(); send(&mut writer, "250 OK\r\n").await?; continue; } if upper.starts_with("RCPT TO:") { let to = extract_addr(cmd); if to.is_empty() { send(&mut writer, "501 bad address\r\n").await?; continue; } if !is_local_recipient(&to, &daemon.domain) { send(&mut writer, &format!("550 not local: {to}\r\n")).await?; continue; } if state.mail_from.is_none() { send(&mut writer, "503 need MAIL FROM first\r\n").await?; continue; } state.rcpt_to.push(to); send(&mut writer, "250 OK\r\n").await?; continue; } if upper == "DATA" { if state.mail_from.is_none() || state.rcpt_to.is_empty() { send(&mut writer, "503 need MAIL FROM and RCPT TO\r\n").await?; continue; } send(&mut writer, "354 end data with .\r\n").await?; let raw = read_data(&mut reader, max_message_bytes).await?; tracing::debug!( peer = %peer, from = ?state.mail_from, rcpts = ?state.rcpt_to, bytes = raw.len(), "SMTP DATA received", ); match process_message( &store, &daemon.domain, state.mail_from.as_deref().unwrap_or(""), &state.rcpt_to, &raw, ) .await { Ok(_) => send(&mut writer, "250 OK message queued\r\n").await?, Err(e) if e.downcast_ref::().is_some() => { tracing::warn!(error = %e, "rejecting message permanently"); send(&mut writer, &format!("550 {e}\r\n")).await?; } Err(e) => { tracing::warn!(error = %e, "process_message failed"); send(&mut writer, "451 temporary failure\r\n").await?; } } state.rcpt_to.clear(); continue; } send(&mut writer, "500 command unrecognized\r\n").await?; } } async fn send(writer: &mut tokio::net::tcp::OwnedWriteHalf, msg: &str) -> anyhow::Result<()> { writer.write_all(msg.as_bytes()).await?; writer.flush().await?; Ok(()) } fn extract_addr(cmd: &str) -> String { let start = cmd.find(':').map(|i| i + 1).unwrap_or(cmd.len()); let s = &cmd[start..]; let s = s.trim(); // ESMTP parameters may follow the address: `MAIL FROM: BODY=8BITMIME`. // Pull out the angle-bracketed addr; if unbracketed, take the first token. // Use the *last* `<` so a display name wrapped in the envelope — which some // clients emit as `>` — yields the inner addr-spec // (`a@b`) instead of `Display Name `). For a plain // `` there is only one `<`, so this is identical to the old behavior. let s = match (s.rfind('<'), s.find('>')) { (Some(l), Some(r)) if l < r => &s[l + 1..r], _ => s.split_whitespace().next().unwrap_or(""), }; s.trim().to_lowercase() } /// Decode an RFC 4616 SASL PLAIN response (base64 of /// `authzid \0 authcid \0 passwd`). Returns (authcid, passwd); the authzid /// is ignored (single-user model). `None` on a bad base64 or missing fields. fn decode_auth_plain(b64: &str) -> Option<(String, String)> { use base64::Engine; let bytes = base64::engine::general_purpose::STANDARD .decode(b64.trim()) .ok()?; let s = String::from_utf8_lossy(&bytes); // RFC 4616: authzid \0 authcid \0 passwd. authzid may be empty. let mut parts = s.splitn(3, '\0'); let _authzid = parts.next(); let authcid = parts.next()?.trim().to_string(); let passwd = parts.next()?.to_string(); if authcid.is_empty() || passwd.is_empty() { return None; } Some((authcid, passwd)) } fn is_local_recipient(addr: &str, domain: &str) -> bool { addr.split('@') .nth(1) .map(|h| h.eq_ignore_ascii_case(domain) || h.ends_with(&format!(".{}", domain))) .unwrap_or(false) } async fn read_data(reader: &mut R, max_bytes: usize) -> anyhow::Result> where R: tokio::io::AsyncBufRead + Unpin, { use tokio::io::{AsyncBufReadExt, AsyncReadExt}; let mut buf = Vec::with_capacity(8192); let mut line: Vec = Vec::new(); loop { line.clear(); // Cap each read: a newline-less stream would otherwise buffer without // limit (read_until, like read_line, only returns at EOF). The +8 // covers the CRLF/terminator slop so the post-line max check is the // authoritative bound. let budget = (max_bytes.saturating_sub(buf.len()) + 8) as u64; let n = (&mut *reader).take(budget).read_until(b'\n', &mut line).await?; if n == 0 { anyhow::bail!("client disconnected during DATA"); } if !line.ends_with(b"\n") { anyhow::bail!("message exceeds max size"); } // Strip the trailing newline (CRLF or bare LF), then re-add CRLF so // stored .emls are canonical: the IMAP BODY[TEXT] path assumes CRLF. line.pop(); if line.last() == Some(&b'\r') { line.pop(); } // Terminator check *after* newline strip: a lone "." line ends DATA. if line == b"." { break; } // RFC 5321 §4.5.2 dot-unstuffing: a leading ".." is an escaped ".". let content = unstuff_line(&line); if buf.len() + content.len() + 2 > max_bytes { anyhow::bail!("message exceeds max size"); } buf.extend_from_slice(content); buf.extend_from_slice(b"\r\n"); } Ok(buf) } /// Reverse RFC 5321 §4.5.2 transparency: a body line beginning with ".." is /// an escaped single dot. The terminator "." is handled by the caller. fn unstuff_line(line: &[u8]) -> &[u8] { if line.starts_with(b"..") { &line[1..] } else { line } } async fn process_message( store: &Store, domain: &str, from: &str, recipients: &[String], raw: &[u8], ) -> anyhow::Result<()> { let headers = crate::mailparse::parse_headers(raw)?; let subject = headers .subject .clone() .unwrap_or_else(|| "(no subject)".to_string()); // Pick the first local recipient that parses as model@provider.domain. let (provider_name, model_id) = recipients .iter() .find_map(|addr| parse_recipient(addr, domain)) .ok_or_else(|| { Permanent::new(format!( "no valid model recipient (got: {}); address mail to @.{domain}", recipients.join(", ") )) })?; let user = store .find_user_by_email(from) .await .ok_or_else(|| Permanent::new(format!("sender not found: {from}")))?; let _provider = store .find_provider_by_name(&provider_name) .await .ok_or_else(|| Permanent::new(format!("unknown provider: {provider_name}")))?; // Resolve thread via In-Reply-To, then References ids newest-first. // Both are whitespace-separated message-id lists (RFC 5322 §3.6.4). let mut candidates: Vec<&str> = Vec::new(); if let Some(irt) = headers.in_reply_to.as_deref() { candidates.extend(irt.split_whitespace()); } if let Some(refs) = headers.references.as_deref() { candidates.extend(refs.split_whitespace().rev()); } let mut thread_seed = None; for id in candidates { if let Some(m) = store.find_message_by_message_id(&user.email, id)? { thread_seed = Some(m); break; } } let thread_id = match &thread_seed { Some(seed) => Store::thread_id_for( &user.email, &provider_name, &model_id, &seed.subject, ), None => Store::thread_id_for(&user.email, &provider_name, &model_id, &subject), }; let message_id = headers .message_id .clone() .unwrap_or_else(|| format!("", uuid::Uuid::new_v4(), domain)); // A client retrying a timed-out DATA re-sends the same Message-ID; the // maildir would gain an orphan .eml per attempt. Ack duplicates before // touching disk. We check by re-reading the maildir: same Message-ID // already in the user's maildir means a duplicate. if store .find_message_by_message_id(&user.email, &message_id)? .is_some() { tracing::info!(message_id, "duplicate delivery; acking without requeue"); return Ok(()); } // The stored .eml must carry the Message-ID we record (so the LLM // reply's In-Reply-To/References resolve) and a Date (so clients can // sort the thread). Inject whichever the sender omitted, prepending to // the header block. Keep `message_id` authoritative so the on-disk // and in-memory views agree. let mut prepend = String::new(); if headers.message_id.is_none() { prepend.push_str(&format!("Message-ID: {}\r\n", message_id)); } if headers.date.is_none() { prepend.push_str(&format!("Date: {}\r\n", chrono::Utc::now().to_rfc2822())); } let final_raw: Vec = if prepend.is_empty() { raw.to_vec() } else { let mut out = prepend.into_bytes(); out.extend_from_slice(raw); out }; let raw_path = store .maildir() .deliver(&user.email, &final_raw) .with_context(|| "deliver to maildir")?; store.increment_jmap_state(&user.email).await; // Enqueue the LLM job. The worker reads the message back from the // maildir by walking the user's thread chain. let job = store .create_job(thread_id, &message_id) .await?; tracing::info!( thread_id, job_id = job.id, message_id, provider = %provider_name, model = %model_id, path = %raw_path.display(), "queued LLM reply job" ); Ok(()) } /// Parse `model-id@provider-id.example.com` where the daemon domain is /// `example.com`. Returns (provider, model_id). /// /// The local part maps `+` to `:` so `model+tag` becomes `model:tag`, /// matching the OpenAI-compatible `model:tag` syntax (e.g. an Ollama quant): /// `glm-4.7-flash+q8_0@ollama.local.llm` -> `glm-4.7-flash:q8_0`. fn parse_recipient(addr: &str, domain: &str) -> Option<(String, String)> { let parts: Vec<&str> = addr.split('@').collect(); if parts.len() != 2 { return None; } let local = parts[0]; let host = parts[1]; let suffix = format!(".{}", domain); if !host.eq_ignore_ascii_case(domain) && !host.ends_with(&suffix) { return None; } let provider = if host.eq_ignore_ascii_case(domain) { // model@domain -> no provider specified; reject for now. return None; } else { host.trim_end_matches(&suffix).to_string() }; if provider.is_empty() || local.is_empty() { return None; } Some((provider, local.replace('+', ":"))) } #[cfg(test)] mod tests { use super::*; use base64::Engine; use tokio::io::BufReader; /// Drive `read_data` against an in-memory byte buffer (CRLF line norms /// preserved), the same path the TCP reader takes. async fn read_data_bytes(input: &[u8], max_bytes: usize) -> anyhow::Result> { let mut reader = BufReader::new(input); read_data(&mut reader, max_bytes).await } #[test] fn extract_addr_bracketed_with_esmtp_params() { // ESMTP params after the address must not glue onto it. assert_eq!( extract_addr("MAIL FROM: BODY=8BITMIME SIZE=1234"), "a@b.com" ); assert_eq!(extract_addr("RCPT TO: NOTIFY=NEVER"), "x@y.com"); } #[test] fn extract_addr_unbracketed_takes_first_token() { assert_eq!(extract_addr("MAIL FROM: a@b.com"), "a@b.com"); assert_eq!(extract_addr("MAIL FROM:a@b.com BODY=8BITMIME"), "a@b.com"); } #[test] fn extract_addr_empty() { assert_eq!(extract_addr("MAIL FROM:<>"), ""); assert_eq!(extract_addr("MAIL FROM: "), ""); } #[test] fn decode_auth_plain_inline() { // RFC 4616: base64(`authzid \0 authcid \0 passwd`). authzid may be // empty; we ignore it (single-user model). Two NUL separators. let parts: [&[u8]; 3] = [b"", b"me@example.com", b"hunter2"]; let plain = parts.join(&b'\0'); let blob = base64::engine::general_purpose::STANDARD.encode(&plain); let (u, p) = decode_auth_plain(&blob).expect("valid PLAIN"); assert_eq!(u, "me@example.com"); assert_eq!(p, "hunter2"); } #[test] fn decode_auth_plain_rejects_missing_fields() { // No password -> not a complete PLAIN response. let blob = base64::engine::general_purpose::STANDARD.encode(b"\0user\0"); assert!(decode_auth_plain(&blob).is_none()); } #[test] fn decode_auth_plain_rejects_bad_base64() { assert!(decode_auth_plain("!!!not base64!!!").is_none()); } #[test] fn extract_addr_unwraps_display_name_in_envelope() { // A client that wraps a display name in the envelope as // `>` must still yield the inner addr-spec, not // `Display Name `). The bare `` form is // unchanged since it has a single `<`. assert_eq!( extract_addr("RCPT TO:>"), "gemma4-31b@ollama.clee.wtf" ); assert_eq!( extract_addr("RCPT TO:> NOTIFY=NEVER"), "gemma4-31b@ollama.clee.wtf" ); // Unbracketed display name + bracketed addr also resolves to the addr. assert_eq!( extract_addr("RCPT TO:gemma4-31b "), "gemma4-31b@ollama.clee.wtf" ); } #[tokio::test] async fn read_data_basic_crlf_message() { let input = b"From: a@b\r\nTo: c@d\r\n\r\nbody\r\n.\r\n"; let out = read_data_bytes(input, 1024).await.unwrap(); // Every line normalized to CRLF; terminator consumed; dot stripped. assert_eq!(out, b"From: a@b\r\nTo: c@d\r\n\r\nbody\r\n"); } #[tokio::test] async fn read_data_accepts_8bit_latin1_body() { // Raw 0xE9 (é in latin-1) and other high bytes — would have killed the // old String/read_line path with a UTF-8 error mid-DATA. let input = b"Subject: caf\xe9\r\n\r\nbody with \xff byte\r\n.\r\n"; let out = read_data_bytes(input, 1024).await.unwrap(); assert_eq!(out, b"Subject: caf\xe9\r\n\r\nbody with \xff byte\r\n"); } #[tokio::test] async fn read_data_normalizes_bare_lf_to_crlf() { // Bare-LF message: must be stored canonical CRLF or BODY[TEXT] comes // back empty on the IMAP side. let input = b"From: a@b\n\nbody line 1\nbody line 2\n.\n"; let out = read_data_bytes(input, 1024).await.unwrap(); assert_eq!(out, b"From: a@b\r\n\r\nbody line 1\r\nbody line 2\r\n"); } #[tokio::test] async fn read_data_dot_unstuffs_doubled_leading_dot() { // A body line of ".." is an escaped "."; a line of "..foo" -> ".foo". let input = b"..\r\n..foo\r\n.\r\n"; let out = read_data_bytes(input, 1024).await.unwrap(); assert_eq!(out, b".\r\n.foo\r\n"); } #[tokio::test] async fn read_data_accepts_message_without_final_newline() { // No trailing CRLF before the terminator — still terminates cleanly. let input = b"From: a@b\r\n\r\nno final newline\r\n.\r\n"; let out = read_data_bytes(input, 1024).await.unwrap(); assert_eq!(out, b"From: a@b\r\n\r\nno final newline\r\n"); } #[tokio::test] async fn read_data_rejects_oversized_message() { // 200 bytes of body but only 50 allowed. let body = "x".repeat(200); let input = format!("From: a@b\r\n\r\n{body}\r\n.\r\n"); let err = read_data_bytes(input.as_bytes(), 50).await.unwrap_err(); assert!(err.to_string().contains("exceeds max size")); } #[tokio::test] async fn read_data_rejects_unbounded_line_no_newline() { // A 50KB line with no newline must hit the per-line budget cap, not // buffer forever. max_bytes well below the line length. let line = "x".repeat(50_000); let input = format!("{line}"); // no newline at all let err = read_data_bytes(input.as_bytes(), 1024).await.unwrap_err(); assert!(err.to_string().contains("exceeds max size")); } #[test] fn unstuff_line_strips_one_leading_dot() { assert_eq!(unstuff_line(b".."), b"."); assert_eq!(unstuff_line(b"..foo"), b".foo"); assert_eq!(unstuff_line(b"..."), b".."); assert_eq!(unstuff_line(b"normal"), b"normal"); assert_eq!(unstuff_line(b"."), b"."); // not doubled; caller handles terminator } #[test] fn permanent_error_downcasts_at_reply_site() { // The reply arm classifies via downcast_ref::(); a Permanent // must downcast successfully where a plain anyhow does not. let perm: anyhow::Error = Permanent::new("no valid model recipient").into(); assert!(perm.downcast_ref::().is_some()); let transient = anyhow::anyhow!("connection refused"); assert!(transient.downcast_ref::().is_none()); } #[test] fn permanent_renders_its_message() { let e: anyhow::Error = Permanent::new("unknown provider: ollama").into(); // 550 reply uses {e} (Display): the message must come through clean. assert_eq!(e.to_string(), "unknown provider: ollama"); } #[test] fn parse_recipient_plus_translates_to_tag() { // `+` in the local part becomes `:` for the OpenAI `model:tag` syntax. let (p, m) = parse_recipient("glm-4.7-flash+q8_0@ollama.local.llm", "local.llm").unwrap(); assert_eq!(p, "ollama"); assert_eq!(m, "glm-4.7-flash:q8_0"); } #[test] fn parse_recipient_translates_every_plus_for_multi_segment_models() { // Real providers use multi-segment names like `hf:zai-org/GLM-5.2`; // every `+` must invert to `:` so the address round-trips. This is // the symmetric pair to model_to_local's replace-all. let (p, m) = parse_recipient("hf+zai-org/GLM-5.2@ollama.local.llm", "local.llm").unwrap(); assert_eq!(p, "ollama"); assert_eq!(m, "hf:zai-org/GLM-5.2"); } #[test] fn parse_recipient_without_plus_is_unchanged() { let (p, m) = parse_recipient("llama@llama-cpp.local.llm", "local.llm").unwrap(); assert_eq!(p, "llama-cpp"); assert_eq!(m, "llama"); } #[test] fn parse_recipient_rejects_bare_domain() { // model@domain with no provider subdomain is not routable. assert!(parse_recipient("llama@local.llm", "local.llm").is_none()); } }