Oregami
Repositories/oxedyne/fe2o3

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
21use oxedyne_fe2o3_core::{
22 prelude::*,
23 rand::Rand,
24};
25
26use 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)]
37pub 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)]
50pub struct OutboundSpool {
51 /// Spool directory. Must exist.
52 pub root: Arc<PathBuf>,
53}
54
55impl 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
185fn 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}