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
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:
@@ -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
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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,16 +342,47 @@ try {
|
|||||||
await sleep(100);
|
await sleep(100);
|
||||||
}
|
}
|
||||||
|
|
||||||
console.log('LOCK_PROBE_RESULT ' + JSON.stringify({
|
// Re-entrant-safe acquire (acceptance criterion: concurrent acquire() on
|
||||||
baseline,
|
// the same lock instance is re-entrant-safe — the in-flight acquire is
|
||||||
heldCount,
|
// memoized so exactly one connection is checked out and no locked
|
||||||
secondFailsFast,
|
// connection leaks). A fresh pool so the connection counts are exact: two
|
||||||
observedWait,
|
// concurrent acquire() calls must share one in-flight acquire — one
|
||||||
secondAcquiredAfterRelease,
|
// connection checked out (pool.totalCount stays 1), the lock granted once,
|
||||||
countAfterRelease,
|
// and nothing left held after release.
|
||||||
heldWhileSessionOpen,
|
let reentrantPool = null;
|
||||||
thirdAcquiresAfterSessionEnd,
|
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 {
|
} 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/, '');
|
||||||
|
|||||||
Reference in New Issue
Block a user