use sqlx::PgPool;
use crate::application::service::trace_write_service::SMS_FAILURE_CODES;
use crate::infrastructure::persistence::mailing_send_repository::MailingSendRepository;
use crate::infrastructure::persistence::trace_repository::{SmsTraceVerdict, TraceRepository};
const TRACE_BATCH: i64 = 5000;
const CANDIDATE_BATCH: i64 = 500;
const BOUNCE_CLASS_CODES: &[&str] = &["sms_invalid_destination", "sms_not_allowed", "sms_rejected"];
#[derive(Debug, Clone, Default, PartialEq)]
pub struct SmsPumpOutcome {
pub traces_advanced: usize,
pub trace_skips: usize,
pub candidates: usize,
pub mailings_completed: usize,
}
enum VerdictAction {
Skip,
Process,
Pending,
Sent,
Bounce(&'static str),
Failed(&'static str),
Canceled(&'static str),
}
fn map_verdict(tracker_state: &str, failure_type: Option<&str>) -> VerdictAction {
let code_or = |fallback: &'static str| -> &'static str {
match failure_type {
Some(code) => match SMS_FAILURE_CODES.iter().copied().find(|m| *m == code) {
Some(member) => member,
None => fallback,
},
None => fallback,
}
};
match tracker_state {
"ready" => VerdictAction::Skip,
"process" => VerdictAction::Process,
"pending" => VerdictAction::Pending,
"sent" => VerdictAction::Sent,
"bounce" => VerdictAction::Bounce(code_or("sms_not_delivered")),
"exception" => match failure_type {
Some(code) if BOUNCE_CLASS_CODES.contains(&code) => {
VerdictAction::Bounce(code_or("sms_not_delivered"))
}
_ => VerdictAction::Failed(code_or("sms_server")),
},
"canceled" => VerdictAction::Canceled(code_or("sms_blacklist")),
other => {
tracing::warn!(
tracker_state = other,
"unknown delivery-tracker state; leaving the trace transient"
);
VerdictAction::Skip
}
}
}
pub struct SmsDeliveryPumpService {
pool: PgPool,
}
impl SmsDeliveryPumpService {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn pump_once(&self) -> Result<SmsPumpOutcome, sqlx::Error> {
let mut out = SmsPumpOutcome::default();
let verdicts = {
let mut conn = self.pool.acquire().await?;
TraceRepository::sms_traces_with_tracker_verdicts(&mut conn).await?
};
for verdict in &verdicts {
if self.apply_verdict(verdict).await? {
out.traces_advanced += 1;
} else {
out.trace_skips += 1;
}
}
let candidates = {
let mut conn = self.pool.acquire().await?;
MailingSendRepository::sms_done_inference_candidates(&mut conn, CANDIDATE_BATCH).await?
};
out.candidates = candidates.len();
for mailing_id in candidates {
let mut tx = self.pool.begin().await?;
let Some(locked) =
MailingSendRepository::lock_mailing_for_done(&mut tx, mailing_id).await?
else {
continue;
};
let remaining = TraceRepository::count_transient_sms_traces(&mut tx, locked).await?;
let mut completed = false;
if remaining == 0 {
completed = MailingSendRepository::complete_mailing(&mut tx, locked).await?;
}
tx.commit().await?;
if completed {
tracing::info!(
mailing_id = %locked,
"sms mailing completed by the delivery-tracker pump"
);
out.mailings_completed += 1;
}
}
Ok(out)
}
async fn apply_verdict(&self, verdict: &SmsTraceVerdict) -> Result<bool, sqlx::Error> {
let mut tx = self.pool.begin().await?;
let moved = match map_verdict(&verdict.tracker_state, verdict.failure_type.as_deref()) {
VerdictAction::Skip => {
tx.commit().await?;
return Ok(false);
}
VerdictAction::Process => {
TraceRepository::set_process(&mut tx, verdict.trace_id).await?
}
VerdictAction::Pending => {
TraceRepository::set_pending(&mut tx, verdict.trace_id).await?
}
VerdictAction::Sent => TraceRepository::set_sent(&mut tx, verdict.trace_id).await?,
VerdictAction::Bounce(code) => {
TraceRepository::set_bounced_sms(
&mut tx,
verdict.trace_id,
code,
verdict.failure_reason.as_deref(),
)
.await?
}
VerdictAction::Failed(code) => {
TraceRepository::set_failed(
&mut tx,
verdict.trace_id,
code,
verdict.failure_reason.as_deref(),
)
.await?
}
VerdictAction::Canceled(code) => {
TraceRepository::set_canceled(&mut tx, verdict.trace_id, code).await?
}
};
tx.commit().await?;
Ok(moved)
}
}