1use std::{error::Error as StdError, time::Duration};
2
3use crate::config::program_config::ProgramConfig;
4use crate::prelude::*;
5use crate::setup_tracing;
6use anyhow::{Context, Result};
7use chrono::{DateTime, Utc};
8use futures::{FutureExt, StreamExt};
9use headless_lms_models::email_deliveries::{
10 Email, EmailDeliveryErrorInsert, FETCH_LIMIT, RETRY_WINDOW_SECS, fetch_emails,
11 increment_retry_and_mark_non_retryable, increment_retry_and_schedule,
12 insert_email_delivery_error, mark_as_sent, maybe_purge_expired_recipient_addresses,
13};
14use headless_lms_models::email_templates::EmailTemplateType;
15use headless_lms_models::user_email_codes::UserEmailCodePurpose;
16use headless_lms_models::user_passwords::get_unused_reset_password_token_with_user_id;
17use headless_lms_utils::backoff;
18use headless_lms_utils::email_processor::{self, BlockAttributes, EmailGutenbergBlock};
19use lettre::transport::smtp::Error as SmtpError;
20use lettre::transport::smtp::authentication::Credentials;
21use lettre::{
22 Message, SmtpTransport, Transport,
23 message::{MultiPart, SinglePart, header},
24};
25use once_cell::sync::Lazy;
26use secrecy::ExposeSecret;
27use sqlx::{Connection, PgConnection, PgPool};
28use std::collections::HashMap;
29use uuid::Uuid;
30const BATCH_SIZE: usize = FETCH_LIMIT as usize;
31
32const RECIPIENT_ADDRESS_PURGE_INTERVAL: Duration = Duration::from_secs(60 * 60);
35
36const BASE_BACKOFF_SECS: i64 = 60;
37const MAX_BACKOFF_SECS: i64 = 24 * 60 * 60;
38const JITTER_SECS: i64 = 30;
39
40static SMTP_FROM: Lazy<String> = Lazy::new(|| {
41 ProgramConfig::required("SMTP_FROM").expect("No moocfi email found in the env variables.")
42});
43static BASE_URL: Lazy<String> = Lazy::new(|| {
47 ProgramConfig::required("BASE_URL").expect("No BASE_URL found in the env variables.")
48});
49static SMTP_HOST: Lazy<String> = Lazy::new(|| {
50 ProgramConfig::required("SMTP_HOST").expect("No email relay found in the env variables.")
51});
52static DB_URL: Lazy<String> = Lazy::new(|| {
53 ProgramConfig::required("DATABASE_URL").expect("No db url found in the env variables.")
54});
55static SMTP_MESSAGE_ID_DOMAIN: Lazy<String> = Lazy::new(|| {
56 ProgramConfig::optional("SMTP_MESSAGE_ID_DOMAIN")
57 .and_then(|value| {
58 let trimmed = value.trim();
59 if trimmed.is_empty() {
60 None
61 } else {
62 Some(trimmed.to_string())
63 }
64 })
65 .or_else(|| infer_email_domain(SMTP_FROM.as_str()))
66 .unwrap_or_else(|| "courses.mooc.fi".to_string())
67});
68static SMTP_USER: Lazy<String> = Lazy::new(|| {
69 ProgramConfig::required("SMTP_USER").expect("No smtp user found in env variables.")
70});
71static SMTP_PASS: Lazy<String> = Lazy::new(|| {
72 ProgramConfig::required("SMTP_PASS").expect("No smtp password found in env variables.")
73});
74
75pub async fn mail_sender(pool: &PgPool, mailer: &SmtpTransport) -> Result<()> {
76 let mut conn = pool.acquire().await?;
77
78 let emails = fetch_emails(&mut conn).await?;
79
80 let mut futures = tokio_stream::iter(emails)
81 .map(|email| {
82 let email_id = email.id;
83 send_message(email, mailer, pool.clone()).inspect(move |r| {
84 if let Err(err) = r {
85 tracing::error!("Failed to send email {}: {}", email_id, err)
86 }
87 })
88 })
89 .buffer_unordered(BATCH_SIZE);
90
91 while futures.next().await.is_some() {}
92
93 Ok(())
94}
95
96pub async fn send_message(email: Email, mailer: &SmtpTransport, pool: PgPool) -> Result<()> {
97 let mut conn = pool.acquire().await?;
98 tracing::info!("Email send messages...");
99
100 let now = Utc::now();
101 let attempt = email.retry_count + 1;
102 if retry_window_expired(email.first_failed_at, now) {
103 tracing::warn!(
104 "Retry window expired for email {} (first_failed_at={:?})",
105 email.id,
106 email.first_failed_at
107 );
108 record_non_retryable_failure(
109 &mut conn,
110 email.id,
111 attempt,
112 "retry_window_expired",
113 format!(
114 "Retry window expired before send attempt (email_id={}, user_id={:?}, template={:?}, first_failed_at={:?})",
115 email.id, email.user_id, email.template_type, email.first_failed_at
116 ),
117 )
118 .await?;
119 return Ok(());
120 }
121
122 let mut email_block: Vec<EmailGutenbergBlock> =
123 match email.body.as_ref().context("No body").and_then(|value| {
124 serde_json::from_value(value.clone()).context("Failed to parse email body JSON")
125 }) {
126 Ok(blocks) => blocks,
127 Err(err) => {
128 record_message_build_failure(&mut conn, &email, attempt, &err).await?;
129 return Ok(());
130 }
131 };
132
133 if let Some(template_type) = email.template_type {
134 let template_result = apply_email_template_replacements(
135 &mut conn,
136 template_type,
137 email.id,
138 email.user_id,
139 email.placeholders.as_ref(),
140 email_block,
141 attempt,
142 )
143 .await?;
144 match template_result {
145 TemplateApplyResult::Ready(blocks) => email_block = blocks,
146 TemplateApplyResult::Abandoned => return Ok(()),
147 }
148 }
149
150 let msg_as_plaintext = email_processor::process_content_to_plaintext(&email_block);
151 let msg_as_html = email_processor::process_content_to_html(&email_block);
152
153 let msg = match build_email_message(&email, attempt, msg_as_plaintext, msg_as_html) {
154 Ok(msg) => msg,
155 Err(err) => {
156 record_message_build_failure(&mut conn, &email, attempt, &err).await?;
157 return Ok(());
158 }
159 };
160
161 match mailer.send(&msg) {
162 Ok(_) => {
163 tracing::info!("Email sent successfully {}", email.id);
164 mark_as_sent(&mut conn, email.id)
165 .await
166 .context("Couldn't mark as sent")?;
167 }
168 Err(err) => {
169 let is_transient = is_transient_smtp_error(&err);
170 let (error_code, smtp_response, smtp_response_code) = extract_smtp_error_details(&err);
171
172 tracing::error!(
173 "SMTP send failed for {} (attempt {}, transient={}): {:?}",
174 email.id,
175 attempt,
176 is_transient,
177 err
178 );
179
180 let mut tx = (*conn)
181 .begin()
182 .await
183 .context("Couldn't start email failure transaction")?;
184
185 insert_email_delivery_error(
186 &mut tx,
187 EmailDeliveryErrorInsert {
188 email_delivery_id: email.id,
189 attempt,
190 error_message: err.to_string(),
191 error_code,
192 smtp_response,
193 smtp_response_code,
194 is_transient,
195 },
196 )
197 .await
198 .context("Couldn't insert email delivery error history")?;
199
200 if is_transient {
201 if retry_window_expired(Some(email.first_failed_at.unwrap_or(now)), now) {
202 increment_retry_and_mark_non_retryable(&mut tx, email.id)
203 .await
204 .context("Couldn't close expired retryable email")?;
205 } else {
206 let next_retry_at = compute_next_retry_at(now, email.retry_count);
209 increment_retry_and_schedule(&mut tx, email.id, Some(next_retry_at))
210 .await
211 .context("Couldn't schedule retry")?;
212 }
213 } else {
214 increment_retry_and_mark_non_retryable(&mut tx, email.id)
215 .await
216 .context("Couldn't close non-retryable email")?;
217 }
218
219 tx.commit()
220 .await
221 .context("Couldn't commit email failure transaction")?;
222 }
223 };
224
225 Ok(())
226}
227
228enum TemplateApplyResult {
229 Ready(Vec<EmailGutenbergBlock>),
230 Abandoned,
231}
232
233async fn apply_email_template_replacements(
234 conn: &mut PgConnection,
235 template_type: EmailTemplateType,
236 email_id: Uuid,
237 user_id: Option<Uuid>,
238 placeholders: Option<&serde_json::Value>,
239 blocks: Vec<EmailGutenbergBlock>,
240 attempt: i32,
241) -> anyhow::Result<TemplateApplyResult> {
242 let mut replacements = HashMap::new();
243
244 if template_type == EmailTemplateType::Generic {
245 return Ok(TemplateApplyResult::Ready(blocks));
246 }
247
248 if template_type.uses_placeholder_bag() {
249 let replacements = placeholder_bag_replacements(placeholders);
250 return Ok(TemplateApplyResult::Ready(insert_placeholders(
251 blocks,
252 &replacements,
253 )));
254 }
255
256 let Some(user_id) = user_id else {
258 let msg = format!(
259 "Template {template_type:?} requires a user but the delivery is addressed to a raw address"
260 );
261 record_non_retryable_failure(conn, email_id, attempt, "template", msg).await?;
262 return Ok(TemplateApplyResult::Abandoned);
263 };
264
265 match template_type {
266 EmailTemplateType::ResetPasswordEmail => {
267 if let Some(token_str) =
268 get_unused_reset_password_token_with_user_id(conn, user_id).await?
269 {
270 let reset_url = format!(
271 "{}/reset-user-password/{}",
272 BASE_URL.trim_end_matches('/'),
273 token_str.token
274 );
275
276 replacements.insert("RESET_LINK".to_string(), reset_url);
277 } else {
278 let msg = anyhow::anyhow!("No reset token found for user {}", user_id);
279 record_non_retryable_failure(conn, email_id, attempt, "template", msg.to_string())
280 .await?;
281 return Ok(TemplateApplyResult::Abandoned);
282 }
283 }
284 EmailTemplateType::DeleteUserEmail => {
285 if let Some(code) =
286 headless_lms_models::user_email_codes::get_unused_user_email_code_with_user_id(
287 conn,
288 user_id,
289 UserEmailCodePurpose::AccountDeletion,
290 )
291 .await?
292 {
293 replacements.insert("CODE".to_string(), code.code.expose_secret().to_string());
294 } else {
295 let msg = anyhow::anyhow!("No deletion code found for user {}", user_id);
296 record_non_retryable_failure(conn, email_id, attempt, "template", msg.to_string())
297 .await?;
298 return Ok(TemplateApplyResult::Abandoned);
299 }
300 }
301 EmailTemplateType::ConfirmEmailCode => {
302 if let Some(code) =
303 headless_lms_models::user_email_codes::get_unused_user_email_code_with_user_id(
304 conn,
305 user_id,
306 UserEmailCodePurpose::AdminLogin,
307 )
308 .await?
309 {
310 replacements.insert("CODE".to_string(), code.code.expose_secret().to_string());
311 } else {
312 let msg = anyhow::anyhow!("No verification code found for user {}", user_id);
313 record_non_retryable_failure(conn, email_id, attempt, "template", msg.to_string())
314 .await?;
315 return Ok(TemplateApplyResult::Abandoned);
316 }
317 }
318 EmailTemplateType::VerifyEmailAddress => {
319 if let Some(code) =
320 headless_lms_models::user_email_codes::get_unused_user_email_code_with_user_id(
321 conn,
322 user_id,
323 UserEmailCodePurpose::EmailOwnershipVerification,
324 )
325 .await?
326 {
327 replacements.insert("CODE".to_string(), code.code.expose_secret().to_string());
328 } else {
329 let msg = anyhow::anyhow!(
330 "No email ownership verification code found for user {}",
331 user_id
332 );
333 record_non_retryable_failure(conn, email_id, attempt, "template", msg.to_string())
334 .await?;
335 return Ok(TemplateApplyResult::Abandoned);
336 }
337 }
338 EmailTemplateType::Generic
340 | EmailTemplateType::CreditRegistrationAccountLinking
341 | EmailTemplateType::CreditRegistrationActionNeeded
342 | EmailTemplateType::CreditRegistrationRegistered => {}
343 }
344
345 Ok(TemplateApplyResult::Ready(insert_placeholders(
346 blocks,
347 &replacements,
348 )))
349}
350
351fn placeholder_bag_replacements(
354 placeholders: Option<&serde_json::Value>,
355) -> HashMap<String, String> {
356 let Some(serde_json::Value::Object(bag)) = placeholders else {
357 return HashMap::new();
358 };
359 bag.iter()
360 .filter_map(|(key, value)| {
361 let rendered = match value {
362 serde_json::Value::String(s) => s.clone(),
363 serde_json::Value::Number(n) => n.to_string(),
364 serde_json::Value::Bool(b) => b.to_string(),
365 _ => return None,
366 };
367 Some((key.clone(), rendered))
368 })
369 .collect()
370}
371
372fn insert_placeholders(
373 blocks: Vec<EmailGutenbergBlock>,
374 replacements: &HashMap<String, String>,
375) -> Vec<EmailGutenbergBlock> {
376 blocks
377 .into_iter()
378 .map(|mut block| {
379 if let BlockAttributes::Paragraph {
380 content,
381 drop_cap,
382 rest,
383 } = block.attributes
384 {
385 let replaced_content = replacements.iter().fold(content, |acc, (key, value)| {
386 acc.replace(&format!("{{{{{}}}}}", key), value)
387 });
388
389 block.attributes = BlockAttributes::Paragraph {
390 content: replaced_content,
391 drop_cap,
392 rest,
393 };
394 }
395 block
396 })
397 .collect()
398}
399
400fn build_email_message(
401 email: &Email,
402 attempt: i32,
403 msg_as_plaintext: String,
404 msg_as_html: String,
405) -> Result<Message> {
406 Message::builder()
407 .from(SMTP_FROM.parse()?)
408 .to(email
409 .to
410 .parse()
411 .with_context(|| format!("Invalid recipient address for email_id {}", email.id))?)
412 .subject(email.subject.clone().context("No subject")?)
413 .message_id(Some(format!(
414 "<{}-{}@{}>",
415 email.id,
416 attempt,
417 SMTP_MESSAGE_ID_DOMAIN.as_str()
418 )))
419 .multipart(
420 MultiPart::alternative()
421 .singlepart(
422 SinglePart::builder()
423 .header(header::ContentType::TEXT_PLAIN)
424 .body(msg_as_plaintext),
425 )
426 .singlepart(
427 SinglePart::builder()
428 .header(header::ContentType::TEXT_HTML)
429 .body(msg_as_html),
430 ),
431 )
432 .context("Failed to build email message")
433}
434
435fn infer_email_domain(value: &str) -> Option<String> {
436 let candidate = value
437 .rsplit('<')
438 .next()
439 .unwrap_or(value)
440 .trim()
441 .trim_end_matches('>')
442 .trim();
443 let (_, domain) = candidate.rsplit_once('@')?;
444 let domain = domain.trim();
445 if domain.is_empty() || domain.contains(' ') {
446 None
447 } else {
448 Some(domain.to_string())
449 }
450}
451
452async fn record_message_build_failure(
453 conn: &mut PgConnection,
454 email: &Email,
455 attempt: i32,
456 err: &anyhow::Error,
457) -> Result<()> {
458 tracing::error!(
459 "Message construction failed for email {} (attempt {}): {:#}",
460 email.id,
461 attempt,
462 err
463 );
464 record_non_retryable_failure(
465 conn,
466 email.id,
467 attempt,
468 "message_build",
469 format!("Message construction failed: {err:#}"),
470 )
471 .await
472}
473
474pub async fn main() -> anyhow::Result<()> {
475 dotenvy::dotenv().ok();
476 ProgramConfig::ensure_default_rust_log_for_workers();
477 setup_tracing()?;
478 tracing::info!("Email sender starting up...");
479
480 if ProgramConfig::optional("SMTP_USER").is_none()
481 || ProgramConfig::optional("SMTP_PASS").is_none()
482 {
483 tracing::warn!("SMTP user or password is missing or incorrect");
484 }
485
486 let pool = PgPool::connect(&DB_URL.to_string()).await?;
487 let creds = Credentials::new(SMTP_USER.to_string(), SMTP_PASS.to_string());
488
489 let mailer = match SmtpTransport::relay(&SMTP_HOST) {
490 Ok(builder) => builder.credentials(creds).build(),
491 Err(e) => {
492 tracing::error!("Could not configure SMTP transport: {}", e);
493 return Err(e.into());
494 }
495 };
496
497 let mut interval = tokio::time::interval(Duration::from_secs(10));
498 let mut last_purge_attempt = tokio::time::Instant::now();
501 loop {
502 interval.tick().await;
503 mail_sender(&pool, &mailer).await?;
504
505 if last_purge_attempt.elapsed() >= RECIPIENT_ADDRESS_PURGE_INTERVAL {
508 last_purge_attempt = tokio::time::Instant::now();
509 let purged = async {
510 let mut conn = pool.acquire().await?;
511 maybe_purge_expired_recipient_addresses(&mut conn).await
512 }
513 .await;
514 match purged {
517 Ok(purged) if purged > 0 => {
518 tracing::info!(
519 "Purged retained recipient addresses from {purged} email deliveries"
520 )
521 }
522 Ok(_) => {}
523 Err(err) => tracing::error!("Failed to purge retained recipient addresses: {err}"),
524 }
525 }
526 }
527}
528
529fn retry_window_expired(first_failed_at: Option<DateTime<Utc>>, now: DateTime<Utc>) -> bool {
530 backoff::window_expired(first_failed_at, now, RETRY_WINDOW_SECS)
531}
532
533fn compute_next_retry_at(now: DateTime<Utc>, retry_count: i32) -> DateTime<Utc> {
534 let delay = backoff::exponential_backoff_secs(BASE_BACKOFF_SECS, MAX_BACKOFF_SECS, retry_count);
535 backoff::next_attempt_at(now, delay, JITTER_SECS)
536}
537
538fn is_transient_smtp_error(err: &SmtpError) -> bool {
539 if err.is_transient() {
540 return true;
541 }
542 if err.is_timeout() || err.is_transport_shutdown() {
543 return true;
544 }
545 has_io_error(err)
546}
547
548fn has_io_error(err: &SmtpError) -> bool {
549 let mut source = err.source();
550 while let Some(inner) = source {
551 if inner.is::<std::io::Error>() {
552 return true;
553 }
554 source = inner.source();
555 }
556 false
557}
558
559fn extract_smtp_error_details(err: &SmtpError) -> (Option<String>, Option<String>, Option<i32>) {
560 let smtp_response_code = err.status().map(|code| i32::from(u16::from(code)));
561
562 let error_code = if err.is_transient() {
563 Some("transient".to_string())
564 } else if err.is_permanent() {
565 Some("permanent".to_string())
566 } else if err.is_timeout() {
567 Some("timeout".to_string())
568 } else if has_io_error(err) {
569 Some("network_io".to_string())
570 } else if err.is_transport_shutdown() {
571 Some("transport_shutdown".to_string())
572 } else if err.is_response() {
573 Some("response".to_string())
574 } else if err.is_client() {
575 Some("client".to_string())
576 } else {
577 None
578 };
579
580 let smtp_response = err.source().map(|source| source.to_string());
581
582 (error_code, smtp_response, smtp_response_code)
583}
584
585async fn record_non_retryable_failure(
586 conn: &mut PgConnection,
587 email_id: Uuid,
588 attempt: i32,
589 error_code: &'static str,
590 message: String,
591) -> Result<()> {
592 let mut tx = (*conn)
593 .begin()
594 .await
595 .context("Couldn't start template failure transaction")?;
596
597 insert_email_delivery_error(
598 &mut tx,
599 EmailDeliveryErrorInsert {
600 email_delivery_id: email_id,
601 attempt,
602 error_message: message,
603 error_code: Some(error_code.to_string()),
604 smtp_response: None,
605 smtp_response_code: None,
606 is_transient: false,
607 },
608 )
609 .await
610 .context("Couldn't insert email delivery error history")?;
611
612 increment_retry_and_mark_non_retryable(&mut tx, email_id)
613 .await
614 .context("Couldn't mark email as non-retryable for template error")?;
615
616 tx.commit()
617 .await
618 .context("Couldn't commit template failure transaction")?;
619
620 Ok(())
621}