Skip to main content

headless_lms_server/programs/
email_deliver.rs

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
30/// How often the sender asks the model to purge expired recipient addresses. The model then rolls its
31/// own 1 in 10 chance, so a sweep lands roughly every ten hours.
32const 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});
41/// The same variable the server reads into `ApplicationConfiguration::base_url`: a mailed link has to
42/// point at the environment that minted its token. Not `FRONTEND_BASE_URL`, which no overlay sets, so
43/// reading that sends dev and test links to production.
44static 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                    // `retry_count` is pre-increment from the claimed row; using it here keeps
205                    // backoff aligned with the next failed-attempt number.
206                    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    // Every remaining template derives its values from an account.
255    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        // Handled above. Listed rather than caught by `_` so a new template type is a compile error.
337        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
350/// Turns a delivery's placeholder bag into `{{ KEY }}` substitutions. Nested objects and arrays are
351/// skipped: there is no sensible rendering for them in body text.
352fn 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    // Startup counts as an attempt: pods restart often, and firing on the first tick would turn every
497    // restart into another sweep.
498    let mut last_purge_attempt = tokio::time::Instant::now();
499    loop {
500        interval.tick().await;
501        mail_sender(&pool, &mailer).await?;
502
503        // An elapsed check rather than a second interval: another `tick().await` in this loop would
504        // stall the 10 second send cycle until the hour was up and stop mail going out.
505        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            // Logged, not propagated: an error out of main() restarts the pod, which would stop mail
513            // going out over a retention sweep.
514            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    // Saturating math + MAX_BACKOFF_SECS clamp intentionally handles outlier values safely.
536    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}