Files
inbuxa-server/crates/smtp/src/reporting/dmarc.rs
T
jcoffey-dev 5dde9793eb
ci / fork-checks (pull_request) Successful in 43s
ci / build (pull_request) Successful in 17m25s
Every node records DMARC and TLS results for the aggregate reports
The report scheduler dropped DMARC and TLS events on a node whose role
lacks outboundMta (upstream never started it there, so they sat in a
channel nobody read). Mail received on a front node therefore never
reached an aggregate report, which is meant to cover all of a domain's
inbound mail, whichever node received it. In rehearsal, five messages
received on port 25 on a front node were missing from every report.

- The report scheduler records on every node. Recording is a store write
  the nodes already share, so it needs nothing from the outbound MTA.
  Building and sending a report (the DmarcReport and TlsReport tasks) stay
  with outboundMta nodes, as the task manager already enforces.
- More nodes now append to one report at once. Appends already guard the
  report's versioned primary key; a write that loses now retries up to ten
  times after a short random pause, not three times at once.
- The node sending a report deletes it only if it is unchanged since it
  was read, and reads it again otherwise, so a record another node appends
  meanwhile goes out with the report instead of being deleted unsent.

Test: cluster::front_reports (PostgreSQL and MySQL). A front node's
results appear in the report the MTA node sends, alongside eight appended
at once from both nodes, and the front node never runs the report task.
It fails on main: the front node's results are never recorded.
2026-09-24 18:28:00 -07:00

713 lines
27 KiB
Rust

