[E00-S03-T04] Advisory lock prevents concurrent migration runners #393
@@ -21,9 +21,11 @@
|
||||
* (`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), and released implicitly by the
|
||||
* server when the holding session ends — the issue's rollback note: "the
|
||||
* lock releases when the runner exits".
|
||||
* 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
|
||||
@@ -133,7 +135,12 @@ export class MigrationLock {
|
||||
* 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.
|
||||
* Idempotent: releasing an instance that does not hold the lock is a no-op.
|
||||
* 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;
|
||||
@@ -145,8 +152,15 @@ export class MigrationLock {
|
||||
'SELECT pg_advisory_unlock(hashtextextended($1, 0))',
|
||||
[MIGRATION_LOCK_KEY],
|
||||
);
|
||||
} finally {
|
||||
client.release();
|
||||
} 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();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,13 @@
|
||||
* 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.
|
||||
@@ -169,6 +176,40 @@ function assertSessionScopedRelease(src) {
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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(
|
||||
@@ -388,6 +429,83 @@ try {
|
||||
}
|
||||
`;
|
||||
|
||||
/**
|
||||
* 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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -414,6 +532,10 @@ test('release() unlocks the holding session and returns the connection to the po
|
||||
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(
|
||||
@@ -590,6 +712,68 @@ test('concurrent migration runners are serialized by the advisory lock (real sta
|
||||
}
|
||||
});
|
||||
|
||||
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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -622,6 +806,15 @@ test('removing pg_advisory_unlock makes the session-release criterion fail (muta
|
||||
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/, '');
|
||||
|
||||
Reference in New Issue
Block a user