Merge pull request 'DLP at DATA: block, warn and override over SMTP and JMAP' (#104) from feature/dlp-data-stage into main
ci / fork-checks (push) Successful in 15s
ci / build (push) Canceled after 16s

This commit was merged in pull request #104.
This commit is contained in:
2026-09-29 01:16:45 +00:00
10 changed files with 899 additions and 22 deletions
+5 -3
View File
@@ -46,7 +46,7 @@ pub struct Recipient<'a> {
pub struct Attachment<'a> {
pub name: Option<&'a str>,
/// Declared type, or detected where the caller knows better.
pub content_type: &'a str,
pub content_type: Cow<'a, str>,
pub size: u64,
pub extracted: Extracted,
}
@@ -606,13 +606,15 @@ mod tests {
attachments: vec![
Attachment {
name: Some("plan.docx"),
content_type: "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
content_type:
"application/vnd.openxmlformats-officedocument.wordprocessingml.document"
.into(),
size: 40_000,
extracted: Extracted::Text("IBAN GB29 NWBK 6016 1331 9268 19".into()),
},
Attachment {
name: Some("scan.pdf"),
content_type: "application/pdf",
content_type: "application/pdf".into(),
size: 900_000,
extracted: Extracted::NotInspectable(Why::Pdf),
},
+28
View File
@@ -47,6 +47,18 @@ struct SetErrorInner<P: Property> {
#[serde(skip_serializing_if = "Vec::is_empty")]
#[serde(rename = "validationErrors")]
validation_errors: Vec<ValidationError>,
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.5): each rule that
// warned or blocked, with its notice
#[serde(skip_serializing_if = "Vec::is_empty")]
rules: Vec<DlpRule>,
}
/// inbuxa: a DLP rule named in an `inbuxa:dlpWarning` or `inbuxa:dlpBlocked`.
#[derive(Debug, Clone, serde::Serialize)]
pub struct DlpRule {
pub name: String,
pub notice: String,
}
#[derive(Debug, Clone)]
@@ -127,6 +139,12 @@ pub enum SetErrorType {
// inbuxa: a create that couldn't run (ai-explain spec: busy, timeout, …)
#[serde(rename = "serverFail")]
ServerFail,
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.5): a warning the
// sender may answer with inbuxa:dlpOverride, and a block
#[serde(rename = "inbuxa:dlpWarning")]
DlpWarning,
#[serde(rename = "inbuxa:dlpBlocked")]
DlpBlocked,
}
impl SetErrorType {
@@ -166,6 +184,8 @@ impl SetErrorType {
SetErrorType::PrimaryKeyViolation => "primaryKeyViolation",
SetErrorType::ValidationFailed => "validationFailed",
SetErrorType::ServerFail => "serverFail",
SetErrorType::DlpWarning => "inbuxa:dlpWarning",
SetErrorType::DlpBlocked => "inbuxa:dlpBlocked",
}
}
}
@@ -180,9 +200,16 @@ impl<T: Property> SetError<T> {
object_id: None,
linked_objects: Vec::new(),
validation_errors: Vec::new(),
rules: Vec::new(),
}))
}
/// inbuxa: the DLP rules behind a warning or block.
pub fn with_dlp_rules(mut self, rules: Vec<DlpRule>) -> Self {
self.0.rules = rules;
self
}
pub fn with_description(mut self, description: impl Into<Cow<'static, str>>) -> Self {
self.0.description = description.into().into();
self
@@ -353,6 +380,7 @@ impl From<PatchError> for SetError<registry::schema::properties::Property> {
object_id: None,
linked_objects: Vec::new(),
validation_errors: Vec::new(),
rules: Vec::new(),
}))
}
}
@@ -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::{
@@ -40,6 +42,9 @@ pub enum EmailSubmissionProperty {
Displayed,
DsnBlobIds,
MdnBlobIds,
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.5): `{"reason": ...}`
// to send despite a warning
DlpOverride,
Pointer(JsonPointer<EmailSubmissionProperty>),
}
@@ -90,6 +95,7 @@ impl Property for EmailSubmissionProperty {
EmailSubmissionProperty::Id => "id",
EmailSubmissionProperty::IdentityId => "identityId",
EmailSubmissionProperty::MdnBlobIds => "mdnBlobIds",
EmailSubmissionProperty::DlpOverride => "inbuxa:dlpOverride",
EmailSubmissionProperty::SendAt => "sendAt",
EmailSubmissionProperty::ThreadId => "threadId",
EmailSubmissionProperty::UndoStatus => "undoStatus",
@@ -181,6 +187,7 @@ impl EmailSubmissionProperty {
"displayed" => EmailSubmissionProperty::Displayed,
"dsnBlobIds" => EmailSubmissionProperty::DsnBlobIds,
"mdnBlobIds" => EmailSubmissionProperty::MdnBlobIds,
"inbuxa:dlpOverride" => EmailSubmissionProperty::DlpOverride,
)
.or_else(|| {
if allow_patch && value.contains('/') {
+46 -1
View File
@@ -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 common::{
@@ -16,7 +18,7 @@ use email::{
submission::{Address, Delivered, DeliveryStatus, EmailSubmission, UndoStatus},
};
use jmap_proto::{
error::set::{SetError, SetErrorType},
error::set::{DlpRule, SetError, SetErrorType},
method::set::{SetRequest, SetResponse},
object::email_submission::{self, EmailSubmissionProperty, EmailSubmissionValue},
references::resolve::ResolveCreatedReference,
@@ -379,6 +381,8 @@ impl EmailSubmissionSet for Server {
};
let mut mail_from: Option<MailFrom<Cow<'_, str>>> = None;
let mut rcpt_to: Vec<RcptTo<Cow<'_, str>>> = Vec::new();
// inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.5)
let mut dlp_override: Option<String> = None;
for (property, mut value) in object.into_expanded_object() {
if let Err(err) = response.resolve_self_references(&mut value, 0, false) {
@@ -493,6 +497,25 @@ impl EmailSubmissionSet for Server {
(Key::Property(EmailSubmissionProperty::UndoStatus), Value::Element(_)) => {
continue;
}
// inbuxa: the sender's reason to send despite a DLP warning
(Key::Property(EmailSubmissionProperty::DlpOverride), Value::Object(value)) => {
let reason = value
.iter()
.find(|(key, _)| key.to_string() == "reason")
.and_then(|(_, value)| value.as_str().map(|r| r.trim().to_string()))
.filter(|r| !r.is_empty());
match reason {
Some(reason) => dlp_override = Some(reason.chars().take(500).collect()),
None => {
return Ok(Err(SetError::invalid_properties()
.with_property(EmailSubmissionProperty::DlpOverride)
.with_description("An override needs a reason.")));
}
}
}
(Key::Property(EmailSubmissionProperty::DlpOverride), Value::Null) => {
continue;
}
_ => {
return Ok(Err(SetError::invalid_properties()
.with_property(property.into_owned())
@@ -700,6 +723,7 @@ impl EmailSubmissionSet for Server {
0,
),
);
session.data.dlp_override = dlp_override;
// Spawn SMTP session to avoid overflowing the stack
let handle = tokio::spawn(async move {
@@ -730,6 +754,27 @@ impl EmailSubmissionSet for Server {
let response = session.queue_message().await;
if let smtp::core::State::Accepted(queue_id) = session.state {
Ok((responses, Some(queue_id)))
} else if let Some(refusal) = session.data.dlp_refusal.take() {
// inbuxa: DLP (§2.5): which rules, and what they say
let description = refusal
.rules
.iter()
.map(|(_, notice)| notice.as_str())
.collect::<Vec<_>>()
.join(" ");
Err(SetError::new(if refusal.blocked {
SetErrorType::DlpBlocked
} else {
SetErrorType::DlpWarning
})
.with_description(description)
.with_dlp_rules(
refusal
.rules
.into_iter()
.map(|(name, notice)| DlpRule { name, notice })
.collect(),
))
} else {
Err(
SetError::new(SetErrorType::ForbiddenToSend).with_description(format!(
+20
View File
@@ -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,
}
}
}
+14
View File
@@ -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);
+427
View File
@@ -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(&notice) {
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);
}
}
+3
View File
@@ -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;