diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 55bff9759..37977a454 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -6562,6 +6562,65 @@ mod tests { .await; } + /// Migrations must outlive the runtime caps — an index build or an + /// `ACCESS EXCLUSIVE` wait routinely exceeds them, and startup treats a + /// migration failure as fatal — and the relaxed session must not survive + /// into the pool afterwards. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn migrations_ignore_runtime_timeouts_and_leak_no_relaxed_session() { + const TIGHT: &str = "50ms"; + + let admin = PgPool::connect(&admin_url().await) + .await + .expect("connect admin"); + let name = format!("migration_timeouts_{}", Uuid::new_v4().simple()); + sqlx::query(sqlx::AssertSqlSafe(format!("CREATE DATABASE {name}"))) + .execute(&admin) + .await + .expect("create scratch db"); + let base = admin_url().await; + let idx = base.rfind('/').expect("db url has a path segment"); + let scratch_url = format!("{}/{}", &base[..idx], name); + + // Far shorter than the migration suite needs, ample for a pooled query. + let db = Db::new(&DbConfig { + database_url: scratch_url, + max_connections: 2, + min_connections: 2, + statement_timeout: TIGHT.to_string(), + lock_timeout: TIGHT.to_string(), + ..DbConfig::default() + }) + .await + .expect("connect Db against the unmigrated scratch db"); + + db.migrate() + .await + .expect("migrations must not inherit the runtime caps"); + + // Hold every connection at once so a leaked relaxed session cannot hide + // behind a freshly dialed one. + let mut held = Vec::new(); + for _ in 0..2 { + let mut connection = db.pool.acquire().await.expect("acquire pooled connection"); + let statement_timeout: String = sqlx::query_scalar("SHOW statement_timeout") + .fetch_one(&mut *connection) + .await + .expect("SHOW statement_timeout"); + let lock_timeout: String = sqlx::query_scalar("SHOW lock_timeout") + .fetch_one(&mut *connection) + .await + .expect("SHOW lock_timeout"); + assert_eq!(statement_timeout, TIGHT); + assert_eq!(lock_timeout, TIGHT); + held.push(connection); + } + drop(held); + + drop_scratch_db(&admin, db.pool.clone(), &name).await; + } + /// Insert identical community + channel rows into a database so the same /// (community, channel) ids resolve in both writer and replica. async fn seed_community_channel( diff --git a/crates/buzz-db/src/migration.rs b/crates/buzz-db/src/migration.rs index 30eee9df3..09a993518 100644 --- a/crates/buzz-db/src/migration.rs +++ b/crates/buzz-db/src/migration.rs @@ -38,18 +38,31 @@ pub async fn run_migrations(pool: &PgPool) -> Result<()> { async fn run_migrator_without_runtime_timeouts(pool: &PgPool) -> Result<()> { let mut connection = pool.acquire().await?; + lift_runtime_timeouts(&mut connection).await?; + let migrated = MIGRATOR.run(&mut *connection).await; + // Retire the connection either way; report the migration outcome first so a + // close failure cannot mask it. + let retired = retire_connection(connection).await; + migrated?; + retired?; + Ok(()) +} + +/// Remove both runtime limits from one connection's session. +async fn lift_runtime_timeouts(connection: &mut sqlx::PgConnection) -> Result<()> { crate::apply_runtime_connection_timeouts( - &mut connection, + connection, crate::TIMEOUT_DISABLED, crate::TIMEOUT_DISABLED, ) .await?; - let migrated = MIGRATOR.run(&mut *connection).await; - // Retire the connection either way; report the migration outcome first so a - // close failure cannot mask it. - let closed = connection.detach().close().await; - migrated?; - closed?; + Ok(()) +} + +/// Close a connection instead of returning it to the pool, so a session that +/// carries lifted limits can never serve traffic. +async fn retire_connection(connection: sqlx::pool::PoolConnection) -> Result<()> { + connection.detach().close().await?; Ok(()) } @@ -1149,6 +1162,69 @@ mod tests { .expect("read applied migrations") } + /// The migration connection must have both limits lifted, and it must not + /// come back to the pool afterwards — a session with no statement timeout + /// serving traffic is the failure this exemption trades against. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn migration_connection_is_unbounded_and_is_retired_not_reused() { + const TIGHT: &str = "50ms"; + + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| TEST_DB_URL.to_owned()); + // One slot: a reused connection would be handed straight back below. + let pool = sqlx::postgres::PgPoolOptions::new() + .max_connections(1) + .after_connect(|connection, _meta| { + Box::pin(crate::apply_runtime_connection_timeouts( + connection, TIGHT, TIGHT, + )) + }) + .connect(&database_url) + .await + .expect("connect to test DB"); + + let mut connection = pool.acquire().await.expect("acquire"); + assert_eq!( + show_timeout(&mut connection, "statement_timeout").await, + TIGHT + ); + + lift_runtime_timeouts(&mut connection) + .await + .expect("lift runtime timeouts"); + for setting in ["statement_timeout", "lock_timeout"] { + assert_eq!( + show_timeout(&mut connection, setting).await, + "0", + "{setting} must be lifted for the migrator" + ); + } + + retire_connection(connection) + .await + .expect("retire migration connection"); + + let mut fresh = pool.acquire().await.expect("re-acquire"); + for setting in ["statement_timeout", "lock_timeout"] { + assert_eq!( + show_timeout(&mut fresh, setting).await, + TIGHT, + "the pool must not hand out the relaxed migration session" + ); + } + drop(fresh); + pool.close().await; + } + + async fn show_timeout(connection: &mut sqlx::PgConnection, setting: &str) -> String { + sqlx::query_scalar(sqlx::AssertSqlSafe(format!("SHOW {setting}"))) + .fetch_one(connection) + .await + .unwrap_or_else(|error| panic!("SHOW {setting}: {error}")) + } + #[tokio::test] #[ignore = "requires Postgres"] async fn pre_0007_ambiguous_nip_rs_data_blocks_without_mutation_and_allows_retry() {