#![allow(clippy::unwrap_used, clippy::panic)] #![allow(clippy::print_stderr)]
use sqlx::PgPool;
use super::{MIGRATION_LOCK_KEY, lock_prefixed, run_migration};
const PROBE_DDL: &str = "\
CREATE TABLE IF NOT EXISTS _fraiseql_migration_lock_probe (
pk_probe BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
tenant_id TEXT,
label TEXT NOT NULL,
seen_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE UNIQUE INDEX IF NOT EXISTS uq_migration_lock_probe_label_per_space
ON _fraiseql_migration_lock_probe (label, tenant_id) NULLS NOT DISTINCT;
CREATE INDEX IF NOT EXISTS idx_migration_lock_probe_seen_at
ON _fraiseql_migration_lock_probe (seen_at);
ALTER TABLE _fraiseql_migration_lock_probe ENABLE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS p_migration_lock_probe_tenant ON _fraiseql_migration_lock_probe;
CREATE POLICY p_migration_lock_probe_tenant ON _fraiseql_migration_lock_probe
USING (tenant_id IS NOT DISTINCT FROM current_setting('fraiseql.tenant_id', true));
";
const RACERS: usize = 4;
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 race(pool: &PgPool) -> Vec<String> {
let mut set = tokio::task::JoinSet::new();
for _ in 0..RACERS {
let pool = pool.clone();
set.spawn(async move { run_migration(&pool, PROBE_DDL).await });
}
let mut errors = Vec::new();
while let Some(joined) = set.join_next().await {
if let Err(error) = joined.unwrap() {
errors.push(error.to_string());
}
}
errors
}
#[test]
fn the_lock_is_the_scripts_first_statement() {
let script = lock_prefixed("CREATE TABLE IF NOT EXISTS t (a INT);");
let first = script.split_once(';').expect("the composed script is statement-terminated").0;
assert_eq!(
first,
format!("SELECT pg_advisory_xact_lock({MIGRATION_LOCK_KEY})"),
"the advisory lock must be the script's first statement, carrying the shared key"
);
assert!(script.ends_with("CREATE TABLE IF NOT EXISTS t (a INT);"));
}
#[tokio::test]
async fn concurrent_migrations_race_on_neither_the_catalogue_nor_a_relation_lock() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!(
"SKIP concurrent_migrations_race_on_neither_the_catalogue_nor_a_relation_lock: no \
postgres (set DATABASE_URL or enable fraiseql-test-support/local-testcontainers)"
);
return;
};
sqlx::query("DROP TABLE IF EXISTS _fraiseql_migration_lock_probe CASCADE")
.execute(&pool)
.await
.unwrap();
let cold = race(&pool).await;
assert!(
cold.is_empty(),
"a cold start must not race on the catalogue; {} of {RACERS} runners failed: {cold:?}",
cold.len()
);
let warm = race(&pool).await;
assert!(
warm.is_empty(),
"a warm start must not deadlock on the relation lock; {} of {RACERS} runners failed: \
{warm:?}",
warm.len()
);
}