From e7f532ec6dd10f84c05bcd111913db9edd710fea Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Mon, 25 May 2026 15:54:40 +0200 Subject: [PATCH] Fix MTA: Always update next DSN notify times --- CHANGELOG.md | 2 + crates/common/src/network/dns/update.rs | 3 +- .../src/registry/mapping/queued_message.rs | 2 +- crates/services/src/task_manager/dns.rs | 5 +- crates/smtp/src/inbound/data.rs | 7 +- crates/smtp/src/inbound/ehlo.rs | 2 +- crates/smtp/src/inbound/hooks/message.rs | 1 + crates/smtp/src/inbound/mail.rs | 2 +- crates/smtp/src/inbound/milter/message.rs | 3 + crates/smtp/src/inbound/rcpt.rs | 2 +- crates/smtp/src/inbound/spawn.rs | 2 +- crates/smtp/src/queue/dsn.rs | 76 +++++++++---------- crates/smtp/src/queue/spool.rs | 31 +++++++- 13 files changed, 88 insertions(+), 50 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e98292da..5aed8583 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,8 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If ## Changed ## Fixed +- Log rejected messages to tracing store. +- MTA: Always update next DSN notify times. ## [0.16.6] - 2026-05-20 diff --git a/crates/common/src/network/dns/update.rs b/crates/common/src/network/dns/update.rs index 36226ba6..89f796e2 100644 --- a/crates/common/src/network/dns/update.rs +++ b/crates/common/src/network/dns/update.rs @@ -704,7 +704,6 @@ impl DnsUpdater { }), DnsServer::Inwx(server) => { let password = server.password.secret().await?.into_owned(); - let shared_secret = server.shared_secret.secret().await?.map(|c| c.into_owned()); Ok(DnsUpdater { polling_interval: server.polling_interval.into_inner(), propagation_timeout: server.propagation_timeout.into_inner(), @@ -714,7 +713,6 @@ impl DnsUpdater { updater: dns_update::DnsUpdater::new_inwx( server.username, password, - shared_secret, server.sandbox, server.timeout.into_inner().into(), ) @@ -1014,6 +1012,7 @@ impl DnsUpdater { updater: dns_update::DnsUpdater::new_transip( server.username.as_str(), server.private_key_pem.secret().await?, + true, server.timeout.into_inner().into(), ) .map_err(|err| format!("Failed to build DNS updater: {}", err))?, diff --git a/crates/jmap/src/registry/mapping/queued_message.rs b/crates/jmap/src/registry/mapping/queued_message.rs index e5a89c77..35cb7494 100644 --- a/crates/jmap/src/registry/mapping/queued_message.rs +++ b/crates/jmap/src/registry/mapping/queued_message.rs @@ -81,7 +81,7 @@ pub(crate) async fn queued_message_set( } // Process patches - let prev_event = archived_message.inner.next_delivery_event(None); + let prev_event = archived_message.inner.next_event(None); let mut message = map_message(archived_message.inner); message.next_retry = None; for (key, value) in value.into_expanded_object() { diff --git a/crates/services/src/task_manager/dns.rs b/crates/services/src/task_manager/dns.rs index 1a2093ba..242b2d2f 100644 --- a/crates/services/src/task_manager/dns.rs +++ b/crates/services/src/task_manager/dns.rs @@ -10,7 +10,8 @@ use dns_update::{DnsRecord, DnsRecordType}; use registry::schema::structs::{ DnsManagement, Domain, Task, TaskDnsManagement, TaskDomainManagement, TaskStatus, }; -use std::{collections::HashMap, fmt::Write}; +use std::fmt::Write; +use store::ahash::AHashMap; pub(crate) trait DnsManagementTask: Sync + Send { fn dns_management(&self, task: &TaskDnsManagement) -> impl Future + Send; @@ -61,7 +62,7 @@ async fn dns_management(server: &Server, task: &TaskDnsManagement) -> trc::Resul .await?; // Group records by (name, type) so each RRSet is published in one call. - let mut by_owner: HashMap<(String, DnsRecordType), Vec> = HashMap::new(); + let mut by_owner: AHashMap<(String, DnsRecordType), Vec> = AHashMap::new(); for record in records { by_owner .entry((record.name, record.record.as_type())) diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index 15a9bdd5..a1eb5e57 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -456,6 +456,7 @@ impl Session { trc::event!( Spam(SpamEvent::Classify), SpanId = self.data.session_id, + QueueId = message_id, Result = "discard", Reason = "Message discarded due to excessive spam score.", ); @@ -467,6 +468,7 @@ impl Session { trc::event!( Spam(SpamEvent::Classify), SpanId = self.data.session_id, + QueueId = message_id, Result = "reject", Reason = "Message rejected due to excessive spam score.", ); @@ -481,7 +483,10 @@ impl Session { // Run Milter filters let mut modifications = Vec::new(); - match self.run_milters(Stage::Data, (&auth_message).into()).await { + match self + .run_milters(Stage::Data, (&auth_message).into(), message_id.into()) + .await + { Ok(modifications_) => { if !modifications_.is_empty() { modifications = modifications_; diff --git a/crates/smtp/src/inbound/ehlo.rs b/crates/smtp/src/inbound/ehlo.rs index 57343e73..d82d46b8 100644 --- a/crates/smtp/src/inbound/ehlo.rs +++ b/crates/smtp/src/inbound/ehlo.rs @@ -115,7 +115,7 @@ impl Session { } // Milter filtering - if let Err(message) = self.run_milters(Stage::Ehlo, None).await { + if let Err(message) = self.run_milters(Stage::Ehlo, None, None).await { self.data.mail_from = None; self.data.helo_domain = prev_helo_domain; self.data.spf_ehlo = None; diff --git a/crates/smtp/src/inbound/hooks/message.rs b/crates/smtp/src/inbound/hooks/message.rs index 3df9b89c..31bff929 100644 --- a/crates/smtp/src/inbound/hooks/message.rs +++ b/crates/smtp/src/inbound/hooks/message.rs @@ -61,6 +61,7 @@ impl Session { Action::Quarantine => MtaHookEvent::ActionQuarantine, }), SpanId = self.data.session_id, + QueueId = queue_id, Id = mta_hook.id.to_string(), Elapsed = time.elapsed(), ); diff --git a/crates/smtp/src/inbound/mail.rs b/crates/smtp/src/inbound/mail.rs index 301d5cc5..d732e1e1 100644 --- a/crates/smtp/src/inbound/mail.rs +++ b/crates/smtp/src/inbound/mail.rs @@ -187,7 +187,7 @@ impl Session { } // Milter filtering - if let Err(message) = self.run_milters(Stage::Mail, None).await { + if let Err(message) = self.run_milters(Stage::Mail, None, None).await { self.data.mail_from = None; return self.write(message.message.as_bytes()).await; } diff --git a/crates/smtp/src/inbound/milter/message.rs b/crates/smtp/src/inbound/milter/message.rs index 955a85f7..e910c2ba 100644 --- a/crates/smtp/src/inbound/milter/message.rs +++ b/crates/smtp/src/inbound/milter/message.rs @@ -8,6 +8,7 @@ use super::{Action, Error, Macros, Modification}; use crate::{ core::{Session, SessionAddress, SessionData}, inbound::{FilterResponse, milter::MilterClient}, + queue::QueueId, }; use common::{ DAEMON_NAME, @@ -31,6 +32,7 @@ impl Session { &self, stage: Stage, message: Option<&AuthenticatedMessage<'_>>, + queue_id: Option, ) -> Result, FilterResponse> { let milters = &self.server.core.smtp.session.milters; if milters.is_empty() { @@ -88,6 +90,7 @@ impl Session { Action::Accept | Action::Continue => unreachable!(), }), SpanId = self.data.session_id, + QueueId = queue_id, Id = milter.id.to_string(), Elapsed = time.elapsed(), ); diff --git a/crates/smtp/src/inbound/rcpt.rs b/crates/smtp/src/inbound/rcpt.rs index 75778e52..15a68919 100644 --- a/crates/smtp/src/inbound/rcpt.rs +++ b/crates/smtp/src/inbound/rcpt.rs @@ -140,7 +140,7 @@ impl Session { } // Milter filtering - if let Err(message) = self.run_milters(Stage::Rcpt, None).await { + if let Err(message) = self.run_milters(Stage::Rcpt, None, None).await { self.data.rcpt_to.pop(); return self.write(message.message.as_bytes()).await; } diff --git a/crates/smtp/src/inbound/spawn.rs b/crates/smtp/src/inbound/spawn.rs index 4b80df23..78492244 100644 --- a/crates/smtp/src/inbound/spawn.rs +++ b/crates/smtp/src/inbound/spawn.rs @@ -98,7 +98,7 @@ impl Session { } // Milter filtering - if let Err(message) = self.run_milters(Stage::Connect, None).await { + if let Err(message) = self.run_milters(Stage::Connect, None, None).await { let _ = self.write(message.message.as_bytes()).await; return false; } diff --git a/crates/smtp/src/queue/dsn.rs b/crates/smtp/src/queue/dsn.rs index a2375367..cabf9c25 100644 --- a/crates/smtp/src/queue/dsn.rs +++ b/crates/smtp/src/queue/dsn.rs @@ -62,6 +62,9 @@ impl SendDsn for Server { // Handle double bounce message.handle_double_bounce(); } + + // Update next DSN notify times + message.update_next_dsn(self).await; } async fn log_dsn(&self, message: &MessageWrapper) { @@ -186,7 +189,6 @@ impl MessageWrapper { dsn.push_str("\r\n"); } - // Build text response let txt_len = txt_success.len() + txt_delay.len() + txt_failed.len(); if txt_len == 0 { return None; @@ -250,44 +252,6 @@ impl MessageWrapper { txt.push_str("\r\n"); } - // Update next delay notification time - if has_delay { - let mut changes = Vec::new(); - for (rcpt_idx, rcpt) in self.message.recipients.iter().enumerate() { - if matches!( - &rcpt.status, - Status::TemporaryFailure(_) | Status::Scheduled - ) && rcpt.notify.due <= now - { - let envelope = QueueEnvelope::new(&self.message, rcpt); - - let queue_id = server - .eval_if::( - &server.core.smtp.queue.queue, - &envelope, - self.span_id, - ) - .await - .unwrap_or_else(|| "default".to_string()); - let queue = server.get_queue_or_default(&queue_id, self.span_id); - - if let Some(next_notify) = - queue.notify.get((rcpt.notify.inner + 1) as usize).copied() - { - changes.push((rcpt_idx, 1, now + next_notify)); - } else { - changes.push((rcpt_idx, 0, u64::MAX)); - } - } - } - - for (rcpt_idx, inner, due) in changes { - let rcpt = &mut self.message.recipients[rcpt_idx]; - rcpt.notify.inner += inner; - rcpt.notify.due = due; - } - } - // Obtain hostname and sender addresses let from_name = server .eval_if(&config.dsn.name, &self.message, self.span_id) @@ -393,6 +357,40 @@ impl MessageWrapper { .into() } + pub async fn update_next_dsn(&mut self, server: &Server) { + let now = now(); + let mut notify_changes = Vec::new(); + for (rcpt_idx, rcpt) in self.message.recipients.iter().enumerate() { + if matches!( + &rcpt.status, + Status::TemporaryFailure(_) | Status::Scheduled + ) && rcpt.notify.due <= now + { + let envelope = QueueEnvelope::new(&self.message, rcpt); + + let queue_id = server + .eval_if::(&server.core.smtp.queue.queue, &envelope, self.span_id) + .await + .unwrap_or_else(|| "default".to_string()); + let queue = server.get_queue_or_default(&queue_id, self.span_id); + + if let Some(next_notify) = + queue.notify.get((rcpt.notify.inner + 1) as usize).copied() + { + notify_changes.push((rcpt_idx, 1, now + next_notify)); + } else { + notify_changes.push((rcpt_idx, 0, u64::MAX)); + } + } + } + + for (rcpt_idx, inner, due) in notify_changes { + let rcpt = &mut self.message.recipients[rcpt_idx]; + rcpt.notify.inner += inner; + rcpt.notify.due = due; + } + } + fn handle_double_bounce(&mut self) { let mut is_double_bounce = Vec::with_capacity(0); let now = now(); diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index 6bee013a..bee4e364 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -14,7 +14,7 @@ use crate::queue::{ FROM_UNAUTHENTICATED_DMARC, MessageWrapper, }; use ahash::AHashSet; -use common::config::smtp::queue::QueueName; +use common::config::smtp::queue::{ArchivedQueueExpiry, QueueName}; use common::ipc::QueueEvent; use common::{KV_LOCK_QUEUE_MESSAGE, Server}; use registry::schema::prelude::{ObjectType, Property}; @@ -803,6 +803,35 @@ impl ArchivedMessage { next_delivery } + pub fn next_event(&self, queue: Option) -> Option { + let created = self.created.to_native(); + let mut next_event = None; + + for rcpt in self.recipients.iter().filter(|d| { + matches!( + d.status, + ArchivedStatus::Scheduled | ArchivedStatus::TemporaryFailure(_) + ) && queue.is_none_or(|q| d.queue == q) + }) { + let mut earlier_event = + std::cmp::min(rcpt.retry.due.to_native(), rcpt.notify.due.to_native()); + + if let ArchivedQueueExpiry::Ttl(ttl) = &rcpt.expires { + earlier_event = std::cmp::min(earlier_event, created + ttl.to_native()); + } + + if let Some(next_event) = &mut next_event { + if earlier_event < *next_event { + *next_event = earlier_event; + } + } else { + next_event = Some(earlier_event); + } + } + + next_event + } + pub fn next_notify_event(&self, queue: Option) -> Option { let mut next_notify = None;