/* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ use super::AuthResult; use crate::{ core::{Session, SessionAddress, State}, inbound::{dkim::DkimSign, milter::Modification}, queue::{ self, Message, MessageSource, MessageWrapper, QueueEnvelope, RCPT_SPAM_MASK, quota::HasQueueQuota, rcpt_spam_flag, spool::QueueParams, }, reporting::analysis::{AnalyzeReport, ReportData}, scripts::ScriptResult, }; use common::{ config::{ mailstore::spamfilter::{SpamFilterAction, spam_status}, smtp::{ auth::VerifyStrategy, queue::{QueueExpiry, QueueName}, session::Stage, }, }, network::SessionStream, scripts::ScriptModification, }; use mail_auth::{ AuthenticatedMessage, AuthenticationResults, Dkim2Result, DkimResult, DmarcResult, ReceivedSpf, SpfOutput, SpfResult, common::{ crypto::Algorithm, headers::{Header, HeaderWriter}, verify::VerifySignature, }, dkim::DkimError, dkim2::{Dkim2Dsn, Dkim2DsnFailure, Envelope as Dkim2Envelope}, dmarc::{self, verify::DmarcParameters}, }; use mail_builder::headers::{date::Date, message_id::generate_message_id_header}; use mail_parser::{MessageParser, MimeHeaders, parsers::fields::thread::thread_name}; use registry::schema::structs::Rate; use sieve::runtime::Variable; use smtp_proto::{ MAIL_BY_RETURN, RCPT_NOTIFY_DELAY, RCPT_NOTIFY_FAILURE, RCPT_NOTIFY_NEVER, RCPT_NOTIFY_SUCCESS, }; use std::{ borrow::Cow, time::{Instant, SystemTime}, }; use trc::{SmtpEvent, SpamEvent}; use utils::DomainPart; impl Session { pub async fn queue_message(&mut self) -> Cow<'static, [u8]> { // Parse message let raw_message = std::mem::take(&mut self.data.message); let parsed_message = match MessageParser::new() .parse(&raw_message) .filter(|p| p.headers().iter().any(|h| !h.name.is_other())) { Some(parsed_message) => parsed_message, None => { trc::event!( Smtp(SmtpEvent::MessageParseFailed), SpanId = self.data.session_id, ); return (&b"550 5.7.7 Failed to parse message.\r\n"[..]).into(); } }; // Authenticate message let mut auth_message = AuthenticatedMessage::from_parsed( &parsed_message, &raw_message, self.server.core.smtp.mail_auth.dkim.strict, ); let has_date_header = auth_message.has_date_header(); let has_message_id_header = auth_message.has_message_id_header(); // Loop detection let dc = &self.server.core.smtp.session.data; let ac = &self.server.core.smtp.mail_auth; let rc = &self.server.core.smtp.report; if auth_message.received_headers_count() > self .server .eval_if(&dc.max_received_headers, self, self.data.session_id) .await .unwrap_or(50) { trc::event!( Smtp(SmtpEvent::LoopDetected), SpanId = self.data.session_id, Total = auth_message.received_headers_count(), ); return (&b"450 4.4.6 Too many Received headers. Possible loop detected.\r\n"[..]) .into(); } // Verify DKIM let dkim = self .server .eval_if(&ac.dkim.verify, self, self.data.session_id) .await .unwrap_or(VerifyStrategy::Relaxed); let dmarc = self .server .eval_if(&ac.dmarc.verify, self, self.data.session_id) .await .unwrap_or(VerifyStrategy::Relaxed); let dkim_output = if dkim.verify() || dmarc.verify() { // Remove insecure DKIM signatures before verification auth_message.dkim_headers.retain(|header| { let signature = &header.header; if signature.algorithm() == Algorithm::RsaSha1 || (signature.algorithm() == Algorithm::RsaSha256 && signature.b.len() < 128) { auth_message.errors.push(Header { name: header.name, value: header.value, header: mail_auth::Error::Dkim(DkimError::UnsupportedAlgorithm), }); auth_message.has_dkim_errors = true; false } else { true } }); let time = Instant::now(); let dkim_output = self .server .core .smtp .resolvers .dns .verify_dkim(self.server.inner.cache.build_auth_parameters(&auth_message)) .await; let pass = dkim_output .iter() .any(|d| matches!(d.result(), DkimResult::Pass)); let strict = dkim.is_strict(); let rejected = strict && !pass; // Send reports for failed signatures if let Some(rate) = self .server .eval_if::(&rc.dkim.send, self, self.data.session_id) .await { for output in &dkim_output { if let Some(rcpt) = output.failure_report_addr() { self.send_dkim_report(rcpt, &auth_message, &rate, rejected, output) .await; } } } trc::event!( Smtp(if pass { SmtpEvent::DkimPass } else { SmtpEvent::DkimFail }), SpanId = self.data.session_id, Strict = strict, Result = dkim_output.iter().map(trc::Error::from).collect::>(), Elapsed = time.elapsed(), ); if rejected { // 'Strict' mode violates the advice of Section 6.1 of RFC6376 return if dkim_output .iter() .any(|d| matches!(d.result(), DkimResult::TempError(_))) { (&b"451 4.7.20 No passing DKIM signatures found.\r\n"[..]).into() } else { (&b"550 5.7.20 No passing DKIM signatures found.\r\n"[..]).into() }; } dkim_output } else { vec![] }; // Verify ARC let arc = self .server .eval_if(&ac.arc.verify, self, self.data.session_id) .await .unwrap_or(VerifyStrategy::Relaxed); let arc_output = if arc.verify() { let time = Instant::now(); let arc_output = self .server .core .smtp .resolvers .dns .verify_arc(self.server.inner.cache.build_auth_parameters(&auth_message)) .await; let strict = arc.is_strict(); let pass = matches!(arc_output.result(), DkimResult::Pass | DkimResult::None); trc::event!( Smtp(if pass { SmtpEvent::ArcPass } else { SmtpEvent::ArcFail }), SpanId = self.data.session_id, Strict = strict, Result = trc::Error::from(arc_output.result()), Elapsed = time.elapsed(), ); if strict && !pass { return if matches!(arc_output.result(), DkimResult::TempError(_)) { (&b"451 4.7.29 ARC validation failed.\r\n"[..]).into() } else { (&b"550 5.7.29 ARC validation failed.\r\n"[..]).into() }; } arc_output.into() } else { None }; // Verify DKIM2 let mail_from = self.data.mail_from.as_ref().unwrap(); let dkim2_output = if (dkim.verify() || dmarc.verify()) && (!auth_message.dkim2_signatures.is_empty() || auth_message.has_dkim2_errors) { // Discard forged DKIM2-signed delivery status notifications if dkim.verify() && !auth_message.dkim2_signatures.is_empty() && let Some(dsn) = parse_dkim2_dsn(&parsed_message, &auth_message, raw_message.as_slice()) && let Err(failure) = self .server .core .smtp .resolvers .dns .verify_dkim2_dsn( self.server.inner.cache.build_auth_parameters(&dsn), Dkim2Envelope { mail_from: &mail_from.address, rcpt_to: self.data.rcpt_to.iter().map(|r| r.address.as_str()), }, ) .await && matches!( failure, Dkim2DsnFailure::DsnChainFailed | Dkim2DsnFailure::ReturnedChainFailed | Dkim2DsnFailure::NotAligned ) { trc::event!( Smtp(SmtpEvent::Dkim2DsnDiscarded), SpanId = self.data.session_id, From = mail_from.address.to_string(), To = self .data .rcpt_to .iter() .map(|rcpt| trc::Value::from(rcpt.address.to_string())) .collect::>(), Reason = failure.to_string(), ); self.data.messages_sent += 1; return (b"550 5.7.1 DSN rejected due to DKIM2 verification failure.\r\n"[..]) .into(); } let time = Instant::now(); let output = self .server .core .smtp .resolvers .dns .verify_dkim2( self.server.inner.cache.build_auth_parameters(&auth_message), Dkim2Envelope { mail_from: &mail_from.address, rcpt_to: self.data.rcpt_to.iter().map(|r| r.address.as_str()), }, ) .await; if !matches!(output.result(), Dkim2Result::None) { trc::event!( Smtp(if matches!(output.result(), Dkim2Result::Pass) { SmtpEvent::Dkim2Pass } else { SmtpEvent::Dkim2Fail }), SpanId = self.data.session_id, Result = trc::Error::from(&output), Elapsed = time.elapsed(), ); } Some(output) } else { None }; // Build authentication results header let mut auth_results = AuthenticationResults::new(&self.hostname); if !dkim_output.is_empty() { auth_results = auth_results.with_dkim_results(&dkim_output, auth_message.from()) } if let Some(dkim2_output) = &dkim2_output && (!matches!(dkim2_output.result(), Dkim2Result::None) || !dkim2_output.chain().is_empty()) { auth_results = auth_results.with_dkim2_result(dkim2_output); } if let Some(spf_ehlo) = &self.data.spf_ehlo { auth_results = auth_results.with_spf_ehlo_result( spf_ehlo, self.data.remote_ip, &self.data.helo_domain, ); } if let Some(spf_mail_from) = &self.data.spf_mail_from { auth_results = auth_results.with_spf_mailfrom_result( spf_mail_from, self.data.remote_ip, &mail_from.address, &self.data.helo_domain, ); } if let Some(iprev) = &self.data.iprev { auth_results = auth_results.with_iprev_result(iprev, self.data.remote_ip); } // Verify DMARC let is_report = !self.is_authenticated() && self.is_report(); let (dmarc_result, dmarc_policy) = if dmarc.verify() { { let synthetic_spf; let spf_output = match &self.data.spf_mail_from { Some(spf_output) => spf_output, None => { synthetic_spf = SpfOutput::new(String::new()).with_result(SpfResult::None); &synthetic_spf } }; let time = Instant::now(); let dmarc_output = self.server .core .smtp .resolvers .dns .verify_dmarc(self.server.inner.cache.build_auth_parameters( DmarcParameters { message: &auth_message, dkim_output: &dkim_output, dkim2_output: dkim2_output.as_ref(), rfc5321_mail_from_domain: if !mail_from.domain.is_empty() { &mail_from.domain } else { &self.data.helo_domain }, spf_output, }, )) .await; let dmarc_result = dmarc_output.result(); let pass = dmarc_result == DmarcResult::Pass; let strict = dmarc.is_strict(); let is_temp_fail = matches!(dmarc_result, DmarcResult::TempError(_)); let rejected = strict && dmarc_output.policy() == dmarc::Policy::Reject && (is_temp_fail || matches!(dmarc_result, DmarcResult::Fail(_))); // Add to DMARC output to the Authentication-Results header auth_results = auth_results.with_dmarc_result(&dmarc_output); let dmarc_policy = dmarc_output.policy(); trc::event!( Smtp(if pass { SmtpEvent::DmarcPass } else { SmtpEvent::DmarcFail }), SpanId = self.data.session_id, Strict = strict, Domain = dmarc_output.domain().to_string(), Policy = dmarc_policy.to_string(), Result = trc::Error::from(&dmarc_result), Elapsed = time.elapsed(), ); // Send DMARC report if dmarc_output.requested_reports() && !is_report && !(rejected && is_temp_fail) { self.send_dmarc_report( &auth_message, &auth_results, rejected, dmarc_output, &dkim_output, dkim2_output.as_ref(), &arc_output, ) .await; } if rejected { return if is_temp_fail { (&b"451 4.7.1 Email temporarily rejected per DMARC policy.\r\n"[..]).into() } else { (&b"550 5.7.1 Email rejected per DMARC policy.\r\n"[..]).into() }; } (dmarc_result.into(), dmarc_policy.into()) } } else { (None, None) }; // Analyze reports if is_report && ReportData::is_present(&parsed_message) { if !rc.analysis.forward { self.data .rcpt_to .retain(|rcpt| !rc.analysis.is_report_address(rcpt.report_address())); } if self.data.rcpt_to.is_empty() { self.server.analyze_report( mail_parser::Message { html_body: parsed_message.html_body, text_body: parsed_message.text_body, attachments: parsed_message.attachments, parts: parsed_message .parts .into_iter() .map(|p| p.into_owned()) .collect(), raw_message: b"".into(), }, self.data.session_id, ); self.data.messages_sent += 1; return (b"250 2.0.0 Message queued for delivery.\r\n"[..]).into(); } else { self.server.analyze_report( mail_parser::Message { html_body: parsed_message.html_body.clone(), text_body: parsed_message.text_body.clone(), attachments: parsed_message.attachments.clone(), parts: parsed_message .parts .iter() .map(|p| p.clone().into_owned()) .collect(), raw_message: b"".into(), }, self.data.session_id, ); } } // Add Received header let message_id = self.server.inner.data.queue_id_gen.generate(); let mut headers = Vec::with_capacity(64); if self .server .eval_if(&dc.add_received, self, self.data.session_id) .await .unwrap_or(true) { self.write_received(&mut headers, message_id) } // Add authentication results header if self .server .eval_if(&dc.add_auth_results, self, self.data.session_id) .await .unwrap_or(true) { auth_results.write_header(&mut headers); } // Add Received-SPF header if let Some(spf_output) = &self.data.spf_mail_from && self .server .eval_if(&dc.add_received_spf, self, self.data.session_id) .await .unwrap_or(true) { ReceivedSpf::new( spf_output, self.data.remote_ip, &self.data.helo_domain, &mail_from.address_lcase, &self.hostname, ) .write_header(&mut headers); } // Run SPAM filter let mut train_spam = None; let mut spam_result = None; if self.server.core.spam.enabled && self .server .eval_if(&dc.spam_filter, self, self.data.session_id) .await .unwrap_or(true) { match self .spam_classify( &parsed_message, &dkim_output, dkim2_output.as_ref(), (&arc_output).into(), dmarc_result.as_ref(), dmarc_policy.as_ref(), ) .await { SpamFilterAction::Allow(score) => { // Add headers headers.extend_from_slice(score.headers.as_bytes()); train_spam = score.train_spam.map(|is_spam| { ( is_spam, thread_name(parsed_message.subject().unwrap_or_default()).to_string(), ) }); let scores = &self.server.core.spam.scores; spam_result = Some((score.score, scores.spam_percentage(score.score))); // Add scores for local recipients for (user_score, recipient) in score.results.into_iter().zip(self.data.rcpt_to.iter_mut()) { recipient.flags = (recipient.flags & !RCPT_SPAM_MASK) | rcpt_spam_flag(scores.spam_percentage(user_score)); } } SpamFilterAction::Discard => { trc::event!( Spam(SpamEvent::Classify), SpanId = self.data.session_id, QueueId = message_id, Result = "discard", Reason = "Message discarded due to excessive spam score.", ); self.data.messages_sent += 1; return (b"250 2.0.0 Message queued for delivery.\r\n"[..]).into(); } SpamFilterAction::Reject => { trc::event!( Spam(SpamEvent::Classify), SpanId = self.data.session_id, QueueId = message_id, Result = "reject", Reason = "Message rejected due to excessive spam score.", ); self.data.messages_sent += 1; return (b"550 5.7.1 Message rejected due to excessive spam score.\r\n"[..]) .into(); } SpamFilterAction::Disabled => {} } } // Run Milter filters let mut modifications = Vec::new(); match self .run_milters(Stage::Data, (&auth_message).into(), message_id.into()) .await { Ok(modifications_) => { if !modifications_.is_empty() { modifications = modifications_; } } Err(response) => { return response.into_bytes(); } }; // Run MTA Hooks match self .run_mta_hooks(Stage::Data, (&auth_message).into(), message_id.into()) .await { Ok(modifications_) => { if !modifications_.is_empty() { modifications.retain(|m| !matches!(m, Modification::ReplaceBody { .. })); modifications.extend(modifications_); } } Err(response) => { return response.into_bytes(); } }; // Apply modifications let mut edited_message = if !modifications.is_empty() { self.data .apply_milter_modifications(modifications, &auth_message) } else { None }; // Sieve filtering if let Some((script, script_id)) = self .server .eval_if::(&dc.script, self, self.data.session_id) .await .and_then(|name| { self.server .get_trusted_sieve_script(&name, self.data.session_id) .map(|s| (s, name)) }) { let mut params = self .build_script_parameters("data") .with_auth_headers(&headers); if let Some((score, percentage)) = spam_result { params = params .with_spam_status(spam_status(Some(percentage))) .set_variable("spam.score", score as f64) .set_variable("spam.is_spam", self.server.core.spam.scores.is_spam(score)); } let params = params .set_variable( "arc.result", arc_output .as_ref() .map(|a| a.result().as_str()) .unwrap_or_default(), ) .set_variable( "dkim.result", dkim_output .iter() .find(|r| matches!(r.result(), DkimResult::Pass)) .or_else(|| dkim_output.first()) .map(|r| r.result().as_str()) .unwrap_or_default(), ) .set_variable( "dkim.domains", dkim_output .iter() .filter_map(|r| { if matches!(r.result(), DkimResult::Pass) { r.signature() .map(|s| Variable::from(s.domain().to_lowercase())) } else { None } }) .collect::>(), ) .set_variable( "dmarc.result", dmarc_result .as_ref() .map(|a| a.as_str()) .unwrap_or_default(), ) .set_variable( "dmarc.policy", dmarc_policy .as_ref() .map(|a| a.as_str()) .unwrap_or_default(), ) .with_message(parsed_message); let modifications = match self.run_script(script_id, script.clone(), params).await { ScriptResult::Accept { modifications } => modifications, ScriptResult::Replace { message, modifications, } => { edited_message = message.into(); modifications } ScriptResult::Reject(message) => { return message.as_bytes().to_vec().into(); } ScriptResult::Discard => { return (b"250 2.0.0 Message queued for delivery.\r\n"[..]).into(); } }; // Apply modifications for modification in modifications { match modification { ScriptModification::AddHeader { name, value } => { headers.extend_from_slice(name.as_bytes()); headers.extend_from_slice(b": "); headers.extend_from_slice(value.as_bytes()); if !value.ends_with('\n') { headers.extend_from_slice(b"\r\n"); } } ScriptModification::SetEnvelope { name, value } => { self.data.apply_envelope_modification(name, value); } } } } // Build message let mail_from = self.data.mail_from.clone().unwrap(); let rcpt_to = std::mem::take(&mut self.data.rcpt_to); let source = if !self.is_authenticated() { let dmarc_pass = dmarc_result.is_some_and(|result| result == DmarcResult::Pass); #[cfg(feature = "test_mode")] { MessageSource::Unauthenticated { dmarc_pass: dmarc_pass || mail_from.address.starts_with("dmarc-"), } } #[cfg(not(feature = "test_mode"))] { MessageSource::Unauthenticated { dmarc_pass } } } else { MessageSource::Authenticated }; let mut message = self .build_message(mail_from, rcpt_to, source, message_id, self.data.session_id) .await; // Add Return-Path if self .server .eval_if(&dc.add_return_path, self, self.data.session_id) .await .unwrap_or(true) { headers.extend_from_slice(b"Return-Path: <"); headers.extend_from_slice(message.message.return_path.as_bytes()); headers.extend_from_slice(b">\r\n"); } // Add any missing headers if !has_date_header && self .server .eval_if(&dc.add_date, self, self.data.session_id) .await .unwrap_or(true) { headers.extend_from_slice(b"Date: "); headers.extend_from_slice(Date::now().to_rfc822().as_bytes()); headers.extend_from_slice(b"\r\n"); } if !has_message_id_header && self .server .eval_if(&dc.add_message_id, self, self.data.session_id) .await .unwrap_or(true) { headers.extend_from_slice(b"Message-ID: "); generate_message_id_header(&mut headers, &self.hostname); headers.extend_from_slice(b"\r\n"); } // Update size let original_message = raw_message.as_slice(); let raw_message = edited_message.as_deref().unwrap_or(raw_message.as_slice()); message.message.size = (raw_message.len() + headers.len()) as u64; // Verify queue quota if let Some(metadata) = self.server.has_quota(&mut message).await { // Queue message let queue_id = message.queue_id; let dkim_signers = self .server .eval_signers(&ac.dkim.sign, self, self.data.session_id) .await; if message .queue( QueueParams::new(raw_message, self.data.session_id, &self.server) .with_train_spam(train_spam) .with_raw_headers(&headers) .with_dkim_signers(dkim_signers) .with_original_raw_message(original_message) .with_original_authenticated_message(auth_message) .with_metadata(metadata), ) .await { self.state = State::Accepted(queue_id); self.data.messages_sent += 1; format!("250 2.0.0 Message queued with id {queue_id:x}.\r\n") .into_bytes() .into() } else { (b"451 4.3.5 Unable to accept message at this time.\r\n"[..]).into() } } else { (b"452 4.3.1 Mail system full, try again later.\r\n"[..]).into() } } pub async fn build_message( &self, mail_from: SessionAddress, mut rcpt_to: Vec, source: MessageSource, queue_id: u64, span_id: u64, ) -> MessageWrapper { // Build message let created = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) .map_or(0, |d| d.as_secs()); let mut message = Message { created, return_path: mail_from .address .to_lowercase_address(false) .into_boxed_str(), recipients: Vec::with_capacity(rcpt_to.len()), flags: mail_from.flags | source.flags(), priority: self.data.priority, size: 0, env_id: mail_from.dsn_info.map(|i| i.into_boxed_str()), blob_hash: Default::default(), metadata: Default::default(), received_from_ip: self.data.remote_ip, received_via_port: self.data.local_port, }; // Add recipients let future_release = self.data.future_release; rcpt_to.sort_unstable(); for rcpt in rcpt_to { message.recipients.push( queue::Recipient::new(rcpt.address) .with_flags( if rcpt.flags & (RCPT_NOTIFY_DELAY | RCPT_NOTIFY_FAILURE | RCPT_NOTIFY_SUCCESS | RCPT_NOTIFY_NEVER) != 0 { rcpt.flags } else { rcpt.flags | RCPT_NOTIFY_DELAY | RCPT_NOTIFY_FAILURE }, ) .with_orcpt(rcpt.dsn_info.map(|v| v.into_boxed_str())), ); let envelope = QueueEnvelope::new(&message, message.recipients.last().unwrap()); // Set next retry time let retry = if self.data.future_release == 0 { queue::Schedule::now() } else { queue::Schedule::later(future_release) }; // Resolve queue let queue = self.server.get_queue_or_default( &self .server .eval_if::( &self.server.core.smtp.queue.queue, &envelope, self.data.session_id, ) .await .unwrap_or_else(|| "default".to_string()), self.data.session_id, ); // Set expiration and notification times let num_intervals = std::cmp::max(queue.notify.len(), 1); let next_notify = queue.notify.first().copied().unwrap_or(86400); let (notify, expires) = if self.data.delivery_by == 0 { ( queue::Schedule::later(future_release + next_notify), match queue.expiry { QueueExpiry::Ttl(time) => QueueExpiry::Ttl(future_release + time), QueueExpiry::Attempts(count) => QueueExpiry::Attempts(count), }, ) } else if (message.flags & MAIL_BY_RETURN) != 0 { ( queue::Schedule::later(future_release + next_notify), QueueExpiry::Ttl(self.data.delivery_by as u64), ) } else { let (notify, expires) = match queue.expiry { QueueExpiry::Ttl(expire_secs) => ( (if self.data.delivery_by.is_positive() { let notify_at = self.data.delivery_by as u64; if expire_secs > notify_at { notify_at } else { next_notify } } else { let notify_at = -self.data.delivery_by as u64; if expire_secs > notify_at { expire_secs - notify_at } else { next_notify } }), QueueExpiry::Ttl(expire_secs), ), QueueExpiry::Attempts(_) => ( next_notify, QueueExpiry::Ttl(self.data.delivery_by.unsigned_abs()), ), }; let mut notify = queue::Schedule::later(future_release + notify); notify.inner = (num_intervals - 1) as u32; // Disable further notification attempts (notify, expires) }; // Update recipient let recipient = message.recipients.last_mut().unwrap(); recipient.retry = retry; recipient.notify = notify; recipient.expires = expires; recipient.queue = queue.virtual_queue; } MessageWrapper { queue_id, queue_name: QueueName::default(), is_multi_queue: false, span_id, message, } } pub async fn can_send_data(&mut self) -> Option<&'static [u8]> { if self.data.mail_from.is_none() { trc::event!( Smtp(SmtpEvent::MailFromMissing), SpanId = self.data.session_id, ); Some(b"503 5.5.1 MAIL is required first.\r\n") } else if self.data.rcpt_to.is_empty() { trc::event!( Smtp(SmtpEvent::RcptToMissing), SpanId = self.data.session_id, ); Some(b"503 5.5.1 RCPT is required first.\r\n") } else if self.data.messages_sent < self .server .eval_if( &self.server.core.smtp.session.data.max_messages, self, self.data.session_id, ) .await .unwrap_or(10) { None } else { trc::event!( Smtp(SmtpEvent::TooManyMessages), SpanId = self.data.session_id, Limit = self.data.messages_sent ); Some(b"452 4.4.5 Maximum number of messages per session exceeded.\r\n") } } fn write_received(&self, headers: &mut Vec, id: u64) { headers.extend_from_slice(b"Received: from "); headers.extend_from_slice(self.data.helo_domain.as_bytes()); headers.extend_from_slice(b" ("); headers.extend_from_slice( self.data .iprev .as_ref() .and_then(|ir| ir.ptr.as_ref()) .and_then(|ptr| ptr.first().map(|s| s.strip_suffix('.').unwrap_or(s))) .unwrap_or("unknown") .as_bytes(), ); headers.extend_from_slice(b" ["); headers.extend_from_slice(self.data.remote_ip.to_string().as_bytes()); headers.extend_from_slice(b"]"); if self.data.asn_geo_data.asn.is_some() || self.data.asn_geo_data.country.is_some() { headers.extend_from_slice(b" ("); if let Some(asn) = &self.data.asn_geo_data.asn { headers.extend_from_slice(b"AS"); headers.extend_from_slice(asn.id.to_string().as_bytes()); if let Some(name) = &asn.name { headers.extend_from_slice(b" "); headers.extend_from_slice(name.as_bytes()); } } if let Some(country) = &self.data.asn_geo_data.country { if self.data.asn_geo_data.asn.is_some() { headers.extend_from_slice(b", "); } headers.extend_from_slice(country.as_bytes()); } headers.extend_from_slice(b")"); } headers.extend_from_slice(b")\r\n\t"); if self.stream.is_tls() { let (version, cipher) = self.stream.tls_version_and_cipher(); headers.extend_from_slice(b"(using "); headers.extend_from_slice(version.as_bytes()); headers.extend_from_slice(b" with cipher "); headers.extend_from_slice(cipher.as_bytes()); headers.extend_from_slice(b")\r\n\t"); } headers.extend_from_slice(b"by "); headers.extend_from_slice(self.hostname.as_bytes()); headers.extend_from_slice(concat!(" (", types::brand!(), " SMTP) with ").as_bytes()); headers.extend_from_slice(match (self.stream.is_tls(), !self.is_authenticated()) { (true, true) => b"ESMTPS", (true, false) => b"ESMTPSA", (false, true) => b"ESMTP", (false, false) => b"ESMTPA", }); headers.extend_from_slice(b" id "); headers.extend_from_slice(format!("{id:X}").as_bytes()); headers.extend_from_slice(b";\r\n\t"); headers.extend_from_slice(Date::now().to_rfc822().as_bytes()); headers.extend_from_slice(b"\r\n"); } } fn parse_dkim2_dsn<'x, 'r>( parsed_message: &mail_parser::Message<'x>, raw: &'r AuthenticatedMessage<'x>, raw_message: &'x [u8], ) -> Option, AuthenticatedMessage<'x>>> { if !parsed_message.content_type().is_some_and(|ct| { ct.ctype().eq_ignore_ascii_case("multipart") && ct .subtype() .is_some_and(|subtype| subtype.eq_ignore_ascii_case("report")) && ct .attribute("report-type") .is_some_and(|report_type| report_type.eq_ignore_ascii_case("delivery-status")) }) { return None; } let mail_parser::PartType::Multipart(children) = &parsed_message.root_part().body else { return None; }; let mut returned = &b""[..]; let mut returned_full = false; for child in children { let part = parsed_message.parts.get(*child as usize)?; if part.is_content_type("message", "rfc822") { returned = raw_message.get(part.offset_body as usize..part.offset_end as usize)?; returned_full = true; } else if part.is_content_type("text", "rfc822-headers") { returned = raw_message.get(part.offset_body as usize..part.offset_end as usize)?; returned_full = false; } } if !returned.is_empty() { Some(Dkim2Dsn::new( raw, AuthenticatedMessage::parse(returned)?, returned_full, )) } else { None } }