use std::{future::Future, pin::Pin};
use fraiseql_error::Result;
use sqlx::PgPool;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum SuppressionReason {
HardBounce,
ChallengeUnanswered,
Unsubscribe,
}
impl SuppressionReason {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
SuppressionReason::HardBounce => "hard_bounce",
SuppressionReason::ChallengeUnanswered => "challenge_unanswered",
SuppressionReason::Unsubscribe => "unsubscribe",
}
}
#[must_use]
pub fn parse(token: &str) -> Option<Self> {
match token {
"hard_bounce" => Some(SuppressionReason::HardBounce),
"challenge_unanswered" => Some(SuppressionReason::ChallengeUnanswered),
"unsubscribe" => Some(SuppressionReason::Unsubscribe),
_ => None,
}
}
#[must_use]
pub fn default_ttl(
self,
now: chrono::DateTime<chrono::Utc>,
) -> Option<chrono::DateTime<chrono::Utc>> {
match self {
SuppressionReason::HardBounce | SuppressionReason::Unsubscribe => None,
SuppressionReason::ChallengeUnanswered => Some(now + chrono::Duration::days(30)),
}
}
}
impl std::fmt::Display for SuppressionReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RecordedSend {
pub message_id: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub struct SentRecord<'a> {
pub send_id: &'a str,
pub tenant: Option<&'a str>,
pub recipient: &'a str,
pub sending_address: &'a str,
pub message_id: Option<&'a str>,
}
pub trait SendTracker: Send + Sync {
fn suppression_reason<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<SuppressionReason>>> + Send + 'a>>;
fn recorded_send<'a>(
&'a self,
tenant: Option<&'a str>,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<RecordedSend>>> + Send + 'a>>;
fn record_sent<'a>(
&'a self,
record: SentRecord<'a>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CorrelatedSend {
pub send_id: String,
pub tenant: Option<String>,
pub recipient: String,
}
pub trait SendCorrelator: Send + Sync {
fn find_by_send_id<'a>(
&'a self,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<CorrelatedSend>>> + Send + 'a>>;
fn find_by_message_id<'a>(
&'a self,
message_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<CorrelatedSend>>> + Send + 'a>>;
fn mark_bounced<'a>(
&'a self,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
fn bump_challenge<'a>(
&'a self,
send_id: &'a str,
tenant: Option<&'a str>,
recipient: &'a str,
) -> Pin<Box<dyn Future<Output = Result<i64>> + Send + 'a>>;
fn mark_replied<'a>(
&'a self,
send_id: &'a str,
tenant: Option<&'a str>,
recipient: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
fn record_signal<'a>(
&'a self,
send_id: &'a str,
signal: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
fn suppress<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
reason: SuppressionReason,
ttl: Option<chrono::DateTime<chrono::Utc>>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
fn lift_suppression<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
reason: SuppressionReason,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
}
fn db_err(context: &str, error: &sqlx::Error) -> fraiseql_error::FraiseQLError {
fraiseql_error::FraiseQLError::database(format!("send tracking: {context}: {error}"))
}
pub struct PgSendTracker {
pool: PgPool,
}
impl PgSendTracker {
#[must_use]
pub const fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn init(&self) -> fraiseql_error::Result<()> {
sqlx::raw_sql(fraiseql_functions::migrations::send_tracking_migration_sql())
.execute(&self.pool)
.await
.map_err(|error| db_err("init", &error))?;
Ok(())
}
}
impl SendTracker for PgSendTracker {
fn suppression_reason<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<SuppressionReason>>> + Send + 'a>> {
Box::pin(async move {
let row: Option<(String,)> = sqlx::query_as(
"SELECT reason FROM _fraiseql_suppression \
WHERE address_hash = $1 AND tenant_id IS NOT DISTINCT FROM $2 \
AND (ttl IS NULL OR ttl > now()) \
LIMIT 1",
)
.bind(address_hash)
.bind(tenant)
.fetch_optional(&self.pool)
.await
.map_err(|error| db_err("suppression lookup", &error))?;
Ok(row.map(|(reason,)| {
SuppressionReason::parse(&reason).unwrap_or(SuppressionReason::Unsubscribe)
}))
})
}
fn recorded_send<'a>(
&'a self,
tenant: Option<&'a str>,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<RecordedSend>>> + Send + 'a>> {
Box::pin(async move {
let row: Option<(Option<String>,)> = sqlx::query_as(
"SELECT message_id FROM _fraiseql_send_status \
WHERE send_id = $1 AND tenant_id IS NOT DISTINCT FROM $2 \
LIMIT 1",
)
.bind(send_id)
.bind(tenant)
.fetch_optional(&self.pool)
.await
.map_err(|error| db_err("recorded-send lookup", &error))?;
Ok(row.map(|(message_id,)| RecordedSend { message_id }))
})
}
fn record_sent<'a>(
&'a self,
record: SentRecord<'a>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"INSERT INTO _fraiseql_send_status \
(send_id, tenant_id, recipient, sending_address, status, message_id, \
sent_at, updated_at) \
VALUES ($1, $2, $3, $4, 'Sent', $5, now(), now()) \
ON CONFLICT (COALESCE(tenant_id, ''), send_id) DO NOTHING",
)
.bind(record.send_id)
.bind(record.tenant)
.bind(record.recipient)
.bind(record.sending_address)
.bind(record.message_id)
.execute(&self.pool)
.await
.map_err(|error| db_err("record sent", &error))?;
Ok(())
})
}
}
type SendRow = (String, Option<String>, String);
fn to_correlated((send_id, tenant, recipient): SendRow) -> CorrelatedSend {
CorrelatedSend {
send_id,
tenant,
recipient,
}
}
impl SendCorrelator for PgSendTracker {
fn find_by_send_id<'a>(
&'a self,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<CorrelatedSend>>> + Send + 'a>> {
Box::pin(async move {
let row: Option<SendRow> = sqlx::query_as(
"SELECT send_id, tenant_id, recipient FROM _fraiseql_send_status \
WHERE send_id = $1 LIMIT 1",
)
.bind(send_id)
.fetch_optional(&self.pool)
.await
.map_err(|error| db_err("find by send-id", &error))?;
Ok(row.map(to_correlated))
})
}
fn find_by_message_id<'a>(
&'a self,
message_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Option<CorrelatedSend>>> + Send + 'a>> {
Box::pin(async move {
let row: Option<SendRow> = sqlx::query_as(
"SELECT send_id, tenant_id, recipient FROM _fraiseql_send_status \
WHERE message_id = $1 LIMIT 1",
)
.bind(message_id)
.fetch_optional(&self.pool)
.await
.map_err(|error| db_err("find by message-id", &error))?;
Ok(row.map(to_correlated))
})
}
fn mark_bounced<'a>(
&'a self,
send_id: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"UPDATE _fraiseql_send_status \
SET status = 'Bounced', last_signal = 'bounce', updated_at = now() \
WHERE send_id = $1",
)
.bind(send_id)
.execute(&self.pool)
.await
.map_err(|error| db_err("mark bounced", &error))?;
Ok(())
})
}
fn bump_challenge<'a>(
&'a self,
send_id: &'a str,
tenant: Option<&'a str>,
recipient: &'a str,
) -> Pin<Box<dyn Future<Output = Result<i64>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"UPDATE _fraiseql_send_status \
SET status = 'ChallengePending', challenge_count = challenge_count + 1, \
last_signal = 'challenge', updated_at = now() \
WHERE send_id = $1",
)
.bind(send_id)
.execute(&self.pool)
.await
.map_err(|error| db_err("bump challenge", &error))?;
let (count,): (i64,) = sqlx::query_as(
"SELECT count(*) FROM _fraiseql_send_status \
WHERE recipient = $1 AND tenant_id IS NOT DISTINCT FROM $2 \
AND status = 'ChallengePending'",
)
.bind(recipient)
.bind(tenant)
.fetch_one(&self.pool)
.await
.map_err(|error| db_err("count pending challenges", &error))?;
Ok(count)
})
}
fn mark_replied<'a>(
&'a self,
send_id: &'a str,
tenant: Option<&'a str>,
recipient: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"UPDATE _fraiseql_send_status \
SET status = 'Replied', challenge_count = 0, last_signal = 'reply', \
updated_at = now() \
WHERE tenant_id IS NOT DISTINCT FROM $2 \
AND (send_id = $1 OR (recipient = $3 AND status = 'ChallengePending'))",
)
.bind(send_id)
.bind(tenant)
.bind(recipient)
.execute(&self.pool)
.await
.map_err(|error| db_err("mark replied", &error))?;
Ok(())
})
}
fn record_signal<'a>(
&'a self,
send_id: &'a str,
signal: &'a str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"UPDATE _fraiseql_send_status SET last_signal = $2, updated_at = now() \
WHERE send_id = $1",
)
.bind(send_id)
.bind(signal)
.execute(&self.pool)
.await
.map_err(|error| db_err("record signal", &error))?;
Ok(())
})
}
fn suppress<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
reason: SuppressionReason,
ttl: Option<chrono::DateTime<chrono::Utc>>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"INSERT INTO _fraiseql_suppression (tenant_id, address_hash, reason, ttl) \
VALUES ($1, $2, $3, $4) \
ON CONFLICT (COALESCE(tenant_id, ''), address_hash) DO UPDATE \
SET reason = EXCLUDED.reason, ttl = EXCLUDED.ttl, \
since = now(), updated_at = now() \
WHERE _fraiseql_suppression.ttl IS NOT NULL",
)
.bind(tenant)
.bind(address_hash)
.bind(reason.as_str())
.bind(ttl)
.execute(&self.pool)
.await
.map_err(|error| db_err("suppress", &error))?;
Ok(())
})
}
fn lift_suppression<'a>(
&'a self,
tenant: Option<&'a str>,
address_hash: &'a str,
reason: SuppressionReason,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
Box::pin(async move {
sqlx::query(
"DELETE FROM _fraiseql_suppression \
WHERE address_hash = $1 AND tenant_id IS NOT DISTINCT FROM $2 AND reason = $3",
)
.bind(address_hash)
.bind(tenant)
.bind(reason.as_str())
.execute(&self.pool)
.await
.map_err(|error| db_err("lift suppression", &error))?;
Ok(())
})
}
}
#[cfg(test)]
mod tests;