Mail flow rules: carry out the transport actions
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 <body> or before </body> 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.
This commit is contained in:
@@ -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);
|
||||
|
||||
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<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,
|
||||
/// where to, which rules and each detector's count, what happened, and
|
||||
/// an override's reason. Never the matched text.
|
||||
|
||||
Reference in New Issue
Block a user