Skip to content

Commit c30d58d

Browse files
authored
fix(logs): record a run's cost before it reads finished, and hold sync execute responses until the log is final (#8473)
* fix(logs): record a run's cost before it reads finished, and hold sync execute responses until the log is final * fix(logs): stop the ledger-order reader when a completion throws * fix(logs): keep a completion failure visible when the ledger-order reader also fails * fix(logs): let the ledger-order reader observe every finished run * fix(logs): assert ledger-before-terminal ordering deterministically at the ledger lock * fix(logs): scope the ledger-lock waiter to the execution and surface early completion failures
1 parent 294c505 commit c30d58d

5 files changed

Lines changed: 404 additions & 166 deletions

File tree

‎apps/sim/app/api/workflows/[id]/execute/route.test.ts‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import { flushMacrotask } from '@sim/testing/helpers/async'
2+
import { createDeferred } from '@sim/testing/helpers/deferred'
13
import { createRouteContext } from '@sim/testing/helpers/http'
24
import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock'
35
import {
@@ -551,6 +553,29 @@ describe('workflow execute async route', () => {
551553
expect(executionOptions.snapshot.input).not.toHaveProperty(PRIVATE_SECRET_PROVENANCE_FIELD)
552554
})
553555

556+
it('holds a synchronous response until the run log and its cost are finalized', async () => {
557+
configureExecutionCaller(EXECUTION_CALLERS[4])
558+
const finalizer = createDeferred<void>()
559+
loggingSessionMockFns.mockWaitForPostExecution.mockReturnValue(finalizer.promise)
560+
561+
let responded = false
562+
const pending = POST(
563+
createInternalProvenanceRequest(),
564+
createRouteContext({ id: 'workflow-1' })
565+
)
566+
void pending.then(() => {
567+
responded = true
568+
})
569+
await vi.waitFor(() => {
570+
expect(loggingSessionMockFns.mockWaitForPostExecution).toHaveBeenCalled()
571+
})
572+
await flushMacrotask()
573+
expect(responded).toBe(false)
574+
575+
finalizer.resolve()
576+
expect((await pending).status).toBe(200)
577+
})
578+
554579
it('queues authenticated workflow input provenance without exposing the private sidecar as input', async () => {
555580
configureExecutionCaller(EXECUTION_CALLERS[4])
556581

‎apps/sim/app/api/workflows/[id]/execute/route.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1636,6 +1636,12 @@ async function handleExecutePost(
16361636
reqLogger.error('Failed to cleanup base64 cache', { error })
16371637
})
16381638
}
1639+
/**
1640+
* The sync response is the run's receipt: callers read its log and cost as soon
1641+
* as it lands. The core finalizes both in the background, so hold the response
1642+
* until they are durable.
1643+
*/
1644+
await loggingSession.waitForPostExecution()
16391645
}
16401646
}
16411647

Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,185 @@
1+
/**
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.
4+
*/
5+
import { db } from '@sim/db'
6+
import {
7+
usageLog,
8+
user,
9+
workflow,
10+
workflowExecutionLogs,
11+
workflowExecutionSnapshots,
12+
workspace,
13+
} from '@sim/db/schema'
14+
import { createDeferred } from '@sim/testing'
15+
import { generateId } from '@sim/utils/id'
16+
import { eq, sql } from 'drizzle-orm'
17+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
18+
import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
19+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
20+
import { buildCostLedger } from '@/lib/logs/cost-ledger'
21+
import { executionLogger } from '@/lib/logs/execution/logger'
22+
import { calculateCostSummary } from '@/lib/logs/execution/logging-factory'
23+
import type { WorkflowState } from '@/lib/logs/types'
24+
25+
const ids = {
26+
owner: `ledger-order-owner-${generateId()}`,
27+
workspace: generateId(),
28+
workflow: generateId(),
29+
}
30+
31+
/** The advisory lock the usage ledger write takes for its execution before inserting. */
32+
const USAGE_RECONCILE_LOCK = 'execution_usage_reconcile'
33+
const EXECUTION_FEE = 0.005
34+
35+
const workflowState: WorkflowState = {
36+
blocks: {
37+
start: {
38+
id: 'start',
39+
type: 'starter',
40+
name: 'Start',
41+
position: { x: 0, y: 0 },
42+
subBlocks: {},
43+
outputs: {},
44+
enabled: true,
45+
},
46+
},
47+
edges: [],
48+
loops: {},
49+
parallels: {},
50+
}
51+
52+
async function startExecution(executionId: string) {
53+
await executionLogger.startWorkflowExecution({
54+
workflowId: ids.workflow,
55+
workspaceId: ids.workspace,
56+
executionId,
57+
trigger: { type: 'api', source: 'api', timestamp: new Date().toISOString() },
58+
environment: {
59+
variables: {},
60+
workflowId: ids.workflow,
61+
executionId,
62+
userId: ids.owner,
63+
workspaceId: ids.workspace,
64+
},
65+
workflowState,
66+
})
67+
}
68+
69+
async function logRow(executionId: string) {
70+
const [row] = await db
71+
.select({ status: workflowExecutionLogs.status, costTotal: workflowExecutionLogs.costTotal })
72+
.from(workflowExecutionLogs)
73+
.where(eq(workflowExecutionLogs.executionId, executionId))
74+
return row
75+
}
76+
77+
beforeAll(async () => {
78+
const now = new Date()
79+
await db.insert(user).values({
80+
id: ids.owner,
81+
name: 'Ledger Order',
82+
email: `${ids.owner}@ledger-order.test`,
83+
emailVerified: true,
84+
createdAt: now,
85+
updatedAt: now,
86+
})
87+
await db.insert(workspace).values({
88+
id: ids.workspace,
89+
name: 'Ledger Order',
90+
ownerId: ids.owner,
91+
billedAccountUserId: ids.owner,
92+
})
93+
await db.insert(workflow).values({
94+
id: ids.workflow,
95+
userId: ids.owner,
96+
workspaceId: ids.workspace,
97+
name: 'Ledger Order',
98+
lastSynced: now,
99+
createdAt: now,
100+
updatedAt: now,
101+
})
102+
})
103+
104+
afterAll(async () => {
105+
await db.delete(usageLog).where(eq(usageLog.workflowId, ids.workflow))
106+
await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow))
107+
await db
108+
.delete(workflowExecutionSnapshots)
109+
.where(eq(workflowExecutionSnapshots.workflowId, ids.workflow))
110+
await db.delete(workspace).where(eq(workspace.id, ids.workspace))
111+
await db.delete(user).where(eq(user.id, ids.owner))
112+
})
113+
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+
`)
126+
return Boolean(rows[0]?.waiting)
127+
}
128+
129+
describe('completeWorkflowExecution', () => {
130+
it('writes the cost ledger before the run reads finished', async () => {
131+
const billingAttribution = await resolveBillingAttribution({
132+
actorUserId: ids.owner,
133+
workspaceId: ids.workspace,
134+
})
135+
const executionId = generateId()
136+
await startExecution(executionId)
137+
138+
/** Holds the ledger write at its lock, so the log can be read while it waits there. */
139+
const lockHeld = createDeferred<void>()
140+
const releaseLock = createDeferred<void>()
141+
const holder = db.transaction(async (tx) => {
142+
await acquireAdvisoryXactLock(tx, USAGE_RECONCILE_LOCK, executionId)
143+
lockHeld.resolve()
144+
await releaseLock.promise
145+
})
146+
await lockHeld.promise
147+
148+
const completion = executionLogger.completeWorkflowExecution({
149+
executionId,
150+
endedAt: new Date().toISOString(),
151+
totalDurationMs: 5,
152+
costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }),
153+
finalOutput: {},
154+
traceSpans: [],
155+
status: 'completed',
156+
actorUserId: ids.owner,
157+
billingAttribution,
158+
})
159+
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+
165+
let statusWhileLedgerBlocked: string | undefined
166+
try {
167+
await Promise.race([
168+
vi.waitFor(async () => {
169+
expect(await isLedgerLockAwaited(executionId)).toBe(true)
170+
}),
171+
settledWithoutBlocking,
172+
])
173+
statusWhileLedgerBlocked = (await logRow(executionId))?.status
174+
} finally {
175+
releaseLock.resolve()
176+
await holder
177+
}
178+
await completion
179+
180+
expect(statusWhileLedgerBlocked).toBe('running')
181+
expect(await logRow(executionId)).toMatchObject({ status: 'completed' })
182+
expect((await buildCostLedger(executionId))?.total).toBeCloseTo(EXECUTION_FEE, 8)
183+
expect(Number((await logRow(executionId))?.costTotal)).toBeCloseTo(EXECUTION_FEE, 8)
184+
})
185+
})

‎apps/sim/lib/logs/execution/logger.test.ts‎

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -219,7 +219,10 @@ describe('ExecutionLogger', () => {
219219
vi.spyOn(logger as any, 'applyPiiRedaction').mockImplementation(
220220
async (_workspaceId: unknown, payload: unknown) => payload
221221
)
222-
vi.spyOn(logger as any, 'recordExecutionUsage').mockResolvedValue(0)
222+
vi.spyOn(logger as any, 'recordExecutionUsage').mockResolvedValue({
223+
recordedIncrement: 0,
224+
costTotalRefined: false,
225+
})
223226

224227
const result = await logger.completeWorkflowExecution({
225228
executionId: 'execution-1',
@@ -293,15 +296,21 @@ describe('ExecutionLogger', () => {
293296
])
294297
const internals = logger as unknown as {
295298
applyPiiRedaction: (workspaceId: string, payload: Record<string, unknown>) => unknown
296-
recordExecutionUsage: () => Promise<number>
299+
recordExecutionUsage: () => Promise<{
300+
recordedIncrement: number
301+
costTotalRefined: boolean
302+
}>
297303
}
298304
vi.spyOn(internals, 'applyPiiRedaction').mockImplementation(
299305
async (_workspaceId: string, payload: Record<string, unknown>) =>
300306
Object.hasOwn(params, 'redactedState')
301307
? { ...payload, executionState: params.redactedState }
302308
: payload
303309
)
304-
vi.spyOn(internals, 'recordExecutionUsage').mockResolvedValue(0)
310+
vi.spyOn(internals, 'recordExecutionUsage').mockResolvedValue({
311+
recordedIncrement: 0,
312+
costTotalRefined: false,
313+
})
305314

306315
await logger.completeWorkflowExecution({
307316
executionId: 'execution-1',
@@ -827,7 +836,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
827836
}),
828837
])
829838
// Returns the amount recorded at this boundary (drives threshold-email math).
830-
expect(recorded).toBeCloseTo(1.005, 8)
839+
expect(recorded.recordedIncrement).toBeCloseTo(1.005, 8)
831840
// cost_total is refined to the exact ledger sum inside the locked tx.
832841
expect(dbChainMockFns.update).toHaveBeenCalledTimes(1)
833842
})
@@ -872,7 +881,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
872881
expect(lastEntries()).not.toContainEqual(
873882
expect.objectContaining({ category: 'model', description: 'mothership' })
874883
)
875-
expect(recorded).toBeCloseTo(1.005, 8)
884+
expect(recorded.recordedIncrement).toBeCloseTo(1.005, 8)
876885
expect(setCostTotalMock).toHaveBeenCalledWith({ costTotal: '1.505' })
877886
})
878887

@@ -913,7 +922,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
913922
'user-1'
914923
)
915924

916-
expect(recorded).toBe(0)
925+
expect(recorded.recordedIncrement).toBe(0)
917926
expect(recordUsage).not.toHaveBeenCalled()
918927
})
919928

0 commit comments

Comments
 (0)