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