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..1d5c8eeafd9 --- /dev/null +++ b/apps/sim/lib/workflows/executor/pause-persistence.integration.ts @@ -0,0 +1,216 @@ +/** + * 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, + user, + workflow, + workflowExecutionLogs, + 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, vi } 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(), +} + +const CONTEXT_ID = 'approval' + +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( + executionId: string, + 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, + userId: ids.owner, + billingAttribution, + }, + }), + triggerIds: [], + }, + } +} + +/** 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, executionId)) + return row?.status +} + +function resume(executionId: string) { + return PauseResumeManager.enqueueOrStartResume({ + executionId, + 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 () => { + // 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 + .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 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() + loggingSession.setPostExecutionPromise( + finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] })) + ) + + const persistPauseResult = PauseResumeManager.persistPauseResult + const logFinalizedAtPublish: boolean[] = [] + 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') + 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, + }) + + expect(await logStatus(executionId)).toBe('failed') + await expect(resume(executionId)).rejects.toMatchObject({ + name: 'ResumeAdmissionError', + statusCode: 404, + }) + }) +}) diff --git a/apps/sim/lib/workflows/executor/pause-persistence.ts b/apps/sim/lib/workflows/executor/pause-persistence.ts index 552263f490f..f7a83da52c4 100644 --- a/apps/sim/lib/workflows/executor/pause-persistence.ts +++ b/apps/sim/lib/workflows/executor/pause-persistence.ts @@ -21,9 +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 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, @@ -33,7 +39,11 @@ export async function handlePostExecutionPauseState({ loggingSession, }: HandlePostExecutionPauseStateArgs): Promise { if (result.status === 'paused') { - if (!result.snapshotSeed) { + await loggingSession.waitForPostExecution() + 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 {