fraiseql-functions 2.16.0

Serverless functions runtime for FraiseQL — WASM and Deno backends
Documentation
//! Pure email normalization: RFC 5322 / MIME → [`InboundMessage`].
//!
//! The shared, high-value layer of the email adapter, sitting *above* the
//! poll-IMAP transport. Given the raw bytes of one message it derives everything
//! an `after:ingest:email` function needs — a stable dedup key, a conversation
//! [`thread_key`](InboundMessage::thread_key), the sender and recipients, the
//! text and HTML bodies, the headers, and a [`Classification`] (human /
//! out-of-office / bounce / challenge / auto-generated) that reply-awareness keys
//! on — with **no network or database I/O**.
//!
//! Attachments are the one thing the pure layer cannot finish: writing bytes to
//! `[storage]` is I/O. [`normalize_email`] therefore returns each attachment as a
//! [`PendingAttachment`] (filename + content type + decoded bytes) and leaves
//! [`InboundMessage::attachments`] empty; the caller streams each blob to storage
//! and fills in the [`StorageRef`](crate::StorageRef)s. This keeps the whole
//! module deterministic and unit-testable in the fast lib leg.
//!
//! ## Classification taxonomy
//!
//! Reply-awareness must treat a person's reply differently from the several kinds
//! of automated mail. The layer resolves exactly one [`Classification`] per
//! message, most-specific first:
//!
//! | Class | Signal |
//! |-------|--------|
//! | [`Bounce`](Classification::Bounce) | RFC 3464 `multipart/report; report-type=delivery-status`, a `MAILER-DAEMON` / `postmaster` sender, or an `X-Failed-Recipients` header. |
//! | [`OutOfOffice`](Classification::OutOfOffice) | `Auto-Submitted: auto-replied`, `Precedence: auto_reply`, or a vacation-responder header. |
//! | [`Challenge`](Classification::Challenge) | A challenge-response header, or an automated message whose subject asks the sender to confirm themselves. |
//! | [`AutoGenerated`](Classification::AutoGenerated) | `Auto-Submitted: auto-generated`, `Precedence: bulk`/`list`/`junk`, or list headers (`List-Id` / `List-Unsubscribe`). |
//! | [`Human`](Classification::Human) | none of the above — the only class that advances reply-awareness. |

use std::collections::BTreeMap;

use mail_parser::{Address, HeaderValue, Message, MessageParser, MessagePart, MimeHeaders};
use sha2::{Digest, Sha256};

use super::{Classification, InboundMessage, IngestError, IngestSource};

/// An attachment extracted during normalization but not yet persisted.
///
/// The pure layer decodes the bytes; the caller streams them into `[storage]` and
/// records the resulting [`StorageRef`](crate::StorageRef) on the message.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PendingAttachment {
    /// Filename the sender supplied, or a synthesized fallback.
    pub filename:     String,
    /// MIME type (`type/subtype`, lower-cased).
    pub content_type: String,
    /// Transfer-decoded attachment bytes.
    pub bytes:        Vec<u8>,
}

/// The result of normalizing one raw email: the [`InboundMessage`] plus the
/// attachments the caller still has to stream into storage.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParsedEmail {
    /// The normalized message, with [`attachments`](InboundMessage::attachments)
    /// left empty for the caller to fill in.
    pub message:     InboundMessage,
    /// Attachments decoded from the message, pending a storage write.
    pub attachments: Vec<PendingAttachment>,
}

