From ee376e66f75a6e8be0f6e6d3da6705266a880bf3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:18:36 -0700 Subject: [PATCH 1/6] fix(logs): record a run's cost before it reads finished, and hold sync execute responses until the log is final --- .../api/workflows/[id]/execute/route.test.ts | 25 ++ .../app/api/workflows/[id]/execute/route.ts | 6 + .../completion-ledger-order.integration.ts | 155 ++++++++ apps/sim/lib/logs/execution/logger.test.ts | 21 +- apps/sim/lib/logs/execution/logger.ts | 333 +++++++++--------- 5 files changed, 374 insertions(+), 166 deletions(-) create mode 100644 apps/sim/lib/logs/execution/completion-ledger-order.integration.ts diff --git a/apps/sim/app/api/workflows/[id]/execute/route.test.ts b/apps/sim/app/api/workflows/[id]/execute/route.test.ts index e9ee7726393..2deaadc5b43 100644 --- a/apps/sim/app/api/workflows/[id]/execute/route.test.ts +++ b/apps/sim/app/api/workflows/[id]/execute/route.test.ts @@ -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 { @@ -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() + 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]) diff --git a/apps/sim/app/api/workflows/[id]/execute/route.ts b/apps/sim/app/api/workflows/[id]/execute/route.ts index f40f0d6330a..efaf8c241fd 100644 --- a/apps/sim/app/api/workflows/[id]/execute/route.ts +++ b/apps/sim/app/api/workflows/[id]/execute/route.ts @@ -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() } } diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts new file mode 100644 index 00000000000..f7529195319 --- /dev/null +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -0,0 +1,155 @@ +/** + * Completion ordering against real PostgreSQL: a run's usage ledger is durable 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 { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution' +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(), +} + +/** Enough completions that an ordering gap is observed on every run where one exists. */ +const COMPLETIONS = 20 +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)) +}) + +describe('completeWorkflowExecution', () => { + it('never exposes a finished run before its cost ledger', async () => { + const billingAttribution = await resolveBillingAttribution({ + actorUserId: ids.owner, + workspaceId: ids.workspace, + }) + let finishedWithoutLedger = 0 + + for (let i = 0; i < COMPLETIONS; i++) { + const executionId = generateId() + await startExecution(executionId) + + /** The ledger only grows, so the first read that finds the run finished is the one that counts. */ + let completing = true + const reader = (async () => { + while (completing) { + if ((await logRow(executionId))?.status !== 'completed') continue + if ((await buildCostLedger(executionId)) === null) finishedWithoutLedger++ + return + } + })() + + await executionLogger.completeWorkflowExecution({ + executionId, + endedAt: new Date().toISOString(), + totalDurationMs: 5, + costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }), + finalOutput: {}, + traceSpans: [], + status: 'completed', + actorUserId: ids.owner, + billingAttribution, + }) + completing = false + await reader + + const ledger = await buildCostLedger(executionId) + expect(ledger?.total).toBeCloseTo(EXECUTION_FEE, 8) + expect(Number((await logRow(executionId))?.costTotal)).toBeCloseTo(EXECUTION_FEE, 8) + } + + expect(finishedWithoutLedger).toBe(0) + }) +}) diff --git a/apps/sim/lib/logs/execution/logger.test.ts b/apps/sim/lib/logs/execution/logger.test.ts index cbd20a79913..07f4b4b8bb8 100644 --- a/apps/sim/lib/logs/execution/logger.test.ts +++ b/apps/sim/lib/logs/execution/logger.test.ts @@ -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', @@ -293,7 +296,10 @@ describe('ExecutionLogger', () => { ]) const internals = logger as unknown as { applyPiiRedaction: (workspaceId: string, payload: Record) => unknown - recordExecutionUsage: () => Promise + recordExecutionUsage: () => Promise<{ + recordedIncrement: number + costTotalRefined: boolean + }> } vi.spyOn(internals, 'applyPiiRedaction').mockImplementation( async (_workspaceId: string, payload: Record) => @@ -301,7 +307,10 @@ describe('ExecutionLogger', () => { ? { ...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', @@ -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) }) @@ -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' }) }) @@ -913,7 +922,7 @@ describe('recordExecutionUsage boundary-delta reconciliation', () => { 'user-1' ) - expect(recorded).toBe(0) + expect(recorded.recordedIncrement).toBe(0) expect(recordUsage).not.toHaveBeenCalled() }) diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index 650481bc68d..7313c0fdcc8 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -105,6 +105,39 @@ const STATE_SNAPSHOT_FOREIGN_KEY = 'workflow_execution_logs_state_snapshot_id_wo type ExecutionData = WorkflowExecutionLog['executionData'] +/** What one completion boundary wrote to the usage ledger. */ +interface ExecutionUsageRecording { + /** Billable cost recorded at this boundary — the increment, not the run total. */ + recordedIncrement: number + /** Whether the ledger write also set `cost_total` to the exact reconciled sum. */ + costTotalRefined: boolean +} + +const NO_USAGE_RECORDED: ExecutionUsageRecording = { recordedIncrement: 0, costTotalRefined: false } + +/** + * The payer's usage before a boundary records its increment, read for the threshold + * email so usage after = before + increment never counts the boundary twice. + */ +type UsageThresholdEmailContext = + | { + scope: 'user' + userId: string + userEmail: string + userName: string | null + planName: string + periodStart: Date + before: Awaited> + } + | { + scope: 'organization' + organizationId: string + planName: string + periodStart: Date + orgLimit: number + orgUsageBefore: number + } + function getJsonByteSize( value: unknown, maxBytes = MAX_EXECUTION_DATA_BYTES + 1 @@ -1129,6 +1162,101 @@ export class ExecutionLogger implements IExecutionLoggerService { } const completedExecutionLargeValueKeys = collectLargeValueReferenceKeys(storedExecutionData) + const exactBillingContext = billingAttribution + ? toBillingContext(billingAttribution) + : undefined + + /** + * The usage ledger is written before the terminal status commits, so a reader that + * sees a finished run also sees its itemized cost: `buildCostLedger` reads a run with + * no ledger rows as one that has no ledger at all. Skipped without a log row, whose + * completion below throws before this boundary could bill anything. + */ + let usageRecording = NO_USAGE_RECORDED + let emailContext: UsageThresholdEmailContext | undefined + if (existingLog) { + try { + // Skip workflow lookup if workflow was deleted. + const wf = existingLog.workflowId + ? (await db.select().from(workflow).where(eq(workflow.id, existingLog.workflowId)))[0] + : undefined + + const payerContactUserId = billingAttribution?.billedAccountUserId ?? actorUserId + const usr = + wf && payerContactUserId + ? ( + await db + .select({ id: userTable.id, email: userTable.email, name: userTable.name }) + .from(userTable) + .where(eq(userTable.id, payerContactUserId)) + .limit(1) + )[0] + : undefined + + /** + * The pre-increment usage for the threshold email is read BEFORE recording. The + * organization read is the soft one: the email is level-triggered and claimed + * once per period, so a lagging sum only delays it. + */ + if ( + billingAttribution?.billingEntity.type === 'organization' && + billingAttribution.payerSubscription && + exactBillingContext + ) { + const organizationId = billingAttribution.billingEntity.id + const payerSubscription = billingAttribution.payerSubscription + const [{ getDisplayPlanName }, { limit: orgLimit }, orgUsageBefore] = await Promise.all([ + import('@/lib/billing/plan-helpers'), + getOrgUsageLimit(organizationId, payerSubscription.plan, payerSubscription.seats), + readSoftGateUsageCost( + billingAttribution.billingEntity, + exactBillingContext.billingPeriod + ), + ]) + emailContext = { + scope: 'organization', + organizationId, + planName: getDisplayPlanName(payerSubscription.plan), + periodStart: exactBillingContext.billingPeriod.start, + orgLimit, + orgUsageBefore, + } + } else if ( + billingAttribution?.billingEntity.type === 'user' && + exactBillingContext && + usr?.email + ) { + const sub = await getHighestPriorityPersonalSubscription(usr.id) + const { getDisplayPlanName } = await import('@/lib/billing/plan-helpers') + emailContext = { + scope: 'user', + userId: usr.id, + userEmail: usr.email, + userName: usr.name, + planName: getDisplayPlanName(sub?.plan), + periodStart: exactBillingContext.billingPeriod.start, + before: await checkResolvedUsageStatus(usr.id, sub, exactBillingContext), + } + } + } catch (e) { + execLog.warn('Usage threshold notification check failed (non-fatal)', { error: e }) + } + + // Record usage exactly once for every path; a failed threshold read above must + // never leave the run unbilled. The recorded increment is the amount billed at + // this boundary, not the cumulative run total — so resumed runs don't + // double-count pre-pause cost in the threshold email. + usageRecording = await this.recordExecutionUsage( + existingLog.workflowId, + costSummary, + existingLog.trigger as ExecutionTrigger['type'], + executionId, + actorUserId, + exactBillingContext, + status !== 'pending' + ) + } + const { updatedLog, completionPersisted } = await execDb.transaction(async (tx) => { await setExecutionLogWriteTimeouts(tx) @@ -1147,8 +1275,14 @@ export class ExecutionLogger implements IExecutionLoggerService { // resumes into an empty-span error/cancel/cost-only fallback produces a // base-only summary. GREATEST keeps the higher cumulative cost_total, // and models_used is overwritten only when this boundary actually has - // models — so both stay == SUM(usage_log) on every monotonic path. - costTotal: sql`GREATEST(COALESCE(${workflowExecutionLogs.costTotal}, 0), ${costSummary.totalCost.toString()}::numeric)`, + // models — so both stay == SUM(usage_log) on every monotonic path. When + // this boundary's ledger write already set the exact reconciled sum, that + // value stands. + ...(usageRecording.costTotalRefined + ? {} + : { + costTotal: sql`GREATEST(COALESCE(${workflowExecutionLogs.costTotal}, 0), ${costSummary.totalCost.toString()}::numeric)`, + }), ...(Object.keys(costSummary.models).length > 0 ? { modelsUsed: Object.keys(costSummary.models) } : {}), @@ -1201,161 +1335,38 @@ export class ExecutionLogger implements IExecutionLoggerService { }) if (progressMarkers !== null) void clearProgressMarkers(executionId) - const exactBillingContext = billingAttribution - ? toBillingContext(billingAttribution) - : undefined - try { - // Skip workflow lookup if workflow was deleted. - const wf = updatedLog.workflowId - ? (await db.select().from(workflow).where(eq(workflow.id, updatedLog.workflowId)))[0] - : undefined - - const payerContactUserId = billingAttribution?.billedAccountUserId ?? actorUserId - const usr = - wf && payerContactUserId - ? ( - await db - .select({ id: userTable.id, email: userTable.email, name: userTable.name }) - .from(userTable) - .where(eq(userTable.id, payerContactUserId)) - .limit(1) - )[0] - : undefined - - /** - * The billing context and pre-increment usage for the threshold email are read BEFORE - * recording, so usage after = before + costDelta doesn't double-count this boundary's own - * increment. The organization read is the soft one: the email is level-triggered and - * claimed once per period, so a lagging sum only delays it. - */ - type EmailContext = - | { - scope: 'user' - userId: string - userEmail: string - userName: string | null - planName: string - periodStart: Date - before: Awaited> - } - | { - scope: 'organization' - organizationId: string - planName: string - periodStart: Date - orgLimit: number - orgUsageBefore: number - } - const billingContext = exactBillingContext - let emailContext: EmailContext | undefined - - if ( - billingAttribution?.billingEntity.type === 'organization' && - billingAttribution.payerSubscription && - exactBillingContext - ) { - const organizationId = billingAttribution.billingEntity.id - const payerSubscription = billingAttribution.payerSubscription - const { getDisplayPlanName } = await import('@/lib/billing/plan-helpers') - const { limit: orgLimit } = await getOrgUsageLimit( - organizationId, - payerSubscription.plan, - payerSubscription.seats - ) - emailContext = { - scope: 'organization', - organizationId, - planName: getDisplayPlanName(payerSubscription.plan), - periodStart: exactBillingContext.billingPeriod.start, - orgLimit, - orgUsageBefore: await readSoftGateUsageCost( - billingAttribution.billingEntity, - exactBillingContext.billingPeriod - ), - } - } else if ( - billingAttribution?.billingEntity.type === 'user' && - exactBillingContext && - usr?.email - ) { - const sub = await getHighestPriorityPersonalSubscription(usr.id) - const { getDisplayPlanName } = await import('@/lib/billing/plan-helpers') - emailContext = { - scope: 'user', - userId: usr.id, - userEmail: usr.email, - userName: usr.name, - planName: getDisplayPlanName(sub?.plan), - periodStart: exactBillingContext.billingPeriod.start, - before: await checkResolvedUsageStatus(usr.id, sub, exactBillingContext), - } - } - - // Record usage exactly once for every path. costDelta is the amount - // actually recorded at this boundary (the increment), not the cumulative - // run total — so resumed runs don't double-count pre-pause cost below. - const costDelta = await this.recordExecutionUsage( - updatedLog.workflowId, - costSummary, - updatedLog.trigger as ExecutionTrigger['type'], - executionId, - actorUserId, - billingContext, - status !== 'pending' - ) - - // Best-effort usage-threshold email. - if (emailContext?.scope === 'user') { - await maybeSendUsageThresholdEmail({ - scope: 'user', - userId: emailContext.userId, - userEmail: emailContext.userEmail, - userName: emailContext.userName || undefined, - planName: emailContext.planName, - periodStart: emailContext.periodStart, - workspaceId: updatedLog.workspaceId, - usageBefore: emailContext.before.currentUsage, - costDelta, - limit: emailContext.before.limit, - }) - } else if (emailContext?.scope === 'organization') { - await maybeSendUsageThresholdEmail({ - scope: 'organization', - organizationId: emailContext.organizationId, - planName: emailContext.planName, - periodStart: emailContext.periodStart, - workspaceId: updatedLog.workspaceId, - usageBefore: emailContext.orgUsageBefore, - costDelta, - limit: emailContext.orgLimit, - }) - } - } catch (e) { - // Safety net: if a step above threw BEFORE the single record call, ensure - // the run is still billed. Reconciliation is idempotent, so re-recording - // after a successful call is a no-op. + if (emailContext) { + const costDelta = usageRecording.recordedIncrement try { - await this.recordExecutionUsage( - updatedLog.workflowId, - costSummary, - updatedLog.trigger as ExecutionTrigger['type'], - executionId, - actorUserId, - exactBillingContext, - status !== 'pending' - ) - } catch (recordError) { - /* The safety net is the last thing between a completed run and an unbilled - one. Swallowing it left the only emitted line saying a notification check - had failed and was non-fatal. */ - execLog.error('Failed to record execution usage — this run may be unbilled', { - error: recordError, - executionId, - workflowId: updatedLog.workflowId, - }) + if (emailContext.scope === 'user') { + await maybeSendUsageThresholdEmail({ + scope: 'user', + userId: emailContext.userId, + userEmail: emailContext.userEmail, + userName: emailContext.userName || undefined, + planName: emailContext.planName, + periodStart: emailContext.periodStart, + workspaceId: updatedLog.workspaceId, + usageBefore: emailContext.before.currentUsage, + costDelta, + limit: emailContext.before.limit, + }) + } else { + await maybeSendUsageThresholdEmail({ + scope: 'organization', + organizationId: emailContext.organizationId, + planName: emailContext.planName, + periodStart: emailContext.periodStart, + workspaceId: updatedLog.workspaceId, + usageBefore: emailContext.orgUsageBefore, + costDelta, + limit: emailContext.orgLimit, + }) + } + } catch (e) { + execLog.warn('Usage threshold notification check failed (non-fatal)', { error: e }) } - execLog.warn('Usage threshold notification check failed (non-fatal)', { error: e }) } if (completionPersisted) { @@ -1489,7 +1500,7 @@ export class ExecutionLogger implements IExecutionLoggerService { * boundary, where the summary already holds the run's cumulative tokens. */ isTerminalBoundary = true - ): Promise { + ): Promise { const statsLog = logger.withMetadata({ workflowId: workflowId ?? undefined, executionId }) // The usage ledger (recordUsage below) is written regardless of @@ -1501,10 +1512,11 @@ export class ExecutionLogger implements IExecutionLoggerService { if (!workflowId) { statsLog.debug('Workflow was deleted, skipping usage recording') - return 0 + return NO_USAGE_RECORDED } let recordedIncrement = 0 + let costTotalRefined = false try { const [workflowRecord] = await db .select() @@ -1514,7 +1526,7 @@ export class ExecutionLogger implements IExecutionLoggerService { if (!workflowRecord) { statsLog.error('Workflow not found for usage recording') - return 0 + return NO_USAGE_RECORDED } const userId = actorUserId?.trim() || null @@ -1522,7 +1534,7 @@ export class ExecutionLogger implements IExecutionLoggerService { statsLog.error('Missing actor in execution context; skipping usage recording', { trigger, }) - return 0 + return NO_USAGE_RECORDED } // Build the run's *cumulative* target ledger lines from the cost summary. @@ -1623,7 +1635,7 @@ export class ExecutionLogger implements IExecutionLoggerService { // error for a charge that does not exist. if (targets.length === 0 && !canRecordUnbilled) { statsLog.debug('No cost to record') - return 0 + return NO_USAGE_RECORDED } if (workflowRecord.workspaceId && !billingContext) { @@ -1775,6 +1787,7 @@ export class ExecutionLogger implements IExecutionLoggerService { .update(workflowExecutionLogs) .set({ costTotal: displayedCostTotal.toString() }) .where(eq(workflowExecutionLogs.executionId, executionId)) + costTotalRefined = true } } }) @@ -1819,7 +1832,7 @@ export class ExecutionLogger implements IExecutionLoggerService { ) } - return recordedIncrement + return { recordedIncrement, costTotalRefined } } /** From 7e0af7b654320e480622961216e927738ef4cfcd Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:27:22 -0700 Subject: [PATCH 2/6] fix(logs): stop the ledger-order reader when a completion throws --- .../completion-ledger-order.integration.ts | 29 ++++++++++--------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts index f7529195319..d8f149ee200 100644 --- a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -131,19 +131,22 @@ describe('completeWorkflowExecution', () => { } })() - await executionLogger.completeWorkflowExecution({ - executionId, - endedAt: new Date().toISOString(), - totalDurationMs: 5, - costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }), - finalOutput: {}, - traceSpans: [], - status: 'completed', - actorUserId: ids.owner, - billingAttribution, - }) - completing = false - await reader + try { + await executionLogger.completeWorkflowExecution({ + executionId, + endedAt: new Date().toISOString(), + totalDurationMs: 5, + costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }), + finalOutput: {}, + traceSpans: [], + status: 'completed', + actorUserId: ids.owner, + billingAttribution, + }) + } finally { + completing = false + await reader + } const ledger = await buildCostLedger(executionId) expect(ledger?.total).toBeCloseTo(EXECUTION_FEE, 8) From dca04658f140867f981ea80f96cfb666a9203eb9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:33:07 -0700 Subject: [PATCH 3/6] fix(logs): keep a completion failure visible when the ledger-order reader also fails --- .../lib/logs/execution/completion-ledger-order.integration.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts index d8f149ee200..4edcdb14ab6 100644 --- a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -145,8 +145,10 @@ describe('completeWorkflowExecution', () => { }) } finally { completing = false - await reader + // Settle the reader without letting its error replace a completion failure. + await reader.catch(() => {}) } + await reader const ledger = await buildCostLedger(executionId) expect(ledger?.total).toBeCloseTo(EXECUTION_FEE, 8) From ec69d07bb263637911a0b48077b355c006ea3fbd Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:44:34 -0700 Subject: [PATCH 4/6] fix(logs): let the ledger-order reader observe every finished run --- .../completion-ledger-order.integration.ts | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts index 4edcdb14ab6..ac53308a435 100644 --- a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -121,13 +121,20 @@ describe('completeWorkflowExecution', () => { const executionId = generateId() await startExecution(executionId) - /** The ledger only grows, so the first read that finds the run finished is the one that counts. */ + /** + * The ledger only grows, so the first read that finds the run finished is the one that + * counts. Settlement is captured before each read, so a read already in flight when the + * completion lands cannot end the loop before the finished run is observed. + */ let completing = true const reader = (async () => { - while (completing) { - if ((await logRow(executionId))?.status !== 'completed') continue - if ((await buildCostLedger(executionId)) === null) finishedWithoutLedger++ - return + for (;;) { + const settledBeforeRead = !completing + if ((await logRow(executionId))?.status === 'completed') { + if ((await buildCostLedger(executionId)) === null) finishedWithoutLedger++ + return + } + if (settledBeforeRead) return } })() From 1ece8aa4cd05d522883f2ccb77c15df5220a0dbf Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:56:21 -0700 Subject: [PATCH 5/6] fix(logs): assert ledger-before-terminal ordering deterministically at the ledger lock --- .../completion-ledger-order.integration.ts | 103 +++++++++--------- 1 file changed, 53 insertions(+), 50 deletions(-) diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts index ac53308a435..6d5351c4cb4 100644 --- a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -1,6 +1,6 @@ /** - * Completion ordering against real PostgreSQL: a run's usage ledger is durable before - * its log reads terminal, so a reader that sees a finished run always sees its cost. + * 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 { @@ -11,10 +11,12 @@ import { workflowExecutionSnapshots, workspace, } from '@sim/db/schema' +import { createDeferred } from '@sim/testing' import { generateId } from '@sim/utils/id' -import { eq } from 'drizzle-orm' -import { afterAll, beforeAll, describe, expect, it } from 'vitest' +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' @@ -26,8 +28,8 @@ const ids = { workflow: generateId(), } -/** Enough completions that an ordering gap is observed on every run where one exists. */ -const COMPLETIONS = 20 +/** 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 = { @@ -109,59 +111,60 @@ afterAll(async () => { await db.delete(user).where(eq(user.id, ids.owner)) }) +/** Whether a session is waiting on an advisory lock — this file's database has no other traffic. */ +async function hasAdvisoryLockWaiter() { + const rows = await db.execute<{ waiting: boolean }>( + sql`SELECT EXISTS (SELECT 1 FROM pg_locks WHERE locktype = 'advisory' AND NOT granted) AS waiting` + ) + return Boolean(rows[0]?.waiting) +} + describe('completeWorkflowExecution', () => { - it('never exposes a finished run before its cost ledger', async () => { + it('writes the cost ledger before the run reads finished', async () => { const billingAttribution = await resolveBillingAttribution({ actorUserId: ids.owner, workspaceId: ids.workspace, }) - let finishedWithoutLedger = 0 + const executionId = generateId() + await startExecution(executionId) - for (let i = 0; i < COMPLETIONS; i++) { - const executionId = generateId() - await startExecution(executionId) - - /** - * The ledger only grows, so the first read that finds the run finished is the one that - * counts. Settlement is captured before each read, so a read already in flight when the - * completion lands cannot end the loop before the finished run is observed. - */ - let completing = true - const reader = (async () => { - for (;;) { - const settledBeforeRead = !completing - if ((await logRow(executionId))?.status === 'completed') { - if ((await buildCostLedger(executionId)) === null) finishedWithoutLedger++ - return - } - if (settledBeforeRead) return - } - })() + /** Holds the ledger write at its lock, so the log can be read while it waits there. */ + const lockHeld = createDeferred() + const releaseLock = createDeferred() + const holder = db.transaction(async (tx) => { + await acquireAdvisoryXactLock(tx, USAGE_RECONCILE_LOCK, executionId) + lockHeld.resolve() + await releaseLock.promise + }) + await lockHeld.promise - try { - await executionLogger.completeWorkflowExecution({ - executionId, - endedAt: new Date().toISOString(), - totalDurationMs: 5, - costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }), - finalOutput: {}, - traceSpans: [], - status: 'completed', - actorUserId: ids.owner, - billingAttribution, - }) - } finally { - completing = false - // Settle the reader without letting its error replace a completion failure. - await reader.catch(() => {}) - } - await reader + const completion = executionLogger.completeWorkflowExecution({ + executionId, + endedAt: new Date().toISOString(), + totalDurationMs: 5, + costSummary: calculateCostSummary([], { baseExecutionCharge: EXECUTION_FEE }), + finalOutput: {}, + traceSpans: [], + status: 'completed', + actorUserId: ids.owner, + billingAttribution, + }) - const ledger = await buildCostLedger(executionId) - expect(ledger?.total).toBeCloseTo(EXECUTION_FEE, 8) - expect(Number((await logRow(executionId))?.costTotal)).toBeCloseTo(EXECUTION_FEE, 8) + let statusWhileLedgerBlocked: string | undefined + try { + await vi.waitFor(async () => { + expect(await hasAdvisoryLockWaiter()).toBe(true) + }) + statusWhileLedgerBlocked = (await logRow(executionId))?.status + } finally { + releaseLock.resolve() + await holder } + await completion - expect(finishedWithoutLedger).toBe(0) + 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) }) }) From 4a3b2c0f5d23adc2dfa7f385dc7d4fd6db2ee062 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 11:02:27 -0700 Subject: [PATCH 6/6] fix(logs): scope the ledger-lock waiter to the execution and surface early completion failures --- .../completion-ledger-order.integration.ts | 31 ++++++++++++++----- 1 file changed, 23 insertions(+), 8 deletions(-) diff --git a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts index 6d5351c4cb4..111e618a82d 100644 --- a/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts +++ b/apps/sim/lib/logs/execution/completion-ledger-order.integration.ts @@ -111,11 +111,18 @@ afterAll(async () => { await db.delete(user).where(eq(user.id, ids.owner)) }) -/** Whether a session is waiting on an advisory lock — this file's database has no other traffic. */ -async function hasAdvisoryLockWaiter() { - const rows = await db.execute<{ waiting: boolean }>( - sql`SELECT EXISTS (SELECT 1 FROM pg_locks WHERE locktype = 'advisory' AND NOT granted) AS waiting` - ) +/** + * 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) } @@ -150,11 +157,19 @@ describe('completeWorkflowExecution', () => { 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 vi.waitFor(async () => { - expect(await hasAdvisoryLockWaiter()).toBe(true) - }) + await Promise.race([ + vi.waitFor(async () => { + expect(await isLedgerLockAwaited(executionId)).toBe(true) + }), + settledWithoutBlocking, + ]) statusWhileLedgerBlocked = (await logRow(executionId))?.status } finally { releaseLock.resolve()