From 3cc348726bfd81462716ae43df5ebc29bbaf510a Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Sun, 19 Jul 2026 23:09:31 +0200 Subject: [PATCH] Fix MTA: `queue_name` variable not available in rate limiter expressions --- CHANGELOG.md | 4 +- crates/smtp/src/outbound/delivery.rs | 2 +- crates/smtp/src/queue/mod.rs | 47 ++++++++++++++ tests/src/smtp/outbound/throttle.rs | 91 ++++++++++++++++++++++++++++ 4 files changed, 142 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2f0365b3..3c512edf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,7 +22,9 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If - `Email/import` does not send push notifications for imported messages. - `CalendarEvent/set` silently ignores `ifInState`. - CalDAV: `calendar-query` REPORT returns empty calendar-data for JMAP-created events. -- MTA: DMARC is skipped when MAIL FROM SPF is unavailable. +- MTA: + - DMARC is skipped when MAIL FROM SPF is unavailable. + - `queue_name` variable not available in rate limiter expressions. - Calendar: - No expanded occurrences are returned for a daily recurrences crossing DST. - Uppercase `MAILTO` calendar addresses become invalid SMTP recipients. diff --git a/crates/smtp/src/outbound/delivery.rs b/crates/smtp/src/outbound/delivery.rs index 5d62a2e5..f1468081 100644 --- a/crates/smtp/src/outbound/delivery.rs +++ b/crates/smtp/src/outbound/delivery.rs @@ -185,7 +185,7 @@ impl QueuedMessage { // Throttle sender for throttle in &server.core.smtp.queue.outbound_limiters.sender { if let Err(retry_at) = server - .is_allowed(throttle, &message.message, message.span_id) + .is_allowed(throttle, &message, message.span_id) .await { trc::event!( diff --git a/crates/smtp/src/queue/mod.rs b/crates/smtp/src/queue/mod.rs index 9cf933a2..ee1a52cb 100644 --- a/crates/smtp/src/queue/mod.rs +++ b/crates/smtp/src/queue/mod.rs @@ -379,6 +379,53 @@ impl ResolveVariable for Message { } } +impl ResolveVariable for MessageWrapper { + fn resolve_variable(&self, variable: ExpressionVariable) -> expr::Variable<'_> { + match variable { + ExpressionVariable::Sender => self.message.return_path.as_ref().into(), + ExpressionVariable::SenderDomain => self.message.return_path.domain_part().into(), + ExpressionVariable::Recipients => self + .message + .recipients + .iter() + .map(|r| Variable::from(r.address.as_ref())) + .collect::>() + .into(), + ExpressionVariable::Priority => self.message.priority.into(), + ExpressionVariable::QueueName => self.queue_name.as_str().into(), + ExpressionVariable::QueueAge => { + now().saturating_sub(self.message.created).into() + } + ExpressionVariable::Source => if (self.message.flags & FROM_AUTHENTICATED) != 0 { + "authenticated" + } else if (self.message.flags & FROM_UNAUTHENTICATED_DMARC) != 0 { + "dmarc_pass" + } else if (self.message.flags & FROM_UNAUTHENTICATED) != 0 { + "unauthenticated" + } else if (self.message.flags & FROM_DSN) != 0 { + "dsn" + } else if (self.message.flags & FROM_REPORT) != 0 { + "report" + } else if (self.message.flags & FROM_AUTOGENERATED) != 0 { + "autogenerated" + } else { + "unknown" + } + .into(), + ExpressionVariable::ReceivedFromIp => { + self.message.received_from_ip.to_compact_string().into() + } + ExpressionVariable::ReceivedViaPort => self.message.received_via_port.into(), + ExpressionVariable::Size => self.message.size.into(), + _ => "".into(), + } + } + + fn resolve_global(&self, _: &str) -> Variable<'_> { + Variable::Integer(0) + } +} + pub struct RecipientDomain<'x>(&'x str); impl<'x> RecipientDomain<'x> { diff --git a/tests/src/smtp/outbound/throttle.rs b/tests/src/smtp/outbound/throttle.rs index c601135b..6b9e0c56 100644 --- a/tests/src/smtp/outbound/throttle.rs +++ b/tests/src/smtp/outbound/throttle.rs @@ -25,6 +25,7 @@ use registry::{ }, types::{list::List, map::Map}, }; +use common::config::smtp::queue::QueueName; use smtp::queue::{Message, QueueEnvelope, Recipient, throttle::IsAllowed}; use std::{ net::{IpAddr, Ipv4Addr}, @@ -304,6 +305,96 @@ async fn throttle_outbound() { assert!(due > 0, "Due: {}", due); } +#[tokio::test] +async fn throttle_outbound_queue_name() { + let mut local = TestServerBuilder::new("smtp_throttle_outbound_queue_name") + .await + .with_http_listener(19033) + .await + .disable_services() + .capture_queue() + .build() + .await; + + let admin = local.account("admin"); + let queue_id = admin + .registry_create_object(MtaVirtualQueue { + name: "default".into(), + threads_per_node: 25, + description: None, + }) + .await; + admin + .registry_create_object(MtaDeliverySchedule { + name: "default".into(), + retry: MtaDeliveryScheduleIntervalsOrDefault::Default, + notify: MtaDeliveryScheduleIntervalsOrDefault::Default, + expiry: MtaDeliveryExpiration::Ttl(MtaDeliveryExpirationTtl { + expire: 3_600_000u64.into(), + }), + queue_id, + description: None, + }) + .await; + admin + .registry_create_object(MtaOutboundStrategy { + schedule: Expression { + else_: "'default'".into(), + ..Default::default() + }, + ..Default::default() + }) + .await; + admin + .registry_create_object(MtaOutboundThrottle { + enable: true, + key: Map::new(vec![MtaOutboundThrottleKey::SenderDomain]), + match_: Expression { + else_: "queue_name == 'default'".into(), + ..Default::default() + }, + rate: Rate { + count: 1, + period: (30u64 * 60 * 1000).into(), + }, + description: "queue_name throttle".into(), + }) + .await; + admin.reload_settings().await; + local.reload_core(); + local.expect_reload_settings().await; + + let core = local.server.core.clone(); + let throttle = &core.smtp.queue.outbound_limiters; + assert_eq!(throttle.sender.len(), 1); + assert!(throttle.rcpt.is_empty()); + assert!(throttle.remote.is_empty()); + + let matching = new_message(0); + assert_eq!(matching.queue_name, QueueName::default()); + local + .server + .is_allowed(&throttle.sender[0], &matching, 0) + .await + .unwrap(); + assert!( + local + .server + .is_allowed(&throttle.sender[0], &matching, 0) + .await + .is_err(), + "sender-bucket throttle failed to resolve queue_name and never engaged" + ); + + let mut other_queue = new_message(0); + other_queue.queue_name = QueueName::new("remote").unwrap(); + local + .server + .is_allowed(&throttle.sender[0], &other_queue, 0) + .await + .unwrap(); +} + pub trait TestQueueEnvelope<'x> { fn test(message: &'x Message, rcpt: &'x Recipient, mx: &'x str) -> Self; }