1use 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
26pub 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#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct AcceptResult {
41 pub accepted: bool,
42 pub rejected: Option<String>,
43 pub stored: Option<StoredMail>,
44}
45
46pub 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
116pub 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
171pub 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#[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 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 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
265fn 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 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 let trimmed = line.trim_end_matches(['\r', '\n']);
334 if trimmed == "." {
335 break;
336 }
337 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 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
386fn 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 subject = line[line.len() - rest.len()..].trim().to_string();
400 break;
401 }
402 }
403 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 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}