From ac35656d5886def2922b7122b1eee05a283c24b1 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 19:50:40 -0700 Subject: [PATCH 1/3] fix(execution): publish a pause only after its run log is finalized --- .../executor/pause-persistence.integration.ts | 187 ++++++++++++++++++ .../workflows/executor/pause-persistence.ts | 6 + 2 files changed, 193 insertions(+) create mode 100644 apps/sim/lib/workflows/executor/pause-persistence.integration.ts diff --git a/apps/sim/lib/workflows/executor/pause-persistence.integration.ts b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts new file mode 100644 index 00000000000..5d5b8dc74fa --- /dev/null +++ b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts @@ -0,0 +1,187 @@ +/** + * Pause publication against real PostgreSQL: a paused run becomes resumable only after + * its log has been finalized out of `running`, so an immediate resume finds a claimable log. + */ +import { db } from '@sim/db' +import { + pausedExecutions, + resumeQueue, + user, + workflow, + workflowExecutionLogs, + workflowExecutionSnapshots, + workspace, +} from '@sim/db/schema' +import { createDeferred } from '@sim/testing' +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { + type BillingAttributionSnapshot, + resolveBillingAttribution, +} from '@/lib/billing/core/billing-attribution' +import { LoggingSession } from '@/lib/logs/execution/logging-session' +import type { WorkflowState } from '@/lib/logs/types' +import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager' +import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence' +import type { ExecutionResult } from '@/executor/types' + +const ids = { + owner: `pause-publish-owner-${generateId()}`, + workspace: generateId(), + workflow: generateId(), + execution: generateId(), +} + +const CONTEXT_ID = 'approval' + +/** Long enough for an unguarded publish to commit; a guarded one never resolves while held. */ +const UNGUARDED_PUBLISH_WINDOW_MS = 250 + +const workflowState: WorkflowState = { + blocks: { + start: { + id: 'start', + type: 'starter', + name: 'Start', + position: { x: 0, y: 0 }, + subBlocks: {}, + outputs: {}, + enabled: true, + }, + }, + edges: [], + loops: {}, + parallels: {}, +} + +function pausedResult(billingAttribution: BillingAttributionSnapshot): ExecutionResult { + return { + success: true, + output: {}, + status: 'paused', + pausePoints: [ + { + contextId: CONTEXT_ID, + blockId: CONTEXT_ID, + response: {}, + registeredAt: new Date().toISOString(), + resumeStatus: 'paused', + snapshotReady: true, + pauseKind: 'human', + }, + ], + snapshotSeed: { + snapshot: JSON.stringify({ + metadata: { + workflowId: ids.workflow, + workspaceId: ids.workspace, + executionId: ids.execution, + userId: ids.owner, + billingAttribution, + }, + }), + triggerIds: [], + }, + } +} + +async function logStatus() { + const [row] = await db + .select({ status: workflowExecutionLogs.status }) + .from(workflowExecutionLogs) + .where(eq(workflowExecutionLogs.executionId, ids.execution)) + return row?.status +} + +function resume() { + return PauseResumeManager.enqueueOrStartResume({ + executionId: ids.execution, + workflowId: ids.workflow, + contextId: CONTEXT_ID, + resumeInput: {}, + userId: ids.owner, + }) +} + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: ids.owner, + name: 'Pause Publish', + email: `${ids.owner}@pause-publish.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: ids.workspace, + name: 'Pause Publish', + ownerId: ids.owner, + billedAccountUserId: ids.owner, + }) + await db.insert(workflow).values({ + id: ids.workflow, + userId: ids.owner, + workspaceId: ids.workspace, + name: 'Pause Publish', + lastSynced: now, + createdAt: now, + updatedAt: now, + }) +}) + +afterAll(async () => { + await db.delete(resumeQueue).where(eq(resumeQueue.parentExecutionId, ids.execution)) + await db.delete(pausedExecutions).where(eq(pausedExecutions.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('handlePostExecutionPauseState', () => { + it('publishes a pause only after the run log is finalized, so an immediate resume finds a claimable log', async () => { + const billingAttribution = await resolveBillingAttribution({ + actorUserId: ids.owner, + workspaceId: ids.workspace, + }) + const loggingSession = new LoggingSession(ids.workflow, ids.execution, 'api', 'pause-publish') + await loggingSession.safeStart({ + userId: ids.owner, + workspaceId: ids.workspace, + billingAttribution, + workflowState, + }) + expect(await logStatus()).toBe('running') + + /** Holds the core's background log finalizer open, as a slow trace projection would. */ + const finalizer = createDeferred() + loggingSession.setPostExecutionPromise( + finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] })) + ) + + const publish = handlePostExecutionPauseState({ + result: pausedResult(billingAttribution), + workflowId: ids.workflow, + executionId: ids.execution, + loggingSession, + }) + await Promise.race([publish, sleep(UNGUARDED_PUBLISH_WINDOW_MS)]) + + expect(await logStatus()).toBe('running') + await expect(resume()).rejects.toMatchObject({ + name: 'ResumeAdmissionError', + statusCode: 404, + }) + + finalizer.resolve() + await publish + + expect(await logStatus()).toBe('pending') + await expect(resume()).resolves.toMatchObject({ status: 'starting' }) + }) +}) diff --git a/apps/sim/lib/workflows/executor/pause-persistence.ts b/apps/sim/lib/workflows/executor/pause-persistence.ts index 552263f490f..75092c65449 100644 --- a/apps/sim/lib/workflows/executor/pause-persistence.ts +++ b/apps/sim/lib/workflows/executor/pause-persistence.ts @@ -24,6 +24,11 @@ interface HandlePostExecutionPauseStateArgs { * - If execution is paused with a valid snapshot: persists to `paused_executions` table * - If execution is paused without a snapshot: marks execution as failed * - If execution is not paused: processes any queued resume entries + * + * A pause is published only after the core's post-execution logging settles. The + * core finalizes the run log in the background, and a resume claims that log only + * once it has left `running`, so publishing first lets an immediate resume be + * rejected as no longer resumable. */ export async function handlePostExecutionPauseState({ result, @@ -33,6 +38,7 @@ export async function handlePostExecutionPauseState({ loggingSession, }: HandlePostExecutionPauseStateArgs): Promise { if (result.status === 'paused') { + await loggingSession.waitForPostExecution() if (!result.snapshotSeed) { logger.error('Missing snapshot seed for paused execution', { executionId }) await loggingSession.markAsFailed('Missing snapshot seed for paused execution') From 10bc116347024ad2b447d22bac12d52fca0e0fa5 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 20:01:17 -0700 Subject: [PATCH 2/3] fix(execution): fail a paused run whose log was never finalized instead of publishing it --- .../executor/pause-persistence.integration.ts | 107 +++++++++++------- .../workflows/executor/pause-persistence.ts | 16 ++- 2 files changed, 77 insertions(+), 46 deletions(-) diff --git a/apps/sim/lib/workflows/executor/pause-persistence.integration.ts b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts index 5d5b8dc74fa..f4592e36875 100644 --- a/apps/sim/lib/workflows/executor/pause-persistence.integration.ts +++ b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts @@ -1,11 +1,10 @@ /** - * Pause publication against real PostgreSQL: a paused run becomes resumable only after - * its log has been finalized out of `running`, so an immediate resume finds a claimable log. + * Pause publication against real PostgreSQL: a paused run becomes resumable only once its + * log has been finalized out of `running`, so an immediate resume finds a claimable log. */ import { db } from '@sim/db' import { pausedExecutions, - resumeQueue, user, workflow, workflowExecutionLogs, @@ -13,10 +12,9 @@ import { workspace, } from '@sim/db/schema' import { createDeferred } from '@sim/testing' -import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import { eq } from 'drizzle-orm' -import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest' import { type BillingAttributionSnapshot, resolveBillingAttribution, @@ -31,14 +29,10 @@ const ids = { owner: `pause-publish-owner-${generateId()}`, workspace: generateId(), workflow: generateId(), - execution: generateId(), } const CONTEXT_ID = 'approval' -/** Long enough for an unguarded publish to commit; a guarded one never resolves while held. */ -const UNGUARDED_PUBLISH_WINDOW_MS = 250 - const workflowState: WorkflowState = { blocks: { start: { @@ -56,7 +50,10 @@ const workflowState: WorkflowState = { parallels: {}, } -function pausedResult(billingAttribution: BillingAttributionSnapshot): ExecutionResult { +function pausedResult( + executionId: string, + billingAttribution: BillingAttributionSnapshot +): ExecutionResult { return { success: true, output: {}, @@ -77,7 +74,7 @@ function pausedResult(billingAttribution: BillingAttributionSnapshot): Execution metadata: { workflowId: ids.workflow, workspaceId: ids.workspace, - executionId: ids.execution, + executionId, userId: ids.owner, billingAttribution, }, @@ -87,17 +84,34 @@ function pausedResult(billingAttribution: BillingAttributionSnapshot): Execution } } -async function logStatus() { +/** Starts a run whose log is `running`, as the core leaves it when execution returns. */ +async function startRun() { + const executionId = generateId() + const billingAttribution = await resolveBillingAttribution({ + actorUserId: ids.owner, + workspaceId: ids.workspace, + }) + const loggingSession = new LoggingSession(ids.workflow, executionId, 'api', 'pause-publish') + await loggingSession.safeStart({ + userId: ids.owner, + workspaceId: ids.workspace, + billingAttribution, + workflowState, + }) + return { executionId, loggingSession, result: pausedResult(executionId, billingAttribution) } +} + +async function logStatus(executionId: string) { const [row] = await db .select({ status: workflowExecutionLogs.status }) .from(workflowExecutionLogs) - .where(eq(workflowExecutionLogs.executionId, ids.execution)) + .where(eq(workflowExecutionLogs.executionId, executionId)) return row?.status } -function resume() { +function resume(executionId: string) { return PauseResumeManager.enqueueOrStartResume({ - executionId: ids.execution, + executionId, workflowId: ids.workflow, contextId: CONTEXT_ID, resumeInput: {}, @@ -132,8 +146,12 @@ beforeAll(async () => { }) }) +afterEach(() => { + vi.restoreAllMocks() +}) + afterAll(async () => { - await db.delete(resumeQueue).where(eq(resumeQueue.parentExecutionId, ids.execution)) + // Deleting a paused execution cascades to its resume queue entries. await db.delete(pausedExecutions).where(eq(pausedExecutions.workflowId, ids.workflow)) await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow)) await db @@ -144,19 +162,8 @@ afterAll(async () => { }) describe('handlePostExecutionPauseState', () => { - it('publishes a pause only after the run log is finalized, so an immediate resume finds a claimable log', async () => { - const billingAttribution = await resolveBillingAttribution({ - actorUserId: ids.owner, - workspaceId: ids.workspace, - }) - const loggingSession = new LoggingSession(ids.workflow, ids.execution, 'api', 'pause-publish') - await loggingSession.safeStart({ - userId: ids.owner, - workspaceId: ids.workspace, - billingAttribution, - workflowState, - }) - expect(await logStatus()).toBe('running') + it('publishes a pause only after its log is finalized, so an immediate resume finds a claimable log', async () => { + const { executionId, loggingSession, result } = await startRun() /** Holds the core's background log finalizer open, as a slow trace projection would. */ const finalizer = createDeferred() @@ -164,24 +171,44 @@ describe('handlePostExecutionPauseState', () => { finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] })) ) + const persistPauseResult = PauseResumeManager.persistPauseResult + const logFinalizedAtPublish: boolean[] = [] + vi.spyOn(PauseResumeManager, 'persistPauseResult').mockImplementation((args) => { + logFinalizedAtPublish.push(loggingSession.hasCompleted()) + return persistPauseResult.call(PauseResumeManager, args) + }) + const publish = handlePostExecutionPauseState({ - result: pausedResult(billingAttribution), + result, workflowId: ids.workflow, - executionId: ids.execution, + executionId, + loggingSession, + }) + finalizer.resolve() + await publish + + expect(logFinalizedAtPublish).toEqual([true]) + expect(await logStatus(executionId)).toBe('pending') + await expect(resume(executionId)).resolves.toMatchObject({ status: 'starting' }) + }) + + it('fails the run instead of publishing a pause whose log was never finalized', async () => { + const { executionId, loggingSession, result } = await startRun() + + /** The core's finalizer swallows its own failures, so a lost pause write still settles. */ + loggingSession.setPostExecutionPromise(Promise.resolve()) + + await handlePostExecutionPauseState({ + result, + workflowId: ids.workflow, + executionId, loggingSession, }) - await Promise.race([publish, sleep(UNGUARDED_PUBLISH_WINDOW_MS)]) - expect(await logStatus()).toBe('running') - await expect(resume()).rejects.toMatchObject({ + expect(await logStatus(executionId)).toBe('failed') + await expect(resume(executionId)).rejects.toMatchObject({ name: 'ResumeAdmissionError', statusCode: 404, }) - - finalizer.resolve() - await publish - - expect(await logStatus()).toBe('pending') - await expect(resume()).resolves.toMatchObject({ status: 'starting' }) }) }) diff --git a/apps/sim/lib/workflows/executor/pause-persistence.ts b/apps/sim/lib/workflows/executor/pause-persistence.ts index 75092c65449..f7a83da52c4 100644 --- a/apps/sim/lib/workflows/executor/pause-persistence.ts +++ b/apps/sim/lib/workflows/executor/pause-persistence.ts @@ -21,14 +21,15 @@ interface HandlePostExecutionPauseStateArgs { * Every caller of `executeWorkflowCore` must call this after execution completes * to ensure HITL pause state is persisted to the database and queued resumes are drained. * - * - If execution is paused with a valid snapshot: persists to `paused_executions` table + * - If execution is paused but its log was never finalized: marks execution as failed * - If execution is paused without a snapshot: marks execution as failed + * - If execution is paused with a valid snapshot: persists to `paused_executions` table * - If execution is not paused: processes any queued resume entries * - * A pause is published only after the core's post-execution logging settles. The - * core finalizes the run log in the background, and a resume claims that log only - * once it has left `running`, so publishing first lets an immediate resume be - * rejected as no longer resumable. + * A pause is published only after the core's post-execution logging has persisted + * the paused log. A resume claims that log only once it has left `running`, so a + * pause published before then (or without it ever happening) is rejected as no + * longer resumable. */ export async function handlePostExecutionPauseState({ result, @@ -39,7 +40,10 @@ export async function handlePostExecutionPauseState({ }: HandlePostExecutionPauseStateArgs): Promise { if (result.status === 'paused') { await loggingSession.waitForPostExecution() - if (!result.snapshotSeed) { + if (!loggingSession.hasCompleted()) { + logger.error('Paused execution log was not finalized', { executionId }) + await loggingSession.markAsFailed('Failed to record paused execution') + } else if (!result.snapshotSeed) { logger.error('Missing snapshot seed for paused execution', { executionId }) await loggingSession.markAsFailed('Missing snapshot seed for paused execution') } else { From 81b2e99fe42f70a9911f70506642438f44dd7cf9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 20:08:54 -0700 Subject: [PATCH 3/3] fix(execution): restore the publish spy explicitly instead of in a hook --- .../executor/pause-persistence.integration.ts | 38 ++++++++++--------- 1 file changed, 20 insertions(+), 18 deletions(-) diff --git a/apps/sim/lib/workflows/executor/pause-persistence.integration.ts b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts index f4592e36875..1d5c8eeafd9 100644 --- a/apps/sim/lib/workflows/executor/pause-persistence.integration.ts +++ b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts @@ -14,7 +14,7 @@ import { import { createDeferred } from '@sim/testing' import { generateId } from '@sim/utils/id' import { eq } from 'drizzle-orm' -import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' import { type BillingAttributionSnapshot, resolveBillingAttribution, @@ -146,10 +146,6 @@ beforeAll(async () => { }) }) -afterEach(() => { - vi.restoreAllMocks() -}) - afterAll(async () => { // Deleting a paused execution cascades to its resume queue entries. await db.delete(pausedExecutions).where(eq(pausedExecutions.workflowId, ids.workflow)) @@ -173,19 +169,25 @@ describe('handlePostExecutionPauseState', () => { const persistPauseResult = PauseResumeManager.persistPauseResult const logFinalizedAtPublish: boolean[] = [] - vi.spyOn(PauseResumeManager, 'persistPauseResult').mockImplementation((args) => { - logFinalizedAtPublish.push(loggingSession.hasCompleted()) - return persistPauseResult.call(PauseResumeManager, args) - }) - - const publish = handlePostExecutionPauseState({ - result, - workflowId: ids.workflow, - executionId, - loggingSession, - }) - finalizer.resolve() - await publish + const publishSpy = vi + .spyOn(PauseResumeManager, 'persistPauseResult') + .mockImplementation((args) => { + logFinalizedAtPublish.push(loggingSession.hasCompleted()) + return persistPauseResult.call(PauseResumeManager, args) + }) + + try { + const publish = handlePostExecutionPauseState({ + result, + workflowId: ids.workflow, + executionId, + loggingSession, + }) + finalizer.resolve() + await publish + } finally { + publishSpy.mockRestore() + } expect(logFinalizedAtPublish).toEqual([true]) expect(await logStatus(executionId)).toBe('pending')