use chrono::{DateTime, Utc};
use sqlx::PgConnection;
use uuid::Uuid;
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct AudienceRecipient {
pub contact_id: Uuid,
pub email: String,
pub name: Option<String>,
}
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct SubscriptionRow {
pub id: Uuid,
pub contact_id: Uuid,
pub mailing_audience_id: Uuid,
pub opt_out: bool,
pub opt_out_datetime: Option<DateTime<Utc>>,
pub opt_out_reason_id: Option<Uuid>,
}
pub struct SubscriptionRepository;
impl SubscriptionRepository {
pub async fn subscribe(
conn: &mut PgConnection,
contact_id: Uuid,
audience_id: Uuid,
) -> Result<SubscriptionRow, sqlx::Error> {
sqlx::query_as::<_, SubscriptionRow>(
r#"INSERT INTO mailing.mailing_subscriptions (contact_id, mailing_audience_id, opt_out)
VALUES ($1, $2, FALSE)
ON CONFLICT (contact_id, mailing_audience_id)
WHERE (metadata->>'deleted_at') IS NULL
DO UPDATE SET opt_out = FALSE,
opt_out_datetime = NULL,
opt_out_reason_id = NULL
RETURNING id, contact_id, mailing_audience_id, opt_out,
opt_out_datetime, opt_out_reason_id"#,
)
.bind(contact_id)
.bind(audience_id)
.fetch_one(&mut *conn)
.await
}
pub async fn opt_out(
conn: &mut PgConnection,
contact_id: Uuid,
audience_id: Uuid,
reason_id: Option<Uuid>,
) -> Result<Option<SubscriptionRow>, sqlx::Error> {
sqlx::query_as::<_, SubscriptionRow>(
r#"UPDATE mailing.mailing_subscriptions
SET opt_out = TRUE,
opt_out_datetime = now(),
opt_out_reason_id = $3
WHERE contact_id = $1 AND mailing_audience_id = $2
AND NOT opt_out
AND (metadata->>'deleted_at') IS NULL
RETURNING id, contact_id, mailing_audience_id, opt_out,
opt_out_datetime, opt_out_reason_id"#,
)
.bind(contact_id)
.bind(audience_id)
.bind(reason_id)
.fetch_optional(&mut *conn)
.await
}
pub async fn opted_out_anywhere(
conn: &mut PgConnection,
email: &str,
) -> Result<bool, sqlx::Error> {
sqlx::query_scalar::<_, bool>(
r#"SELECT EXISTS (
SELECT 1
FROM mailing.mailing_subscriptions s
JOIN mailing.mailing_contacts c ON c.id = s.contact_id
WHERE lower(c.email) = lower($1)
AND s.opt_out
AND (s.metadata->>'deleted_at') IS NULL
AND (c.metadata->>'deleted_at') IS NULL
)"#,
)
.bind(email)
.fetch_one(&mut *conn)
.await
}
pub async fn find_live(
conn: &mut PgConnection,
contact_id: Uuid,
audience_id: Uuid,
) -> Result<Option<SubscriptionRow>, sqlx::Error> {
sqlx::query_as::<_, SubscriptionRow>(
r#"SELECT id, contact_id, mailing_audience_id, opt_out,
opt_out_datetime, opt_out_reason_id
FROM mailing.mailing_subscriptions
WHERE contact_id = $1 AND mailing_audience_id = $2
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(contact_id)
.bind(audience_id)
.fetch_optional(&mut *conn)
.await
}
pub async fn audience_recipients(
conn: &mut PgConnection,
audience_id: Uuid,
limit: i64,
after_email: Option<&str>,
) -> Result<Vec<AudienceRecipient>, sqlx::Error> {
sqlx::query_as::<_, AudienceRecipient>(
r#"SELECT c.id AS contact_id, c.email, c.name
FROM mailing.mailing_subscriptions s
JOIN mailing.mailing_contacts c ON c.id = s.contact_id
WHERE s.mailing_audience_id = $1
AND NOT s.opt_out
AND (s.metadata->>'deleted_at') IS NULL
AND (c.metadata->>'deleted_at') IS NULL
AND ($2::text IS NULL OR c.email > $2)
ORDER BY c.email
LIMIT $3"#,
)
.bind(audience_id)
.bind(after_email)
.bind(limit)
.fetch_all(&mut *conn)
.await
}
pub async fn audience_in_use_by_mailings(
conn: &mut PgConnection,
audience_id: Uuid,
) -> Result<Vec<Uuid>, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"SELECT m.id
FROM mailing.mailings m,
jsonb_array_elements(m.mailing_domain) term,
jsonb_each(term) kv
WHERE m.state IN ('draft', 'in_queue', 'sending')
AND (m.metadata->>'deleted_at') IS NULL
AND kv.value #>> '{}' = $1::text
GROUP BY m.id"#,
)
.bind(audience_id.to_string())
.fetch_all(&mut *conn)
.await
}
pub async fn remove(
conn: &mut PgConnection,
contact_id: Uuid,
audience_id: Uuid,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_subscriptions
SET metadata = metadata
|| jsonb_build_object('deleted_at', to_jsonb(now()))
WHERE contact_id = $1 AND mailing_audience_id = $2
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(contact_id)
.bind(audience_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn ensure_default_reasons(conn: &mut PgConnection) -> Result<(), sqlx::Error> {
for (idx, (name, is_feedback)) in [
("Not Interested", false),
("Unsubscribed", false),
("Too many emails", true),
("Never subscribed", true),
("Other", true),
]
.into_iter()
.enumerate()
{
sqlx::query(
r#"INSERT INTO mailing.mailing_opt_out_reasons (name, sequence, is_feedback)
SELECT $1, $2, $3
WHERE NOT EXISTS (
SELECT 1 FROM mailing.mailing_opt_out_reasons
WHERE name = $1 AND (metadata->>'deleted_at') IS NULL
)"#,
)
.bind(name)
.bind((idx as i32 + 1) * 10)
.bind(is_feedback)
.execute(&mut *conn)
.await?;
}
Ok(())
}
pub async fn default_reason_id(conn: &mut PgConnection) -> Result<Option<Uuid>, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"SELECT id FROM mailing.mailing_opt_out_reasons
WHERE name = 'Unsubscribed' AND (metadata->>'deleted_at') IS NULL
ORDER BY sequence LIMIT 1"#,
)
.fetch_optional(&mut *conn)
.await
}
pub async fn soft_delete_audience(
conn: &mut PgConnection,
audience_id: Uuid,
) -> Result<bool, sqlx::Error> {
sqlx::query(
r#"UPDATE mailing.mailing_audiences
SET metadata = metadata
|| jsonb_build_object('deleted_at', to_jsonb(now()))
WHERE id = $1 AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(audience_id)
.execute(&mut *conn)
.await
.map(|r| r.rows_affected() > 0)
}
pub async fn active_member_count(
conn: &mut PgConnection,
audience_id: Uuid,
) -> Result<i64, sqlx::Error> {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM mailing.mailing_subscriptions s
JOIN mailing.mailing_contacts c ON c.id = s.contact_id
WHERE s.mailing_audience_id = $1
AND NOT s.opt_out
AND (s.metadata->>'deleted_at') IS NULL
AND (c.metadata->>'deleted_at') IS NULL"#,
)
.bind(audience_id)
.fetch_one(&mut *conn)
.await
}
pub async fn contact_id_by_email(
conn: &mut PgConnection,
email: &str,
) -> Result<Option<Uuid>, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"SELECT id FROM mailing.mailing_contacts
WHERE lower(email) = lower($1)
AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(email)
.fetch_optional(&mut *conn)
.await
}
pub async fn audience_is_public(
conn: &mut PgConnection,
audience_id: Uuid,
) -> Result<Option<bool>, sqlx::Error> {
sqlx::query_scalar::<_, bool>(
r#"SELECT is_public FROM mailing.mailing_audiences
WHERE id = $1 AND (metadata->>'deleted_at') IS NULL"#,
)
.bind(audience_id)
.fetch_optional(&mut *conn)
.await
}
pub async fn upsert_contact(
conn: &mut PgConnection,
email: &str,
name: Option<&str>,
) -> Result<Uuid, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_contacts (email, name)
VALUES ($1, $2)
ON CONFLICT (email) WHERE (metadata->>'deleted_at') IS NULL
DO UPDATE SET name = COALESCE(EXCLUDED.name, mailing.mailing_contacts.name)
RETURNING id"#,
)
.bind(email)
.bind(name)
.fetch_one(&mut *conn)
.await
}
pub async fn insert_audience(
conn: &mut PgConnection,
id: Uuid,
name: &str,
is_public: bool,
) -> Result<Uuid, sqlx::Error> {
sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO mailing.mailing_audiences (id, name, is_public)
VALUES ($1, $2, $3) RETURNING id"#,
)
.bind(id)
.bind(name)
.bind(is_public)
.fetch_one(&mut *conn)
.await
}
}