use uuid::Uuid;
use crate::infrastructure::persistence::trace_repository::{TraceRepository, TraceRow};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TraceTransition {
Moved,
Skipped,
Missing,
}
#[derive(Debug, thiserror::Error)]
pub enum TraceWriteError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("invalid: {0}")]
Invalid(String),
}
impl TraceWriteError {
pub fn code(&self) -> &'static str {
match self {
Self::Db(_) => "mailing_db_error",
Self::Invalid(_) => "invalid_input",
}
}
pub fn http_status(&self) -> u16 {
match self {
Self::Db(_) => 500,
Self::Invalid(_) => 422,
}
}
}
pub struct TraceWriteService {
pool: sqlx::PgPool,
}
impl TraceWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
pub async fn find(&self, trace_id: Uuid) -> Result<Option<TraceRow>, TraceWriteError> {
let mut tx = self.pool.begin().await?;
let row = TraceRepository::find_live(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(row)
}
pub async fn set_process(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_process(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_pending(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_pending(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_bounced_sms(
&self,
trace_id: Uuid,
failure_type: &str,
failure_reason: Option<&str>,
) -> Result<TraceTransition, TraceWriteError> {
if !is_sms_failure_code(failure_type) {
return Err(TraceWriteError::Invalid(format!(
"not an SMS failure code: {failure_type}"
)));
}
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved =
TraceRepository::set_bounced_sms(&mut tx, trace_id, failure_type, failure_reason)
.await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn attach_sms_uuid(
&self,
trace_id: Uuid,
sms_uuid: &str,
) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
let moved = TraceRepository::attach_sms_uuid(&mut tx, trace_id, sms_uuid).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_sent(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_sent(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_opened(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_opened(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_replied(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_replied(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_bounced(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_bounced(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_clicked(&self, trace_id: Uuid) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
if TraceRepository::find_live(&mut tx, trace_id).await?.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_clicked(&mut tx, trace_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_failed(
&self,
trace_id: Uuid,
failure_type: &str,
failure_reason: Option<&str>,
) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
let row = TraceRepository::find_live(&mut tx, trace_id).await?;
if row.is_none() {
tx.commit().await?;
return Ok(TraceTransition::Missing);
}
let moved = TraceRepository::set_failed(&mut tx, trace_id, failure_type, failure_reason)
.await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_canceled(
&self,
trace_id: Uuid,
failure_type: &str,
) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
let moved = TraceRepository::set_canceled(&mut tx, trace_id, failure_type).await?;
tx.commit().await?;
Ok(classify(moved))
}
pub async fn set_message_id(
&self,
trace_id: Uuid,
message_id: &str,
) -> Result<TraceTransition, TraceWriteError> {
let mut tx = self.pool.begin().await?;
let moved = TraceRepository::set_message_id(&mut tx, trace_id, message_id).await?;
tx.commit().await?;
Ok(classify(moved))
}
}
fn classify(moved: bool) -> TraceTransition {
if moved {
TraceTransition::Moved
} else {
TraceTransition::Skipped
}
}
pub const SMS_FAILURE_CODES: &[&str] = &[
"sms_number_missing",
"sms_number_format",
"sms_country_not_supported",
"sms_registration_needed",
"sms_credit",
"sms_server",
"sms_acc",
"sms_blacklist",
"sms_duplicate",
"sms_optout",
"sms_expired",
"sms_invalid_destination",
"sms_not_allowed",
"sms_not_delivered",
"sms_rejected",
];
pub fn is_sms_failure_code(code: &str) -> bool {
SMS_FAILURE_CODES.contains(&code)
}