diff --git a/crates/common/src/config/smtp/queue.rs b/crates/common/src/config/smtp/queue.rs index 23856ca7..287bb094 100644 --- a/crates/common/src/config/smtp/queue.rs +++ b/crates/common/src/config/smtp/queue.rs @@ -34,6 +34,7 @@ use utils::config::{Config, utils::ParseValue}; rkyv::Archive, serde::Deserialize, )] +#[rkyv(derive(Debug, Clone, Copy, PartialEq), compare(PartialEq))] #[repr(transparent)] pub struct QueueName([u8; 8]); @@ -433,13 +434,13 @@ fn parse_route(config: &mut Config, id: &str) -> Option { "local" => RoutingStrategy::Local.into(), "mx" => RoutingStrategy::Mx(MxConfig { max_mx: config - .property_require(("queue.route", id, "limits.mx")) + .property(("queue.route", id, "limits.mx")) .unwrap_or(5), max_multi_homed: config - .property_require(("queue.route", id, "limits.multihomed")) + .property(("queue.route", id, "limits.multihomed")) .unwrap_or(2), ip_lookup_strategy: config - .property_require(("queue.route", id, "ip-lookup")) + .property(("queue.route", id, "ip-lookup")) .unwrap_or(IpLookupStrategy::Ipv4thenIpv6), }) .into(), @@ -475,22 +476,22 @@ fn parse_tls_strategies(config: &mut Config) -> AHashMap { fn parse_tls(config: &mut Config, id: &str) -> Option { Some(TlsStrategy { dane: config - .property_require::(("queue.tls", id, "dane")) + .property::(("queue.tls", id, "dane")) .unwrap_or(RequireOptional::Optional), mta_sts: config - .property_require::(("queue.tls", id, "mta-sts")) + .property::(("queue.tls", id, "mta-sts")) .unwrap_or(RequireOptional::Optional), tls: config - .property_require::(("queue.tls", id, "starttls")) + .property::(("queue.tls", id, "starttls")) .unwrap_or(RequireOptional::Optional), allow_invalid_certs: config - .property_require::(("queue.tls", id, "allow-invalid-certs")) + .property::(("queue.tls", id, "allow-invalid-certs")) .unwrap_or(false), timeout_tls: config - .property_require::(("queue.tls", id, "timeout.tls")) + .property::(("queue.tls", id, "timeout.tls")) .unwrap_or(Duration::from_secs(3 * 60)), timeout_mta_sts: config - .property_require::(("queue.tls", id, "timeout.mta-sts")) + .property::(("queue.tls", id, "timeout.mta-sts")) .unwrap_or(Duration::from_secs(5 * 60)), }) } @@ -538,22 +539,22 @@ fn parse_connection(config: &mut Config, id: &str) -> Option source_ipv6, ehlo_hostname: config.property::(("queue.connection", id, "ehlo-hostname")), timeout_connect: config - .property_require::(("queue.connection", id, "timeout.connect")) + .property::(("queue.connection", id, "timeout.connect")) .unwrap_or(Duration::from_secs(5 * 60)), timeout_greeting: config - .property_require::(("queue.connection", id, "timeout.greeting")) + .property::(("queue.connection", id, "timeout.greeting")) .unwrap_or(Duration::from_secs(5 * 60)), timeout_ehlo: config - .property_require::(("queue.connection", id, "timeout.ehlo")) + .property::(("queue.connection", id, "timeout.ehlo")) .unwrap_or(Duration::from_secs(5 * 60)), timeout_mail: config - .property_require::(("queue.connection", id, "timeout.mail-from")) + .property::(("queue.connection", id, "timeout.mail-from")) .unwrap_or(Duration::from_secs(5 * 60)), timeout_rcpt: config - .property_require::(("queue.connection", id, "timeout.rcpt-to")) + .property::(("queue.connection", id, "timeout.rcpt-to")) .unwrap_or(Duration::from_secs(5 * 60)), timeout_data: config - .property_require::(("queue.connection", id, "timeout.data")) + .property::(("queue.connection", id, "timeout.data")) .unwrap_or(Duration::from_secs(10 * 60)), }) } @@ -935,6 +936,14 @@ impl QueueName { } } +impl ArchivedQueueName { + pub fn as_str(&self) -> &str { + std::str::from_utf8(self.0.as_ref()) + .unwrap_or_default() + .trim_end_matches('\0') + } +} + impl Default for QueueName { fn default() -> Self { DEFAULT_QUEUE_NAME @@ -955,7 +964,13 @@ impl ParseValue for QueueName { impl Display for QueueName { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{}", self.as_str()) + self.as_str().fmt(f) + } +} + +impl Display for ArchivedQueueName { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.as_str().fmt(f) } } diff --git a/crates/common/src/expr/tokenizer.rs b/crates/common/src/expr/tokenizer.rs index 0d8677f8..91c3181a 100644 --- a/crates/common/src/expr/tokenizer.rs +++ b/crates/common/src/expr/tokenizer.rs @@ -367,8 +367,14 @@ impl TokenMap { V_QUEUE_EXPIRES_IN, V_QUEUE_LAST_STATUS, V_QUEUE_LAST_ERROR, + V_QUEUE_NAME, + V_QUEUE_AGE, V_ASN, V_COUNTRY, + V_RECEIVED_FROM_IP, + V_RECEIVED_VIA_PORT, + V_SOURCE, + V_SIZE, ]) } diff --git a/crates/common/src/manager/boot.rs b/crates/common/src/manager/boot.rs index 4a479006..f765fb95 100644 --- a/crates/common/src/manager/boot.rs +++ b/crates/common/src/manager/boot.rs @@ -29,7 +29,7 @@ use store::{ }; use tokio::sync::{Notify, mpsc}; use utils::{ - Semver, UnwrapFailure, + UnwrapFailure, config::{Config, ConfigKey}, failed, }; @@ -75,6 +75,106 @@ enum StoreOp { None, } +pub const DEFAULT_SETTINGS: &[(&str, &str)] = &[ + ("queue.quota.size.messages", "100000"), + ("queue.quota.size.size", "10737418240"), + ("queue.quota.size.enable", "true"), + ("queue.limiter.inbound.ip.key", "remote_ip"), + ("queue.limiter.inbound.ip.rate", "5/1s"), + ("queue.limiter.inbound.ip.enable", "true"), + ("queue.limiter.inbound.sender.key.0", "sender_domain"), + ("queue.limiter.inbound.sender.key.1", "rcpt"), + ("queue.limiter.inbound.sender.rate", "25/1h"), + ("queue.limiter.inbound.sender.enable", "true"), + ("report.analysis.addresses", "postmaster@*"), + ("queue.virtual.local.threads-per-node", "25"), + ("queue.virtual.local.description", "Local delivery queue"), + ("queue.virtual.remote.threads-per-node", "50"), + ("queue.virtual.remote.description", "Remote delivery queue"), + ("queue.virtual.dsn.threads-per-node", "5"), + ( + "queue.virtual.dsn.description", + "Delivery Status Notification delivery queue", + ), + ("queue.virtual.report.threads-per-node", "5"), + ( + "queue.virtual.report.description", + "DMARC and TLS report delivery queue", + ), + ("queue.schedule.local.queue-name", "local"), + ("queue.schedule.local.retry.0", "2m"), + ("queue.schedule.local.retry.1", "5m"), + ("queue.schedule.local.retry.2", "10m"), + ("queue.schedule.local.retry.3", "15m"), + ("queue.schedule.local.retry.4", "30m"), + ("queue.schedule.local.retry.5", "1h"), + ("queue.schedule.local.retry.6", "2h"), + ("queue.schedule.local.notify.0", "1d"), + ("queue.schedule.local.notify.1", "3d"), + ("queue.schedule.local.expire-type", "ttl"), + ("queue.schedule.local.expire", "3d"), + ( + "queue.schedule.local.description", + "Local delivery schedule", + ), + ("queue.schedule.remote.queue-name", "remote"), + ("queue.schedule.remote.retry.0", "2m"), + ("queue.schedule.remote.retry.1", "5m"), + ("queue.schedule.remote.retry.2", "10m"), + ("queue.schedule.remote.retry.3", "15m"), + ("queue.schedule.remote.retry.4", "30m"), + ("queue.schedule.remote.retry.5", "1h"), + ("queue.schedule.remote.retry.6", "2h"), + ("queue.schedule.remote.notify.0", "1d"), + ("queue.schedule.remote.notify.1", "3d"), + ("queue.schedule.remote.expire-type", "ttl"), + ("queue.schedule.remote.expire", "3d"), + ( + "queue.schedule.remote.description", + "Remote delivery schedule", + ), + ("queue.schedule.dsn.queue-name", "dsn"), + ("queue.schedule.dsn.retry.0", "15m"), + ("queue.schedule.dsn.retry.1", "30m"), + ("queue.schedule.dsn.retry.2", "1h"), + ("queue.schedule.dsn.retry.3", "2h"), + ("queue.schedule.dsn.expire-type", "attempts"), + ("queue.schedule.dsn.max-attempts", "10"), + ( + "queue.schedule.dsn.description", + "Delivery Status Notification delivery schedule", + ), + ("queue.schedule.report.queue-name", "report"), + ("queue.schedule.report.retry.0", "30m"), + ("queue.schedule.report.retry.1", "1h"), + ("queue.schedule.report.retry.2", "2h"), + ("queue.schedule.report.expire-type", "attempts"), + ("queue.schedule.report.max-attempts", "8"), + ( + "queue.schedule.report.description", + "DMARC and TLS report delivery schedule", + ), + ("queue.tls.invalid-tls.allow-invalid-certs", "true"), + ( + "queue.tls.invalid-tls.description", + "Allow invalid TLS certificates", + ), + ("queue.tls.default.allow-invalid-certs", "false"), + ("queue.tls.default.description", "Default TLS settings"), + ("queue.route.local.type", "local"), + ("queue.route.local.description", "Local delivery route"), + ("queue.route.mx.type", "mx"), + ("queue.route.mx.limits.multihomed", "2"), + ("queue.route.mx.limits.mx", "5"), + ("queue.route.mx.ip-lookup", "ipv4_then_ipv6"), + ("queue.route.mx.description", "MX delivery route"), + ("queue.connection.default.timeout.connect", "5m"), + ( + "queue.connection.default.description", + "Default connection settings", + ), +]; + impl BootManager { pub async fn init() -> Self { let mut config_path = std::env::var("CONFIG_PATH").ok(); @@ -244,84 +344,29 @@ impl BootManager { } // Download Spam filter rules if missing - // TODO remove this check in 1.0 - let update_webadmin = match config.value("version.spam-filter").and_then(|v| { - if !v.is_empty() { - Some(Semver::try_from(v)) - } else { - None - } - }) { - Some(Err(_)) => { - let _ = manager.clear_prefix("lookup.spam-").await; - let _ = manager - .clear_prefix("sieve.trusted.scripts.spam-filter") - .await; - let _ = manager - .clear_prefix("sieve.trusted.scripts.track-replies") - .await; - let _ = manager.clear_prefix("sieve.trusted.scripts.greylist").await; - let _ = manager.clear_prefix("sieve.trusted.scripts.train").await; - let _ = manager.clear("version.spam-filter").await; - - match manager.fetch_spam_rules().await { - Ok(external_config) => { - trc::event!( - Config(trc::ConfigEvent::ImportExternal), - Version = external_config.version.to_string(), - Id = "spam-filter" - ); - insert_keys.extend(external_config.keys); - } - Err(err) => { - config.new_build_error( - "*", - format!("Failed to fetch spam filter: {err}"), - ); - } + if config.value("version.spam-filter").is_none() { + match manager.fetch_spam_rules().await { + Ok(external_config) => { + trc::event!( + Config(trc::ConfigEvent::ImportExternal), + Version = external_config.version.to_string(), + Id = "spam-filter" + ); + insert_keys.extend(external_config.keys); } - - true - } - Some(Ok(_)) => false, - None => { - match manager.fetch_spam_rules().await { - Ok(external_config) => { - trc::event!( - Config(trc::ConfigEvent::ImportExternal), - Version = external_config.version.to_string(), - Id = "spam-filter" - ); - insert_keys.extend(external_config.keys); - } - Err(err) => { - config.new_build_error( - "*", - format!("Failed to fetch spam filter: {err}"), - ); - } + Err(err) => { + config.new_build_error( + "*", + format!("Failed to fetch spam filter: {err}"), + ); } - - // Add default settings - for key in [ - ("queue.quota.size.messages", "100000"), - ("queue.quota.size.size", "10737418240"), - ("queue.quota.size.enable", "true"), - ("queue.limiter.inbound.ip.key", "remote_ip"), - ("queue.limiter.inbound.ip.rate", "5/1s"), - ("queue.limiter.inbound.ip.enable", "true"), - ("queue.limiter.inbound.sender.key.0", "sender_domain"), - ("queue.limiter.inbound.sender.key.1", "rcpt"), - ("queue.limiter.inbound.sender.rate", "25/1h"), - ("queue.limiter.inbound.sender.enable", "true"), - ("report.analysis.addresses", "postmaster@*"), - ] { - insert_keys.push(ConfigKey::from(key)); - } - - false } - }; + + // Add default settings + for key in DEFAULT_SETTINGS { + insert_keys.push(ConfigKey::from(*key)); + } + } // Download webadmin if missing if let Some(blob_store) = config @@ -393,10 +438,9 @@ impl BootManager { ); // Webadmin auto-update - if update_webadmin - || config - .property_or_default::("webadmin.auto-update", "false") - .unwrap_or_default() + if config + .property_or_default::("webadmin.auto-update", "false") + .unwrap_or_default() { if let Err(err) = data.webadmin.update(&core).await { trc::event!( diff --git a/crates/http/src/management/queue.rs b/crates/http/src/management/queue.rs index d14a5c2d..3bc2d459 100644 --- a/crates/http/src/management/queue.rs +++ b/crates/http/src/management/queue.rs @@ -67,9 +67,8 @@ pub struct Message { #[derive(Debug, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub struct Recipient { pub address: String, - + pub queue: String, pub status: Status, - pub retry_num: u32, #[serde(skip_serializing_if = "Option::is_none")] @@ -295,7 +294,7 @@ impl QueueManagement for Server { Status::Scheduled | Status::TemporaryFailure(_) ) && item .as_ref() - .is_none_or(|item| recipient.address_lcase.contains(item)) + .is_none_or(|item| recipient.address.contains(item)) { recipient.retry.due = time; if recipient @@ -384,7 +383,7 @@ impl QueueManagement for Server { if let Some(item) = params.get("filter") { // Cancel delivery for all recipients that match for rcpt in &mut message.message.recipients { - if rcpt.address_lcase.contains(item) { + if rcpt.address.contains(item) { rcpt.status = Status::PermanentFailure(ErrorDetails { entity: "localhost".to_string(), details: queue::Error::Io("Delivery canceled.".to_string()), @@ -587,6 +586,7 @@ impl Message { .iter() .map(|rcpt| Recipient { address: rcpt.address.to_string(), + queue: rcpt.queue.to_string(), status: match &rcpt.status { ArchivedStatus::Scheduled => Status::Scheduled, ArchivedStatus::Completed(status) => { @@ -631,6 +631,7 @@ async fn fetch_queued_messages( params: &UrlParams<'_>, tenant_domains: &Option>, ) -> trc::Result { + let queue = params.get("queue").and_then(QueueName::new); let text = params.get("text"); let from = params.get("from"); let to = params.get("to"); @@ -655,8 +656,12 @@ async fn fetch_queued_messages( }; let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(range_start))); let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(range_end))); - let has_filters = - text.is_some() || from.is_some() || to.is_some() || before.is_some() || after.is_some(); + let has_filters = text.is_some() + || from.is_some() + || to.is_some() + || before.is_some() + || after.is_some() + || queue.is_some(); let mut offset = page.saturating_sub(1) * limit; let mut total_returned = 0; @@ -680,27 +685,28 @@ async fn fetch_queued_messages( .as_ref() .map(|text| { message.return_path.contains(text) - || message - .recipients - .iter() - .any(|r| r.address_lcase.contains(text)) + || message.recipients.iter().any(|r| r.address.contains(text)) }) .unwrap_or_else(|| { from.as_ref() .is_none_or(|from| message.return_path.contains(from)) && to.as_ref().is_none_or(|to| { - message - .recipients - .iter() - .any(|r| r.address_lcase.contains(to)) + message.recipients.iter().any(|r| r.address.contains(to)) }) }) - && before + && before.as_ref().is_none_or(|before| { + message + .next_delivery_event(queue) + .is_some_and(|next| next < *before) + }) + && after.as_ref().is_none_or(|after| { + message + .next_delivery_event(queue) + .is_some_and(|next| next > *after) + }) + && queue .as_ref() - .is_none_or(|before| message.next_delivery_event() < *before) - && after - .as_ref() - .is_none_or(|after| message.next_delivery_event() > *after))); + .is_none_or(|q| message.recipients.iter().any(|r| &r.queue == q)))); if matches { if offset == 0 { diff --git a/crates/jmap/src/submission/get.rs b/crates/jmap/src/submission/get.rs index 62c73be3..06cc58cb 100644 --- a/crates/jmap/src/submission/get.rs +++ b/crates/jmap/src/submission/get.rs @@ -113,7 +113,7 @@ impl EmailSubmissionGet for Server { .unarchive::() .caused_by(trc::location!())?; for rcpt in queued_message.recipients.iter() { - *delivery_status.get_mut_or_insert(rcpt.address_lcase.to_string()) = + *delivery_status.get_mut_or_insert(rcpt.address.to_string()) = DeliveryStatus { smtp_reply: match &rcpt.status { ArchivedStatus::Completed(reply) => { diff --git a/crates/migration/src/lib.rs b/crates/migration/src/lib.rs index 77784f42..6b6b7970 100644 --- a/crates/migration/src/lib.rs +++ b/crates/migration/src/lib.rs @@ -10,14 +10,16 @@ use crate::{ tasks::migrate_tasks_v011, }; use changelog::reset_changelog; -use common::{DATABASE_SCHEMA_VERSION, KV_LOCK_HOUSEKEEPER, Server}; +use common::{ + DATABASE_SCHEMA_VERSION, KV_LOCK_HOUSEKEEPER, Server, manager::boot::DEFAULT_SETTINGS, +}; use jmap_proto::types::{collection::Collection, property::Property}; use principal::{migrate_principal, migrate_principals}; use report::migrate_reports; use std::time::Duration; use store::{ Deserialize, IterateParams, SUBSPACE_PROPERTY, SUBSPACE_QUEUE_MESSAGE, SUBSPACE_REPORT_IN, - SUBSPACE_REPORT_OUT, SerializeInfallible, U32_LEN, Value, ValueKey, + SUBSPACE_REPORT_OUT, SUBSPACE_SETTINGS, SerializeInfallible, U32_LEN, Value, ValueKey, dispatch::{DocumentSet, lookup::KeyValue}, rand::{self, seq::SliceRandom}, write::{AnyClass, AnyKey, BatchBuilder, ValueClass, key::DeserializeBigEndian}, @@ -65,7 +67,7 @@ pub async fn try_migrate(server: &Server) -> trc::Result<()> { return Ok(()); } - match server + let add_v013_config = match server .store() .get_value::(AnyKey { subspace: SUBSPACE_PROPERTY, @@ -81,11 +83,13 @@ pub async fn try_migrate(server: &Server) -> trc::Result<()> { migrate_v0_12(server, true) .await .caused_by(trc::location!())?; + true } Some(2) => { migrate_v0_12(server, false) .await .caused_by(trc::location!())?; + true } Some(version) => { panic!( @@ -96,9 +100,12 @@ pub async fn try_migrate(server: &Server) -> trc::Result<()> { _ => { if !is_new_install(server).await.caused_by(trc::location!())? { migrate_v0_11(server).await.caused_by(trc::location!())?; + true + } else { + false } } - } + }; let mut batch = BatchBuilder::new(); batch.set( @@ -108,6 +115,24 @@ pub async fn try_migrate(server: &Server) -> trc::Result<()> { }), DATABASE_SCHEMA_VERSION.serialize(), ); + + if add_v013_config { + for (key, value) in DEFAULT_SETTINGS { + if key + .strip_prefix("queue.") + .is_some_and(|s| !s.starts_with("limiter.") && !s.starts_with("quota.")) + { + batch.set( + ValueClass::Any(AnyClass { + subspace: SUBSPACE_SETTINGS, + key: key.as_bytes().to_vec(), + }), + value.as_bytes().to_vec(), + ); + } + } + } + server .store() .write(batch.build_all()) diff --git a/crates/migration/src/queue.rs b/crates/migration/src/queue.rs index 5020d46c..44c19630 100644 --- a/crates/migration/src/queue.rs +++ b/crates/migration/src/queue.rs @@ -244,17 +244,14 @@ where Message { created: message.created, blob_hash: message.blob_hash, - return_path: message.return_path, - return_path_lcase: message.return_path_lcase, - return_path_domain: message.return_path_domain, + return_path: message.return_path_lcase, recipients: message .recipients .into_iter() .map(|r| { let domain = &domains[r.domain_idx.as_u64() as usize]; Recipient { - address: r.address, - address_lcase: r.address_lcase, + address: r.address_lcase, status: match r.status { Status::Scheduled => match &domain.status { Status::Scheduled | Status::Completed(_) => Status::Scheduled, diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index 9a54ae8c..dc98e643 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -709,9 +709,7 @@ impl Session { .map_or(0, |d| d.as_secs()); let mut message = Message { created, - return_path: mail_from.address, - return_path_lcase: mail_from.address_lcase, - return_path_domain: mail_from.domain, + return_path: mail_from.address_lcase, recipients: Vec::with_capacity(rcpt_to.len()), flags: mail_from.flags, priority: self.data.priority, @@ -728,8 +726,7 @@ impl Session { rcpt_to.sort_unstable(); for rcpt in rcpt_to { message.recipients.push(queue::Recipient { - address: rcpt.address, - address_lcase: rcpt.address_lcase, + address: rcpt.address_lcase, status: queue::Status::Scheduled, flags: if rcpt.flags & (RCPT_NOTIFY_DELAY diff --git a/crates/smtp/src/outbound/delivery.rs b/crates/smtp/src/outbound/delivery.rs index fb87042b..38863edc 100644 --- a/crates/smtp/src/outbound/delivery.rs +++ b/crates/smtp/src/outbound/delivery.rs @@ -74,7 +74,7 @@ impl QueuedMessage { Status::Scheduled | Status::TemporaryFailure(_) ) && r.queue == message.queue_name { - Some(trc::Value::String(r.address_lcase.as_str().into())) + Some(trc::Value::String(r.address.as_str().into())) } else { None } @@ -236,7 +236,7 @@ impl QueuedMessage { ); routes - .entry((rcpt.address_lcase.domain_part(), route)) + .entry((rcpt.address.domain_part(), route)) .or_default() .push(rcpt_idx); } @@ -1325,7 +1325,7 @@ impl MessageWrapper { SpanId = self.span_id, QueueId = self.queue_id, QueueName = self.queue_name.as_str().to_string(), - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Reason = from_error_details(&err.details), Details = trc::Value::Timestamp(now), Expires = rcpt @@ -1344,7 +1344,7 @@ impl MessageWrapper { SpanId = self.span_id, QueueId = self.queue_id, QueueName = self.queue_name.as_str().to_string(), - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Reason = "Message expired without any delivery attempts made.", Details = trc::Value::Timestamp(now), Expires = rcpt @@ -1355,7 +1355,7 @@ impl MessageWrapper { ); rcpt.status = Status::PermanentFailure(ErrorDetails { - entity: rcpt.address_lcase.domain_part().to_string(), + entity: rcpt.address.domain_part().to_string(), details: Error::Io( "Message expired without any delivery attempts made.".into(), ), diff --git a/crates/smtp/src/outbound/local.rs b/crates/smtp/src/outbound/local.rs index 919af0d2..04d03f21 100644 --- a/crates/smtp/src/outbound/local.rs +++ b/crates/smtp/src/outbound/local.rs @@ -7,9 +7,9 @@ use crate::{ outbound::DeliveryResult, queue::{ - DomainPart, Error, ErrorDetails, FROM_AUTHENTICATED, FROM_UNAUTHENTICATED_DMARC, - HostResponse, MessageSource, MessageWrapper, Status, UnexpectedResponse, - quota::HasQueueQuota, spool::SmtpSpool, + Error, ErrorDetails, FROM_AUTHENTICATED, FROM_UNAUTHENTICATED_DMARC, HostResponse, + MessageSource, MessageWrapper, Status, UnexpectedResponse, quota::HasQueueQuota, + spool::SmtpSpool, }, reporting::SmtpReporting, }; @@ -29,7 +29,7 @@ impl MessageWrapper { let mut pending_recipients = Vec::new(); let mut recipient_addresses = Vec::new(); for &rcpt_idx in rcpt_idxs { - let rcpt_addr = &self.message.recipients[rcpt_idx].address_lcase; + let rcpt_addr = &self.message.recipients[rcpt_idx].address; recipient_addresses.push(rcpt_addr.clone()); pending_recipients.push((rcpt_idx, rcpt_addr)); } @@ -37,7 +37,7 @@ impl MessageWrapper { // Deliver message let delivery_result = server .deliver_message(IngestMessage { - sender_address: self.message.return_path_lcase.clone(), + sender_address: self.message.return_path.clone(), sender_authenticated: self.message.flags & (FROM_UNAUTHENTICATED_DMARC | FROM_AUTHENTICATED) != 0, @@ -93,15 +93,8 @@ impl MessageWrapper { // Process autogenerated messages for autogenerated in delivery_result.autogenerated { - let from_addr_lcase = autogenerated.sender_address.to_lowercase(); - let from_addr_domain = from_addr_lcase.domain_part().to_string(); - - let mut message = server.new_message( - autogenerated.sender_address, - from_addr_lcase, - from_addr_domain, - self.span_id, - ); + let mut message = + server.new_message(autogenerated.sender_address.to_lowercase(), self.span_id); for rcpt in autogenerated.recipients { message.add_recipient(rcpt, server).await; } @@ -132,12 +125,12 @@ impl MessageWrapper { trc::event!( Sieve(SieveEvent::QuotaExceeded), SpanId = self.span_id, - From = message.message.return_path_lcase, + From = message.message.return_path, To = message .message .recipients .into_iter() - .map(|r| trc::Value::from(r.address_lcase)) + .map(|r| trc::Value::from(r.address)) .collect::>(), ); } diff --git a/crates/smtp/src/queue/dsn.rs b/crates/smtp/src/queue/dsn.rs index 2198974a..035bc5fe 100644 --- a/crates/smtp/src/queue/dsn.rs +++ b/crates/smtp/src/queue/dsn.rs @@ -37,13 +37,9 @@ impl SendDsn for Server { if !message.message.return_path.is_empty() { // Build DSN if let Some(dsn) = message.build_dsn(self).await { - let mut dsn_message = self.new_message("", "", "", message.span_id); + let mut dsn_message = self.new_message("", message.span_id); dsn_message - .add_recipient_parts( - message.message.return_path.as_str(), - message.message.return_path_lcase.as_str(), - self, - ) + .add_recipient_parts(message.message.return_path.as_str(), self) .await; // Sign message @@ -81,7 +77,7 @@ impl SendDsn for Server { trc::event!( Delivery(trc::DeliveryEvent::DsnSuccess), SpanId = message.span_id, - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Hostname = response.hostname.clone(), Code = response.response.code, Details = response.response.message.to_string(), @@ -91,7 +87,7 @@ impl SendDsn for Server { trc::event!( Delivery(trc::DeliveryEvent::DsnTempFail), SpanId = message.span_id, - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Hostname = response.entity.clone(), Details = response.details.to_string(), NextRetry = trc::Value::Timestamp(rcpt.retry.due), @@ -105,7 +101,7 @@ impl SendDsn for Server { trc::event!( Delivery(trc::DeliveryEvent::DsnPermFail), SpanId = message.span_id, - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Hostname = response.entity.clone(), Details = response.details.to_string(), Total = rcpt.retry.inner, @@ -115,7 +111,7 @@ impl SendDsn for Server { trc::event!( Delivery(trc::DeliveryEvent::DsnTempFail), SpanId = message.span_id, - To = rcpt.address_lcase.clone(), + To = rcpt.address.clone(), Details = "Concurrency limited", NextRetry = trc::Value::Timestamp(rcpt.retry.due), Expires = rcpt diff --git a/crates/smtp/src/queue/mod.rs b/crates/smtp/src/queue/mod.rs index dde15e78..81f0f2b3 100644 --- a/crates/smtp/src/queue/mod.rs +++ b/crates/smtp/src/queue/mod.rs @@ -53,14 +53,12 @@ pub struct Message { pub created: u64, pub blob_hash: BlobHash, + pub return_path: String, + pub recipients: Vec, + pub received_from_ip: IpAddr, pub received_via_port: u16, - pub return_path: String, - pub return_path_lcase: String, - pub return_path_domain: String, - pub recipients: Vec, - pub flags: u64, pub env_id: Option, pub priority: i16, @@ -105,7 +103,6 @@ pub enum QuotaKey { )] pub struct Recipient { pub address: String, - pub address_lcase: String, pub retry: Schedule, pub notify: Schedule, @@ -268,7 +265,7 @@ impl<'x> QueueEnvelope<'x> { pub fn new(message: &'x Message, rcpt: &'x Recipient) -> Self { Self { message, - domain: rcpt.address_lcase.domain_part(), + domain: rcpt.address.domain_part(), rcpt, mx: "", remote_ip: IpAddr::V4(Ipv4Addr::new(0, 0, 0, 0)), @@ -280,22 +277,24 @@ impl<'x> QueueEnvelope<'x> { impl<'x> ResolveVariable for QueueEnvelope<'x> { fn resolve_variable(&self, variable: u32) -> expr::Variable<'x> { match variable { - V_SENDER => self.message.return_path_lcase.as_str().into(), - V_SENDER_DOMAIN => self.message.return_path_domain.as_str().into(), + V_SENDER => self.message.return_path.as_str().into(), + V_SENDER_DOMAIN => self.message.return_path.domain_part().into(), V_RECIPIENT_DOMAIN => self.domain.into(), - V_RECIPIENT => self.rcpt.address_lcase.as_str().into(), + V_RECIPIENT => self.rcpt.address.as_str().into(), V_RECIPIENTS => self .message .recipients .iter() - .map(|r| Variable::from(r.address_lcase.as_str())) + .map(|r| Variable::from(r.address.as_str())) .collect::>() .into(), V_QUEUE_RETRY_NUM => self.rcpt.retry.inner.into(), V_QUEUE_NOTIFY_NUM => self.rcpt.notify.inner.into(), V_QUEUE_EXPIRES_IN => match &self.rcpt.expires { QueueExpiry::Ttl(time) => (*time + self.message.created).saturating_sub(now()), - QueueExpiry::Attempts(count) => (*count) as u64, + QueueExpiry::Attempts(count) => { + (count.saturating_sub(self.rcpt.retry.inner)) as u64 + } } .into(), V_QUEUE_LAST_STATUS => self.rcpt.status.to_compact_string().into(), @@ -353,12 +352,12 @@ impl<'x> ResolveVariable for QueueEnvelope<'x> { impl ResolveVariable for Message { fn resolve_variable(&self, variable: u32) -> expr::Variable<'_> { match variable { - V_SENDER => self.return_path_lcase.as_str().into(), - V_SENDER_DOMAIN => self.return_path_domain.as_str().into(), + V_SENDER => self.return_path.as_str().into(), + V_SENDER_DOMAIN => self.return_path.domain_part().into(), V_RECIPIENTS => self .recipients .iter() - .map(|r| Variable::from(r.address_lcase.as_str())) + .map(|r| Variable::from(r.address.as_str())) .collect::>() .into(), V_PRIORITY => self.priority.into(), diff --git a/crates/smtp/src/queue/quota.rs b/crates/smtp/src/queue/quota.rs index 5363f40e..cb9e21de 100644 --- a/crates/smtp/src/queue/quota.rs +++ b/crates/smtp/src/queue/quota.rs @@ -64,7 +64,7 @@ impl HasQueueQuota for Server { let mut seen_domains = AHashSet::new(); for quota in &self.core.smtp.queue.quota.rcpt_domain { for (rcpt_idx, rcpt) in message.message.recipients.iter().enumerate() { - if seen_domains.insert(rcpt.address_lcase.domain_part()) + if seen_domains.insert(rcpt.address.domain_part()) && !self .check_quota( quota, @@ -192,7 +192,7 @@ impl MessageWrapper { &rcpt.status, Status::Completed(_) | Status::PermanentFailure(_) ) { - if seen_domains.insert(rcpt.address_lcase.domain_part()) { + if seen_domains.insert(rcpt.address.domain_part()) { quota_ids.push(((pos + 1) as u64) << 32); } quota_ids.push((pos + 1) as u64); diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index d3c1ec16..19cab484 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -39,13 +39,7 @@ pub struct QueuedMessages { } pub trait SmtpSpool: Sync + Send { - fn new_message( - &self, - return_path: impl Into, - return_path_lcase: impl Into, - return_path_domain: impl Into, - span_id: u64, - ) -> MessageWrapper; + fn new_message(&self, return_path: impl Into, span_id: u64) -> MessageWrapper; fn next_event(&self, queue: &mut Queue) -> impl Future + Send; @@ -74,13 +68,7 @@ pub trait SmtpSpool: Sync + Send { } impl SmtpSpool for Server { - fn new_message( - &self, - return_path: impl Into, - return_path_lcase: impl Into, - return_path_domain: impl Into, - span_id: u64, - ) -> MessageWrapper { + fn new_message(&self, return_path: impl Into, span_id: u64) -> MessageWrapper { let created = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) .map_or(0, |d| d.as_secs()); @@ -93,8 +81,6 @@ impl SmtpSpool for Server { message: Message { created, return_path: return_path.into(), - return_path_lcase: return_path_lcase.into(), - return_path_domain: return_path_domain.into(), recipients: Vec::with_capacity(1), flags: 0, env_id: None, @@ -389,7 +375,7 @@ impl MessageWrapper { .message .recipients .iter() - .map(|r| trc::Value::String(r.address_lcase.as_str().into())) + .map(|r| trc::Value::String(r.address.as_str().into())) .collect::>(), Size = self.message.size, NextRetry = self @@ -492,16 +478,10 @@ impl MessageWrapper { true } - pub async fn add_recipient_parts( - &mut self, - rcpt: impl Into, - rcpt_lcase: impl Into, - server: &Server, - ) { + pub async fn add_recipient_parts(&mut self, rcpt: impl Into, server: &Server) { // Resolve queue self.message.recipients.push(Recipient { address: rcpt.into(), - address_lcase: rcpt_lcase.into(), status: Status::Scheduled, flags: 0, orcpt: None, @@ -530,10 +510,9 @@ impl MessageWrapper { recipient.queue = queue.virtual_queue; } - pub async fn add_recipient(&mut self, rcpt: impl Into, server: &Server) { - let rcpt = rcpt.into(); - let rcpt_lcase = rcpt.to_lowercase(); - self.add_recipient_parts(rcpt, rcpt_lcase, server).await; + pub async fn add_recipient(&mut self, rcpt: impl AsRef, server: &Server) { + let rcpt = rcpt.as_ref().to_lowercase(); + self.add_recipient_parts(rcpt, server).await; } pub async fn save_changes(mut self, server: &Server, prev_event: Option) -> bool { @@ -691,7 +670,7 @@ impl MessageWrapper { pub fn has_domain(&self, domains: &[String]) -> bool { self.message.recipients.iter().any(|r| { - let domain = r.address_lcase.domain_part(); + let domain = r.address.domain_part(); domains.iter().any(|dd| dd == domain) }) || self .message @@ -704,7 +683,7 @@ impl MessageWrapper { impl ArchivedMessage { pub fn has_domain(&self, domains: &[String]) -> bool { self.recipients.iter().any(|r| { - let domain = r.address_lcase.domain_part(); + let domain = r.address.domain_part(); domains.iter().any(|dd| dd == domain) }) || self .return_path @@ -712,22 +691,22 @@ impl ArchivedMessage { .is_some_and(|(_, domain)| domains.iter().any(|dd| dd == domain)) } - pub fn next_delivery_event(&self) -> u64 { - let mut next_delivery = now(); + pub fn next_delivery_event(&self, queue: Option) -> Option { + let mut next_delivery = None; - for (pos, rcpt) in self - .recipients - .iter() - .filter(|d| { - matches!( - d.status, - ArchivedStatus::Scheduled | ArchivedStatus::TemporaryFailure(_) - ) - }) - .enumerate() - { - if pos == 0 || rcpt.retry.due < next_delivery { - next_delivery = rcpt.retry.due.into(); + for rcpt in self.recipients.iter().filter(|d| { + matches!( + d.status, + ArchivedStatus::Scheduled | ArchivedStatus::TemporaryFailure(_) + ) && queue.is_none_or(|q| d.queue == q) + }) { + let retry_due = rcpt.retry.due.to_native(); + if let Some(next_delivery) = &mut next_delivery { + if retry_due < *next_delivery { + *next_delivery = retry_due; + } + } else { + next_delivery = Some(retry_due); } } diff --git a/crates/smtp/src/reporting/mod.rs b/crates/smtp/src/reporting/mod.rs index b82a7a40..ef1aa953 100644 --- a/crates/smtp/src/reporting/mod.rs +++ b/crates/smtp/src/reporting/mod.rs @@ -7,7 +7,7 @@ use crate::{ core::Session, inbound::DkimSign, - queue::{DomainPart, MessageSource, MessageWrapper, spool::SmtpSpool}, + queue::{MessageSource, MessageWrapper, spool::SmtpSpool}, }; use common::{ Server, USER_AGENT, @@ -83,8 +83,8 @@ pub trait SmtpReporting: Sync + Send { fn send_autogenerated( &self, - from_addr: impl Into + Sync + Send, - rcpts: impl Iterator + Sync + Send> + Sync + Send, + from_addr: impl AsRef + Sync + Send, + rcpts: impl Iterator + Sync + Send> + Sync + Send, raw_message: Vec, sign_config: Option<&IfBlock>, parent_session_id: u64, @@ -114,14 +114,7 @@ impl SmtpReporting for Server { parent_session_id: u64, ) { // Build message - let from_addr_lcase = from_addr.to_lowercase(); - let from_addr_domain = from_addr_lcase.domain_part().to_string(); - let mut message = self.new_message( - from_addr, - from_addr_lcase, - from_addr_domain, - parent_session_id, - ); + let mut message = self.new_message(from_addr.to_lowercase(), parent_session_id); for rcpt_ in rcpts { message.add_recipient(rcpt_.as_ref(), self).await; } @@ -161,22 +154,14 @@ impl SmtpReporting for Server { async fn send_autogenerated( &self, - from_addr: impl Into + Sync + Send, - rcpts: impl Iterator + Sync + Send> + Sync + Send, + from_addr: impl AsRef + Sync + Send, + rcpts: impl Iterator + Sync + Send> + Sync + Send, raw_message: Vec, sign_config: Option<&IfBlock>, parent_session_id: u64, ) { // Build message - let from_addr = from_addr.into(); - let from_addr_lcase = from_addr.to_lowercase(); - let from_addr_domain = from_addr_lcase.domain_part().to_string(); - let mut message = self.new_message( - from_addr, - from_addr_lcase, - from_addr_domain, - parent_session_id, - ); + let mut message = self.new_message(from_addr.as_ref().to_lowercase(), parent_session_id); for rcpt in rcpts { message.add_recipient(rcpt, self).await; } diff --git a/crates/smtp/src/scripts/event_loop.rs b/crates/smtp/src/scripts/event_loop.rs index 40042636..8d1cf947 100644 --- a/crates/smtp/src/scripts/event_loop.rs +++ b/crates/smtp/src/scripts/event_loop.rs @@ -4,10 +4,11 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use std::{borrow::Cow, future::Future, sync::Arc, time::Instant}; - +use crate::{ + inbound::DkimSign, + queue::{MessageSource, quota::HasQueueQuota, spool::SmtpSpool}, +}; use common::{Server, config::smtp::queue::QueueExpiry, scripts::plugins::PluginContext}; - use mail_auth::common::headers::HeaderWriter; use mail_parser::{Encoding, Message, MessagePart, PartType}; use sieve::{ @@ -18,13 +19,9 @@ use smtp_proto::{ MAIL_BY_TRACE, MAIL_RET_FULL, MAIL_RET_HDRS, RCPT_NOTIFY_DELAY, RCPT_NOTIFY_FAILURE, RCPT_NOTIFY_NEVER, RCPT_NOTIFY_SUCCESS, }; +use std::{borrow::Cow, future::Future, sync::Arc, time::Instant}; use trc::SieveEvent; -use crate::{ - inbound::DkimSign, - queue::{DomainPart, MessageSource, quota::HasQueueQuota, spool::SmtpSpool}, -}; - use super::{ScriptModification, ScriptParameters, ScriptResult}; pub trait RunScript: Sync + Send { @@ -161,14 +158,8 @@ impl RunScript for Server { message_id, } => { // Build message - let return_path_lcase = params.return_path.to_lowercase(); - let return_path_domain = return_path_lcase.domain_part().to_string(); - let mut message = self.new_message( - params.return_path.clone(), - return_path_lcase, - return_path_domain, - session_id, - ); + let mut message = + self.new_message(params.return_path.to_lowercase(), session_id); match recipient { Recipient::Address(rcpt) => { message.add_recipient(rcpt, self).await; @@ -331,12 +322,12 @@ impl RunScript for Server { Sieve(SieveEvent::QuotaExceeded), SpanId = session_id, Id = script_id.clone(), - From = message.message.return_path_lcase, + From = message.message.return_path, To = message .message .recipients .into_iter() - .map(|r| trc::Value::from(r.address_lcase)) + .map(|r| trc::Value::from(r.address)) .collect::>(), ); } diff --git a/tests/src/smtp/lookup/utils.rs b/tests/src/smtp/lookup/utils.rs index a77efec1..cdfa1aa0 100644 --- a/tests/src/smtp/lookup/utils.rs +++ b/tests/src/smtp/lookup/utils.rs @@ -147,11 +147,8 @@ async fn strategies() { received_from_ip: "1.2.3.4".parse().unwrap(), received_via_port: 7911, return_path: "test@example.com".to_string(), - return_path_lcase: "test@example.com".to_string(), - return_path_domain: "example.com".to_string(), recipients: vec![Recipient { address: "recipient@foobar.com".to_string(), - address_lcase: "recipient@foobar.com".to_string(), retry: Schedule::now(), notify: Schedule::now(), expires: QueueExpiry::Ttl(3600), diff --git a/tests/src/smtp/outbound/throttle.rs b/tests/src/smtp/outbound/throttle.rs index bf66284b..34e28415 100644 --- a/tests/src/smtp/outbound/throttle.rs +++ b/tests/src/smtp/outbound/throttle.rs @@ -68,7 +68,7 @@ async fn throttle_outbound() { // Build test message let mut test_message = new_message(0).message; - test_message.return_path_domain = "foobar.org".into(); + test_message.return_path = "test@foobar.org".into(); test_message .recipients .push(build_rcpt("bill@test.org", 0, 0, 0)); @@ -110,7 +110,7 @@ async fn throttle_outbound() { local.queue_receiver.read_event().await.assert_on_hold();*/ // Expect rate limit throttle for sender domain 'foobar.net' - test_message.return_path_domain = "foobar.net".into(); + test_message.return_path = "test@foobar.net".into(); for t in &throttle.sender { core.is_allowed( t, @@ -136,7 +136,7 @@ async fn throttle_outbound() { assert!(due > 0, "Due: {}", due); // Expect concurrency throttle for recipient domain 'example.org' - test_message.return_path_domain = "test.net".into(); + test_message.return_path = "test@test.net".into(); test_message .recipients .push(build_rcpt("test@example.org", 0, 0, 0)); @@ -286,7 +286,7 @@ impl<'x> TestQueueEnvelope<'x> for QueueEnvelope<'x> { mx, remote_ip: IpAddr::V4(Ipv4Addr::new(0, 0, 0, 0)), local_ip: IpAddr::V4(Ipv4Addr::new(0, 0, 0, 0)), - domain: rcpt.address_lcase.domain_part(), + domain: rcpt.address.domain_part(), rcpt, } } diff --git a/tests/src/smtp/queue/dsn.rs b/tests/src/smtp/queue/dsn.rs index ef3b1f67..46fe1ba3 100644 --- a/tests/src/smtp/queue/dsn.rs +++ b/tests/src/smtp/queue/dsn.rs @@ -64,11 +64,8 @@ async fn generate_dsn() { .duration_since(SystemTime::UNIX_EPOCH) .map_or(0, |d| d.as_secs()), return_path: "sender@foobar.org".into(), - return_path_lcase: "".into(), - return_path_domain: "foobar.org".into(), recipients: vec![Recipient { address: "foobar@example.org".into(), - address_lcase: "foobar@example.org".into(), status: Status::PermanentFailure(ErrorDetails { entity: "mx.example.org".into(), details: Error::UnexpectedResponse(UnexpectedResponse { @@ -125,7 +122,6 @@ async fn generate_dsn() { // Success DSN message.message.recipients.push(Recipient { address: "jane@example.org".into(), - address_lcase: "jane@example.org".into(), status: Status::Completed(HostResponse { hostname: "mx2.example.org".into(), response: Response { @@ -148,7 +144,6 @@ async fn generate_dsn() { // Delay DSN message.message.recipients.push(Recipient { address: "john.doe@example.org".into(), - address_lcase: "john.doe@example.org".into(), status: Status::TemporaryFailure(ErrorDetails { entity: "mx.domain.org".into(), details: Error::ConnectionError("Connection timeout".into()), diff --git a/tests/src/smtp/queue/manager.rs b/tests/src/smtp/queue/manager.rs index 1fe3ff0f..8bc9d9e0 100644 --- a/tests/src/smtp/queue/manager.rs +++ b/tests/src/smtp/queue/manager.rs @@ -171,8 +171,6 @@ pub fn new_message(queue_id: u64) -> MessageWrapper { size: 0, created: now(), return_path: "sender@foobar.org".into(), - return_path_lcase: "".into(), - return_path_domain: "foobar.org".into(), recipients: vec![], flags: 0, env_id: None, @@ -222,14 +220,14 @@ impl TestMessage for Message { fn rcpt(&self, name: &str) -> &Recipient { self.recipients .iter() - .find(|d| d.address_lcase == name) + .find(|d| d.address == name) .unwrap_or_else(|| panic!("Expected rcpt {name} not found in {:?}", self.recipients)) } fn rcpt_mut(&mut self, name: &str) -> &mut Recipient { self.recipients .iter_mut() - .find(|d| d.address_lcase == name) + .find(|d| d.address == name) .unwrap() } } diff --git a/tests/src/smtp/queue/mod.rs b/tests/src/smtp/queue/mod.rs index 418edac1..41c3d869 100644 --- a/tests/src/smtp/queue/mod.rs +++ b/tests/src/smtp/queue/mod.rs @@ -24,7 +24,6 @@ pub mod virtualq; pub fn build_rcpt(address: &str, retry: u64, notify: u64, expires: u64) -> Recipient { Recipient { address: address.to_string(), - address_lcase: address.to_string(), retry: Schedule::later(retry), notify: Schedule::later(notify), expires: QueueExpiry::Ttl(expires),