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