use std::collections::BTreeMap;
use mail_parser::{Address, HeaderValue, Message, MessageParser, MessagePart, MimeHeaders};
use sha2::{Digest, Sha256};
use super::{Classification, InboundMessage, IngestError, IngestSource};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PendingAttachment {
pub filename: String,
pub content_type: String,
pub bytes: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParsedEmail {
pub message: InboundMessage,
pub attachments: Vec<PendingAttachment>,
}
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,
})
}
#[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))
}
#[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
}
fn header<'a>(headers: &'a BTreeMap<String, String>, name: &str) -> Option<&'a str> {
headers.get(name).map(String::as_str)
}
fn local_part(address: &str) -> &str {
address.trim().rsplit_once('@').map_or(address.trim(), |(local, _)| local)
}
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"
})
}
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())
}
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"))
}
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
}
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")
}
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"))
})
}
fn dedup_key(message_id: Option<&str>, raw: &[u8]) -> String {
message_id
.map_or_else(|| format!("sha256:{}", hex::encode(Sha256::digest(raw))), str::to_string)
}
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()
}
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(),
}
}
fn split_ids(text: &str) -> Vec<String> {
text.split_whitespace()
.map(strip_message_id)
.filter(|id| !id.is_empty())
.collect()
}
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()
})
}
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
}
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()
}
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;