/// Normalize the raw bytes of one email into a [`ParsedEmail`].
///
/// `source` is normally [`IngestSource::Email`]; it is a parameter so a future
/// email-shaped source (a mail-provider inbound-parse webhook) can reuse the same
/// normalization. `received_at` is when the adapter fetched the message.
///
/// The [`idempotency_key`](InboundMessage::idempotency_key) is the `Message-ID`
/// (angle brackets stripped); a message with no `Message-ID` falls back to a
/// `sha256:` digest of the raw bytes, so a re-fetch of the same message still
/// deduplicates on the spine.
///
/// # Errors
///
/// Returns [`IngestError`] if the bytes cannot be parsed as an RFC 5322 / MIME
/// message. Normalization fails loud rather than fabricating a partial message.
pub fn normalize_email(
    raw: &[u8],
    source: IngestSource,
    received_at: chrono::DateTime<chrono::Utc>,
) -> Result<ParsedEmail, IngestError> {
    let parsed = MessageParser::default()
        .parse(raw)
        .ok_or_else(|| IngestError::new("email is not parseable as an RFC 5322/MIME message"))?;

    let headers = collect_headers(raw, &parsed);
    let message_id = parsed.message_id().map(strip_message_id).filter(|id| !id.is_empty());
    let idempotency_key = dedup_key(message_id.as_deref(), raw);

    let in_reply_to = header_value_ids(parsed.in_reply_to());
    let references = header_value_ids(parsed.references());
    let thread_key = derive_thread_key(message_id.as_deref(), &in_reply_to, &references);

    let from = address_list(parsed.from()).into_iter().next();
    let to = address_list(parsed.to());
    let cc = address_list(parsed.cc());
    let subject = parsed.subject().map(str::to_string);
    let body_text = parsed.body_text(0).map(|body| body.into_owned());
    let body_html = parsed.body_html(0).map(|body| body.into_owned());

    let classification = classify(
        &headers,
        from.as_deref(),
        is_delivery_status_report(&parsed),
        subject.as_deref(),
    );

    let attachments = parsed.attachments().filter_map(pending_attachment).collect();

    let mut message = InboundMessage::new(source, idempotency_key, received_at);
    message.thread_key = thread_key;
    message.from = from;
    message.to = to;
    message.cc = cc;
    message.subject = subject;
    message.body_text = body_text;
    message.body_html = body_html;
    message.headers = headers;
    message.classification = Some(classification);

    Ok(ParsedEmail {
        message,
        attachments,
    })
}

/// Derive the conversation [`thread_key`](InboundMessage::thread_key) for a
/// message from its threading headers.
///
/// The key is the conversation *root*: the first entry of `References` (the
/// oldest ancestor), else the `In-Reply-To` parent, else the message's own
/// `Message-ID` (it starts a new thread). Every reply in a chain quotes the same
/// `References` root, so the whole chain collapses to one key. All ids are
/// expected already stripped of angle brackets.
#[must_use]
pub fn derive_thread_key(
    message_id: Option<&str>,
    in_reply_to: &[String],
    references: &[String],
) -> Option<String> {
    references
        .first()
        .or_else(|| in_reply_to.first())
        .cloned()
        .or_else(|| message_id.map(str::to_string))
}

/// Classify a message from its headers, sender, and structure.
///
/// Resolves exactly one [`Classification`], most-specific first (see the module
/// docs for the full taxonomy). `is_delivery_status` is whether the top-level
/// MIME structure is an RFC 3464 delivery-status report; `subject` is only
/// consulted for the conservative challenge heuristic.
#[must_use]
pub fn classify(
    headers: &BTreeMap<String, String>,
    from: Option<&str>,
    is_delivery_status: bool,
    subject: Option<&str>,
) -> Classification {
    if is_bounce(headers, from, is_delivery_status) {
        return Classification::Bounce;
    }
    let auto_submitted = auto_submitted_keyword(headers);
    if auto_submitted.as_deref() == Some("auto-replied") || is_vacation_responder(headers) {
        return Classification::OutOfOffice;
    }
    if is_challenge(headers, subject, auto_submitted.as_deref()) {
        return Classification::Challenge;
    }
    if is_auto_generated(headers, auto_submitted.as_deref()) {
        return Classification::AutoGenerated;
    }
    Classification::Human
}

/// Look up a header by its lower-cased name.
fn header<'a>(headers: &'a BTreeMap<String, String>, name: &str) -> Option<&'a str> {
    headers.get(name).map(String::as_str)
}

