|
| 1 | +/** |
| 2 | + * Pause publication against real PostgreSQL: a paused run becomes resumable only once its |
| 3 | + * log has been finalized out of `running`, so an immediate resume finds a claimable log. |
| 4 | + */ |
| 5 | +import { db } from '@sim/db' |
| 6 | +import { |
| 7 | + pausedExecutions, |
| 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 } from 'drizzle-orm' |
| 17 | +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' |
| 18 | +import { |
| 19 | + type BillingAttributionSnapshot, |
| 20 | + resolveBillingAttribution, |
| 21 | +} from '@/lib/billing/core/billing-attribution' |
| 22 | +import { LoggingSession } from '@/lib/logs/execution/logging-session' |
| 23 | +import type { WorkflowState } from '@/lib/logs/types' |
| 24 | +import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager' |
| 25 | +import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence' |
| 26 | +import type { ExecutionResult } from '@/executor/types' |
| 27 | + |
| 28 | +const ids = { |
| 29 | + owner: `pause-publish-owner-${generateId()}`, |
| 30 | + workspace: generateId(), |
| 31 | + workflow: generateId(), |
| 32 | +} |
| 33 | + |
| 34 | +const CONTEXT_ID = 'approval' |
| 35 | + |
| 36 | +const workflowState: WorkflowState = { |
| 37 | + blocks: { |
| 38 | + start: { |
| 39 | + id: 'start', |
| 40 | + type: 'starter', |
| 41 | + name: 'Start', |
| 42 | + position: { x: 0, y: 0 }, |
| 43 | + subBlocks: {}, |
| 44 | + outputs: {}, |
| 45 | + enabled: true, |
| 46 | + }, |
| 47 | + }, |
| 48 | + edges: [], |
| 49 | + loops: {}, |
| 50 | + parallels: {}, |
| 51 | +} |
| 52 | + |
| 53 | +function pausedResult( |
| 54 | + executionId: string, |
| 55 | + billingAttribution: BillingAttributionSnapshot |
| 56 | +): ExecutionResult { |
| 57 | + return { |
| 58 | + success: true, |
| 59 | + output: {}, |
| 60 | + status: 'paused', |
| 61 | + pausePoints: [ |
| 62 | + { |
| 63 | + contextId: CONTEXT_ID, |
| 64 | + blockId: CONTEXT_ID, |
| 65 | + response: {}, |
| 66 | + registeredAt: new Date().toISOString(), |
| 67 | + resumeStatus: 'paused', |
| 68 | + snapshotReady: true, |
| 69 | + pauseKind: 'human', |
| 70 | + }, |
| 71 | + ], |
| 72 | + snapshotSeed: { |
| 73 | + snapshot: JSON.stringify({ |
| 74 | + metadata: { |
| 75 | + workflowId: ids.workflow, |
| 76 | + workspaceId: ids.workspace, |
| 77 | + executionId, |
| 78 | + userId: ids.owner, |
| 79 | + billingAttribution, |
| 80 | + }, |
| 81 | + }), |
| 82 | + triggerIds: [], |
| 83 | + }, |
| 84 | + } |
| 85 | +} |
| 86 | + |
| 87 | +/** Starts a run whose log is `running`, as the core leaves it when execution returns. */ |
| 88 | +async function startRun() { |
| 89 | + const executionId = generateId() |
| 90 | + const billingAttribution = await resolveBillingAttribution({ |
| 91 | + actorUserId: ids.owner, |
| 92 | + workspaceId: ids.workspace, |
| 93 | + }) |
| 94 | + const loggingSession = new LoggingSession(ids.workflow, executionId, 'api', 'pause-publish') |
| 95 | + await loggingSession.safeStart({ |
| 96 | + userId: ids.owner, |
| 97 | + workspaceId: ids.workspace, |
| 98 | + billingAttribution, |
| 99 | + workflowState, |
| 100 | + }) |
| 101 | + return { executionId, loggingSession, result: pausedResult(executionId, billingAttribution) } |
| 102 | +} |
| 103 | + |
| 104 | +async function logStatus(executionId: string) { |
| 105 | + const [row] = await db |
| 106 | + .select({ status: workflowExecutionLogs.status }) |
| 107 | + .from(workflowExecutionLogs) |
| 108 | + .where(eq(workflowExecutionLogs.executionId, executionId)) |
| 109 | + return row?.status |
| 110 | +} |
| 111 | + |
| 112 | +function resume(executionId: string) { |
| 113 | + return PauseResumeManager.enqueueOrStartResume({ |
| 114 | + executionId, |
| 115 | + workflowId: ids.workflow, |
| 116 | + contextId: CONTEXT_ID, |
| 117 | + resumeInput: {}, |
| 118 | + userId: ids.owner, |
| 119 | + }) |
| 120 | +} |
| 121 | + |
| 122 | +beforeAll(async () => { |
| 123 | + const now = new Date() |
| 124 | + await db.insert(user).values({ |
| 125 | + id: ids.owner, |
| 126 | + name: 'Pause Publish', |
| 127 | + email: `${ids.owner}@pause-publish.test`, |
| 128 | + emailVerified: true, |
| 129 | + createdAt: now, |
| 130 | + updatedAt: now, |
| 131 | + }) |
| 132 | + await db.insert(workspace).values({ |
| 133 | + id: ids.workspace, |
| 134 | + name: 'Pause Publish', |
| 135 | + ownerId: ids.owner, |
| 136 | + billedAccountUserId: ids.owner, |
| 137 | + }) |
| 138 | + await db.insert(workflow).values({ |
| 139 | + id: ids.workflow, |
| 140 | + userId: ids.owner, |
| 141 | + workspaceId: ids.workspace, |
| 142 | + name: 'Pause Publish', |
| 143 | + lastSynced: now, |
| 144 | + createdAt: now, |
| 145 | + updatedAt: now, |
| 146 | + }) |
| 147 | +}) |
| 148 | + |
| 149 | +afterAll(async () => { |
| 150 | + // Deleting a paused execution cascades to its resume queue entries. |
| 151 | + await db.delete(pausedExecutions).where(eq(pausedExecutions.workflowId, ids.workflow)) |
| 152 | + await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow)) |
| 153 | + await db |
| 154 | + .delete(workflowExecutionSnapshots) |
| 155 | + .where(eq(workflowExecutionSnapshots.workflowId, ids.workflow)) |
| 156 | + await db.delete(workspace).where(eq(workspace.id, ids.workspace)) |
| 157 | + await db.delete(user).where(eq(user.id, ids.owner)) |
| 158 | +}) |
| 159 | + |
| 160 | +describe('handlePostExecutionPauseState', () => { |
| 161 | + it('publishes a pause only after its log is finalized, so an immediate resume finds a claimable log', async () => { |
| 162 | + const { executionId, loggingSession, result } = await startRun() |
| 163 | + |
| 164 | + /** Holds the core's background log finalizer open, as a slow trace projection would. */ |
| 165 | + const finalizer = createDeferred<void>() |
| 166 | + loggingSession.setPostExecutionPromise( |
| 167 | + finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] })) |
| 168 | + ) |
| 169 | + |
| 170 | + const persistPauseResult = PauseResumeManager.persistPauseResult |
| 171 | + const logFinalizedAtPublish: boolean[] = [] |
| 172 | + const publishSpy = vi |
| 173 | + .spyOn(PauseResumeManager, 'persistPauseResult') |
| 174 | + .mockImplementation((args) => { |
| 175 | + logFinalizedAtPublish.push(loggingSession.hasCompleted()) |
| 176 | + return persistPauseResult.call(PauseResumeManager, args) |
| 177 | + }) |
| 178 | + |
| 179 | + try { |
| 180 | + const publish = handlePostExecutionPauseState({ |
| 181 | + result, |
| 182 | + workflowId: ids.workflow, |
| 183 | + executionId, |
| 184 | + loggingSession, |
| 185 | + }) |
| 186 | + finalizer.resolve() |
| 187 | + await publish |
| 188 | + } finally { |
| 189 | + publishSpy.mockRestore() |
| 190 | + } |
| 191 | + |
| 192 | + expect(logFinalizedAtPublish).toEqual([true]) |
| 193 | + expect(await logStatus(executionId)).toBe('pending') |
| 194 | + await expect(resume(executionId)).resolves.toMatchObject({ status: 'starting' }) |
| 195 | + }) |
| 196 | + |
| 197 | + it('fails the run instead of publishing a pause whose log was never finalized', async () => { |
| 198 | + const { executionId, loggingSession, result } = await startRun() |
| 199 | + |
| 200 | + /** The core's finalizer swallows its own failures, so a lost pause write still settles. */ |
| 201 | + loggingSession.setPostExecutionPromise(Promise.resolve()) |
| 202 | + |
| 203 | + await handlePostExecutionPauseState({ |
| 204 | + result, |
| 205 | + workflowId: ids.workflow, |
| 206 | + executionId, |
| 207 | + loggingSession, |
| 208 | + }) |
| 209 | + |
| 210 | + expect(await logStatus(executionId)).toBe('failed') |
| 211 | + await expect(resume(executionId)).rejects.toMatchObject({ |
| 212 | + name: 'ResumeAdmissionError', |
| 213 | + statusCode: 404, |
| 214 | + }) |
| 215 | + }) |
| 216 | +}) |
0 commit comments