use std::sync::Arc;
use sqlx::PgPool;
use uuid::Uuid;
use backbone_mail::application::service::phone_blacklist_write_service::{
PhoneBlacklistError, PhoneBlacklistWriteService,
};
use backbone_mail::application::service::phone_validation_service::E164Number;
use crate::infrastructure::persistence::trace_repository::TraceRepository;
pub const SMS_BLACKLIST_FAILURE_TYPE: &str = "sms_blacklist";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PhoneRecipient {
pub recipient_id: Uuid,
pub email: String,
pub number: E164Number,
}
#[derive(Debug, Clone)]
pub struct ClaimedSmsContext {
pub mailing_id: Uuid,
pub campaign_id: Option<Uuid>,
pub recipient_model: String,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct SuppressionOutcome {
pub checked: usize,
pub suppressed: Vec<Suppressed>,
pub send_set: Vec<PhoneRecipient>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Suppressed {
pub recipient_id: Uuid,
pub email: String,
pub number: E164Number,
pub trace_id: Uuid,
}
#[derive(Debug, thiserror::Error)]
pub enum SmsSuppressionError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("blacklist seam: {0}")]
Blacklist(#[from] PhoneBlacklistError),
}
impl SmsSuppressionError {
pub fn code(&self) -> &'static str {
match self {
Self::Db(_) | Self::Blacklist(_) => "mailing_db_error",
}
}
pub fn http_status(&self) -> u16 {
500
}
}
pub struct SmsSuppressionService {
pool: PgPool,
blacklist: Arc<PhoneBlacklistWriteService>,
}
impl SmsSuppressionService {
pub fn new(pool: PgPool, blacklist: Arc<PhoneBlacklistWriteService>) -> Self {
Self { pool, blacklist }
}
pub async fn suppress_at_claim(
&self,
ctx: &ClaimedSmsContext,
recipients: Vec<PhoneRecipient>,
) -> Result<SuppressionOutcome, SmsSuppressionError> {
let numbers: Vec<E164Number> = recipients.iter().map(|r| r.number.clone()).collect();
let blacklisted = self.blacklist.listed_among(&numbers).await?;
let mut outcome = SuppressionOutcome { checked: recipients.len(), ..Default::default() };
let mut to_cancel: Vec<&PhoneRecipient> = Vec::new();
for recipient in &recipients {
if blacklisted.contains(&recipient.number) {
to_cancel.push(recipient);
} else {
outcome.send_set.push(recipient.clone());
}
}
if !to_cancel.is_empty() {
let mut tx = self.pool.begin().await?;
for recipient in &to_cancel {
let trace_id = Uuid::new_v4();
TraceRepository::mint_trace_channel(
&mut tx,
trace_id,
ctx.mailing_id,
ctx.campaign_id,
&ctx.recipient_model,
recipient.recipient_id,
&recipient.email,
Some(&recipient.number.to_string()),
"sms",
"cancel",
Some(SMS_BLACKLIST_FAILURE_TYPE),
false,
)
.await?;
outcome.suppressed.push(Suppressed {
recipient_id: recipient.recipient_id,
email: recipient.email.clone(),
number: recipient.number.clone(),
trace_id,
});
}
tx.commit().await?;
}
Ok(outcome)
}
}
#[cfg(test)]
mod tests {
}