/// The local part of an address (`local@domain` → `local`).
fn local_part(address: &str) -> &str {
    address.trim().rsplit_once('@').map_or(address.trim(), |(local, _)| local)
}

/// A delivery-status notification: an explicit report structure, a
/// `X-Failed-Recipients` header, or a daemon sender.
fn is_bounce(
    headers: &BTreeMap<String, String>,
    from: Option<&str>,
    is_delivery_status: bool,
) -> bool {
    if is_delivery_status || headers.contains_key("x-failed-recipients") {
        return true;
    }
    from.is_some_and(|address| {
        let local = local_part(address).to_ascii_lowercase();
        local == "mailer-daemon" || local == "postmaster"
    })
}

/// The leading keyword of the `Auto-Submitted` header (`auto-replied`,
/// `auto-generated`, `no`), lower-cased and stripped of any parameters.
fn auto_submitted_keyword(headers: &BTreeMap<String, String>) -> Option<String> {
    header(headers, "auto-submitted")
        .map(|value| value.split([';', '(']).next().unwrap_or(value).trim().to_ascii_lowercase())
}

/// A vacation / out-of-office responder, identified by a responder header or
/// `Precedence: auto_reply`.
fn is_vacation_responder(headers: &BTreeMap<String, String>) -> bool {
    [
        "x-autoreply",
        "x-autorespond",
        "x-vacation-message",
        "x-autoreply-domain",
    ]
    .iter()
    .any(|name| headers.contains_key(*name))
        || header(headers, "precedence")
            .is_some_and(|p| p.trim().eq_ignore_ascii_case("auto_reply"))
}

/// A challenge-response verification.
///
/// Deliberately conservative: it fires only on an explicit challenge header, or
/// on an automated message whose subject carries a confirm/verify phrase — so a
/// human's "please confirm the meeting" reply is never misread as a challenge.
fn is_challenge(
    headers: &BTreeMap<String, String>,
    subject: Option<&str>,
    auto_submitted: Option<&str>,
) -> bool {
    if headers.contains_key("x-challenge") || headers.contains_key("x-antiphishing-challenge") {
        return true;
    }
    let automated = auto_submitted.is_some_and(|keyword| !keyword.is_empty() && keyword != "no")
        || headers.contains_key("list-id");
    let phrase = subject.is_some_and(|subject| {
        let subject = subject.to_ascii_lowercase();
        subject.contains("please confirm")
            || subject.contains("confirm your")
            || subject.contains("verify you")
            || subject.contains("challenge-response")
    });
    automated && phrase
}

/// Machine-generated bulk/list mail that is not a human reply.
fn is_auto_generated(headers: &BTreeMap<String, String>, auto_submitted: Option<&str>) -> bool {
    if auto_submitted.is_some_and(|keyword| !keyword.is_empty() && keyword != "no") {
        return true;
    }
    if header(headers, "precedence").is_some_and(|p| {
        let p = p.trim().to_ascii_lowercase();
        p == "bulk" || p == "list" || p == "junk"
    }) {
        return true;
    }
    headers.contains_key("list-id") || headers.contains_key("list-unsubscribe")
}

/// Whether the top-level MIME structure is an RFC 3464 delivery-status report.
fn is_delivery_status_report(message: &Message<'_>) -> bool {
    message.content_type().is_some_and(|content_type| {
        let report = content_type.ctype().eq_ignore_ascii_case("multipart")
            && content_type.subtype().is_some_and(|s| s.eq_ignore_ascii_case("report"))
            && content_type
                .attribute("report-type")
                .is_some_and(|value| value.eq_ignore_ascii_case("delivery-status"));
        report
            || content_type
                .subtype()
                .is_some_and(|s| s.eq_ignore_ascii_case("delivery-status"))
    })
}