/*
* 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 super::AggregateTimestamp;
use super::shared::{MAX_WRITE_RETRIES, Revisioned, write_retry_pause};
use crate::{
core::Session,
queue::RecipientDomain,
reporting::{index::InternalReportIndex, send::MtaReportSend},
};
use common::{
Server,
config::smtp::report::AggregateFrequency,
ipc::{DmarcEvent, ToHash},
network::SessionStream,
};
use compact_str::ToCompactString;
use mail_auth::{
ArcOutput, AuthenticatedMessage, AuthenticationResults, DkimOutput, DkimResult, DmarcOutput,
DmarcResult, SpfResult,
common::verify::VerifySignature,
dkim2::Dkim2Output,
dmarc::{self},
report::{AuthFailureType, IdentityAlignment, PolicyPublished, Record, SPFDomainScope},
};
use registry::{
schema::{
enums::FailureReportingOption,
prelude::{ObjectType, Property},
structs::{DmarcInternalReport, DmarcReport, DmarcReportRecord, Rate},
},
types::{EnumImpl, ObjectImpl, datetime::UTCDateTime, map::Map},
};
use std::{borrow::Cow, future::Future};
use store::{
SerializeInfallible, U64_LEN, ValueKey,
registry::ObjectIdVersioned,
write::{BatchBuilder, RegistryClass, ValueClass, assert::AssertValue, key::KeySerializer},
};
use trc::{AddContext, OutgoingReportEvent};
use utils::DomainPart;
impl<T: SessionStream> Session<T> {
#[allow(clippy::too_many_arguments)]
pub async fn send_dmarc_report(
&self,
message: &AuthenticatedMessage<'_>,
auth_results: &AuthenticationResults<'_>,
rejected: bool,
dmarc_output: DmarcOutput,
dkim_output: &[DkimOutput<'_>],
dkim2_output: Option<&Dkim2Output<'_>>,
arc_output: &Option<ArcOutput<'_>>,
) {
let dmarc_record = dmarc_output.dmarc_record_cloned().unwrap();
let config = &self.server.core.smtp.report.dmarc;
if self
.server
.is_local_report_domain(dmarc_output.domain(), self.data.session_id)
.await
{
return;
}
// Send failure report. RFC 9991 Section 2: report generators MUST NOT
// honor "ruf" for policy records published with "psd=y".
if !matches!(dmarc_record.psd, dmarc::Psd::Yes)
&& let (Some(failure_rate), Some(report_options)) = (
self.server
.eval_if::<Rate, _>(&config.send, self, self.data.session_id)
.await,
dmarc_output.failure_report(),
)
{
// Verify that any external reporting addresses are authorized
let rcpts = match self
.server
.core
.smtp
.resolvers
.dns
.verify_dmarc_report_address(
dmarc_output.domain(),
dmarc_record.ruf(),
Some(&self.server.inner.cache.dns_txt),
)
.await
{
Some(rcpts) => {
if !rcpts.is_empty() {
let mut new_rcpts = Vec::with_capacity(rcpts.len());
for rcpt in rcpts {
if self.throttle_rcpt(rcpt.uri(), &failure_rate, "dmarc").await {
new_rcpts.push(rcpt.uri());
}
}
new_rcpts
} else {
if !dmarc_record.ruf().is_empty() {
trc::event!(
OutgoingReport(OutgoingReportEvent::UnauthorizedReportingAddress),
SpanId = self.data.session_id,
Url = dmarc_record
.ruf()
.iter()
.map(|u| trc::Value::String(u.uri().to_compact_string()))
.collect::<Vec<_>>(),
);
}
vec![]
}
}
None => {
trc::event!(
OutgoingReport(OutgoingReportEvent::ReportingAddressValidationError),
SpanId = self.data.session_id,
Url = dmarc_record
.ruf()
.iter()
.map(|u| trc::Value::String(u.uri().to_compact_string()))
.collect::<Vec<_>>(),
);
vec![]
}
};
// Throttle recipient
if !rcpts.is_empty() {
let mut report = Vec::with_capacity(128);
let from_addr = self
.server
.eval_if(&config.address, self, self.data.session_id)
.await
.unwrap_or_else(|| "MAILER-DAEMON@localhost".to_compact_string());
let mut auth_failure = self
.new_auth_failure(AuthFailureType::Dmarc, rejected)
.with_authentication_results(auth_results.to_string())
.with_headers(std::str::from_utf8(message.raw_headers()).unwrap_or_default());
let dkim_aligned = matches!(dmarc_output.dkim_result(), DmarcResult::Pass);
let spf_aligned = matches!(dmarc_output.spf_result(), DmarcResult::Pass);
// Report the first failed signature
if let (
dmarc::Report::Dkim
| dmarc::Report::DkimSpf
| dmarc::Report::All
| dmarc::Report::Any,
Some(signature),
) = (
&report_options,
if !dkim_aligned {
dkim_output
.iter()
.find_map(|o| {
let s = o.signature()?;
if !matches!(o.result(), DkimResult::Pass) {
Some(s)
} else {
None
}
})
.or_else(|| dkim_output.iter().find_map(|o| o.signature()))
} else {
None
},
) {
auth_failure = auth_failure
.with_dkim_domain(signature.domain())
.with_dkim_selector(signature.selector())
.with_dkim_identity(signature.identity());
}
// Report SPF failure
if let (
dmarc::Report::Spf
| dmarc::Report::DkimSpf
| dmarc::Report::All
| dmarc::Report::Any,
Some(output),
) = (
&report_options,
if !spf_aligned {
self.data
.spf_ehlo
.as_ref()
.and_then(|s| {
if s.result() != SpfResult::Pass {
s.into()
} else {
None
}
})
.or_else(|| {
self.data.spf_mail_from.as_ref().and_then(|s| {
if s.result() != SpfResult::Pass {
s.into()
} else {
None
}
})
})
.or(self.data.spf_mail_from.as_ref())
} else {
None
},
) {
auth_failure =
auth_failure.with_spf_dns(format!("txt : {} : v=SPF1", output.domain()));
// TODO use DNS record
}
auth_failure
.with_identity_alignment(match (dkim_aligned, spf_aligned) {
(false, false) => IdentityAlignment::DkimSpf,
(false, true) => IdentityAlignment::Dkim,
(true, false) => IdentityAlignment::Spf,
(true, true) => IdentityAlignment::None,
})
.write_rfc5322(
(
self.server
.eval_if(&config.name, self, self.data.session_id)
.await
.unwrap_or_else(|| "Mail Delivery Subsystem".to_compact_string())
.as_str(),
from_addr.as_str(),
),
&rcpts.join(", "),
&self
.server
.eval_if(&config.subject, self, self.data.session_id)
.await
.unwrap_or_else(|| "DMARC Report".to_compact_string()),
&mut report,
)
.ok();
trc::event!(
OutgoingReport(OutgoingReportEvent::DmarcReport),
SpanId = self.data.session_id,
From = from_addr.to_string(),
To = rcpts
.iter()
.map(|a| trc::Value::String(a.to_compact_string()))
.collect::<Vec<_>>(),
);
// Send report
self.server
.send_report(
&from_addr,
rcpts.into_iter(),
report,
&config.sign,
true,
self.data.session_id,
)
.await;
} else {
trc::event!(
OutgoingReport(OutgoingReportEvent::DmarcRateLimited),
SpanId = self.data.session_id,
Limit = vec![
trc::Value::from(failure_rate.count),
trc::Value::from(failure_rate.period.into_inner())
],
);
}
}
// Send aggregate reports
let interval = self
.server
.eval_if(
&self.server.core.smtp.report.dmarc_aggregate.send,
self,
self.data.session_id,
)
.await
.unwrap_or(AggregateFrequency::Never);
if matches!(interval, AggregateFrequency::Never) || dmarc_record.rua().is_empty() {
return;
}
// Report the same identifier forms that were used for alignment
let message_from = message.from();
let header_from = message_from.domain_part();
let header_from = header_from
.to_ascii_domain()
.unwrap_or(Cow::Borrowed(header_from));
let envelope_from = self
.data
.mail_from
.as_ref()
.map(|mf| mf.domain.as_str())
.unwrap_or_else(|| self.data.helo_domain.as_str());
let envelope_from = envelope_from
.to_ascii_domain()
.unwrap_or(Cow::Borrowed(envelope_from));
// Create DMARC report record
let mut report_record = Record::new()
.with_dmarc_output(&dmarc_output)
.with_dkim_output(dkim_output)
.with_source_ip(self.data.remote_ip)
.with_header_from(header_from.as_ref())
.with_envelope_from(envelope_from.as_ref());
if let Some(dkim2_output) = dkim2_output {
report_record = report_record.with_dkim2_output(dkim2_output);
}
if let Some(spf_mail_from) = &self.data.spf_mail_from {
report_record = report_record.with_spf_output(spf_mail_from, SPFDomainScope::MailFrom);
}
if let Some(arc_output) = arc_output {
report_record = report_record.with_arc_output(arc_output);
}
// Submit DMARC report event
self.server
.schedule_report(DmarcEvent {
domain: dmarc_output.into_domain(),
report_record,
dmarc_record,
interval,
span_id: self.data.session_id,
})
.await;
}
}
pub trait DmarcReporting: Sync + Send {
fn send_dmarc_aggregate_report(
&self,
report_id: u64,
) -> impl Future<Output = trc::Result<()>> + Send;
fn schedule_dmarc(&self, event: Box<DmarcEvent>) -> impl Future<Output = ()> + Send;
}
impl DmarcReporting for Server {
async fn send_dmarc_aggregate_report(&self, item_id: u64) -> trc::Result<()> {
let object_id = ObjectType::DmarcInternalReport.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
// Delete report. inbuxa: only the version read here, so a record
// another node appends meanwhile is sent with it rather than lost
let mut attempt = 0;
let report = loop {
let Some(Revisioned {
revision,
value: report,
}) = self
.store()
.get_value::<Revisioned<DmarcInternalReport>>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(());
};
let mut batch = BatchBuilder::new();
batch
.assert_value(key.clone(), AssertValue::Hash(revision))
.clear(key.clone())
.clear(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: KeySerializer::new(report.domain.len() + U64_LEN)
.write(&report.domain)
.write(report.policy_identifier)
.finalize(),
});
match self.store().write(batch.build_all()).await {
Ok(_) => break report,
Err(err) if err.is_assertion_failure() && attempt < MAX_WRITE_RETRIES => {
attempt += 1;
write_retry_pause(attempt).await;
}
Err(err) => return Err(err.caused_by(trc::location!())),
}
};
let span_id = self.inner.data.span_id_gen.generate();
let event_from = report.report.date_range_begin.timestamp() as u64;
let event_to = report.report.date_range_end.timestamp() as u64;
trc::event!(
OutgoingReport(OutgoingReportEvent::DmarcAggregateReport),
SpanId = span_id,
ReportId = event_from,
Domain = report.domain.clone(),
RangeFrom = trc::Value::Timestamp(event_from),
RangeTo = trc::Value::Timestamp(event_to),
);
// Verify external reporting addresses
let rua = match self
.core
.smtp
.resolvers
.dns
.verify_dmarc_report_address(
&report.domain,
report.rua.as_slice(),
Some(&self.inner.cache.dns_txt),
)
.await
{
Some(rcpts) => {
if !rcpts.is_empty() {
rcpts
} else {
trc::event!(
OutgoingReport(OutgoingReportEvent::UnauthorizedReportingAddress),
SpanId = span_id,
Url = report
.rua
.into_iter()
.map(|u| trc::Value::String(u.into()))
.collect::<Vec<_>>(),
);
return Ok(());
}
}
None => {
trc::event!(
OutgoingReport(OutgoingReportEvent::ReportingAddressValidationError),
SpanId = span_id,
Url = report
.rua
.into_iter()
.map(|u| trc::Value::String(u.into()))
.collect::<Vec<_>>(),
);
return Ok(());
}
};
// Serialize report
let config = &self.core.smtp.report.dmarc_aggregate;
let from_addr = self
.eval_if(
&config.address,
&RecipientDomain::new(report.domain.as_str()),
span_id,
)
.await
.unwrap_or_else(|| "MAILER-DAEMON@localhost".to_compact_string());
let mut message = Vec::with_capacity(2048);
let _ = mail_auth::report::Report::from(report.report).write_rfc5322(
&self
.eval_if(
&self.core.smtp.report.submitter,
&RecipientDomain::new(report.domain.as_str()),
span_id,
)
.await
.unwrap_or_else(|| "localhost".to_compact_string()),
(
self.eval_if(
&config.name,
&RecipientDomain::new(report.domain.as_str()),
span_id,
)
.await
.unwrap_or_else(|| "Mail Delivery Subsystem".to_compact_string())
.as_str(),
from_addr.as_str(),
),
rua.iter().map(|a| a.as_str()),
&mut message,
);
// Send report
self.send_report(
&from_addr,
rua.iter(),
message,
&config.sign,
false,
span_id,
)
.await;
Ok(())
}
async fn schedule_dmarc(&self, event: Box<DmarcEvent>) {
let object_id = ObjectType::DmarcInternalReport.to_id();
let policy_hash = event.dmarc_record.to_hash();
let pk = ValueClass::Registry(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: KeySerializer::new(event.domain.len() + U64_LEN)
.write(&event.domain)
.write(policy_hash)
.finalize(),
});
let mut rety_count = 0;
loop {
// Find the report by domain name
let mut batch = BatchBuilder::new();
let report = match self
.store()
.get_value::<ObjectIdVersioned>(ValueKey::from(pk.clone()))
.await
{
Ok(Some(object_id_v)) => {
match self
.store()
.get_value::<DmarcInternalReport>(ValueKey::from(ValueClass::Registry(
RegistryClass::Item {
object_id,
item_id: object_id_v.object_id.id().id(),
},
)))
.await
{
Ok(Some(report)) => Some((object_id_v, report)),
Ok(None) => {
trc::event!(
OutgoingReport(OutgoingReportEvent::NotFound),
Id = object_id_v.object_id.id().id(),
CausedBy = trc::location!(),
Details = "Failed to find DMARC report for domain"
);
return;
}
Err(err) => {
trc::error!(
err.caused_by(trc::location!())
.details("Failed to query registry for DMARC report")
);
return;
}
}
}
Ok(None) => None,
Err(err) => {
trc::error!(
err.caused_by(trc::location!())
.details("Failed to query registry for DMARC report")
);
return;
}
};
// Create report if missing
let config = &self.core.smtp.report.dmarc_aggregate;
let (item_id, mut report) = if let Some((mut object_id_v, report)) = report {
batch.assert_value(pk.clone(), AssertValue::U32(object_id_v.version));
object_id_v.version += 1;
batch.set(pk.clone(), object_id_v.serialize());
(object_id_v.object_id.id().id(), report)
} else {
let item_id = self.inner.data.queue_id_gen.generate();
let date_range_begin = UTCDateTime::now();
let date_range_end = UTCDateTime::from_timestamp(
date_range_begin.timestamp() + event.interval.as_secs() as i64,
);
let policy =
PolicyPublished::from_record(event.domain.clone(), &event.dmarc_record);
let report = DmarcInternalReport {
created_at: date_range_begin,
deliver_at: date_range_end,
domain: event.domain.clone(),
report: DmarcReport {
report_id: format!("{}_{policy_hash}", date_range_begin.timestamp()),
date_range_begin,
date_range_end,
email: self
.eval_if(
&config.address,
&RecipientDomain::new(event.domain.as_str()),
event.span_id,
)
.await
.unwrap_or_else(|| "MAILER-DAEMON@localhost".to_string()),
extra_contact_info: self
.eval_if::<String, _>(
&config.contact_info,
&RecipientDomain::new(event.domain.as_str()),
event.span_id,
)
.await,
org_name: self
.eval_if::<String, _>(
&config.org_name,
&RecipientDomain::new(event.domain.as_str()),
event.span_id,
)
.await
.unwrap_or_default(),
policy_adkim: policy.adkim.into(),
policy_aspf: policy.aspf.into(),
policy_disposition: policy.p.into(),
policy_domain: policy.domain,
policy_failure_reporting_options: match event.dmarc_record.fo {
dmarc::Report::All => vec![FailureReportingOption::All],
dmarc::Report::Any => vec![FailureReportingOption::Any],
dmarc::Report::Dkim => vec![FailureReportingOption::DkimFailure],
dmarc::Report::Spf => vec![FailureReportingOption::SpfFailure],
dmarc::Report::DkimSpf => vec![
FailureReportingOption::DkimFailure,
FailureReportingOption::SpfFailure,
],
}
.into(),
policy_subdomain_disposition: policy.sp.into(),
policy_np: policy.np.into(),
policy_discovery_method: policy.discovery_method.into(),
policy_testing_mode: policy.testing,
policy_version: None,
version: 1.0.into(),
..Default::default()
},
policy_identifier: policy_hash,
rua: Map::new(
event
.dmarc_record
.rua()
.iter()
.map(|u| u.uri.clone())
.collect(),
),
};
report.write_ops(&mut batch, item_id, true);
(item_id, report)
};
// Add record
let mut record = DmarcReportRecord::from(event.report_record.clone());
if let Some(idx) = report
.report
.records
.0
.inner
.iter()
.position(|d| d.value.eq_except_count(&record))
{
report.report.records.0.inner[idx].value.count += 1;
} else {
record.count = 1;
report.report.records.push(record);
}
// Write entry
let report_bytes = report.to_pickled_vec();
let max_report_size = self
.eval_if(
&config.max_size,
&RecipientDomain::new(&event.domain),
event.span_id,
)
.await
.unwrap_or(5 * 1024 * 1024);
if max_report_size != 0 && report_bytes.len() > max_report_size {
trc::event!(
OutgoingReport(OutgoingReportEvent::MaxSizeExceeded),
SpanId = event.span_id,
Domain = event.domain.clone(),
Details = report_bytes.len(),
Limit = max_report_size,
);
return;
}
batch.set(
ValueClass::Registry(RegistryClass::Item { object_id, item_id }),
report_bytes,
);
match self.core.storage.data.write(batch.build_all()).await {
Ok(_) => {
break;
}
Err(err) => {
// inbuxa: another node appended first; try again
// after a short pause
if err.is_assertion_failure() && rety_count < MAX_WRITE_RETRIES {
rety_count += 1;
write_retry_pause(rety_count).await;
continue;
}
trc::error!(
err.caused_by(trc::location!())
.details("Failed to write DMARC report")
);
break;
}
}
}
}
}