use std::time::Duration;
use fraiseql_functions::{
EmailTransport, IngestSource, SendContext, SendEmailRequest, SenderIdentity, normalize_email,
parse_recipient,
};
use super::imap::MailboxFetcher;
const PROBE_TAG_PREFIX: &str = "probe-";
const PROBE_FETCH_BATCH: u32 = 50;
#[must_use]
pub fn probe_recipient(local_part: &str, domain: &str, nonce: &str) -> String {
format!("{local_part}+{PROBE_TAG_PREFIX}{nonce}@{domain}")
}
#[must_use]
pub fn message_carries_probe(message: &fraiseql_functions::InboundMessage, nonce: &str) -> bool {
let want = format!("{PROBE_TAG_PREFIX}{nonce}");
message
.to
.iter()
.map(String::as_str)
.filter_map(parse_recipient)
.filter_map(|recipient| recipient.tag)
.any(|tag| tag == want)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProbeOutcome {
Confirmed,
NotObserved,
}
pub async fn run_return_path_probe(
transport: &dyn EmailTransport,
fetcher: &dyn MailboxFetcher,
sender: &SenderIdentity,
probe_to: &str,
nonce: &str,
timeout: Duration,
poll_interval: Duration,
) -> fraiseql_error::Result<ProbeOutcome> {
let request = SendEmailRequest {
to: probe_to.to_string(),
subject: "fraiseql Return-Path probe".to_string(),
text: Some(
"Automated probe verifying VERP plus-addressing. Safe to ignore.".to_string(),
),
html: None,
reply_to: None,
};
transport.send(sender, &request, SendContext::default()).await?;
let deadline = tokio::time::Instant::now() + timeout;
loop {
let batch = fetcher
.fetch(None, PROBE_FETCH_BATCH)
.await
.map_err(|error| fraiseql_error::FraiseQLError::database(error.to_string()))?;
for message in &batch.messages {
if let Ok(parsed) =
normalize_email(&message.raw, IngestSource::Email, chrono::Utc::now())
{
if message_carries_probe(&parsed.message, nonce) {
return Ok(ProbeOutcome::Confirmed);
}
}
}
if tokio::time::Instant::now() >= deadline {
return Ok(ProbeOutcome::NotObserved);
}
tokio::time::sleep(poll_interval).await;
}
}
#[cfg(test)]
mod tests;