From 7f22006e979a57326e35f76cb86a97042ada6104 Mon Sep 17 00:00:00 2001 From: John Coffey Date: Mon, 28 Sep 2026 18:08:40 -0700 Subject: [PATCH] Mail flow rules: carry out the transport actions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phase 2g of the DLP and mail flow rules spec: transport rules now act, on outgoing and incoming mail. - features/mailflow/rewrite.rs: add or remove a header, prefix or set the subject (an RFC 2047 word when not ASCII), add a disclaimer. A disclaimer edits the message's main text and HTML bodies only, each decoded, changed and written back as UTF-8 quoted-printable with its other headers kept, top or bottom (after or before in HTML); attachments and attached messages are left alone, and a disclaimer already present isn't added again. - smtp/inbound/mailflow.rs: the check runs for incoming mail too (transport rules only; DLP stays outgoing). After DLP passes, each matched transport rule's actions run in order: message edits, add-recipient and redirect (envelope changes DATA applies), route (a per-message queue ahead of the queue strategy), refuse (550 5.7.1 with the rule's text). The override tag is stripped with the same subject writer, so a non-ASCII subject stays valid. - Audit: refusals and changes to where mail goes are recorded (sender, or system:mail-flow for incoming mail); wording and header changes aren't, or a banner rule would record every message (spec §2.7). Tests: rewrite unit tests (headers, encoded subjects, disclaimers on a single part and on multipart/alternative with an attachment, once only); mail_rules_tests gains the actions end to end: disclaimer, header and subject prefix on a delivered message, a redirect, a refusal, a banner on incoming LMTP mail that outgoing rules leave alone, and which of those are audited. --- Cargo.lock | 2 + crates/features/Cargo.toml | 2 + crates/features/src/mailflow/mod.rs | 4 +- crates/features/src/mailflow/rewrite.rs | 306 ++++++++++++++++++ crates/smtp/src/core/mod.rs | 4 + crates/smtp/src/inbound/data.rs | 33 +- crates/smtp/src/inbound/mailflow.rs | 194 ++++++++--- docs/spec/features/dlp-and-mail-flow-rules.md | 9 +- tests/src/system/mail_rules.rs | 274 +++++++++++++++- 9 files changed, 780 insertions(+), 48 deletions(-) create mode 100644 crates/features/src/mailflow/rewrite.rs diff --git a/Cargo.lock b/Cargo.lock index d14ec2d..1c26a28 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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", diff --git a/crates/features/Cargo.toml b/crates/features/Cargo.toml index 3054682..4051c50 100644 --- a/crates/features/Cargo.toml +++ b/crates/features/Cargo.toml @@ -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"] } diff --git a/crates/features/src/mailflow/mod.rs b/crates/features/src/mailflow/mod.rs index 1516faa..d07dc1f 100644 --- a/crates/features/src/mailflow/mod.rs +++ b/crates/features/src/mailflow/mod.rs @@ -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; diff --git a/crates/features/src/mailflow/rewrite.rs b/crates/features/src/mailflow/rewrite.rs new file mode 100644 index 0000000..a99c93b --- /dev/null +++ b/crates/features/src/mailflow/rewrite.rs @@ -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 { + 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> { + 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 { + 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> { + 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', "
\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("').map(|end| at + end + 1)) + { + Some(at) => format!("{}{html}{}", &body[..at], &body[at..]), + None => format!("{html}{body}"), + }, + Position::Bottom => match lower.rfind("") { + 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> { + let parsed = MessageParser::new().parse(message)?; + let html = html + .map(str::to_string) + .unwrap_or_else(|| format!("

{}

", escape_html(text.trim()))); + let marker = text.trim(); + + let mut body_parts: Vec = 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)> = 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: a@example.com\r\nTo: b@elsewhere.org\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: a@example.com\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: a@example.com\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("b@elsewhere.org") + ); + // 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("

Confidential.

"), + 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("

Confidential.

See you.

"), + "{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")); + } +} diff --git a/crates/smtp/src/core/mod.rs b/crates/smtp/src/core/mod.rs index 866e6cc..dedad18 100644 --- a/crates/smtp/src/core/mod.rs +++ b/crates/smtp/src/core/mod.rs @@ -100,6 +100,8 @@ pub struct SessionData { // message, for the submission to report pub dlp_override: Option, pub dlp_refusal: Option, + // inbuxa: a mail flow rule's route for this message + pub mailflow_queue: Option, } /// 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, } } } diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index f6970ca..4d13fc0 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -745,7 +745,26 @@ impl Session { .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 Session { }; // 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::( &self.server.core.smtp.queue.queue, @@ -927,8 +948,10 @@ impl Session { ) .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); diff --git a/crates/smtp/src/inbound/mailflow.rs b/crates/smtp/src/inbound/mailflow.rs index 7229a1f..a84704b 100644 --- a/crates/smtp/src/inbound/mailflow.rs +++ b/crates/smtp/src/inbound/mailflow.rs @@ -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), + /// Go on, with a changed message (the override tag taken out, a + /// disclaimer, headers, the subject) and envelope. + Changed { + message: Option>, + envelope: Vec, + }, /// Refuse, with this SMTP reply, and for a JMAP submission, why. Refuse(Vec, Option), } +/// What a transport rule changes about where a message goes. +pub enum EnvelopeChange { + AddRecipient(String), + Redirect(Vec), + 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 Session { /// 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 Session { ); } }; - if !rules.applies_to(true) { + if !rules.applies_to(outgoing) { return Checked::Accept; } @@ -197,7 +212,12 @@ impl Session { .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 Session { .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 Session { let domains = domains.join(", "); drop(envelope); - self.record_dlp( - account_id, - &account, - &outcome, - &decision, - override_reason.as_deref(), - &domains, - ) - .await; + if let Some((account_id, account)) = &sender { + self.record_dlp( + *account_id, + account, + &outcome, + &decision, + override_reason.as_deref(), + &domains, + ) + .await; + } match decision { Decision::Block(rules) | Decision::Hold { rules, .. } => { @@ -310,27 +334,119 @@ impl Session { .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> = + 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 = 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)>, + 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, /// where to, which rules and each detector's count, what happened, and /// an override's reason. Never the matched text. diff --git a/docs/spec/features/dlp-and-mail-flow-rules.md b/docs/spec/features/dlp-and-mail-flow-rules.md index b017619..00489dd 100644 --- a/docs/spec/features/dlp-and-mail-flow-rules.md +++ b/docs/spec/features/dlp-and-mail-flow-rules.md @@ -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 diff --git a/tests/src/system/mail_rules.rs b/tests/src/system/mail_rules.rs index 5f69cb1..3eae538 100644 --- a/tests/src/system/mail_rules.rs +++ b/tests/src/system/mail_rules.rs @@ -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("admin@example.com"); + let sender = admin + .create_user_account( + "flow-sender@example.com", + "flow-sender-secret-6602", + "Flow sender", + &[], + vec![], + ) + .await; + let other = admin + .create_user_account( + "flow-other@example.com", + "flow-other-secret-6603", + "Flow other", + &[], + vec![], + ) + .await; + let (_, response) = call( + &sender, + "Identity/set", + json!({"create": {"i": {"name": "Sender", "email": "flow-sender@example.com"}}}), + ) + .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": ["flow-sender@example.com"]}], + "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": ["flow-other@example.com"]}]}, + "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, + &["flow-sender@example.com"], + "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, + &["flow-sender@example.com"], + "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, + &["flow-sender@example.com"], + "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( + "someone@elsewhere.org", + &["flow-other@example.com"], + "From: someone@elsewhere.org\r\nTo: flow-other@example.com\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 = 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 flow-other@example.com" + )), + "{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 = 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(); }