/** * 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/, ); });