use sqlx::postgres::PgPoolOptions;
use sqlx::PgPool;
use std::time::Duration;
pub const DEFAULT_DB_URL: &str = "postgres://root:password@localhost:5432/backbone_mail_test";
pub async fn test_pool() -> Option<PgPool> {
let url = std::env::var("DATABASE_URL").unwrap_or_else(|_| DEFAULT_DB_URL.into());
let pool = PgPoolOptions::new()
.max_connections(8)
.acquire_timeout(Duration::from_secs(5))
.connect(&url)
.await
.ok()?;
backbone_outbox::outbox::migrate(&pool, "messaging").await.ok()?;
Some(pool)
}
pub fn skipped(marker: &str) {
eprintln!("SKIPPED-DB: {marker} — no live Postgres reachable, not faking results");
}
pub static DRAIN_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
pub async fn sweep_queues(pool: &sqlx::PgPool) {
sqlx::query(
r#"DELETE FROM messaging.sms_trackers
WHERE sms_uuid IN (SELECT uuid FROM messaging.sms WHERE state = 'outgoing'::sms_state)"#,
)
.execute(pool)
.await
.ok();
sqlx::query("DELETE FROM messaging.sms WHERE state = 'outgoing'::sms_state")
.execute(pool)
.await
.ok();
sqlx::query("DELETE FROM messaging.mails WHERE state = 'outgoing'::mail_state")
.execute(pool)
.await
.ok();
sqlx::query(
r#"DELETE FROM messaging.outbox_events
WHERE event_type IN ('SmsCreated', 'MailQueued', 'MailDispatchRequested')
AND aggregate_id NOT IN (
SELECT id::text FROM messaging.sms
UNION SELECT id::text FROM messaging.mails)"#,
)
.execute(pool)
.await
.ok();
}
pub async fn cleanup(pool: &PgPool, tables: &[(&str, &[uuid::Uuid])]) {
for (table, ids) in tables {
if ids.is_empty() {
continue;
}
let _ = sqlx::query(&format!("DELETE FROM messaging.{table} WHERE id = ANY($1)"))
.bind(ids)
.execute(pool)
.await;
}
}
pub async fn seed_subtype(pool: &PgPool, name: &str) -> uuid::Uuid {
let id = uuid::Uuid::new_v4();
sqlx::query(
r#"INSERT INTO messaging.mail_message_subtypes (id, name)
VALUES ($1, $2)"#,
)
.bind(id)
.bind(name)
.execute(pool)
.await
.expect("seed subtype");
id
}