Mail flow rules: carry out the transport actions #106

Merged
jcoffey-dev merged 1 commits from feature/mailflow-actions into main 2026-09-29 01:17:04 +00:00
9 changed files with 780 additions and 48 deletions
Generated
+2
View File
@@ -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",
+2
View File
@@ -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"] }
+3 -1
View File
@@ -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;
+306
View File
@@ -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('&', "&amp;")
.replace('<', "&lt;")
.replace('>', "&gt;")
.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"));
}
}
+4
View File
@@ -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,
}
}
}
+28 -5
View File
@@ -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);
+148 -32
View File
@@ -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
+273 -1
View File
@@ -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();
}