From 3f80063303560e10cb569831277f8a0d0099ae7e Mon Sep 17 00:00:00 2001 From: implementer Date: Sun, 30 Aug 2026 00:08:57 +0000 Subject: [PATCH 1/4] feat: add migration advisory lock to database-postgres (E00-S03-T04) --- .gitignore | 5 +- docs/development/non-container.md | 2 +- packages/database-postgres/package.json | 2 +- packages/database-postgres/src/index.ts | 9 +- packages/database-postgres/src/lock.ts | 132 ++++++++++++++++++++++++ 5 files changed, 142 insertions(+), 8 deletions(-) create mode 100644 packages/database-postgres/src/lock.ts 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..88a3def --- /dev/null +++ b/packages/database-postgres/src/lock.ts @@ -0,0 +1,132 @@ +/** + * 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), 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; + + /** @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. 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; + 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. + * 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], + ); + } finally { + client.release(); + } + } +} From a2b98439e92aecf9cdc48548963c0fc6ee86cb3b Mon Sep 17 00:00:00 2001 From: implementer Date: Sun, 30 Aug 2026 00:13:46 +0000 Subject: [PATCH 2/4] test: lock in the migration advisory lock with static + real-stack probes (E00-S03-T04) --- .gitea/workflows/ci.yml | 26 ++ tests/database-postgres-lock.test.mjs | 536 ++++++++++++++++++++++++++ 2 files changed, 562 insertions(+) create mode 100644 tests/database-postgres-lock.test.mjs diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 77e8bcc..22652f1 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -77,6 +77,32 @@ 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, 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; 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/tests/database-postgres-lock.test.mjs b/tests/database-postgres-lock.test.mjs new file mode 100644 index 0000000..8170a7b --- /dev/null +++ b/tests/database-postgres-lock.test.mjs @@ -0,0 +1,536 @@ +/** + * 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"). + * - 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 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 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); + } + + console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ + baseline, + heldCount, + secondFailsFast, + observedWait, + secondAcquiredAfterRelease, + countAfterRelease, + heldWhileSessionOpen, + thirdAcquiresAfterSessionEnd, + })); +} finally { + await pool.end(); +} +`; + +// --------------------------------------------------------------------------- +// 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('release() unlocks the holding session and returns the connection to the pool', () => { + assertSessionScopedRelease(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)', + ); + } finally { + compose(['down', '-v'], { timeout: 120_000 }); + 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('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/, + ); +}); From dac3679e33364dcb8453282873b3c424099e6b4f Mon Sep 17 00:00:00 2001 From: implementer Date: Sun, 30 Aug 2026 00:29:20 +0000 Subject: [PATCH 3/4] fix: memoize in-flight acquire so concurrent acquire() is re-entrant-safe (E00-S03-T04) Security-review finding F2: two concurrent acquire() calls on the same MigrationLock instance could each check out a connection; the second pg_advisory_lock would overwrite this.client, leaking the first locked connection until session end. acquire() now memoizes the in-flight acquire in acquireInFlight and returns it on re-entry, so exactly one connection is checked out and no locked connection leaks. The memo is cleared once the acquire settles. tryAcquire()/release() paths unchanged. Locked in by: - static criterion test: acquireInFlight field, re-entry guard returns the in-flight acquire, memo cleared on settle - mutation probe: removing the re-entry guard fails the criterion - real-stack probe: two concurrent acquire() calls on one instance leave pool.totalCount at 1 (exactly one connection), the lock granted once, nothing left after release; probe fails (hangs) on the pre-fix code - CI job comment updated to reflect the re-entrancy criterion --- .gitea/workflows/ci.yml | 16 ++-- packages/database-postgres/src/lock.ts | 26 ++++- tests/database-postgres-lock.test.mjs | 128 +++++++++++++++++++++++-- 3 files changed, 150 insertions(+), 20 deletions(-) diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 22652f1..9d4a05e 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -80,13 +80,15 @@ jobs: # 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, 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; 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). + # 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 diff --git a/packages/database-postgres/src/lock.ts b/packages/database-postgres/src/lock.ts index 88a3def..413b352 100644 --- a/packages/database-postgres/src/lock.ts +++ b/packages/database-postgres/src/lock.ts @@ -47,6 +47,8 @@ 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) { @@ -61,12 +63,30 @@ export class MigrationLock { /** * 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. 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). + * 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( diff --git a/tests/database-postgres-lock.test.mjs b/tests/database-postgres-lock.test.mjs index 8170a7b..1b0dabe 100644 --- a/tests/database-postgres-lock.test.mjs +++ b/tests/database-postgres-lock.test.mjs @@ -20,6 +20,14 @@ * 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 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. @@ -113,6 +121,33 @@ function assertNonBlockingAcquire(src) { 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 @@ -307,16 +342,47 @@ try { await sleep(100); } - console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ - baseline, - heldCount, - secondFailsFast, - observedWait, - secondAcquiredAfterRelease, - countAfterRelease, - heldWhileSessionOpen, - thirdAcquiresAfterSessionEnd, - })); + // 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(); } @@ -340,6 +406,10 @@ test('tryAcquire() fails fast while another runner holds the lock (a second runn 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)); }); @@ -483,6 +553,37 @@ test('concurrent migration runners are serialized by the advisory lock (real sta 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 }); @@ -521,6 +622,13 @@ test('removing pg_advisory_unlock makes the session-release criterion fail (muta assert.throws(() => assertSessionScopedRelease(withoutUnlock), /pg_advisory_unlock/); }); +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/, ''); From cb1b161ce95f7d18fd390a255946a22008fe4a90 Mon Sep 17 00:00:00 2001 From: implementer Date: Sun, 30 Aug 2026 00:49:56 +0000 Subject: [PATCH 4/4] fix: destroy the connection on unlock failure so a still-locked session is never reused (E00-S03-T04) release() previously returned the connection to the pool in a finally even when the pg_advisory_unlock statement failed, so a pooled connection could be reused while its session still held the migration advisory lock - the next borrower would block every other runner (reviewer finding F4). On unlock failure the connection is now destroyed (client.release(error) removes the client from the pool, ending the session and its lock); the plain client.release() is kept only on the success path. Locked in by a static assertion, a mutation probe, and a deterministic stub-pool behavioral probe of the committed release() control flow. --- packages/database-postgres/src/lock.ts | 26 +++- tests/database-postgres-lock.test.mjs | 193 +++++++++++++++++++++++++ 2 files changed, 213 insertions(+), 6 deletions(-) diff --git a/packages/database-postgres/src/lock.ts b/packages/database-postgres/src/lock.ts index 413b352..d00e52d 100644 --- a/packages/database-postgres/src/lock.ts +++ b/packages/database-postgres/src/lock.ts @@ -21,9 +21,11 @@ * (`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), and released implicitly by the - * server when the holding session ends — the issue's rollback note: "the - * lock releases when the runner exits". + * 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 @@ -133,7 +135,12 @@ export class MigrationLock { * 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. - * Idempotent: releasing an instance that does not hold the lock is a no-op. + * 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; @@ -145,8 +152,15 @@ export class MigrationLock { 'SELECT pg_advisory_unlock(hashtextextended($1, 0))', [MIGRATION_LOCK_KEY], ); - } finally { - client.release(); + } 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 index 1b0dabe..016c1ed 100644 --- a/tests/database-postgres-lock.test.mjs +++ b/tests/database-postgres-lock.test.mjs @@ -28,6 +28,13 @@ * 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. @@ -169,6 +176,40 @@ function assertSessionScopedRelease(src) { ); } +/** + * 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( @@ -388,6 +429,83 @@ try { } `; +/** + * 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 // --------------------------------------------------------------------------- @@ -414,6 +532,10 @@ test('release() unlocks the holding session and returns the connection to the po 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( @@ -590,6 +712,68 @@ test('concurrent migration runners are serialized by the advisory lock (real sta } }); +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 // --------------------------------------------------------------------------- @@ -622,6 +806,15 @@ test('removing pg_advisory_unlock makes the session-release criterion fail (muta 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/, '');