| 1 | //! `df-db` — connection pooling, migrations, and shared row types. |
| 2 | //! |
| 3 | //! Every other crate talks to Postgres through this one so that pool |
| 4 | //! configuration, migration locking, and the enum mappings live in one place. |
| 5 | |
| 6 | use std::time::Duration; |
| 7 | |
| 8 | use sqlx::postgres::{PgConnectOptions, PgPoolOptions}; |
| 9 | use sqlx::{ConnectOptions, Connection}; |
| 10 | |
| 11 | pub mod ids; |
| 12 | pub mod models; |
| 13 | |
| 14 | pub use sqlx; |
| 15 | pub use sqlx::PgPool; |
| 16 | pub use sqlx::PgPool as Pool; |
| 17 | |
| 18 | /// Advisory lock key for schema migrations. |
| 19 | /// |
| 20 | /// `web` runs migrations at startup (spec §10) and more than one instance may |
| 21 | /// start at once, so the lock serialises them. The constant is arbitrary but |
| 22 | /// must never change — a different value would let two versions migrate |
| 23 | /// concurrently. |
| 24 | const MIGRATION_LOCK_KEY: i64 = 0x0D06_F00D_0000_0001u64 as i64; |
| 25 | |
| 26 | /// How long a liveness `ping()` may take before we treat the connection as dead. |
| 27 | /// |
| 28 | /// Comfortably longer than a healthy round trip to the pooler and well inside |
| 29 | /// `acquire_timeout`, so discarding a dead connection still leaves most of the |
| 30 | /// acquire budget to open a replacement. |
| 31 | const PING_TIMEOUT: Duration = Duration::from_secs(2); |
| 32 | |
| 33 | /// Open the pool. |
| 34 | /// |
| 35 | /// Statement logging is disabled below WARN: at INFO sqlx logs every statement, |
| 36 | /// which on the push path means logging the full contents of an indexing |
| 37 | /// transaction on every push. |
| 38 | pub async fn connect(url: &str, max_connections: u32) -> anyhow::Result<PgPool> { |
| 39 | let opts: PgConnectOptions = url |
| 40 | .parse::<PgConnectOptions>()? |
| 41 | .log_statements(tracing::log::LevelFilter::Debug) |
| 42 | .log_slow_statements(tracing::log::LevelFilter::Warn, Duration::from_millis(500)) |
| 43 | .application_name("dogfood"); |
| 44 | |
| 45 | let pool = PgPoolOptions::new() |
| 46 | .max_connections(max_connections) |
| 47 | .min_connections(1) |
| 48 | .acquire_timeout(Duration::from_secs(10)) |
| 49 | // No `idle_timeout` or `max_lifetime`, deliberately. |
| 50 | // |
| 51 | // Both are implemented by closing the connection from |
| 52 | // `PoolConnection::return_to_pool`, which awaits `PgConnection::close()` |
| 53 | // with no timeout of its own — unlike the `close_on_drop` path, which |
| 54 | // bounds it. `close()` sends Terminate and then awaits the stream |
| 55 | // shutdown. Against the pooler on :6543 the Terminate lands but the |
| 56 | // shutdown never completes, so the spawned task hangs forever still |
| 57 | // holding its pool slot, and the socket is never closed. |
| 58 | // |
| 59 | // With `max_lifetime` set to 30 minutes that leaked one slot every 30 |
| 60 | // minutes, so every process died a fixed `max_connections * 30min` |
| 61 | // after boot: on 2026-08-03 the worker booted at 01:01 and started |
| 62 | // failing every acquire at 06:02. sqlx 0.8 has no keepalive option and |
| 63 | // `test_before_acquire` cannot rescue it either, because `ping()` is |
| 64 | // likewise unbounded and just burns the whole `acquire_timeout`. |
| 65 | // |
| 66 | // Recycling was never load-bearing: the instance allows 100 |
| 67 | // connections and all three services together hold ~23. |
| 68 | // |
| 69 | // The built-in `test_before_acquire` is off for the same reason: its |
| 70 | // `ping()` is unbounded, so a connection the pooler dropped silently |
| 71 | // hangs the acquire instead of failing it. The hooks below do the same |
| 72 | // liveness check under `PING_TIMEOUT`. |
| 73 | // |
| 74 | // They must report a dead connection as `Err`, never `Ok(false)`: |
| 75 | // sqlx answers `Err` with `close_hard()`, which only awaits |
| 76 | // `stream.shutdown()`, but answers `Ok(false)` with the same unbounded |
| 77 | // `close()` that stranded slots above. |
| 78 | .test_before_acquire(false) |
| 79 | .before_acquire(|conn, _meta| Box::pin(async move { ping_bounded(conn).await.map(|()| true) })) |
| 80 | // `return_to_pool` pings unconditionally on release, and that ping is |
| 81 | // unbounded too — a connection that died while checked out would hang |
| 82 | // the spawned release task and strand the slot for good. Pinging here |
| 83 | // first means the unbounded one runs only against a connection that |
| 84 | // answered milliseconds ago. |
| 85 | .after_release(|conn, _meta| Box::pin(async move { ping_bounded(conn).await.map(|()| true) })) |
| 86 | .connect_with(opts) |
| 87 | .await?; |
| 88 | |
| 89 | Ok(pool) |
| 90 | } |
| 91 | |
| 92 | /// `ping()` under a timeout, reporting a timeout as an error. |
| 93 | /// |
| 94 | /// Returning `Err` matters: it is what routes the connection to sqlx's bounded |
| 95 | /// `close_hard()` rather than the `close()` that hangs against the pooler. |
| 96 | async fn ping_bounded(conn: &mut sqlx::PgConnection) -> Result<(), sqlx::Error> { |
| 97 | match tokio::time::timeout(PING_TIMEOUT, conn.ping()).await { |
| 98 | Ok(result) => result, |
| 99 | Err(_elapsed) => Err(sqlx::Error::Io(std::io::Error::new( |
| 100 | std::io::ErrorKind::TimedOut, |
| 101 | "connection ping timed out", |
| 102 | ))), |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | /// Apply pending migrations under an advisory lock. |
| 107 | /// |
| 108 | /// sqlx has its own locking, but it is per-migration; taking one session-level |
| 109 | /// lock around the whole run means a second starting instance waits for the |
| 110 | /// entire migration set rather than interleaving with it. |
| 111 | pub async fn migrate(pool: &PgPool) -> anyhow::Result<()> { |
| 112 | // A dedicated connection: the lock is session-scoped, so it must be held on |
| 113 | // one connection for the duration and released explicitly. |
| 114 | let mut conn = pool.acquire().await?; |
| 115 | |
| 116 | tracing::info!("acquiring migration advisory lock"); |
| 117 | sqlx::query("SELECT pg_advisory_lock($1)") |
| 118 | .bind(MIGRATION_LOCK_KEY) |
| 119 | .execute(&mut *conn) |
| 120 | .await?; |
| 121 | |
| 122 | let result = sqlx::migrate!("../../migrations").run(&mut *conn).await; |
| 123 | |
| 124 | // Release the lock whether or not migration succeeded, so a failed deploy |
| 125 | // does not wedge every subsequent start. |
| 126 | let unlock = sqlx::query("SELECT pg_advisory_unlock($1)") |
| 127 | .bind(MIGRATION_LOCK_KEY) |
| 128 | .execute(&mut *conn) |
| 129 | .await; |
| 130 | |
| 131 | result?; |
| 132 | unlock?; |
| 133 | |
| 134 | tracing::info!("migrations up to date"); |
| 135 | Ok(()) |
| 136 | } |
| 137 | |
| 138 | /// Liveness check for `/readyz`. |
| 139 | pub async fn ping(pool: &PgPool) -> anyhow::Result<()> { |
| 140 | sqlx::query_scalar::<_, i32>("SELECT 1") |
| 141 | .fetch_one(pool) |
| 142 | .await?; |
| 143 | Ok(()) |
| 144 | } |
144 lines · Rust