fix: memoize in-flight acquire so concurrent acquire() is re-entrant-safe (E00-S03-T04)
CI / Frozen lockfile install (pull_request) Successful in 49s
CI / Secrets not embedded (E00-S02-T08) (pull_request) Successful in 26s
CI / Database-postgres import isolation (E00-S03-T02) (pull_request) Successful in 26s
CI / Migration ledger (E00-S03-T03) (pull_request) Successful in 45s
CI / Migration advisory lock (E00-S03-T04) (pull_request) Successful in 50s
CI / Compose config (E00-S03-T01) (pull_request) Successful in 25s

Security-review finding F2: two concurrent acquire() calls on the same
MigrationLock instance could each check out a connection; the second
pg_advisory_lock would overwrite this.client, leaking the first locked
connection until session end.

acquire() now memoizes the in-flight acquire in acquireInFlight and
returns it on re-entry, so exactly one connection is checked out and no
locked connection leaks. The memo is cleared once the acquire settles.
tryAcquire()/release() paths unchanged.

Locked in by:
- static criterion test: acquireInFlight field, re-entry guard returns
  the in-flight acquire, memo cleared on settle
- mutation probe: removing the re-entry guard fails the criterion
- real-stack probe: two concurrent acquire() calls on one instance leave
  pool.totalCount at 1 (exactly one connection), the lock granted once,
  nothing left after release; probe fails (hangs) on the pre-fix code
- CI job comment updated to reflect the re-entrancy criterion
This commit is contained in:
implementer
2026-08-30 00:29:20 +00:00
parent a2b98439e9
commit dac3679e33
3 changed files with 150 additions and 20 deletions
+9 -7
View File
@@ -80,13 +80,15 @@ jobs:
# E00-S03-T04: the static assertions of tests/database-postgres-lock.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- # 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 # 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 # dedicated connection, re-entrant-safe in-flight acquire so concurrent
# the docker-gated real-stack concurrent probe (a second runner waits or # acquire() calls share one connection, driver-boundary re-export) with
# fails while the first holds the lock; the lock releases when the holding # mutation probes, and the docker-gated real-stack concurrent probe (a
# session ends) runs where a Docker daemon is available and skips cleanly # second runner waits or fails while the first holds the lock; concurrent
# otherwise. The job installs the frozen workspace because the real-stack # acquire() checks out exactly one connection; the lock releases when the
# probe executes the committed lock module from the host (it imports `pg` # holding session ends) runs where a Docker daemon is available and skips
# through the package's own links). # 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: database-postgres-lock:
name: Migration advisory lock (E00-S03-T04) name: Migration advisory lock (E00-S03-T04)
runs-on: ubuntu-latest runs-on: ubuntu-latest
+23 -3
View File
@@ -47,6 +47,8 @@ export class MigrationLock {
private readonly pool: Pool; private readonly pool: Pool;
private client: PoolClient | null = null; private client: PoolClient | null = null;
private held = false; 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`). */ /** @param pool The package-owned PostgreSQL pool (`pg.Pool`). */
constructor(pool: Pool) { constructor(pool: Pool) {
@@ -61,12 +63,30 @@ export class MigrationLock {
/** /**
* Acquires the migration advisory lock, blocking until it is free — a * Acquires the migration advisory lock, blocking until it is free — a
* second runner waits here while the first holds the lock. Idempotent: * 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 * acquiring an already-held instance is a no-op. Re-entrant-safe:
* dedicated connection until `release()` (or until the session ends — e.g. * concurrent `acquire()` calls on the same instance share the single
* the runner exits and the pool closes its connections). * 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> { async acquire(): Promise<void> {
if (this.held) return; 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(); const client = await this.pool.connect();
try { try {
await client.query( await client.query(
+108
View File
@@ -20,6 +20,14 @@
* session that takes the lock and then closes (the runner exits) * session that takes the lock and then closes (the runner exits)
* releases it — the issue's rollback note ("the lock releases when the * releases it — the issue's rollback note ("the lock releases when the
* runner exits"). * 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 * - the lock is part of the driver boundary: `src/index.ts` re-exports
* `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the * `MigrationLock` + `MIGRATION_LOCK_KEY`, so no other package needs the
* `pg` driver to lock migration state. * `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'); 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 * Asserts the session-scoped release: `release()` unlocks the holding session
* (`SELECT pg_advisory_unlock(hashtextextended($1, 0))`) BEFORE returning the * (`SELECT pg_advisory_unlock(hashtextextended($1, 0))`) BEFORE returning the
@@ -307,6 +342,27 @@ try {
await sleep(100); 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({ console.log('LOCK_PROBE_RESULT ' + JSON.stringify({
baseline, baseline,
heldCount, heldCount,
@@ -316,7 +372,17 @@ try {
countAfterRelease, countAfterRelease,
heldWhileSessionOpen, heldWhileSessionOpen,
thirdAcquiresAfterSessionEnd, thirdAcquiresAfterSessionEnd,
reentrantPoolTotal,
reentrantPoolIdle,
reentrantHeld,
reentrantHeldCount,
reentrantCountAfterRelease,
})); }));
} finally {
if (reentrantPool !== null) {
await reentrantPool.end().catch(() => {});
}
}
} finally { } finally {
await pool.end(); await pool.end();
} }
@@ -340,6 +406,10 @@ test('tryAcquire() fails fast while another runner holds the lock (a second runn
assertNonBlockingAcquire(read(LOCK_SRC)); 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', () => { test('release() unlocks the holding session and returns the connection to the pool', () => {
assertSessionScopedRelease(read(LOCK_SRC)); assertSessionScopedRelease(read(LOCK_SRC));
}); });
@@ -483,6 +553,37 @@ test('concurrent migration runners are serialized by the advisory lock (real sta
true, true,
'the lock must be acquirable again after the holding session closes (the runner exits)', '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 { } finally {
compose(['down', '-v'], { timeout: 120_000 }); compose(['down', '-v'], { timeout: 120_000 });
rmSync(probeFile, { force: true }); 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/); 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)', () => { test('dropping the lock re-export from the boundary fails the boundary criterion (mutation probe)', () => {
const src = read(INDEX_SRC); const src = read(INDEX_SRC);
const withoutReexport = src.replace(/export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js';\n/, ''); const withoutReexport = src.replace(/export \{ MigrationLock, MIGRATION_LOCK_KEY \} from '\.\/lock\.js';\n/, '');