oxedyne/fe2o3/fe2o3_steel/src/srv/mail.rs
6.4 KiB, 54 runs
created by r1870400018:10040, 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 | //! Mail listener spawner. |
| 2 | //! |
| 3 | //! Boots the SMTP receive (port 25), SMTP submission (port 587) and |
| 4 | //! IMAP (port 993) listeners alongside the HTTPS server, sharing the |
| 5 | //! same rustls server config so a single ACME-issued certificate |
| 6 | //! covers every protocol. |
| 7 | //! |
| 8 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 9 | //! Anthropic Claude |
| 10 | |
| 11 | use crate::app::mail::AppMailHandler; |
| 12 | |
| 13 | use oxedyne_fe2o3_core::prelude::*; |
| 14 | use oxedyne_fe2o3_mail::passwd::PasswdFileUserStore; |
| 15 | use oxedyne_fe2o3_net::{ |
| 16 | imap::server::ImapServer, |
| 17 | smtp::server::{ |
| 18 | SmtpMode, |
| 19 | SmtpServer, |
| 20 | }, |
| 21 | tls::{ |
| 22 | BoundedTlsAcceptor, |
| 23 | Handshake, |
| 24 | }, |
| 25 | }; |
| 26 | |
| 27 | use std::{ |
| 28 | net::SocketAddr, |
| 29 | sync::{ |
| 30 | Arc, |
| 31 | atomic::{ |
| 32 | AtomicU32, |
| 33 | Ordering, |
| 34 | }, |
| 35 | }, |
| 36 | }; |
| 37 | |
| 38 | use tokio::net::TcpListener; |
| 39 | |
| 40 | |
| 41 | /// How many mail listeners bound, against how many the configuration asked for, |
| 42 | /// surfaced in the health body as `mail_down`. |
| 43 | /// |
| 44 | /// Counted up front from the configuration, before anything that can fail, so a |
| 45 | /// set-up that dies before it binds anything still shows every listener as down |
| 46 | /// rather than showing no mail server at all. |
| 47 | #[derive(Debug, Default)] |
| 48 | pub struct ListenerTally { |
| 49 | wanted: AtomicU32, |
| 50 | bound: AtomicU32, |
| 51 | } |
| 52 | |
| 53 | impl ListenerTally { |
| 54 | pub fn new_shared() -> Arc<Self> { |
| 55 | Arc::new(Self::default()) |
| 56 | } |
| 57 | |
| 58 | pub fn want(&self) { |
| 59 | self.wanted.fetch_add(1, Ordering::Relaxed); |
| 60 | } |
| 61 | |
| 62 | pub fn bound(&self) { |
| 63 | self.bound.fetch_add(1, Ordering::Relaxed); |
| 64 | } |
| 65 | |
| 66 | /// Listeners asked for and not bound, or `None` when none were asked for. |
| 67 | pub fn down(&self) -> Option<u32> { |
| 68 | let wanted = self.wanted.load(Ordering::Relaxed); |
| 69 | if wanted == 0 { |
| 70 | return None; |
| 71 | } |
| 72 | Some(wanted.saturating_sub(self.bound.load(Ordering::Relaxed))) |
| 73 | } |
| 74 | } |
| 75 | |
| 76 | |
| 77 | /// Runs the accept loop forever. An error on an individual accept is logged and |
| 78 | /// swallowed, so a single bad connection cannot kill the whole port. |
| 79 | pub async fn run_smtp_listener( |
| 80 | addr: SocketAddr, |
| 81 | server: SmtpServer<AppMailHandler, PasswdFileUserStore>, |
| 82 | tally: Option<Arc<ListenerTally>>, |
| 83 | ) |
| 84 | -> Outcome<()> |
| 85 | { |
| 86 | let listener = match TcpListener::bind(&addr).await { |
| 87 | Ok(l) => l, |
| 88 | Err(e) => return Err(err!(e, |
| 89 | "Binding SMTP listener on {}.", addr; |
| 90 | IO, Network, Init)), |
| 91 | }; |
| 92 | if let Some(t) = &tally { |
| 93 | t.bound(); |
| 94 | } |
| 95 | info!("SMTP {:?} listening on {} (mode={:?})", |
| 96 | server.mode, addr, server.mode); |
| 97 | loop { |
| 98 | let (stream, peer) = match listener.accept().await { |
| 99 | Ok(p) => p, |
| 100 | Err(e) => { |
| 101 | error!(err!(e, |
| 102 | "SMTP accept error on {}.", addr; |
| 103 | IO, Network)); |
| 104 | continue; |
| 105 | } |
| 106 | }; |
| 107 | info!("SMTP {:?} {}: connection from {}", server.mode, addr, peer); |
| 108 | let server = server.clone(); |
| 109 | tokio::spawn(async move { |
| 110 | if let Err(e) = server.run(stream, peer).await { |
| 111 | warn!("SMTP session error from {}: {}", peer, e); |
| 112 | } |
| 113 | }); |
| 114 | } |
| 115 | } |
| 116 | |
| 117 | /// Implicit TLS: each accepted connection completes a TLS handshake before |
| 118 | /// entering the IMAP state machine. |
| 119 | pub async fn run_imap_listener( |
| 120 | addr: SocketAddr, |
| 121 | tls_acceptor: BoundedTlsAcceptor, |
| 122 | server: ImapServer< |
| 123 | oxedyne_fe2o3_mail::maildir::MaildirStore, |
| 124 | PasswdFileUserStore, |
| 125 | >, |
| 126 | tally: Option<Arc<ListenerTally>>, |
| 127 | ) |
| 128 | -> Outcome<()> |
| 129 | { |
| 130 | let listener = match TcpListener::bind(&addr).await { |
| 131 | Ok(l) => l, |
| 132 | Err(e) => return Err(err!(e, |
| 133 | "Binding IMAP listener on {}.", addr; |
| 134 | IO, Network, Init)), |
| 135 | }; |
| 136 | if let Some(t) = &tally { |
| 137 | t.bound(); |
| 138 | } |
| 139 | info!("IMAP listening on {}", addr); |
| 140 | loop { |
| 141 | let (stream, peer) = match listener.accept().await { |
| 142 | Ok(p) => p, |
| 143 | Err(e) => { |
| 144 | error!(err!(e, |
| 145 | "IMAP accept error on {}.", addr; |
| 146 | IO, Network)); |
| 147 | continue; |
| 148 | } |
| 149 | }; |
| 150 | info!("IMAP {}: connection from {}", addr, peer); |
| 151 | let acceptor = tls_acceptor.clone(); |
| 152 | let server = server.clone(); |
| 153 | tokio::spawn(async move { |
| 154 | let tls = match acceptor.accept(stream).await { |
| 155 | Handshake::Ok(t) => t, |
| 156 | Handshake::TimedOut => { |
| 157 | warn!("IMAP TLS handshake from {} timed out.", peer); |
| 158 | return; |
| 159 | } |
| 160 | Handshake::Failed(e) => { |
| 161 | warn!("IMAP TLS handshake from {} failed: {}", peer, e); |
| 162 | return; |
| 163 | } |
| 164 | }; |
| 165 | if let Err(e) = server.run(tls, peer).await { |
| 166 | warn!("IMAP session error from {}: {}", peer, e); |
| 167 | } |
| 168 | }); |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | pub type AppSmtpServer = SmtpServer<AppMailHandler, PasswdFileUserStore>; |
| 173 | |
| 174 | pub type AppImapServer = ImapServer< |
| 175 | oxedyne_fe2o3_mail::maildir::MaildirStore, |
| 176 | PasswdFileUserStore, |
| 177 | >; |
| 178 | |
| 179 | /// Builds both SMTP servers -- receive on 25, submission on 587 -- sharing one |
| 180 | /// handler. Both get the TLS acceptor: submission needs it for the STARTTLS |
| 181 | /// upgrade, and the receive server offers it so a modern peer can |
| 182 | /// opportunistically encrypt. |
| 183 | pub fn build_smtp_servers( |
| 184 | handler: AppMailHandler, |
| 185 | users: PasswdFileUserStore, |
| 186 | tls_acceptor: Option<BoundedTlsAcceptor>, |
| 187 | hostname: Arc<String>, |
| 188 | ) |
| 189 | -> (AppSmtpServer, AppSmtpServer) |
| 190 | { |
| 191 | let receive = SmtpServer { |
| 192 | handler: handler.clone(), |
| 193 | users: users.clone(), |
| 194 | tls_acceptor: tls_acceptor.clone(), |
| 195 | hostname: hostname.clone(), |
| 196 | mode: SmtpMode::Receive, |
| 197 | }; |
| 198 | let submission = SmtpServer { |
| 199 | handler, |
| 200 | users, |
| 201 | tls_acceptor, |
| 202 | hostname, |
| 203 | mode: SmtpMode::Submission, |
| 204 | }; |
| 205 | (receive, submission) |
| 206 | } |
| 207 | |
| 208 | |
| 209 | #[cfg(test)] |
| 210 | mod tests { |
| 211 | use super::*; |
| 212 | |
| 213 | #[test] |
| 214 | fn listener_tally_reports_only_what_was_asked_for() { |
| 215 | let t = ListenerTally::default(); |
| 216 | assert_eq!(t.down(), None, "a host with no mail server reports no mail field"); |
| 217 | t.want(); |
| 218 | t.want(); |
| 219 | t.want(); |
| 220 | t.bound(); |
| 221 | t.bound(); |
| 222 | assert_eq!(t.down(), Some(1), "three asked for, two bound"); |
| 223 | t.bound(); |
| 224 | assert_eq!(t.down(), Some(0)); |
| 225 | } |
| 226 | } |