use std::{collections::HashMap, future::Future, pin::Pin, time::Duration};
use fraiseql_error::{FraiseQLError, Result};
use fraiseql_functions::{
EmailTransport, SendContext, SendEmailRequest, SendEmailResponse, SenderIdentity,
};
use lettre::{
Address, AsyncSmtpTransport, AsyncTransport, Message, Tokio1Executor,
address::Envelope,
message::{Mailbox, MultiPart, SinglePart},
transport::smtp::authentication::Credentials,
};
use tracing::warn;
use super::{
config::{MailboxSmtpConfig, SmtpTlsMode},
tracking::{SendTracker, SentRecord},
};
const GREYLISTING_BACKOFF_SECS: u64 = 300;
#[must_use]
pub fn build_email_transport<S: std::hash::BuildHasher>(
mailboxes: &HashMap<String, super::MailboxConfig, S>,
get_env: impl Fn(&str) -> Option<String>,
tracker: Option<std::sync::Arc<dyn SendTracker>>,
address_hash_key: Option<std::sync::Arc<[u8]>>,
) -> Option<std::sync::Arc<dyn EmailTransport>> {
let accounts = mailboxes
.iter()
.filter_map(|(name, mailbox)| mailbox.smtp.as_ref().map(|smtp| (name.as_str(), smtp)));
let mut transport = SmtpMailboxTransport::build(accounts, get_env)?;
if let Some(tracker) = tracker {
transport = transport.with_tracker(tracker, address_hash_key);
}
Some(std::sync::Arc::new(transport) as std::sync::Arc<dyn EmailTransport>)
}
struct SmtpAccount {
transport: AsyncSmtpTransport<Tokio1Executor>,
verp_local_part: String,
verp_domain: String,
}
impl SmtpAccount {
fn verp_from(&self, send_id: &str) -> Result<Address> {
Address::new(format!("{}+{send_id}", self.verp_local_part), &self.verp_domain).map_err(
|error| FraiseQLError::Validation {
message: format!(
"invalid VERP Return-Path {}+{send_id}@{}: {error}",
self.verp_local_part, self.verp_domain
),
path: None,
},
)
}
}
pub struct SmtpMailboxTransport {
accounts: HashMap<String, SmtpAccount>,
counter: Option<std::sync::Arc<dyn super::warming::SendCounter>>,
tracker: Option<std::sync::Arc<dyn SendTracker>>,
address_hash_key: Option<std::sync::Arc<[u8]>>,
}
impl SmtpMailboxTransport {
#[must_use]
pub fn build<'a>(
mailboxes: impl Iterator<Item = (&'a str, &'a MailboxSmtpConfig)>,
get_env: impl Fn(&str) -> Option<String>,
) -> Option<Self> {
let mut accounts = HashMap::new();
for (name, cfg) in mailboxes {
let Some(password) = get_env(&cfg.password_env) else {
warn!(
mailbox = %name,
password_env = %cfg.password_env,
"SMTP send not enabled for mailbox: password env is unset"
);
continue;
};
let verp_domain = cfg.return_path_domain().to_string();
if verp_domain != cfg.sending_domain() {
warn!(
mailbox = %name,
sending_domain = %cfg.sending_domain(),
return_path_domain = %verp_domain,
"VERP Return-Path domain differs from the sending domain — SPF/DMARC \
alignment is broken and deliverability of tracked sends may degrade"
);
}
match build_account_transport(cfg, password) {
Ok(transport) => {
accounts.insert(
cfg.address.clone(),
SmtpAccount {
transport,
verp_local_part: cfg.return_path_local_part().to_string(),
verp_domain,
},
);
},
Err(error) => warn!(
mailbox = %name,
%error,
"SMTP send not enabled for mailbox: relay build failed"
),
}
}
if accounts.is_empty() {
None
} else {
Some(Self {
accounts,
counter: None,
tracker: None,
address_hash_key: None,
})
}
}
#[must_use]
pub fn with_send_counter(
mut self,
counter: std::sync::Arc<dyn super::warming::SendCounter>,
) -> Self {
self.counter = Some(counter);
self
}
#[must_use]
pub fn with_tracker(
mut self,
tracker: std::sync::Arc<dyn SendTracker>,
address_hash_key: Option<std::sync::Arc<[u8]>>,
) -> Self {
self.tracker = Some(tracker);
self.address_hash_key = address_hash_key;
self
}
#[must_use]
pub fn account_count(&self) -> usize {
self.accounts.len()
}
}
impl EmailTransport for SmtpMailboxTransport {
fn send<'a>(
&'a self,
sender: &'a SenderIdentity,
request: &'a SendEmailRequest,
context: SendContext<'a>,
) -> Pin<Box<dyn Future<Output = Result<SendEmailResponse>> + Send + 'a>> {
Box::pin(async move {
let Some(account) = self.accounts.get(&sender.address) else {
return Err(FraiseQLError::Validation {
message: format!(
"no connected SMTP account for sending address {:?}",
sender.address
),
path: None,
});
};
if let Some(tracker) = self.tracker.as_ref() {
if let Some(key) = self.address_hash_key.as_ref() {
let recipient_hash = fraiseql_observers::hash_address(key, &request.to);
if let Some(reason) =
tracker.suppression_reason(context.tenant, &recipient_hash).await?
{
return Err(FraiseQLError::Validation {
message: format!(
"recipient is suppressed ({reason}) — refusing to send"
),
path: None,
});
}
}
if let Some(send_id) = context.send_id {
if let Some(recorded) = tracker.recorded_send(context.tenant, send_id).await? {
return Ok(SendEmailResponse {
message_id: recorded.message_id,
accepted: true,
});
}
}
}
if let Some(counter) = self.counter.as_ref() {
if let Some(state) = counter.state(&sender.address).await? {
if !state.within_cap() {
return Err(FraiseQLError::RateLimited {
message: format!(
"sending address {:?} is at its warming daily cap",
sender.address
),
retry_after_secs: 86_400,
});
}
}
}
let verp_from =
context.send_id.map(|send_id| account.verp_from(send_id)).transpose()?;
let message = build_message(sender, request, verp_from)?;
match account.transport.send(message).await {
Ok(response) => {
let message_id = response.first_line().map(ToString::to_string);
if let (Some(tracker), Some(send_id)) = (self.tracker.as_ref(), context.send_id)
{
let record = SentRecord {
send_id,
tenant: context.tenant,
recipient: &request.to,
sending_address: &sender.address,
message_id: message_id.as_deref(),
};
if let Err(error) = tracker.record_sent(record).await {
warn!(%send_id, %error, "failed to record Sent status after relay");
}
}
if let Some(counter) = self.counter.as_ref() {
if let Err(error) = counter.record_send(&sender.address).await {
warn!(address = %sender.address, %error, "failed to record send for warming");
}
}
Ok(SendEmailResponse {
message_id,
accepted: true,
})
},
Err(error) if error.is_permanent() => Err(FraiseQLError::Validation {
message: format!("SMTP permanent error: {error}"),
path: None,
}),
Err(error) => Err(FraiseQLError::ServiceUnavailable {
message: format!("SMTP transient error: {error}"),
retry_after: Some(GREYLISTING_BACKOFF_SECS),
}),
}
})
}
}
fn build_account_transport(
cfg: &MailboxSmtpConfig,
password: String,
) -> Result<AsyncSmtpTransport<Tokio1Executor>> {
let mut builder = match cfg.tls {
SmtpTlsMode::StartTls => AsyncSmtpTransport::<Tokio1Executor>::starttls_relay(&cfg.host)
.map_err(|error| FraiseQLError::Configuration {
message: format!("cannot build STARTTLS relay to {}: {error}", cfg.host),
})?,
SmtpTlsMode::Tls => {
AsyncSmtpTransport::<Tokio1Executor>::relay(&cfg.host).map_err(|error| {
FraiseQLError::Configuration {
message: format!("cannot build TLS relay to {}: {error}", cfg.host),
}
})?
},
SmtpTlsMode::None => AsyncSmtpTransport::<Tokio1Executor>::builder_dangerous(&cfg.host),
};
builder = builder
.port(cfg.port)
.timeout(Some(Duration::from_secs(cfg.timeout_secs)))
.credentials(Credentials::new(cfg.username.clone(), password));
Ok(builder.build())
}
fn build_message(
sender: &SenderIdentity,
request: &SendEmailRequest,
verp_from: Option<Address>,
) -> Result<Message> {
let to_address = request.to.parse::<Address>().map_err(|error| FraiseQLError::Validation {
message: format!("invalid email address {:?}: {error}", request.to),
path: None,
})?;
let from = mailbox(&sender.address, sender.display_name.as_deref())?;
let to = Mailbox::new(None, to_address.clone());
let mut builder = Message::builder().from(from).to(to).subject(request.subject.clone());
if let Some(reply_to) = request.reply_to.as_deref() {
builder = builder.reply_to(mailbox(reply_to, None)?);
}
if let Some(verp) = verp_from {
let envelope = Envelope::new(Some(verp), vec![to_address]).map_err(|error| {
FraiseQLError::Validation {
message: format!("failed to build VERP envelope: {error}"),
path: None,
}
})?;
builder = builder.envelope(envelope);
}
let built = match (request.text.as_deref(), request.html.as_deref()) {
(Some(text), Some(html)) => {
builder.multipart(MultiPart::alternative_plain_html(text.to_owned(), html.to_owned()))
},
(Some(text), None) => builder.singlepart(SinglePart::plain(text.to_owned())),
(None, Some(html)) => builder.singlepart(SinglePart::html(html.to_owned())),
(None, None) => builder.singlepart(SinglePart::plain(String::new())),
};
built.map_err(|error| FraiseQLError::Validation {
message: format!("failed to build email message: {error}"),
path: None,
})
}
fn mailbox(address: &str, display_name: Option<&str>) -> Result<Mailbox> {
let parsed = address.parse::<Address>().map_err(|error| FraiseQLError::Validation {
message: format!("invalid email address {address:?}: {error}"),
path: None,
})?;
Ok(Mailbox::new(display_name.map(ToOwned::to_owned), parsed))
}
#[cfg(test)]
mod tests;