Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_net/src/smtp/server.rs

51.0 KiB, 125 runs

created by r1870400018:9882, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Server-side SMTP session state machine.
2//!
3//! Implements RFC 5321 (SMTP), RFC 3207 (STARTTLS), RFC 4954 (AUTH PLAIN
4//! and LOGIN) at the level needed to host a single Hematite mailbox: the
5//! receive path (port 25, no auth, accepts mail destined to local users)
6//! and the submission path (port 587, AUTH required after STARTTLS, may
7//! relay anywhere).
8//!
9//! # Submission binds the sender to the account
10//!
11//! Until 2026-09-23 an authenticated account could send as any address: the
12//! envelope and the header were never compared with the login, and the relay
13//! DKIM-signs what it forwards, so one leaked password sent validly signed mail
14//! as anybody in the relay's domain. Now `MAIL FROM` and every address in the
15//! header `From` (and `Sender`) must be the account's own or one its entry lists
16//! under `send_as` (see [`MailUser::may_send_as`]), or the command is refused
17//! with `550 5.7.1`. The header is bound as well as the envelope because DMARC,
18//! and the reader, go by the header. A null `MAIL FROM:<>` is allowed, since it
19//! names nobody and the header is bound all the same.
20//!
21//! The session loop runs over an enum-based `MaybeTls` stream so that a
22//! plain TCP connection can be transparently swapped to TLS in response
23//! to `STARTTLS` without duplicating the rest of the state machine.
24//!
25//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
26//! Anthropic Claude
27
28use crate::{
29 email::header::{
30 addresses,
31 header_fields,
32 split_headers_body,
33 },
34 smtp::{
35 cmd::SmtpCommand,
36 codes::SmtpResponseCode,
37 handler::{
38 HandlerOutcome,
39 SmtpHandler,
40 SmtpTransaction,
41 },
42 },
43 mail::{
44 store::MailUser,
45 user::UserStore,
46 },
47};
48
49use oxedyne_fe2o3_core::prelude::*;
50
51use std::{
52 net::SocketAddr,
53 pin::Pin,
54 sync::Arc,
55 task::{
56 Context,
57 Poll,
58 },
59};
60
61use base64;
62use tokio::{
63 io::{
64 AsyncRead,
65 AsyncReadExt,
66 AsyncWrite,
67 AsyncWriteExt,
68 ReadBuf,
69 },
70 net::TcpStream,
71};
72use tokio_rustls::server::TlsStream;
73
74use crate::tls::{
75 BoundedTlsAcceptor,
76 Handshake,
77};
78
79
80//// Session ceilings.
81// The size limit mirrors Postfix's `message_size_limit = 20480000` default. A wire line must not
82// exceed 1000 bytes including CRLF (RFC 5321 §4.5.3.1.6), and 512 for a command line, so the read
83// budget carries slack over both.
84pub const SMTP_MESSAGE_SIZE_LIMIT: usize = 20_480_000; // bytes, inside `DATA`
85pub const SMTP_MAX_RCPT: usize = 100; // recipients per transaction
86pub const SMTP_MAX_LINE: usize = 4_096; // bytes per line read
87
88
89// ┌───────────────────────────────────────────────────────────────────────────┐
90// │ MAYBE TLS │
91// │ │
92// │ Enum-based stream that is either a plain TcpStream or a TlsStream<TcpStream│
93// │ after a successful STARTTLS upgrade. Implements AsyncRead/AsyncWrite by │
94// │ delegating, with explicit pin projection on each variant. │
95// └───────────────────────────────────────────────────────────────────────────┘
96
97/// Either a plain TCP stream or a TLS-wrapped TCP stream.
98///
99/// Used so the SMTP and IMAP session loops can be written once over a
100/// single `S: AsyncRead + AsyncWrite` and dynamically swap in a TLS
101/// upgrade in response to STARTTLS.
102pub enum MaybeTls {
103 Plain(TcpStream), // before STARTTLS, or where TLS is not negotiated
104 Tls(Box<TlsStream<TcpStream>>), // after STARTTLS, or on an implicit-TLS port
105}
106
107impl MaybeTls {
108 pub fn is_tls(&self) -> bool {
109 matches!(self, MaybeTls::Tls(_))
110 }
111
112 /// `None` once the connection has been upgraded, which is what the STARTTLS path wants.
113 pub fn into_plain(self) -> Option<TcpStream> {
114 match self {
115 MaybeTls::Plain(s) => Some(s),
116 MaybeTls::Tls(_) => None,
117 }
118 }
119}
120
121impl AsyncRead for MaybeTls {
122 fn poll_read(
123 self: Pin<&mut Self>,
124 cx: &mut Context<'_>,
125 buf: &mut ReadBuf<'_>,
126 )
127 -> Poll<std::io::Result<()>>
128 {
129 // Safe pin projection by hand: we never move the inner stream
130 // ourselves, and Pin::new is sound for both variants because
131 // TcpStream and TlsStream are Unpin.
132 match self.get_mut() {
133 MaybeTls::Plain(s) => Pin::new(s).poll_read(cx, buf),
134 MaybeTls::Tls(s) => Pin::new(s.as_mut()).poll_read(cx, buf),
135 }
136 }
137}
138
139impl AsyncWrite for MaybeTls {
140 fn poll_write(
141 self: Pin<&mut Self>,
142 cx: &mut Context<'_>,
143 buf: &[u8],
144 )
145 -> Poll<std::io::Result<usize>>
146 {
147 match self.get_mut() {
148 MaybeTls::Plain(s) => Pin::new(s).poll_write(cx, buf),
149 MaybeTls::Tls(s) => Pin::new(s.as_mut()).poll_write(cx, buf),
150 }
151 }
152
153 fn poll_flush(
154 self: Pin<&mut Self>,
155 cx: &mut Context<'_>,
156 )
157 -> Poll<std::io::Result<()>>
158 {
159 match self.get_mut() {
160 MaybeTls::Plain(s) => Pin::new(s).poll_flush(cx),
161 MaybeTls::Tls(s) => Pin::new(s.as_mut()).poll_flush(cx),
162 }
163 }
164
165 fn poll_shutdown(
166 self: Pin<&mut Self>,
167 cx: &mut Context<'_>,
168 )
169 -> Poll<std::io::Result<()>>
170 {
171 match self.get_mut() {
172 MaybeTls::Plain(s) => Pin::new(s).poll_shutdown(cx),
173 MaybeTls::Tls(s) => Pin::new(s.as_mut()).poll_shutdown(cx),
174 }
175 }
176}
177
178
179// ┌───────────────────────────────────────────────────────────────────────────┐
180// │ SMTP SERVER MODE │
181// └───────────────────────────────────────────────────────────────────────────┘
182
183/// Which port and policy this listener is serving.
184#[derive(Clone, Copy, Debug, Eq, PartialEq)]
185pub enum SmtpMode {
186 Receive, // MX inbound on 25: no `AUTH`, recipients must resolve locally
187 Submission, // MSA on 587: STARTTLS then `AUTH`, and relay anywhere
188}
189
190
191// ┌───────────────────────────────────────────────────────────────────────────┐
192// │ SESSION STATE │
193// └───────────────────────────────────────────────────────────────────────────┘
194
195/// Where in the SMTP transaction we currently are; each variant names what has arrived, and so
196/// what may arrive next.
197#[derive(Clone, Debug, Eq, PartialEq)]
198enum SmtpPhase {
199 Greeted, // just connected, expecting HELO or EHLO
200 Ehlo, // ready for `MAIL FROM`, or another extension
201 MailFrom, // expecting `RCPT TO`
202 Rcpt, // one recipient at least, expecting more or `DATA`
203}
204
205/// Per-connection mutable state. Reset by `RSET` and after a successful
206/// `DATA` transaction.
207struct SmtpSession {
208 phase: SmtpPhase,
209 helo_domain: String,
210 mail_from: String,
211 rcpt_to: Vec<String>,
212 auth_user: Option<crate::mail::store::MailUser>,
213}
214
215impl SmtpSession {
216 fn new() -> Self {
217 Self {
218 phase: SmtpPhase::Greeted,
219 helo_domain: String::new(),
220 mail_from: String::new(),
221 rcpt_to: Vec::new(),
222 auth_user: None,
223 }
224 }
225
226 /// The envelope goes; the HELO domain and the authenticated identity stay.
227 fn reset_transaction(&mut self) {
228 if !self.helo_domain.is_empty() {
229 self.phase = SmtpPhase::Ehlo;
230 } else {
231 self.phase = SmtpPhase::Greeted;
232 }
233 self.mail_from.clear();
234 self.rcpt_to.clear();
235 }
236}
237
238
239// ┌───────────────────────────────────────────────────────────────────────────┐
240// │ SMTP SERVER │
241// └───────────────────────────────────────────────────────────────────────────┘
242
243/// One SMTP listener configuration.
244///
245/// Cheaply cloneable -- the inner state (handler, user store, TLS
246/// acceptor, hostname) is wrapped in `Arc`s by the caller so a single
247/// listener can fan out across every accept loop.
248#[derive(Clone)]
249pub struct SmtpServer<H: SmtpHandler, U: UserStore> {
250 pub handler: H,
251 pub users: U, // read by the AUTH path
252 // `None` disables STARTTLS; a submission listener must always have one. The
253 // acceptor is bounded, so a STARTTLS handshake shares the same concurrency
254 // cap and deadline as every other TLS listener on the host.
255 pub tls_acceptor: Option<BoundedTlsAcceptor>,
256 // Advertised in the 220 banner and the `Received:` header, so it should be the public MX
257 // hostname.
258 pub hostname: Arc<String>,
259 pub mode: SmtpMode,
260}
261
262impl<H: SmtpHandler, U: UserStore> SmtpServer<H, U> {
263
264 /// Runs to the end of the session, an in-session STARTTLS upgrade included.
265 pub async fn run(
266 &self,
267 plain: TcpStream,
268 peer: SocketAddr,
269 )
270 -> Outcome<()>
271 {
272 let mut stream = MaybeTls::Plain(plain);
273 let mut session = SmtpSession::new();
274
275 // 220 banner.
276 let banner = fmt!(
277 "{} ESMTP Hematite Steel ready",
278 self.hostname,
279 );
280 res!(write_response(&mut stream, SmtpResponseCode::ServiceReady, &banner).await);
281
282 // Inner loop. Returns when the client quits, errors out, or
283 // requests a STARTTLS upgrade. STARTTLS triggers a re-entry
284 // with a TLS-wrapped stream.
285 loop {
286 let next = res!(self.run_loop(&mut stream, &mut session, peer).await);
287 match next {
288 LoopExit::Quit => break,
289 LoopExit::StartTls => {
290 // Pull the plain stream back out and run the TLS
291 // handshake. The acceptor must be configured -- if
292 // it is not, run_loop refused the STARTTLS request.
293 let acceptor = match self.tls_acceptor.clone() {
294 Some(a) => a,
295 None => return Err(err!(
296 "STARTTLS requested but no TLS acceptor configured.";
297 Init, Missing)),
298 };
299 let plain = match stream.into_plain() {
300 Some(s) => s,
301 None => return Err(err!(
302 "STARTTLS requested on a stream that was already TLS.";
303 Invalid, Bug)),
304 };
305 let tls = match acceptor.accept(plain).await {
306 Handshake::Ok(t) => t,
307 Handshake::TimedOut => return Err(err!(
308 "STARTTLS handshake timed out for {:?}.", peer;
309 IO, Network, Init)),
310 Handshake::Failed(e) => return Err(err!(e,
311 "STARTTLS handshake failed for {:?}.", peer;
312 IO, Network, Init)),
313 };
314 stream = MaybeTls::Tls(Box::new(tls));
315 // RFC 3207 §4.2: a successful STARTTLS resets the
316 // session to immediately after the greeting.
317 session = SmtpSession::new();
318 // No new banner -- the greeting was sent before TLS.
319 }
320 }
321 }
322
323 // Best-effort shutdown.
324 let _ = stream.shutdown().await;
325 Ok(())
326 }
327
328 async fn run_loop(
329 &self,
330 stream: &mut MaybeTls,
331 session: &mut SmtpSession,
332 peer: SocketAddr,
333 )
334 -> Outcome<LoopExit>
335 {
336 loop {
337 let line = match read_line(stream).await {
338 Ok(Some(l)) => l,
339 Ok(None) => return Ok(LoopExit::Quit),
340 Err(e) => return Err(err!(e,
341 "Reading SMTP command line from {:?}.", peer;
342 IO, Network, Read)),
343 };
344
345 // Empty lines are tolerated.
346 if line.trim().is_empty() {
347 continue;
348 }
349
350 let cmd = match SmtpCommand::from_str(&line) {
351 Ok(c) => c,
352 Err(e) => {
353 warn!("SMTP {:?} {}: unparseable command {:?}: {}",
354 self.mode, peer, line, e);
355 res!(write_response(
356 stream,
357 SmtpResponseCode::CommandUnrecognized,
358 "Command not recognised",
359 ).await);
360 continue;
361 }
362 };
363
364 match cmd {
365 SmtpCommand::Helo(domain) => {
366 session.helo_domain = domain.as_str().to_string();
367 session.phase = SmtpPhase::Ehlo;
368 res!(write_response(
369 stream,
370 SmtpResponseCode::RequestedMailActionOkayCompleted,
371 &fmt!("{} Hello {}", self.hostname, domain.as_str()),
372 ).await);
373 }
374 SmtpCommand::Ehlo(domain) => {
375 session.helo_domain = domain.as_str().to_string();
376 session.phase = SmtpPhase::Ehlo;
377 res!(self.write_ehlo(stream, domain.as_str()).await);
378 }
379 SmtpCommand::StartTls => {
380 if stream.is_tls() {
381 res!(write_response(
382 stream,
383 SmtpResponseCode::BadSequenceOfCommands,
384 "STARTTLS not allowed inside TLS",
385 ).await);
386 continue;
387 }
388 if self.tls_acceptor.is_none() {
389 res!(write_response(
390 stream,
391 SmtpResponseCode::CommandNotImplemented,
392 "STARTTLS not available",
393 ).await);
394 continue;
395 }
396 res!(write_response(
397 stream,
398 SmtpResponseCode::ServiceReady,
399 "Ready to start TLS",
400 ).await);
401 return Ok(LoopExit::StartTls);
402 }
403 SmtpCommand::Auth(arg) => {
404 if self.mode == SmtpMode::Receive {
405 res!(write_response(
406 stream,
407 SmtpResponseCode::CommandNotImplemented,
408 "AUTH not available on this port",
409 ).await);
410 continue;
411 }
412 if !stream.is_tls() {
413 res!(write_response(
414 stream,
415 SmtpResponseCode::EncryptionRequiredForAuthentication,
416 "Must issue STARTTLS before AUTH",
417 ).await);
418 continue;
419 }
420 match self.handle_auth(stream, &arg).await {
421 Ok(Some(user)) => {
422 session.auth_user = Some(user);
423 res!(write_response(
424 stream,
425 SmtpResponseCode::AuthenticationSuccessful,
426 "Authentication successful",
427 ).await);
428 }
429 Ok(None) => {
430 res!(write_response(
431 stream,
432 SmtpResponseCode::AuthenticationCredentialsInvalid,
433 "Authentication credentials invalid",
434 ).await);
435 }
436 Err(e) => {
437 warn!("SMTP AUTH transport error from {:?}: {}", peer, e);
438 res!(write_response(
439 stream,
440 SmtpResponseCode::TransactionFailed,
441 "Authentication aborted",
442 ).await);
443 }
444 }
445 }
446 SmtpCommand::MailFrom(addr) => {
447 if self.mode == SmtpMode::Submission && session.auth_user.is_none() {
448 res!(write_response(
449 stream,
450 SmtpResponseCode::AuthenticationRequired,
451 "Authentication required",
452 ).await);
453 continue;
454 }
455 if session.phase != SmtpPhase::Ehlo {
456 res!(write_response(
457 stream,
458 SmtpResponseCode::BadSequenceOfCommands,
459 "Bad command sequence",
460 ).await);
461 continue;
462 }
463 if let (SmtpMode::Submission, Some(user)) = (self.mode, &session.auth_user) {
464 if let Some(why) = envelope_refusal(user, &addr) {
465 warn!("SMTP submission from {:?}: {} refused: {}",
466 peer, user.address(), why);
467 res!(write_response(
468 stream,
469 SmtpResponseCode::MailboxUnavailableOrAccessDenied,
470 &fmt!("5.7.1 {}", why),
471 ).await);
472 continue;
473 }
474 }
475 session.mail_from = addr;
476 session.phase = SmtpPhase::MailFrom;
477 res!(write_response(
478 stream,
479 SmtpResponseCode::RequestedMailActionOkayCompleted,
480 "Sender OK",
481 ).await);
482 }
483 SmtpCommand::RcptTo(addr) => {
484 if !matches!(session.phase, SmtpPhase::MailFrom | SmtpPhase::Rcpt) {
485 res!(write_response(
486 stream,
487 SmtpResponseCode::BadSequenceOfCommands,
488 "MAIL FROM required first",
489 ).await);
490 continue;
491 }
492 if session.rcpt_to.len() >= SMTP_MAX_RCPT {
493 res!(write_response(
494 stream,
495 SmtpResponseCode::ExceededStorageAllocation,
496 "Too many recipients",
497 ).await);
498 continue;
499 }
500 if self.mode == SmtpMode::Receive
501 && !self.handler.rcpt_acceptable(&addr)
502 {
503 res!(write_response(
504 stream,
505 SmtpResponseCode::MailboxUnavailableOrAccessDenied,
506 "No such user",
507 ).await);
508 continue;
509 }
510 session.rcpt_to.push(addr);
511 session.phase = SmtpPhase::Rcpt;
512 res!(write_response(
513 stream,
514 SmtpResponseCode::RequestedMailActionOkayCompleted,
515 "Recipient OK",
516 ).await);
517 }
518 SmtpCommand::Data => {
519 if session.phase != SmtpPhase::Rcpt {
520 res!(write_response(
521 stream,
522 SmtpResponseCode::BadSequenceOfCommands,
523 "Need MAIL FROM and RCPT TO first",
524 ).await);
525 continue;
526 }
527 res!(write_response(
528 stream,
529 SmtpResponseCode::StartMailInput,
530 "End data with <CR><LF>.<CR><LF>",
531 ).await);
532 let body = match read_data(stream).await {
533 Ok(b) => b,
534 Err(e) => {
535 return Err(err!(e,
536 "Reading SMTP DATA body from {:?}.", peer;
537 IO, Network, Read));
538 }
539 };
540 if let (SmtpMode::Submission, Some(user)) = (self.mode, &session.auth_user) {
541 if let Some(why) = header_refusal(user, &body) {
542 warn!("SMTP submission from {:?}: {} refused: {}",
543 peer, user.address(), why);
544 res!(write_response(
545 stream,
546 SmtpResponseCode::MailboxUnavailableOrAccessDenied,
547 &fmt!("5.7.1 {}", why),
548 ).await);
549 session.reset_transaction();
550 continue;
551 }
552 }
553 let txn = SmtpTransaction {
554 mail_from: session.mail_from.clone(),
555 rcpt_to: session.rcpt_to.clone(),
556 helo_domain: session.helo_domain.clone(),
557 auth_user: session.auth_user.clone(),
558 peer,
559 tls: stream.is_tls(),
560 raw_message: body,
561 };
562 let outcome = match self.mode {
563 SmtpMode::Receive => self.handler.deliver_inbound(txn),
564 SmtpMode::Submission => self.handler.submit_outbound(txn),
565 };
566 match outcome {
567 Ok(HandlerOutcome::Accepted(qid)) => {
568 res!(write_response(
569 stream,
570 SmtpResponseCode::RequestedMailActionOkayCompleted,
571 &fmt!("OK queued as {}", qid),
572 ).await);
573 }
574 Ok(HandlerOutcome::RejectPermanent(reason)) => {
575 res!(write_response(
576 stream,
577 SmtpResponseCode::TransactionFailed,
578 &reason,
579 ).await);
580 }
581 Ok(HandlerOutcome::RejectTemporary(reason)) => {
582 res!(write_response(
583 stream,
584 SmtpResponseCode::LocalErrorInProcessing,
585 &reason,
586 ).await);
587 }
588 Err(e) => {
589 error!(err!(e,
590 "SMTP handler error from {:?}.", peer;
591 IO));
592 res!(write_response(
593 stream,
594 SmtpResponseCode::LocalErrorInProcessing,
595 "Local error processing message",
596 ).await);
597 }
598 }
599 session.reset_transaction();
600 }
601 SmtpCommand::Rset => {
602 session.reset_transaction();
603 res!(write_response(
604 stream,
605 SmtpResponseCode::RequestedMailActionOkayCompleted,
606 "Reset OK",
607 ).await);
608 }
609 SmtpCommand::Noop => {
610 res!(write_response(
611 stream,
612 SmtpResponseCode::RequestedMailActionOkayCompleted,
613 "OK",
614 ).await);
615 }
616 SmtpCommand::Vrfy(_) => {
617 res!(write_response(
618 stream,
619 SmtpResponseCode::CannotVerifyUserButWillAttemptDelivery,
620 "Cannot VRFY user, but will attempt delivery",
621 ).await);
622 }
623 SmtpCommand::Expn(_) => {
624 res!(write_response(
625 stream,
626 SmtpResponseCode::CommandNotImplemented,
627 "EXPN not supported",
628 ).await);
629 }
630 SmtpCommand::Help(_) => {
631 res!(write_response(
632 stream,
633 SmtpResponseCode::HelpMessage,
634 "HELO EHLO MAIL RCPT DATA RSET NOOP QUIT STARTTLS AUTH",
635 ).await);
636 }
637 SmtpCommand::Quit => {
638 res!(write_response(
639 stream,
640 SmtpResponseCode::ServiceClosingTransmissionChannel,
641 &fmt!("{} closing connection", self.hostname),
642 ).await);
643 return Ok(LoopExit::Quit);
644 }
645 _ => {
646 res!(write_response(
647 stream,
648 SmtpResponseCode::CommandNotImplemented,
649 "Not implemented",
650 ).await);
651 }
652 }
653 }
654 }
655
656 async fn write_ehlo(
657 &self,
658 stream: &mut MaybeTls,
659 peer: &str,
660 )
661 -> Outcome<()>
662 {
663 // Build the multi-line capability list. Lines other than the
664 // last use `250-`, the last uses `250 `.
665 let mut lines: Vec<String> = Vec::new();
666 lines.push(fmt!("{} Hello {}", self.hostname, peer));
667 lines.push(fmt!("PIPELINING"));
668 lines.push(fmt!("8BITMIME"));
669 lines.push(fmt!("SIZE {}", SMTP_MESSAGE_SIZE_LIMIT));
670 lines.push(fmt!("ENHANCEDSTATUSCODES"));
671 if !stream.is_tls() && self.tls_acceptor.is_some() {
672 lines.push(fmt!("STARTTLS"));
673 }
674 if self.mode == SmtpMode::Submission && stream.is_tls() {
675 lines.push(fmt!("AUTH PLAIN LOGIN"));
676 }
677 lines.push(fmt!("HELP"));
678
679 let last = lines.len() - 1;
680 for (i, line) in lines.iter().enumerate() {
681 let sep = if i == last { ' ' } else { '-' };
682 let frame = fmt!("250{}{}\r\n", sep, line);
683 if let Err(e) = stream.write_all(frame.as_bytes()).await {
684 return Err(err!(e, "Writing EHLO line."; IO, Network, Write));
685 }
686 }
687 if let Err(e) = stream.flush().await {
688 return Err(err!(e, "Flushing EHLO."; IO, Network, Write));
689 }
690 Ok(())
691 }
692
693 async fn handle_auth(
694 &self,
695 stream: &mut MaybeTls,
696 arg: &str,
697 )
698 -> Outcome<Option<crate::mail::store::MailUser>>
699 {
700 let mut parts = arg.splitn(2, char::is_whitespace);
701 let mech = match parts.next() {
702 Some(m) => m.to_uppercase(),
703 None => return Ok(None),
704 };
705 let initial = parts.next().map(|s| s.trim().to_string());
706
707 match mech.as_str() {
708 "PLAIN" => {
709 // RFC 4616: base64( authzid \0 authcid \0 passwd ).
710 let payload = match initial {
711 Some(p) if !p.is_empty() => p,
712 _ => {
713 res!(write_response(
714 stream,
715 SmtpResponseCode::AuthInputData,
716 "",
717 ).await);
718 match res!(read_line(stream).await) {
719 Some(l) => l.trim().to_string(),
720 None => return Ok(None),
721 }
722 }
723 };
724 let raw = match base64::decode(payload.as_bytes()) {
725 Ok(b) => b,
726 Err(_) => return Ok(None),
727 };
728 let (user, pass) = match parse_plain(&raw) {
729 Some(p) => p,
730 None => return Ok(None),
731 };
732 self.users.authenticate(&user, &pass)
733 }
734 "LOGIN" => {
735 // RFC ietf-sasl-login: server prompts for username then
736 // password, both base64. The challenges may or may not
737 // be expected by the client; safest is to send the
738 // standard "Username:" / "Password:" prompts.
739 res!(write_response(
740 stream,
741 SmtpResponseCode::AuthInputData,
742 &base64::encode(b"Username:"),
743 ).await);
744 let user_b64 = match res!(read_line(stream).await) {
745 Some(l) => l.trim().to_string(),
746 None => return Ok(None),
747 };
748 let user_bytes = match base64::decode(user_b64.as_bytes()) {
749 Ok(b) => b,
750 Err(_) => return Ok(None),
751 };
752 res!(write_response(
753 stream,
754 SmtpResponseCode::AuthInputData,
755 &base64::encode(b"Password:"),
756 ).await);
757 let pass_b64 = match res!(read_line(stream).await) {
758 Some(l) => l.trim().to_string(),
759 None => return Ok(None),
760 };
761 let pass_bytes = match base64::decode(pass_b64.as_bytes()) {
762 Ok(b) => b,
763 Err(_) => return Ok(None),
764 };
765 let user = String::from_utf8_lossy(&user_bytes).to_string();
766 let pass = String::from_utf8_lossy(&pass_bytes).to_string();
767 self.users.authenticate(&user, &pass)
768 }
769 _ => {
770 Ok(None)
771 }
772 }
773 }
774}
775
776
777/// The two ways the inner read loop ends.
778#[derive(Clone, Copy, Debug)]
779enum LoopExit {
780 Quit, // QUIT, or a clean close
781 // STARTTLS was issued: the caller handshakes on the underlying TCP stream and re-enters the
782 // read loop.
783 StartTls,
784}
785
786
787// ┌───────────────────────────────────────────────────────────────────────────┐
788// │ WIRE HELPERS │
789// └───────────────────────────────────────────────────────────────────────────┘
790
791/// `NNN<space>text<CRLF>`, which is the single-line form.
792pub async fn write_response<S: AsyncWrite + Unpin>(
793 stream: &mut S,
794 code: SmtpResponseCode,
795 text: &str,
796)
797 -> Outcome<()>
798{
799 let line = fmt!("{} {}\r\n", code, text);
800 if let Err(e) = stream.write_all(line.as_bytes()).await {
801 return Err(err!(e, "Writing SMTP response."; IO, Network, Write));
802 }
803 if let Err(e) = stream.flush().await {
804 return Err(err!(e, "Flushing SMTP response."; IO, Network, Write));
805 }
806 Ok(())
807}
808
809/// `None` on a clean EOF. The trailing `\r\n`, or a bare `\n`, is trimmed, and the line is capped
810/// at [`SMTP_MAX_LINE`] bytes.
811pub async fn read_line<S: AsyncRead + Unpin>(
812 stream: &mut S,
813)
814 -> Outcome<Option<String>>
815{
816 let mut buf = Vec::with_capacity(128);
817 let mut byte = [0u8; 1];
818 loop {
819 let n = match stream.read(&mut byte).await {
820 Ok(n) => n,
821 Err(e) => return Err(err!(e, "Reading SMTP line byte."; IO, Network, Read)),
822 };
823 if n == 0 {
824 if buf.is_empty() {
825 return Ok(None);
826 }
827 break;
828 }
829 buf.push(byte[0]);
830 if byte[0] == b'\n' {
831 break;
832 }
833 if buf.len() >= SMTP_MAX_LINE {
834 return Err(err!(
835 "SMTP line exceeded {} bytes.", SMTP_MAX_LINE;
836 Invalid, Input, Excessive));
837 }
838 }
839 // Trim trailing CRLF / LF.
840 while buf.last() == Some(&b'\n') || buf.last() == Some(&b'\r') {
841 buf.pop();
842 }
843 Ok(Some(String::from_utf8_lossy(&buf).into_owned()))
844}
845
846/// Every byte up to and excluding the terminating `<CRLF>.<CRLF>`, dot-unstuffed: a line whose
847/// first character is `.` loses that dot, RFC 5321 §4.5.2. CR/LF is preserved as `\r\n`, so what
848/// comes out has canonical RFC 5322 line endings.
849pub async fn read_data<S: AsyncRead + Unpin>(
850 stream: &mut S,
851)
852 -> Outcome<Vec<u8>>
853{
854 let mut out = Vec::with_capacity(4096);
855 let mut at_line_start = true;
856 let mut just_saw_cr = false;
857 let mut byte = [0u8; 1];
858
859 loop {
860 let n = match stream.read(&mut byte).await {
861 Ok(n) => n,
862 Err(e) => return Err(err!(e,
863 "Reading SMTP DATA byte."; IO, Network, Read)),
864 };
865 if n == 0 {
866 return Err(err!(
867 "Connection closed mid-DATA.";
868 IO, Network, Read, Missing));
869 }
870 let b = byte[0];
871
872 // Detect the dot-stuff terminator: a line consisting solely of
873 // `.` ends the DATA. We track at_line_start to identify the
874 // first byte after a CRLF.
875 if at_line_start && b == b'.' {
876 // Peek ahead: read the next byte. If it is CR (followed by
877 // LF) the dot terminates the message; otherwise it is a
878 // dot-stuffed body line and we strip the leading dot.
879 let mut peek = [0u8; 1];
880 let m = match stream.read(&mut peek).await {
881 Ok(n) => n,
882 Err(e) => return Err(err!(e,
883 "Reading SMTP DATA dot peek."; IO, Network, Read)),
884 };
885 if m == 0 {
886 return Err(err!(
887 "Connection closed mid-DATA after lone dot.";
888 IO, Network, Read, Missing));
889 }
890 if peek[0] == b'\r' {
891 // Expect LF next.
892 let mut tail = [0u8; 1];
893 let k = match stream.read(&mut tail).await {
894 Ok(n) => n,
895 Err(e) => return Err(err!(e,
896 "Reading SMTP DATA terminator LF."; IO, Network, Read)),
897 };
898 if k == 0 || tail[0] != b'\n' {
899 return Err(err!(
900 "Malformed DATA terminator (.<CR> without <LF>).";
901 Invalid, Input));
902 }
903 return Ok(out);
904 }
905 // Not the terminator -- it was a dot-stuffed line. The
906 // leading dot has been consumed and discarded; the peeked
907 // byte is the actual first body byte of the line.
908 out.push(peek[0]);
909 at_line_start = peek[0] == b'\n';
910 just_saw_cr = peek[0] == b'\r';
911 if out.len() > SMTP_MESSAGE_SIZE_LIMIT {
912 return Err(err!(
913 "SMTP DATA exceeded size limit of {} bytes.",
914 SMTP_MESSAGE_SIZE_LIMIT;
915 Invalid, Input, Excessive));
916 }
917 continue;
918 }
919
920 out.push(b);
921 if out.len() > SMTP_MESSAGE_SIZE_LIMIT {
922 return Err(err!(
923 "SMTP DATA exceeded size limit of {} bytes.",
924 SMTP_MESSAGE_SIZE_LIMIT;
925 Invalid, Input, Excessive));
926 }
927
928 // Update line-start tracking. A new line starts after CRLF or
929 // a bare LF.
930 if just_saw_cr && b == b'\n' {
931 at_line_start = true;
932 just_saw_cr = false;
933 } else if b == b'\n' {
934 at_line_start = true;
935 just_saw_cr = false;
936 } else if b == b'\r' {
937 just_saw_cr = true;
938 at_line_start = false;
939 } else {
940 at_line_start = false;
941 just_saw_cr = false;
942 }
943 }
944}
945
946/// Why this account may not use `mail_from` as its envelope sender, or `None` when it may.
947///
948/// The null reverse-path is allowed: it names nobody, and the header is bound all the same.
949fn envelope_refusal(user: &MailUser, mail_from: &str) -> Option<String> {
950 if mail_from.is_empty() || user.may_send_as(mail_from) {
951 return None;
952 }
953 Some(fmt!("<{}> is not an address {} may send as", mail_from, user.address()))
954}
955
956/// Why this account may not send `message` under its header, or `None` when it may.
957///
958/// There must be exactly one `From`, and every address in it, and in `Sender` where there is
959/// one, must be one the account may send as. A second `From` is refused rather than checked,
960/// since which one a reader or a verifier heeds is theirs to choose. A redirect, which keeps
961/// the author's `From` and adds a `Resent-From`, is refused with the rest: `Resent-From` is not
962/// what DMARC or the reader goes by. A header section or an address list that cannot be read
963/// as RFC 5322 reads it is refused, never passed, so the addresses checked are the ones a
964/// receiver's parser reports.
965fn header_refusal(user: &MailUser, message: &[u8]) -> Option<String> {
966 let (head, _) = split_headers_body(message);
967 let fields = match header_fields(head) {
968 Ok(f) => f,
969 Err(e) => {
970 debug!("SMTP submission: the header section cannot be read: {}", e);
971 return Some(fmt!("the header section cannot be read as RFC 5322"));
972 },
973 };
974 let from: Vec<&(String, String)> = fields.iter()
975 .filter(|(n, _)| n.eq_ignore_ascii_case("From"))
976 .collect();
977 match from.len() {
978 0 => return Some(fmt!("the message has no From field")),
979 1 => {},
980 n => return Some(fmt!("the message has {} From fields, and may have one", n)),
981 }
982 for (name, value) in fields.iter()
983 .filter(|(n, _)| n.eq_ignore_ascii_case("From") || n.eq_ignore_ascii_case("Sender"))
984 {
985 let list = match addresses(value) {
986 Ok(l) => l,
987 Err(_) => return Some(fmt!("the {} field cannot be read", name)),
988 };
989 if list.is_empty() {
990 return Some(fmt!("the {} field names nobody", name));
991 }
992 for a in list {
993 if !user.may_send_as(&a) {
994 return Some(fmt!("{}: <{}> is not an address {} may send as",
995 name, a, user.address()));
996 }
997 }
998 }
999 None
1000}
1001
1002/// The payload is `authzid \0 authcid \0 passwd`, and the `authzid` is ignored.
1003fn parse_plain(raw: &[u8]) -> Option<(String, String)> {
1004 let mut nuls = raw.iter().enumerate().filter_map(|(i, b)| {
1005 if *b == 0u8 { Some(i) } else { None }
1006 });
1007 let first = ok!(nuls.next());
1008 let second = ok!(nuls.next());
1009 let user = ok!(String::from_utf8(raw[first + 1..second].to_vec()).ok());
1010 let pass = ok!(String::from_utf8(raw[second + 1..].to_vec()).ok());
1011 Some((user, pass))
1012}
1013
1014
1015// ┌───────────────────────────────────────────────────────────────────────────┐
1016// │ TESTS │
1017// └───────────────────────────────────────────────────────────────────────────┘
1018
1019#[cfg(test)]
1020mod tests {
1021 use super::*;
1022 use crate::{
1023 imap::client::Security,
1024 smtp::client::{
1025 OutboundClient,
1026 SubmissionConfig,
1027 },
1028 tls::BoundedTlsAcceptor,
1029 };
1030
1031 use std::sync::Mutex;
1032
1033 use tokio::net::TcpListener;
1034 use tokio_rustls::{
1035 rustls::{
1036 ClientConfig,
1037 RootCertStore,
1038 ServerConfig,
1039 pki_types::{
1040 CertificateDer,
1041 PrivateKeyDer,
1042 PrivatePkcs8KeyDer,
1043 },
1044 },
1045 TlsAcceptor,
1046 };
1047
1048 const HOST: &str = "mail.test.local";
1049
1050 /// An account that lists one identity besides its own.
1051 fn hello() -> MailUser {
1052 MailUser {
1053 local: fmt!("hello"),
1054 domain: fmt!("oxegen.test"),
1055 delivery_key: fmt!("oxegen.test/hello"),
1056 send_as: vec![fmt!("news@oxegen.test")],
1057 }
1058 }
1059
1060 #[test]
1061 fn the_envelope_is_bound_to_the_account_00() {
1062 let u = hello();
1063 assert_eq!(envelope_refusal(&u, "hello@oxegen.test"), None);
1064 assert_eq!(envelope_refusal(&u, "Hello@OXEGEN.test"), None, "case does not matter");
1065 assert_eq!(envelope_refusal(&u, "news@oxegen.test"), None, "a listed identity");
1066 assert_eq!(envelope_refusal(&u, ""), None, "the null reverse-path names nobody");
1067 match envelope_refusal(&u, "ceo@oxegen.test") {
1068 Some(why) => assert!(why.contains("ceo@oxegen.test") && why.contains("hello@oxegen.test"),
1069 "the refusal must name the address and the account: {}", why),
1070 None => panic!("another address in the relay's own domain was let through"),
1071 }
1072 }
1073
1074 fn message(headers: &str) -> Vec<u8> {
1075 fmt!("{}\r\nSubject: s\r\n\r\nbody\r\n", headers).into_bytes()
1076 }
1077
1078 #[test]
1079 fn the_header_from_is_bound_to_the_account_00() {
1080 let u = hello();
1081 for fine in [
1082 "From: hello@oxegen.test",
1083 "From: Jason Hoogland <hello@oxegen.test>",
1084 "From: Jason\r\n <hello@oxegen.test>",
1085 "From:\r\n\t\"Hoogland, Jason\"\r\n <hello@oxegen.test>",
1086 "From: \"Oxegen News\" <News@oxegen.test>",
1087 "From: \"The CEO <ceo@oxegen.test>: urgent\" <hello@oxegen.test>",
1088 "From : hello@oxegen.test",
1089 "from: hello@oxegen.test\r\nSender: news@oxegen.test",
1090 ] {
1091 assert_eq!(header_refusal(&u, &message(fine)), None, "{:?} was refused", fine);
1092 }
1093 for (bad, says) in [
1094 ("From: The CEO <ceo@oxegen.test>", "ceo@oxegen.test"),
1095 ("From: hello@oxegen.test, ceo@oxegen.test", "ceo@oxegen.test"),
1096 ("From: hello@oxegen.test\r\nSender: ceo@oxegen.test", "Sender"),
1097 ("From: hello@oxegen.test\r\nFrom: ceo@oxegen.test", "2 From"),
1098 ("To: bob@example.org", "no From"),
1099 ("From: Jason Hoogland", "cannot be read"),
1100 ("From: undisclosed-recipients:;", "names nobody"),
1101 // A redirect keeps the author's From, which DMARC and the reader go by.
1102 ("From: ceo@oxegen.test\r\nResent-From: hello@oxegen.test", "ceo@oxegen.test"),
1103 // D-06 audit F4-1: forms a receiver's parser reads otherwise than the old checker did.
1104 ("From: ceo@oxegen.test:hello@oxegen.test;", "From field cannot be read"),
1105 ("From: ceo@oxegen.test\rX: hello@oxegen.test", "header section cannot be read"),
1106 ("From: ceo@oxegen.test\r<hello@oxegen.test>", "header section cannot be read"),
1107 ("From\r\n : ceo@oxegen.test\r\nFrom: hello@oxegen.test", "2 From"),
1108 ("From: ceo@oxegen.test <hello@oxegen.test>", "From field cannot be read"),
1109 ("From: hello@oxegen.test\nFrom: ceo@oxegen.test", "header section cannot be read"),
1110 ("\r\nFrom: hello@oxegen.test", "header section cannot be read"),
1111 ("From: hello@oxegen.test\r\nSender: ceo@oxegen.test:hello@oxegen.test;",
1112 "Sender field cannot be read"),
1113 ] {
1114 match header_refusal(&u, &message(bad)) {
1115 Some(why) => assert!(why.contains(says), "{:?} refused as {:?}", bad, why),
1116 None => panic!("{:?} was let through", bad),
1117 }
1118 }
1119 }
1120
1121 /// One account, with the password "pw".
1122 #[derive(Clone)]
1123 struct Accounts;
1124
1125 impl UserStore for Accounts {
1126 fn authenticate(&self, address: &str, password: &str) -> Outcome<Option<MailUser>> {
1127 if password != "pw" {
1128 return Ok(None);
1129 }
1130 self.lookup(address)
1131 }
1132
1133 fn lookup(&self, address: &str) -> Outcome<Option<MailUser>> {
1134 Ok(if address.eq_ignore_ascii_case("hello@oxegen.test") { Some(hello()) } else { None })
1135 }
1136 }
1137
1138 /// Takes every submitted message, so the test can see what got past the server.
1139 #[derive(Clone, Default)]
1140 struct Taken(Arc<Mutex<Vec<SmtpTransaction>>>);
1141
1142 impl SmtpHandler for Taken {
1143 fn deliver_inbound(&self, _txn: SmtpTransaction) -> Outcome<HandlerOutcome> {
1144 Ok(HandlerOutcome::RejectPermanent(fmt!("not a receiving server")))
1145 }
1146
1147 fn submit_outbound(&self, txn: SmtpTransaction) -> Outcome<HandlerOutcome> {
1148 match self.0.lock() {
1149 Ok(mut g) => g.push(txn),
1150 Err(_) => return Err(err!("The test handler's lock was poisoned."; Test, Lock)),
1151 }
1152 Ok(HandlerOutcome::Accepted(fmt!("T1")))
1153 }
1154
1155 fn rcpt_acceptable(&self, _address: &str) -> bool {
1156 false
1157 }
1158 }
1159
1160 /// A submission listener on loopback, with a certificate for [`HOST`], and a client TLS
1161 /// configuration that trusts it.
1162 async fn submission_server(taken: Taken)
1163 -> Outcome<(std::net::SocketAddr, Arc<ClientConfig>)>
1164 {
1165 crate::tls::ensure_crypto_provider();
1166 let cert = res!(rcgen::generate_simple_self_signed(vec![HOST.to_string()]), Init);
1167 let der = res!(cert.serialize_der(), Init);
1168 let key = PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(cert.serialize_private_key_der()));
1169 let server_tls = res!(ServerConfig::builder()
1170 .with_no_client_auth()
1171 .with_single_cert(vec![CertificateDer::from(der.clone())], key), Init);
1172 let mut roots = RootCertStore::empty();
1173 res!(roots.add(CertificateDer::from(der)), Init);
1174 let client_tls = ClientConfig::builder().with_root_certificates(roots).with_no_client_auth();
1175
1176 let server = SmtpServer {
1177 handler: taken,
1178 users: Accounts,
1179 tls_acceptor: Some(BoundedTlsAcceptor::unbounded(
1180 TlsAcceptor::from(Arc::new(server_tls)))),
1181 hostname: Arc::new(HOST.to_string()),
1182 mode: SmtpMode::Submission,
1183 };
1184 let listener = res!(TcpListener::bind("127.0.0.1:0").await, IO, Network);
1185 let addr = res!(listener.local_addr(), IO, Network);
1186 tokio::spawn(async move {
1187 loop {
1188 let (sock, peer) = match listener.accept().await {
1189 Ok(x) => x,
1190 Err(_) => return,
1191 };
1192 let s = server.clone();
1193 tokio::spawn(async move {
1194 let _ = s.run(sock, peer).await;
1195 });
1196 }
1197 });
1198 Ok((addr, Arc::new(client_tls)))
1199 }
1200
1201 /// End to end, the real client against the real server over STARTTLS and AUTH: an account
1202 /// sends as itself, as the identity it lists, and under the null envelope, and a forged
1203 /// sender is refused `550 5.7.1` in the envelope and in the header alike, before the handler,
1204 /// and so the signer, sees it.
1205 #[tokio::test]
1206 async fn a_submission_may_send_only_as_its_account_00() -> Outcome<()> {
1207 let taken = Taken::default();
1208 let (addr, client_tls) = res!(submission_server(taken.clone()).await);
1209 let client = OutboundClient {
1210 hostname: Arc::new(fmt!("client.test.local")),
1211 tls_config: client_tls,
1212 };
1213 let cfg = SubmissionConfig::new(
1214 HOST.to_string(), addr.port(), Security::StartTls,
1215 fmt!("hello@oxegen.test"), fmt!("pw")).with_addr(addr);
1216 let rcpt = vec![fmt!("bob@example.org")];
1217 let body = |from: &str| fmt!(
1218 "From: {}\r\nTo: bob@example.org\r\nSubject: s\r\n\r\nhi\r\n", from);
1219
1220 res!(client.submit(&cfg, "hello@oxegen.test", &rcpt,
1221 body("Jason <hello@oxegen.test>").as_bytes()).await);
1222 res!(client.submit(&cfg, "news@oxegen.test", &rcpt,
1223 body("news@oxegen.test").as_bytes()).await);
1224 res!(client.submit(&cfg, "", &rcpt,
1225 body("\"The CEO <ceo@oxegen.test>: urgent\" <hello@oxegen.test>").as_bytes()).await);
1226
1227 for (envelope, from, which) in [
1228 ("ceo@oxegen.test", "hello@oxegen.test", "envelope"),
1229 ("hello@oxegen.test", "The CEO <ceo@oxegen.test>", "header"),
1230 ] {
1231 match client.submit(&cfg, envelope, &rcpt, body(from).as_bytes()).await {
1232 Ok(q) => return Err(err!(
1233 "A forged {} sender was accepted, queued as {}.", which, q; Test)),
1234 Err(e) => {
1235 let s = fmt!("{}", e);
1236 assert!(s.contains("550") && s.contains("5.7.1") && s.contains("ceo@oxegen.test"),
1237 "the {} refusal must be a 550 5.7.1 naming the address: {}", which, s);
1238 },
1239 }
1240 }
1241 // D-06 audit F4-1: the header a receiver's parser reads as from ceo@.
1242 match client.submit(&cfg, "hello@oxegen.test", &rcpt,
1243 body("ceo@oxegen.test:hello@oxegen.test;").as_bytes()).await
1244 {
1245 Ok(q) => return Err(err!(
1246 "A From the checker cannot read was accepted, queued as {}.", q; Test)),
1247 Err(e) => {
1248 let s = fmt!("{}", e);
1249 assert!(s.contains("550") && s.contains("5.7.1"),
1250 "an unreadable From must be refused 550 5.7.1: {}", s);
1251 },
1252 }
1253 let got = match taken.0.lock() {
1254 Ok(g) => g.iter().map(|t| t.mail_from.clone()).collect::<Vec<_>>(),
1255 Err(_) => Vec::new(),
1256 };
1257 assert_eq!(got, vec![fmt!("hello@oxegen.test"), fmt!("news@oxegen.test"), fmt!("")],
1258 "only the three honest messages may reach the handler, the null envelope's included");
1259 Ok(())
1260 }
1261}