diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 77e8bcc..22652f1 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -77,6 +77,32 @@ jobs: - name: Run migration ledger test suite run: node --test tests/database-postgres-ledger.test.mjs + # E00-S03-T04: the static assertions of tests/database-postgres-lock.test.mjs + # gate every PR — the suite locks in the migration advisory lock (session- + # scoped pg_advisory_lock/pg_try_advisory_lock over a stable keyed hash on a + # dedicated connection, driver-boundary re-export) with mutation probes, and + # the docker-gated real-stack concurrent probe (a second runner waits or + # fails while the first holds the lock; the lock releases when the holding + # session ends) runs where a Docker daemon is available and skips cleanly + # otherwise. The job installs the frozen workspace because the real-stack + # probe executes the committed lock module from the host (it imports `pg` + # through the package's own links). + database-postgres-lock: + name: Migration advisory lock (E00-S03-T04) + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - name: Install Node.js 24 + uses: actions/setup-node@v4 + with: + node-version: '24' + - name: Enable pnpm (corepack, pinned to 11.23.0 via packageManager) + run: corepack enable + - name: Install dependencies (frozen lockfile) + run: pnpm install --frozen-lockfile + - name: Run migration advisory lock test suite + run: node --test tests/database-postgres-lock.test.mjs + # E00-S03-T01: the static assertions of tests/compose-config.test.mjs (db # image pinned to postgres:18.6-bookworm, health gate, volume persistence, # build platforms) gate every PR (the docker-gated real-stack probes inside diff --git a/tests/database-postgres-lock.test.mjs b/tests/database-postgres-lock.test.mjs new file mode 100644 index 0000000..8170a7b --- /dev/null +++ b/tests/database-postgres-lock.test.mjs @@ -0,0 +1,536 @@ +/** + * Migration advisory lock test — locks in the [E00-S03-T04] advisory lock + * that prevents concurrent migration runners. + * + * Acceptance criteria covered (each test fails without the committed state): + * - "advisory lock prevents concurrent migration runners" → the committed + * `packages/database-postgres/src/lock.ts` defines a `MigrationLock` + * backed by PostgreSQL session-scoped advisory locks keyed by a stable + * constant (`MIGRATION_LOCK_KEY` → `hashtextextended($1, 0)`): `acquire()` + * blocks (`pg_advisory_lock`) and `tryAcquire()` fails fast + * (`pg_try_advisory_lock`), both on a dedicated pool connection; the + * docker-gated real-stack probe starts a real database and proves two + * runners cannot hold the lock at once (the issue's test plan: "run a + * concurrent migration lock test"). + * - "a second runner waits or fails while the first holds the lock" → the + * real-stack probe holds the lock in the first runner, shows the second + * runner's `tryAcquire()` returns false (fails fast) and its blocking + * `acquire()` waits (the server reports `wait_event advisory` for the + * waiting session) until the first runner releases, then acquires; a + * session that takes the lock and then closes (the runner exits) + * releases it — the issue's rollback note ("the lock releases when the + * runner exits"). + * - the lock is part of the driver boundary: `src/index.ts` re-exports + * `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the + * `pg` driver to lock migration state. + * + * Run: `node --test tests/database-postgres-lock.test.mjs` + * (node:test — built into Node >= 18; no dependencies, lockfile untouched.) + */ + +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { readFileSync, existsSync, writeFileSync, rmSync } from 'node:fs'; +import { spawnSync } from 'node:child_process'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; + +const REPO_ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..'); + +const read = (relPath) => readFileSync(path.join(REPO_ROOT, relPath), 'utf8'); + +/** The committed advisory-lock module and the driver boundary that re-exports it. */ +const LOCK_SRC = 'packages/database-postgres/src/lock.ts'; +const INDEX_SRC = 'packages/database-postgres/src/index.ts'; + +/** The root test glob (root `scripts.test`, E00-S01-T12) that runs every suite. */ +const ROOT_TEST_GLOB = 'tests/**/*.test.mjs'; + +/** The CI job that gates the advisory-lock criterion on every PR. */ +const CI_JOB = 'database-postgres-lock'; + +/** The committed advisory-lock key (stable namespace, bound as $1). */ +const LOCK_KEY = 'personal-blog-migrations'; + +// --------------------------------------------------------------------------- +// Static assertions on the committed lock source +// --------------------------------------------------------------------------- + +/** + * Asserts the lock module exists with the documented shape: a `MigrationLock` + * class and a stable `MIGRATION_LOCK_KEY` constant. Fails fast on a missing + * module; the mutation probes below prove the assertions are non-vacuous. + */ +function assertLockModule(src) { + assert.match(src, /export class MigrationLock/, 'the lock module must export the MigrationLock class'); + assert.match( + src, + /export const MIGRATION_LOCK_KEY = 'personal-blog-migrations'/, + 'the lock module must declare the MIGRATION_LOCK_KEY constant', + ); + assert.match( + src, + /get isHeld\(\)/, + 'the lock module must expose isHeld so callers can tell whether the instance holds the lock', + ); +} + +/** + * Asserts the blocking path: `acquire()` takes a session-scoped advisory lock + * (`SELECT pg_advisory_lock(hashtextextended($1, 0))`) on a dedicated + * connection (`pool.connect()`), keyed by the stable constant bound as a + * parameter — a second runner calling `acquire()` waits while the first holds + * the lock. The key is never interpolated into the SQL. + */ +function assertBlockingAcquire(src) { + assert.match( + src, + /SELECT pg_advisory_lock\(hashtextextended\(\$1, 0\)\)/, + 'acquire() must take the blocking session-scoped advisory lock (SELECT pg_advisory_lock(hashtextextended($1, 0)))', + ); + assert.match(src, /pool\.connect\(\)/, 'the lock must be held on a dedicated connection (pool.connect())'); + assert.match(src, /\[MIGRATION_LOCK_KEY\]/, 'the lock key must be passed as the bound-parameter array ([MIGRATION_LOCK_KEY])'); + assert.doesNotMatch( + src, + /\$\{MIGRATION_LOCK_KEY\}/, + 'the lock key must never be interpolated into the SQL', + ); +} + +/** + * Asserts the non-blocking path: `tryAcquire()` uses + * `SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired` and + * resolves `false` when another runner holds the lock — a second runner fails + * fast while the first holds the lock. + */ +function assertNonBlockingAcquire(src) { + assert.match( + src, + /SELECT pg_try_advisory_lock\(hashtextextended\(\$1, 0\)\) AS acquired/, + 'tryAcquire() must use the non-blocking advisory lock (SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired)', + ); + assert.match(src, /result\.rows\[0\]\?\.acquired === true/, 'tryAcquire() must resolve true only when the lock was acquired'); + assert.match(src, /return false;/, 'tryAcquire() must resolve false when another runner holds the lock'); +} + +/** + * Asserts the session-scoped release: `release()` unlocks the holding session + * (`SELECT pg_advisory_unlock(hashtextextended($1, 0))`) BEFORE returning the + * connection to the pool (`client.release()`), so a pooled connection is + * never reused while still holding the lock; releasing an unheld instance is + * a no-op. + */ +function assertSessionScopedRelease(src) { + const releaseBody = src.slice(src.indexOf('async release()')); + assert.match( + releaseBody, + /SELECT pg_advisory_unlock\(hashtextextended\(\$1, 0\)\)/, + 'release() must unlock the holding session (SELECT pg_advisory_unlock(hashtextextended($1, 0)))', + ); + assert.match(releaseBody, /client\.release\(\)/, 'release() must return the connection to the pool (client.release())'); + assert.ok( + releaseBody.indexOf('pg_advisory_unlock') < releaseBody.indexOf('client.release()'), + 'release() must unlock the session before returning the connection to the pool', + ); +} + +/** Asserts the driver boundary re-exports the lock from the package entrypoint. */ +function assertBoundaryReexport(src) { + assert.match( + src, + /export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js'/, + 'the driver boundary must re-export the lock (src/index.ts → ./lock.js)', + ); +} + +// --------------------------------------------------------------------------- +// Docker probe helpers (integration test skips cleanly without Docker) +// --------------------------------------------------------------------------- + +function run(cmd, args, opts = {}) { + return spawnSync(cmd, args, { + encoding: 'utf8', + timeout: 600_000, + ...opts, + }); +} + +/** True when the `docker` CLI with the Compose plugin is on PATH. */ +function dockerComposeAvailable() { + try { + return run('docker', ['compose', 'version'], { timeout: 15_000 }).status === 0; + } catch { + return false; + } +} + +/** True when a reachable Docker daemon exists. */ +function dockerDaemonAvailable() { + try { + return run('docker', ['info'], { timeout: 15_000 }).status === 0; + } catch { + return false; + } +} + +/** Parses `docker compose ps --format json` (JSON array or one object per line). */ +function parsePsJson(stdout) { + const text = String(stdout).trim(); + if (!text) return []; + try { + const parsed = JSON.parse(text); + return Array.isArray(parsed) ? parsed : [parsed]; + } catch { + return text + .split('\n') + .map((line) => line.trim()) + .filter(Boolean) + .map((line) => JSON.parse(line)); + } +} + +/** Tolerant field lookup across compose ps JSON shapes. */ +function field(container, ...names) { + for (const name of names) { + if (container[name] !== undefined) return container[name]; + } + return undefined; +} + +/** + * Polls `docker compose ps` until the db container of the given project + * reports healthy (or the deadline passes), so the probe never races a + * still-booting database. + */ +function waitForDbHealthy(project, deadlineMs = 60_000) { + const deadline = Date.now() + deadlineMs; + let last = ''; + while (Date.now() < deadline) { + const ps = run('docker', ['compose', '-p', project, 'ps', '--format', 'json'], { + cwd: REPO_ROOT, + timeout: 15_000, + }); + if (ps.status === 0) { + last = ps.stdout; + const db = parsePsJson(ps.stdout).find((c) => field(c, 'Service', 'service') === 'db'); + if (db && /healthy/i.test(String(field(db, 'Health', 'health') ?? ''))) return; + } + run(process.execPath, ['-e', 'setTimeout(() => {}, 1000)']); // db still booting — retry + } + throw new Error(`the db container did not become healthy within ${deadlineMs}ms (last ps: "${last.trim()}")`); +} + +/** + * The host-side probe body: exercises the committed lock module against a + * real database (DATABASE_URL) exactly as concurrent migration runners will — + * first runner holds the lock, second runner fails fast AND waits, acquires + * after release, and a session that closes (the runner exits) releases the + * lock. Written to a temp file inside `packages/database-postgres/` so `pg` + * resolves through the package's own dependency links, then removed. + */ +const PROBE_SOURCE = ` +import { Pool, Client } from 'pg'; +import { MigrationLock, MIGRATION_LOCK_KEY } from './src/lock.ts'; + +const pool = new Pool({ connectionString: process.env.DATABASE_URL }); + +const advisoryCount = async () => { + const r = await pool.query( + "SELECT count(*)::int AS n FROM pg_locks WHERE locktype = 'advisory' AND granted", + ); + return r.rows[0].n; +}; + +const advisoryWaitEvents = async () => { + const r = await pool.query( + \`SELECT wait_event_type, wait_event FROM pg_stat_activity + WHERE datname = current_database() AND pid <> pg_backend_pid()\`, + ); + return r.rows + .filter((row) => row.wait_event_type !== null && row.wait_event !== null) + .map((row) => row.wait_event_type + ':' + row.wait_event); +}; + +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + +try { + const baseline = await advisoryCount(); + + // First runner takes the lock: the server must hold one granted advisory + // lock more than before. + const lockA = new MigrationLock(pool); + await lockA.acquire(); + const heldCount = await advisoryCount(); + + // A second runner fails fast while the first holds the lock. + const lockB = new MigrationLock(pool); + const secondFailsFast = await lockB.tryAcquire(); + + // A second runner waits while the first holds the lock: the blocking + // acquire() must not resolve, and the server must report the waiting + // session (wait_event 'advisory'). + const waitPromise = lockB.acquire(); + let observedWait = false; + const waitDeadline = Date.now() + 15_000; + while (Date.now() < waitDeadline) { + if ((await advisoryWaitEvents()).includes('Lock:advisory')) { + observedWait = true; + break; + } + await sleep(100); + } + + // First runner finishes -> the waiting runner acquires, then releases. + await lockA.release(); + await waitPromise; + const secondAcquiredAfterRelease = lockB.isHeld; + await lockB.release(); + const countAfterRelease = await advisoryCount(); + + // Session exit releases the lock (rollback note: "the lock releases when + // the runner exits"): a dedicated connection takes the same advisory lock + // and then closes; the server must release the lock with the session. + const exitClient = new Client({ connectionString: process.env.DATABASE_URL }); + await exitClient.connect(); + await exitClient.query('SELECT pg_advisory_lock(hashtextextended($1, 0))', [MIGRATION_LOCK_KEY]); + const heldWhileSessionOpen = await advisoryCount(); + await exitClient.end(); // the runner exits -> session closes -> lock releases + let thirdAcquiresAfterSessionEnd = false; + const retryDeadline = Date.now() + 5_000; + while (Date.now() < retryDeadline) { + const lockC = new MigrationLock(pool); + if (await lockC.tryAcquire()) { + thirdAcquiresAfterSessionEnd = true; + await lockC.release(); + break; + } + await sleep(100); + } + + console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ + baseline, + heldCount, + secondFailsFast, + observedWait, + secondAcquiredAfterRelease, + countAfterRelease, + heldWhileSessionOpen, + thirdAcquiresAfterSessionEnd, + })); +} finally { + await pool.end(); +} +`; + +// --------------------------------------------------------------------------- +// Criterion tests +// --------------------------------------------------------------------------- + +test('the lock module exists in the driver-owner package and defines the migration lock key', () => { + assert.ok(existsSync(path.join(REPO_ROOT, LOCK_SRC)), `committed ${LOCK_SRC} must exist`); + const src = read(LOCK_SRC); + assertLockModule(src); +}); + +test('acquire() takes the blocking session-scoped advisory lock keyed by a bound parameter (a second runner waits)', () => { + assertBlockingAcquire(read(LOCK_SRC)); +}); + +test('tryAcquire() fails fast while another runner holds the lock (a second runner fails)', () => { + assertNonBlockingAcquire(read(LOCK_SRC)); +}); + +test('release() unlocks the holding session and returns the connection to the pool', () => { + assertSessionScopedRelease(read(LOCK_SRC)); +}); + +test('the driver boundary re-exports the lock from the package entrypoint', () => { + const src = read(INDEX_SRC); + assert.match( + src, + /from '\.\/lock\.js'/, + 'the driver boundary must import the lock module (from \'./lock.js\')', + ); + assertBoundaryReexport(src); +}); + +test('the advisory-lock criterion is enforced in CI', () => { + // Picked up by the root test command (root `scripts.test` glob). + const scripts = JSON.parse(read('package.json')).scripts ?? {}; + assert.equal( + scripts.test, + `node --test "${ROOT_TEST_GLOB}"`, + `root scripts.test must run the "${ROOT_TEST_GLOB}" glob so this suite runs with the rest`, + ); + // And a dedicated CI job gates it on every PR. + const workflow = read('.gitea/workflows/ci.yml'); + assert.ok( + workflow.includes(`node --test tests/database-postgres-lock.test.mjs`), + `CI must run the advisory-lock suite (job "${CI_JOB}") on every PR`, + ); +}); + +// --------------------------------------------------------------------------- +// Real-stack probe — the issue's test plan: "run a concurrent migration +// lock test" against a real database +// --------------------------------------------------------------------------- + +const DOCKER_COMPOSE = dockerComposeAvailable(); +const DOCKER_DAEMON = dockerDaemonAvailable(); + +/** True when this Node can execute the committed `.ts` lock module (>= 23.6, type stripping). */ +const TS_STRIPPING = (() => { + const [major, minor] = process.versions.node.split('.').map(Number); + return major > 23 || (major === 23 && minor >= 6); +})(); + +// An isolated compose project + non-default host port so this probe never +// collides with the compose-config suite's default-project containers, the +// ledger probe's project/port, or the default 5432 binding when several +// suites run on the same host. +const COMPOSE_PROJECT = 'eppp-lock-probe'; +const POSTGRES_HOST_PORT = '55433'; +const DATABASE_URL = `postgres://eppp:eppp@127.0.0.1:${POSTGRES_HOST_PORT}/eppp`; + +test('concurrent migration runners are serialized by the advisory lock (real stack)', { skip: !DOCKER_COMPOSE || !DOCKER_DAEMON || !TS_STRIPPING }, () => { + // The issue's test plan: "run a concurrent migration lock test". The probe + // starts the committed compose `db` service (its own project + host port), + // runs the committed lock module against it with two runner instances, and + // asserts the first runner holds the lock, the second fails fast AND waits + // until the first releases (server-reported wait_event 'advisory'), and a + // session that closes — the runner exits — releases the lock (rollback + // note). Rollback: `docker compose -p eppp-lock-probe down`. + const composeEnv = { ...process.env, POSTGRES_PORT: POSTGRES_HOST_PORT }; + const compose = (args, opts = {}) => + run('docker', ['compose', '-p', COMPOSE_PROJECT, ...args], { cwd: REPO_ROOT, env: composeEnv, ...opts }); + + const probeFile = path.join(REPO_ROOT, 'packages/database-postgres', `.lock-probe-${process.pid}.mjs`); + + try { + const up = compose(['up', '-d', 'db'], { timeout: 180_000 }); + assert.equal( + up.status, + 0, + `"docker compose up -d db" must exit 0:\n${(up.stdout || '')}\n${(up.stderr || '')}`.trim(), + ); + waitForDbHealthy(COMPOSE_PROJECT); + + // Run the committed lock module against the real database. + writeFileSync(probeFile, PROBE_SOURCE); + const probe = run(process.execPath, [probeFile], { + cwd: REPO_ROOT, + env: { ...process.env, DATABASE_URL }, + timeout: 60_000, + }); + assert.equal( + probe.status, + 0, + `the lock probe must exit 0:\n${(probe.stdout || '')}\n${(probe.stderr || '')}`.trim(), + ); + const resultLine = (probe.stdout || '') + .split('\n') + .map((l) => l.trim()) + .find((l) => l.startsWith('LOCK_PROBE_RESULT')); + assert.ok(resultLine, `the lock probe must report a result line (stdout: "${(probe.stdout || '').trim()}")`); + const result = JSON.parse(resultLine.slice('LOCK_PROBE_RESULT'.length).trim()); + + // "advisory lock prevents concurrent migration runners": the first runner + // holds exactly one granted advisory lock more than the empty baseline. + assert.equal( + result.heldCount, + result.baseline + 1, + `acquire() must hold one granted advisory lock (baseline ${result.baseline}, got ${result.heldCount})`, + ); + + // "a second runner fails while the first holds the lock". + assert.equal( + result.secondFailsFast, + false, + 'tryAcquire() must fail fast while the first runner holds the lock', + ); + + // "a second runner waits while the first holds the lock": the blocking + // acquire() does not resolve and the server reports the waiting session. + assert.equal( + result.observedWait, + true, + 'the blocking acquire() must wait (server-reported wait_event advisory) while the first runner holds the lock', + ); + + // The waiting runner acquires once the first runner releases. + assert.equal( + result.secondAcquiredAfterRelease, + true, + 'the waiting runner must acquire the lock once the first runner releases it', + ); + + // Nothing is left held after both runners release. + assert.equal( + result.countAfterRelease, + result.baseline, + `no advisory lock may remain granted after both runners release (got ${result.countAfterRelease}, expected ${result.baseline})`, + ); + + // Rollback note: "the lock releases when the runner exits" — a session + // that holds the lock and then closes releases it. + assert.equal( + result.heldWhileSessionOpen, + result.baseline + 1, + `the dedicated session must hold the advisory lock while open (baseline ${result.baseline}, got ${result.heldWhileSessionOpen})`, + ); + assert.equal( + result.thirdAcquiresAfterSessionEnd, + true, + 'the lock must be acquirable again after the holding session closes (the runner exits)', + ); + } finally { + compose(['down', '-v'], { timeout: 120_000 }); + rmSync(probeFile, { force: true }); + } +}); + +// --------------------------------------------------------------------------- +// Non-vacuous probes — the assertions above really do fail on violations +// --------------------------------------------------------------------------- + +test('removing pg_advisory_lock makes the blocking-acquire criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutBlocking = src.replace(/SELECT pg_advisory_lock\(/, 'SELECT pg_advisory_xact_lock('); + assert.notEqual(withoutBlocking, src, 'the mutation must actually replace the blocking lock call'); + assert.throws(() => assertBlockingAcquire(withoutBlocking), /pg_advisory_lock/); +}); + +test('removing pg_try_advisory_lock makes the non-blocking criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutTry = src.replace(/SELECT pg_try_advisory_lock\(/, 'SELECT pg_advisory_lock('); + assert.notEqual(withoutTry, src, 'the mutation must actually replace the non-blocking lock call'); + assert.throws(() => assertNonBlockingAcquire(withoutTry), /pg_try_advisory_lock/); +}); + +test('interpolating the lock key into the SQL makes the keyed-parameter criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const interpolated = src.replace(/hashtextextended\(\$1, 0\)/g, "hashtextextended('${MIGRATION_LOCK_KEY}', 0)"); + assert.notEqual(interpolated, src, 'the mutation must actually replace the $1 placeholder'); + assert.throws(() => assertBlockingAcquire(interpolated), /\$1/); +}); + +test('removing pg_advisory_unlock makes the session-release criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutUnlock = src.replace(/SELECT pg_advisory_unlock\(/, 'SELECT pg_advisory_lock('); + assert.notEqual(withoutUnlock, src, 'the mutation must actually replace the unlock call'); + assert.throws(() => assertSessionScopedRelease(withoutUnlock), /pg_advisory_unlock/); +}); + +test('dropping the lock re-export from the boundary fails the boundary criterion (mutation probe)', () => { + const src = read(INDEX_SRC); + const withoutReexport = src.replace(/export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js';\n/, ''); + assert.notEqual(withoutReexport, src, 'the mutation must actually remove the lock re-export'); + assert.throws(() => assertBoundaryReexport(withoutReexport), /re-export/); +}); + +test('a placeholder lock module fails the advisory-lock criterion (mutation probe)', () => { + assert.throws( + () => assertLockModule('export class MigrationLock {}\n'), + /MIGRATION_LOCK_KEY/, + ); +});