Skip to content
25 changes: 25 additions & 0 deletions apps/sim/app/api/workflows/[id]/execute/route.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { flushMacrotask } from '@sim/testing/helpers/async'
import { createDeferred } from '@sim/testing/helpers/deferred'
import { createRouteContext } from '@sim/testing/helpers/http'
import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock'
import {
Expand Down Expand Up @@ -551,6 +553,29 @@ describe('workflow execute async route', () => {
expect(executionOptions.snapshot.input).not.toHaveProperty(PRIVATE_SECRET_PROVENANCE_FIELD)
})

it('holds a synchronous response until the run log and its cost are finalized', async () => {
configureExecutionCaller(EXECUTION_CALLERS[4])
const finalizer = createDeferred<void>()
loggingSessionMockFns.mockWaitForPostExecution.mockReturnValue(finalizer.promise)

let responded = false
const pending = POST(
createInternalProvenanceRequest(),
createRouteContext({ id: 'workflow-1' })
)
void pending.then(() => {
responded = true
})
await vi.waitFor(() => {
expect(loggingSessionMockFns.mockWaitForPostExecution).toHaveBeenCalled()
})
await flushMacrotask()
expect(responded).toBe(false)

finalizer.resolve()
expect((await pending).status).toBe(200)
})

it('queues authenticated workflow input provenance without exposing the private sidecar as input', async () => {
configureExecutionCaller(EXECUTION_CALLERS[4])

Expand Down
6 changes: 6 additions & 0 deletions apps/sim/app/api/workflows/[id]/execute/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1636,6 +1636,12 @@ async function handleExecutePost(
reqLogger.error('Failed to cleanup base64 cache', { error })
})
}
/**
* The sync response is the run's receipt: callers read its log and cost as soon
* as it lands. The core finalizes both in the background, so hold the response
* until they are durable.
*/
await loggingSession.waitForPostExecution()
Comment thread
waleedlatif1 marked this conversation as resolved.
}
}

Expand Down
185 changes: 185 additions & 0 deletions apps/sim/lib/logs/execution/completion-ledger-order.integration.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
/**
* Completion ordering against real PostgreSQL: a run's usage ledger is written before its
* log reads terminal, so a reader that sees a finished run always sees its cost.
*/
import { db } from '@sim/db'
import {
usageLog,
user,
workflow,
workflowExecutionLogs,
workflowExecutionSnapshots,
workspace,
} from '@sim/db/schema'
import { createDeferred } from '@sim/testing'
import { generateId } from '@sim/utils/id'
import { eq, sql } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
import { buildCostLedger } from '@/lib/logs/cost-ledger'
import { executionLogger } from '@/lib/logs/execution/logger'
import { calculateCostSummary } from '@/lib/logs/execution/logging-factory'
import type { WorkflowState } from '@/lib/logs/types'

const ids = {
owner: `ledger-order-owner-${generateId()}`,
workspace: generateId(),
workflow: generateId(),
}

/** The advisory lock the usage ledger write takes for its execution before inserting. */
const USAGE_RECONCILE_LOCK = 'execution_usage_reconcile'
const EXECUTION_FEE = 0.005

const workflowState: WorkflowState = {
blocks: {
start: {
id: 'start',
type: 'starter',
name: 'Start',
position: { x: 0, y: 0 },
subBlocks: {},
outputs: {},
enabled: true,
},
},
edges: [],
loops: {},
parallels: {},
}

async function startExecution(executionId: string) {
await executionLogger.startWorkflowExecution({
workflowId: ids.workflow,
workspaceId: ids.workspace,
executionId,
trigger: { type: 'api', source: 'api', timestamp: new Date().toISOString() },
environment: {
variables: {},
workflowId: ids.workflow,
executionId,
userId: ids.owner,
workspaceId: ids.workspace,
},
workflowState,
})
}

async function logRow(executionId: string) {
const [row] = await db
.select({ status: workflowExecutionLogs.status, costTotal: workflowExecutionLogs.costTotal })
.from(workflowExecutionLogs)
.where(eq(workflowExecutionLogs.executionId, executionId))
return row
}

beforeAll(async () => {
const now = new Date()
await db.insert(user).values({
id: ids.owner,
name: 'Ledger Order',
email: `${ids.owner}@ledger-order.test`,
emailVerified: true,
createdAt: now,
updatedAt: now,
})
await db.insert(workspace).values({
id: ids.workspace,
name: 'Ledger Order',
ownerId: ids.owner,
billedAccountUserId: ids.owner,
})
await db.insert(workflow).values({
id: ids.workflow,
userId: ids.owner,
workspaceId: ids.workspace,
name: 'Ledger Order',
lastSynced: now,
createdAt: now,
updatedAt: now,
})
})

afterAll(async () => {
await db.delete(usageLog).where(eq(usageLog.workflowId, ids.workflow))
await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow))
await db
.delete(workflowExecutionSnapshots)
.where(eq(workflowExecutionSnapshots.workflowId, ids.workflow))
await db.delete(workspace).where(eq(workspace.id, ids.workspace))
await db.delete(user).where(eq(user.id, ids.owner))
})

