oxedyne/fe2o3/fe2o3_mail/src/outbound.rs
6.2 KiB, 1 run
created by r1870400018:9816, 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 | //! On-disk spool for outbound mail awaiting delivery. |
| 2 | //! |
| 3 | //! Each accepted submission is written to a single file in the spool |
| 4 | //! directory. A background worker reads the spool, attempts delivery |
| 5 | //! via [`oxedyne_fe2o3_net::smtp::client::OutboundClient`], and removes |
| 6 | //! the file on success. The spool format is a tiny text envelope |
| 7 | //! followed by a blank line and the raw RFC 5322 message: |
| 8 | //! |
| 9 | //! ```text |
| 10 | //! From: postmaster@example.com |
| 11 | //! Rcpt: alice@example.org |
| 12 | //! Rcpt: bob@example.net |
| 13 | //! |
| 14 | //! <RFC 5322 message bytes> |
| 15 | //! ``` |
| 16 | //! |
| 17 | //! Filenames are `<unix>.<usec>.<pid>.<rand>.eml`. There is no retry |
| 18 | //! schedule beyond "the worker tries again on the next sweep" -- that |
| 19 | //! is enough for an MVP whose volume is a handful of messages a day. |
| 20 | |
| 21 | use oxedyne_fe2o3_core::{ |
| 22 | prelude::*, |
| 23 | rand::Rand, |
| 24 | }; |
| 25 | |
| 26 | use std::{ |
| 27 | fs::{self, File}, |
| 28 | io::{Read, Write}, |
| 29 | path::{Path, PathBuf}, |
| 30 | sync::Arc, |
| 31 | time::{SystemTime, UNIX_EPOCH}, |
| 32 | }; |
| 33 | |
| 34 | |
| 35 | /// One queued message ready for delivery. |
| 36 | #[derive(Clone, Debug)] |
| 37 | pub struct SpooledMessage { |
| 38 | /// Spool filename (no path). |
| 39 | pub filename: String, |
| 40 | /// Envelope sender. |
| 41 | pub mail_from: String, |
| 42 | /// Envelope recipients. |
| 43 | pub rcpt_to: Vec<String>, |
| 44 | /// Raw RFC 5322 message bytes. |
| 45 | pub body: Vec<u8>, |
| 46 | } |
| 47 | |
| 48 | /// Filesystem-backed outbound spool. |
| 49 | #[derive(Clone, Debug)] |
| 50 | pub struct OutboundSpool { |
| 51 | /// Spool directory. Must exist. |
| 52 | pub root: Arc<PathBuf>, |
| 53 | } |
| 54 | |
| 55 | impl OutboundSpool { |
| 56 | /// Build a spool rooted at `root`. Creates the directory if it |
| 57 | /// does not yet exist. |
| 58 | pub fn new(root: PathBuf) -> Outcome<Self> { |
| 59 | if !root.exists() { |
| 60 | if let Err(e) = fs::create_dir_all(&root) { |
| 61 | return Err(err!(e, |
| 62 | "Creating spool dir {:?}.", root; |
| 63 | IO, File, Init)); |
| 64 | } |
| 65 | } |
| 66 | Ok(Self { root: Arc::new(root) }) |
| 67 | } |
| 68 | |
| 69 | /// Append a new message to the spool. Returns the assigned queue |
| 70 | /// id (the filename without the `.eml` extension) so the SMTP |
| 71 | /// handler can echo it back to the client in the `250 OK` line. |
| 72 | pub fn enqueue( |
| 73 | &self, |
| 74 | mail_from: &str, |
| 75 | rcpt_to: &[String], |
| 76 | body: &[u8], |
| 77 | ) |
| 78 | -> Outcome<String> |
| 79 | { |
| 80 | let now = SystemTime::now() |
| 81 | .duration_since(UNIX_EPOCH) |
| 82 | .map(|d| d) |
| 83 | .unwrap_or_default(); |
| 84 | let suffix = Rand::generate_random_string(6, "abcdefghijklmnopqrstuvwxyz0123456789"); |
| 85 | let qid = fmt!("{}.{}.{}.{}", now.as_secs(), now.subsec_micros(), |
| 86 | std::process::id(), suffix); |
| 87 | let filename = fmt!("{}.eml", qid); |
| 88 | let path = self.root.join(&filename); |
| 89 | |
| 90 | let mut buf: Vec<u8> = Vec::with_capacity(body.len() + 256); |
| 91 | buf.extend_from_slice(fmt!("From: {}\n", mail_from).as_bytes()); |
| 92 | for r in rcpt_to { |
| 93 | buf.extend_from_slice(fmt!("Rcpt: {}\n", r).as_bytes()); |
| 94 | } |
| 95 | buf.extend_from_slice(b"\n"); |
| 96 | buf.extend_from_slice(body); |
| 97 | |
| 98 | let mut f = match File::create(&path) { |
| 99 | Ok(f) => f, |
| 100 | Err(e) => return Err(err!(e, |
| 101 | "Creating spool file {:?}.", path; |
| 102 | IO, File, Write)), |
| 103 | }; |
| 104 | if let Err(e) = f.write_all(&buf) { |
| 105 | return Err(err!(e, |
| 106 | "Writing spool file {:?}.", path; |
| 107 | IO, File, Write)); |
| 108 | } |
| 109 | Ok(qid) |
| 110 | } |
| 111 | |
| 112 | /// Read the entire spool, returning every message currently |
| 113 | /// queued. Used by the delivery worker on each sweep. |
| 114 | pub fn list(&self) -> Outcome<Vec<SpooledMessage>> { |
| 115 | let mut out: Vec<SpooledMessage> = Vec::new(); |
| 116 | let rd = match fs::read_dir(self.root.as_path()) { |
| 117 | Ok(r) => r, |
| 118 | Err(e) => return Err(err!(e, |
| 119 | "Reading spool dir {:?}.", self.root; |
| 120 | IO, File, Read)), |
| 121 | }; |
| 122 | for entry in rd.flatten() { |
| 123 | let path = entry.path(); |
| 124 | if !path.is_file() { continue; } |
| 125 | let name = entry.file_name().to_string_lossy().into_owned(); |
| 126 | if !name.ends_with(".eml") { continue; } |
| 127 | match Self::read_one(&path) { |
| 128 | Ok((from, rcpt, body)) => out.push(SpooledMessage { |
| 129 | filename: name, |
| 130 | mail_from: from, |
| 131 | rcpt_to: rcpt, |
| 132 | body, |
| 133 | }), |
| 134 | Err(e) => warn!("Skipping spool file {:?}: {}", path, e), |
| 135 | } |
| 136 | } |
| 137 | Ok(out) |
| 138 | } |
| 139 | |
| 140 | /// Remove a successfully-delivered message from the spool. |
| 141 | pub fn remove(&self, filename: &str) -> Outcome<()> { |
| 142 | let path = self.root.join(filename); |
| 143 | if let Err(e) = fs::remove_file(&path) { |
| 144 | return Err(err!(e, |
| 145 | "Removing spool file {:?}.", path; |
| 146 | IO, File, Write)); |
| 147 | } |
| 148 | Ok(()) |
| 149 | } |
| 150 | |
| 151 | /// Parse a single on-disk spool entry. |
| 152 | fn read_one(path: &Path) -> Outcome<(String, Vec<String>, Vec<u8>)> { |
| 153 | let mut bytes = Vec::new(); |
| 154 | let mut f = match File::open(path) { |
| 155 | Ok(f) => f, |
| 156 | Err(e) => return Err(err!(e, |
| 157 | "Opening {:?}.", path; IO, File, Read)), |
| 158 | }; |
| 159 | if let Err(e) = f.read_to_end(&mut bytes) { |
| 160 | return Err(err!(e, |
| 161 | "Reading {:?}.", path; IO, File, Read)); |
| 162 | } |
| 163 | // Find the blank line separating envelope from body. |
| 164 | let split = match find_double_lf(&bytes) { |
| 165 | Some(i) => i, |
| 166 | None => return Err(err!( |
| 167 | "Spool file {:?} missing envelope/body separator.", path; |
| 168 | Invalid, Input, Decode)), |
| 169 | }; |
| 170 | let header = String::from_utf8_lossy(&bytes[..split]).into_owned(); |
| 171 | let body = bytes[split + 2..].to_vec(); |
| 172 | let mut from = String::new(); |
| 173 | let mut rcpt: Vec<String> = Vec::new(); |
| 174 | for line in header.lines() { |
| 175 | if let Some(v) = line.strip_prefix("From: ") { |
| 176 | from = v.to_string(); |
| 177 | } else if let Some(v) = line.strip_prefix("Rcpt: ") { |
| 178 | rcpt.push(v.to_string()); |
| 179 | } |
| 180 | } |
| 181 | Ok((from, rcpt, body)) |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | fn find_double_lf(bytes: &[u8]) -> Option<usize> { |
| 186 | for i in 0..bytes.len().saturating_sub(1) { |
| 187 | if bytes[i] == b'\n' && bytes[i + 1] == b'\n' { |
| 188 | return Some(i); |
| 189 | } |
| 190 | } |
| 191 | None |
| 192 | } |