use uuid::Uuid;
use crate::application::service::mailing_write_service::parse_domain;
use crate::application::service::subscription_write_service::SubscriptionWriteService;
use crate::application::service::trace_click_ports::{
CodeAttribution, TraceClickError, TraceClickPort, TraceClickSlot,
};
use crate::application::service::trace_write_service::{TraceWriteError, TraceWriteService};
use crate::infrastructure::persistence::mailing_send_repository::MailingSendRepository;
use crate::infrastructure::persistence::trace_repository::TraceRepository;
#[derive(Debug, thiserror::Error)]
pub enum TraceRouteError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("unknown trace")]
UnknownTrace,
#[error("inconsistent (code, trace) pair")]
Inconsistent,
#[error("trace click port not composed")]
NotComposed(String),
#[error("backend: {0}")]
Backend(String),
#[error("invalid stored data: {0}")]
Invalid(String),
}
fn click_err(e: TraceClickError) -> TraceRouteError {
match e {
TraceClickError::UnknownCode => TraceRouteError::Inconsistent,
TraceClickError::AttributionMismatch => TraceRouteError::Inconsistent,
TraceClickError::NotComposed { detail } => TraceRouteError::NotComposed(detail),
TraceClickError::Backend(detail) => TraceRouteError::Backend(detail),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AudienceOptOut {
pub audience_id: Uuid,
pub changed: bool,
pub opted_out: bool,
}
pub struct TraceRouteService {
pool: sqlx::PgPool,
clicks: TraceClickSlot,
trace_write: TraceWriteService,
subscription_write: SubscriptionWriteService,
}
impl TraceRouteService {
pub fn new(pool: sqlx::PgPool, clicks: TraceClickSlot) -> Self {
Self {
trace_write: TraceWriteService::new(pool.clone()),
subscription_write: SubscriptionWriteService::new(pool.clone()),
pool,
clicks,
}
}
pub async fn record_click(
&self,
code: &str,
trace_id: Uuid,
ip: Option<&str>,
country_code: Option<&str>,
) -> Result<String, TraceRouteError> {
let row = self.route_trace(trace_id).await?;
let resolved = self
.clicks
.resolve_click(code, row.campaign_id, ip, country_code)
.await
.map_err(click_err)?;
self.trace_write.set_opened(trace_id).await?;
self.trace_write.set_clicked(trace_id).await?;
Ok(resolved.url)
}
pub async fn record_open(&self, code: &str, trace_id: Uuid) -> Result<(), TraceRouteError> {
let row = self.route_trace(trace_id).await?;
self.check_consistency(code, row.campaign_id).await?;
self.trace_write.set_opened(trace_id).await?;
Ok(())
}
pub async fn unsubscribe(
&self,
code: &str,
trace_id: Uuid,
reason_id: Option<Uuid>,
) -> Result<Vec<AudienceOptOut>, TraceRouteError> {
let row = self.route_trace(trace_id).await?;
self.check_consistency(code, row.campaign_id).await?;
let mut tx = self.pool.begin().await?;
let domain_json = MailingSendRepository::find_live_mailing_domain(&mut tx, row.mailing_id)
.await?
.ok_or_else(|| TraceRouteError::Invalid(format!(
"trace {} points at mailing {} which no longer resolves a live row",
row.id, row.mailing_id
)))?;
tx.commit().await?;
let audiences = audience_ids_of(&domain_json)?;
let mut out = Vec::with_capacity(audiences.len());
for audience_id in audiences {
let (changed, opted_out) = match self
.subscription_write
.unsubscribe_by_email(&row.recipient_email, audience_id, reason_id)
.await
{
Ok((row, changed)) => (changed, row.opt_out),
Err(
crate::application::service::subscription_write_service::SubscriptionWriteError::NotFound(_),
) => (false, false),
Err(e) => return Err(e.into()),
};
out.push(AudienceOptOut { audience_id, changed, opted_out });
}
Ok(out)
}
async fn route_trace(
&self,
trace_id: Uuid,
) -> Result<crate::infrastructure::persistence::trace_repository::RouteTraceRow, TraceRouteError> {
let mut tx = self.pool.begin().await?;
let row = TraceRepository::find_route_trace(&mut tx, trace_id).await?;
tx.commit().await?;
row.ok_or(TraceRouteError::UnknownTrace)
}
async fn check_consistency(
&self,
code: &str,
trace_campaign_id: Option<Uuid>,
) -> Result<CodeAttribution, TraceRouteError> {
let attribution = self.clicks.attribution(code).await.map_err(click_err)?;
if attribution.campaign_id != trace_campaign_id {
return Err(TraceRouteError::Inconsistent);
}
Ok(attribution)
}
}
impl From<crate::application::service::subscription_write_service::SubscriptionWriteError>
for TraceRouteError
{
fn from(
e: crate::application::service::subscription_write_service::SubscriptionWriteError,
) -> Self {
use crate::application::service::subscription_write_service::SubscriptionWriteError as E;
match e {
E::Db(e) => TraceRouteError::Db(e),
E::NotFound(d) => TraceRouteError::Backend(format!("subscription: {d}")),
E::Conflict(d) => TraceRouteError::Backend(format!("subscription: {d}")),
E::Invalid(d) => TraceRouteError::Invalid(format!("subscription: {d}")),
E::AudienceNotPublic => {
TraceRouteError::Backend("subscription: audience is not open".to_string())
}
}
}
}
impl From<TraceWriteError> for TraceRouteError {
fn from(e: TraceWriteError) -> Self {
match e {
TraceWriteError::Db(e) => TraceRouteError::Db(e),
TraceWriteError::Invalid(d) => TraceRouteError::Invalid(format!("trace verb: {d}")),
}
}
}
fn audience_ids_of(domain: &serde_json::Value) -> Result<Vec<Uuid>, TraceRouteError> {
use crate::infrastructure::persistence::mailing_send_repository::{
DomainField, DomainOp,
};
let compiled = parse_domain(domain)
.map_err(|e| TraceRouteError::Invalid(format!("stored mailing domain: {e}")))?;
let mut ids = Vec::new();
for term in &compiled.terms {
if term.field == DomainField::MailingAudienceId
&& matches!(term.op, DomainOp::Eq | DomainOp::In)
{
for raw in &term.values {
let id = Uuid::parse_str(raw).map_err(|_| {
TraceRouteError::Invalid(format!(
"stored mailing domain carries a non-uuid audience value: {raw}"
))
})?;
if !ids.contains(&id) {
ids.push(id);
}
}
}
}
Ok(ids)
}
#[cfg(test)]
mod tests {
}