DLP at DATA: block, warn and override over SMTP and JMAP
Phase 2f of the DLP and mail flow rules spec: the rules now run on mail
an authenticated sender submits, after the DATA system script and
before headers and DKIM signing (§2.1).
- smtp/inbound/mailflow.rs: builds what the rules look at from the
message (subject, the text version of each body, one level of attached
messages, attachment text via the extractor, 10 MB of text at most)
and the envelope (sender's groups and tenant; each recipient local or
not, and its groups). Skipped entirely when no enabled rule applies to
outgoing mail. Rules that can't be loaded refuse with a 451: nothing
unchecked leaves.
- Block: 550 5.7.1 with the rule's notice. Warn: 550 5.7.1 with the
notice and how to override: "[override: reason]" at the start of the
subject, taken out before the message goes on (settled answer 1).
Until phase 3, a hold rule blocks rather than let mail through.
- JMAP: EmailSubmission takes inbuxa:dlpOverride {reason}; a refusal
comes back as inbuxa:dlpWarning or inbuxa:dlpBlocked with each rule's
name and notice (description too, for older clients).
- Audit: one record per DLP match, the sender as actor, action create,
target a message: the recipient domains, each rule with its detectors'
counts, the outcome, an override's reason. Never the matched text. No
new audit action: an older node that meets one fails its daily
clean-up, which would make rolling back unsafe (spec §2.7 updated).
Tests: mail_rules_tests gains the DLP flow over JMAP (no rules, warning
with rule and notice, local recipient not warned, override with a
reason, block that no reason passes, the subject tag stripped from the
delivered message, audit records with no card or key text). smtp
inbound tests pass; system_tests passed twice after one timeout in the
email delivery tests that didn't recur.
This commit is contained in:
@@ -2,6 +2,8 @@
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||
*
|
||||
* Modified by Coffey Labs in 2026 for INBUXA.
|
||||
*/
|
||||
|
||||
use crate::{inbound::auth::SaslToken, queue::QueueId};
|
||||
@@ -92,6 +94,20 @@ pub struct SessionData {
|
||||
pub spf_ehlo: Option<SpfOutput>,
|
||||
pub spf_mail_from: Option<SpfOutput>,
|
||||
pub dnsbl_error: Option<Vec<u8>>,
|
||||
|
||||
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.5): the reason a JMAP
|
||||
// sender gave to send despite a warning, and why DATA refused a
|
||||
// message, for the submission to report
|
||||
pub dlp_override: Option<String>,
|
||||
pub dlp_refusal: Option<DlpRefusal>,
|
||||
}
|
||||
|
||||
/// inbuxa: a DATA refusal by DLP rules: blocked, or a warning the sender
|
||||
/// may override, with each rule's name and notice.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct DlpRefusal {
|
||||
pub blocked: bool,
|
||||
pub rules: Vec<(String, String)>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
@@ -168,6 +184,8 @@ impl SessionData {
|
||||
spf_ehlo: None,
|
||||
spf_mail_from: None,
|
||||
dnsbl_error: None,
|
||||
dlp_override: None,
|
||||
dlp_refusal: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -291,6 +309,8 @@ impl SessionData {
|
||||
spf_ehlo: None,
|
||||
spf_mail_from: None,
|
||||
dnsbl_error: None,
|
||||
dlp_override: None,
|
||||
dlp_refusal: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -738,6 +738,20 @@ impl<T: SessionStream> Session<T> {
|
||||
}
|
||||
}
|
||||
|
||||
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.1): after the system
|
||||
// script, before headers and signing
|
||||
match self
|
||||
.check_mail_rules(edited_message.as_deref().unwrap_or(raw_message.as_slice()))
|
||||
.await
|
||||
{
|
||||
super::mailflow::Checked::Accept => {}
|
||||
super::mailflow::Checked::Replace(message) => edited_message = Some(message),
|
||||
super::mailflow::Checked::Refuse(reply, refusal) => {
|
||||
self.data.dlp_refusal = refusal;
|
||||
return reply.into();
|
||||
}
|
||||
}
|
||||
|
||||
// Build message
|
||||
let mail_from = self.data.mail_from.clone().unwrap();
|
||||
let rcpt_to = std::mem::take(&mut self.data.rcpt_to);
|
||||
|
||||
@@ -0,0 +1,427 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2026 Coffey Labs
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only
|
||||
*/
|
||||
|
||||
//! inbuxa: DLP at DATA (dlp-and-mail-flow-rules spec, §2.1, §2.4–§2.7).
|
||||
//!
|
||||
//! Runs after the DATA system script and before headers and DKIM signing,
|
||||
//! on mail an authenticated sender submits over SMTP or JMAP. The rules and
|
||||
//! the detectors are `inbuxa_features::mailflow`; this is the glue: build
|
||||
//! what they look at from the message, apply the decision, record it.
|
||||
|
||||
use crate::core::{DlpRefusal, Session};
|
||||
use common::network::SessionStream;
|
||||
use inbuxa_features::{
|
||||
audit::{Action, Actor, Outcome, Record, Target},
|
||||
mailflow::{
|
||||
cache,
|
||||
engine::{
|
||||
Attachment, Content, Decision, Envelope, Outcome as RulesOutcome, Recipient, RuleRef,
|
||||
},
|
||||
extract::{self, Extracted, Limits},
|
||||
rules::Kind,
|
||||
},
|
||||
};
|
||||
use mail_parser::{HeaderName, Message, MessageParser, MimeHeaders, PartType};
|
||||
use std::{borrow::Cow, time::SystemTime};
|
||||
|
||||
/// How much text one message is read for; past it, the rest counts as
|
||||
/// "can't be inspected" (§2.3).
|
||||
const INSPECTION_LIMIT: usize = 10 * 1024 * 1024;
|
||||
|
||||
/// What the check decided.
|
||||
pub enum Checked {
|
||||
/// Go on, with the message unchanged.
|
||||
Accept,
|
||||
/// Go on with this message instead (the override tag taken out).
|
||||
Replace(Vec<u8>),
|
||||
/// Refuse, with this SMTP reply, and for a JMAP submission, why.
|
||||
Refuse(Vec<u8>, Option<DlpRefusal>),
|
||||
}
|
||||
|
||||
/// `[override: reason]` at the start of a subject: the reason, and the
|
||||
/// subject without it.
|
||||
pub fn override_tag(subject: &str) -> Option<(String, String)> {
|
||||
let trimmed = subject.trim_start();
|
||||
let head = trimmed.get(..10)?;
|
||||
if !head.eq_ignore_ascii_case("[override:") {
|
||||
return None;
|
||||
}
|
||||
let close = trimmed.find(']')?;
|
||||
let reason = trimmed[10..close].trim();
|
||||
if reason.is_empty() {
|
||||
return None;
|
||||
}
|
||||
Some((
|
||||
reason.chars().take(500).collect(),
|
||||
trimmed[close + 1..].trim_start().to_string(),
|
||||
))
|
||||
}
|
||||
|
||||
/// One line of an SMTP reply: no line breaks, a sane length.
|
||||
fn reply_text(text: &str) -> String {
|
||||
text.split_whitespace()
|
||||
.collect::<Vec<_>>()
|
||||
.join(" ")
|
||||
.chars()
|
||||
.take(400)
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn notices(rules: &[RuleRef]) -> String {
|
||||
let mut seen = Vec::new();
|
||||
for rule in rules {
|
||||
let notice = reply_text(&rule.notice);
|
||||
if !seen.contains(¬ice) {
|
||||
seen.push(notice);
|
||||
}
|
||||
}
|
||||
seen.join(" ")
|
||||
}
|
||||
|
||||
fn refusal(blocked: bool, rules: &[RuleRef]) -> DlpRefusal {
|
||||
DlpRefusal {
|
||||
blocked,
|
||||
rules: rules
|
||||
.iter()
|
||||
.map(|r| (r.name.clone(), r.notice.clone()))
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
|
||||
/// A message and the messages attached to it, one level down.
|
||||
fn collect<'x>(message: &'x Message<'x>, content: &mut Content<'x>, budget: &mut usize, depth: u8) {
|
||||
let add = |text: Cow<'x, str>, content: &mut Content<'x>, budget: &mut usize| {
|
||||
if *budget == 0 {
|
||||
content.truncated = true;
|
||||
return;
|
||||
}
|
||||
if text.len() > *budget {
|
||||
let mut cut = *budget;
|
||||
while !text.is_char_boundary(cut) {
|
||||
cut -= 1;
|
||||
}
|
||||
content.bodies.push(Cow::Owned(text[..cut].to_string()));
|
||||
content.truncated = true;
|
||||
*budget = 0;
|
||||
} else {
|
||||
*budget -= text.len();
|
||||
content.bodies.push(text);
|
||||
}
|
||||
};
|
||||
// The text version of each body (an HTML-only one converted), not both
|
||||
// versions of the same alternative, so words aren't counted twice
|
||||
for part in message.text_bodies() {
|
||||
match &part.body {
|
||||
PartType::Text(text) => add(Cow::Borrowed(text.as_ref()), content, budget),
|
||||
PartType::Html(html) => add(
|
||||
Cow::Owned(mail_parser::decoders::html::html_to_text(html)),
|
||||
content,
|
||||
budget,
|
||||
),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
for part in message.attachments() {
|
||||
if let (Some(inner), true) = (part.message(), depth == 0) {
|
||||
if let Some(subject) = inner.subject() {
|
||||
add(Cow::Borrowed(subject), content, budget);
|
||||
}
|
||||
collect(inner, content, budget, depth + 1);
|
||||
continue;
|
||||
}
|
||||
let content_type = part
|
||||
.content_type()
|
||||
.map(|ct| match ct.subtype() {
|
||||
Some(sub) => format!("{}/{}", ct.ctype(), sub),
|
||||
None => ct.ctype().to_string(),
|
||||
})
|
||||
.unwrap_or_default();
|
||||
let bytes = part.contents();
|
||||
let mut extracted = extract::extract(
|
||||
&content_type,
|
||||
part.attachment_name(),
|
||||
bytes,
|
||||
&Limits::default(),
|
||||
);
|
||||
if let Extracted::Text(text) = &extracted {
|
||||
if text.len() > *budget {
|
||||
extracted = Extracted::NotInspectable(extract::Why::TooLarge);
|
||||
} else {
|
||||
*budget -= text.len();
|
||||
}
|
||||
}
|
||||
content.attachments.push(Attachment {
|
||||
name: part.attachment_name(),
|
||||
content_type: content_type.into(),
|
||||
size: bytes.len() as u64,
|
||||
extracted,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
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());
|
||||
let rules = match cache::compiled(self.server.store()).await {
|
||||
Ok(rules) => rules,
|
||||
Err(err) => {
|
||||
trc::error!(
|
||||
err.span_id(self.data.session_id)
|
||||
.caused_by(trc::location!())
|
||||
.details("Failed to load mail rules")
|
||||
);
|
||||
// Fail closed: a message nobody could check doesn't leave
|
||||
return Checked::Refuse(
|
||||
b"451 4.3.0 This message couldn't be checked against the server's rules. Try again later.\r\n"
|
||||
.to_vec(),
|
||||
None,
|
||||
);
|
||||
}
|
||||
};
|
||||
if !rules.applies_to(true) {
|
||||
return Checked::Accept;
|
||||
}
|
||||
|
||||
let parsed = MessageParser::new().parse(message);
|
||||
let subject = parsed
|
||||
.as_ref()
|
||||
.and_then(|m| m.subject())
|
||||
.unwrap_or_default();
|
||||
let jmap_override = self.data.dlp_override.clone();
|
||||
let tag = override_tag(subject);
|
||||
let checked_subject = tag.as_ref().map_or(subject, |(_, rest)| rest.as_str());
|
||||
|
||||
let mut content = Content {
|
||||
subject: checked_subject,
|
||||
size: message.len() as u64,
|
||||
..Default::default()
|
||||
};
|
||||
let mut budget = INSPECTION_LIMIT;
|
||||
match &parsed {
|
||||
Some(parsed) => {
|
||||
content.headers = parsed
|
||||
.headers()
|
||||
.iter()
|
||||
.filter_map(|h| h.value.as_text().map(|v| (h.name.as_str(), v)))
|
||||
.collect();
|
||||
collect(parsed, &mut content, &mut budget, 0);
|
||||
}
|
||||
// Nothing a rule could read: say so, rather than pass it
|
||||
None => content.truncated = true,
|
||||
}
|
||||
let mut recipient_groups = Vec::with_capacity(self.data.rcpt_to.len());
|
||||
for rcpt in &self.data.rcpt_to {
|
||||
let local = self
|
||||
.server
|
||||
.domain(&rcpt.domain)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.is_some();
|
||||
let groups = if local {
|
||||
match self
|
||||
.server
|
||||
.account_id_from_email(&rcpt.address_lcase, false)
|
||||
.await
|
||||
{
|
||||
Ok(Some(id)) => self
|
||||
.server
|
||||
.account(id)
|
||||
.await
|
||||
.map(|a| a.id_member_of.to_vec())
|
||||
.unwrap_or_default(),
|
||||
_ => Vec::new(),
|
||||
}
|
||||
} else {
|
||||
Vec::new()
|
||||
};
|
||||
recipient_groups.push((local, groups));
|
||||
}
|
||||
let sender_address = self
|
||||
.data
|
||||
.mail_from
|
||||
.as_ref()
|
||||
.map(|m| m.address_lcase.clone())
|
||||
.unwrap_or_default();
|
||||
let envelope = Envelope {
|
||||
outgoing: true,
|
||||
sender: &sender_address,
|
||||
sender_groups: &account.id_member_of,
|
||||
sender_tenant: account.id_tenant,
|
||||
recipients: self
|
||||
.data
|
||||
.rcpt_to
|
||||
.iter()
|
||||
.zip(&recipient_groups)
|
||||
.map(|(rcpt, (local, groups))| Recipient {
|
||||
address: &rcpt.address_lcase,
|
||||
local: *local,
|
||||
groups,
|
||||
})
|
||||
.collect(),
|
||||
};
|
||||
|
||||
let outcome = rules.evaluate(&envelope, &content);
|
||||
let override_reason =
|
||||
jmap_override.or_else(|| tag.as_ref().map(|(reason, _)| reason.clone()));
|
||||
let decision = outcome.decision(override_reason.is_some());
|
||||
let mut domains: Vec<&str> = self
|
||||
.data
|
||||
.rcpt_to
|
||||
.iter()
|
||||
.map(|r| r.domain.as_str())
|
||||
.collect();
|
||||
domains.sort_unstable();
|
||||
domains.dedup();
|
||||
let domains = domains.join(", ");
|
||||
drop(envelope);
|
||||
|
||||
self.record_dlp(
|
||||
account_id,
|
||||
&account,
|
||||
&outcome,
|
||||
&decision,
|
||||
override_reason.as_deref(),
|
||||
&domains,
|
||||
)
|
||||
.await;
|
||||
|
||||
match decision {
|
||||
Decision::Block(rules) | Decision::Hold { rules, .. } => {
|
||||
// Hold for review is phase 3: until then a hold rule blocks,
|
||||
// rather than let the message through unreviewed
|
||||
let refusal = refusal(true, &rules);
|
||||
Checked::Refuse(format!("550 5.7.1 {}\r\n", notices(&rules)).into_bytes(), Some(refusal))
|
||||
}
|
||||
Decision::Warn(rules) => Checked::Refuse(
|
||||
format!(
|
||||
"550 5.7.1 {} To send anyway, start the subject with [override: your reason]\r\n",
|
||||
notices(&rules)
|
||||
)
|
||||
.into_bytes(),
|
||||
Some(refusal(false, &rules)),
|
||||
),
|
||||
Decision::Pass => match (tag, &parsed) {
|
||||
// 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,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// 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.
|
||||
async fn record_dlp(
|
||||
&self,
|
||||
account_id: u32,
|
||||
account: &common::auth::AccountCache,
|
||||
outcome: &RulesOutcome,
|
||||
decision: &Decision,
|
||||
override_reason: Option<&str>,
|
||||
domains: &str,
|
||||
) {
|
||||
let dlp: Vec<_> = outcome
|
||||
.matched
|
||||
.iter()
|
||||
.filter(|m| m.kind == Kind::Dlp)
|
||||
.collect();
|
||||
if dlp.is_empty() {
|
||||
return;
|
||||
}
|
||||
let rules = dlp
|
||||
.iter()
|
||||
.map(|m| {
|
||||
let counts = m
|
||||
.counts
|
||||
.iter()
|
||||
.map(|(id, n)| format!("{id} {n}"))
|
||||
.collect::<Vec<_>>()
|
||||
.join(", ");
|
||||
if counts.is_empty() {
|
||||
format!("\"{}\"", m.name)
|
||||
} else {
|
||||
format!("\"{}\" ({counts})", m.name)
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
.join("; ");
|
||||
let (what, outcome, reason) = match decision {
|
||||
Decision::Block(_) | Decision::Hold { .. } => {
|
||||
("blocked", Outcome::refused("inbuxa:dlpBlocked", None), None)
|
||||
}
|
||||
Decision::Warn(_) => ("warned", Outcome::refused("inbuxa:dlpWarning", None), None),
|
||||
Decision::Pass => (
|
||||
"sent after a warning",
|
||||
Outcome::success(),
|
||||
override_reason.map(str::to_string),
|
||||
),
|
||||
};
|
||||
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: Actor::account(account_id, account.name.to_string(), account.id_tenant),
|
||||
via: None,
|
||||
remote_ip: Some(self.data.remote_ip),
|
||||
action: Action::Create,
|
||||
target: Target {
|
||||
kind: "message".into(),
|
||||
id: None,
|
||||
name: None,
|
||||
account_id: Some(account_id),
|
||||
tenant_id: account.id_tenant,
|
||||
},
|
||||
changes: vec![],
|
||||
details: Some(format!("DLP {what}, to {domains}: {rules}")),
|
||||
reason,
|
||||
outcome,
|
||||
})
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::override_tag;
|
||||
|
||||
#[test]
|
||||
fn override_tags() {
|
||||
assert_eq!(
|
||||
override_tag("[override: client asked for it] Card details"),
|
||||
Some(("client asked for it".into(), "Card details".into()))
|
||||
);
|
||||
assert_eq!(
|
||||
override_tag(" [OVERRIDE:yes]x"),
|
||||
Some(("yes".into(), "x".into()))
|
||||
);
|
||||
assert_eq!(override_tag("[override: ] x"), None);
|
||||
assert_eq!(override_tag("Re: [override: no] x"), None);
|
||||
assert_eq!(override_tag("[override: unclosed"), None);
|
||||
assert_eq!(override_tag(""), None);
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,8 @@
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||
*
|
||||
* Modified by Coffey Labs in 2026 for INBUXA.
|
||||
*/
|
||||
|
||||
use mail_auth::{DkimResult, DmarcResult, IprevResult, SpfResult, dmarc::Policy};
|
||||
@@ -13,6 +15,7 @@ pub mod dkim;
|
||||
pub mod ehlo;
|
||||
pub mod hooks;
|
||||
pub mod mail;
|
||||
pub mod mailflow; // inbuxa: DLP and mail flow rules
|
||||
pub mod milter;
|
||||
pub mod rcpt;
|
||||
pub mod session;
|
||||
|
||||
Reference in New Issue
Block a user