Oregami
Repositories/oxedyne/fe2o3

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

80.4 KiB, 354 runs

created by r1870400018:9854, 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//! Client-side SMTP, for the two different conversations a sender can have.
2//!
3//! **Delivery** ([`OutboundClient::deliver`]) is what a mail server does: look up the recipient
4//! domain's MX, connect to the best-preference exchange on port 25, EHLO, opportunistic STARTTLS,
5//! then MAIL/RCPT/DATA. Nobody authenticates -- the receiving server accepts the mail because it is
6//! responsible for the recipient, not because it knows the sender.
7//!
8//! **Submission** ([`OutboundClient::submit`]) is what a mail *client* does, and it is a different
9//! conversation with a different party: connect to the account holder's own provider on the
10//! submission port, and prove who you are before the provider will carry anything. Without it a
11//! sender can only talk to servers that already wanted the message; with it, a sender can post mail
12//! through the account it holds a password for, which is how every desktop mail client works.
13//!
14//! No queue, no retry policy, no exponential backoff -- the caller is
15//! expected to drive retries itself by enqueueing the message in a
16//! spool directory and re-invoking the client. Keeps the abstraction
17//! useful for both a "fire and forget" path and a real queue runner.
18//!
19//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
20//! Anthropic Claude
21
22use crate::{
23 dns_resolver,
24 imap::client::Security,
25 smtp::server::read_line,
26 tls::{
27 self,
28 ClientStream,
29 },
30};
31
32use oxedyne_fe2o3_core::prelude::*;
33
34use std::{
35 net::{
36 IpAddr,
37 SocketAddr,
38 },
39 sync::Arc,
40 time::Duration,
41};
42
43use tokio::{
44 io::AsyncWriteExt,
45 net::TcpStream,
46 time::timeout,
47};
48use tokio_rustls::rustls::ClientConfig;
49
50
51// Generous, because some receiving MX hosts greylist or impose multi-second waits before 220.
52// It bounds each step of a conversation, not the whole of it.
53pub const SMTP_CLIENT_TIMEOUT: Duration = Duration::from_secs(60);
54
55// The reply to the final "." may take this long whatever the step deadline, as RFC 5321
56// §4.5.3.2.6 allows. A sender that gives up while the receiver is still accepting sends the
57// message again on its next try, and the receiver ends up with both.
58const DATA_DONE_TIMEOUT: Duration = Duration::from_secs(600);
59
60// Once the message is accepted, QUIT is a courtesy, and no step of it waits longer than this. A
61// whole step spent on it could run out a caller's own deadline and report an accepted message as
62// failed, which is a false alarm and, where the caller retries, a duplicate.
63const QUIT_TIMEOUT: Duration = Duration::from_secs(2);
64
65// The message goes over in pieces of this size, each held to the deadline on its own, so a large
66// message on a slow link is bounded by its progress rather than cut off for its size.
67const SEND_PIECE: usize = 64 * 1024;
68
69
70/// Where to post a message, and how to prove you may.
71///
72/// The provider's submission service, not a recipient's MX: the host is the one the account lives
73/// on, and the credential is the account's own. Port 587 conventionally starts in the clear and
74/// upgrades with `STARTTLS`; port 465 is TLS from the first byte.
75#[derive(Clone, Debug)]
76pub struct SubmissionConfig {
77 pub host: String, // also the name the certificate is validated against
78 pub port: u16, // conventionally 587 (STARTTLS) or 465 (implicit TLS)
79 pub security: Security,
80 pub user: String, // usually, but not always, the address being sent from
81 // For a provider with two-factor authentication this is an application password, not the
82 // password the human types into a browser.
83 pub password: String,
84 // Each step: the connect, the TLS handshake, each command and each reply. The reply to the
85 // message itself may take DATA_DONE_TIMEOUT where that is longer, and QUIT no more than
86 // QUIT_TIMEOUT.
87 pub timeout: Duration,
88 // Dialled instead of resolving `host`. The certificate is still validated against `host`, so
89 // pinning the address weakens nothing -- and a server connecting on behalf of a user must vet
90 // the address it dials rather than hand the name to the resolver twice.
91 pub addr: Option<SocketAddr>,
92}
93
94impl SubmissionConfig {
95
96 /// The conventional deadline, and no pinned address.
97 pub fn new(
98 host: impl Into<String>,
99 port: u16,
100 security: Security,
101 user: impl Into<String>,
102 password: impl Into<String>,
103 )
104 -> Self
105 {
106 Self {
107 host: host.into(),
108 port,
109 security,
110 user: user.into(),
111 password: password.into(),
112 timeout: SMTP_CLIENT_TIMEOUT,
113 addr: None,
114 }
115 }
116
117 pub fn with_addr(mut self, addr: SocketAddr) -> Self {
118 self.addr = Some(addr);
119 self
120 }
121
122 pub fn with_timeout(mut self, timeout: Duration) -> Self {
123 self.timeout = timeout;
124 self
125 }
126}
127
128
129/// One outbound delivery target after MX resolution, sorted into preference order by
130/// [`OutboundClient::deliver`].
131#[derive(Clone, Debug)]
132struct DeliveryTarget {
133 host: String, // MX exchange
134 addr: IpAddr,
135 // Always 25 for a resolved exchange, which is the only port mail is delivered on. It is a
136 // field rather than a literal so a fixture can stand in as an exchange on a loopback port;
137 // nothing reads it from configuration and nothing should.
138 port: u16,
139 preference: u16, // as the MX record gave it
140}
141
142/// Per-process outbound SMTP client.
143///
144/// Holds a rustls `ClientConfig` initialised with the system trust
145/// anchors so STARTTLS to any public MX validates correctly. Cheap to
146/// clone -- the inner config is in an `Arc`.
147#[derive(Clone)]
148pub struct OutboundClient {
149 // Sent in EHLO. Should be the public hostname of the sending server, the one that owns the
150 // IP whose PTR lines up.
151 pub hostname: Arc<String>,
152 pub tls_config: Arc<ClientConfig>, // for STARTTLS, built once
153}
154
155impl OutboundClient {
156
157 /// STARTTLS validation goes against the system CA bundle. A caller needing a custom root
158 /// store builds the `ClientConfig` itself.
159 pub fn with_system_roots(hostname: impl Into<String>) -> Outcome<Self> {
160 let cfg = res!(Self::default_tls_config());
161 Ok(Self {
162 hostname: Arc::new(hostname.into()),
163 tls_config: Arc::new(cfg),
164 })
165 }
166
167 /// The host's CA bundle, through [`crate::tls::default_client_config`], which every protocol
168 /// client in this crate shares.
169 pub fn default_tls_config() -> Outcome<ClientConfig> {
170 tls::default_client_config()
171 }
172
173 /// Each MX in preference order until one succeeds. The queue id is the first accepting
174 /// server's; where every host failed, the error is the last one's.
175 pub async fn deliver(
176 &self,
177 mail_from: &str,
178 rcpt_to: &[String],
179 body: &[u8],
180 )
181 -> Outcome<String>
182 {
183 if rcpt_to.is_empty() {
184 return Err(err!(
185 "OutboundClient::deliver called with no recipients.";
186 Invalid, Input, Missing));
187 }
188
189 // Group recipients by domain so each domain delivery is one
190 // SMTP transaction. The MVP only handles the common case of
191 // every recipient sharing one domain.
192 let domain = res!(extract_domain(&rcpt_to[0]));
193 for r in rcpt_to.iter().skip(1) {
194 let other = res!(extract_domain(r));
195 if !other.eq_ignore_ascii_case(&domain) {
196 return Err(err!(
197 "OutboundClient::deliver: multi-domain delivery is \
198 not supported in the MVP (got '{}' and '{}').",
199 domain, other;
200 Invalid, Input));
201 }
202 }
203
204 // MX lookup, then resolve each MX host to an A record.
205 let mxs = res!(
206 tokio::task::spawn_blocking(move || dns_resolver::lookup_mx(&domain)).await
207 .map_err(|e| err!("MX lookup task join failure: {}.", e;
208 IO, Network, Init))
209 );
210 let mxs = res!(mxs);
211
212 let mut targets: Vec<DeliveryTarget> = Vec::new();
213 for mx in &mxs {
214 let exchange = mx.exchange.clone();
215 let pref = mx.preference;
216 let addrs_outcome = tokio::task::spawn_blocking(move || {
217 dns_resolver::lookup_a(&exchange)
218 }).await;
219 let addrs = match addrs_outcome {
220 Ok(Ok(v)) => v,
221 _ => continue,
222 };
223 for ip in addrs {
224 targets.push(DeliveryTarget {
225 host: mx.exchange.clone(),
226 addr: IpAddr::V4(ip),
227 port: 25,
228 preference: pref,
229 });
230 }
231 }
232 self.deliver_to_exchanges(&targets, mail_from, rcpt_to, body, SMTP_CLIENT_TIMEOUT).await
233 }
234
235 /// The delivery loop itself, given the exchanges rather than resolving them.
236 ///
237 /// Split out of [`Self::deliver`] for one reason: everything above this point needs a live MX
238 /// lookup, so the loop below -- preference order, which failures are worth another exchange, and
239 /// whether the collapsed error is permanent -- could not be exercised at all. It carries
240 /// jarrah's outbound mail and had no test until 2026-08-17. Private, and takes the exchanges as
241 /// an argument rather than reading them from anywhere: this is not a way to configure where mail
242 /// goes, it is a way for a fixture to stand in as an exchange. `deadline` bounds each step of
243 /// each conversation, and is an argument so a fixture need not wait a minute to see one pass.
244 async fn deliver_to_exchanges(
245 &self,
246 targets: &[DeliveryTarget],
247 mail_from: &str,
248 rcpt_to: &[String],
249 body: &[u8],
250 deadline: Duration,
251 )
252 -> Outcome<String>
253 {
254 if targets.is_empty() {
255 return Err(err!(
256 "No reachable MX hosts for any of the configured recipients.";
257 IO, Network, Missing));
258 }
259 let mut targets: Vec<DeliveryTarget> = targets.to_vec();
260 targets.sort_by_key(|t| t.preference);
261
262 let mut last_err: Option<String> = None;
263 // Whether any exchange refused this recipient with a 5xx. A permanent rejection -- an unknown
264 // mailbox, a blocked sender -- is authoritative for the domain, so retrying another exchange or a
265 // later sweep will not cure it, and the caller wants to suppress the address rather than keep
266 // trying. A 4xx, a timeout or a connection error is transient and carries no such tag.
267 let mut permanent = false;
268 for tgt in &targets {
269 match self.try_one(tgt, mail_from, rcpt_to, body, deadline).await {
270 Ok(qid) => return Ok(qid),
271 Err(e) => {
272 if is_permanent(&e) {
273 permanent = true;
274 }
275 let msg = fmt!("MX {} ({}): {}", tgt.host, tgt.addr, e);
276 warn!("Outbound SMTP attempt failed: {}", msg);
277 last_err = Some(msg);
278 }
279 }
280 }
281 let last = last_err.unwrap_or_else(|| "(none)".to_string());
282 // The permanence is carried on the collapsed error so a single failed `deliver` still tells the
283 // caller whether to retry the address or give up on it, without unpicking the wrapped causes.
284 if permanent {
285 Err(err!(
286 "All MX delivery attempts failed; last error: {}", last;
287 IO, Network, Permanent))
288 } else {
289 Err(err!(
290 "All MX delivery attempts failed; last error: {}", last;
291 IO, Network))
292 }
293 }
294
295 /// The conversation a mail client has, not the one a mail server has: the provider
296 /// carries the message because the sender proved they hold the account, so the credential is
297 /// not optional and neither is the encryption under it. The client refuses to send the password
298 /// over a connection it could not secure -- a provider that offers no TLS on its submission
299 /// port is not one a password may be spoken to, and failing loudly is the only safe answer.
300 ///
301 /// Unlike delivery, `rcpt_to` may span any number of domains: the provider, not this client,
302 /// works out where each one goes. What comes back is whatever the provider said on accepting
303 /// the message, which usually carries its queue id.
304 pub async fn submit(
305 &self,
306 cfg: &SubmissionConfig,
307 mail_from: &str,
308 rcpt_to: &[String],
309 body: &[u8],
310 )
311 -> Outcome<String>
312 {
313 if rcpt_to.is_empty() {
314 return Err(err!(
315 "OutboundClient::submit called with no recipients.";
316 Invalid, Input, Missing));
317 }
318
319 let addr = match cfg.addr {
320 Some(a) => a,
321 None => {
322 let host = cfg.host.clone();
323 let ips = res!(
324 tokio::task::spawn_blocking(move || dns_resolver::lookup_a(&host)).await
325 .map_err(|e| err!("Submission host lookup task join failure: {}.", e;
326 IO, Network, Init))
327 );
328 let ips = res!(ips);
329 match ips.first() {
330 Some(ip) => SocketAddr::new(IpAddr::V4(*ip), cfg.port),
331 None => return Err(err!(
332 "The submission host {} resolves to no address.", cfg.host;
333 IO, Network, Missing)),
334 }
335 },
336 };
337
338 let mut conv = res!(Conversation::connect(addr, &cfg.host, cfg.timeout).await);
339
340 // TLS from the first byte, or in the clear until STARTTLS lifts it.
341 if cfg.security == Security::ImplicitTls {
342 conv = res!(conv.upgrade(self.tls_config.clone()).await);
343 }
344
345 let banner = res!(conv.reply().await);
346 if banner.code != 220 {
347 return Err(err!(
348 "Expected a 220 banner from {}, got {} {}", cfg.host, banner.code, banner.text;
349 IO, Network, Wire));
350 }
351
352 let mut ehlo = res!(self.ehlo(&mut conv).await);
353
354 if cfg.security == Security::StartTls {
355 let offered = ehlo.text.lines().any(|l| l.trim().eq_ignore_ascii_case("STARTTLS"));
356 if !offered {
357 return Err(err!(
358 "{} does not offer STARTTLS, so the account password cannot be sent to it \
359 without being readable on the wire.", cfg.host;
360 IO, Network, Invalid));
361 }
362 res!(conv.command("STARTTLS").await);
363 let resp = res!(conv.reply().await);
364 if resp.code != 220 {
365 return Err(err!(
366 "{} refused STARTTLS: {} {}", cfg.host, resp.code, resp.text;
367 IO, Network, Wire));
368 }
369 conv = res!(conv.upgrade(self.tls_config.clone()).await);
370 // The extension list before the upgrade cannot be trusted, and AUTH is usually only
371 // offered after it, so ask again inside TLS.
372 ehlo = res!(self.ehlo(&mut conv).await);
373 }
374
375 if cfg.security == Security::Plain {
376 warn!("Submitting to {} without TLS: the account password will cross the wire in \
377 the clear. Only a loopback test server should ever be reached this way.", cfg.host);
378 }
379
380 res!(authenticate(&mut conv, &ehlo, &cfg.user, &cfg.password).await);
381 let queue_id = res!(transact(&mut conv, mail_from, rcpt_to, body).await);
382
383 conv.quit().await;
384 Ok(queue_id)
385 }
386
387 async fn ehlo(&self, conv: &mut Conversation) -> Outcome<SmtpResponse> {
388 res!(conv.command(&fmt!("EHLO {}", self.hostname)).await);
389 let resp = res!(conv.reply().await);
390 if resp.code != 250 {
391 return Err(err!(
392 "EHLO rejected: {} {}", resp.code, resp.text;
393 IO, Network, Wire));
394 }
395 Ok(resp)
396 }
397
398 async fn try_one(
399 &self,
400 tgt: &DeliveryTarget,
401 mail_from: &str,
402 rcpt_to: &[String],
403 body: &[u8],
404 deadline: Duration,
405 )
406 -> Outcome<String>
407 {
408 let addr = SocketAddr::new(tgt.addr, tgt.port);
409 let mut conv = res!(Conversation::connect(addr, &tgt.host, deadline).await);
410
411 // Read the 220 banner.
412 let banner = res!(conv.reply().await);
413 if banner.code != 220 {
414 return Err(err!(
415 "Expected 220 banner, got {} {}", banner.code, banner.text;
416 IO, Network, Wire));
417 }
418
419 // EHLO, then look at extensions.
420 let ehlo = res!(self.ehlo(&mut conv).await);
421 let supports_starttls = ehlo.text.lines().any(|l| {
422 l.trim().eq_ignore_ascii_case("STARTTLS")
423 });
424
425 // Opportunistic STARTTLS.
426 if supports_starttls {
427 res!(conv.command("STARTTLS").await);
428 let resp = res!(conv.reply().await);
429 if resp.code == 220 {
430 conv = res!(conv.upgrade(self.tls_config.clone()).await);
431
432 // Re-issue EHLO inside TLS.
433 res!(conv.command(&fmt!("EHLO {}", self.hostname)).await);
434 let _ = res!(conv.reply().await);
435 }
436 }
437
438 let queue_id = res!(transact(&mut conv, mail_from, rcpt_to, body).await);
439
440 conv.quit().await;
441 Ok(queue_id)
442 }
443}
444
445/// One SMTP conversation: the stream, and the deadline every step on it is held to.
446///
447/// Every read, write and handshake goes through here, so no step can wait for ever. Until
448/// 2026-09-24 only the TCP connect was timed, and a server that took the connection and never
449/// spoke held `submit` for as long as the socket lived, whatever its timeout said -- and with it
450/// the text an alert sends beside its mail (D-06 audit A2).
451struct Conversation {
452 stream: ClientStream,
453 peer: String, // named in errors, and the name a certificate is validated against
454 timeout: Duration, // per step
455}
456
457impl Conversation {
458
459 async fn connect(addr: SocketAddr, peer: &str, deadline: Duration) -> Outcome<Self> {
460 let plain = match timeout(deadline, TcpStream::connect(addr)).await {
461 Ok(Ok(s)) => s,
462 Ok(Err(e)) => return Err(err!(e,
463 "Connecting to {} at {}.", peer, addr; IO, Network)),
464 Err(_) => return Err(err!(
465 "Timeout connecting to {} at {}.", peer, addr; IO, Network, Timeout)),
466 };
467 Ok(Self { stream: ClientStream::Plain(plain), peer: peer.to_string(), timeout: deadline })
468 }
469
470 /// TLS over the plain stream, from the first byte or once STARTTLS has been agreed.
471 async fn upgrade(self, tls_config: Arc<ClientConfig>) -> Outcome<Self> {
472 let plain = match self.stream.into_plain() {
473 Some(s) => s,
474 None => return Err(err!(
475 "A TLS upgrade was asked of the stream to {}, which is already encrypted.",
476 self.peer; Invalid, Bug)),
477 };
478 let stream = res!(tls::upgrade(plain, &self.peer, tls_config, self.timeout).await);
479 Ok(Self { stream, peer: self.peer, timeout: self.timeout })
480 }
481
482 /// All of one reply, however many lines it runs to, within the step deadline.
483 async fn reply(&mut self) -> Outcome<SmtpResponse> {
484 let deadline = self.timeout;
485 self.reply_within(deadline).await
486 }
487
488 async fn reply_within(&mut self, deadline: Duration) -> Outcome<SmtpResponse> {
489 match timeout(deadline, read_smtp_response(&mut self.stream)).await {
490 Ok(r) => r,
491 Err(_) => Err(err!(
492 "{} sent no reply within {:?}.", self.peer, deadline;
493 IO, Network, Read, Timeout)),
494 }
495 }
496
497 /// Ends the conversation once the message is accepted. The message is the receiver's by now,
498 /// so a QUIT that stalls or fails is logged and changes nothing.
499 async fn quit(mut self) {
500 self.timeout = self.timeout.min(QUIT_TIMEOUT);
501 let said = match self.command("QUIT").await {
502 Ok(()) => self.reply().await.map(|_| ()),
503 Err(e) => Err(e),
504 };
505 if let Err(e) = said {
506 debug!("{} accepted the message and then did not see QUIT through; the message \
507 stands: {}", self.peer, e);
508 }
509 }
510
511 /// CRLF-terminated, and flushed. The command is never named in an error, since AUTH carries
512 /// the credential.
513 async fn command(&mut self, cmd: &str) -> Outcome<()> {
514 self.send(fmt!("{}\r\n", cmd).as_bytes(), "a command").await
515 }
516
517 /// Written and flushed, each piece within the deadline.
518 async fn send(&mut self, bytes: &[u8], what: &str) -> Outcome<()> {
519 for piece in bytes.chunks(SEND_PIECE) {
520 match timeout(self.timeout, self.stream.write_all(piece)).await {
521 Ok(Ok(())) => (),
522 Ok(Err(e)) => return Err(err!(e,
523 "Writing {} to {}.", what, self.peer; IO, Network, Write)),
524 Err(_) => return Err(err!(
525 "{} took no more of {} within {:?}.", self.peer, what, self.timeout;
526 IO, Network, Write, Timeout)),
527 }
528 }
529 match timeout(self.timeout, self.stream.flush()).await {
530 Ok(Ok(())) => Ok(()),
531 Ok(Err(e)) => Err(err!(e,
532 "Flushing {} to {}.", what, self.peer; IO, Network, Write)),
533 Err(_) => Err(err!(
534 "{} took no more of {} within {:?}.", self.peer, what, self.timeout;
535 IO, Network, Write, Timeout)),
536 }
537 }
538}
539
540/// Walk MAIL/RCPT/DATA on a stream that is already open, secured and (where the server demands it)
541/// authenticated. Delivery and submission differ in how they reach this point and not at all in
542/// what they do once they are here, so they share the transaction rather than each keeping a copy
543/// of it -- the second copy is where the dot-stuffing gets forgotten.
544async fn transact(
545 conv: &mut Conversation,
546 mail_from: &str,
547 rcpt_to: &[String],
548 body: &[u8],
549)
550 -> Outcome<String>
551{
552 res!(conv.command(&fmt!("MAIL FROM:<{}>", mail_from)).await);
553 let resp = res!(conv.reply().await);
554 if resp.code / 100 != 2 {
555 // A 5xx is a permanent refusal, tagged so the caller can suppress rather than retry; a 4xx is
556 // transient and carries no such tag.
557 if resp.code / 100 == 5 {
558 return Err(err!(
559 "MAIL FROM rejected: {} {}", resp.code, resp.text;
560 IO, Network, Wire, Permanent));
561 }
562 return Err(err!(
563 "MAIL FROM rejected: {} {}", resp.code, resp.text;
564 IO, Network, Wire));
565 }
566 for r in rcpt_to {
567 res!(conv.command(&fmt!("RCPT TO:<{}>", r)).await);
568 let resp = res!(conv.reply().await);
569 if resp.code / 100 != 2 {
570 // A 5xx here is the no-such-mailbox case: permanent for this recipient, so it is tagged for
571 // suppression. A 4xx (greylisting, a full mailbox) is transient and is not.
572 if resp.code / 100 == 5 {
573 return Err(err!(
574 "RCPT TO:<{}> rejected: {} {}", r, resp.code, resp.text;
575 IO, Network, Wire, Permanent));
576 }
577 return Err(err!(
578 "RCPT TO:<{}> rejected: {} {}", r, resp.code, resp.text;
579 IO, Network, Wire));
580 }
581 }
582 res!(conv.command("DATA").await);
583 let resp = res!(conv.reply().await);
584 if resp.code != 354 {
585 return Err(err!(
586 "DATA rejected: {} {}", resp.code, resp.text;
587 IO, Network, Wire));
588 }
589
590 // A line of the body that begins with a full stop would otherwise end the message, and the
591 // terminator goes on a line of its own however the body ends.
592 let mut stuffed = dot_stuff(body);
593 if !body.ends_with(b"\r\n") {
594 stuffed.extend_from_slice(b"\r\n");
595 }
596 stuffed.extend_from_slice(b".\r\n");
597 res!(conv.send(&stuffed, "the message").await);
598
599 // The receiver may be filtering the message before it answers, and giving up on it here is
600 // how it comes to be delivered twice.
601 let done = conv.timeout.max(DATA_DONE_TIMEOUT);
602 let resp = res!(conv.reply_within(done).await);
603 if resp.code / 100 != 2 {
604 // A 5xx on the message itself -- refused content, a policy block -- will not be cured by resending
605 // the same message, so it is tagged permanent for the caller to suppress on.
606 if resp.code / 100 == 5 {
607 return Err(err!(
608 "Server rejected message: {} {}", resp.code, resp.text;
609 IO, Network, Wire, Permanent));
610 }
611 return Err(err!(
612 "Server rejected message: {} {}", resp.code, resp.text;
613 IO, Network, Wire));
614 }
615 Ok(resp.text)
616}
617
618/// Is this a permanent rejection -- a 5xx from the receiving server -- that retrying will not
619/// cure?
620///
621/// The one predicate a caller needs to tell "this address is bad, suppress it" from "the network
622/// hiccupped, try again later". A permanent failure is tagged [`ErrTag::Permanent`] where the 5xx is
623/// read off the wire in [`transact`], so this reads the tag rather than parse a status code out of a
624/// message.
625///
626/// The tag is read from the whole chain, as [`Error::tags`] gathers it. `res!` wraps a cause in an
627/// `Error::Upstream` carrying **no tags of its own**, so reading only the outer frame, as this did
628/// until 2026-08-17, made the predicate answer `false` to every permanent failure there has ever
629/// been: `submit` and `try_one` each pass `transact`'s error through one `res!`. Nothing was
630/// suppressed, and `fe2o3_steel`'s subscriber list kept mailing addresses their servers had refused
631/// outright.
632pub fn is_permanent(e: &Error<ErrTag>) -> bool {
633 e.tags().contains(&ErrTag::Permanent)
634}
635
636/// Prove to the provider that the sender holds the account.
637///
638/// `PLAIN` is preferred and `LOGIN` accepted, because between them they are what every provider
639/// worth submitting through offers. Both hand over the password in base64, which is an encoding and
640/// not a protection -- the only thing keeping it safe is the TLS underneath, which is why the
641/// caller establishes that first and refuses to proceed without it. `ehlo` is the extension list
642/// the server advertised inside TLS.
643async fn authenticate(
644 conv: &mut Conversation,
645 ehlo: &SmtpResponse,
646 user: &str,
647 password: &str,
648)
649 -> Outcome<()>
650{
651 let mut mechanisms: Vec<String> = Vec::new();
652 for line in ehlo.text.lines() {
653 let l = line.trim();
654 if l.len() >= 4 && l[..4].eq_ignore_ascii_case("AUTH") {
655 for m in l[4..].split_whitespace() {
656 mechanisms.push(m.to_uppercase());
657 }
658 }
659 }
660 if mechanisms.is_empty() {
661 return Err(err!(
662 "The server offers no AUTH mechanism, so there is no way to prove the account is \
663 ours and it will not carry the message. It advertised: {}",
664 ehlo.text.replace('\n', " | ");
665 IO, Network, Missing));
666 }
667
668 if mechanisms.iter().any(|m| m == "PLAIN") {
669 // RFC 4616: an authorisation identity we leave empty, then the account, then the password,
670 // each separated by a NUL.
671 let raw = fmt!("\0{}\0{}", user, password);
672 let cmd = fmt!("AUTH PLAIN {}", base64::encode(raw.as_bytes()));
673 res!(conv.command(&cmd).await);
674 let resp = res!(conv.reply().await);
675 return check_auth(&resp);
676 }
677
678 if mechanisms.iter().any(|m| m == "LOGIN") {
679 res!(conv.command("AUTH LOGIN").await);
680 let resp = res!(conv.reply().await);
681 if resp.code != 334 {
682 return Err(err!(
683 "AUTH LOGIN was refused before the username: {} {}", resp.code, resp.text;
684 IO, Network, Wire));
685 }
686 res!(conv.command(&base64::encode(user.as_bytes())).await);
687 let resp = res!(conv.reply().await);
688 if resp.code != 334 {
689 return Err(err!(
690 "The server rejected the username: {} {}", resp.code, resp.text;
691 IO, Network, Wire));
692 }
693 res!(conv.command(&base64::encode(password.as_bytes())).await);
694 let resp = res!(conv.reply().await);
695 return check_auth(&resp);
696 }
697
698 Err(err!(
699 "The server offers only {}, and this client can prove itself with PLAIN or LOGIN.",
700 mechanisms.join(", ");
701 IO, Network, Unimplemented))
702}
703
704/// Read the server's verdict on a login attempt, saying what a rejection usually means rather than
705/// only that it happened. A wrong password and a password the provider will not accept from a
706/// program look identical on the wire, and the second is the common case.
707fn check_auth(resp: &SmtpResponse) -> Outcome<()> {
708 if resp.code / 100 == 2 {
709 return Ok(());
710 }
711 if resp.code == 535 || resp.code == 534 {
712 return Err(err!(
713 "The provider rejected the credential ({} {}). If the account has two-factor \
714 authentication, an ordinary password will always be refused here and an application \
715 password is required.", resp.code, resp.text;
716 Invalid, Input, Unauthorised));
717 }
718 Err(err!(
719 "Authentication failed: {} {}", resp.code, resp.text;
720 IO, Network, Wire))
721}
722
723/// One parsed SMTP server response, potentially multi-line.
724#[derive(Clone, Debug)]
725struct SmtpResponse {
726 code: u16,
727 text: String, // the text lines, joined by '\n'
728}
729
730/// A response ends at the line whose fourth byte is a space rather than a hyphen. Untimed, so it
731/// is read only through [`Conversation::reply`].
732async fn read_smtp_response(stream: &mut ClientStream) -> Outcome<SmtpResponse> {
733 let mut text = String::new();
734 let mut code: u16 = 0;
735 loop {
736 let line = match res!(read_line(stream).await) {
737 Some(l) => l,
738 None => return Err(err!(
739 "Connection closed while reading SMTP response.";
740 IO, Network, Read)),
741 };
742 if line.len() < 4 {
743 return Err(err!(
744 "SMTP response line too short: '{}'.", line;
745 Invalid, Input, Decode));
746 }
747 let code_str = &line[..3];
748 let sep = line.as_bytes()[3];
749 let parsed: u16 = match code_str.parse() {
750 Ok(n) => n,
751 Err(_) => return Err(err!(
752 "SMTP response code '{}' not numeric.", code_str;
753 Invalid, Input, Decode)),
754 };
755 if code == 0 {
756 code = parsed;
757 }
758 if !text.is_empty() {
759 text.push('\n');
760 }
761 text.push_str(&line[4..]);
762 if sep == b' ' {
763 break;
764 }
765 if sep != b'-' {
766 return Err(err!(
767 "SMTP response line has invalid separator: '{}'.", line;
768 Invalid, Input, Decode));
769 }
770 }
771 Ok(SmtpResponse { code, text })
772}
773
774/// RFC 5321 §4.5.2: any line whose first character is `.` gets a second one prepended, so the
775/// receiver does not mistake it for the message terminator.
776fn dot_stuff(body: &[u8]) -> Vec<u8> {
777 let mut out = Vec::with_capacity(body.len() + body.len() / 64);
778 let mut at_line_start = true;
779 for &b in body {
780 if at_line_start && b == b'.' {
781 out.push(b'.');
782 }
783 out.push(b);
784 at_line_start = b == b'\n';
785 }
786 out
787}
788
789/// Lowercased, from the last `@`.
790fn extract_domain(addr: &str) -> Outcome<String> {
791 match addr.rfind('@') {
792 Some(i) => Ok(addr[i + 1..].to_lowercase()),
793 None => Err(err!(
794 "Address '{}' has no '@'.", addr;
795 Invalid, Input, Mismatch)),
796 }
797}
798
799
800// ┌───────────────────────────────────────────────────────────────────────────┐
801// │ TESTS │
802// └───────────────────────────────────────────────────────────────────────────┘
803
804// These drive the whole conversation against a stand-in provider on loopback, not
805// the pieces of it in isolation. They live in the module rather than in `tests/`
806// because `tests/main.rs` gates every case behind a single `filter` string, and on
807// 2026-08-17 that filter was `"dns"` -- so the four submission cases in
808// `tests/smtp_submit.rs` had never once run, while the harness reported them green.
809// A test that cannot be switched off by editing one word somewhere else is worth
810// more than a tidier home for it.
811#[cfg(test)]
812mod tests {
813 use super::*;
814
815 use std::sync::Mutex;
816
817 use tokio::{
818 io::{
819 AsyncBufReadExt,
820 AsyncReadExt,
821 BufReader,
822 },
823 net::TcpListener,
824 };
825
826
827 const USER: &str = "alice@example.com";
828 const PASS: &str = "app-password-not-a-real-one";
829 const EHLO: &str = "sender.test";
830 const WAIT: Duration = Duration::from_secs(10); // each step, where nothing should stall
831
832
833 /// What the stand-in server does when it is spoken to.
834 ///
835 /// It poses as a submission provider for [`OutboundClient::submit`] and as an MX exchange for
836 /// [`OutboundClient::deliver_to_exchanges`], which is honest: SMTP is the same protocol on both
837 /// sides of the difference, and the difference is entirely in what the client does with it.
838 #[derive(Clone, Copy)]
839 struct Provider {
840 mechs: &'static str, // advertised after `AUTH`; empty for no `AUTH` line at all
841 starttls: bool, // advertise `STARTTLS` in the EHLO reply
842 // Whether a `STARTTLS` command is answered 220. There is no TLS behind this fixture, so a
843 // test that advertises the extension answers 454 -- which is the opportunistic case worth
844 // covering anyway: a client offered TLS and refused it must deliver in the clear rather
845 // than give up on the exchange.
846 starttls_ok: bool,
847 auth_ok: bool,
848 banner: u16, // 220 is ready for mail; 421 refuses the connection
849 rcpt_code: u16, // 250 accepts the recipient
850 data_code: u16, // 250 accepts the message
851 stall: &'static str, // the command after which it goes silent; empty for never
852 slow_done: Duration, // how long the reply to the final "." is held back
853 }
854
855 impl Provider {
856 /// Offers both mechanisms this client can speak, and takes everything.
857 fn accepting() -> Self {
858 Self {
859 mechs: "PLAIN LOGIN",
860 starttls: false,
861 starttls_ok: false,
862 auth_ok: true,
863 banner: 220,
864 rcpt_code: 250,
865 data_code: 250,
866 stall: "",
867 slow_done: Duration::ZERO,
868 }
869 }
870
871 /// An exchange, which advertises no `AUTH`: a delivering client is not expected to prove
872 /// anything, and must not try.
873 fn exchange() -> Self {
874 Self { mechs: "", ..Self::accepting() }
875 }
876
877 fn ehlo_reply(&self) -> String {
878 let mut out = fmt!("250-provider.example.com\r\n250-SIZE 35882577\r\n");
879 if self.starttls {
880 out.push_str("250-STARTTLS\r\n");
881 }
882 if !self.mechs.is_empty() {
883 out.push_str(&fmt!("250-AUTH {}\r\n", self.mechs));
884 }
885 out.push_str("250 8BITMIME\r\n");
886 out
887 }
888 }
889
890 /// Every line the provider was sent, in order, including the DATA body.
891 type Transcript = Arc<Mutex<Vec<String>>>;
892
893 async fn provider(p: Provider) -> Outcome<(SocketAddr, Transcript)> {
894 let listener = res!(TcpListener::bind("127.0.0.1:0").await
895 .map_err(|e| err!(e, "Binding the stand-in provider."; IO, Network)));
896 let addr = res!(listener.local_addr()
897 .map_err(|e| err!(e, "Reading the stand-in provider's address."; IO, Network)));
898
899 let seen: Transcript = Arc::new(Mutex::new(Vec::new()));
900 let log = seen.clone();
901
902 tokio::spawn(async move {
903 let (sock, _) = match listener.accept().await {
904 Ok(x) => x,
905 Err(_) => return,
906 };
907 let (r, mut w) = sock.into_split();
908 let mut lines = BufReader::new(r).lines();
909
910 let _ = w.write_all(fmt!("{} provider.example.com ESMTP\r\n",
911 p.banner).as_bytes()).await;
912
913 let mut in_data = false;
914 let mut await_user = false;
915 let mut await_pass = false;
916 let mut silent = false;
917 let verdict = if p.auth_ok {
918 &b"235 2.7.0 Accepted\r\n"[..]
919 } else {
920 &b"535 5.7.8 Username and Password not accepted\r\n"[..]
921 };
922
923 while let Ok(Some(line)) = lines.next_line().await {
924 if let Ok(mut g) = log.lock() {
925 g.push(line.clone());
926 }
927 if !p.stall.is_empty() && line.to_uppercase().starts_with(p.stall) {
928 silent = true;
929 }
930 if silent {
931 continue; // reads on, and answers nothing
932 }
933 if in_data {
934 if line == "." {
935 in_data = false;
936 tokio::time::sleep(p.slow_done).await;
937 let _ = w.write_all(fmt!("{} 2.0.0 Ok: queued as STANDIN1\r\n",
938 p.data_code).as_bytes()).await;
939 }
940 continue;
941 }
942 if await_user {
943 await_user = false;
944 await_pass = true;
945 let _ = w.write_all(b"334 UGFzc3dvcmQ6\r\n").await; // "Password:"
946 continue;
947 }
948 if await_pass {
949 await_pass = false;
950 let _ = w.write_all(verdict).await;
951 continue;
952 }
953
954 let up = line.to_uppercase();
955 if up.starts_with("EHLO") {
956 let _ = w.write_all(p.ehlo_reply().as_bytes()).await;
957 } else if up.starts_with("STARTTLS") {
958 let _ = w.write_all(if p.starttls_ok {
959 &b"220 Go ahead\r\n"[..]
960 } else {
961 &b"454 4.7.0 TLS not available at the moment\r\n"[..]
962 }).await;
963 } else if up.starts_with("AUTH PLAIN") {
964 let _ = w.write_all(verdict).await;
965 } else if up.starts_with("AUTH LOGIN") {
966 await_user = true;
967 let _ = w.write_all(b"334 VXNlcm5hbWU6\r\n").await; // "Username:"
968 } else if up.starts_with("MAIL FROM") {
969 let _ = w.write_all(b"250 2.1.0 Ok\r\n").await;
970 } else if up.starts_with("RCPT TO") {
971 let _ = w.write_all(fmt!("{} recipient\r\n", p.rcpt_code).as_bytes()).await;
972 } else if up.starts_with("DATA") {
973 in_data = true;
974 let _ = w.write_all(b"354 End data with <CR><LF>.<CR><LF>\r\n").await;
975 } else if up.starts_with("QUIT") {
976 let _ = w.write_all(b"221 2.0.0 Bye\r\n").await;
977 return;
978 } else {
979 let _ = w.write_all(b"250 2.0.0 Ok\r\n").await;
980 }
981 }
982 });
983
984 Ok((addr, seen))
985 }
986
987 fn cfg(addr: SocketAddr, security: Security) -> SubmissionConfig {
988 SubmissionConfig::new("provider.example.com", addr.port(), security, USER, PASS)
989 .with_addr(addr)
990 .with_timeout(Duration::from_secs(10))
991 }
992
993 /// One line deliberately begins with a full stop, so dot-stuffing is under
994 /// test in every case that gets as far as the body.
995 fn body() -> Vec<u8> {
996 let mut s = String::new();
997 s.push_str("From: Alice <alice@example.com>\r\n");
998 s.push_str("To: Bob <bob@example.net>\r\n");
999 s.push_str("Subject: Hello\r\n");
1000 s.push_str("\r\n");
1001 s.push_str("A line.\r\n");
1002 s.push_str(".A line that begins with a full stop.\r\n");
1003 s.into_bytes()
1004 }
1005
1006 fn lines_of(t: &Transcript) -> Outcome<Vec<String>> {
1007 match t.lock() {
1008 Ok(g) => Ok(g.clone()),
1009 Err(_) => Err(err!("The provider's transcript was poisoned."; Lock, Poisoned)),
1010 }
1011 }
1012
1013 /// Whether the password reached the provider, in the clear or base64, on any
1014 /// line. The only safe answer for a conversation that never authenticated.
1015 fn password_crossed(lines: &[String]) -> bool {
1016 for l in lines {
1017 if l.contains(PASS) {
1018 return true;
1019 }
1020 if let Ok(raw) = base64::decode(l.trim()) {
1021 if String::from_utf8_lossy(&raw).contains(PASS) {
1022 return true;
1023 }
1024 }
1025 // AUTH PLAIN carries it inside the command.
1026 if let Some(b64) = l.split_whitespace().last() {
1027 if let Ok(raw) = base64::decode(b64) {
1028 if String::from_utf8_lossy(&raw).contains(PASS) {
1029 return true;
1030 }
1031 }
1032 }
1033 }
1034 false
1035 }
1036
1037 async fn client() -> Outcome<OutboundClient> {
1038 OutboundClient::with_system_roots(EHLO)
1039 }
1040
1041 // ── The submission conversation ───────────────────────────────
1042
1043 /// The whole exchange, in order, with the credential in it: this is the shape
1044 /// every caller of `submit` depends on and none of them can see.
1045 #[tokio::test]
1046 async fn test_submission_speaks_the_conversation_in_order_00() -> Outcome<()> {
1047 let (addr, seen) = res!(provider(Provider::accepting()).await);
1048 let c = res!(client().await);
1049 let qid = res!(c.submit(
1050 &cfg(addr, Security::Plain),
1051 "alice@example.com",
1052 &[fmt!("bob@example.net")],
1053 &body(),
1054 ).await);
1055 req!(true, qid.contains("STANDIN1"), "the provider's queue id must come back");
1056
1057 let lines = res!(lines_of(&seen));
1058 // Position matters: EHLO before AUTH, AUTH before MAIL FROM, and the
1059 // terminator after the body. A client that authenticates after offering
1060 // the envelope has told the provider who it is too late.
1061 let at = |want: &str| -> Option<usize> {
1062 lines.iter().position(|l| l.to_uppercase().starts_with(want))
1063 };
1064 let ehlo = res!(at("EHLO").ok_or_else(|| err!("No EHLO was sent."; Test, Missing)));
1065 let auth = res!(at("AUTH").ok_or_else(|| err!("No AUTH was sent."; Test, Missing)));
1066 let mail = res!(at("MAIL FROM").ok_or_else(|| err!("No MAIL FROM."; Test, Missing)));
1067 let rcpt = res!(at("RCPT TO").ok_or_else(|| err!("No RCPT TO."; Test, Missing)));
1068 let data = res!(at("DATA").ok_or_else(|| err!("No DATA."; Test, Missing)));
1069 req!(true, ehlo < auth, "AUTH came before EHLO");
1070 req!(true, auth < mail, "the envelope was offered before the login");
1071 req!(true, mail < rcpt, "RCPT TO came before MAIL FROM");
1072 req!(true, rcpt < data, "DATA came before RCPT TO");
1073 req!(true, lines.iter().any(|l| l == "."), "the DATA terminator never arrived");
1074 req!(true, lines.iter().any(|l| l.to_uppercase().starts_with("QUIT")),
1075 "the client did not say QUIT");
1076 Ok(())
1077 }
1078
1079 /// `PLAIN` is offered first and must be the one chosen, with an empty
1080 /// authorisation identity, the account, then the password, NUL-separated.
1081 #[tokio::test]
1082 async fn test_auth_plain_encodes_the_credential_as_rfc4616_asks_00() -> Outcome<()> {
1083 let (addr, seen) = res!(provider(Provider::accepting()).await);
1084 let c = res!(client().await);
1085 res!(c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1086 &body()).await);
1087
1088 let lines = res!(lines_of(&seen));
1089 let auth = res!(lines.iter()
1090 .find(|l| l.to_uppercase().starts_with("AUTH PLAIN"))
1091 .ok_or_else(|| err!("The client never sent AUTH PLAIN: {:?}", lines;
1092 Test, Missing)));
1093 let b64 = auth["AUTH PLAIN ".len()..].trim().to_string();
1094 let raw = res!(base64::decode(&b64));
1095 req!(fmt!("\0{}\0{}", USER, PASS).into_bytes(), raw);
1096 Ok(())
1097 }
1098
1099 /// A provider offering only `LOGIN` must be spoken to, not given up on: the
1100 /// account and the password go over on their own lines, base64 and nothing
1101 /// else.
1102 #[tokio::test]
1103 async fn test_auth_login_is_the_fallback_when_plain_is_absent_00() -> Outcome<()> {
1104 let p = Provider { mechs: "LOGIN", ..Provider::accepting() };
1105 let (addr, seen) = res!(provider(p).await);
1106 let c = res!(client().await);
1107 let qid = res!(c.submit(&cfg(addr, Security::Plain), USER,
1108 &[fmt!("bob@example.net")], &body()).await);
1109 req!(true, qid.contains("STANDIN1"));
1110
1111 let lines = res!(lines_of(&seen));
1112 req!(true, lines.iter().any(|l| l.to_uppercase().starts_with("AUTH LOGIN")));
1113 req!(true, lines.iter().any(|l| base64::decode(l.trim())
1114 .map(|b| b == USER.as_bytes()).unwrap_or(false)),
1115 "the account never went over");
1116 req!(true, lines.iter().any(|l| base64::decode(l.trim())
1117 .map(|b| b == PASS.as_bytes()).unwrap_or(false)),
1118 "the password never went over");
1119 Ok(())
1120 }
1121
1122 /// RFC 5321 §4.5.2 on the wire, not in a unit test of `dot_stuff`: a body line
1123 /// beginning with a full stop must arrive doubled, or it ends the message and
1124 /// the rest of it is read as commands.
1125 #[tokio::test]
1126 async fn test_a_leading_full_stop_arrives_doubled_00() -> Outcome<()> {
1127 let (addr, seen) = res!(provider(Provider::accepting()).await);
1128 let c = res!(client().await);
1129 res!(c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1130 &body()).await);
1131
1132 let lines = res!(lines_of(&seen));
1133 req!(true, lines.iter().any(|l| l.starts_with("..A line that begins")),
1134 "the line was not stuffed: {:?}", lines);
1135 // And the single-dot terminator is still exactly one dot.
1136 req!(1, lines.iter().filter(|l| *l == ".").count());
1137 Ok(())
1138 }
1139
1140 /// A body that does not end in CRLF still gets a terminator on a line of its
1141 /// own, rather than one glued to the last line of the message.
1142 #[tokio::test]
1143 async fn test_a_body_without_a_trailing_crlf_still_terminates_00() -> Outcome<()> {
1144 let (addr, seen) = res!(provider(Provider::accepting()).await);
1145 let c = res!(client().await);
1146 res!(c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1147 b"Subject: x\r\n\r\nno trailing newline").await);
1148
1149 let lines = res!(lines_of(&seen));
1150 req!(true, lines.iter().any(|l| l == "no trailing newline"),
1151 "the last body line was mangled: {:?}", lines);
1152 req!(true, lines.iter().any(|l| l == "."), "no terminator on its own line");
1153 Ok(())
1154 }
1155
1156 /// Submission may span domains -- the provider works out where each goes --
1157 /// and every recipient must get its own RCPT TO.
1158 #[tokio::test]
1159 async fn test_submission_offers_every_recipient_across_domains_00() -> Outcome<()> {
1160 let (addr, seen) = res!(provider(Provider::accepting()).await);
1161 let c = res!(client().await);
1162 res!(c.submit(
1163 &cfg(addr, Security::Plain),
1164 USER,
1165 &[fmt!("bob@example.net"), fmt!("carol@elsewhere.example"), fmt!("dan@third.test")],
1166 &body(),
1167 ).await);
1168
1169 let lines = res!(lines_of(&seen));
1170 for who in ["bob@example.net", "carol@elsewhere.example", "dan@third.test"] {
1171 req!(true, lines.iter().any(|l| l == &fmt!("RCPT TO:<{}>", who)),
1172 "{} was not offered", who);
1173 }
1174 req!(3, lines.iter().filter(|l| l.to_uppercase().starts_with("RCPT TO")).count());
1175 Ok(())
1176 }
1177
1178 // ── Refusals, and what must not have happened first ───────────
1179
1180 /// A refused credential is an error, and the message must not be offered
1181 /// anyway. The provider said no; sending on regardless is how a message is
1182 /// reported sent and never arrives.
1183 #[tokio::test]
1184 async fn test_a_refused_credential_stops_the_conversation_00() -> Outcome<()> {
1185 let p = Provider { auth_ok: false, ..Provider::accepting() };
1186 let (addr, seen) = res!(provider(p).await);
1187 let c = res!(client().await);
1188 let out = c.submit(&cfg(addr, Security::Plain), USER,
1189 &[fmt!("bob@example.net")], &body()).await;
1190 if out.is_ok() {
1191 return Err(err!(
1192 "The provider refused the credential and the client reported success.";
1193 Test, Invalid));
1194 }
1195 // The 535 case names the application-password fix, because a wrong
1196 // password and a password the provider will not take from a program are
1197 // indistinguishable on the wire and the second is the common one.
1198 let msg = match out {
1199 Err(e) => fmt!("{}", e),
1200 Ok(_) => String::new(),
1201 };
1202 req!(true, msg.to_lowercase().contains("application"),
1203 "the refusal did not name the fix: {}", msg);
1204 req!(false, msg.contains(PASS), "the password leaked into the error: {}", msg);
1205
1206 let lines = res!(lines_of(&seen));
1207 req!(false, lines.iter().any(|l| l.to_uppercase().starts_with("MAIL FROM")),
1208 "the client offered the envelope after its login was rejected");
1209 Ok(())
1210 }
1211
1212 /// A provider offering only a mechanism this client cannot speak gets no
1213 /// password at all -- not an attempt, and not a plaintext fallback.
1214 #[tokio::test]
1215 async fn test_an_unspeakable_mechanism_never_sees_the_password_00() -> Outcome<()> {
1216 let p = Provider { mechs: "XOAUTH2 GSSAPI", ..Provider::accepting() };
1217 let (addr, seen) = res!(provider(p).await);
1218 let c = res!(client().await);
1219 let out = c.submit(&cfg(addr, Security::Plain), USER,
1220 &[fmt!("bob@example.net")], &body()).await;
1221 if out.is_ok() {
1222 return Err(err!(
1223 "The client claimed to submit through a provider whose only mechanisms \
1224 it cannot speak."; Test, Invalid));
1225 }
1226 let lines = res!(lines_of(&seen));
1227 req!(false, password_crossed(&lines),
1228 "the password crossed the wire anyway: {:?}", lines);
1229 Ok(())
1230 }
1231
1232 /// A provider advertising no `AUTH` line at all is refused for the same
1233 /// reason, and the message with it.
1234 #[tokio::test]
1235 async fn test_a_provider_offering_no_auth_is_refused_00() -> Outcome<()> {
1236 let p = Provider { mechs: "", ..Provider::accepting() };
1237 let (addr, seen) = res!(provider(p).await);
1238 let c = res!(client().await);
1239 let out = c.submit(&cfg(addr, Security::Plain), USER,
1240 &[fmt!("bob@example.net")], &body()).await;
1241 req!(true, out.is_err(), "a provider that cannot be authenticated to was used");
1242 let lines = res!(lines_of(&seen));
1243 req!(false, password_crossed(&lines));
1244 req!(false, lines.iter().any(|l| l.to_uppercase().starts_with("MAIL FROM")));
1245 Ok(())
1246 }
1247
1248 /// `STARTTLS` was asked for and the provider does not offer it. The password
1249 /// would cross in the clear, so nothing crosses at all -- and the error says
1250 /// so, rather than only that something failed.
1251 #[tokio::test]
1252 async fn test_starttls_absent_means_the_password_is_withheld_00() -> Outcome<()> {
1253 let p = Provider { starttls: false, ..Provider::accepting() };
1254 let (addr, seen) = res!(provider(p).await);
1255 let c = res!(client().await);
1256 let out = c.submit(&cfg(addr, Security::StartTls), USER,
1257 &[fmt!("bob@example.net")], &body()).await;
1258 let msg = match out {
1259 Err(e) => fmt!("{}", e),
1260 Ok(_) => return Err(err!(
1261 "The client submitted a password to a provider offering no TLS.";
1262 Test, Invalid)),
1263 };
1264 req!(true, msg.contains("STARTTLS"), "the error did not name STARTTLS: {}", msg);
1265 let lines = res!(lines_of(&seen));
1266 req!(false, password_crossed(&lines),
1267 "the password crossed a connection that could not be secured: {:?}", lines);
1268 req!(false, lines.iter().any(|l| l.to_uppercase().starts_with("AUTH")),
1269 "the client began to authenticate anyway");
1270 Ok(())
1271 }
1272
1273 // ── Permanence, which is what a caller retries or suppresses on ──
1274
1275 /// A 5xx on a recipient is authoritative: the caller suppresses the address
1276 /// rather than sweeping it again forever. This is the only thing
1277 /// [`is_permanent`] is for, and the only way a caller can tell the two apart.
1278 #[tokio::test]
1279 async fn test_a_5xx_recipient_refusal_is_permanent_00() -> Outcome<()> {
1280 let p = Provider { rcpt_code: 550, ..Provider::accepting() };
1281 let (addr, _) = res!(provider(p).await);
1282 let c = res!(client().await);
1283 match c.submit(&cfg(addr, Security::Plain), USER,
1284 &[fmt!("nobody@example.net")], &body()).await
1285 {
1286 Ok(_) => Err(err!("A 550 on RCPT TO was reported as a send."; Test, Invalid)),
1287 Err(e) => {
1288 req!(true, is_permanent(&e),
1289 "a 550 was not tagged permanent, so the caller will retry it forever: {}", e);
1290 Ok(())
1291 },
1292 }
1293 }
1294
1295 /// A 4xx is the greylisting case, and carries no such tag: retried later it
1296 /// usually succeeds, and suppressing the address would lose the mail.
1297 #[tokio::test]
1298 async fn test_a_4xx_recipient_refusal_is_transient_00() -> Outcome<()> {
1299 let p = Provider { rcpt_code: 451, ..Provider::accepting() };
1300 let (addr, _) = res!(provider(p).await);
1301 let c = res!(client().await);
1302 match c.submit(&cfg(addr, Security::Plain), USER,
1303 &[fmt!("bob@example.net")], &body()).await
1304 {
1305 Ok(_) => Err(err!("A 451 on RCPT TO was reported as a send."; Test, Invalid)),
1306 Err(e) => {
1307 req!(false, is_permanent(&e),
1308 "a 451 was tagged permanent, so the caller will suppress a good \
1309 address: {}", e);
1310 Ok(())
1311 },
1312 }
1313 }
1314
1315 /// And the same distinction on the message itself, where a policy block is
1316 /// permanent and a full mailbox is not.
1317 #[tokio::test]
1318 async fn test_a_refused_message_carries_its_permanence_00() -> Outcome<()> {
1319 let (addr, _) = res!(provider(Provider { data_code: 552, ..Provider::accepting() }).await);
1320 let c = res!(client().await);
1321 match c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1322 &body()).await
1323 {
1324 Ok(_) => return Err(err!("A 552 was reported as a send."; Test, Invalid)),
1325 Err(e) => req!(true, is_permanent(&e), "a 552 on the message was not permanent: {}", e),
1326 }
1327
1328 let (addr, _) = res!(provider(Provider { data_code: 452, ..Provider::accepting() }).await);
1329 match c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1330 &body()).await
1331 {
1332 Ok(_) => Err(err!("A 452 was reported as a send."; Test, Invalid)),
1333 Err(e) => {
1334 req!(false, is_permanent(&e), "a 452 was tagged permanent: {}", e);
1335 Ok(())
1336 },
1337 }
1338 }
1339
1340 // ── Wire shape ────────────────────────────────────────────────
1341
1342 /// A multi-line greeting or EHLO reply is one response, and the client must
1343 /// read to the line whose fourth byte is a space rather than stopping at the
1344 /// first.
1345 #[tokio::test]
1346 async fn test_a_multiline_ehlo_is_read_as_one_response_00() -> Outcome<()> {
1347 // `Provider::ehlo_reply` sends four lines, three continued. Reaching AUTH
1348 // at all proves the mechanism list was read out of the last-but-one.
1349 let (addr, seen) = res!(provider(Provider::accepting()).await);
1350 let c = res!(client().await);
1351 res!(c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1352 &body()).await);
1353 let lines = res!(lines_of(&seen));
1354 req!(true, lines.iter().any(|l| l.to_uppercase().starts_with("AUTH PLAIN")),
1355 "the mechanism list was not read out of the continued EHLO reply");
1356 Ok(())
1357 }
1358
1359 /// The EHLO name the caller configured is the one that goes over: it is what
1360 /// a receiving server checks against the sending IP's PTR.
1361 #[tokio::test]
1362 async fn test_the_configured_ehlo_name_is_the_one_sent_00() -> Outcome<()> {
1363 let (addr, seen) = res!(provider(Provider::accepting()).await);
1364 let c = res!(client().await);
1365 res!(c.submit(&cfg(addr, Security::Plain), USER, &[fmt!("bob@example.net")],
1366 &body()).await);
1367 let lines = res!(lines_of(&seen));
1368 req!(true, lines.iter().any(|l| l == &fmt!("EHLO {}", EHLO)),
1369 "EHLO did not carry the configured name: {:?}", lines);
1370 Ok(())
1371 }
1372
1373 /// The envelope sender is the address the caller gave, angle-bracketed, and
1374 /// not the login -- the two differ whenever a person sends as an alias.
1375 #[tokio::test]
1376 async fn test_the_envelope_sender_is_not_the_login_00() -> Outcome<()> {
1377 let (addr, seen) = res!(provider(Provider::accepting()).await);
1378 let c = res!(client().await);
1379 res!(c.submit(&cfg(addr, Security::Plain), "alias@example.com",
1380 &[fmt!("bob@example.net")], &body()).await);
1381 let lines = res!(lines_of(&seen));
1382 req!(true, lines.iter().any(|l| l == "MAIL FROM:<alias@example.com>"),
1383 "the envelope sender was rewritten: {:?}", lines);
1384 Ok(())
1385 }
1386
1387 /// Nothing at all happens without a recipient: no connection, no banner read,
1388 /// no credential offered to a conversation that cannot carry a message.
1389 #[tokio::test]
1390 async fn test_no_recipient_is_refused_before_dialling_00() -> Outcome<()> {
1391 let (addr, seen) = res!(provider(Provider::accepting()).await);
1392 let c = res!(client().await);
1393 let out = c.submit(&cfg(addr, Security::Plain), USER, &[], &body()).await;
1394 req!(true, out.is_err(), "a message with no recipient was submitted");
1395 req!(true, res!(lines_of(&seen)).is_empty(),
1396 "the client dialled a provider for a message it could not send");
1397 Ok(())
1398 }
1399
1400 /// A greeting that is not a 220 is not a provider ready to take mail, and the
1401 /// client must stop there rather than talk over it.
1402 #[tokio::test]
1403 async fn test_a_refused_banner_stops_before_ehlo_00() -> Outcome<()> {
1404 let listener = res!(TcpListener::bind("127.0.0.1:0").await
1405 .map_err(|e| err!(e, "Binding a refusing provider."; IO, Network)));
1406 let addr = res!(listener.local_addr()
1407 .map_err(|e| err!(e, "Reading its address."; IO, Network)));
1408 let seen: Transcript = Arc::new(Mutex::new(Vec::new()));
1409 let log = seen.clone();
1410 tokio::spawn(async move {
1411 if let Ok((sock, _)) = listener.accept().await {
1412 let (r, mut w) = sock.into_split();
1413 let _ = w.write_all(b"554 no service here\r\n").await;
1414 let mut buf = Vec::new();
1415 let mut r = r;
1416 let _ = r.read_to_end(&mut buf).await;
1417 if let Ok(mut g) = log.lock() {
1418 g.push(String::from_utf8_lossy(&buf).into_owned());
1419 }
1420 }
1421 });
1422
1423 let c = res!(client().await);
1424 let out = c.submit(&cfg(addr, Security::Plain), USER,
1425 &[fmt!("bob@example.net")], &body()).await;
1426 let msg = match out {
1427 Err(e) => fmt!("{}", e),
1428 Ok(_) => return Err(err!(
1429 "A 554 greeting was treated as a provider ready for mail."; Test, Invalid)),
1430 };
1431 req!(true, msg.contains("220"), "the error did not say what was expected: {}", msg);
1432 let said = res!(lines_of(&seen)).join("");
1433 req!(false, said.to_uppercase().contains("EHLO"),
1434 "the client talked over a refusing banner: {:?}", said);
1435 Ok(())
1436 }
1437
1438 // ── The pieces, where the whole conversation cannot reach them ──
1439
1440 #[test]
1441 fn test_dot_stuff_doubles_only_at_a_line_start_00() -> Outcome<()> {
1442 req!(b"..a\r\n".to_vec(), dot_stuff(b".a\r\n"));
1443 req!(b"a.b\r\n".to_vec(), dot_stuff(b"a.b\r\n"));
1444 req!(b"x\r\n..y\r\n".to_vec(), dot_stuff(b"x\r\n.y\r\n"));
1445 // Two dots on their own line become three: the receiver unstuffs one.
1446 req!(b"...\r\n".to_vec(), dot_stuff(b"..\r\n"));
1447 Ok(())
1448 }
1449
1450 #[test]
1451 fn test_extract_domain_takes_the_last_at_00() -> Outcome<()> {
1452 req!(fmt!("example.com"), res!(extract_domain("a@example.com")));
1453 req!(fmt!("example.com"), res!(extract_domain("A@Example.COM")));
1454 // A quoted local part may itself contain an '@'.
1455 req!(fmt!("example.com"), res!(extract_domain("\"odd@name\"@example.com")));
1456 req!(true, extract_domain("no-at-sign").is_err());
1457 Ok(())
1458 }
1459
1460 /// The tag survives being wrapped. `res!` builds an `Error::Upstream` with no
1461 /// tags of its own, so a predicate reading one frame answers `false` however
1462 /// many 5xx there were underneath -- which is what made this inert.
1463 #[test]
1464 fn test_permanence_survives_a_wrapping_frame_00() -> Outcome<()> {
1465 let inner: Error<ErrTag> = err!(
1466 "RCPT TO:<nobody@example.net> rejected: 550 no such mailbox";
1467 IO, Network, Wire, Permanent);
1468 req!(true, is_permanent(&inner), "the tag is not read at the innermost frame");
1469
1470 let once = Error::Upstream(std::sync::Arc::new(inner), ErrMsg {
1471 tags: &[],
1472 msg: errmsg!(),
1473 });
1474 req!(true, is_permanent(&once), "one wrapping frame hid the tag");
1475
1476 let twice = Error::Upstream(std::sync::Arc::new(once), ErrMsg {
1477 tags: &[],
1478 msg: errmsg!(),
1479 });
1480 req!(true, is_permanent(&twice), "two wrapping frames hid the tag");
1481
1482 // And a transient failure stays transient however deep it is, or every
1483 // greylisted address would be suppressed.
1484 let soft: Error<ErrTag> = err!(
1485 "RCPT TO:<bob@example.net> rejected: 451 try later";
1486 IO, Network, Wire);
1487 let wrapped = Error::Upstream(std::sync::Arc::new(soft), ErrMsg {
1488 tags: &[],
1489 msg: errmsg!(),
1490 });
1491 req!(false, is_permanent(&wrapped), "a 451 was read as permanent");
1492 Ok(())
1493 }
1494
1495 /// Delivery is one SMTP transaction per domain, and the MVP does one domain.
1496 /// Refusing loudly beats delivering to the first domain and dropping the rest.
1497 #[tokio::test]
1498 async fn test_delivery_refuses_a_mixed_domain_envelope_00() -> Outcome<()> {
1499 let c = res!(client().await);
1500 req!(true, c.deliver("a@example.com", &[], b"x").await.is_err(),
1501 "delivery with no recipient was accepted");
1502 req!(true, c.deliver("a@example.com",
1503 &[fmt!("b@one.example"), fmt!("c@two.example")], b"x").await.is_err(),
1504 "delivery accepted two domains in one transaction");
1505 Ok(())
1506 }
1507
1508 // ── Delivery to an exchange, which carries jarrah's outbound mail ──
1509
1510 /// A stand-in exchange, and the target that points at it. `preference` is what
1511 /// the MX record would have said.
1512 async fn exchange_at(p: Provider, preference: u16) -> Outcome<(DeliveryTarget, Transcript)> {
1513 let (addr, seen) = res!(provider(p).await);
1514 Ok((DeliveryTarget {
1515 host: fmt!("mx{}.example.net", preference),
1516 addr: addr.ip(),
1517 port: addr.port(),
1518 preference,
1519 }, seen))
1520 }
1521
1522 /// The conversation a mail *server* has: no credential, because the receiving
1523 /// exchange takes the message for being responsible for the recipient. A
1524 /// delivering client that tries to log in is doing a mail client's job.
1525 #[tokio::test]
1526 async fn test_delivery_never_authenticates_00() -> Outcome<()> {
1527 // The exchange advertises AUTH anyway, which real ones commonly do.
1528 let (tgt, seen) = res!(exchange_at(Provider::accepting(), 10).await);
1529 let c = res!(client().await);
1530 let qid = res!(c.deliver_to_exchanges(&[tgt], "postmaster@example.com",
1531 &[fmt!("bob@example.net")], &body(), WAIT).await);
1532 req!(true, qid.contains("STANDIN1"));
1533
1534 let lines = res!(lines_of(&seen));
1535 req!(false, lines.iter().any(|l| l.to_uppercase().starts_with("AUTH")),
1536 "a delivering client tried to authenticate: {:?}", lines);
1537 req!(false, password_crossed(&lines));
1538 req!(true, lines.iter().any(|l| l == "MAIL FROM:<postmaster@example.com>"));
1539 req!(true, lines.iter().any(|l| l == "RCPT TO:<bob@example.net>"));
1540 req!(true, lines.iter().any(|l| l == "."), "the message was never terminated");
1541 // Dot-stuffing is the same on this path, and it is a different call site.
1542 req!(true, lines.iter().any(|l| l.starts_with("..A line that begins")));
1543 Ok(())
1544 }
1545
1546 /// Preference order, which is the whole point of holding a list: the lowest
1547 /// number is tried first.
1548 ///
1549 /// The preferred exchange refuses at `RCPT TO` rather than at the banner,
1550 /// deliberately. A banner-refusing fixture records nothing, so "its transcript
1551 /// is empty" is satisfied by *never having been dialled* as much as by having
1552 /// been -- and an assertion that cannot fail proves nothing. Written that way
1553 /// first, this case passed with `sort_by_key` deleted.
1554 #[tokio::test]
1555 async fn test_exchanges_are_tried_in_preference_order_00() -> Outcome<()> {
1556 let dead = Provider { rcpt_code: 550, ..Provider::exchange() };
1557 let c = res!(client().await);
1558
1559 // The refusing exchange is preferred, so it must be spoken to first and
1560 // the message must still land on the second.
1561 let (bad, saw_bad) = res!(exchange_at(dead, 10).await);
1562 let (good, saw_good) = res!(exchange_at(Provider::exchange(), 20).await);
1563 let qid = res!(c.deliver_to_exchanges(&[good.clone(), bad.clone()],
1564 "a@example.com", &[fmt!("bob@example.net")], &body(), WAIT).await);
1565 req!(true, qid.contains("STANDIN1"));
1566 req!(true, res!(lines_of(&saw_bad)).iter().any(|l| l.starts_with("RCPT TO")),
1567 "the preferred exchange was skipped: it was never offered the recipient");
1568 req!(true, res!(lines_of(&saw_good)).iter().any(|l| l == "."),
1569 "the message did not reach the second exchange");
1570
1571 // With the preferences swapped the good one is used first, and the
1572 // refusing one is never dialled at all.
1573 let (good, saw_good) = res!(exchange_at(Provider::exchange(), 10).await);
1574 let (bad, saw_bad) = res!(exchange_at(dead, 20).await);
1575 res!(c.deliver_to_exchanges(&[bad, good], "a@example.com",
1576 &[fmt!("bob@example.net")], &body(), WAIT).await);
1577 req!(true, res!(lines_of(&saw_good)).iter().any(|l| l == "."));
1578 req!(true, res!(lines_of(&saw_bad)).is_empty(),
1579 "a less-preferred exchange was used while a better one worked");
1580 Ok(())
1581 }
1582
1583 /// A 5xx recipient refusal is authoritative for the domain: it survives the
1584 /// collapse of every per-exchange error into one, so the caller suppresses the
1585 /// address instead of sweeping it forever. This is the path `deliver`'s own
1586 /// re-tagging was written for, and which never worked.
1587 #[tokio::test]
1588 async fn test_a_permanent_refusal_survives_the_collapse_00() -> Outcome<()> {
1589 let dead = Provider { rcpt_code: 550, ..Provider::exchange() };
1590 let (a, _) = res!(exchange_at(dead, 10).await);
1591 let (b, _) = res!(exchange_at(dead, 20).await);
1592 let c = res!(client().await);
1593 match c.deliver_to_exchanges(&[a, b], "a@example.com",
1594 &[fmt!("nobody@example.net")], &body(), WAIT).await
1595 {
1596 Ok(_) => Err(err!("Two 550s were reported as a delivery."; Test, Invalid)),
1597 Err(e) => {
1598 req!(true, is_permanent(&e),
1599 "a 550 from every exchange was not permanent, so the address is \
1600 retried forever: {}", e);
1601 Ok(())
1602 },
1603 }
1604 }
1605
1606 /// Greylisting is the common case and must not suppress anybody: a 4xx from
1607 /// every exchange collapses to a transient error.
1608 #[tokio::test]
1609 async fn test_a_transient_refusal_stays_transient_through_the_collapse_00() -> Outcome<()> {
1610 let busy = Provider { rcpt_code: 450, ..Provider::exchange() };
1611 let (a, _) = res!(exchange_at(busy, 10).await);
1612 let (b, _) = res!(exchange_at(busy, 20).await);
1613 let c = res!(client().await);
1614 match c.deliver_to_exchanges(&[a, b], "a@example.com",
1615 &[fmt!("bob@example.net")], &body(), WAIT).await
1616 {
1617 Ok(_) => Err(err!("Two 450s were reported as a delivery."; Test, Invalid)),
1618 Err(e) => {
1619 req!(false, is_permanent(&e),
1620 "greylisting was read as permanent, so a good address is \
1621 suppressed: {}", e);
1622 Ok(())
1623 },
1624 }
1625 }
1626
1627 /// One exchange refusing permanently is enough, even where another failed for
1628 /// a reason that carries no verdict. The address is bad; which exchange said
1629 /// so does not change that, and the flag must survive a later attempt.
1630 #[tokio::test]
1631 async fn test_one_permanent_refusal_among_failures_is_enough_00() -> Outcome<()> {
1632 let (dead, _) = res!(exchange_at(
1633 Provider { rcpt_code: 550, ..Provider::exchange() }, 10).await);
1634 let (refusing, _) = res!(exchange_at(
1635 Provider { banner: 421, ..Provider::exchange() }, 20).await);
1636 let c = res!(client().await);
1637 match c.deliver_to_exchanges(&[dead, refusing], "a@example.com",
1638 &[fmt!("nobody@example.net")], &body(), WAIT).await
1639 {
1640 Ok(_) => Err(err!("A 550 and a 421 were reported as a delivery."; Test, Invalid)),
1641 Err(e) => {
1642 req!(true, is_permanent(&e),
1643 "the permanent refusal was lost behind a later transient one: {}", e);
1644 Ok(())
1645 },
1646 }
1647 }
1648
1649 /// Opportunistic means opportunistic: an exchange that advertises `STARTTLS`
1650 /// and then will not do it still gets the mail, in the clear. Delivery has no
1651 /// credential to protect, and refusing here would silently stop mail to any
1652 /// exchange having a bad day with its certificate.
1653 #[tokio::test]
1654 async fn test_a_refused_starttls_still_delivers_in_the_clear_00() -> Outcome<()> {
1655 let p = Provider { starttls: true, starttls_ok: false, ..Provider::exchange() };
1656 let (tgt, seen) = res!(exchange_at(p, 10).await);
1657 let c = res!(client().await);
1658 let qid = res!(c.deliver_to_exchanges(&[tgt], "a@example.com",
1659 &[fmt!("bob@example.net")], &body(), WAIT).await);
1660 req!(true, qid.contains("STANDIN1"), "a refused STARTTLS stopped the delivery");
1661
1662 let lines = res!(lines_of(&seen));
1663 req!(true, lines.iter().any(|l| l.to_uppercase() == "STARTTLS"),
1664 "the offer was advertised and not taken up: {:?}", lines);
1665 req!(true, lines.iter().any(|l| l == "."), "the message never arrived");
1666 Ok(())
1667 }
1668
1669 /// No exchanges is a named error, not a silent success. A resolver that found
1670 /// nothing must not look like a delivery.
1671 #[tokio::test]
1672 async fn test_no_exchange_is_a_named_failure_00() -> Outcome<()> {
1673 let c = res!(client().await);
1674 let msg = match c.deliver_to_exchanges(&[], "a@example.com",
1675 &[fmt!("bob@example.net")], &body(), WAIT).await
1676 {
1677 Err(e) => fmt!("{}", e),
1678 Ok(_) => return Err(err!(
1679 "Delivery with nowhere to deliver reported success."; Test, Invalid)),
1680 };
1681 req!(true, msg.contains("No reachable MX"), "the error did not say why: {}", msg);
1682 Ok(())
1683 }
1684
1685 /// Where every exchange failed, the caller gets the last one's words -- and
1686 /// the exchange they came from, because "delivery failed" without a host is
1687 /// not something an operator can act on.
1688 #[tokio::test]
1689 async fn test_a_collapsed_error_names_an_exchange_00() -> Outcome<()> {
1690 let (a, _) = res!(exchange_at(
1691 Provider { rcpt_code: 550, ..Provider::exchange() }, 10).await);
1692 let host = a.host.clone();
1693 let c = res!(client().await);
1694 let msg = match c.deliver_to_exchanges(&[a], "a@example.com",
1695 &[fmt!("nobody@example.net")], &body(), WAIT).await
1696 {
1697 Err(e) => fmt!("{}", e),
1698 Ok(_) => return Err(err!("A 550 was a delivery."; Test, Invalid)),
1699 };
1700 req!(true, msg.contains(&host), "the failing exchange was not named: {}", msg);
1701 req!(true, msg.contains("550"), "the server's code was dropped: {}", msg);
1702 Ok(())
1703 }
1704
1705 // ── A peer that stops talking ─────────────────────────────────
1706
1707 /// A server that takes every connection and never says a word, as a wedged relay does while
1708 /// its kernel still completes the handshake. Each connection is held open, and silent.
1709 async fn mute_server() -> Outcome<SocketAddr> {
1710 let listener = res!(TcpListener::bind("127.0.0.1:0").await
1711 .map_err(|e| err!(e, "Binding the mute server."; IO, Network)));
1712 let addr = res!(listener.local_addr()
1713 .map_err(|e| err!(e, "Reading the mute server's address."; IO, Network)));
1714 tokio::spawn(async move {
1715 let mut held = Vec::new();
1716 while let Ok((sock, _)) = listener.accept().await {
1717 held.push(sock);
1718 }
1719 });
1720 Ok(addr)
1721 }
1722
1723 /// Runs one conversation against a peer that stalls, and requires it to fail, to fail
1724 /// within a few deadlines rather than hang, and to say that it timed out.
1725 async fn fails_in_time<F>(what: &str, conversation: F) -> Outcome<()>
1726 where
1727 F: std::future::Future<Output = Outcome<String>>,
1728 {
1729 let start = std::time::Instant::now();
1730 let out = match timeout(Duration::from_secs(10), conversation).await {
1731 Ok(o) => o,
1732 Err(_) => return Err(err!(
1733 "{}: still waiting after 10s, on a deadline of {:?}.", what, STALL;
1734 Test, Timeout)),
1735 };
1736 let took = start.elapsed();
1737 let msg = match out {
1738 Ok(_) => return Err(err!(
1739 "{}: a peer that stopped talking was reported to have taken the message.", what;
1740 Test, Invalid)),
1741 Err(e) => fmt!("{}", e),
1742 };
1743 req!(true, took < Duration::from_secs(5),
1744 "{}: took {:?} to fail on a deadline of {:?}", what, took, STALL);
1745 req!(true, msg.contains("within"), "{}: the error did not say it timed out: {}", what, msg);
1746 Ok(())
1747 }
1748
1749 const STALL: Duration = Duration::from_millis(500);
1750
1751 /// THE HANG THAT WAS D-06 A2: a relay that takes the connection and never greets. `submit`
1752 /// waited on the banner for as long as the socket lived, whatever its timeout said, and an
1753 /// alert's text waited behind it. The same peer met with TLS from the first byte stalls the
1754 /// handshake instead, which is timed by the same deadline.
1755 #[tokio::test]
1756 async fn test_a_server_that_never_speaks_fails_submission_in_time_00() -> Outcome<()> {
1757 let addr = res!(mute_server().await);
1758 let c = res!(client().await);
1759 for security in [Security::Plain, Security::ImplicitTls] {
1760 let cfg = cfg(addr, security).with_timeout(STALL);
1761 res!(fails_in_time(&fmt!("{:?}", security),
1762 c.submit(&cfg, USER, &[fmt!("bob@example.net")], &body())).await);
1763 }
1764 Ok(())
1765 }
1766
1767 /// A server that greets, takes the login and then goes silent is held to the same deadline:
1768 /// every reply is timed, not only the first.
1769 #[tokio::test]
1770 async fn test_a_server_that_stops_mid_conversation_fails_in_time_00() -> Outcome<()> {
1771 let (addr, seen) = res!(provider(
1772 Provider { stall: "MAIL FROM", ..Provider::accepting() }).await);
1773 let c = res!(client().await);
1774 let cfg = cfg(addr, Security::Plain).with_timeout(STALL);
1775 res!(fails_in_time("stalled at MAIL FROM",
1776 c.submit(&cfg, USER, &[fmt!("bob@example.net")], &body())).await);
1777 let lines = res!(lines_of(&seen));
1778 req!(true, lines.iter().any(|l| l.to_uppercase().starts_with("AUTH")),
1779 "the stall came before the login, so a later step was not tested: {:?}", lines);
1780 req!(false, lines.iter().any(|l| l.to_uppercase().starts_with("RCPT TO")),
1781 "the client talked on past a reply that never came: {:?}", lines);
1782 Ok(())
1783 }
1784
1785 /// Delivery meets the same mute peer as an exchange, and is held to its deadline too.
1786 #[tokio::test]
1787 async fn test_an_exchange_that_never_speaks_fails_delivery_in_time_00() -> Outcome<()> {
1788 let addr = res!(mute_server().await);
1789 let tgt = DeliveryTarget {
1790 host: fmt!("mx10.example.net"),
1791 addr: addr.ip(),
1792 port: addr.port(),
1793 preference: 10,
1794 };
1795 let c = res!(client().await);
1796 res!(fails_in_time("delivery",
1797 c.deliver_to_exchanges(&[tgt], "a@example.com", &[fmt!("bob@example.net")],
1798 &body(), STALL)).await);
1799 Ok(())
1800 }
1801
1802 // ── The ends of the transaction, which must not be cut short ──
1803
1804 /// A receiver that filters the message before it answers the final "." is waited for,
1805 /// beyond the step deadline. Cut off at the step deadline, as it was from 4f0e16c until
1806 /// D-06 audit R1, a receiver that accepted late was sent the message again on every retry.
1807 #[tokio::test]
1808 async fn test_a_slow_acceptance_is_waited_for_00() -> Outcome<()> {
1809 let slow = Provider { slow_done: STALL * 3, ..Provider::exchange() };
1810 let (tgt, seen) = res!(exchange_at(slow, 10).await);
1811 let c = res!(client().await);
1812 let qid = res!(c.deliver_to_exchanges(&[tgt], "a@example.com",
1813 &[fmt!("bob@example.net")], &body(), STALL).await);
1814 req!(true, qid.contains("STANDIN1"), "the late acceptance was not read: {}", qid);
1815 req!(1, res!(lines_of(&seen)).iter().filter(|l| *l == ".").count());
1816
1817 let (addr, _) = res!(provider(
1818 Provider { slow_done: STALL * 3, ..Provider::accepting() }).await);
1819 let cfg = cfg(addr, Security::Plain).with_timeout(STALL);
1820 let qid = res!(c.submit(&cfg, USER, &[fmt!("bob@example.net")], &body()).await);
1821 req!(true, qid.contains("STANDIN1"), "the late acceptance was not read: {}", qid);
1822 Ok(())
1823 }
1824
1825 /// Once the message is accepted, a receiver that goes silent at QUIT cannot turn the send
1826 /// into a failure, nor hold it for a whole step: the caller would report an accepted message
1827 /// as lost, and send it again.
1828 #[tokio::test]
1829 async fn test_a_silent_quit_does_not_fail_an_accepted_message_00() -> Outcome<()> {
1830 let mute_at_quit = Provider { stall: "QUIT", ..Provider::exchange() };
1831 let c = res!(client().await);
1832
1833 let (tgt, seen) = res!(exchange_at(mute_at_quit, 10).await);
1834 let start = std::time::Instant::now();
1835 let qid = res!(c.deliver_to_exchanges(&[tgt], "a@example.com",
1836 &[fmt!("bob@example.net")], &body(), WAIT).await);
1837 let took = start.elapsed();
1838 req!(true, qid.contains("STANDIN1"));
1839 req!(true, res!(lines_of(&seen)).iter().any(|l| l.to_uppercase() == "QUIT"),
1840 "the client never said QUIT");
1841 req!(true, took < WAIT / 2, "delivery waited {:?} on a silent QUIT", took);
1842
1843 let (addr, _) = res!(provider(
1844 Provider { stall: "QUIT", ..Provider::accepting() }).await);
1845 let start = std::time::Instant::now();
1846 let qid = res!(c.submit(&cfg(addr, Security::Plain), USER,
1847 &[fmt!("bob@example.net")], &body()).await);
1848 let took = start.elapsed();
1849 req!(true, qid.contains("STANDIN1"));
1850 req!(true, took < WAIT / 2, "submission waited {:?} on a silent QUIT", took);
1851 Ok(())
1852 }
1853}