use std::collections::HashSet;
use std::sync::Arc;
use std::time::{Duration, Instant};
use hmac::{Hmac, Mac};
use sha2::Sha256;
use uuid::Uuid;
use crate::infrastructure::persistence::mailing_send_repository::{
ClaimedMailing, CompiledDomain, DomainField, DomainOp, DomainTerm,
MailingSendRepository, ResolvedRecipient,
};
use crate::infrastructure::persistence::trace_repository::TraceRepository;
use crate::application::service::sms_delivery_pump_service::SmsDeliveryPumpService;
use backbone_mail::application::service::mail_queue_write_service::{
MailQueueWriteService, MailQueueError,
};
use backbone_mail::application::service::message_write_service::{
MessagePostCommand, MessageWriteService,
};
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
pub enum DomainInvalid {
#[error("domain must be a JSON array of terms")]
NotAnArray,
#[error("term {index}: expected [field, op, value] or {{field, op, value}}")]
MalformedTerm { index: usize },
#[error("term {index}: field must be a string")]
FieldNotString { index: usize },
#[error("field '{0}' is not in the domain whitelist (email, name, first_name, last_name, company_name, country_code, mailing_audience_id)")]
UnknownField(String),
#[error("term {index}: op must be one of =, !=, in, not in, like, not like")]
UnknownOp { index: usize },
#[error("term {index}: 'in'/'not in' need an array of strings")]
ValueNotArray { index: usize },
#[error("term {index}: '='/'!='/'like'/'not like' need exactly one value")]
ValueNotScalar { index: usize },
#[error("term {index}: audience values must be uuids ({raw})")]
BadAudienceUuid { index: usize, raw: String },
}
#[derive(Debug, thiserror::Error)]
pub enum MailingWriteError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("domain invalid: {0}")]
Domain(#[from] DomainInvalid),
#[error("not found: {0}")]
NotFound(String),
#[error("invalid: {0}")]
Invalid(String),
#[error("conflict: {0}")]
Conflict(String),
#[error("mail seam: {0}")]
MailSeam(String),
}
impl MailingWriteError {
pub fn code(&self) -> &'static str {
match self {
Self::Db(_) => "mailing_db_error",
Self::Domain(_) => "domain_invalid",
Self::NotFound(_) => "not_found",
Self::Invalid(_) => "invalid_input",
Self::Conflict(_) => "state_conflict",
Self::MailSeam(_) => "mail_seam_error",
}
}
pub fn http_status(&self) -> u16 {
match self {
Self::Db(_) => 500,
Self::Domain(_) => 422,
Self::NotFound(_) => 404,
Self::Invalid(_) => 422,
Self::Conflict(_) => 409,
Self::MailSeam(_) => 502,
}
}
}
impl From<MailQueueError> for MailingWriteError {
fn from(e: MailQueueError) -> Self {
Self::MailSeam(e.to_string())
}
}
#[derive(Debug, Clone)]
pub struct MailingSendConfig {
pub claim_batch: i64,
pub recipients_per_pass: i64,
pub mail_enqueue_batch: usize,
pub inter_batch_delay_ms: u64,
pub max_enqueues_per_second: Option<u64>,
pub resolver_hard_cap: i64,
}
impl Default for MailingSendConfig {
fn default() -> Self {
Self {
claim_batch: 16,
recipients_per_pass: 500,
mail_enqueue_batch: 100,
inter_batch_delay_ms: 0,
max_enqueues_per_second: None,
resolver_hard_cap: 100_000,
}
}
}
#[derive(Debug, Clone)]
pub struct AutoBlacklistConfig {
pub max_bounces: i64,
pub window_weeks: i32,
pub spread_days: i32,
}
impl Default for AutoBlacklistConfig {
fn default() -> Self {
Self {
max_bounces: 5,
window_weeks: 13,
spread_days: 7,
}
}
}
pub fn default_domain_for(target_model: &str) -> Option<serde_json::Value> {
match target_model {
"crm_lead" | "crm_deal" => None,
"event_registration" => None,
"selling_customer" => Some(serde_json::json!([
{"field": "email", "op": "!=", "value": ""}
])),
_ => None,
}
}
fn domain_to_store(
parsed: &CompiledDomain,
target_model: &str,
) -> Result<CompiledDomain, DomainInvalid> {
if parsed.terms.is_empty() {
if let Some(default_raw) = default_domain_for(target_model) {
return parse_domain(&default_raw);
}
}
Ok(parsed.clone())
}
#[async_trait::async_trait]
pub trait PartyRecipientResolver: Send + Sync {
async fn resolve(&self, domain: &CompiledDomain) -> Result<Vec<ResolvedRecipient>, String>;
}
#[async_trait::async_trait]
pub trait TargetRecipientResolver: Send + Sync {
fn target_model(&self) -> &'static str;
async fn resolve(&self, domain: &CompiledDomain) -> Result<Vec<ResolvedRecipient>, String>;
}
struct PartyResolverAdapter {
inner: Arc<dyn PartyRecipientResolver>,
}
#[async_trait::async_trait]
impl TargetRecipientResolver for PartyResolverAdapter {
fn target_model(&self) -> &'static str {
"party"
}
async fn resolve(&self, domain: &CompiledDomain) -> Result<Vec<ResolvedRecipient>, String> {
self.inner.resolve(domain).await
}
}
pub trait MailingEventSink: Send + Sync {
fn mailing_launched(&self, mailing_id: Uuid, campaign_id: Option<Uuid>, schedule_type: &str);
fn recipient_suppressed(&self, mailing_id: Uuid, trace_id: Uuid, email: &str, cause: &str);
fn auto_blacklist_added(&self, email: &str);
fn ab_test_winner_promoted(&self, ab_test_id: Uuid, winner_mailing_id: Uuid);
}
#[derive(Debug, Default, Clone, Copy)]
pub struct TracingEventSink;
impl MailingEventSink for TracingEventSink {
fn mailing_launched(&self, mailing_id: Uuid, campaign_id: Option<Uuid>, schedule_type: &str) {
tracing::info!(
mailing_id = %mailing_id,
campaign_id = ?campaign_id,
schedule_type,
"MailingQueuedForSend"
);
}
fn recipient_suppressed(&self, mailing_id: Uuid, trace_id: Uuid, email: &str, cause: &str) {
tracing::info!(
mailing_id = %mailing_id,
trace_id = %trace_id,
recipient_email = email,
cause,
"RecipientSuppressedAtSend"
);
}
fn auto_blacklist_added(&self, email: &str) {
tracing::info!(email, "auto_blacklist_added");
}
fn ab_test_winner_promoted(&self, ab_test_id: Uuid, winner_mailing_id: Uuid) {
tracing::info!(
ab_test_id = %ab_test_id,
winner_mailing_id = %winner_mailing_id,
"ab_test_winner_promoted"
);
}
}
#[derive(Debug, Clone, Default)]
pub struct MailingUpsertCommand {
pub subject: String,
pub preview: Option<String>,
pub body_html: String,
pub email_from: String,
pub reply_to: Option<String>,
pub mailing_domain_raw: serde_json::Value,
pub target_model: String,
pub schedule_type: String,
pub schedule_date: Option<chrono::DateTime<chrono::Utc>>,
pub use_exclusion_list: bool,
pub campaign_id: Option<Uuid>,
pub source_id: Option<Uuid>,
pub ab_test_id: Option<Uuid>,
pub ab_testing_enabled: bool,
pub ab_testing_pc: i32,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct SweepOutcome {
pub claimed: usize,
pub parked: usize,
pub completed: usize,
pub completed_empty: usize,
pub still_sending: usize,
pub recipients_resolved: usize,
pub minted: usize,
pub repaired: usize,
pub enqueued: usize,
pub fence_skips: usize,
pub suppressed_blacklist: usize,
pub suppressed_optout: usize,
pub suppressed_dup: usize,
pub skipped_seen: usize,
pub skipped_ab_nonmember: usize,
pub reconcile_sent: usize,
pub reconcile_failed: usize,
pub auto_blacklisted: u64,
pub ab_promotions: usize,
pub ab_seam_refusals: usize,
pub sms_walks_marked: usize,
pub sms_traces_advanced: usize,
pub sms_trace_skips: usize,
pub sms_mailings_completed: usize,
}
pub struct MailingWriteService {
pool: sqlx::PgPool,
cfg: MailingSendConfig,
blacklist_cfg: AutoBlacklistConfig,
target_resolvers: std::collections::HashMap<String, Arc<dyn TargetRecipientResolver>>,
phone_resolvers:
std::collections::HashMap<String, Arc<dyn crate::application::service::sms_targeting_service::TargetPhoneResolver>>,
invoiced_amounts: Arc<dyn crate::application::service::sale_invoiced_amount_port::SaleInvoicedAmountPort>,
events: Arc<dyn MailingEventSink>,
messages: MessageWriteService,
queue: MailQueueWriteService,
sms_pump: SmsDeliveryPumpService,
}
impl MailingWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self {
messages: MessageWriteService::new(pool.clone()),
queue: MailQueueWriteService::new(pool.clone()),
sms_pump: SmsDeliveryPumpService::new(pool.clone()),
pool,
cfg: MailingSendConfig::default(),
blacklist_cfg: AutoBlacklistConfig::default(),
target_resolvers: std::collections::HashMap::new(),
phone_resolvers: std::collections::HashMap::new(),
invoiced_amounts: Arc::new(
crate::application::service::sale_invoiced_amount_port::RefusingSaleInvoicedAmount,
),
events: Arc::new(TracingEventSink),
}
}
pub fn with_send_config(mut self, cfg: MailingSendConfig) -> Self {
self.cfg = cfg;
self
}
pub fn with_auto_blacklist_config(mut self, cfg: AutoBlacklistConfig) -> Self {
self.blacklist_cfg = cfg;
self
}
pub fn with_party_resolver(mut self, r: Arc<dyn PartyRecipientResolver>) -> Self {
self.target_resolvers
.insert("party".to_string(), Arc::new(PartyResolverAdapter { inner: r }));
self
}
pub fn with_target_resolver(
mut self,
r: Arc<dyn TargetRecipientResolver>,
) -> Self {
self.target_resolvers
.insert(r.target_model().to_string(), r);
self
}
pub fn with_event_registration_resolver(
self,
r: Arc<
dyn crate::application::service::event_registration_target_port::EventRegistrationTargetResolver,
>,
) -> Self {
self.with_target_resolver(Arc::new(
crate::application::service::event_registration_target_port::EventRegistrationTargetAdapter::new(r),
))
}
pub fn with_phone_resolver(
mut self,
r: Arc<dyn crate::application::service::sms_targeting_service::TargetPhoneResolver>,
) -> Self {
self.phone_resolvers.insert(r.target_model().to_string(), r);
self
}
pub fn phone_resolver_for(
&self,
target_model: &str,
) -> Option<Arc<dyn crate::application::service::sms_targeting_service::TargetPhoneResolver>> {
self.phone_resolvers.get(target_model).cloned()
}
pub fn with_sale_invoiced_amount_port(
mut self,
port: Arc<dyn crate::application::service::sale_invoiced_amount_port::SaleInvoicedAmountPort>,
) -> Self {
self.invoiced_amounts = port;
self
}
pub fn with_sale_invoiced_amount_slot(
mut self,
slot: &crate::application::service::sale_invoiced_amount_port::SaleInvoicedAmountSlot,
) -> Self {
let slot = slot.clone();
self.invoiced_amounts = Arc::new(slot);
self
}
pub fn with_event_sink(mut self, sink: Arc<dyn MailingEventSink>) -> Self {
self.events = sink;
self
}
pub async fn create_mailing(&self, cmd: &MailingUpsertCommand) -> Result<Uuid, MailingWriteError> {
if cmd.subject.trim().is_empty() {
return Err(MailingWriteError::Invalid("subject is required".into()));
}
if cmd.body_html.trim().is_empty() {
return Err(MailingWriteError::Invalid("body_html is required".into()));
}
if !cmd.email_from.contains('@') {
return Err(MailingWriteError::Invalid(format!(
"email_from must be an address: {}",
cmd.email_from
)));
}
let domain = parse_domain(&cmd.mailing_domain_raw)?;
let target_model = match cmd.target_model.parse::<crate::domain::entity::MailingTargetModel>() {
Ok(_) => cmd.target_model.clone(),
Err(_) => {
return Err(MailingWriteError::Invalid(format!(
"target_model must be one of the closed enum's values \
(mailing_contact, party, crm_lead, crm_deal, selling_customer, \
event_registration), not {}",
cmd.target_model
)))
}
};
let domain = domain_to_store(&domain, &target_model)?;
let id = Uuid::new_v4();
let mut tx = self.pool.begin().await?;
MailingSendRepository::insert_mailing(
&mut tx,
id,
&cmd.subject,
cmd.preview.as_deref(),
&cmd.body_html,
&cmd.email_from,
cmd.reply_to.as_deref(),
&canonical_domain_json(&domain),
&target_model,
&cmd.schedule_type,
cmd.schedule_date,
cmd.use_exclusion_list,
cmd.campaign_id,
cmd.source_id,
cmd.ab_testing_enabled,
cmd.ab_testing_pc,
cmd.ab_test_id,
)
.await?;
tx.commit().await?;
Ok(id)
}
pub async fn update_mailing(
&self,
id: Uuid,
cmd: &MailingUpsertCommand,
) -> Result<(), MailingWriteError> {
let domain = parse_domain(&cmd.mailing_domain_raw)?;
let mut tx = self.pool.begin().await?;
let updated = MailingSendRepository::update_draft(
&mut tx,
id,
&cmd.subject,
cmd.preview.as_deref(),
&cmd.body_html,
&cmd.email_from,
cmd.reply_to.as_deref(),
&canonical_domain_json(&domain),
&cmd.schedule_type,
cmd.schedule_date,
cmd.use_exclusion_list,
)
.await?;
tx.commit().await?;
if !updated {
match self.live_state(id).await? {
None => Err(MailingWriteError::NotFound(format!("mailing {id}"))),
Some(s) => Err(MailingWriteError::Conflict(format!(
"mailing {id} is {s}, only draft is editable"
))),
}
} else {
Ok(())
}
}
pub async fn launch(
&self,
id: Uuid,
schedule_type: &str,
schedule_date: Option<chrono::DateTime<chrono::Utc>>,
) -> Result<(), MailingWriteError> {
let (schedule_type, schedule_date) = match schedule_type {
"immediate" => ("immediate".to_string(), None),
"scheduled" => {
let d = schedule_date.ok_or_else(|| {
MailingWriteError::Invalid("scheduled launch needs schedule_date".into())
})?;
("scheduled".to_string(), Some(d))
}
other => {
return Err(MailingWriteError::Invalid(format!(
"schedule_type must be immediate or scheduled, not {other}"
)))
}
};
let mut tx = self.pool.begin().await?;
let launch_shape = MailingSendRepository::mailing_launch_shape(&mut tx, id).await?;
match launch_shape {
None => {
tx.rollback().await?;
return Err(MailingWriteError::NotFound(format!("mailing {id}")));
}
Some((mailing_type, body_plaintext)) => {
if mailing_type == "sms"
&& body_plaintext
.as_deref()
.map(str::trim)
.filter(|b| !b.is_empty())
.is_none()
{
tx.rollback().await?;
return Err(MailingWriteError::Invalid(format!(
"sms mailing {id} needs body_plaintext before launch"
)));
}
}
}
let queued = MailingSendRepository::queue_mailing(
&mut tx,
id,
&schedule_type,
schedule_date,
)
.await?;
tx.commit().await?;
if !queued {
return Err(match self.live_state(id).await? {
None => MailingWriteError::NotFound(format!("mailing {id}")),
Some(s) => MailingWriteError::Conflict(format!(
"mailing {id} is {s}, only draft can launch"
)),
});
}
let campaign_id = {
let mut tx = self.pool.begin().await?;
let cid = MailingSendRepository::mailing_campaign_id(&mut tx, id).await?;
tx.commit().await?;
cid
};
self.events
.mailing_launched(id, campaign_id, &schedule_type);
Ok(())
}
pub async fn cancel(&self, id: Uuid) -> Result<(), MailingWriteError> {
let mut tx = self.pool.begin().await?;
let canceled = MailingSendRepository::cancel_mailing(&mut tx, id).await?;
tx.commit().await?;
if !canceled {
return Err(match self.live_state(id).await? {
None => MailingWriteError::NotFound(format!("mailing {id}")),
Some(s) => MailingWriteError::Conflict(format!(
"mailing {id} is {s} and cannot be canceled from this state"
)),
});
}
Ok(())
}
pub async fn requeue_after_domain_fix(&self, id: Uuid) -> Result<(), MailingWriteError> {
let mut tx = self.pool.begin().await?;
let unparked = MailingSendRepository::unpark_mailing(&mut tx, id).await?;
tx.commit().await?;
if !unparked {
return match self.live_state(id).await? {
None => Err(MailingWriteError::NotFound(format!("mailing {id}"))),
Some(_) => Err(MailingWriteError::Conflict(format!(
"mailing {id} is not parked (no send_error to clear)"
))),
};
}
Ok(())
}
pub async fn retry_failed(&self, id: Uuid) -> Result<u64, MailingWriteError> {
let mut tx = self.pool.begin().await?;
let failed = MailingSendRepository::count_failed_traces(&mut tx, id).await?;
if failed == 0 {
tx.rollback().await?;
return Err(MailingWriteError::Invalid(format!(
"mailing {id} has no failed (error/bounce) traces to retry"
)));
}
let soft_deleted = MailingSendRepository::soft_delete_failed_traces(&mut tx, id, 1000).await?;
let requeued = MailingSendRepository::requeue_mailing(&mut tx, id).await?;
tx.commit().await?;
if !requeued {
return Err(MailingWriteError::Conflict(format!(
"mailing {id} is not done — retry_failed only exits done"
)));
}
Ok(soft_deleted)
}
async fn live_state(&self, id: Uuid) -> Result<Option<String>, MailingWriteError> {
let mut tx = self.pool.begin().await?;
let row = MailingSendRepository::find_live_state(&mut tx, id).await?;
tx.commit().await?;
Ok(row.map(|(_, state, _)| state))
}
pub async fn send_queue_sweep(&self) -> Result<SweepOutcome, MailingWriteError> {
let mut out = SweepOutcome::default();
let mut parked_ids = HashSet::new();
let claims = {
let mut tx = self.pool.begin().await?;
let claimed =
MailingSendRepository::claim_due_mailings(&mut tx, self.cfg.claim_batch).await?;
for m in &claimed {
match parse_domain(&m.mailing_domain) {
Ok(_) => {
MailingSendRepository::flip_sending(&mut tx, m.id).await?;
}
Err(e) => {
MailingSendRepository::park_mailing(&mut tx, m.id, &e.to_string()).await?;
parked_ids.insert(m.id);
out.parked += 1;
tracing::warn!(
mailing_id = %m.id,
error = %e,
"mailing parked: domain failed to parse at send time"
);
}
}
}
tx.commit().await?;
claimed
};
out.claimed = claims.len();
for mailing in &claims {
if parked_ids.contains(&mailing.id) {
continue;
}
self.drive_mailing(mailing, &mut out).await?;
}
{
let mut tx = self.pool.begin().await?;
let settled = TraceRepository::settled_mails_for_reconcile(&mut tx).await?;
for (trace_id, state, failure_type) in settled {
let moved = match state.as_str() {
"sent" => TraceRepository::set_sent(&mut tx, trace_id).await?,
_ => {
let ftype = failure_type.unwrap_or_else(|| "mail_smtp".into());
TraceRepository::set_failed(&mut tx, trace_id, &ftype, None).await?
}
};
if moved {
if state == "sent" {
out.reconcile_sent += 1;
} else {
out.reconcile_failed += 1;
}
}
}
tx.commit().await?;
}
{
let mut tx = self.pool.begin().await?;
out.auto_blacklisted = MailingSendRepository::auto_blacklist_sweep(
&mut tx,
self.blacklist_cfg.window_weeks,
self.blacklist_cfg.max_bounces,
self.blacklist_cfg.spread_days,
)
.await?;
tx.commit().await?;
}
self.promote_due_ab_tests(&mut out).await?;
{
let pumped = self.sms_pump.pump_once().await?;
out.sms_traces_advanced += pumped.traces_advanced;
out.sms_trace_skips += pumped.trace_skips;
out.sms_mailings_completed += pumped.mailings_completed;
}
Ok(out)
}
async fn drive_mailing(
&self,
m: &ClaimedMailing,
out: &mut SweepOutcome,
) -> Result<(), MailingWriteError> {
if m.mailing_type == "sms" {
if m.sms_walk_done {
return Ok(());
}
let mut tx = self.pool.begin().await?;
MailingSendRepository::park_mailing(
&mut tx,
m.id,
"sms-type mailing reached the mail send walk — the sms send walk is not composed at this seam",
)
.await?;
tx.commit().await?;
out.parked += 1;
tracing::warn!(
mailing_id = %m.id,
"sms mailing parked: the mail send walk refuses to route an sms audience"
);
return Ok(());
}
const REPAIR_GRACE_MINUTES: i32 = 5;
{
let mut tx = self.pool.begin().await?;
let orphans =
TraceRepository::outgoing_without_mail(&mut tx, m.id, REPAIR_GRACE_MINUTES)
.await?;
tx.commit().await?;
for (trace_id, _rid, email) in orphans {
match self.enqueue_one(m, trace_id, &email).await {
Ok(()) => out.repaired += 1,
Err(e) => {
tracing::error!(
trace_id = %trace_id,
error = %e,
"orphan repair enqueue failed"
);
}
}
}
}
let domain = parse_domain(&m.mailing_domain)?; let recipients = match m.target_model.as_str() {
"mailing_contact" => {
let mut tx = self.pool.begin().await?;
let r = MailingSendRepository::resolve_contact_recipients(
&mut tx,
&domain,
self.cfg.resolver_hard_cap,
)
.await?;
tx.commit().await?;
r
}
"party" | "crm_lead" | "crm_deal" | "selling_customer" | "event_registration" => {
match self.target_resolvers.get(m.target_model.as_str()) {
Some(resolver) => resolver.resolve(&domain).await.map_err(|e| {
MailingWriteError::Conflict(format!(
"{} resolver failed: {e}",
m.target_model
))
})?,
None => {
let mut tx = self.pool.begin().await?;
MailingSendRepository::park_mailing(
&mut tx,
m.id,
&format!(
"target_model='{}' but no {} resolver is composed",
m.target_model, m.target_model
),
)
.await?;
tx.commit().await?;
out.parked += 1;
return Ok(());
}
}
}
other => {
let mut tx = self.pool.begin().await?;
MailingSendRepository::park_mailing(
&mut tx,
m.id,
&format!("unknown target_model '{other}'"),
)
.await?;
tx.commit().await?;
out.parked += 1;
return Ok(());
}
};
if recipients.len() as i64 >= self.cfg.resolver_hard_cap {
return Err(MailingWriteError::Conflict(format!(
"mailing {} resolved {} recipients (>= hard cap {}) — refuse to silently truncate",
m.id,
recipients.len(),
self.cfg.resolver_hard_cap
)));
}
out.recipients_resolved += recipients.len();
if recipients.is_empty() {
let mut tx = self.pool.begin().await?;
let flipped = MailingSendRepository::complete_empty(&mut tx, m.id).await?;
tx.commit().await?;
if flipped {
out.completed_empty += 1;
}
return Ok(());
}
let mut tx = self.pool.begin().await?;
let disposition = MailingSendRepository::recipient_disposition(
&mut tx,
m.id,
m.campaign_id.filter(|_| m.ab_test_id.is_some()),
)
.await?;
let emails: Vec<String> = recipients.iter().map(|r| r.email.to_lowercase()).collect();
let blacklisted = if m.use_exclusion_list {
MailingSendRepository::blacklisted_emails(&mut tx, &emails).await?
} else {
Vec::new()
};
let opted_out = MailingSendRepository::opted_out_emails(&mut tx, &emails).await?;
tx.commit().await?;
let blacklisted: HashSet<String> = blacklisted.into_iter().collect();
let opted_out: HashSet<String> = opted_out.into_iter().collect();
let mut seen_active: HashSet<Uuid> = HashSet::new();
let mut suppressed_here: HashSet<Uuid> = HashSet::new();
for (rid, active, here) in disposition {
if active {
seen_active.insert(rid);
}
if here {
suppressed_here.insert(rid);
}
}
let ab_seed = if m.ab_testing_enabled {
match m.ab_test_id {
Some(ab_id) => {
let mut tx = self.pool.begin().await?;
let seed = MailingSendRepository::sampling_seed(&mut tx, ab_id).await?;
tx.commit().await?;
match seed {
Some(s) => Some(s),
None => {
let mut tx = self.pool.begin().await?;
MailingSendRepository::park_mailing(
&mut tx,
m.id,
"ab_testing_enabled but the bound A/B test has no sampling seed",
)
.await?;
tx.commit().await?;
out.parked += 1;
return Ok(());
}
}
}
None => {
let mut tx = self.pool.begin().await?;
MailingSendRepository::park_mailing(
&mut tx,
m.id,
"ab_testing_enabled but no ab_test_id is bound",
)
.await?;
tx.commit().await?;
out.parked += 1;
return Ok(());
}
}
} else {
None
};
let mut seen_emails: HashSet<String> = HashSet::new();
let mut to_send: Vec<&ResolvedRecipient> = Vec::new();
let mut cancel_work: Vec<(&ResolvedRecipient, &'static str)> = Vec::new();
for r in &recipients {
let email_lc = r.email.to_lowercase();
if seen_active.contains(&r.recipient_id) {
out.skipped_seen += 1;
seen_emails.insert(email_lc);
continue;
}
if suppressed_here.contains(&r.recipient_id) {
out.skipped_seen += 1;
seen_emails.insert(email_lc);
continue;
}
if !seen_emails.insert(email_lc.clone()) {
cancel_work.push((r, "mail_dup"));
out.suppressed_dup += 1;
continue;
}
if blacklisted.contains(&email_lc) {
cancel_work.push((r, "mail_bl"));
out.suppressed_blacklist += 1;
continue;
}
if opted_out.contains(&email_lc) {
cancel_work.push((r, "mail_optout"));
out.suppressed_optout += 1;
continue;
}
if let Some(seed) = &ab_seed {
if !ab_member(seed, &email_lc, m.ab_testing_pc) {
out.skipped_ab_nonmember += 1;
continue;
}
}
to_send.push(r);
}
for (r, cause) in &cancel_work {
let trace_id = Uuid::new_v4();
let mut tx = self.pool.begin().await?;
match TraceRepository::mint_trace(
&mut tx,
trace_id,
m.id,
m.campaign_id,
&m.target_model,
r.recipient_id,
&r.email,
"cancel",
Some(cause),
false,
)
.await
{
Ok(_) => {
tx.commit().await?;
self.events
.recipient_suppressed(m.id, trace_id, &r.email, cause);
}
Err(e) if is_unique_violation(&e) => {
tx.rollback().await?;
out.fence_skips += 1;
}
Err(e) => return Err(e.into()),
}
}
let budget = self.cfg.recipients_per_pass as usize;
let remaining = to_send.len().saturating_sub(budget);
let work: Vec<&ResolvedRecipient> = to_send.into_iter().take(budget).collect();
let mut pacer = Pacer::new(self.cfg.max_enqueues_per_second);
for batch in work.chunks(self.cfg.mail_enqueue_batch.max(1)) {
let mut tx = self.pool.begin().await?;
let mut minted: Vec<(Uuid, &ResolvedRecipient)> = Vec::new();
for r in batch {
let r: &ResolvedRecipient = r;
let trace_id = Uuid::new_v4();
let inserted = TraceRepository::mint_trace_fenced(
&mut tx,
trace_id,
m.id,
m.campaign_id,
&m.target_model,
r.recipient_id,
&r.email,
"outgoing",
None,
false,
)
.await?;
if inserted {
minted.push((trace_id, r));
out.minted += 1;
} else {
tracing::warn!(
mailing_id = %m.id,
recipient_id = %r.recipient_id,
"mint fence: live trace already exists, skipping"
);
out.fence_skips += 1;
}
}
tx.commit().await?;
for (trace_id, r) in minted {
match self.enqueue_one(m, trace_id, &r.email).await {
Ok(()) => out.enqueued += 1,
Err(e) => {
let mut tx = self.pool.begin().await?;
TraceRepository::set_failed(
&mut tx,
trace_id,
"mail_smtp",
Some(&e.to_string()),
)
.await?;
tx.commit().await?;
tracing::error!(
trace_id = %trace_id,
error = %e,
"enqueue failed at the mail seam"
);
}
}
pacer.wait().await;
}
if self.cfg.inter_batch_delay_ms > 0 {
tokio::time::sleep(Duration::from_millis(self.cfg.inter_batch_delay_ms)).await;
}
}
if remaining > 0 {
out.still_sending += 1;
} else if m.mailing_type == "sms" {
let mut tx = self.pool.begin().await?;
let marked = MailingSendRepository::mark_sms_walk_complete(&mut tx, m.id).await?;
tx.commit().await?;
if marked {
out.sms_walks_marked += 1;
}
} else {
let mut tx = self.pool.begin().await?;
let flipped = MailingSendRepository::complete_mailing(&mut tx, m.id).await?;
tx.commit().await?;
if flipped {
out.completed += 1;
}
}
Ok(())
}
async fn enqueue_one(
&self,
m: &ClaimedMailing,
trace_id: Uuid,
email: &str,
) -> Result<(), MailingWriteError> {
let posted = self
.messages
.message_post(MessagePostCommand {
body: m.body_html.clone(),
subject: Some(m.subject.clone()),
message_type: "email".into(),
subtype_id: None,
subtype_name: None,
is_internal: false,
author_id: None,
author_guest_id: None,
email_from: Some(m.email_from.clone()),
reply_to: m.reply_to.clone(),
model: Some("mailing".into()),
res_id: Some(m.id),
record_name: Some(m.subject.clone()),
recipients: Vec::new(),
})
.await
.map_err(|e| MailingWriteError::MailSeam(e.to_string()))?;
let mail_id = self
.queue
.enqueue(
posted.message_id,
email,
None,
m.reply_to.as_deref(),
None,
None,
Some("mailing"),
Some(m.id),
)
.await?;
let mut tx = self.pool.begin().await?;
TraceRepository::attach_mail_id(&mut tx, trace_id, mail_id).await?;
tx.commit().await?;
Ok(())
}
async fn promote_due_ab_tests(&self, out: &mut SweepOutcome) -> Result<(), MailingWriteError> {
let mut tx = self.pool.begin().await?;
let due = MailingSendRepository::ab_tests_due_for_promotion(&mut tx).await?;
tx.commit().await?;
for (ab_test_id, selection) in due {
if selection == "manual" {
continue;
}
let winner = if selection == "sale_invoiced_amount" {
match self.rank_variants_by_invoiced_amount(ab_test_id).await {
RankBySeamOutcome::Ranked(winner) => Some(winner),
RankBySeamOutcome::Wait => None,
RankBySeamOutcome::SeamRefused(detail) => {
out.ab_seam_refusals += 1;
tracing::warn!(
ab_test_id = %ab_test_id,
detail = %detail,
"A/B promotion skipped: the invoiced-amount seam refused"
);
continue;
}
}
} else {
let mut tx = self.pool.begin().await?;
let ranked =
MailingSendRepository::rank_variants_by_metric(&mut tx, ab_test_id, &selection)
.await?;
tx.commit().await?;
ranked.first().map(|(id, _)| *id)
};
let Some(winner_mailing_id) = winner else {
continue; };
let mut tx = self.pool.begin().await?;
let promoted_id = Uuid::new_v4();
MailingSendRepository::promote_winner(&mut tx, promoted_id, winner_mailing_id, ab_test_id)
.await?;
MailingSendRepository::queue_promoted_winner(&mut tx, promoted_id).await?;
let stamped =
MailingSendRepository::complete_ab_test(&mut tx, ab_test_id, winner_mailing_id)
.await?;
tx.commit().await?;
if stamped {
out.ab_promotions += 1;
self.events
.ab_test_winner_promoted(ab_test_id, winner_mailing_id);
}
}
Ok(())
}
async fn rank_variants_by_invoiced_amount(
&self,
ab_test_id: Uuid,
) -> RankBySeamOutcome {
use crate::application::service::sale_invoiced_amount_port::SaleInvoicedAmountPort;
let mut tx = match self.pool.begin().await {
Ok(tx) => tx,
Err(e) => return RankBySeamOutcome::SeamRefused(e.to_string()),
};
let variants = match MailingSendRepository::variant_source_ids(&mut tx, ab_test_id).await {
Ok(v) => v,
Err(e) => {
let _ = tx.rollback().await;
return RankBySeamOutcome::SeamRefused(e.to_string());
}
};
tx.commit().await.ok();
let mut per_source = std::collections::HashMap::new();
for (_, source) in variants.iter().filter_map(|(id, s)| s.map(|s| (*id, s))) {
if !per_source.contains_key(&source) {
match self.invoiced_amounts.invoiced_amount_for_source(source).await {
Ok(amount) => {
per_source.insert(source, amount.amount_untaxed_total);
}
Err(e) => return RankBySeamOutcome::SeamRefused(e.to_string()),
}
}
}
let mut ranked: Vec<(Uuid, Option<rust_decimal::Decimal>)> = variants
.into_iter()
.map(|(id, source)| {
let amount = source.and_then(|s| per_source.get(&s).copied());
(id, amount)
})
.collect();
if ranked.is_empty() {
return RankBySeamOutcome::Wait;
}
ranked.sort_by(|a, b| {
b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0))
});
if ranked.first().map(|(_, amount)| amount.is_none()).unwrap_or(true) {
return RankBySeamOutcome::Wait;
}
RankBySeamOutcome::Ranked(ranked[0].0)
}
}
enum RankBySeamOutcome {
Ranked(Uuid),
Wait,
SeamRefused(String),
}
fn field_of(raw: &str) -> Option<DomainField> {
match raw {
"email" => Some(DomainField::Email),
"name" => Some(DomainField::Name),
"first_name" => Some(DomainField::FirstName),
"last_name" => Some(DomainField::LastName),
"company_name" => Some(DomainField::CompanyName),
"country_code" => Some(DomainField::CountryCode),
"mailing_audience_id" => Some(DomainField::MailingAudienceId),
_ => None,
}
}
pub fn parse_domain(raw: &serde_json::Value) -> Result<CompiledDomain, DomainInvalid> {
let terms_json = raw.as_array().ok_or(DomainInvalid::NotAnArray)?;
let mut terms = Vec::with_capacity(terms_json.len());
for (index, tj) in terms_json.iter().enumerate() {
let (field_raw, op_raw, value): (String, String, serde_json::Value) = if let Some(arr) =
tj.as_array()
{
if arr.len() != 3 {
return Err(DomainInvalid::MalformedTerm { index });
}
let f = arr[0]
.as_str()
.ok_or(DomainInvalid::FieldNotString { index })?
.to_string();
let o = match arr[1].as_str() {
Some(s) => s.to_string(),
None => return Err(DomainInvalid::UnknownOp { index }),
};
(f, o, arr[2].clone())
} else if let Some(obj) = tj.as_object() {
let f = obj
.get("field")
.and_then(|v| v.as_str())
.ok_or(DomainInvalid::FieldNotString { index })?
.to_string();
let o = obj
.get("op")
.and_then(|v| v.as_str())
.ok_or(DomainInvalid::UnknownOp { index })?
.to_string();
let v = obj
.get("value")
.cloned()
.ok_or(DomainInvalid::MalformedTerm { index })?;
(f, o, v)
} else {
return Err(DomainInvalid::MalformedTerm { index });
};
let field = field_of(&field_raw).ok_or(DomainInvalid::UnknownField(field_raw.clone()))?;
let op = match op_raw.as_str() {
"=" => DomainOp::Eq,
"!=" => DomainOp::Ne,
"in" => DomainOp::In,
"not in" => DomainOp::NotIn,
"like" => DomainOp::Like,
"not like" => DomainOp::NotLike,
_ => return Err(DomainInvalid::UnknownOp { index }),
};
let values: Vec<String> = match (&op, &value) {
(DomainOp::In | DomainOp::NotIn, serde_json::Value::Array(a)) => {
let mut vs = Vec::with_capacity(a.len());
for v in a {
vs.push(
v.as_str()
.ok_or(DomainInvalid::ValueNotArray { index })?
.to_string(),
);
}
vs
}
(DomainOp::In | DomainOp::NotIn, _) => {
return Err(DomainInvalid::ValueNotArray { index })
}
(_, serde_json::Value::String(s)) => vec![s.clone()],
(_, serde_json::Value::Array(a)) if a.len() == 1 => vec![a[0]
.as_str()
.ok_or(DomainInvalid::ValueNotScalar { index })?
.to_string()],
_ => return Err(DomainInvalid::ValueNotScalar { index }),
};
if field == DomainField::MailingAudienceId {
for v in &values {
if Uuid::parse_str(v).is_err() {
return Err(DomainInvalid::BadAudienceUuid {
index,
raw: v.clone(),
});
}
}
}
terms.push(DomainTerm { field, op, values });
}
Ok(CompiledDomain { terms })
}
pub fn canonical_domain_json(domain: &CompiledDomain) -> serde_json::Value {
serde_json::Value::Array(
domain
.terms
.iter()
.map(|t| {
let value = if t.values.len() == 1
&& !matches!(t.op, DomainOp::In | DomainOp::NotIn)
{
serde_json::Value::String(t.values[0].clone())
} else {
serde_json::Value::Array(
t.values
.iter()
.map(|v| serde_json::Value::String(v.clone()))
.collect(),
)
};
serde_json::json!({
"field": field_name(t.field),
"op": op_name(t.op),
"value": value,
})
})
.collect(),
)
}
fn field_name(f: DomainField) -> &'static str {
match f {
DomainField::Email => "email",
DomainField::Name => "name",
DomainField::FirstName => "first_name",
DomainField::LastName => "last_name",
DomainField::CompanyName => "company_name",
DomainField::CountryCode => "country_code",
DomainField::MailingAudienceId => "mailing_audience_id",
}
}
fn op_name(op: DomainOp) -> &'static str {
match op {
DomainOp::Eq => "=",
DomainOp::Ne => "!=",
DomainOp::In => "in",
DomainOp::NotIn => "not in",
DomainOp::Like => "like",
DomainOp::NotLike => "not like",
}
}
pub fn ab_member(seed: &str, recipient_email: &str, pc: i32) -> bool {
let mut mac = match Hmac::<Sha256>::new_from_slice(seed.as_bytes()) {
Ok(m) => m,
Err(_) => return false,
};
mac.update(recipient_email.as_bytes());
let digest = mac.finalize().into_bytes();
let bucket = u64::from_be_bytes([
digest[0], digest[1], digest[2], digest[3],
digest[4], digest[5], digest[6], digest[7],
]);
(bucket % 100) < pc as u64
}
fn is_unique_violation(err: &sqlx::Error) -> bool {
matches!(err, sqlx::Error::Database(db) if db.code().as_deref() == Some("23505"))
}
struct Pacer {
min_interval: Option<Duration>,
last: Option<Instant>,
}
impl Pacer {
fn new(max_per_second: Option<u64>) -> Self {
Self {
min_interval: max_per_second
.filter(|eps| *eps > 0)
.map(|eps| Duration::from_millis(1000 / eps)),
last: None,
}
}
async fn wait(&mut self) {
if let (Some(min), Some(last)) = (self.min_interval, self.last) {
let elapsed = last.elapsed();
if elapsed < min {
tokio::time::sleep(min - elapsed).await;
}
}
self.last = Some(Instant::now());
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_accepts_odoo_triples() {
let raw = serde_json::json!([["email", "like", "acme"], ["country_code", "=", "ID"]]);
let d = parse_domain(&raw).unwrap();
assert_eq!(d.terms.len(), 2);
assert_eq!(canonical_domain_json(&d)[0]["op"], "like");
}
#[test]
fn parse_accepts_canonical_objects() {
let raw = serde_json::json!([
{"field": "email", "op": "in", "value": ["a@x.id", "b@x.id"]}
]);
let d = parse_domain(&raw).unwrap();
assert_eq!(d.terms[0].op, DomainOp::In);
}
#[test]
fn parse_refuses_unknown_fields_loudly() {
let raw = serde_json::json!([["1; DROP TABLE", "=", "x"]]);
assert!(matches!(
parse_domain(&raw),
Err(DomainInvalid::UnknownField(_))
));
}
#[test]
fn parse_refuses_unknown_ops_loudly() {
let raw = serde_json::json!([["email", "ilike", "x"]]);
assert!(matches!(parse_domain(&raw), Err(DomainInvalid::UnknownOp { .. })));
}
#[test]
fn parse_refuses_non_array_loudly() {
let raw = serde_json::json!({"field": "email"});
assert!(matches!(parse_domain(&raw), Err(DomainInvalid::NotAnArray)));
}
#[test]
fn parse_refuses_bad_audience_uuid_loudly() {
let raw = serde_json::json!([["mailing_audience_id", "=", "not-a-uuid"]]);
assert!(matches!(
parse_domain(&raw),
Err(DomainInvalid::BadAudienceUuid { .. })
));
}
#[test]
fn empty_domain_is_legal_matches_all() {
let d = parse_domain(&serde_json::json!([])).unwrap();
assert!(d.terms.is_empty());
}
#[test]
#[expect(clippy::expect_used, reason = "unit test asserts the compiled-domain contract holds here")]
fn default_domains_are_declarative_and_whitelist_clean() {
for target in ["crm_lead", "crm_deal", "selling_customer"] {
match default_domain_for(target) {
Some(raw) => {
let d = parse_domain(&raw)
.unwrap_or_else(|e| panic!("{target} default must parse: {e}"));
assert!(!d.terms.is_empty(), "{target} default is non-empty");
let stored = domain_to_store(&CompiledDomain::default(), target)
.expect("empty domain adopts the default");
assert_eq!(stored, d, "{target}: adoption is exact");
}
None => {
let stored = domain_to_store(&CompiledDomain::default(), target).unwrap();
assert!(stored.terms.is_empty(), "{target} must adopt nothing");
}
}
}
let authored = parse_domain(&serde_json::json!([["email", "like", "acme"]])).unwrap();
let stored = domain_to_store(&authored, "selling_customer").unwrap();
assert_eq!(stored, authored);
}
#[test]
fn ab_membership_is_deterministic_and_roughly_proportional() {
let seed = "cafebabedeadbeef";
let mut members = 0u32;
let n = 2000;
for i in 0..n {
let email = format!("user{i}@example.id");
assert_eq!(ab_member(seed, &email, 20), ab_member(seed, &email, 20));
if ab_member(seed, &email, 20) {
members += 1;
}
}
let pct = members as f64 / n as f64 * 100.0;
assert!((12.0..=28.0).contains(&pct), "membership pct was {pct}");
let flipped = (0..100)
.filter(|i| {
ab_member(seed, &format!("u{i}@x.id"), 50)
!= ab_member("other-seed", &format!("u{i}@x.id"), 50)
})
.count();
assert!(flipped > 0);
}
}