Merge pull request '[E00-S03-T04] Advisory lock prevents concurrent migration runners' (#393) from feature/179 into main
CI / Frozen lockfile install (push) Successful in 43s
CI / Secrets not embedded (E00-S02-T08) (push) Successful in 24s
CI / Database-postgres import isolation (E00-S03-T02) (push) Successful in 35s
CI / Migration ledger (E00-S03-T03) (push) Successful in 41s
CI / Migration advisory lock (E00-S03-T04) (push) Successful in 43s
CI / Compose config (E00-S03-T01) (push) Successful in 29s

This commit was merged in pull request #393.
This commit is contained in:
2026-08-30 01:00:39 +00:00
7 changed files with 1041 additions and 8 deletions
+28
View File
@@ -77,6 +77,34 @@ jobs:
- name: Run migration ledger test suite
run: node --test tests/database-postgres-ledger.test.mjs
# E00-S03-T04: the static assertions of tests/database-postgres-lock.test.mjs
# gate every PR — the suite locks in the migration advisory lock (session-
# scoped pg_advisory_lock/pg_try_advisory_lock over a stable keyed hash on a
# dedicated connection, re-entrant-safe in-flight acquire so concurrent
# acquire() calls share one connection, driver-boundary re-export) with
# mutation probes, and the docker-gated real-stack concurrent probe (a
# second runner waits or fails while the first holds the lock; concurrent
# acquire() checks out exactly one connection; the lock releases when the
# holding session ends) runs where a Docker daemon is available and skips
# cleanly otherwise. The job installs the frozen workspace because the
# real-stack probe executes the committed lock module from the host (it
# imports `pg` through the package's own links).
database-postgres-lock:
name: Migration advisory lock (E00-S03-T04)
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Install Node.js 24
uses: actions/setup-node@v4
with:
node-version: '24'
- name: Enable pnpm (corepack, pinned to 11.23.0 via packageManager)
run: corepack enable
- name: Install dependencies (frozen lockfile)
run: pnpm install --frozen-lockfile
- name: Run migration advisory lock test suite
run: node --test tests/database-postgres-lock.test.mjs
# E00-S03-T01: the static assertions of tests/compose-config.test.mjs (db
# image pinned to postgres:18.6-bookworm, health gate, volume persistence,
# build platforms) gate every PR (the docker-gated real-stack probes inside
+3 -2
View File
@@ -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
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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"
+5 -4
View File
@@ -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';
+166
View File
@@ -0,0 +1,166 @@
/**
* Migration advisory lock — [E00-S03-T04].
*
* The advisory lock is the PostgreSQL-side guarantee that concurrent migration
* runners cannot run at the same time: the migration runner (built on this
* boundary in later stories — failure diagnostic E00-S03-T05, ready gate
* E00-S03-T06) takes the lock before applying migrations, so a second runner
* either waits (`acquire()`) or fails fast (`tryAcquire()`) while the first
* holds it.
*
* The lock lives in `database-postgres` — the single workspace package
* allowed to import the PostgreSQL driver (E00-S03-T02) — and talks to the
* database exclusively through the package-owned `pg` Pool, so no other
* package needs the driver to lock migration state.
*
* Lock model:
* - session-scoped advisory lock (`pg_advisory_lock`/`pg_try_advisory_lock`
* on the `bigint` key `hashtextextended($1, 0)`), so it lives exactly as
* long as the holding database session and no longer;
* - held on a dedicated connection checked out of the caller's pool
* (`pool.connect()`), so the lock never taints a pooled connection that
* other queries share;
* - released explicitly by `release()` (`pg_advisory_unlock` first, then
* the client goes back to the pool; if the unlock statement itself fails,
* the connection is destroyed — `client.release(error)` — so a pooled
* connection is never reused while its session still holds the lock),
* and released implicitly by the server when the holding session ends —
* the issue's rollback note: "the lock releases when the runner exits".
*
* The lock key is always passed as a bound parameter (`$1`) — the SQL never
* interpolates it, so the only interpolated value is the `$1` placeholder
* itself. `hashtextextended($1, 0)` maps the key string to a stable `bigint`
* advisory-lock key: deterministic for a given database, identical across all
* sessions, so every runner contends on the same lock.
*/
import type { Pool, PoolClient } from 'pg';
/** Namespace string for the migration advisory lock. */
export const MIGRATION_LOCK_KEY = 'personal-blog-migrations';
/**
* The migration advisory lock: serializes migration runners so two runners
* can never apply migrations at the same time. Instances are cheap and share
* the caller's pool; the lock performs no ledger writes (migration ledger is
* E00-S03-T03) and no diagnostics (E00-S03-T05).
*/
export class MigrationLock {
private readonly pool: Pool;
private client: PoolClient | null = null;
private held = false;
/** In-flight acquire, memoized so concurrent acquire() calls share one connection (re-entrant-safe). */
private acquireInFlight: Promise<void> | null = null;
/** @param pool The package-owned PostgreSQL pool (`pg.Pool`). */
constructor(pool: Pool) {
this.pool = pool;
}
/** True while this instance holds the advisory lock. */
get isHeld(): boolean {
return this.held;
}
/**
* Acquires the migration advisory lock, blocking until it is free — a
* second runner waits here while the first holds the lock. Idempotent:
* acquiring an already-held instance is a no-op. Re-entrant-safe:
* concurrent `acquire()` calls on the same instance share the single
* in-flight acquire (memoized in `acquireInFlight`), so exactly one
* connection is checked out and no locked connection leaks. The lock is
* held on a dedicated connection until `release()` (or until the session
* ends — e.g. the runner exits and the pool closes its connections).
*/
async acquire(): Promise<void> {
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<void> {
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<boolean> {
if (this.held) return true;
const client = await this.pool.connect();
try {
const result = await client.query<{ acquired: boolean }>(
'SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired',
[MIGRATION_LOCK_KEY],
);
const acquired = result.rows[0]?.acquired === true;
if (acquired) {
this.client = client;
this.held = true;
return true;
}
client.release();
return false;
} catch (error) {
client.release();
throw error;
}
}
/**
* Releases the advisory lock: unlocks it on the holding session
* (`pg_advisory_unlock`) and returns the connection to the pool, so the
* lock is gone before any other query could reuse that pooled connection.
* If the unlock statement fails (e.g. statement timeout or cancellation)
* the connection is destroyed instead — `client.release(error)` tells the
* pool to drop the client, ending the session (and its lock), so a pooled
* connection is never reused while its session still holds the lock; the
* failure is re-thrown. Idempotent: releasing an instance that does not
* hold the lock is a no-op.
*/
async release(): Promise<void> {
if (!this.held || this.client === null) return;
const client = this.client;
this.client = null;
this.held = false;
try {
await client.query(
'SELECT pg_advisory_unlock(hashtextextended($1, 0))',
[MIGRATION_LOCK_KEY],
);
} catch (error) {
// The unlock statement failed while this session may still hold the
// advisory lock. The connection must not go back into the pool in that
// state — the next borrower would block every other runner. Destroy it
// (release(error) removes the client from the pool), so the session —
// and its lock — ends.
client.release(error as Error);
throw error;
}
client.release();
}
}
+837
View File
@@ -0,0 +1,837 @@
/**
* Migration advisory lock test — locks in the [E00-S03-T04] advisory lock
* that prevents concurrent migration runners.
*
* Acceptance criteria covered (each test fails without the committed state):
* - "advisory lock prevents concurrent migration runners" → the committed
* `packages/database-postgres/src/lock.ts` defines a `MigrationLock`
* backed by PostgreSQL session-scoped advisory locks keyed by a stable
* constant (`MIGRATION_LOCK_KEY` → `hashtextextended($1, 0)`): `acquire()`
* blocks (`pg_advisory_lock`) and `tryAcquire()` fails fast
* (`pg_try_advisory_lock`), both on a dedicated pool connection; the
* docker-gated real-stack probe starts a real database and proves two
* runners cannot hold the lock at once (the issue's test plan: "run a
* concurrent migration lock test").
* - "a second runner waits or fails while the first holds the lock" → the
* real-stack probe holds the lock in the first runner, shows the second
* runner's `tryAcquire()` returns false (fails fast) and its blocking
* `acquire()` waits (the server reports `wait_event advisory` for the
* waiting session) until the first runner releases, then acquires; a
* session that takes the lock and then closes (the runner exits)
* releases it — the issue's rollback note ("the lock releases when the
* runner exits").
* - "concurrent acquire() calls on the same lock instance are
* re-entrant-safe" → `acquire()` memoizes the in-flight acquire
* (`acquireInFlight`): a second concurrent `acquire()` on the same
* instance returns the in-flight acquire instead of checking out another
* connection, so exactly one connection is checked out and no locked
* connection leaks; the real-stack probe proves it behaviorally (two
* concurrent acquire() calls → `pool.totalCount` stays 1, one granted
* advisory lock, none left after release).
* - the release() failure path (reviewer finding F4): if the
* `pg_advisory_unlock` statement fails, `release()` destroys the
* connection (`client.release(error)` — the pool drops the client, ending
* the session and its lock) instead of returning it to the pool, so a
* pooled connection is never reused while its session still holds the
* lock; locked in statically, by a mutation probe, and behaviorally by a
* deterministic stub-pool probe of the committed module's control flow.
* - the lock is part of the driver boundary: `src/index.ts` re-exports
* `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the
* `pg` driver to lock migration state.
*
* Run: `node --test tests/database-postgres-lock.test.mjs`
* (node:test — built into Node >= 18; no dependencies, lockfile untouched.)
*/
import test from 'node:test';
import assert from 'node:assert/strict';
import { readFileSync, existsSync, writeFileSync, rmSync } from 'node:fs';
import { spawnSync } from 'node:child_process';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
const REPO_ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..');
const read = (relPath) => readFileSync(path.join(REPO_ROOT, relPath), 'utf8');
/** The committed advisory-lock module and the driver boundary that re-exports it. */
const LOCK_SRC = 'packages/database-postgres/src/lock.ts';
const INDEX_SRC = 'packages/database-postgres/src/index.ts';
/** The root test glob (root `scripts.test`, E00-S01-T12) that runs every suite. */
const ROOT_TEST_GLOB = 'tests/**/*.test.mjs';
/** The CI job that gates the advisory-lock criterion on every PR. */
const CI_JOB = 'database-postgres-lock';
/** The committed advisory-lock key (stable namespace, bound as $1). */
const LOCK_KEY = 'personal-blog-migrations';
// ---------------------------------------------------------------------------
// Static assertions on the committed lock source
// ---------------------------------------------------------------------------
/**
* Asserts the lock module exists with the documented shape: a `MigrationLock`
* class and a stable `MIGRATION_LOCK_KEY` constant. Fails fast on a missing
* module; the mutation probes below prove the assertions are non-vacuous.
*/
function assertLockModule(src) {
assert.match(src, /export class MigrationLock/, 'the lock module must export the MigrationLock class');
assert.match(
src,
/export const MIGRATION_LOCK_KEY = 'personal-blog-migrations'/,
'the lock module must declare the MIGRATION_LOCK_KEY constant',
);
assert.match(
src,
/get isHeld\(\)/,
'the lock module must expose isHeld so callers can tell whether the instance holds the lock',
);
}
/**
* Asserts the blocking path: `acquire()` takes a session-scoped advisory lock
* (`SELECT pg_advisory_lock(hashtextextended($1, 0))`) on a dedicated
* connection (`pool.connect()`), keyed by the stable constant bound as a
* parameter — a second runner calling `acquire()` waits while the first holds
* the lock. The key is never interpolated into the SQL.
*/
function assertBlockingAcquire(src) {
assert.match(
src,
/SELECT pg_advisory_lock\(hashtextextended\(\$1, 0\)\)/,
'acquire() must take the blocking session-scoped advisory lock (SELECT pg_advisory_lock(hashtextextended($1, 0)))',
);
assert.match(src, /pool\.connect\(\)/, 'the lock must be held on a dedicated connection (pool.connect())');
assert.match(src, /\[MIGRATION_LOCK_KEY\]/, 'the lock key must be passed as the bound-parameter array ([MIGRATION_LOCK_KEY])');
assert.doesNotMatch(
src,
/\$\{MIGRATION_LOCK_KEY\}/,
'the lock key must never be interpolated into the SQL',
);
}
/**
* Asserts the non-blocking path: `tryAcquire()` uses
* `SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired` and
* resolves `false` when another runner holds the lock — a second runner fails
* fast while the first holds the lock.
*/
function assertNonBlockingAcquire(src) {
assert.match(
src,
/SELECT pg_try_advisory_lock\(hashtextextended\(\$1, 0\)\) AS acquired/,
'tryAcquire() must use the non-blocking advisory lock (SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired)',
);
assert.match(src, /result\.rows\[0\]\?\.acquired === true/, 'tryAcquire() must resolve true only when the lock was acquired');
assert.match(src, /return false;/, 'tryAcquire() must resolve false when another runner holds the lock');
}
/**
* Asserts the re-entrancy guard: `acquire()` memoizes the in-flight acquire
* (`acquireInFlight`) so a concurrent `acquire()` on the same instance
* returns the same in-flight promise instead of checking out another
* connection — exactly one connection is checked out and no locked
* connection leaks. The memo is cleared once the acquire settles, so later
* acquire() calls behave normally.
*/
function assertReentrantAcquire(src) {
assert.match(
src,
/private acquireInFlight: Promise<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/,
);
});