mailrs_outbound_queue/worker/
delivery.rs1use hickory_resolver::TokioResolver;
5use sqlx::PgPool;
6
7use super::smtp::try_deliver_via_mx;
8use crate::dsn;
9use crate::queue::{self, QueuedMessage};
10use crate::retry::{retry_delay_secs, should_bounce};
11use crate::{DeliveryEvent, DeliveryEventSender};
12
13pub(super) async fn enqueue_dsn(pool: &PgPool, hostname: &str, msg: &QueuedMessage, error: &str) {
15 if msg.sender.is_empty() || msg.sender == "<>" {
16 return; }
18 let dsn_msg = dsn::format_dsn(
19 hostname,
20 &msg.sender,
21 &msg.recipient,
22 error,
23 msg.message_id.as_deref(),
24 );
25 let sender_domain = msg
26 .sender
27 .rsplit_once('@')
28 .map(|(_, d)| d)
29 .unwrap_or("unknown");
30 let now = chrono::Utc::now().timestamp();
31 let _ = queue::enqueue(
32 pool,
33 "<>",
34 &msg.sender,
35 sender_domain,
36 dsn_msg.as_bytes(),
37 None,
38 now,
39 )
40 .await;
41}
42
43#[tracing::instrument(
50 name = "outbound.deliver_domain",
51 skip(resolver, hostname, messages, pool, event_sender),
52 fields(domain, n_messages = messages.len(), max_per_conn),
53)]
54#[allow(clippy::too_many_arguments)]
55pub async fn deliver_domain_static(
56 resolver: &TokioResolver,
57 hostname: &str,
58 domain: &str,
59 messages: Vec<QueuedMessage>,
60 pool: &PgPool,
61 port: u16,
62 max_per_conn: usize,
63 event_sender: Option<&DeliveryEventSender>,
64) {
65 let batch_start = std::time::Instant::now();
70
71 let mut messages = messages;
73 let now_check = chrono::Utc::now().timestamp();
74 {
75 let mut suppressed_ids = Vec::new();
76 for msg in &messages {
77 if queue::is_suppressed(pool, &msg.recipient).await {
78 tracing::info!("skipping suppressed recipient: {}", msg.recipient);
79 let _ = queue::mark_bounced(
80 pool,
81 msg.id,
82 "recipient suppressed (hard bounce history)",
83 now_check,
84 )
85 .await;
86 if let Some(es) = event_sender {
87 es(DeliveryEvent::Bounced {
88 queue_id: msg.id,
89 sender: msg.sender.clone(),
90 });
91 }
92 suppressed_ids.push(msg.id);
93 }
94 }
95 if !suppressed_ids.is_empty() {
96 messages.retain(|msg| !suppressed_ids.contains(&msg.id));
97 }
98 if messages.is_empty() {
99 return;
100 }
101 }
102
103 let mx_records = match mailrs_smtp_client::resolve_mx(resolver, domain).await {
105 Ok(records) => records,
106 Err(e) => {
107 tracing::warn!("MX resolution failed for {domain}: {e}");
108 let now = chrono::Utc::now().timestamp();
109 for msg in &messages {
110 let delay = retry_delay_secs(msg.attempts);
111 if should_bounce(msg.attempts + 1, msg.max_attempts) {
112 let error = format!("MX resolution failed: {e}");
113 let _ = queue::mark_bounced(pool, msg.id, &error, now).await;
114 if queue::is_hard_bounce(&error) {
116 let _ = queue::add_suppression(pool, &msg.recipient, &error, None).await;
117 }
118 enqueue_dsn(pool, hostname, msg, &error).await;
119 if let Some(es) = event_sender {
120 es(DeliveryEvent::Bounced {
121 queue_id: msg.id,
122 sender: msg.sender.clone(),
123 });
124 }
125 } else {
126 let _ = queue::mark_failed(
127 pool,
128 msg.id,
129 &format!("MX resolution failed: {e}"),
130 now + delay as i64,
131 now,
132 )
133 .await;
134 if let Some(es) = event_sender {
135 es(DeliveryEvent::Failed {
136 queue_id: msg.id,
137 domain: domain.to_string(),
138 error: format!("MX resolution failed: {e}"),
139 });
140 }
141 }
142 }
143 return;
144 }
145 };
146
147 let chunks: Vec<&[QueuedMessage]> = messages.chunks(max_per_conn).collect();
149
150 for mx in &mx_records {
152 let mut all_ok = true;
153 for chunk in &chunks {
154 match try_deliver_via_mx(
155 hostname,
156 &mx.exchange,
157 port,
158 domain,
159 chunk,
160 resolver,
161 event_sender,
162 )
163 .await
164 {
165 Ok(()) => {
166 let now = chrono::Utc::now().timestamp();
167 let elapsed = batch_start.elapsed().as_secs_f64();
168 for msg in *chunk {
169 metrics::histogram!(
170 "mailrs_outbound_delivery_seconds",
171 "outcome" => "delivered",
172 )
173 .record(elapsed);
174 let _ = queue::mark_delivered(pool, msg.id, now).await;
175 if let Some(es) = event_sender {
176 es(DeliveryEvent::Success {
177 queue_id: msg.id,
178 domain: domain.to_string(),
179 });
180 }
181 }
182 }
183 Err(e) => {
184 tracing::warn!("delivery to {} via {} failed: {e}", domain, mx.exchange);
185 all_ok = false;
186 break;
187 }
188 }
189 }
190 if all_ok {
191 tracing::info!(
192 "delivered {} messages to {domain} via {}",
193 messages.len(),
194 mx.exchange
195 );
196 return;
197 }
198 }
199
200 let now = chrono::Utc::now().timestamp();
202 for msg in &messages {
203 if let Ok(Some(current)) = queue::get_message(pool, msg.id).await
205 && current.status == crate::queue::QueueStatus::Delivered
206 {
207 continue;
208 }
209 let delay = retry_delay_secs(msg.attempts);
210 if should_bounce(msg.attempts + 1, msg.max_attempts) {
211 let _ = queue::mark_bounced(pool, msg.id, "all MX hosts failed", now).await;
212 if let Some(ref err) = msg.last_error
214 && queue::is_hard_bounce(err)
215 {
216 let _ = queue::add_suppression(pool, &msg.recipient, err, None).await;
217 }
218 enqueue_dsn(pool, hostname, msg, "all MX hosts failed").await;
219 if let Some(es) = event_sender {
220 es(DeliveryEvent::Bounced {
221 queue_id: msg.id,
222 sender: msg.sender.clone(),
223 });
224 }
225 } else {
226 let _ =
227 queue::mark_failed(pool, msg.id, "all MX hosts failed", now + delay as i64, now)
228 .await;
229 if let Some(es) = event_sender {
230 es(DeliveryEvent::Failed {
231 queue_id: msg.id,
232 domain: domain.to_string(),
233 error: "all MX hosts failed".into(),
234 });
235 }
236 }
237 }
238}