Files
PersonalBlog/tests/database-postgres-lock.test.mjs
implementer cb1b161ce9
CI / Frozen lockfile install (pull_request) Successful in 47s
CI / Secrets not embedded (E00-S02-T08) (pull_request) Successful in 26s
CI / Database-postgres import isolation (E00-S03-T02) (pull_request) Successful in 25s
CI / Migration ledger (E00-S03-T03) (pull_request) Successful in 41s
CI / Migration advisory lock (E00-S03-T04) (pull_request) Successful in 44s
CI / Compose config (E00-S03-T01) (pull_request) Successful in 25s
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.
2026-08-30 00:49:56 +00:00

838 lines
34 KiB
JavaScript

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