Skip to main content

qualia_client_core/
mail_inbound.rs

1//! Inbound mail acceptance — the path that makes domain mail **work**.
2//!
3//! Flow:
4//! 1. Envelope `RCPT TO` / inject target → [`crate::domains::resolve_delivery`]
5//! 2. Body + headers → [`crate::mail_rules::evaluate`]
6//! 3. If deliver → [`crate::mail_store::store_delivery`]
7//!
8//! Also runs a **local SMTP receiver** (default `127.0.0.1:2525`) so mail can land without a paid
9//! mailbox product. Public internet delivery still needs DNS MX + a tunnel/VPS to your host
10//! (port 25 is often blocked on residential links) — but the product inbox is yours, local, and
11//! rule-bearing.
12
13use std::io::{BufRead, BufReader, Write};
14use std::net::{TcpListener, TcpStream};
15use std::sync::atomic::{AtomicBool, AtomicU16, Ordering};
16use std::sync::{Mutex, OnceLock};
17use std::thread;
18use std::time::Duration;
19
20use serde::{Deserialize, Serialize};
21
22use crate::domains::{self, DeliveryResolution, ResolutionVia};
23use crate::mail_rules::{self, InboundMessage};
24use crate::mail_store::{self, StoredMail};
25
26/// Default bind for the personal SMTP edge (submission-style port; avoid privileged 25).
27pub const DEFAULT_SMTP_BIND: &str = "127.0.0.1:2525";
28
29static SMTP_RUNNING: AtomicBool = AtomicBool::new(false);
30static SMTP_PORT: AtomicU16 = AtomicU16::new(0);
31static SMTP_STOP: AtomicBool = AtomicBool::new(false);
32static SMTP_BIND_DISPLAY: OnceLock<Mutex<String>> = OnceLock::new();
33
34fn bind_display() -> &'static Mutex<String> {
35    SMTP_BIND_DISPLAY.get_or_init(|| Mutex::new(String::new()))
36}
37
38/// Outcome of accepting one message (for API / UI).
39#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct AcceptResult {
41    pub accepted: bool,
42    pub rejected: Option<String>,
43    pub stored: Option<StoredMail>,
44}
45
46/// Accept a message for a registered domain mailbox — pure product path (no network).
47pub fn accept_message(
48    from: &str,
49    to: &str,
50    subject: &str,
51    body: &str,
52    sender_verified: bool,
53    sender_did: Option<String>,
54) -> AcceptResult {
55    let addresses = domains::list_addresses(None);
56    match domains::resolve_delivery(&addresses, to) {
57        DeliveryResolution::Reject { reason } => AcceptResult {
58            accepted: false,
59            rejected: Some(reason),
60            stored: None,
61        },
62        DeliveryResolution::Deliver { address, via } => {
63            let msg = InboundMessage {
64                from_address: from.to_string(),
65                to_address: to.to_string(),
66                sender_did,
67                sender_verified,
68                subject: subject.to_string(),
69                size_bytes: body.len(),
70            };
71            let mut verdict = mail_rules::evaluate(&address.rules, &msg);
72            let via_str = match via {
73                ResolutionVia::Exact => "exact",
74                ResolutionVia::Catchall => "catchall",
75                ResolutionVia::Unsolicited => "unsolicited",
76            };
77            if matches!(via, ResolutionVia::Catchall) && !verdict.quarantined {
78                verdict.quarantined = true;
79                verdict.reasons.push("catchall intake — quarantined".into());
80            }
81            if !verdict.deliver {
82                return AcceptResult {
83                    accepted: false,
84                    rejected: verdict
85                        .rejected
86                        .or_else(|| Some("rejected by rules".into())),
87                    stored: None,
88                };
89            }
90            match mail_store::store_delivery(
91                from,
92                to,
93                &address.address,
94                subject,
95                body,
96                via_str,
97                verdict.quarantined,
98                verdict.priority,
99                verdict.reasons,
100            ) {
101                Ok(stored) => AcceptResult {
102                    accepted: true,
103                    rejected: None,
104                    stored: Some(stored),
105                },
106                Err(e) => AcceptResult {
107                    accepted: false,
108                    rejected: Some(e),
109                    stored: None,
110                },
111            }
112        }
113    }
114}
115
116/// DNS records a human pastes so the **internet** can find this receiver.
117/// Host is optional public hostname (MX target); when empty, placeholders are used.
118pub fn mail_dns_forms(domain: &str, mx_host: Option<&str>) -> serde_json::Value {
119    let domain = domain.trim().to_lowercase();
120    let mx = mx_host
121        .map(|h| h.trim().to_lowercase())
122        .filter(|h| !h.is_empty())
123        .unwrap_or_else(|| format!("mail.{domain}"));
124    let port = SMTP_PORT.load(Ordering::Relaxed);
125    let port_note = if port == 0 {
126        DEFAULT_SMTP_BIND.to_string()
127    } else {
128        format!("listening on port {port}")
129    };
130    serde_json::json!({
131        "domain": domain,
132        "mx_host": mx,
133        "records": [
134            {
135                "type": "MX",
136                "name": "@",
137                "priority": 10,
138                "value": format!("{mx}."),
139                "hint": "Points senders at your mail host (tunnel, VPS, or home public hostname)."
140            },
141            {
142                "type": "A_or_AAAA",
143                "name": mx.strip_suffix(&format!(".{domain}")).unwrap_or("mail"),
144                "value": "<your public IP or tunnel target>",
145                "hint": "Where the MX name resolves. Cloudflare Tunnel can map hostname → this app."
146            },
147            {
148                "type": "TXT",
149                "name": "@",
150                "value": "v=spf1 mx a -all",
151                "hint": "SPF: only MX/A may send as this domain (tighten later with DKIM)."
152            },
153        ],
154        "local_receiver": {
155            "default_bind": DEFAULT_SMTP_BIND,
156            "status": if SMTP_RUNNING.load(Ordering::Relaxed) { "running" } else { "stopped" },
157            "bind": bind_display().lock().map(|g| g.clone()).unwrap_or_default(),
158            "port_note": port_note,
159            "how": "Start the local SMTP receiver in Talk → Mail. For public mail, forward TCP 25 (or 2525) to that port via tunnel/router. Same machine can use localhost:2525 without DNS."
160        },
161        "plaintext_block": format!(
162            "MX  @  10  {mx}.\n\
163             A/AAAA  {mx}  →  <public IP or tunnel>\n\
164             TXT  @  v=spf1 mx a -all\n\
165             # Local receiver: {port_note}\n\
166             # QDP front-door TXT is separate (_qdp) — identity, not MX."
167        ),
168    })
169}
170
171/// Status of the local SMTP receiver.
172pub fn receiver_status() -> serde_json::Value {
173    serde_json::json!({
174        "running": SMTP_RUNNING.load(Ordering::Relaxed),
175        "port": SMTP_PORT.load(Ordering::Relaxed),
176        "bind": bind_display().lock().map(|g| g.clone()).unwrap_or_default(),
177        "default_bind": DEFAULT_SMTP_BIND,
178    })
179}
180
181/// Start the local SMTP receiver on `bind` (e.g. `127.0.0.1:2525` or `0.0.0.0:2525`).
182/// Idempotent if already running on the same bind.
183#[cfg(not(target_arch = "wasm32"))]
184pub fn start_receiver(bind: &str) -> Result<serde_json::Value, String> {
185    let bind = if bind.trim().is_empty() {
186        DEFAULT_SMTP_BIND.to_string()
187    } else {
188        bind.trim().to_string()
189    };
190    if SMTP_RUNNING.load(Ordering::Relaxed) {
191        return Ok(serde_json::json!({
192            "already_running": true,
193            "bind": bind_display().lock().map(|g| g.clone()).unwrap_or_default(),
194            "port": SMTP_PORT.load(Ordering::Relaxed),
195        }));
196    }
197    let listener = TcpListener::bind(&bind)
198        .map_err(|e| format!("bind {bind} failed: {e} (is another process using the port?)"))?;
199    listener
200        .set_nonblocking(true)
201        .map_err(|e| format!("set_nonblocking: {e}"))?;
202    let port = listener.local_addr().map(|a| a.port()).unwrap_or(0);
203    SMTP_STOP.store(false, Ordering::Relaxed);
204    SMTP_PORT.store(port, Ordering::Relaxed);
205    if let Ok(mut g) = bind_display().lock() {
206        *g = bind.clone();
207    }
208    SMTP_RUNNING.store(true, Ordering::Relaxed);
209
210    thread::Builder::new()
211        .name("webizen-smtp".into())
212        .spawn(move || {
213            while !SMTP_STOP.load(Ordering::Relaxed) {
214                match listener.accept() {
215                    Ok((stream, _peer)) => {
216                        let _ = stream.set_read_timeout(Some(Duration::from_secs(60)));
217                        let _ = stream.set_write_timeout(Some(Duration::from_secs(60)));
218                        // One connection per thread — personal MTA volume is small.
219                        thread::spawn(move || {
220                            if let Err(e) = handle_smtp_session(stream) {
221                                log::debug!("smtp session end: {e}");
222                            }
223                        });
224                    }
225                    Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
226                        thread::sleep(Duration::from_millis(50));
227                    }
228                    Err(_) => {
229                        thread::sleep(Duration::from_millis(100));
230                    }
231                }
232            }
233            SMTP_RUNNING.store(false, Ordering::Relaxed);
234            SMTP_PORT.store(0, Ordering::Relaxed);
235        })
236        .map_err(|e| format!("spawn smtp thread: {e}"))?;
237
238    Ok(serde_json::json!({
239        "started": true,
240        "bind": bind,
241        "port": port,
242        "message": format!("Local SMTP receiver on {bind} — point MX/tunnel here, or send to localhost for tests."),
243    }))
244}
245
246#[cfg(not(target_arch = "wasm32"))]
247pub fn stop_receiver() -> Result<serde_json::Value, String> {
248    SMTP_STOP.store(true, Ordering::Relaxed);
249    // Give the accept loop a moment to notice.
250    thread::sleep(Duration::from_millis(80));
251    Ok(serde_json::json!({
252        "stopped": true,
253        "running": SMTP_RUNNING.load(Ordering::Relaxed),
254    }))
255}
256
257#[cfg(not(target_arch = "wasm32"))]
258fn smtp_write(stream: &mut TcpStream, line: &str) -> Result<(), String> {
259    stream
260        .write_all(line.as_bytes())
261        .map_err(|e| e.to_string())?;
262    stream.flush().map_err(|e| e.to_string())
263}
264
265/// Extract address from `MAIL FROM:<a@b>` / `RCPT TO:<a@b>` / bare forms.
266fn parse_smtp_path(arg: &str) -> String {
267    let s = arg.trim();
268    if let Some(start) = s.find('<') {
269        if let Some(end) = s[start + 1..].find('>') {
270            return s[start + 1..start + 1 + end].trim().to_string();
271        }
272    }
273    s.trim_matches(|c| c == '<' || c == '>').to_string()
274}
275
276#[cfg(not(target_arch = "wasm32"))]
277fn handle_smtp_session(mut stream: TcpStream) -> Result<(), String> {
278    // Blocking reads after accept — switch socket back to blocking for BufReader.
279    stream.set_nonblocking(false).map_err(|e| e.to_string())?;
280    smtp_write(
281        &mut stream,
282        "220 webizen.local ESMTP Qualia semantic mail ready\r\n",
283    )?;
284
285    let mut reader = BufReader::new(stream.try_clone().map_err(|e| e.to_string())?);
286    let mut mail_from = String::new();
287    let mut rcpt_to: Vec<String> = Vec::new();
288    let mut line = String::new();
289
290    loop {
291        line.clear();
292        let n = reader.read_line(&mut line).map_err(|e| e.to_string())?;
293        if n == 0 {
294            break;
295        }
296        let cmd = line.trim_end_matches(['\r', '\n']);
297        let upper = cmd.to_ascii_uppercase();
298
299        if upper.starts_with("EHLO") || upper.starts_with("HELO") {
300            smtp_write(
301                &mut stream,
302                "250-webizen.local\r\n250-PIPELINING\r\n250 8BITMIME\r\n",
303            )?;
304        } else if upper.starts_with("MAIL FROM:") {
305            mail_from = parse_smtp_path(&cmd[10..]);
306            smtp_write(&mut stream, "250 OK\r\n")?;
307        } else if upper.starts_with("RCPT TO:") {
308            let to = parse_smtp_path(&cmd[8..]);
309            let addresses = domains::list_addresses(None);
310            match domains::resolve_delivery(&addresses, &to) {
311                DeliveryResolution::Deliver { .. } => {
312                    rcpt_to.push(to);
313                    smtp_write(&mut stream, "250 OK\r\n")?;
314                }
315                DeliveryResolution::Reject { reason } => {
316                    smtp_write(&mut stream, &format!("550 5.1.1 {reason}\r\n"))?;
317                }
318            }
319        } else if upper == "DATA" {
320            if rcpt_to.is_empty() {
321                smtp_write(&mut stream, "503 5.5.1 No valid recipients\r\n")?;
322                continue;
323            }
324            smtp_write(&mut stream, "354 End data with <CR><LF>.<CR><LF>\r\n")?;
325            let mut data = String::new();
326            loop {
327                line.clear();
328                let n = reader.read_line(&mut line).map_err(|e| e.to_string())?;
329                if n == 0 {
330                    break;
331                }
332                // End of data: line is ".\r\n" or "."
333                let trimmed = line.trim_end_matches(['\r', '\n']);
334                if trimmed == "." {
335                    break;
336                }
337                // RFC 5321 dot-stuffing: lines starting with ".." → strip one dot.
338                if let Some(rest) = line.strip_prefix("..") {
339                    data.push_str(rest);
340                } else {
341                    data.push_str(&line);
342                }
343            }
344            let (subject, body) = split_headers_body(&data);
345            let from = if mail_from.is_empty() {
346                "unknown@invalid".to_string()
347            } else {
348                mail_from.clone()
349            };
350            let mut any_ok = false;
351            let mut last_err = String::new();
352            for to in &rcpt_to {
353                let r = accept_message(&from, to, &subject, &body, false, None);
354                if r.accepted {
355                    any_ok = true;
356                } else if let Some(e) = r.rejected {
357                    last_err = e;
358                }
359            }
360            if any_ok {
361                smtp_write(&mut stream, "250 OK queued\r\n")?;
362            } else {
363                smtp_write(&mut stream, &format!("550 5.7.1 rejected: {last_err}\r\n"))?;
364            }
365            // Reset transaction
366            mail_from.clear();
367            rcpt_to.clear();
368        } else if upper == "RSET" {
369            mail_from.clear();
370            rcpt_to.clear();
371            smtp_write(&mut stream, "250 OK\r\n")?;
372        } else if upper == "NOOP" {
373            smtp_write(&mut stream, "250 OK\r\n")?;
374        } else if upper == "QUIT" {
375            smtp_write(&mut stream, "221 webizen.local closing\r\n")?;
376            break;
377        } else if upper.starts_with("VRFY") || upper.starts_with("EXPN") {
378            smtp_write(&mut stream, "252 Not verified\r\n")?;
379        } else {
380            smtp_write(&mut stream, "502 5.5.2 Command not implemented\r\n")?;
381        }
382    }
383    Ok(())
384}
385
386/// Split RFC822-ish message into subject + body (body includes remaining headers stripped of Subject).
387fn split_headers_body(data: &str) -> (String, String) {
388    let normalized = data.replace("\r\n", "\n");
389    let (head, body) = if let Some(i) = normalized.find("\n\n") {
390        (&normalized[..i], normalized[i + 2..].to_string())
391    } else {
392        return (String::new(), normalized);
393    };
394    let mut subject = String::new();
395    for line in head.lines() {
396        let lower = line.to_ascii_lowercase();
397        if let Some(rest) = lower.strip_prefix("subject:") {
398            // Preserve original casing of value from `line`
399            subject = line[line.len() - rest.len()..].trim().to_string();
400            break;
401        }
402    }
403    // Prefer full original as body for reading; subject extracted for list UI.
404    let full_body = if body.is_empty() {
405        data.to_string()
406    } else {
407        body
408    };
409    (subject, full_body)
410}
411
412#[cfg(test)]
413mod tests {
414    use super::*;
415    use crate::domains::{
416        make_domain, make_purpose_address, upsert_address, upsert_domain, AgentType, DomainOwner,
417        MailRules,
418    };
419
420    #[test]
421    fn parse_smtp_path_angles() {
422        assert_eq!(parse_smtp_path("<bob@alice.example>"), "bob@alice.example");
423        assert_eq!(parse_smtp_path(" bob@alice.example "), "bob@alice.example");
424    }
425
426    #[test]
427    fn split_subject_from_data() {
428        let (s, b) = split_headers_body("From: a@b\r\nSubject: Hello\r\n\r\nBody line\r\n");
429        assert_eq!(s, "Hello");
430        assert!(b.contains("Body line"));
431    }
432
433    #[test]
434    fn accept_message_fail_closed_without_mailbox() {
435        let r = accept_message(
436            "x@y.example",
437            "nobody@not-registered.invalid.example",
438            "hi",
439            "body",
440            false,
441            None,
442        );
443        assert!(!r.accepted);
444        assert!(r.rejected.is_some());
445    }
446
447    #[test]
448    fn accept_message_delivers_when_domain_onboarded() {
449        let domain = format!("inbound-test-{}.example", std::process::id());
450        let d = make_domain(
451            &domain,
452            AgentType::NaturalPerson,
453            DomainOwner::Personal {
454                did: "did:test".into(),
455            },
456            "did:test",
457            "T",
458            None,
459            1,
460        )
461        .unwrap();
462        upsert_domain(d).unwrap();
463        let a = make_purpose_address(
464            &domain,
465            "frontdoor",
466            MailRules {
467                notify: true,
468                ..Default::default()
469            },
470            1,
471        )
472        .unwrap();
473        upsert_address(a).unwrap();
474        let to = format!("frontdoor@{domain}");
475        let r = accept_message(
476            "peer@other.example",
477            &to,
478            "Ping",
479            "hello world",
480            false,
481            None,
482        );
483        assert!(r.accepted, "{:?}", r.rejected);
484        let stored = r.stored.expect("stored");
485        assert_eq!(stored.mailbox, to);
486        assert!(!stored.quarantined);
487        // cleanup
488        let _ = mail_store::delete(&stored.id);
489    }
490
491    #[test]
492    fn mail_dns_forms_include_mx_and_spf() {
493        let v = mail_dns_forms("alice.example", Some("mail.alice.example"));
494        let block = v["plaintext_block"].as_str().unwrap_or("");
495        assert!(block.contains("MX"));
496        assert!(block.contains("spf1"));
497    }
498}