@@ -111,11 +111,18 @@ afterAll(async () => {
111111 await db . delete ( user ) . where ( eq ( user . id , ids . owner ) )
112112} )
113113
114- /** Whether a session is waiting on an advisory lock — this file's database has no other traffic. */
115- async function hasAdvisoryLockWaiter ( ) {
116- const rows = await db . execute < { waiting : boolean } > (
117- sql `SELECT EXISTS (SELECT 1 FROM pg_locks WHERE locktype = 'advisory' AND NOT granted) AS waiting`
118- )
114+ /**
115+ * Whether a session waits on this execution's ledger lock. A bigint advisory key is stored
116+ * split across `classid` (high 32 bits) and `objid` (low 32 bits) with `objsubid = 1`.
117+ */
118+ async function isLedgerLockAwaited ( executionId : string ) {
119+ const rows = await db . execute < { waiting : boolean } > ( sql `
120+ SELECT EXISTS (
121+ SELECT 1 FROM pg_locks
122+ WHERE locktype = 'advisory' AND NOT granted AND objsubid = 1
123+ AND ((classid::bigint << 32) | objid::bigint) = hashtextextended(${ executionId } , 0)
124+ ) AS waiting
125+ ` )
119126 return Boolean ( rows [ 0 ] ?. waiting )
120127}
121128
@@ -150,11 +157,19 @@ describe('completeWorkflowExecution', () => {
150157 billingAttribution,
151158 } )
152159
160+ /** A completion that settles before blocking surfaces its own outcome instead of a timeout. */
161+ const settledWithoutBlocking = completion . then ( ( ) => {
162+ throw new Error ( 'Completion finished without waiting on the ledger lock' )
163+ } )
164+
153165 let statusWhileLedgerBlocked : string | undefined
154166 try {
155- await vi . waitFor ( async ( ) => {
156- expect ( await hasAdvisoryLockWaiter ( ) ) . toBe ( true )
157- } )
167+ await Promise . race ( [
168+ vi . waitFor ( async ( ) => {
169+ expect ( await isLedgerLockAwaited ( executionId ) ) . toBe ( true )
170+ } ) ,
171+ settledWithoutBlocking ,
172+ ] )
158173 statusWhileLedgerBlocked = ( await logRow ( executionId ) ) ?. status
159174 } finally {
160175 releaseLock . resolve ( )
0 commit comments