Skip to main content

mailrs_outbound_queue/worker/
delivery.rs

1//! Per-domain delivery: MX resolution, suppression filtering, retry/bounce
2//! bookkeeping, DSN enqueueing.
3
4use 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
13/// generate DSN bounce and enqueue it back to the original sender
14pub(super) async fn enqueue_dsn(pool: &PgPool, hostname: &str, msg: &QueuedMessage, error: &str) {
15    if msg.sender.is_empty() || msg.sender == "<>" {
16        return; // don't bounce bounces
17    }
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/// deliver messages to a single domain (used by concurrent workers)
44///
45/// `port` is the TCP port used to connect to every MX in the
46/// resolved set. Production wires 25; integration tests pass an
47/// ephemeral mock-server port to drive the full claim → MX-resolve
48/// → SMTP → mark_* lifecycle without a real MTA on the box.
49#[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    // Wall-clock start for the per-domain delivery batch. Used to
66    // emit `mailrs_outbound_delivery_seconds` histogram on each
67    // terminal mark (delivered / failed / bounced) — gives ops a
68    // distribution of end-to-end delivery latency per outcome.
69    let batch_start = std::time::Instant::now();
70
71    // filter out suppressed recipients before delivery
72    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    // resolve MX records
104    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                    // record hard bounce for suppression
115                    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    // split messages into chunks for connection reuse limits
148    let chunks: Vec<&[QueuedMessage]> = messages.chunks(max_per_conn).collect();
149
150    // try each MX in priority order
151    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    // all MX hosts failed — mark remaining undelivered messages
201    let now = chrono::Utc::now().timestamp();
202    for msg in &messages {
203        // skip already delivered messages
204        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            // add to suppression if last error was a hard bounce
213            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}