diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 22652f1..9d4a05e 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -80,13 +80,15 @@ jobs: # E00-S03-T04: the static assertions of tests/database-postgres-lock.test.mjs # gate every PR — the suite locks in the migration advisory lock (session- # scoped pg_advisory_lock/pg_try_advisory_lock over a stable keyed hash on a - # dedicated connection, driver-boundary re-export) with mutation probes, and - # the docker-gated real-stack concurrent probe (a second runner waits or - # fails while the first holds the lock; the lock releases when the holding - # session ends) runs where a Docker daemon is available and skips cleanly - # otherwise. The job installs the frozen workspace because the real-stack - # probe executes the committed lock module from the host (it imports `pg` - # through the package's own links). + # dedicated connection, re-entrant-safe in-flight acquire so concurrent + # acquire() calls share one connection, driver-boundary re-export) with + # mutation probes, and the docker-gated real-stack concurrent probe (a + # second runner waits or fails while the first holds the lock; concurrent + # acquire() checks out exactly one connection; the lock releases when the + # holding session ends) runs where a Docker daemon is available and skips + # cleanly otherwise. The job installs the frozen workspace because the + # real-stack probe executes the committed lock module from the host (it + # imports `pg` through the package's own links). database-postgres-lock: name: Migration advisory lock (E00-S03-T04) runs-on: ubuntu-latest diff --git a/packages/database-postgres/src/lock.ts b/packages/database-postgres/src/lock.ts index 88a3def..413b352 100644 --- a/packages/database-postgres/src/lock.ts +++ b/packages/database-postgres/src/lock.ts @@ -47,6 +47,8 @@ export class MigrationLock { private readonly pool: Pool; private client: PoolClient | null = null; private held = false; + /** In-flight acquire, memoized so concurrent acquire() calls share one connection (re-entrant-safe). */ + private acquireInFlight: Promise | null = null; /** @param pool The package-owned PostgreSQL pool (`pg.Pool`). */ constructor(pool: Pool) { @@ -61,12 +63,30 @@ export class MigrationLock { /** * Acquires the migration advisory lock, blocking until it is free — a * second runner waits here while the first holds the lock. Idempotent: - * acquiring an already-held instance is a no-op. The lock is held on a - * dedicated connection until `release()` (or until the session ends — e.g. - * the runner exits and the pool closes its connections). + * acquiring an already-held instance is a no-op. Re-entrant-safe: + * concurrent `acquire()` calls on the same instance share the single + * in-flight acquire (memoized in `acquireInFlight`), so exactly one + * connection is checked out and no locked connection leaks. The lock is + * held on a dedicated connection until `release()` (or until the session + * ends — e.g. the runner exits and the pool closes its connections). */ async acquire(): Promise { if (this.held) return; + // A second concurrent acquire() on this instance returns the in-flight + // acquire instead of checking out another connection: exactly one + // connection is checked out and no locked connection leaks. + if (this.acquireInFlight !== null) return this.acquireInFlight; + const inFlight = this.doAcquire(); + this.acquireInFlight = inFlight; + try { + await inFlight; + } finally { + this.acquireInFlight = null; + } + } + + /** The memoized acquire body: checks out one dedicated connection and takes the session-scoped lock on it. */ + private async doAcquire(): Promise { const client = await this.pool.connect(); try { await client.query( diff --git a/tests/database-postgres-lock.test.mjs b/tests/database-postgres-lock.test.mjs index 8170a7b..1b0dabe 100644 --- a/tests/database-postgres-lock.test.mjs +++ b/tests/database-postgres-lock.test.mjs @@ -20,6 +20,14 @@ * session that takes the lock and then closes (the runner exits) * releases it — the issue's rollback note ("the lock releases when the * runner exits"). + * - "concurrent acquire() calls on the same lock instance are + * re-entrant-safe" → `acquire()` memoizes the in-flight acquire + * (`acquireInFlight`): a second concurrent `acquire()` on the same + * instance returns the in-flight acquire instead of checking out another + * connection, so exactly one connection is checked out and no locked + * connection leaks; the real-stack probe proves it behaviorally (two + * concurrent acquire() calls → `pool.totalCount` stays 1, one granted + * advisory lock, none left after release). * - the lock is part of the driver boundary: `src/index.ts` re-exports * `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the * `pg` driver to lock migration state. @@ -113,6 +121,33 @@ function assertNonBlockingAcquire(src) { assert.match(src, /return false;/, 'tryAcquire() must resolve false when another runner holds the lock'); } +/** + * Asserts the re-entrancy guard: `acquire()` memoizes the in-flight acquire + * (`acquireInFlight`) so a concurrent `acquire()` on the same instance + * returns the same in-flight promise instead of checking out another + * connection — exactly one connection is checked out and no locked + * connection leaks. The memo is cleared once the acquire settles, so later + * acquire() calls behave normally. + */ +function assertReentrantAcquire(src) { + assert.match( + src, + /private acquireInFlight: Promise \| null = null;/, + 'the lock module must memoize the in-flight acquire (acquireInFlight field)', + ); + const acquireBody = src.slice(src.indexOf('async acquire()'), src.indexOf('async tryAcquire()')); + assert.match( + acquireBody, + /if \(this\.acquireInFlight !== null\) return this\.acquireInFlight;/, + 'a concurrent acquire() on the same instance must return the in-flight acquire (memoized) instead of checking out another connection', + ); + assert.match( + acquireBody, + /this\.acquireInFlight = null;/, + 'the in-flight acquire memo must be cleared once the acquire settles', + ); +} + /** * Asserts the session-scoped release: `release()` unlocks the holding session * (`SELECT pg_advisory_unlock(hashtextextended($1, 0))`) BEFORE returning the @@ -307,16 +342,47 @@ try { await sleep(100); } - console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ - baseline, - heldCount, - secondFailsFast, - observedWait, - secondAcquiredAfterRelease, - countAfterRelease, - heldWhileSessionOpen, - thirdAcquiresAfterSessionEnd, - })); + // Re-entrant-safe acquire (acceptance criterion: concurrent acquire() on + // the same lock instance is re-entrant-safe — the in-flight acquire is + // memoized so exactly one connection is checked out and no locked + // connection leaks). A fresh pool so the connection counts are exact: two + // concurrent acquire() calls must share one in-flight acquire — one + // connection checked out (pool.totalCount stays 1), the lock granted once, + // and nothing left held after release. + let reentrantPool = null; + try { + reentrantPool = new Pool({ connectionString: process.env.DATABASE_URL }); + const lockR = new MigrationLock(reentrantPool); + await Promise.all([lockR.acquire(), lockR.acquire()]); + const reentrantPoolTotal = reentrantPool.totalCount; + const reentrantPoolIdle = reentrantPool.idleCount; + const reentrantHeld = lockR.isHeld; + const reentrantHeldCount = await advisoryCount(); + await lockR.release(); + const reentrantCountAfterRelease = await advisoryCount(); + await reentrantPool.end(); + reentrantPool = null; + + console.log('LOCK_PROBE_RESULT ' + JSON.stringify({ + baseline, + heldCount, + secondFailsFast, + observedWait, + secondAcquiredAfterRelease, + countAfterRelease, + heldWhileSessionOpen, + thirdAcquiresAfterSessionEnd, + reentrantPoolTotal, + reentrantPoolIdle, + reentrantHeld, + reentrantHeldCount, + reentrantCountAfterRelease, + })); + } finally { + if (reentrantPool !== null) { + await reentrantPool.end().catch(() => {}); + } + } } finally { await pool.end(); } @@ -340,6 +406,10 @@ test('tryAcquire() fails fast while another runner holds the lock (a second runn assertNonBlockingAcquire(read(LOCK_SRC)); }); +test('concurrent acquire() on the same instance is re-entrant-safe (in-flight acquire memoized, exactly one connection)', () => { + assertReentrantAcquire(read(LOCK_SRC)); +}); + test('release() unlocks the holding session and returns the connection to the pool', () => { assertSessionScopedRelease(read(LOCK_SRC)); }); @@ -483,6 +553,37 @@ test('concurrent migration runners are serialized by the advisory lock (real sta true, 'the lock must be acquirable again after the holding session closes (the runner exits)', ); + + // "concurrent acquire() on the same lock instance is re-entrant-safe": + // two concurrent acquire() calls share one in-flight acquire — exactly + // one connection is checked out (pool.totalCount stays 1), the lock is + // granted once, and nothing is left held after release (no locked + // connection leaks). + assert.equal( + result.reentrantPoolTotal, + 1, + `two concurrent acquire() calls must check out exactly one connection (pool.totalCount ${result.reentrantPoolTotal})`, + ); + assert.equal( + result.reentrantPoolIdle, + 0, + `the re-entrant acquire must hold its connection checked out while the lock is held (pool.idleCount ${result.reentrantPoolIdle})`, + ); + assert.equal( + result.reentrantHeld, + true, + 'the re-entrant acquire must leave the instance holding the lock', + ); + assert.equal( + result.reentrantHeldCount, + result.baseline + 1, + `two concurrent acquire() calls must grant the advisory lock once (baseline ${result.baseline}, got ${result.reentrantHeldCount})`, + ); + assert.equal( + result.reentrantCountAfterRelease, + result.baseline, + `no advisory lock may remain granted after the re-entrant instance releases (got ${result.reentrantCountAfterRelease})`, + ); } finally { compose(['down', '-v'], { timeout: 120_000 }); rmSync(probeFile, { force: true }); @@ -521,6 +622,13 @@ test('removing pg_advisory_unlock makes the session-release criterion fail (muta assert.throws(() => assertSessionScopedRelease(withoutUnlock), /pg_advisory_unlock/); }); +test('removing the in-flight acquire memo makes the re-entrancy criterion fail (mutation probe)', () => { + const src = read(LOCK_SRC); + const withoutMemo = src.replace(/if \(this\.acquireInFlight !== null\) return this\.acquireInFlight;\n/, ''); + assert.notEqual(withoutMemo, src, 'the mutation must actually remove the in-flight acquire guard'); + assert.throws(() => assertReentrantAcquire(withoutMemo), /in-flight acquire/); +}); + test('dropping the lock re-export from the boundary fails the boundary criterion (mutation probe)', () => { const src = read(INDEX_SRC); const withoutReexport = src.replace(/export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js';\n/, '');