11/**
2- * Completion ordering against real PostgreSQL: a run's usage ledger is durable before
3- * its log reads terminal, so a reader that sees a finished run always sees its cost.
2+ * Completion ordering against real PostgreSQL: a run's usage ledger is written before its
3+ * log reads terminal, so a reader that sees a finished run always sees its cost.
44 */
55import { db } from '@sim/db'
66import {
@@ -11,10 +11,12 @@ import {
1111 workflowExecutionSnapshots ,
1212 workspace ,
1313} from '@sim/db/schema'
14+ import { createDeferred } from '@sim/testing'
1415import { generateId } from '@sim/utils/id'
15- import { eq } from 'drizzle-orm'
16- import { afterAll , beforeAll , describe , expect , it } from 'vitest'
16+ import { eq , sql } from 'drizzle-orm'
17+ import { afterAll , beforeAll , describe , expect , it , vi } from 'vitest'
1718import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
19+ import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
1820import { buildCostLedger } from '@/lib/logs/cost-ledger'
1921import { executionLogger } from '@/lib/logs/execution/logger'
2022import { calculateCostSummary } from '@/lib/logs/execution/logging-factory'
@@ -26,8 +28,8 @@ const ids = {
2628 workflow : generateId ( ) ,
2729}
2830
29- /** Enough completions that an ordering gap is observed on every run where one exists . */
30- const COMPLETIONS = 20
31+ /** The advisory lock the usage ledger write takes for its execution before inserting . */
32+ const USAGE_RECONCILE_LOCK = 'execution_usage_reconcile'
3133const EXECUTION_FEE = 0.005
3234
3335const workflowState : WorkflowState = {
@@ -109,59 +111,60 @@ afterAll(async () => {
109111 await db . delete ( user ) . where ( eq ( user . id , ids . owner ) )
110112} )
111113
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+ )
119+ return Boolean ( rows [ 0 ] ?. waiting )
120+ }
121+
112122describe ( 'completeWorkflowExecution' , ( ) => {
113- it ( 'never exposes a finished run before its cost ledger ' , async ( ) => {
123+ it ( 'writes the cost ledger before the run reads finished ' , async ( ) => {
114124 const billingAttribution = await resolveBillingAttribution ( {
115125 actorUserId : ids . owner ,
116126 workspaceId : ids . workspace ,
117127 } )
118- let finishedWithoutLedger = 0
128+ const executionId = generateId ( )
129+ await startExecution ( executionId )
119130
120- for ( let i = 0 ; i < COMPLETIONS ; i ++ ) {
121- const executionId = generateId ( )
122- await startExecution ( executionId )
123-
124- /**
125- * The ledger only grows, so the first read that finds the run finished is the one that
126- * counts. Settlement is captured before each read, so a read already in flight when the
127- * completion lands cannot end the loop before the finished run is observed.
128- */
129- let completing = true
130- const reader = ( async ( ) => {
131- for ( ; ; ) {
132- const settledBeforeRead = ! completing
133- if ( ( await logRow ( executionId ) ) ?. status === 'completed' ) {
134- if ( ( await buildCostLedger ( executionId ) ) === null ) finishedWithoutLedger ++
135- return
136- }
137- if ( settledBeforeRead ) return
138- }
139- } ) ( )
131+ /** Holds the ledger write at its lock, so the log can be read while it waits there. */
132+ const lockHeld = createDeferred < void > ( )
133+ const releaseLock = createDeferred < void > ( )
134+ const holder = db . transaction ( async ( tx ) => {
135+ await acquireAdvisoryXactLock ( tx , USAGE_RECONCILE_LOCK , executionId )
136+ lockHeld . resolve ( )
137+ await releaseLock . promise
138+ } )
139+ await lockHeld . promise
140140
141- try {
142- await executionLogger . completeWorkflowExecution ( {
143- executionId,
144- endedAt : new Date ( ) . toISOString ( ) ,
145- totalDurationMs : 5 ,
146- costSummary : calculateCostSummary ( [ ] , { baseExecutionCharge : EXECUTION_FEE } ) ,
147- finalOutput : { } ,
148- traceSpans : [ ] ,
149- status : 'completed' ,
150- actorUserId : ids . owner ,
151- billingAttribution,
152- } )
153- } finally {
154- completing = false
155- // Settle the reader without letting its error replace a completion failure.
156- await reader . catch ( ( ) => { } )
157- }
158- await reader
141+ const completion = executionLogger . completeWorkflowExecution ( {
142+ executionId,
143+ endedAt : new Date ( ) . toISOString ( ) ,
144+ totalDurationMs : 5 ,
145+ costSummary : calculateCostSummary ( [ ] , { baseExecutionCharge : EXECUTION_FEE } ) ,
146+ finalOutput : { } ,
147+ traceSpans : [ ] ,
148+ status : 'completed' ,
149+ actorUserId : ids . owner ,
150+ billingAttribution,
151+ } )
159152
160- const ledger = await buildCostLedger ( executionId )
161- expect ( ledger ?. total ) . toBeCloseTo ( EXECUTION_FEE , 8 )
162- expect ( Number ( ( await logRow ( executionId ) ) ?. costTotal ) ) . toBeCloseTo ( EXECUTION_FEE , 8 )
153+ let statusWhileLedgerBlocked : string | undefined
154+ try {
155+ await vi . waitFor ( async ( ) => {
156+ expect ( await hasAdvisoryLockWaiter ( ) ) . toBe ( true )
157+ } )
158+ statusWhileLedgerBlocked = ( await logRow ( executionId ) ) ?. status
159+ } finally {
160+ releaseLock . resolve ( )
161+ await holder
163162 }
163+ await completion
164164
165- expect ( finishedWithoutLedger ) . toBe ( 0 )
165+ expect ( statusWhileLedgerBlocked ) . toBe ( 'running' )
166+ expect ( await logRow ( executionId ) ) . toMatchObject ( { status : 'completed' } )
167+ expect ( ( await buildCostLedger ( executionId ) ) ?. total ) . toBeCloseTo ( EXECUTION_FEE , 8 )
168+ expect ( Number ( ( await logRow ( executionId ) ) ?. costTotal ) ) . toBeCloseTo ( EXECUTION_FEE , 8 )
166169 } )
167170} )
0 commit comments