diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 77e8bcc..9d4a05e 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -77,6 +77,34 @@ jobs: - name: Run migration ledger test suite run: node --test tests/database-postgres-ledger.test.mjs + # E00-S03-T04: the static assertions of tests/database-postgres-lock.test.mjs + # gate every PR — the suite locks in the migration advisory lock (session- + # scoped pg_advisory_lock/pg_try_advisory_lock over a stable keyed hash on a + # dedicated connection, re-entrant-safe in-flight acquire so concurrent + # acquire() calls share one connection, driver-boundary re-export) with + # mutation probes, and the docker-gated real-stack concurrent probe (a + # second runner waits or fails while the first holds the lock; concurrent + # acquire() checks out exactly one connection; the lock releases when the + # holding session ends) runs where a Docker daemon is available and skips + # cleanly otherwise. The job installs the frozen workspace because the + # real-stack probe executes the committed lock module from the host (it + # imports `pg` through the package's own links). + database-postgres-lock: + name: Migration advisory lock (E00-S03-T04) + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - name: Install Node.js 24 + uses: actions/setup-node@v4 + with: + node-version: '24' + - name: Enable pnpm (corepack, pinned to 11.23.0 via packageManager) + run: corepack enable + - name: Install dependencies (frozen lockfile) + run: pnpm install --frozen-lockfile + - name: Run migration advisory lock test suite + run: node --test tests/database-postgres-lock.test.mjs + # E00-S03-T01: the static assertions of tests/compose-config.test.mjs (db # image pinned to postgres:18.6-bookworm, health gate, volume persistence, # build platforms) gate every PR (the docker-gated real-stack probes inside diff --git a/.gitignore b/.gitignore index c908d5d..5a3b2f5 100644 --- a/.gitignore +++ b/.gitignore @@ -13,9 +13,10 @@ coverage/ # Logs *.log -# Transient host-side probe file written by tests/database-postgres-ledger.test.mjs -# into the database-postgres package (removed in its finally block) +# Transient host-side probe files written by the database-postgres test +# suites into the package (removed in their finally blocks) .ledger-probe-*.mjs +.lock-probe-*.mjs # OS / editor .DS_Store diff --git a/docs/development/non-container.md b/docs/development/non-container.md index f86cec7..813a18b 100644 --- a/docs/development/non-container.md +++ b/docs/development/non-container.md @@ -17,7 +17,7 @@ The workspace is a pnpm monorepo with three package groups: | --- | --- | --- | | `apps/` | `apps/server` (`@personal-blog/server`) | Public server application. Serves the application health endpoint (E00-S02-T03); the Fastify 5 application shell lands in a later story. | | `packages/` | `packages/core` (`@personal-blog/core`) | Application core (site identity, content primitives). Bootstrap placeholder. | -| `packages/` | `packages/database-postgres` (`@personal-blog/database-postgres`) | PostgreSQL database adapter package. Single owner of the `pg`/Kysely driver imports (E00-S03-T02); the migration ledger (`schema_migrations`, E00-S03-T03) is implemented here; the migration runner (advisory lock, failure diagnostics) lands in later stories. | +| `packages/` | `packages/database-postgres` (`@personal-blog/database-postgres`) | PostgreSQL database adapter package. Single owner of the `pg`/Kysely driver imports (E00-S03-T02); the migration ledger (`schema_migrations`, E00-S03-T03) and the migration advisory lock (E00-S03-T04) are implemented here; the migration runner (failure diagnostics) lands in later stories. | | `extensions/` | `extensions/example` (`@personal-blog/example-extension`) | Example extension exercising the `extensions/` group. Bootstrap placeholder. | A dependency-boundary rule (`dependency-boundaries.json`, enforced by diff --git a/packages/database-postgres/package.json b/packages/database-postgres/package.json index 223a180..016de03 100644 --- a/packages/database-postgres/package.json +++ b/packages/database-postgres/package.json @@ -3,7 +3,7 @@ "version": "0.0.0", "private": true, "type": "module", - "description": "EPPP PostgreSQL database adapter package. The single workspace package allowed to import the pg driver and Kysely (E00-S03-T02); the migration ledger (E00-S03-T03) is implemented here; the migration runner (advisory lock, diagnostics) lands in later stories.", + "description": "EPPP PostgreSQL database adapter package. The single workspace package allowed to import the pg driver and Kysely (E00-S03-T02); the migration ledger (E00-S03-T03) and the migration advisory lock (E00-S03-T04) are implemented here; the migration runner (failure diagnostics) lands in later stories.", "scripts": { "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit" diff --git a/packages/database-postgres/src/index.ts b/packages/database-postgres/src/index.ts index c25db95..df7ec9e 100644 --- a/packages/database-postgres/src/index.ts +++ b/packages/database-postgres/src/index.ts @@ -9,10 +9,10 @@ * * This module is the driver boundary: it imports the PostgreSQL driver (`pg`) * and Kysely and re-exports the pieces the adapter is built on — the driver - * surface (E00-S03-T02) and the migration ledger (E00-S03-T03). The advisory - * lock (E00-S03-T04) and failure diagnostic (E00-S03-T05) land in later - * stories; until then the re-exports keep the driver reachable only from here - * — the isolation is real, not a placeholder. + * surface (E00-S03-T02), the migration ledger (E00-S03-T03) and the migration + * advisory lock (E00-S03-T04). The failure diagnostic (E00-S03-T05) lands in + * a later story; until then the re-exports keep the driver reachable only + * from here — the isolation is real, not a placeholder. */ import { Pool } from 'pg'; @@ -20,3 +20,4 @@ import { Kysely, PostgresDialect } from 'kysely'; export { Pool, Kysely, PostgresDialect }; export { MigrationLedger, MIGRATION_LEDGER_TABLE } from './ledger.js'; +export { MigrationLock, MIGRATION_LOCK_KEY } from './lock.js'; diff --git a/packages/database-postgres/src/lock.ts b/packages/database-postgres/src/lock.ts new file mode 100644 index 0000000..d00e52d --- /dev/null +++ b/packages/database-postgres/src/lock.ts @@ -0,0 +1,166 @@ +/** + * Migration advisory lock — [E00-S03-T04]. + * + * The advisory lock is the PostgreSQL-side guarantee that concurrent migration + * runners cannot run at the same time: the migration runner (built on this + * boundary in later stories — failure diagnostic E00-S03-T05, ready gate + * E00-S03-T06) takes the lock before applying migrations, so a second runner + * either waits (`acquire()`) or fails fast (`tryAcquire()`) while the first + * holds it. + * + * The lock lives in `database-postgres` — the single workspace package + * allowed to import the PostgreSQL driver (E00-S03-T02) — and talks to the + * database exclusively through the package-owned `pg` Pool, so no other + * package needs the driver to lock migration state. + * + * Lock model: + * - session-scoped advisory lock (`pg_advisory_lock`/`pg_try_advisory_lock` + * on the `bigint` key `hashtextextended($1, 0)`), so it lives exactly as + * long as the holding database session and no longer; + * - held on a dedicated connection checked out of the caller's pool + * (`pool.connect()`), so the lock never taints a pooled connection that + * other queries share; + * - released explicitly by `release()` (`pg_advisory_unlock` first, then + * the client goes back to the pool; if the unlock statement itself fails, + * the connection is destroyed — `client.release(error)` — so a pooled + * connection is never reused while its session still holds the lock), + * and released implicitly by the server when the holding session ends — + * the issue's rollback note: "the lock releases when the runner exits". + * + * The lock key is always passed as a bound parameter (`$1`) — the SQL never + * interpolates it, so the only interpolated value is the `$1` placeholder + * itself. `hashtextextended($1, 0)` maps the key string to a stable `bigint` + * advisory-lock key: deterministic for a given database, identical across all + * sessions, so every runner contends on the same lock. + */ + +import type { Pool, PoolClient } from 'pg'; + +/** Namespace string for the migration advisory lock. */ +export const MIGRATION_LOCK_KEY = 'personal-blog-migrations'; + +/** + * The migration advisory lock: serializes migration runners so two runners + * can never apply migrations at the same time. Instances are cheap and share + * the caller's pool; the lock performs no ledger writes (migration ledger is + * E00-S03-T03) and no diagnostics (E00-S03-T05). + */ +export class MigrationLock { + private readonly pool: Pool; + private client: PoolClient | null = null; + private held = false; + /** In-flight acquire, memoized so concurrent acquire() calls share one connection (re-entrant-safe). */ + private acquireInFlight: Promise | null = null; + + /** @param pool The package-owned PostgreSQL pool (`pg.Pool`). */ + constructor(pool: Pool) { + this.pool = pool; + } + + /** True while this instance holds the advisory lock. */ + get isHeld(): boolean { + return this.held; + } + + /** + * Acquires the migration advisory lock, blocking until it is free — a + * second runner waits here while the first holds the lock. Idempotent: + * acquiring an already-held instance is a no-op. Re-entrant-safe: + * concurrent `acquire()` calls on the same instance share the single + * in-flight acquire (memoized in `acquireInFlight`), so exactly one + * connection is checked out and no locked connection leaks. The lock is + * held on a dedicated connection until `release()` (or until the session + * ends — e.g. the runner exits and the pool closes its connections). + */ + async acquire(): Promise { + if (this.held) return; + // A second concurrent acquire() on this instance returns the in-flight + // acquire instead of checking out another connection: exactly one + // connection is checked out and no locked connection leaks. + if (this.acquireInFlight !== null) return this.acquireInFlight; + const inFlight = this.doAcquire(); + this.acquireInFlight = inFlight; + try { + await inFlight; + } finally { + this.acquireInFlight = null; + } + } + + /** The memoized acquire body: checks out one dedicated connection and takes the session-scoped lock on it. */ + private async doAcquire(): Promise { + const client = await this.pool.connect(); + try { + await client.query( + 'SELECT pg_advisory_lock(hashtextextended($1, 0))', + [MIGRATION_LOCK_KEY], + ); + } catch (error) { + client.release(); + throw error; + } + this.client = client; + this.held = true; + } + + /** + * Non-blocking acquire: resolves `true` when this runner got the lock and + * `false` when another runner holds it — a second runner fails fast while + * the first holds the lock. Also resolves `true` when this instance + * already holds the lock. + */ + async tryAcquire(): Promise { + if (this.held) return true; + const client = await this.pool.connect(); + try { + const result = await client.query<{ acquired: boolean }>( + 'SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired', + [MIGRATION_LOCK_KEY], + ); + const acquired = result.rows[0]?.acquired === true; + if (acquired) { + this.client = client; + this.held = true; + return true; + } + client.release(); + return false; + } catch (error) { + client.release(); + throw error; + } + } + + /** + * Releases the advisory lock: unlocks it on the holding session + * (`pg_advisory_unlock`) and returns the connection to the pool, so the + * lock is gone before any other query could reuse that pooled connection. + * If the unlock statement fails (e.g. statement timeout or cancellation) + * the connection is destroyed instead — `client.release(error)` tells the + * pool to drop the client, ending the session (and its lock), so a pooled + * connection is never reused while its session still holds the lock; the + * failure is re-thrown. Idempotent: releasing an instance that does not + * hold the lock is a no-op. + */ + async release(): Promise { + if (!this.held || this.client === null) return; + const client = this.client; + this.client = null; + this.held = false; + try { + await client.query( + 'SELECT pg_advisory_unlock(hashtextextended($1, 0))', + [MIGRATION_LOCK_KEY], + ); + } catch (error) { + // The unlock statement failed while this session may still hold the + // advisory lock. The connection must not go back into the pool in that + // state — the next borrower would block every other runner. Destroy it + // (release(error) removes the client from the pool), so the session — + // and its lock — ends. + client.release(error as Error); + throw error; + } + client.release(); + } +} diff --git a/tests/database-postgres-lock.test.mjs b/tests/database-postgres-lock.test.mjs new file mode 100644 index 0000000..016c1ed --- /dev/null +++ b/tests/database-postgres-lock.test.mjs @@ -0,0 +1,837 @@ +/** + * Migration advisory lock test — locks in the [E00-S03-T04] advisory lock + * that prevents concurrent migration runners. + * + * Acceptance criteria covered (each test fails without the committed state): + * - "advisory lock prevents concurrent migration runners" → the committed + * `packages/database-postgres/src/lock.ts` defines a `MigrationLock` + * backed by PostgreSQL session-scoped advisory locks keyed by a stable + * constant (`MIGRATION_LOCK_KEY` → `hashtextextended($1, 0)`): `acquire()` + * blocks (`pg_advisory_lock`) and `tryAcquire()` fails fast + * (`pg_try_advisory_lock`), both on a dedicated pool connection; the + * docker-gated real-stack probe starts a real database and proves two + * runners cannot hold the lock at once (the issue's test plan: "run a + * concurrent migration lock test"). + * - "a second runner waits or fails while the first holds the lock" → the + * real-stack probe holds the lock in the first runner, shows the second + * runner's `tryAcquire()` returns false (fails fast) and its blocking + * `acquire()` waits (the server reports `wait_event advisory` for the + * waiting session) until the first runner releases, then acquires; a + * session that takes the lock and then closes (the runner exits) + * releases it — the issue's rollback note ("the lock releases when the + * runner exits"). + * - "concurrent acquire() calls on the same lock instance are + * re-entrant-safe" → `acquire()` memoizes the in-flight acquire + * (`acquireInFlight`): a second concurrent `acquire()` on the same + * instance returns the in-flight acquire instead of checking out another + * connection, so exactly one connection is checked out and no locked + * connection leaks; the real-stack probe proves it behaviorally (two + * concurrent acquire() calls → `pool.totalCount` stays 1, one granted + * advisory lock, none left after release). + * - the release() failure path (reviewer finding F4): if the + * `pg_advisory_unlock` statement fails, `release()` destroys the + * connection (`client.release(error)` — the pool drops the client, ending + * the session and its lock) instead of returning it to the pool, so a + * pooled connection is never reused while its session still holds the + * lock; locked in statically, by a mutation probe, and behaviorally by a + * deterministic stub-pool probe of the committed module's control flow. + * - the lock is part of the driver boundary: `src/index.ts` re-exports + * `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the + * `pg` driver to lock migration state. + * + * Run: `node --test tests/database-postgres-lock.test.mjs` + * (node:test — built into Node >= 18; no dependencies, lockfile untouched.) + */ + +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { readFileSync, existsSync, writeFileSync, rmSync } from 'node:fs'; +import { spawnSync } from 'node:child_process'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; + +const REPO_ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..'); + +const read = (relPath) => readFileSync(path.join(REPO_ROOT, relPath), 'utf8'); + +/** The committed advisory-lock module and the driver boundary that re-exports it. */ +const LOCK_SRC = 'packages/database-postgres/src/lock.ts'; +const INDEX_SRC = 'packages/database-postgres/src/index.ts'; + +/** The root test glob (root `scripts.test`, E00-S01-T12) that runs every suite. */ +const ROOT_TEST_GLOB = 'tests/**/*.test.mjs'; + +/** The CI job that gates the advisory-lock criterion on every PR. */ +const CI_JOB = 'database-postgres-lock'; + +/** The committed advisory-lock key (stable namespace, bound as $1). */ +const LOCK_KEY = 'personal-blog-migrations'; + +// --------------------------------------------------------------------------- +// Static assertions on the committed lock source +// --------------------------------------------------------------------------- + +/** + * Asserts the lock module exists with the documented shape: a `MigrationLock` + * class and a stable `MIGRATION_LOCK_KEY` constant. Fails fast on a missing + * module; the mutation probes below prove the assertions are non-vacuous. + */ +function assertLockModule(src) { + assert.match(src, /export class MigrationLock/, 'the lock module must export the MigrationLock class'); + assert.match( + src, + /export const MIGRATION_LOCK_KEY = 'personal-blog-migrations'/, + 'the lock module must declare the MIGRATION_LOCK_KEY constant', + ); + assert.match( + src, + /get isHeld\(\)/, + 'the lock module must expose isHeld so callers can tell whether the instance holds the lock', + ); +} + +/** + * Asserts the blocking path: `acquire()` takes a session-scoped advisory lock + * (`SELECT pg_advisory_lock(hashtextextended($1, 0))`) on a dedicated + * connection (`pool.connect()`), keyed by the stable constant bound as a + * parameter — a second runner calling `acquire()` waits while the first holds + * the lock. The key is never interpolated into the SQL. + */ +function assertBlockingAcquire(src) { + assert.match( + src, + /SELECT pg_advisory_lock\(hashtextextended\(\$1, 0\)\)/, + 'acquire() must take the blocking session-scoped advisory lock (SELECT pg_advisory_lock(hashtextextended($1, 0)))', + ); + assert.match(src, /pool\.connect\(\)/, 'the lock must be held on a dedicated connection (pool.connect())'); + assert.match(src, /\[MIGRATION_LOCK_KEY\]/, 'the lock key must be passed as the bound-parameter array ([MIGRATION_LOCK_KEY])'); + assert.doesNotMatch( + src, + /\$\{MIGRATION_LOCK_KEY\}/, + 'the lock key must never be interpolated into the SQL', + ); +} + +/** + * Asserts the non-blocking path: `tryAcquire()` uses + * `SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired` and + * resolves `false` when another runner holds the lock — a second runner fails + * fast while the first holds the lock. + */ +function assertNonBlockingAcquire(src) { + assert.match( + src, + /SELECT pg_try_advisory_lock\(hashtextextended\(\$1, 0\)\) AS acquired/, + 'tryAcquire() must use the non-blocking advisory lock (SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired)', + ); + assert.match(src, /result\.rows\[0\]\?\.acquired === true/, 'tryAcquire() must resolve true only when the lock was acquired'); + assert.match(src, /return false;/, 'tryAcquire() must resolve false when another runner holds the lock'); +} + +/** + * Asserts the re-entrancy guard: `acquire()` memoizes the in-flight acquire + * (`acquireInFlight`) so a concurrent `acquire()` on the same instance + * returns the same in-flight promise instead of checking out another + * connection — exactly one connection is checked out and no locked + * connection leaks. The memo is cleared once the acquire settles, so later + * acquire() calls behave normally. + */ +function assertReentrantAcquire(src) { + assert.match( + src, + /private acquireInFlight: Promise \| null = null;/, + 'the lock module must memoize the in-flight acquire (acquireInFlight field)', + ); + const acquireBody = src.slice(src.indexOf('async acquire()'), src.indexOf('async tryAcquire()')); + assert.match( + acquireBody, + /if \(this\.acquireInFlight !== null\) return this\.acquireInFlight;/, + 'a concurrent acquire() on the same instance must return the in-flight acquire (memoized) instead of checking out another connection', + ); + assert.match( + acquireBody, + /this\.acquireInFlight = null;/, + 'the in-flight acquire memo must be cleared once the acquire settles', + ); +} + +/** + * Asserts the session-scoped release: `release()` unlocks the holding session + * (`SELECT pg_advisory_unlock(hashtextextended($1, 0))`) BEFORE returning the + * connection to the pool (`client.release()`), so a pooled connection is + * never reused while still holding the lock; releasing an unheld instance is + * a no-op. + */ +function assertSessionScopedRelease(src) { + const releaseBody = src.slice(src.indexOf('async release()')); + assert.match( + releaseBody, + /SELECT pg_advisory_unlock\(hashtextextended\(\$1, 0\)\)/, + 'release() must unlock the holding session (SELECT pg_advisory_unlock(hashtextextended($1, 0)))', + ); + assert.match(releaseBody, /client\.release\(\)/, 'release() must return the connection to the pool (client.release())'); + assert.ok( + releaseBody.indexOf('pg_advisory_unlock') < releaseBody.indexOf('client.release()'), + 'release() must unlock the session before returning the connection to the pool', + ); +} + +/** + * Asserts the unlock-failure path (F4): when the `pg_advisory_unlock` + * statement fails, `release()` must DESTROY the connection — + * `client.release(error)` tells the pool to drop the client, ending the + * session (and its advisory lock) — so a pooled connection is never reused + * while its session still holds the lock. The plain `client.release()` + * (returning the connection to the pool) must live only on the success path, + * after the unlock try/catch. + */ +function assertReleaseDestroysOnUnlockFailure(src) { + const releaseBody = src.slice(src.indexOf('async release()')); + assert.match( + releaseBody, + /SELECT pg_advisory_unlock\(hashtextextended\(\$1, 0\)\)/, + 'release() must unlock the holding session (SELECT pg_advisory_unlock(hashtextextended($1, 0)))', + ); + assert.match( + releaseBody, + /catch \(error\) \{[\s\S]*?client\.release\(error/, + 'on unlock failure release() must destroy the connection (client.release(error) in the catch of the unlock query)', + ); + const destroyIdx = releaseBody.indexOf('client.release(error'); + const catchIdx = releaseBody.indexOf('catch (error) {'); + const plainReleaseIdx = releaseBody.indexOf('client.release()'); + assert.ok( + destroyIdx !== -1 && + catchIdx !== -1 && + plainReleaseIdx !== -1 && + destroyIdx < plainReleaseIdx && + catchIdx < plainReleaseIdx, + 'the plain client.release() (return to pool) must come after the unlock failure handling (destroy path) — a still-locked connection is never returned to the pool', + ); +} + +/** Asserts the driver boundary re-exports the lock from the package entrypoint. */ +function assertBoundaryReexport(src) { + assert.match( + src, + /export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js'/, + 'the driver boundary must re-export the lock (src/index.ts → ./lock.js)', + ); +} + +// --------------------------------------------------------------------------- +// Docker probe helpers (integration test skips cleanly without Docker) +// --------------------------------------------------------------------------- + +function run(cmd, args, opts = {}) { + return spawnSync(cmd, args, { + encoding: 'utf8', + timeout: 600_000, + ...opts, + }); +} + +/** True when the `docker` CLI with the Compose plugin is on PATH. */ +function dockerComposeAvailable() { + try { + return run('docker', ['compose', 'version'], { timeout: 15_000 }).status === 0; + } catch { + return false; + } +} + +/** True when a reachable Docker daemon exists. */ +function dockerDaemonAvailable() { + try { + return run('docker', ['info'], { timeout: 15_000 }).status === 0; + } catch { + return false; + } +} + +/** Parses `docker compose ps --format json` (JSON array or one object per line). */ +function parsePsJson(stdout) { + const text = String(stdout).trim(); + if (!text) return []; + try { + const parsed = JSON.parse(text); + return Array.isArray(parsed) ? parsed : [parsed]; + } catch { + return text + .split('\n') + .map((line) => line.trim()) + .filter(Boolean) + .map((line) => JSON.parse(line)); + } +} + +/** Tolerant field lookup across compose ps JSON shapes. */ +function field(container, ...names) { + for (const name of names) { + if (container[name] !== undefined) return container[name]; + } + return undefined; +} + +/** + * Polls `docker compose ps` until the db container of the given project + * reports healthy (or the deadline passes), so the probe never races a + * still-booting database. + */ +function waitForDbHealthy(project, deadlineMs = 60_000) { + const deadline = Date.now() + deadlineMs; + let last = ''; + while (Date.now() < deadline) { + const ps = run('docker', ['compose', '-p', project, 'ps', '--format', 'json'], { + cwd: REPO_ROOT, + timeout: 15_000, + }); + if (ps.status === 0) { + last = ps.stdout; + const db = parsePsJson(ps.stdout).find((c) => field(c, 'Service', 'service') === 'db'); + if (db && /healthy/i.test(String(field(db, 'Health', 'health') ?? ''))) return; + } + run(process.execPath, ['-e', 'setTimeout(() => {}, 1000)']); // db still booting — retry + } + throw new Error(`the db container did not become healthy within ${deadlineMs}ms (last ps: "${last.trim()}")`); +} + +/** + * The host-side probe body: exercises the committed lock module against a + * real database (DATABASE_URL) exactly as concurrent migration runners will — + * first runner holds the lock, second runner fails fast AND waits, acquires + * after release, and a session that closes (the runner exits) releases the + * lock. Written to a temp file inside `packages/database-postgres/` so `pg` + * resolves through the package's own dependency links, then removed. + */ +const PROBE_SOURCE = ` +import { Pool, Client } from 'pg'; +import { MigrationLock, MIGRATION_LOCK_KEY } from './src/lock.ts'; + +const pool = new Pool({ connectionString: process.env.DATABASE_URL }); + +const advisoryCount = async () => { + const r = await pool.query( + "SELECT count(*)::int AS n FROM pg_locks WHERE locktype = 'advisory' AND granted", + ); + return r.rows[0].n; +}; + +const advisoryWaitEvents = async () => { + const r = await pool.query( + \`SELECT wait_event_type, wait_event FROM pg_stat_activity + WHERE datname = current_database() AND pid <> pg_backend_pid()\`, + ); + return r.rows + .filter((row) => row.wait_event_type !== null && row.wait_event !== null) + .map((row) => row.wait_event_type + ':' + row.wait_event); +}; + +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + +try { + const baseline = await advisoryCount(); + + // First runner takes the lock: the server must hold one granted advisory + // lock more than before. + const lockA = new MigrationLock(pool); + await lockA.acquire(); + const heldCount = await advisoryCount(); + + // A second runner fails fast while the first holds the lock. + const lockB = new MigrationLock(pool); + const secondFailsFast = await lockB.tryAcquire(); + + // A second runner waits while the first holds the lock: the blocking + // acquire() must not resolve, and the server must report the waiting + // session (wait_event 'advisory'). + const waitPromise = lockB.acquire(); + let observedWait = false; + const waitDeadline = Date.now() + 15_000; + while (Date.now() < waitDeadline) { + if ((await advisoryWaitEvents()).includes('Lock:advisory')) { + observedWait = true; + break; + } + await sleep(100); + } + + // First runner finishes -> the waiting runner acquires, then releases. + await lockA.release(); + await waitPromise; + const secondAcquiredAfterRelease = lockB.isHeld; + await lockB.release(); + const countAfterRelease = await advisoryCount(); + + // Session exit releases the lock (rollback note: "the lock releases when + // the runner exits"): a dedicated connection takes the same advisory lock + // and then closes; the server must release the lock with the session. + const exitClient = new Client({ connectionString: process.env.DATABASE_URL }); + await exitClient.connect(); + await exitClient.query('SELECT pg_advisory_lock(hashtextextended($1, 0))', [MIGRATION_LOCK_KEY]); + const heldWhileSessionOpen = await advisoryCount(); + await exitClient.end(); // the runner exits -> session closes -> lock releases + let thirdAcquiresAfterSessionEnd = false; + const retryDeadline = Date.now() + 5_000; + while (Date.now() < retryDeadline) { + const lockC = new MigrationLock(pool); + if (await lockC.tryAcquire()) { + thirdAcquiresAfterSessionEnd = true; + await lockC.release(); + break; + } + await sleep(100); + } + + // Re-entrant-safe acquire (acceptance criterion: concurrent acquire() on + // the same lock instance is re-entrant-safe — the in-flight acquire is + // memoized so exactly one connection is checked out and no locked + // connection leaks). A fresh pool so the connection counts are exact: two + // concurrent acquire() calls must share one in-flight acquire — one + // connection checked out (pool.totalCount stays 1), the lock granted once, + // and nothing left held after release. + let reentrantPool = null; + try { + reentrantPool = new Pool({ connectionString: process.env.DATABASE_URL }); + const lockR = new MigrationLock(reentrantPool); + await Promise.all([lockR.acquire(), lockR.acquire()]); + const reentrantPoolTotal = reentrantPool.totalCount; + const reentrantPoolIdle = reentrantPool.idleCount; + const reentrantHeld = lockR.isHeld; + const reentrantHeldCount = await advisoryCount(); + await lockR.release(); + const reentrantCountAfterRelease = await advisoryCount(); + await reentrantPool.end(); + reentrantPool = null; + + console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ + baseline, + heldCount, + secondFailsFast, + observedWait, + secondAcquiredAfterRelease, + countAfterRelease, + heldWhileSessionOpen, + thirdAcquiresAfterSessionEnd, + reentrantPoolTotal, + reentrantPoolIdle, + reentrantHeld, + reentrantHeldCount, + reentrantCountAfterRelease, + })); + } finally { + if (reentrantPool !== null) { + await reentrantPool.end().catch(() => {}); + } + } +} finally { + await pool.end(); +} +`; + +/** + * Deterministic behavioral probe of the release() failure path (F4): models + * pg's PoolClient disposal contract — release() returns the client to the + * pool, release(error) destroys it (the pool removes the client, ending its + * session) — and drives the COMMITTED lock module's release() directly with a + * stub pool, so the check is about the module's behavior (which disposal call + * it makes), not its source shape. No database or Docker needed. + */ +const RELEASE_FAILURE_PROBE_SOURCE = ` +import { MigrationLock } from './src/lock.ts'; + +// A stub pg pool whose client fails (or succeeds) the unlock query and +// records how release() disposes of it. +const makeStub = (unlockSucceeds) => { + const events = []; + const client = { + async query() { + if (unlockSucceeds) { + events.push('unlock:ok'); + return { rows: [] }; + } + const err = new Error('unlock statement failed (statement timeout)'); + err.code = '57014'; + events.push('unlock:fail'); + throw err; + }, + release(error) { + events.push(error === undefined ? 'release()' : 'release(error)'); + if (error !== undefined) client.destroyed = true; + }, + }; + const pool = { connect: async () => client }; + return { pool, client, events }; +}; + +const results = {}; + +// Failure path: the unlock query rejects -> release() must DESTROY the +// connection (release(error)), not return it to the pool while its session +// still holds the advisory lock. +{ + const { pool, client, events } = makeStub(false); + const lock = new MigrationLock(pool); + lock.held = true; // drive the committed release() directly (no DB needed) + lock.client = client; + let rejected = false; + try { + await lock.release(); + } catch (error) { + rejected = String(error.message).includes('unlock statement failed'); + } + results.failurePath = { + rejected, + events, + destroyed: client.destroyed === true, + stateCleared: lock.isHeld === false && lock.client === null, + }; +} + +// Success path: the unlock query resolves -> release() returns the connection +// to the pool with the plain release(). +{ + const { pool, client, events } = makeStub(true); + const lock = new MigrationLock(pool); + lock.held = true; + lock.client = client; + await lock.release(); + results.successPath = { + events, + destroyed: client.destroyed === true, + stateCleared: lock.isHeld === false && lock.client === null, + }; +} + +console.log('LOCK_RELEASE_FAILURE_PROBE_RESULT ' + JSON.stringify(results)); +`; + +// --------------------------------------------------------------------------- +// Criterion tests +// --------------------------------------------------------------------------- + +test('the lock module exists in the driver-owner package and defines the migration lock key', () => { + assert.ok(existsSync(path.join(REPO_ROOT, LOCK_SRC)), `committed ${LOCK_SRC} must exist`); + const src = read(LOCK_SRC); + assertLockModule(src); +}); + +test('acquire() takes the blocking session-scoped advisory lock keyed by a bound parameter (a second runner waits)', () => { + assertBlockingAcquire(read(LOCK_SRC)); +}); + +test('tryAcquire() fails fast while another runner holds the lock (a second runner fails)', () => { + assertNonBlockingAcquire(read(LOCK_SRC)); +}); + +test('concurrent acquire() on the same instance is re-entrant-safe (in-flight acquire memoized, exactly one connection)', () => { + assertReentrantAcquire(read(LOCK_SRC)); +}); + +test('release() unlocks the holding session and returns the connection to the pool', () => { + assertSessionScopedRelease(read(LOCK_SRC)); +}); + +test('release() destroys the connection when the unlock statement fails (never reuses a still-locked connection)', () => { + assertReleaseDestroysOnUnlockFailure(read(LOCK_SRC)); +}); + +test('the driver boundary re-exports the lock from the package entrypoint', () => { + const src = read(INDEX_SRC); + assert.match( + src, + /from '\.\/lock\.js'/, + 'the driver boundary must import the lock module (from \'./lock.js\')', + ); + assertBoundaryReexport(src); +}); + +test('the advisory-lock criterion is enforced in CI', () => { + // Picked up by the root test command (root `scripts.test` glob). + const scripts = JSON.parse(read('package.json')).scripts ?? {}; + assert.equal( + scripts.test, + `node --test "${ROOT_TEST_GLOB}"`, + `root scripts.test must run the "${ROOT_TEST_GLOB}" glob so this suite runs with the rest`, + ); + // And a dedicated CI job gates it on every PR. + const workflow = read('.gitea/workflows/ci.yml'); + assert.ok( + workflow.includes(`node --test tests/database-postgres-lock.test.mjs`), + `CI must run the advisory-lock suite (job "${CI_JOB}") on every PR`, + ); +}); + +// --------------------------------------------------------------------------- +// Real-stack probe — the issue's test plan: "run a concurrent migration +// lock test" against a real database +// --------------------------------------------------------------------------- + +const DOCKER_COMPOSE = dockerComposeAvailable(); +const DOCKER_DAEMON = dockerDaemonAvailable(); + +/** True when this Node can execute the committed `.ts` lock module (>= 23.6, type stripping). */ +const TS_STRIPPING = (() => { + const [major, minor] = process.versions.node.split('.').map(Number); + return major > 23 || (major === 23 && minor >= 6); +})(); + +// An isolated compose project + non-default host port so this probe never +// collides with the compose-config suite's default-project containers, the +// ledger probe's project/port, or the default 5432 binding when several +// suites run on the same host. +const COMPOSE_PROJECT = 'eppp-lock-probe'; +const POSTGRES_HOST_PORT = '55433'; +const DATABASE_URL = `postgres://eppp:eppp@127.0.0.1:${POSTGRES_HOST_PORT}/eppp`; + +test('concurrent migration runners are serialized by the advisory lock (real stack)', { skip: !DOCKER_COMPOSE || !DOCKER_DAEMON || !TS_STRIPPING }, () => { + // The issue's test plan: "run a concurrent migration lock test". The probe + // starts the committed compose `db` service (its own project + host port), + // runs the committed lock module against it with two runner instances, and + // asserts the first runner holds the lock, the second fails fast AND waits + // until the first releases (server-reported wait_event 'advisory'), and a + // session that closes — the runner exits — releases the lock (rollback + // note). Rollback: `docker compose -p eppp-lock-probe down`. + const composeEnv = { ...process.env, POSTGRES_PORT: POSTGRES_HOST_PORT }; + const compose = (args, opts = {}) => + run('docker', ['compose', '-p', COMPOSE_PROJECT, ...args], { cwd: REPO_ROOT, env: composeEnv, ...opts }); + + const probeFile = path.join(REPO_ROOT, 'packages/database-postgres', `.lock-probe-${process.pid}.mjs`); + + try { + const up = compose(['up', '-d', 'db'], { timeout: 180_000 }); + assert.equal( + up.status, + 0, + `"docker compose up -d db" must exit 0:\n${(up.stdout || '')}\n${(up.stderr || '')}`.trim(), + ); + waitForDbHealthy(COMPOSE_PROJECT); + + // Run the committed lock module against the real database. + writeFileSync(probeFile, PROBE_SOURCE); + const probe = run(process.execPath, [probeFile], { + cwd: REPO_ROOT, + env: { ...process.env, DATABASE_URL }, + timeout: 60_000, + }); + assert.equal( + probe.status, + 0, + `the lock probe must exit 0:\n${(probe.stdout || '')}\n${(probe.stderr || '')}`.trim(), + ); + const resultLine = (probe.stdout || '') + .split('\n') + .map((l) => l.trim()) + .find((l) => l.startsWith('LOCK_PROBE_RESULT')); + assert.ok(resultLine, `the lock probe must report a result line (stdout: "${(probe.stdout || '').trim()}")`); + const result = JSON.parse(resultLine.slice('LOCK_PROBE_RESULT'.length).trim()); + + // "advisory lock prevents concurrent migration runners": the first runner + // holds exactly one granted advisory lock more than the empty baseline. + assert.equal( + result.heldCount, + result.baseline + 1, + `acquire() must hold one granted advisory lock (baseline ${result.baseline}, got ${result.heldCount})`, + ); + + // "a second runner fails while the first holds the lock". + assert.equal( + result.secondFailsFast, + false, + 'tryAcquire() must fail fast while the first runner holds the lock', + ); + + // "a second runner waits while the first holds the lock": the blocking + // acquire() does not resolve and the server reports the waiting session. + assert.equal( + result.observedWait, + true, + 'the blocking acquire() must wait (server-reported wait_event advisory) while the first runner holds the lock', + ); + + // The waiting runner acquires once the first runner releases. + assert.equal( + result.secondAcquiredAfterRelease, + true, + 'the waiting runner must acquire the lock once the first runner releases it', + ); + + // Nothing is left held after both runners release. + assert.equal( + result.countAfterRelease, + result.baseline, + `no advisory lock may remain granted after both runners release (got ${result.countAfterRelease}, expected ${result.baseline})`, + ); + + // Rollback note: "the lock releases when the runner exits" — a session + // that holds the lock and then closes releases it. + assert.equal( + result.heldWhileSessionOpen, + result.baseline + 1, + `the dedicated session must hold the advisory lock while open (baseline ${result.baseline}, got ${result.heldWhileSessionOpen})`, + ); + assert.equal( + result.thirdAcquiresAfterSessionEnd, + true, + 'the lock must be acquirable again after the holding session closes (the runner exits)', + ); + + // "concurrent acquire() on the same lock instance is re-entrant-safe": + // two concurrent acquire() calls share one in-flight acquire — exactly + // one connection is checked out (pool.totalCount stays 1), the lock is + // granted once, and nothing is left held after release (no locked + // connection leaks). + assert.equal( + result.reentrantPoolTotal, + 1, + `two concurrent acquire() calls must check out exactly one connection (pool.totalCount ${result.reentrantPoolTotal})`, + ); + assert.equal( + result.reentrantPoolIdle, + 0, + `the re-entrant acquire must hold its connection checked out while the lock is held (pool.idleCount ${result.reentrantPoolIdle})`, + ); + assert.equal( + result.reentrantHeld, + true, + 'the re-entrant acquire must leave the instance holding the lock', + ); + assert.equal( + result.reentrantHeldCount, + result.baseline + 1, + `two concurrent acquire() calls must grant the advisory lock once (baseline ${result.baseline}, got ${result.reentrantHeldCount})`, + ); + assert.equal( + result.reentrantCountAfterRelease, + result.baseline, + `no advisory lock may remain granted after the re-entrant instance releases (got ${result.reentrantCountAfterRelease})`, + ); + } finally { + compose(['down', '-v'], { timeout: 120_000 }); + rmSync(probeFile, { force: true }); + } +}); + +test('release() destroys the connection on unlock failure and returns it to the pool on success (behavioral probe)', { skip: !TS_STRIPPING }, () => { + // F4: pg's release(error) destroys the client (the pool removes it, ending + // the session and its lock); the plain release() returns it to the pool. + // The probe drives the committed release() with a stub pool: a failed unlock + // must take the destroy path (release(error)), a successful unlock the + // return path (plain release()) — a still-locked connection is never reused. + const probeFile = path.join(REPO_ROOT, 'packages/database-postgres', `.lock-probe-release-${process.pid}.mjs`); + try { + writeFileSync(probeFile, RELEASE_FAILURE_PROBE_SOURCE); + const probe = run(process.execPath, [probeFile], { cwd: REPO_ROOT, timeout: 30_000 }); + assert.equal( + probe.status, + 0, + `the release-failure probe must exit 0:\n${(probe.stdout || '')}\n${(probe.stderr || '')}`.trim(), + ); + const resultLine = (probe.stdout || '') + .split('\n') + .map((l) => l.trim()) + .find((l) => l.startsWith('LOCK_RELEASE_FAILURE_PROBE_RESULT')); + assert.ok(resultLine, `the release-failure probe must report a result line (stdout: "${(probe.stdout || '').trim()}")`); + const result = JSON.parse(resultLine.slice('LOCK_RELEASE_FAILURE_PROBE_RESULT'.length).trim()); + + // Failure path: the unlock statement failed -> the connection is destroyed + // (release(error)), never returned to the pool while still locked. + assert.equal( + result.failurePath.rejected, + true, + 'release() must reject when the unlock statement fails', + ); + assert.deepEqual( + result.failurePath.events, + ['unlock:fail', 'release(error)'], + 'on unlock failure release() must destroy the connection (client.release(error)), not return it to the pool', + ); + assert.equal( + result.failurePath.destroyed, + true, + 'the failed-unlock connection must be destroyed (removed from the pool)', + ); + assert.equal( + result.failurePath.stateCleared, + true, + 'release() must clear held/client even on the failure path', + ); + + // Success path: the unlock succeeded -> the connection is returned to the + // pool with the plain release(). + assert.deepEqual( + result.successPath.events, + ['unlock:ok', 'release()'], + 'on a successful unlock release() must return the connection to the pool (plain client.release())', + ); + assert.equal( + result.successPath.destroyed, + false, + 'a successful unlock must not destroy the connection', + ); + } finally { + rmSync(probeFile, { force: true }); + } +}); + +// --------------------------------------------------------------------------- +// Non-vacuous probes — the assertions above really do fail on violations +// --------------------------------------------------------------------------- + +test('removing pg_advisory_lock makes the blocking-acquire criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutBlocking = src.replace(/SELECT pg_advisory_lock\(/, 'SELECT pg_advisory_xact_lock('); + assert.notEqual(withoutBlocking, src, 'the mutation must actually replace the blocking lock call'); + assert.throws(() => assertBlockingAcquire(withoutBlocking), /pg_advisory_lock/); +}); + +test('removing pg_try_advisory_lock makes the non-blocking criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutTry = src.replace(/SELECT pg_try_advisory_lock\(/, 'SELECT pg_advisory_lock('); + assert.notEqual(withoutTry, src, 'the mutation must actually replace the non-blocking lock call'); + assert.throws(() => assertNonBlockingAcquire(withoutTry), /pg_try_advisory_lock/); +}); + +test('interpolating the lock key into the SQL makes the keyed-parameter criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const interpolated = src.replace(/hashtextextended\(\$1, 0\)/g, "hashtextextended('${MIGRATION_LOCK_KEY}', 0)"); + assert.notEqual(interpolated, src, 'the mutation must actually replace the $1 placeholder'); + assert.throws(() => assertBlockingAcquire(interpolated), /\$1/); +}); + +test('removing pg_advisory_unlock makes the session-release criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutUnlock = src.replace(/SELECT pg_advisory_unlock\(/, 'SELECT pg_advisory_lock('); + assert.notEqual(withoutUnlock, src, 'the mutation must actually replace the unlock call'); + assert.throws(() => assertSessionScopedRelease(withoutUnlock), /pg_advisory_unlock/); +}); + +test('returning the connection to the pool on unlock failure makes the failure-path criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + // The buggy shape F4 describes: a try/finally that always returns the client + // to the pool, so a failed unlock reuses a still-locked connection. + const withoutDestroy = src.replace('client.release(error as Error);', 'client.release();'); + assert.notEqual(withoutDestroy, src, 'the mutation must actually replace the failure-path destroy with a plain release'); + assert.throws(() => assertReleaseDestroysOnUnlockFailure(withoutDestroy), /client\.release\(error/); +}); + +test('removing the in-flight acquire memo makes the re-entrancy criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutMemo = src.replace(/if \(this\.acquireInFlight !== null\) return this\.acquireInFlight;\n/, ''); + assert.notEqual(withoutMemo, src, 'the mutation must actually remove the in-flight acquire guard'); + assert.throws(() => assertReentrantAcquire(withoutMemo), /in-flight acquire/); +}); + +test('dropping the lock re-export from the boundary fails the boundary criterion (mutation probe)', () => { + const src = read(INDEX_SRC); + const withoutReexport = src.replace(/export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js';\n/, ''); + assert.notEqual(withoutReexport, src, 'the mutation must actually remove the lock re-export'); + assert.throws(() => assertBoundaryReexport(withoutReexport), /re-export/); +}); + +test('a placeholder lock module fails the advisory-lock criterion (mutation probe)', () => { + assert.throws( + () => assertLockModule('export class MigrationLock {}\n'), + /MIGRATION_LOCK_KEY/, + ); +});