use chrono::{DateTime, Utc};
use sqlx::PgConnection;
use uuid::Uuid;
#[derive(Debug, Default, Clone, Copy, PartialEq)]
pub struct MailingTraceCounts {
pub total: i64,
pub sent: i64,
pub delivered: i64,
pub opened: i64,
pub clicked: i64,
pub replied: i64,
pub bounced: i64,
pub errored: i64,
pub canceled: i64,
}
impl MailingTraceCounts {
pub fn opened_ratio(&self) -> Option<i64> {
ratio(self.opened, self.sent)
}
pub fn clicks_ratio(&self) -> Option<i64> {
ratio(self.clicked, self.sent)
}
pub fn replied_ratio(&self) -> Option<i64> {
ratio(self.replied, self.sent)
}
}
fn ratio(part: i64, whole: i64) -> Option<i64> {
(whole > 0).then(|| (part as f64 / whole as f64 * 100.0).round() as i64)
}
pub struct TraceRepository;
impl TraceRepository {
#[allow(clippy::too_many_arguments)]
pub async fn mint_trace(
conn: &mut PgConnection,
id: Uuid,
mailing_id: Uuid,
campaign_id: Option<Uuid>,
recipient_model: &str,
recipient_id: Uuid,
recipient_email: &str,
trace_status: &str,
failure_type: Option<&str>,
is_test_trace: bool,
) -> Result<Uuid, sqlx::Error> {
let out = sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_traces
(id, trace_type, is_test_trace, mailing_id, campaign_id,
recipient_model, recipient_id, recipient_email,
trace_status, failure_type, metadata)
VALUES ($1, 'mail', $2, $3, $4, $5, $6, $7, $8::trace_status,
$9::trace_failure_type,
jsonb_build_object('created_at', to_jsonb(now())))
RETURNING id"#,
)
.bind(id)
.bind(is_test_trace)
.bind(mailing_id)
.bind(campaign_id)
.bind(recipient_model)
.bind(recipient_id)
.bind(recipient_email)
.bind(trace_status)
.bind(failure_type)
.fetch_one(&mut *conn)
.await?;
Ok(out)
}
#[allow(clippy::too_many_arguments)]
pub async fn mint_trace_fenced(
conn: &mut PgConnection,
id: Uuid,
mailing_id: Uuid,
campaign_id: Option<Uuid>,
recipient_model: &str,
recipient_id: Uuid,
recipient_email: &str,
trace_status: &str,
failure_type: Option<&str>,
is_test_trace: bool,
) -> Result<bool, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_traces
(id, trace_type, is_test_trace, mailing_id, campaign_id,
recipient_model, recipient_id, recipient_email,
trace_status, failure_type, metadata)
VALUES ($1, 'mail', $2, $3, $4, $5, $6, $7, $8::trace_status,
$9::trace_failure_type,
jsonb_build_object('created_at', to_jsonb(now())))
ON CONFLICT DO NOTHING
RETURNING id"#,
)
.bind(id)
.bind(is_test_trace)
.bind(mailing_id)
.bind(campaign_id)
.bind(recipient_model)
.bind(recipient_id)
.bind(recipient_email)
.bind(trace_status)
.bind(failure_type)
.fetch_optional(&mut *conn)
.await
.map(|r| r.is_some())
}
#[allow(clippy::too_many_arguments)]
pub async fn mint_trace_channel(
conn: &mut PgConnection,
id: Uuid,
mailing_id: Uuid,
campaign_id: Option<Uuid>,
recipient_model: &str,
recipient_id: Uuid,
recipient_email: &str,
recipient_phone: Option<&str>,
trace_type: &str,
trace_status: &str,
failure_type: Option<&str>,
is_test_trace: bool,
) -> Result<Uuid, sqlx::Error> {
let out = sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_traces
(id, trace_type, is_test_trace, mailing_id, campaign_id,
recipient_model, recipient_id, recipient_email,
recipient_phone, trace_status, failure_type, metadata)
VALUES ($1, $10::trace_type, $2, $3, $4, $5, $6, $7, $8,
$9::trace_status, $11::trace_failure_type,
jsonb_build_object('created_at', to_jsonb(now())))
RETURNING id"#,
)
.bind(id)
.bind(is_test_trace)
.bind(mailing_id)
.bind(campaign_id)
.bind(recipient_model)
.bind(recipient_id)
.bind(recipient_email)
.bind(recipient_phone)
.bind(trace_status)
.bind(trace_type)
.bind(failure_type)
.fetch_one(&mut *conn)
.await?;
Ok(out)
}
#[allow(clippy::too_many_arguments)]
pub async fn mint_trace_channel_fenced(
conn: &mut PgConnection,
id: Uuid,
mailing_id: Uuid,
campaign_id: Option<Uuid>,
recipient_model: &str,
recipient_id: Uuid,
recipient_email: &str,
recipient_phone: Option<&str>,
trace_type: &str,
trace_status: &str,
failure_type: Option<&str>,
is_test_trace: bool,
) -> Result<bool, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_traces
(id, trace_type, is_test_trace, mailing_id, campaign_id,
recipient_model, recipient_id, recipient_email,
recipient_phone, trace_status, failure_type, metadata)
VALUES ($1, $10::trace_type, $2, $3, $4, $5, $6, $7, $8,
$9::trace_status, $11::trace_failure_type,
jsonb_build_object('created_at', to_jsonb(now())))
ON CONFLICT DO NOTHING
RETURNING id"#,
)
.bind(id)
.bind(is_test_trace)
.bind(mailing_id)
.bind(campaign_id)
.bind(recipient_model)
.bind(recipient_id)
.bind(recipient_email)
.bind(recipient_phone)
.bind(trace_status)
.bind(trace_type)
.bind(failure_type)
.fetch_optional(&mut *conn)
.await
.map(|r| r.is_some())
}
pub async fn set_process(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'process'
WHERE id = $1
AND trace_status = 'outgoing'
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_pending(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'pending'
WHERE id = $1
AND trace_status IN ('outgoing', 'process')
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_sent(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'sent', failure_type = NULL, failure_reason = NULL,
sent_datetime = COALESCE(sent_datetime, now())
WHERE id = $1
AND trace_status IN ('outgoing', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_opened(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'open', open_datetime = COALESCE(open_datetime, now())
WHERE id = $1
AND trace_status IN ('outgoing', 'sent', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_replied(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'reply', reply_datetime = now(),
open_datetime = COALESCE(open_datetime, now())
WHERE id = $1
AND trace_status IN ('outgoing', 'sent', 'open', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_bounced_sms(
conn: &mut PgConnection,
trace_id: Uuid,
failure_type: &str,
failure_reason: Option<&str>,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET trace_status = 'bounce', failure_type = $2::trace_failure_type,
failure_reason = $3
WHERE id = $1
AND trace_status IN ('outgoing', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(failure_type)
.bind(failure_reason)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn set_bounced(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
Self::advance(
conn,
trace_id,
r#"UPDATE mailing.mailing_traces
SET trace_status = 'bounce', failure_type = 'mail_bounce',
open_datetime = COALESCE(open_datetime, now())
WHERE id = $1
AND trace_status IN ('outgoing', 'sent', 'open', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.await
}
pub async fn set_failed(
conn: &mut PgConnection,
trace_id: Uuid,
failure_type: &str,
failure_reason: Option<&str>,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET trace_status = 'error', failure_type = $2::trace_failure_type,
failure_reason = $3
WHERE id = $1
AND trace_status IN ('outgoing', 'sent', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(failure_type)
.bind(failure_reason)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn set_canceled(
conn: &mut PgConnection,
trace_id: Uuid,
failure_type: &str,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET trace_status = 'cancel', failure_type = $2::trace_failure_type
WHERE id = $1
AND trace_status = 'outgoing'
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(failure_type)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn set_clicked(conn: &mut PgConnection, trace_id: Uuid) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET links_click_datetime = now()
WHERE id = $1
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn attach_mail_id(
conn: &mut PgConnection,
trace_id: Uuid,
mail_id: Uuid,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET mail_id = $2
WHERE id = $1 AND mail_id IS NULL
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(mail_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn attach_sms_uuid(
conn: &mut PgConnection,
trace_id: Uuid,
sms_uuid: &str,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET sms_uuid = $2
WHERE id = $1 AND sms_uuid IS NULL
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(sms_uuid)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn set_message_id(
conn: &mut PgConnection,
trace_id: Uuid,
message_id: &str,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_traces
SET message_id = $2
WHERE id = $1 AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.bind(message_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
async fn advance(
conn: &mut PgConnection,
trace_id: Uuid,
sql: &'static str,
) -> Result<bool, sqlx::Error> {
sqlx::query(sql)
.bind(trace_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn outgoing_without_mail(
conn: &mut PgConnection,
mailing_id: Uuid,
grace_minutes: i32,
) -> Result<Vec<(Uuid, Uuid, String)>, sqlx::Error> {
sqlx::query_as::<_, (Uuid, Uuid, String)>(
r#"SELECT id, recipient_id, recipient_email
FROM mailing.mailing_traces
WHERE mailing_id = $1 AND trace_status = 'outgoing'
AND trace_type = 'mail'
AND mail_id IS NULL
AND (metadata->>'deleted_at') IS NULL
AND (metadata->>'created_at')::timestamptz
<= now() - make_interval(mins => $2)
ORDER BY recipient_email"#,
)
.bind(mailing_id)
.bind(grace_minutes)
.fetch_all(&mut *conn)
.await
}
pub async fn settled_mails_for_reconcile(
conn: &mut PgConnection,
) -> Result<Vec<(Uuid, String, Option<String>)>, sqlx::Error> {
sqlx::query_as::<_, (Uuid, String, Option<String>)>(
r#"SELECT t.id, m.state::text, m.failure_type::text
FROM mailing.mailing_traces t
JOIN messaging.mails m ON m.id = t.mail_id
WHERE t.trace_status = 'outgoing'
AND m.state IN ('sent', 'exception')
AND (t.metadata->>'deleted_at') IS NULL
ORDER BY t.id
LIMIT 5000"#,
)
.fetch_all(&mut *conn)
.await
}
pub async fn sms_traces_with_tracker_verdicts(
conn: &mut PgConnection,
) -> Result<Vec<SmsTraceVerdict>, sqlx::Error> {
sqlx::query_as::<_, SmsTraceVerdict>(
r#"SELECT t.id AS trace_id,
tr.state::text AS tracker_state,
tr.failure_type::text AS failure_type,
tr.failure_reason
FROM mailing.mailing_traces t
JOIN messaging.sms_trackers tr ON tr.sms_uuid = t.sms_uuid
WHERE t.trace_type = 'sms'
AND t.trace_status IN ('outgoing', 'process', 'pending')
AND (t.metadata->>'deleted_at') IS NULL
ORDER BY t.id
LIMIT 5000"#,
)
.fetch_all(&mut *conn)
.await
}
pub async fn count_transient_sms_traces(
conn: &mut PgConnection,
mailing_id: Uuid,
) -> Result<i64, sqlx::Error> {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*)
FROM mailing.mailing_traces
WHERE mailing_id = $1
AND trace_type = 'sms'
AND trace_status IN ('outgoing', 'process', 'pending')
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(mailing_id)
.fetch_one(&mut *conn)
.await
}
pub async fn find_live(
conn: &mut PgConnection,
trace_id: Uuid,
) -> Result<Option<TraceRow>, sqlx::Error> {
sqlx::query_as::<_, TraceRow>(
r#"SELECT id, trace_status::text AS trace_status, failure_type::text AS failure_type,
sent_datetime, open_datetime, reply_datetime, links_click_datetime, mail_id
FROM mailing.mailing_traces
WHERE id = $1 AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.fetch_optional(&mut *conn)
.await
}
pub async fn find_route_trace(
conn: &mut PgConnection,
trace_id: Uuid,
) -> Result<Option<RouteTraceRow>, sqlx::Error> {
sqlx::query_as::<_, RouteTraceRow>(
r#"SELECT id, trace_status::text AS trace_status, mailing_id, campaign_id,
recipient_email
FROM mailing.mailing_traces
WHERE id = $1 AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(trace_id)
.fetch_optional(&mut *conn)
.await
}
pub async fn counts_for_mailings(
conn: &mut PgConnection,
mailing_ids: &[Uuid],
) -> Result<Vec<(Uuid, MailingTraceCounts)>, sqlx::Error> {
let rows = sqlx::query_as::<_, (Uuid, i64, i64, i64, i64, i64, i64, i64, i64, i64)>(
r#"SELECT mailing_id,
count(*) FILTER (WHERE trace_status <> 'cancel') AS total,
count(*) FILTER (WHERE sent_datetime IS NOT NULL) AS sent,
count(*) FILTER (WHERE trace_status IN ('sent', 'open', 'reply')) AS delivered,
count(*) FILTER (WHERE trace_status IN ('open', 'reply')) AS opened,
count(*) FILTER (WHERE links_click_datetime IS NOT NULL) AS clicked,
count(*) FILTER (WHERE trace_status = 'reply') AS replied,
count(*) FILTER (WHERE trace_status = 'bounce') AS bounced,
count(*) FILTER (WHERE trace_status = 'error') AS errored,
count(*) FILTER (WHERE trace_status = 'cancel') AS canceled
FROM mailing.mailing_traces
WHERE mailing_id = ANY($1)
AND (metadata->>'deleted_at') IS NULL
GROUP BY mailing_id"#,
)
.bind(mailing_ids)
.fetch_all(&mut *conn)
.await?;
Ok(rows
.into_iter()
.map(|(id, total, sent, delivered, opened, clicked, replied, bounced, errored, canceled)| {
(
id,
MailingTraceCounts {
total, sent, delivered, opened, clicked, replied, bounced, errored, canceled,
},
)
})
.collect())
}
#[allow(clippy::type_complexity)]
pub async fn source_grouped_counts(
conn: &mut PgConnection,
mailing_ids: &[Uuid],
) -> Result<Vec<(Uuid, i64, i64, i64, i64, i64, i64, i64)>, sqlx::Error> {
sqlx::query_as::<_, (Uuid, i64, i64, i64, i64, i64, i64, i64)>(
r#"WITH scope AS (
SELECT m.id, m.source_id
FROM mailing.mailings m
WHERE m.id = ANY($1)
AND m.source_id IS NOT NULL
AND (m.metadata->>'deleted_at') IS NULL
),
per_source AS (
SELECT s.source_id,
count(DISTINCT t.mailing_id) AS mailings,
count(*) FILTER (WHERE t.trace_status <> 'cancel') AS total,
count(*) FILTER (WHERE t.sent_datetime IS NOT NULL) AS sent,
count(*) FILTER (WHERE t.trace_status IN ('open', 'reply')) AS opened,
count(*) FILTER (WHERE t.links_click_datetime IS NOT NULL) AS clicked,
count(*) FILTER (WHERE t.trace_status = 'reply') AS replied
FROM scope s
JOIN mailing.mailing_traces t ON t.mailing_id = s.id
AND (t.metadata->>'deleted_at') IS NULL
GROUP BY s.source_id
),
source_campaigns AS (
SELECT m.source_id, count(DISTINCT m.campaign_id) AS campaigns
FROM mailing.mailings m
WHERE m.source_id IN (SELECT source_id FROM per_source)
AND m.campaign_id IS NOT NULL
AND (m.metadata->>'deleted_at') IS NULL
GROUP BY m.source_id
)
SELECT p.source_id, COALESCE(c.campaigns, 0), p.mailings, p.total,
p.sent, p.opened, p.clicked, p.replied
FROM per_source p
LEFT JOIN source_campaigns c ON c.source_id = p.source_id
ORDER BY p.source_id"#,
)
.bind(mailing_ids)
.fetch_all(&mut *conn)
.await
}
}
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct SmsTraceVerdict {
pub trace_id: Uuid,
pub tracker_state: String,
pub failure_type: Option<String>,
pub failure_reason: Option<String>,
}
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct TraceRow {
pub id: Uuid,
pub trace_status: String,
pub failure_type: Option<String>,
pub sent_datetime: Option<DateTime<Utc>>,
pub open_datetime: Option<DateTime<Utc>>,
pub reply_datetime: Option<DateTime<Utc>>,
pub links_click_datetime: Option<DateTime<Utc>>,
pub mail_id: Option<Uuid>,
}
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct RouteTraceRow {
pub id: Uuid,
pub trace_status: String,
pub mailing_id: Uuid,
pub campaign_id: Option<Uuid>,
pub recipient_email: String,
}