Mail flow rules: carry out the transport actions #106
Generated
+2
@@ -3951,6 +3951,8 @@ dependencies = [
|
||||
"base64 0.23.1",
|
||||
"flate2",
|
||||
"jmap_proto",
|
||||
"mail-builder 1.0.0",
|
||||
"mail-parser",
|
||||
"quick-xml 0.41.0",
|
||||
"regex",
|
||||
"registry",
|
||||
|
||||
@@ -26,6 +26,8 @@ regex = "1.13.1"
|
||||
aho-corasick = "1.1"
|
||||
zip = "8.6"
|
||||
quick-xml = "0.41"
|
||||
mail-parser = { version = "0.11", features = ["full_encoding"] }
|
||||
mail-builder = { version = "1.0" }
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { version = "1.53", features = ["macros", "rt"] }
|
||||
|
||||
@@ -16,7 +16,8 @@
|
||||
//! - [`extract`]: the text of an attachment, or why it can't be read;
|
||||
//! - [`rules`]: what a rule is, its checks, and where rules are kept;
|
||||
//! - [`engine`]: rules compiled and run against a message;
|
||||
//! - [`cache`]: each node's compiled copy.
|
||||
//! - [`cache`]: each node's compiled copy;
|
||||
//! - [`rewrite`]: the actions that change a message.
|
||||
//!
|
||||
//! Nothing here writes what it finds anywhere: callers get counts, and the
|
||||
//! matched text never leaves the evaluation (§2.7).
|
||||
@@ -25,5 +26,6 @@ pub mod cache;
|
||||
pub mod detectors;
|
||||
pub mod engine;
|
||||
pub mod extract;
|
||||
pub mod rewrite;
|
||||
pub mod rules;
|
||||
pub mod words;
|
||||
|
||||
@@ -0,0 +1,306 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2026 Coffey Labs
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only
|
||||
*/
|
||||
|
||||
//! Transport actions that change a message (§2.4): headers, the subject,
|
||||
//! disclaimers. Each takes the raw message and returns the new one, or
|
||||
//! `None` when there's nothing to change.
|
||||
//!
|
||||
//! Only what the action names changes. A disclaimer edits the message's
|
||||
//! main text and HTML bodies (not attachments, not attached messages):
|
||||
//! each is decoded, changed and written back as UTF-8 quoted-printable,
|
||||
//! with its other headers kept. A disclaimer already there isn't added
|
||||
//! again, so a reply thread carries it once.
|
||||
|
||||
use base64::{Engine, engine::general_purpose::STANDARD};
|
||||
use mail_builder::encoders::quoted_printable::QuotedPrintableEncoder;
|
||||
use mail_parser::{HeaderName, MessageParser, PartType};
|
||||
|
||||
use super::rules::Position;
|
||||
|
||||
/// A header value, as an RFC 2047 encoded word when it isn't plain ASCII.
|
||||
pub fn header_value(value: &str) -> String {
|
||||
if value.is_ascii() {
|
||||
value.to_string()
|
||||
} else {
|
||||
format!("=?utf-8?B?{}?=", STANDARD.encode(value))
|
||||
}
|
||||
}
|
||||
|
||||
/// `Name: value` added at the top of the message.
|
||||
pub fn add_header(message: &[u8], name: &str, value: &str) -> Vec<u8> {
|
||||
let mut out = Vec::with_capacity(message.len() + name.len() + value.len() + 4);
|
||||
out.extend_from_slice(name.as_bytes());
|
||||
out.extend_from_slice(b": ");
|
||||
out.extend_from_slice(header_value(value).as_bytes());
|
||||
out.extend_from_slice(b"\r\n");
|
||||
out.extend_from_slice(message);
|
||||
out
|
||||
}
|
||||
|
||||
/// Every top-level header called `name` taken out.
|
||||
pub fn remove_header(message: &[u8], name: &str) -> Option<Vec<u8>> {
|
||||
let parsed = MessageParser::new().parse_headers(message)?;
|
||||
let mut ranges: Vec<(usize, usize)> = parsed
|
||||
.headers()
|
||||
.iter()
|
||||
.filter(|h| h.name.as_str().eq_ignore_ascii_case(name))
|
||||
.map(|h| (h.offset_field as usize, h.offset_end as usize))
|
||||
.collect();
|
||||
if ranges.is_empty() {
|
||||
return None;
|
||||
}
|
||||
ranges.sort_unstable();
|
||||
let mut out = Vec::with_capacity(message.len());
|
||||
let mut at = 0;
|
||||
for (start, end) in ranges {
|
||||
out.extend_from_slice(&message[at..start]);
|
||||
at = end;
|
||||
}
|
||||
out.extend_from_slice(&message[at..]);
|
||||
Some(out)
|
||||
}
|
||||
|
||||
/// The Subject header replaced by `subject` (added if there was none).
|
||||
pub fn set_subject(message: &[u8], subject: &str) -> Vec<u8> {
|
||||
let line = format!("Subject: {}\r\n", header_value(subject));
|
||||
let parsed = MessageParser::new().parse_headers(message);
|
||||
match parsed
|
||||
.as_ref()
|
||||
.and_then(|p| p.headers().iter().find(|h| h.name == HeaderName::Subject))
|
||||
{
|
||||
Some(header) => {
|
||||
let mut out = Vec::with_capacity(message.len() + line.len());
|
||||
out.extend_from_slice(&message[..header.offset_field as usize]);
|
||||
out.extend_from_slice(line.as_bytes());
|
||||
out.extend_from_slice(&message[header.offset_end as usize..]);
|
||||
out
|
||||
}
|
||||
None => {
|
||||
let mut out = line.into_bytes();
|
||||
out.extend_from_slice(message);
|
||||
out
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `prefix` put before the subject, unless it's already there.
|
||||
pub fn prefix_subject(message: &[u8], prefix: &str) -> Option<Vec<u8>> {
|
||||
let parsed = MessageParser::new().parse_headers(message)?;
|
||||
let subject = parsed.subject().unwrap_or_default();
|
||||
if subject.trim_start().starts_with(prefix.trim()) {
|
||||
return None;
|
||||
}
|
||||
Some(set_subject(
|
||||
message,
|
||||
&format!("{} {}", prefix.trim(), subject.trim_start()),
|
||||
))
|
||||
}
|
||||
|
||||
fn escape_html(text: &str) -> String {
|
||||
text.replace('&', "&")
|
||||
.replace('<', "<")
|
||||
.replace('>', ">")
|
||||
.replace('\n', "<br>\n")
|
||||
}
|
||||
|
||||
fn with_text_disclaimer(body: &str, text: &str, position: Position) -> String {
|
||||
let text = text.trim_end();
|
||||
match position {
|
||||
Position::Top => format!("{text}\r\n\r\n{body}"),
|
||||
Position::Bottom => format!("{}\r\n\r\n{text}\r\n", body.trim_end()),
|
||||
}
|
||||
}
|
||||
|
||||
fn with_html_disclaimer(body: &str, html: &str, position: Position) -> String {
|
||||
let lower = body.to_ascii_lowercase();
|
||||
match position {
|
||||
Position::Top => match lower
|
||||
.find("<body")
|
||||
.and_then(|at| lower[at..].find('>').map(|end| at + end + 1))
|
||||
{
|
||||
Some(at) => format!("{}{html}{}", &body[..at], &body[at..]),
|
||||
None => format!("{html}{body}"),
|
||||
},
|
||||
Position::Bottom => match lower.rfind("</body>") {
|
||||
Some(at) => format!("{}{html}{}", &body[..at], &body[at..]),
|
||||
None => format!("{body}{html}"),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// The disclaimer added to each main text and HTML body. `html` is the HTML
|
||||
/// version, or the text escaped when there's none.
|
||||
pub fn add_disclaimer(
|
||||
message: &[u8],
|
||||
text: &str,
|
||||
html: Option<&str>,
|
||||
position: Position,
|
||||
) -> Option<Vec<u8>> {
|
||||
let parsed = MessageParser::new().parse(message)?;
|
||||
let html = html
|
||||
.map(str::to_string)
|
||||
.unwrap_or_else(|| format!("<p>{}</p>", escape_html(text.trim())));
|
||||
let marker = text.trim();
|
||||
|
||||
let mut body_parts: Vec<u32> = parsed
|
||||
.text_body
|
||||
.iter()
|
||||
.chain(parsed.html_body.iter())
|
||||
.copied()
|
||||
.collect();
|
||||
body_parts.sort_unstable();
|
||||
body_parts.dedup();
|
||||
|
||||
// (start, end, replacement) for each part, applied from the last
|
||||
let mut edits: Vec<(usize, usize, Vec<u8>)> = Vec::new();
|
||||
for id in body_parts {
|
||||
let Some(part) = parsed.parts.get(id as usize) else {
|
||||
continue;
|
||||
};
|
||||
let (new_body, content_type) = match &part.body {
|
||||
PartType::Text(body) => {
|
||||
if body.contains(marker) {
|
||||
continue;
|
||||
}
|
||||
(with_text_disclaimer(body, text, position), "text/plain")
|
||||
}
|
||||
PartType::Html(body) => {
|
||||
if body.contains(marker) || body.contains(html.as_str()) {
|
||||
continue;
|
||||
}
|
||||
(with_html_disclaimer(body, &html, position), "text/html")
|
||||
}
|
||||
_ => continue,
|
||||
};
|
||||
// The part's own headers, less the two this changes
|
||||
let mut headers = Vec::new();
|
||||
for header in part.headers() {
|
||||
if matches!(
|
||||
header.name,
|
||||
HeaderName::ContentType | HeaderName::ContentTransferEncoding
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
headers.extend_from_slice(
|
||||
&message[header.offset_field as usize..header.offset_end as usize],
|
||||
);
|
||||
}
|
||||
headers.extend_from_slice(
|
||||
format!("Content-Type: {content_type}; charset=utf-8\r\n").as_bytes(),
|
||||
);
|
||||
headers.extend_from_slice(b"Content-Transfer-Encoding: quoted-printable\r\n\r\n");
|
||||
let encoded = QuotedPrintableEncoder::new()
|
||||
.preserve_line_breaks()
|
||||
.encode(new_body.as_bytes())
|
||||
.ok()?;
|
||||
headers.extend_from_slice(&encoded);
|
||||
// A single-part message's headers are the message's: its first
|
||||
// header is where the part starts
|
||||
let start = part.headers().first().map_or(part.offset_header, |h| {
|
||||
h.offset_field.min(part.offset_header)
|
||||
}) as usize;
|
||||
edits.push((start, part.offset_end as usize, headers));
|
||||
}
|
||||
if edits.is_empty() {
|
||||
return None;
|
||||
}
|
||||
edits.sort_by_key(|(start, _, _)| std::cmp::Reverse(*start));
|
||||
let mut out = message.to_vec();
|
||||
for (start, end, replacement) in edits {
|
||||
out.splice(start..end.min(out.len()), replacement);
|
||||
}
|
||||
Some(out)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn parse(message: &[u8]) -> mail_parser::Message<'_> {
|
||||
MessageParser::new().parse(message).expect("parses")
|
||||
}
|
||||
|
||||
const PLAIN: &[u8] = b"From: [email protected]\r\nTo: [email protected]\r\nSubject: Hello\r\nContent-Type: text/plain; charset=iso-8859-1\r\nContent-Transfer-Encoding: quoted-printable\r\n\r\nCaf=E9 at noon.\r\n";
|
||||
|
||||
const ALTERNATIVE: &[u8] = b"From: [email protected]\r\nSubject: Plans\r\nMIME-Version: 1.0\r\nContent-Type: multipart/mixed; boundary=\"outer\"\r\n\r\n--outer\r\nContent-Type: multipart/alternative; boundary=\"inner\"\r\n\r\n--inner\r\nContent-Type: text/plain\r\n\r\nSee you.\r\n--inner\r\nContent-Type: text/html\r\nContent-Transfer-Encoding: base64\r\n\r\nPGh0bWw+PGJvZHk+PHA+U2VlIHlvdS48L3A+PC9ib2R5PjwvaHRtbD4=\r\n--inner--\r\n--outer\r\nContent-Type: text/plain; name=\"notes.txt\"\r\nContent-Disposition: attachment; filename=\"notes.txt\"\r\n\r\nAttachment text.\r\n--outer--\r\n";
|
||||
|
||||
#[test]
|
||||
fn headers() {
|
||||
let added = add_header(PLAIN, "X-Mail-Rule", "External");
|
||||
assert_eq!(
|
||||
parse(&added).header_raw("X-Mail-Rule").map(str::trim),
|
||||
Some("External")
|
||||
);
|
||||
let removed = remove_header(&added, "x-mail-rule").unwrap();
|
||||
assert_eq!(removed, PLAIN);
|
||||
assert!(remove_header(PLAIN, "X-Absent").is_none());
|
||||
let utf8 = add_header(PLAIN, "X-Note", "Überprüft");
|
||||
// An RFC 2047 word: mail readers decode it, the wire stays ASCII
|
||||
assert_eq!(
|
||||
parse(&utf8).header_raw("X-Note").map(str::trim),
|
||||
Some("=?utf-8?B?w5xiZXJwcsO8ZnQ=?=")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn subjects() {
|
||||
let prefixed = prefix_subject(PLAIN, "[External]").unwrap();
|
||||
assert_eq!(parse(&prefixed).subject(), Some("[External] Hello"));
|
||||
assert!(prefix_subject(&prefixed, "[External]").is_none());
|
||||
let accented = set_subject(PLAIN, "Réunion à midi");
|
||||
assert_eq!(parse(&accented).subject(), Some("Réunion à midi"));
|
||||
assert!(accented.is_ascii(), "encoded as an RFC 2047 word");
|
||||
let none = set_subject(b"From: [email protected]\r\n\r\nBody\r\n", "New");
|
||||
assert_eq!(parse(&none).subject(), Some("New"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disclaimer_on_a_single_part() {
|
||||
let out = add_disclaimer(PLAIN, "Sent by Example Co.", None, Position::Bottom).unwrap();
|
||||
let parsed = parse(&out);
|
||||
let body = parsed.body_text(0).unwrap();
|
||||
assert!(body.starts_with("Café at noon."), "{body:?}");
|
||||
assert!(body.trim_end().ends_with("Sent by Example Co."), "{body:?}");
|
||||
assert_eq!(parsed.subject(), Some("Hello"));
|
||||
assert_eq!(
|
||||
parsed.header_raw("To").map(str::trim),
|
||||
Some("[email protected]")
|
||||
);
|
||||
// Once only
|
||||
assert!(add_disclaimer(&out, "Sent by Example Co.", None, Position::Bottom).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disclaimer_on_alternatives_leaves_attachments() {
|
||||
let out = add_disclaimer(
|
||||
ALTERNATIVE,
|
||||
"Confidential.",
|
||||
Some("<p><i>Confidential.</i></p>"),
|
||||
Position::Top,
|
||||
)
|
||||
.unwrap();
|
||||
let parsed = parse(&out);
|
||||
assert!(
|
||||
parsed
|
||||
.body_text(0)
|
||||
.unwrap()
|
||||
.starts_with("Confidential.\r\n\r\nSee you."),
|
||||
"{:?}",
|
||||
parsed.body_text(0)
|
||||
);
|
||||
let html = parsed.body_html(0).unwrap();
|
||||
assert!(
|
||||
html.contains("<body><p><i>Confidential.</i></p><p>See you.</p>"),
|
||||
"{html}"
|
||||
);
|
||||
assert_eq!(parsed.attachment_count(), 1);
|
||||
assert_eq!(
|
||||
parsed.attachment(0).unwrap().text_contents(),
|
||||
Some("Attachment text.")
|
||||
);
|
||||
assert!(!String::from_utf8_lossy(&out).contains("Confidential.\r\n\r\nAttachment"));
|
||||
}
|
||||
}
|
||||
@@ -100,6 +100,8 @@ pub struct SessionData {
|
||||
// message, for the submission to report
|
||||
pub dlp_override: Option<String>,
|
||||
pub dlp_refusal: Option<DlpRefusal>,
|
||||
// inbuxa: a mail flow rule's route for this message
|
||||
pub mailflow_queue: Option<String>,
|
||||
}
|
||||
|
||||
/// inbuxa: a DATA refusal by DLP rules: blocked, or a warning the sender
|
||||
@@ -186,6 +188,7 @@ impl SessionData {
|
||||
dnsbl_error: None,
|
||||
dlp_override: None,
|
||||
dlp_refusal: None,
|
||||
mailflow_queue: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -311,6 +314,7 @@ impl SessionData {
|
||||
dnsbl_error: None,
|
||||
dlp_override: None,
|
||||
dlp_refusal: None,
|
||||
mailflow_queue: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -745,7 +745,26 @@ impl<T: SessionStream> Session<T> {
|
||||
.await
|
||||
{
|
||||
super::mailflow::Checked::Accept => {}
|
||||
super::mailflow::Checked::Replace(message) => edited_message = Some(message),
|
||||
super::mailflow::Checked::Changed { message, envelope } => {
|
||||
if let Some(message) = message {
|
||||
edited_message = Some(message);
|
||||
}
|
||||
for change in envelope {
|
||||
match change {
|
||||
super::mailflow::EnvelopeChange::AddRecipient(address) => {
|
||||
if !self.data.rcpt_to.iter().any(|r| r.address_lcase.eq_ignore_ascii_case(&address)) {
|
||||
self.data.rcpt_to.push(SessionAddress::new(address));
|
||||
}
|
||||
}
|
||||
super::mailflow::EnvelopeChange::Redirect(addresses) => {
|
||||
self.data.rcpt_to = addresses.into_iter().map(SessionAddress::new).collect();
|
||||
}
|
||||
super::mailflow::EnvelopeChange::Route(queue) => {
|
||||
self.data.mailflow_queue = Some(queue);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
super::mailflow::Checked::Refuse(reply, refusal) => {
|
||||
self.data.dlp_refusal = refusal;
|
||||
return reply.into();
|
||||
@@ -917,8 +936,10 @@ impl<T: SessionStream> Session<T> {
|
||||
};
|
||||
|
||||
// Resolve queue
|
||||
let queue = self.server.get_queue_or_default(
|
||||
&self
|
||||
// inbuxa: a mail flow rule's route comes before the strategy
|
||||
let queue_name = match &self.data.mailflow_queue {
|
||||
Some(queue) => queue.clone(),
|
||||
None => self
|
||||
.server
|
||||
.eval_if::<String, _>(
|
||||
&self.server.core.smtp.queue.queue,
|
||||
@@ -927,8 +948,10 @@ impl<T: SessionStream> Session<T> {
|
||||
)
|
||||
.await
|
||||
.unwrap_or_else(|| "default".to_string()),
|
||||
self.data.session_id,
|
||||
);
|
||||
};
|
||||
let queue = self
|
||||
.server
|
||||
.get_queue_or_default(&queue_name, self.data.session_id);
|
||||
|
||||
// Set expiration and notification times
|
||||
let num_intervals = std::cmp::max(queue.notify.len(), 1);
|
||||
|
||||
@@ -21,10 +21,11 @@ use inbuxa_features::{
|
||||
Attachment, Content, Decision, Envelope, Outcome as RulesOutcome, Recipient, RuleRef,
|
||||
},
|
||||
extract::{self, Extracted, Limits},
|
||||
rules::Kind,
|
||||
rewrite,
|
||||
rules::{Action as RuleAction, Kind},
|
||||
},
|
||||
};
|
||||
use mail_parser::{HeaderName, Message, MessageParser, MimeHeaders, PartType};
|
||||
use mail_parser::{Message, MessageParser, MimeHeaders, PartType};
|
||||
use std::{borrow::Cow, time::SystemTime};
|
||||
|
||||
/// How much text one message is read for; past it, the rest counts as
|
||||
@@ -35,12 +36,23 @@ const INSPECTION_LIMIT: usize = 10 * 1024 * 1024;
|
||||
pub enum Checked {
|
||||
/// Go on, with the message unchanged.
|
||||
Accept,
|
||||
/// Go on with this message instead (the override tag taken out).
|
||||
Replace(Vec<u8>),
|
||||
/// Go on, with a changed message (the override tag taken out, a
|
||||
/// disclaimer, headers, the subject) and envelope.
|
||||
Changed {
|
||||
message: Option<Vec<u8>>,
|
||||
envelope: Vec<EnvelopeChange>,
|
||||
},
|
||||
/// Refuse, with this SMTP reply, and for a JMAP submission, why.
|
||||
Refuse(Vec<u8>, Option<DlpRefusal>),
|
||||
}
|
||||
|
||||
/// What a transport rule changes about where a message goes.
|
||||
pub enum EnvelopeChange {
|
||||
AddRecipient(String),
|
||||
Redirect(Vec<String>),
|
||||
Route(String),
|
||||
}
|
||||
|
||||
/// `[override: reason]` at the start of a subject: the reason, and the
|
||||
/// subject without it.
|
||||
pub fn override_tag(subject: &str) -> Option<(String, String)> {
|
||||
@@ -166,11 +178,14 @@ impl<T: SessionStream> Session<T> {
|
||||
/// DLP on an outgoing message (§2.4). `message` is what the DATA stage
|
||||
/// has so far (the script's replacement, if it made one).
|
||||
pub async fn check_mail_rules(&self, message: &[u8]) -> Checked {
|
||||
let Some(sender) = self.data.authenticated_as.as_ref() else {
|
||||
// DLP checks outgoing mail only (settled)
|
||||
return Checked::Accept;
|
||||
};
|
||||
let (account_id, account) = (sender.account_id, sender.account.clone());
|
||||
// Outgoing: an authenticated sender. DLP rules check outgoing mail
|
||||
// only (settled); transport rules may check either
|
||||
let sender = self
|
||||
.data
|
||||
.authenticated_as
|
||||
.as_ref()
|
||||
.map(|s| (s.account_id, s.account.clone()));
|
||||
let outgoing = sender.is_some();
|
||||
let rules = match cache::compiled(self.server.store()).await {
|
||||
Ok(rules) => rules,
|
||||
Err(err) => {
|
||||
@@ -187,7 +202,7 @@ impl<T: SessionStream> Session<T> {
|
||||
);
|
||||
}
|
||||
};
|
||||
if !rules.applies_to(true) {
|
||||
if !rules.applies_to(outgoing) {
|
||||
return Checked::Accept;
|
||||
}
|
||||
|
||||
@@ -197,7 +212,12 @@ impl<T: SessionStream> Session<T> {
|
||||
.and_then(|m| m.subject())
|
||||
.unwrap_or_default();
|
||||
let jmap_override = self.data.dlp_override.clone();
|
||||
let tag = override_tag(subject);
|
||||
// Only a sender of ours can override
|
||||
let tag = if outgoing {
|
||||
override_tag(subject)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let checked_subject = tag.as_ref().map_or(subject, |(_, rest)| rest.as_str());
|
||||
|
||||
let mut content = Content {
|
||||
@@ -253,10 +273,12 @@ impl<T: SessionStream> Session<T> {
|
||||
.map(|m| m.address_lcase.clone())
|
||||
.unwrap_or_default();
|
||||
let envelope = Envelope {
|
||||
outgoing: true,
|
||||
outgoing,
|
||||
sender: &sender_address,
|
||||
sender_groups: &account.id_member_of,
|
||||
sender_tenant: account.id_tenant,
|
||||
sender_groups: sender
|
||||
.as_ref()
|
||||
.map_or(&[][..], |(_, a)| &a.id_member_of[..]),
|
||||
sender_tenant: sender.as_ref().and_then(|(_, a)| a.id_tenant),
|
||||
recipients: self
|
||||
.data
|
||||
.rcpt_to
|
||||
@@ -285,15 +307,17 @@ impl<T: SessionStream> Session<T> {
|
||||
let domains = domains.join(", ");
|
||||
drop(envelope);
|
||||
|
||||
if let Some((account_id, account)) = &sender {
|
||||
self.record_dlp(
|
||||
account_id,
|
||||
&account,
|
||||
*account_id,
|
||||
account,
|
||||
&outcome,
|
||||
&decision,
|
||||
override_reason.as_deref(),
|
||||
&domains,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
match decision {
|
||||
Decision::Block(rules) | Decision::Hold { rules, .. } => {
|
||||
@@ -310,25 +334,117 @@ impl<T: SessionStream> Session<T> {
|
||||
.into_bytes(),
|
||||
Some(refusal(false, &rules)),
|
||||
),
|
||||
Decision::Pass => match (tag, &parsed) {
|
||||
Decision::Pass => {
|
||||
// The tag was an instruction to the server, not part of the
|
||||
// subject: it doesn't go out
|
||||
(Some((_, rest)), Some(parsed)) => parsed
|
||||
.headers()
|
||||
.iter()
|
||||
.find(|h| h.name == HeaderName::Subject)
|
||||
.map_or(Checked::Accept, |header| {
|
||||
let mut out = Vec::with_capacity(message.len());
|
||||
out.extend_from_slice(&message[..header.offset_field as usize]);
|
||||
out.extend_from_slice(b"Subject: ");
|
||||
out.extend_from_slice(rest.as_bytes());
|
||||
out.extend_from_slice(b"\r\n");
|
||||
out.extend_from_slice(&message[header.offset_end as usize..]);
|
||||
Checked::Replace(out)
|
||||
}),
|
||||
_ => Checked::Accept,
|
||||
},
|
||||
let mut current: Option<Vec<u8>> =
|
||||
tag.map(|(_, rest)| rewrite::set_subject(message, &rest));
|
||||
let mut changes = Vec::new();
|
||||
for matched in outcome.matched.iter().filter(|m| m.kind == Kind::Transport) {
|
||||
for action in &matched.actions {
|
||||
let now = current.as_deref().unwrap_or(message);
|
||||
let next = match action {
|
||||
RuleAction::AddDisclaimer { text, html, position } => {
|
||||
rewrite::add_disclaimer(now, text, html.as_deref(), *position)
|
||||
}
|
||||
RuleAction::AddHeader { name, value } => Some(rewrite::add_header(now, name, value)),
|
||||
RuleAction::RemoveHeader { name } => rewrite::remove_header(now, name),
|
||||
RuleAction::PrefixSubject { text } => rewrite::prefix_subject(now, text),
|
||||
RuleAction::AddRecipient { address } => {
|
||||
changes.push(EnvelopeChange::AddRecipient(address.clone()));
|
||||
None
|
||||
}
|
||||
RuleAction::Redirect { addresses } => {
|
||||
changes.push(EnvelopeChange::Redirect(addresses.clone()));
|
||||
None
|
||||
}
|
||||
RuleAction::Route { queue } => {
|
||||
changes.push(EnvelopeChange::Route(queue.clone()));
|
||||
None
|
||||
}
|
||||
RuleAction::Refuse { text } => {
|
||||
self.record_transport(&sender, &matched.name, "refused", &domains).await;
|
||||
return Checked::Refuse(
|
||||
format!("550 5.7.1 {}\r\n", reply_text(text)).into_bytes(),
|
||||
None,
|
||||
);
|
||||
}
|
||||
RuleAction::Block { .. } | RuleAction::Warn { .. } | RuleAction::Hold { .. } => None,
|
||||
};
|
||||
if next.is_some() {
|
||||
current = next;
|
||||
}
|
||||
}
|
||||
// Where mail goes is recorded; wording and headers aren't,
|
||||
// or a banner rule would write a record for every message
|
||||
let routed: Vec<String> = matched
|
||||
.actions
|
||||
.iter()
|
||||
.filter_map(|a| match a {
|
||||
RuleAction::AddRecipient { address } => Some(format!("copied to {address}")),
|
||||
RuleAction::Redirect { addresses } => Some(format!("redirected to {}", addresses.join(", "))),
|
||||
RuleAction::Route { queue } => Some(format!("routed through {queue}")),
|
||||
_ => None,
|
||||
})
|
||||
.collect();
|
||||
if !routed.is_empty() {
|
||||
self.record_transport(&sender, &matched.name, &routed.join(", "), &domains).await;
|
||||
}
|
||||
}
|
||||
if current.is_none() && changes.is_empty() {
|
||||
Checked::Accept
|
||||
} else {
|
||||
Checked::Changed { message: current, envelope: changes }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A transport rule that refused a message or changed where it goes
|
||||
/// (§2.7): who sent it (or the server, for incoming mail), the rule,
|
||||
/// what it did.
|
||||
async fn record_transport(
|
||||
&self,
|
||||
sender: &Option<(u32, std::sync::Arc<common::auth::AccountCache>)>,
|
||||
rule: &str,
|
||||
what: &str,
|
||||
domains: &str,
|
||||
) {
|
||||
let (actor, account_id, tenant_id) = match sender {
|
||||
Some((id, account)) => (
|
||||
Actor::account(*id, account.name.to_string(), account.id_tenant),
|
||||
Some(*id),
|
||||
account.id_tenant,
|
||||
),
|
||||
None => (Actor::system("mail-flow"), None, None),
|
||||
};
|
||||
let at = SystemTime::now()
|
||||
.duration_since(SystemTime::UNIX_EPOCH)
|
||||
.map_or(0, |d| d.as_millis() as u64);
|
||||
self.server
|
||||
.audit_note(Record {
|
||||
at,
|
||||
actor,
|
||||
via: None,
|
||||
remote_ip: Some(self.data.remote_ip),
|
||||
action: Action::Create,
|
||||
target: Target {
|
||||
kind: "message".into(),
|
||||
id: None,
|
||||
name: None,
|
||||
account_id,
|
||||
tenant_id,
|
||||
},
|
||||
changes: vec![],
|
||||
details: Some(format!("Mail flow rule \"{rule}\" {what}, to {domains}")),
|
||||
reason: None,
|
||||
outcome: if what == "refused" {
|
||||
Outcome::refused("forbidden", None)
|
||||
} else {
|
||||
Outcome::success()
|
||||
},
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
/// One audit record per message a DLP rule matched (§2.7): who sent it,
|
||||
|
||||
@@ -312,8 +312,13 @@ the sender's reason. No new audit action was added: an older node reading a
|
||||
record with an action it doesn't know fails its daily clean-up, so a new
|
||||
action would make rolling back unsafe.
|
||||
|
||||
Transport rules that change a message record the rule and action the same way.
|
||||
Unmatched mail writes nothing.
|
||||
Transport rules that refuse a message or change where it goes (redirect, add
|
||||
a recipient, route) record the rule and what it did the same way, the actor
|
||||
being the sender, or `system:mail-flow` for incoming mail. **As built
|
||||
(phase 2g)**, rules that only change wording or headers (a disclaimer, a
|
||||
header, a subject prefix) write nothing: a banner rule would otherwise write
|
||||
a record for every message, kept for the audit log's two years. Unmatched
|
||||
mail writes nothing.
|
||||
|
||||
### 2.8 Permissions and who does what
|
||||
|
||||
|
||||
@@ -11,10 +11,11 @@
|
||||
use crate::utils::{
|
||||
account::Account,
|
||||
server::{TestServer, TestServerBuilder},
|
||||
smtp::SmtpConnection,
|
||||
};
|
||||
use registry::schema::{
|
||||
prelude::{ObjectType, Property},
|
||||
structs::{CustomRoles, Role, UserRoles},
|
||||
structs::{CustomRoles, Expression, MtaStageAuth, Role, UserRoles},
|
||||
};
|
||||
use registry::types::map::Map;
|
||||
use serde_json::{Value, json};
|
||||
@@ -502,6 +503,276 @@ pub async fn dlp(test: &mut TestServer) {
|
||||
.await;
|
||||
}
|
||||
|
||||
/// The subjects, bodies and a header of what `account` received whose
|
||||
/// subject contains `text`, polling until something arrives.
|
||||
async fn received(
|
||||
account: &Account,
|
||||
text: &str,
|
||||
header: &str,
|
||||
drafts: Option<&str>,
|
||||
) -> Vec<(String, String, String)> {
|
||||
for _ in 0..50 {
|
||||
// Not the sender's own draft
|
||||
let filter = match drafts {
|
||||
Some(drafts) => json!({"subject": text, "inMailboxOtherThan": [drafts]}),
|
||||
None => json!({"subject": text}),
|
||||
};
|
||||
let (_, response) = call(account, "Email/query", json!({"filter": filter})).await;
|
||||
let ids = response["ids"].clone();
|
||||
let (_, response) = call(
|
||||
account,
|
||||
"Email/get",
|
||||
json!({"ids": ids, "properties": ["subject", "bodyValues", "textBody", format!("header:{header}:asText")],
|
||||
"fetchTextBodyValues": true}),
|
||||
)
|
||||
.await;
|
||||
let list: Vec<(String, String, String)> = response["list"]
|
||||
.as_array()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.filter(|e| {
|
||||
!e["subject"]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.starts_with("Failed to deliver")
|
||||
})
|
||||
.map(|e| {
|
||||
let part = e["textBody"][0]["partId"]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.to_string();
|
||||
(
|
||||
e["subject"].as_str().unwrap_or_default().to_string(),
|
||||
e["bodyValues"][part.as_str()]["value"]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.to_string(),
|
||||
e[format!("header:{header}:asText").as_str()]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.to_string(),
|
||||
)
|
||||
})
|
||||
.collect();
|
||||
if !list.is_empty() {
|
||||
return list;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
|
||||
}
|
||||
Vec::new()
|
||||
}
|
||||
|
||||
/// Transport rules (§2.4): the actions that change a message, where it
|
||||
/// goes, or refuse it, on outgoing and incoming mail.
|
||||
pub async fn transport(test: &mut TestServer) {
|
||||
println!("Running mail flow action tests...");
|
||||
let admin = test.account("[email protected]");
|
||||
let sender = admin
|
||||
.create_user_account(
|
||||
"[email protected]",
|
||||
"flow-sender-secret-6602",
|
||||
"Flow sender",
|
||||
&[],
|
||||
vec![],
|
||||
)
|
||||
.await;
|
||||
let other = admin
|
||||
.create_user_account(
|
||||
"[email protected]",
|
||||
"flow-other-secret-6603",
|
||||
"Flow other",
|
||||
&[],
|
||||
vec![],
|
||||
)
|
||||
.await;
|
||||
let (_, response) = call(
|
||||
&sender,
|
||||
"Identity/set",
|
||||
json!({"create": {"i": {"name": "Sender", "email": "[email protected]"}}}),
|
||||
)
|
||||
.await;
|
||||
let identity = response["created"]["i"]["id"].as_str().unwrap().to_string();
|
||||
let (_, response) = call(
|
||||
&sender,
|
||||
"Mailbox/set",
|
||||
json!({"create": {"m": {"name": "Flow drafts"}}}),
|
||||
)
|
||||
.await;
|
||||
let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string();
|
||||
|
||||
let (_, response) = call(
|
||||
&admin,
|
||||
"inbuxa:MailRule/set",
|
||||
json!({"create": {
|
||||
"d": {"name": "Footer and tag", "kind": "transport", "direction": "outgoing",
|
||||
"conditions": [{"type": "senderAddress", "addresses": ["[email protected]"]}],
|
||||
"actions": [
|
||||
{"type": "addDisclaimer", "text": "Sent by Example Co.", "position": "bottom"},
|
||||
{"type": "addHeader", "name": "X-Flow", "value": "checked"},
|
||||
{"type": "prefixSubject", "text": "[Example]"}
|
||||
]},
|
||||
"r": {"name": "Redirect projects", "kind": "transport", "direction": "outgoing",
|
||||
"conditions": [{"type": "words", "words": ["project falcon"]}],
|
||||
"actions": [{"type": "redirect", "addresses": ["[email protected]"]}]},
|
||||
"x": {"name": "No invoices", "kind": "transport", "direction": "outgoing",
|
||||
"conditions": [{"type": "words", "words": ["invoice-scam"]}],
|
||||
"actions": [{"type": "refuse", "text": "This kind of message isn't sent from here."}]},
|
||||
"i": {"name": "Outside banner", "kind": "transport", "direction": "incoming",
|
||||
"conditions": [], "actions": [{"type": "prefixSubject", "text": "[External]"}]}
|
||||
}}),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
response["created"].as_object().map(|c| c.len()),
|
||||
Some(4),
|
||||
"{response}"
|
||||
);
|
||||
|
||||
// Disclaimer, header and subject prefix on the sender's own copy
|
||||
let response = submit(
|
||||
&sender,
|
||||
&identity,
|
||||
&mailbox,
|
||||
&["[email protected]"],
|
||||
"Lunch",
|
||||
"See you at noon.",
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
assert!(response["created"].get("s").is_some(), "{response}");
|
||||
let got = received(&sender, "Lunch", "X-Flow", Some(&mailbox)).await;
|
||||
assert_eq!(got.len(), 1, "{got:?}");
|
||||
let (subject, body, header) = &got[0];
|
||||
assert_eq!(subject, "[Example] Lunch");
|
||||
assert!(
|
||||
body.starts_with("See you at noon.") && body.trim_end().ends_with("Sent by Example Co."),
|
||||
"{body:?}"
|
||||
);
|
||||
assert_eq!(header.trim(), "checked");
|
||||
|
||||
// Redirected: the other user gets it, the named recipient doesn't
|
||||
let response = submit(
|
||||
&sender,
|
||||
&identity,
|
||||
&mailbox,
|
||||
&["[email protected]"],
|
||||
"Falcon",
|
||||
"About Project Falcon.",
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
assert!(response["created"].get("s").is_some(), "{response}");
|
||||
assert_eq!(
|
||||
received(&other, "Falcon", "X-Flow", None).await.len(),
|
||||
1,
|
||||
"redirected"
|
||||
);
|
||||
assert!(
|
||||
received(&sender, "Falcon", "X-Flow", Some(&mailbox))
|
||||
.await
|
||||
.is_empty(),
|
||||
"the named recipient didn't get it"
|
||||
);
|
||||
|
||||
// Refused
|
||||
let response = submit(
|
||||
&sender,
|
||||
&identity,
|
||||
&mailbox,
|
||||
&["[email protected]"],
|
||||
"Pay",
|
||||
"invoice-scam inside",
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
let refused = &response["notCreated"]["s"];
|
||||
assert_eq!(refused["type"], "forbiddenToSend", "{response}");
|
||||
assert!(
|
||||
refused["description"]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.contains("isn't sent from here"),
|
||||
"{response}"
|
||||
);
|
||||
|
||||
// Incoming mail: the banner rule, and none of the outgoing ones
|
||||
// LMTP delivery without signing in, as mail from outside arrives
|
||||
admin
|
||||
.registry_create_object(MtaStageAuth {
|
||||
require: Expression {
|
||||
else_: "false".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
})
|
||||
.await;
|
||||
let mut lmtp = SmtpConnection::connect().await;
|
||||
lmtp.ingest(
|
||||
"[email protected]",
|
||||
&["[email protected]"],
|
||||
"From: [email protected]\r\nTo: [email protected]\r\nSubject: Hello from outside\r\n\r\nHi.\r\n",
|
||||
)
|
||||
.await;
|
||||
let got = received(&other, "Hello from outside", "X-Flow", None).await;
|
||||
assert_eq!(
|
||||
got.first().map(|g| g.0.as_str()),
|
||||
Some("[External] Hello from outside"),
|
||||
"{got:?}"
|
||||
);
|
||||
assert!(
|
||||
got[0].2.is_empty(),
|
||||
"outgoing rules left incoming mail alone"
|
||||
);
|
||||
|
||||
// The redirect and the refusal are recorded; the footer isn't
|
||||
let (_, response) = call(
|
||||
&admin,
|
||||
"inbuxa:AuditEvent/query",
|
||||
json!({"filter": {"targetKind": "message", "text": "Mail flow"}}),
|
||||
)
|
||||
.await;
|
||||
let ids = response["ids"].clone();
|
||||
let (_, response) = call(&admin, "inbuxa:AuditEvent/get", json!({"ids": ids})).await;
|
||||
let details: Vec<String> = response["list"]
|
||||
.as_array()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.filter_map(|e| e["details"].as_str().map(str::to_string))
|
||||
.collect();
|
||||
assert!(
|
||||
details.iter().any(|d| d.starts_with(
|
||||
"Mail flow rule \"Redirect projects\" redirected to [email protected]"
|
||||
)),
|
||||
"{details:?}"
|
||||
);
|
||||
assert!(
|
||||
details
|
||||
.iter()
|
||||
.any(|d| d.starts_with("Mail flow rule \"No invoices\" refused")),
|
||||
"{details:?}"
|
||||
);
|
||||
assert!(
|
||||
!details.iter().any(|d| d.contains("Footer and tag")),
|
||||
"{details:?}"
|
||||
);
|
||||
|
||||
// Off again
|
||||
let (_, response) = call(
|
||||
&admin,
|
||||
"inbuxa:MailRule/get",
|
||||
json!({"ids": null, "properties": ["id", "kind"]}),
|
||||
)
|
||||
.await;
|
||||
let transport: Vec<Value> = response["list"]
|
||||
.as_array()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.filter(|r| r["kind"] == "transport")
|
||||
.map(|r| r["id"].clone())
|
||||
.collect();
|
||||
call(&admin, "inbuxa:MailRule/set", json!({"destroy": transport})).await;
|
||||
}
|
||||
|
||||
#[ignore]
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
pub async fn mail_rules_tests() {
|
||||
@@ -515,6 +786,7 @@ pub async fn mail_rules_tests() {
|
||||
test.insert_account(admin);
|
||||
self::test(&mut test).await;
|
||||
self::dlp(&mut test).await;
|
||||
self::transport(&mut test).await;
|
||||
if test.is_reset() {
|
||||
test.temp_dir.delete();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user