#![allow(clippy::unwrap_used)] #![allow(clippy::print_stderr)]
use sqlx::PgPool;
use super::{
cron_migration_sql, dlq_migration_sql, inbound_migration_sql, send_tracking_migration_sql,
};
async fn connect_pool() -> Option<(PgPool, fraiseql_test_support::Service)> {
let svc = fraiseql_test_support::postgres().await?;
let pool = PgPool::connect(svc.url()).await.unwrap();
Some((pool, svc))
}
async fn execute_ddl(pool: &PgPool, ddl: &str) {
for stmt in ddl.split(';') {
let trimmed = stmt.trim();
if !trimmed.is_empty() {
sqlx::query(trimmed).execute(pool).await.unwrap();
}
}
}
#[test]
fn test_cron_migration_ddl_is_valid_sql() {
let ddl = cron_migration_sql();
assert!(
ddl.contains("_fraiseql_cron_state"),
"DDL must create _fraiseql_cron_state table"
);
assert!(ddl.contains("IF NOT EXISTS"), "DDL must use IF NOT EXISTS");
for col in [
"pk_cron_state",
"function_name",
"cron_expr",
"last_fired_at",
"next_fire_at",
"fire_count",
"updated_at",
] {
assert!(ddl.contains(col), "DDL must contain column: {col}");
}
assert!(ddl.contains("idx_cron_state_function"), "DDL must create function_name index");
assert!(ddl.contains("idx_cron_state_next_fire"), "DDL must create next_fire_at index");
assert!(
ddl.contains("GENERATED ALWAYS AS IDENTITY"),
"pk must use GENERATED ALWAYS AS IDENTITY (Trinity pattern)"
);
assert!(
ddl.contains("UNIQUE (function_name, cron_expr)"),
"DDL must have unique constraint on (function_name, cron_expr)"
);
}
#[tokio::test]
async fn test_cron_migration_creates_table() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP test_cron_migration_creates_table: no postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
let ddl = cron_migration_sql();
execute_ddl(&pool, ddl).await;
let (exists,): (bool,) = sqlx::query_as(
"SELECT EXISTS (
SELECT 1 FROM pg_class WHERE relname = '_fraiseql_cron_state'
)",
)
.fetch_one(&pool)
.await
.unwrap();
assert!(exists, "table _fraiseql_cron_state must exist after migration");
}
#[test]
fn test_inbound_migration_ddl_is_valid_sql() {
let ddl = inbound_migration_sql();
assert!(
ddl.contains("_fraiseql_inbound_message"),
"DDL must create _fraiseql_inbound_message table"
);
assert!(ddl.contains("IF NOT EXISTS"), "DDL must use IF NOT EXISTS");
for col in [
"pk_inbound_message",
"source",
"idempotency_key",
"thread_key",
"payload",
"received_at",
"created_at",
] {
assert!(ddl.contains(col), "DDL must contain column: {col}");
}
assert!(
ddl.contains("UNIQUE (source, idempotency_key)"),
"DDL must dedup on (source, idempotency_key)"
);
assert!(
ddl.contains("GENERATED ALWAYS AS IDENTITY"),
"pk must use GENERATED ALWAYS AS IDENTITY (Trinity pattern)"
);
assert!(ddl.contains("idx_inbound_message_thread"), "DDL must create thread_key index");
assert!(
ddl.contains("idx_inbound_message_received"),
"DDL must create received_at index"
);
}
#[tokio::test]
async fn test_inbound_migration_creates_table() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP test_inbound_migration_creates_table: no postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
execute_ddl(&pool, inbound_migration_sql()).await;
let (exists,): (bool,) = sqlx::query_as(
"SELECT EXISTS (
SELECT 1 FROM pg_class WHERE relname = '_fraiseql_inbound_message'
)",
)
.fetch_one(&pool)
.await
.unwrap();
assert!(exists, "table _fraiseql_inbound_message must exist after migration");
}
#[test]
fn test_send_tracking_migration_ddl_is_valid_sql() {
let ddl = send_tracking_migration_sql();
for table in ["_fraiseql_send_status", "_fraiseql_suppression"] {
assert!(ddl.contains(table), "DDL must create {table}");
}
assert!(ddl.contains("IF NOT EXISTS"), "DDL must use IF NOT EXISTS");
for col in [
"send_id",
"tenant_id",
"recipient",
"sending_address",
"status",
"challenge_count",
"last_signal",
"address_hash",
"reason",
] {
assert!(ddl.contains(col), "DDL must contain column: {col}");
}
for key in [
"uq_send_status_per_space ON _fraiseql_send_status (send_id, tenant_id) NULLS NOT DISTINCT",
"uq_suppression_per_space ON _fraiseql_suppression (address_hash, tenant_id) NULLS NOT \
DISTINCT",
] {
assert!(
ddl.split_whitespace().collect::<Vec<_>>().join(" ").contains(key),
"missing {key}"
);
}
assert!(!ddl.contains("COALESCE"), "no expression key is left: {ddl}");
for old in ["uq_send_status_tenant_send", "uq_suppression_tenant_addr"] {
assert!(
ddl.contains(&format!("DROP INDEX IF EXISTS {old};")),
"{old} is dropped in place"
);
}
assert!(ddl.contains("ENABLE ROW LEVEL SECURITY"), "DDL must enable RLS");
assert!(
ddl.contains("current_setting('fraiseql.tenant_id', true)"),
"RLS policy must key on the fraiseql.tenant_id GUC"
);
assert!(ddl.contains("DROP POLICY IF EXISTS"), "policies must be idempotent");
}
#[tokio::test]
async fn test_send_tracking_migration_creates_tables_idempotently() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP test_send_tracking_migration_creates_tables_idempotently: no postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
let ddl = send_tracking_migration_sql();
execute_ddl(&pool, ddl).await;
execute_ddl(&pool, ddl).await;
for table in ["_fraiseql_send_status", "_fraiseql_suppression"] {
let (exists,): (bool,) =
sqlx::query_as("SELECT EXISTS (SELECT 1 FROM pg_class WHERE relname = $1)")
.bind(table)
.fetch_one(&pool)
.await
.unwrap();
assert!(exists, "table {table} must exist after migration");
}
}
#[test]
fn test_dlq_migration_ddl_is_valid_sql() {
let ddl = dlq_migration_sql();
assert!(
ddl.contains("_fraiseql_function_dlq"),
"DDL must create _fraiseql_function_dlq table"
);
assert!(ddl.contains("IF NOT EXISTS"), "DDL must use IF NOT EXISTS");
for col in [
"pk_function_dlq",
"id",
"source",
"function_name",
"trigger_type",
"idempotency_token",
"payload",
"error_message",
"attempts",
"created_at",
] {
assert!(ddl.contains(col), "DDL must contain column: {col}");
}
assert!(
ddl.contains("uq_function_dlq_id"),
"DDL must create a unique index on the dead-letter id"
);
assert!(ddl.contains("idx_function_dlq_created"), "DDL must create a created_at index");
assert!(
ddl.contains("GENERATED ALWAYS AS IDENTITY"),
"pk must use GENERATED ALWAYS AS IDENTITY (Trinity pattern)"
);
}
#[tokio::test]
async fn test_dlq_migration_creates_table_idempotently() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP test_dlq_migration_creates_table_idempotently: no postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
let ddl = dlq_migration_sql();
execute_ddl(&pool, ddl).await;
execute_ddl(&pool, ddl).await;
let (exists,): (bool,) =
sqlx::query_as("SELECT EXISTS (SELECT 1 FROM pg_class WHERE relname = $1)")
.bind("_fraiseql_function_dlq")
.fetch_one(&pool)
.await
.unwrap();
assert!(exists, "table _fraiseql_function_dlq must exist after migration");
}
#[tokio::test]
async fn test_cron_migration_is_idempotent() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP test_cron_migration_is_idempotent: no postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
let ddl = cron_migration_sql();
execute_ddl(&pool, ddl).await;
execute_ddl(&pool, ddl).await;
}