/// The dedup key: `Message-ID` (when present) **plus** a content digest.
///
/// The `Message-ID` header is chosen entirely by the sender, so it must never be
/// the whole key (#775): a forged message carrying a victim's `Message-ID` used
/// to claim the dedup slot and the genuine message was then dropped as a
/// duplicate. Folding in `SHA-256(raw)` means only a byte-identical redelivery
/// deduplicates — which is exactly what dedup is for — while two different
/// messages sharing an id are both processed. (Per-mailbox scoping of the spine
/// key is the other half of the fix; see [`IngestSource::Email`].)
fn dedup_key(message_id: Option<&str>, raw: &[u8]) -> String {
    let digest = hex::encode(Sha256::digest(raw));
    message_id.map_or_else(|| format!("sha256:{digest}"), |id| format!("{id}:sha256:{digest}"))
}

/// Strip surrounding angle brackets and whitespace from a message id.
fn strip_message_id(id: &str) -> String {
    let trimmed = id.trim();
    let inner = trimmed
        .strip_prefix('<')
        .map_or(trimmed, |rest| rest.strip_suffix('>').unwrap_or(rest));
    inner.trim().to_string()
}

/// Extract the message ids from an `In-Reply-To` / `References` header value,
/// each stripped of angle brackets.
fn header_value_ids(value: &HeaderValue<'_>) -> Vec<String> {
    match value {
        HeaderValue::Text(text) => split_ids(text),
        HeaderValue::TextList(list) => list.iter().flat_map(|text| split_ids(text)).collect(),
        _ => Vec::new(),
    }
}

/// Split a whitespace-separated run of `<id>` tokens into stripped ids.
fn split_ids(text: &str) -> Vec<String> {
    text.split_whitespace()
        .map(strip_message_id)
        .filter(|id| !id.is_empty())
        .collect()
}

/// Collect the bare email addresses out of an address header, dropping display
/// names so the result is routable by [`parse_recipient`](crate::parse_recipient).
fn address_list(address: Option<&Address<'_>>) -> Vec<String> {
    address.map_or_else(Vec::new, |address| {
        address
            .iter()
            .filter_map(|addr| addr.address())
            .map(|value| value.trim().to_string())
            .filter(|value| !value.is_empty())
            .collect()
    })
}

/// Collect the top-level headers into a name→value map: names lower-cased,
/// values unfolded, last occurrence winning on a repeated header.
fn collect_headers(raw: &[u8], message: &Message<'_>) -> BTreeMap<String, String> {
    let mut map = BTreeMap::new();
    for header in message.headers() {
        let start = header.offset_start as usize;
        let end = header.offset_end as usize;
        if start > end || end > raw.len() {
            continue;
        }
        let value = unfold(&String::from_utf8_lossy(&raw[start..end]));
        if value.is_empty() {
            continue;
        }
        map.insert(header.name().to_ascii_lowercase(), value);
    }
    map
}

/// Collapse header folding (CRLF + continuation whitespace) and runs of
/// whitespace into single spaces, then trim.
fn unfold(value: &str) -> String {
    let mut out = String::with_capacity(value.len());
    let mut in_whitespace = false;
    for ch in value.chars() {
        if ch.is_whitespace() {
            if !in_whitespace {
                out.push(' ');
                in_whitespace = true;
            }
        } else {
            out.push(ch);
            in_whitespace = false;
        }
    }
    out.trim().to_string()
}

/// Decode one attachment part into a [`PendingAttachment`]; empty parts are
/// dropped.
fn pending_attachment(part: &MessagePart<'_>) -> Option<PendingAttachment> {
    let bytes = part.contents().to_vec();
    if bytes.is_empty() {
        return None;
    }
    let content_type = part.content_type().map_or_else(
        || "application/octet-stream".to_string(),
        |content_type| {
            format!(
                "{}/{}",
                content_type.ctype().to_ascii_lowercase(),
                content_type.subtype().unwrap_or("octet-stream").to_ascii_lowercase()
            )
        },
    );
    let filename = part
        .attachment_name()
        .map_or_else(|| "attachment.bin".to_string(), str::to_string);
    Some(PendingAttachment {
        filename,
        content_type,
        bytes,
    })
}

#[cfg(test)]
mod tests;