/**
* Whether a session waits on this execution's ledger lock. A bigint advisory key is stored
* split across `classid` (high 32 bits) and `objid` (low 32 bits) with `objsubid = 1`.
*/
async function isLedgerLockAwaited(executionId: string) {
const rows = await db.execute<{ waiting: boolean }>(sql`
SELECT EXISTS (
SELECT 1 FROM pg_locks
WHERE locktype = 'advisory' AND NOT granted AND objsubid = 1
AND ((classid::bigint << 32) | objid::bigint) = hashtextextended(${executionId}, 0)
) AS waiting
`)
return Boolean(rows[0]?.waiting)
}

describe('completeWorkflowExecution', () => {
it('writes the cost ledger before the run reads finished', async () => {
const billingAttribution = await resolveBillingAttribution({
actorUserId: ids.owner,
workspaceId: ids.workspace,
})
const executionId = generateId()
await startExecution(executionId)

/** Holds the ledger write at its lock, so the log can be read while it waits there. */
const lockHeld = createDeferred<void>()
const releaseLock = createDeferred<void>()
const holder = db.transaction(async (tx) => {
await acquireAdvisoryXactLock(tx, USAGE_RECONCILE_LOCK, executionId)
lockHeld.resolve()
await releaseLock.promise
})
await lockHeld.promise

const completion = executionLogger.completeWorkflowExecution({
executionId,
endedAt: new Date().toISOString(),
totalDurationMs: 5,
costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }),
finalOutput: {},
traceSpans: [],
status: 'completed',
actorUserId: ids.owner,
billingAttribution,
})

/** A completion that settles before blocking surfaces its own outcome instead of a timeout. */
const settledWithoutBlocking = completion.then(() => {
throw new Error('Completion finished without waiting on the ledger lock')
})

let statusWhileLedgerBlocked: string | undefined
try {
await Promise.race([
vi.waitFor(async () => {
expect(await isLedgerLockAwaited(executionId)).toBe(true)
}),
settledWithoutBlocking,
])
statusWhileLedgerBlocked = (await logRow(executionId))?.status
} finally {
releaseLock.resolve()
await holder
}
await completion

expect(statusWhileLedgerBlocked).toBe('running')
expect(await logRow(executionId)).toMatchObject({ status: 'completed' })
expect((await buildCostLedger(executionId))?.total).toBeCloseTo(EXECUTION_FEE, 8)
expect(Number((await logRow(executionId))?.costTotal)).toBeCloseTo(EXECUTION_FEE, 8)
})
})
21 changes: 15 additions & 6 deletions apps/sim/lib/logs/execution/logger.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -219,7 +219,10 @@ describe('ExecutionLogger', () => {
vi.spyOn(logger as any, 'applyPiiRedaction').mockImplementation(
async (_workspaceId: unknown, payload: unknown) => payload
)
vi.spyOn(logger as any, 'recordExecutionUsage').mockResolvedValue(0)
vi.spyOn(logger as any, 'recordExecutionUsage').mockResolvedValue({
recordedIncrement: 0,
costTotalRefined: false,
})

const result = await logger.completeWorkflowExecution({
executionId: 'execution-1',
Expand Down Expand Up @@ -293,15 +296,21 @@ describe('ExecutionLogger', () => {
])
const internals = logger as unknown as {
applyPiiRedaction: (workspaceId: string, payload: Record<string, unknown>) => unknown
recordExecutionUsage: () => Promise<number>
recordExecutionUsage: () => Promise<{
recordedIncrement: number
costTotalRefined: boolean
}>
}
vi.spyOn(internals, 'applyPiiRedaction').mockImplementation(
async (_workspaceId: string, payload: Record<string, unknown>) =>
Object.hasOwn(params, 'redactedState')
? { ...payload, executionState: params.redactedState }
: payload
)
vi.spyOn(internals, 'recordExecutionUsage').mockResolvedValue(0)
vi.spyOn(internals, 'recordExecutionUsage').mockResolvedValue({
recordedIncrement: 0,
costTotalRefined: false,
})

await logger.completeWorkflowExecution({
executionId: 'execution-1',
Expand Down Expand Up @@ -827,7 +836,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
}),
])
// Returns the amount recorded at this boundary (drives threshold-email math).
expect(recorded).toBeCloseTo(1.005, 8)
expect(recorded.recordedIncrement).toBeCloseTo(1.005, 8)
// cost_total is refined to the exact ledger sum inside the locked tx.
expect(dbChainMockFns.update).toHaveBeenCalledTimes(1)
})
Expand Down Expand Up @@ -872,7 +881,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
expect(lastEntries()).not.toContainEqual(
expect.objectContaining({ category: 'model', description: 'mothership' })
)
expect(recorded).toBeCloseTo(1.005, 8)
expect(recorded.recordedIncrement).toBeCloseTo(1.005, 8)
expect(setCostTotalMock).toHaveBeenCalledWith({ costTotal: '1.505' })
})

Expand Down Expand Up @@ -913,7 +922,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => {
'user-1'
)

expect(recorded).toBe(0)
expect(recorded.recordedIncrement).toBe(0)
expect(recordUsage).not.toHaveBeenCalled()
})

Expand Down
Loading
Loading