use fraiseql_functions::{Classification, InboundMessage, parse_recipient};
use tracing::{info, warn};
use super::tracking::{CorrelatedSend, SendCorrelator, SuppressionReason};
const DELIVERY_HEADERS: [&str; 3] = ["delivered-to", "x-original-to", "envelope-to"];
const REFERENCE_HEADERS: [&str; 2] = ["references", "in-reply-to"];
fn looks_like_send_id(tag: &str) -> bool {
tag.len() == 32 && tag.bytes().all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
}
#[must_use]
pub fn extract_send_id(message: &InboundMessage) -> Option<String> {
let header_addrs = DELIVERY_HEADERS
.iter()
.filter_map(|name| message.headers.get(*name))
.flat_map(|value| value.split([',', ' ']))
.map(str::trim)
.filter(|value| !value.is_empty());
message
.to
.iter()
.map(String::as_str)
.chain(header_addrs)
.filter_map(parse_recipient)
.filter_map(|recipient| recipient.tag)
.find(|tag| looks_like_send_id(tag))
}
#[must_use]
pub fn referenced_message_ids(message: &InboundMessage) -> Vec<String> {
REFERENCE_HEADERS
.iter()
.filter_map(|name| message.headers.get(*name))
.flat_map(|value| value.split_whitespace())
.map(str::trim)
.filter(|token| token.starts_with('<') && token.ends_with('>'))
.map(ToString::to_string)
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CorrelationOutcome {
NoMatch,
Ignored,
Bounced,
Challenge {
pending_count: i64,
suppressed: bool,
},
Replied,
Informational,
}
enum Action {
Bounce,
Challenge,
Reply,
Signal(&'static str),
Ignore,
}
const fn decide(classification: Option<Classification>) -> Action {
match classification {
Some(Classification::Bounce) => Action::Bounce,
Some(Classification::Challenge) => Action::Challenge,
Some(Classification::Human) => Action::Reply,
Some(Classification::OutOfOffice) => Action::Signal("out_of_office"),
Some(Classification::AutoGenerated) => Action::Signal("auto_generated"),
_ => Action::Ignore,
}
}
async fn resolve_send(
correlator: &dyn SendCorrelator,
message: &InboundMessage,
) -> fraiseql_error::Result<Option<CorrelatedSend>> {
if let Some(send_id) = extract_send_id(message) {
if let Some(send) = correlator.find_by_send_id(&send_id).await? {
return Ok(Some(send));
}
}
for message_id in referenced_message_ids(message) {
if let Some(send) = correlator.find_by_message_id(&message_id).await? {
return Ok(Some(send));
}
}
Ok(None)
}
pub async fn correlate(
correlator: &dyn SendCorrelator,
address_hash_key: Option<&[u8]>,
challenge_suppress_after: u32,
now: chrono::DateTime<chrono::Utc>,
message: &InboundMessage,
) -> fraiseql_error::Result<CorrelationOutcome> {
let Some(send) = resolve_send(correlator, message).await? else {
return Ok(CorrelationOutcome::NoMatch);
};
let tenant = send.tenant.as_deref();
let recipient_hash =
address_hash_key.map(|key| fraiseql_observers::hash_address(key, &send.recipient));
match decide(message.classification) {
Action::Bounce => {
correlator.mark_bounced(&send.send_id).await?;
if let Some(hash) = recipient_hash.as_deref() {
correlator.suppress(tenant, hash, SuppressionReason::HardBounce, None).await?;
}
info!(send_id = %send.send_id, "delivery correlation: hard bounce → suppressed");
Ok(CorrelationOutcome::Bounced)
},
Action::Challenge => {
let pending_count =
correlator.bump_challenge(&send.send_id, tenant, &send.recipient).await?;
let suppressed = pending_count >= i64::from(challenge_suppress_after);
if suppressed {
if let Some(hash) = recipient_hash.as_deref() {
correlator
.suppress(
tenant,
hash,
SuppressionReason::ChallengeUnanswered,
SuppressionReason::ChallengeUnanswered.default_ttl(now),
)
.await?;
}
}
warn!(
send_id = %send.send_id,
pending_count,
suppressed,
"delivery correlation: challenge pending — surfaced for review, not auto-solved"
);
Ok(CorrelationOutcome::Challenge {
pending_count,
suppressed,
})
},
Action::Reply => {
correlator.mark_replied(&send.send_id, tenant, &send.recipient).await?;
if let Some(hash) = recipient_hash.as_deref() {
correlator
.lift_suppression(tenant, hash, SuppressionReason::ChallengeUnanswered)
.await?;
}
info!(send_id = %send.send_id, "delivery correlation: reply → engaged");
Ok(CorrelationOutcome::Replied)
},
Action::Signal(signal) => {
correlator.record_signal(&send.send_id, signal).await?;
Ok(CorrelationOutcome::Informational)
},
Action::Ignore => Ok(CorrelationOutcome::Ignored),
}
}
#[cfg(test)]
mod tests;