Skip to content

Commit 2c046fe

Browse files
committed
improvement(db): release advisory lock holders when a lock test fails
1 parent 3984171 commit 2c046fe

1 file changed

Lines changed: 28 additions & 15 deletions

File tree

‎apps/sim/lib/db/advisory-locks.integration.ts‎

Lines changed: 28 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,10 @@ import { acquireAdvisoryXactLock, tryAcquireAdvisoryXactLock } from '@/lib/db/ad
1313
const connection = postgres(readTestDatabaseUrl(), { max: 4, prepare: false, onnotice: () => {} })
1414
const db = drizzle(connection, { schema })
1515

16-
/** Holds `key` in an open transaction until the returned release function is called. */
16+
/**
17+
* Holds `key` in an open transaction until the returned release function is
18+
* called. Rejects if the holder transaction fails before taking the lock.
19+
*/
1720
async function holdLock(tag: string, key: string): Promise<() => Promise<void>> {
1821
let release!: () => void
1922
let acquired!: () => void
@@ -28,7 +31,7 @@ async function holdLock(tag: string, key: string): Promise<() => Promise<void>>
2831
acquired()
2932
await released
3033
})
31-
await held
34+
await Promise.race([held, transaction])
3235
return async () => {
3336
release()
3437
await transaction
@@ -50,9 +53,12 @@ describe('advisory xact locks', () => {
5053
it('blocks a second transaction on the same key until the holder commits', async () => {
5154
const key = `lock-test:${generateId()}`
5255
const release = await holdLock('lock_test', key)
53-
const error = await acquireWithTimeout('lock_test', key, 200).catch((e: unknown) => e)
54-
expect(getPostgresErrorCode(error)).toBe('55P03')
55-
await release()
56+
try {
57+
const error = await acquireWithTimeout('lock_test', key, 200).catch((e: unknown) => e)
58+
expect(getPostgresErrorCode(error)).toBe('55P03')
59+
} finally {
60+
await release()
61+
}
5662
await expect(acquireWithTimeout('lock_test', key, 200)).resolves.toBeUndefined()
5763
})
5864

@@ -70,8 +76,12 @@ describe('advisory xact locks', () => {
7076
it('reports whether the non-blocking variant acquired the lock', async () => {
7177
const key = `lock-test:${generateId()}`
7278
const release = await holdLock('lock_test', key)
73-
const whileHeld = await db.transaction((tx) => tryAcquireAdvisoryXactLock(tx, 'lock_test', key))
74-
await release()
79+
let whileHeld: boolean
80+
try {
81+
whileHeld = await db.transaction((tx) => tryAcquireAdvisoryXactLock(tx, 'lock_test', key))
82+
} finally {
83+
await release()
84+
}
7585
const afterRelease = await db.transaction((tx) =>
7686
tryAcquireAdvisoryXactLock(tx, 'lock_test', key)
7787
)
@@ -83,15 +93,18 @@ describe('advisory xact locks', () => {
8393
const release = await holdLock('lock_test', key)
8494
const waiter = acquireWithTimeout('lock_test_waiter', key, 5_000)
8595
let waitingQuery: string | undefined
86-
for (let attempt = 0; attempt < 50 && !waitingQuery; attempt++) {
87-
const [row] = await connection<{ query: string }[]>`
88-
SELECT query FROM pg_stat_activity
89-
WHERE wait_event_type = 'Lock' AND wait_event = 'advisory' AND query LIKE ${'%lock_test_waiter%'}`
90-
waitingQuery = row?.query
91-
if (!waitingQuery) await sleep(20)
96+
try {
97+
for (let attempt = 0; attempt < 50 && !waitingQuery; attempt++) {
98+
const [row] = await connection<{ query: string }[]>`
99+
SELECT query FROM pg_stat_activity
100+
WHERE wait_event_type = 'Lock' AND wait_event = 'advisory' AND query LIKE ${'%lock_test_waiter%'}`
101+
waitingQuery = row?.query
102+
if (!waitingQuery) await sleep(20)
103+
}
104+
} finally {
105+
await release()
106+
await waiter
92107
}
93-
await release()
94-
await waiter
95108
expect(waitingQuery).toMatch(/pg_advisory_xact_lock\(.*\) \/\*lock='lock_test_waiter'\*\/$/)
96109
})
97110

0 commit comments

Comments
 (0)