Skip to content

Commit 10bc116

Browse files
committed
fix(execution): fail a paused run whose log was never finalized instead of publishing it
1 parent ac35656 commit 10bc116

2 files changed

Lines changed: 77 additions & 46 deletions

File tree

‎apps/sim/lib/workflows/executor/pause-persistence.integration.ts‎

Lines changed: 67 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -1,22 +1,20 @@
11
/**
2-
* Pause publication against real PostgreSQL: a paused run becomes resumable only after
3-
* its log has been finalized out of `running`, so an immediate resume finds a claimable log.
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.
44
*/
55
import { db } from '@sim/db'
66
import {
77
pausedExecutions,
8-
resumeQueue,
98
user,
109
workflow,
1110
workflowExecutionLogs,
1211
workflowExecutionSnapshots,
1312
workspace,
1413
} from '@sim/db/schema'
1514
import { createDeferred } from '@sim/testing'
16-
import { sleep } from '@sim/utils/helpers'
1715
import { generateId } from '@sim/utils/id'
1816
import { eq } from 'drizzle-orm'
19-
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
17+
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
2018
import {
2119
type BillingAttributionSnapshot,
2220
resolveBillingAttribution,
@@ -31,14 +29,10 @@ const ids = {
3129
owner: `pause-publish-owner-${generateId()}`,
3230
workspace: generateId(),
3331
workflow: generateId(),
34-
execution: generateId(),
3532
}
3633

3734
const CONTEXT_ID = 'approval'
3835

39-
/** Long enough for an unguarded publish to commit; a guarded one never resolves while held. */
40-
const UNGUARDED_PUBLISH_WINDOW_MS = 250
41-
4236
const workflowState: WorkflowState = {
4337
blocks: {
4438
start: {
@@ -56,7 +50,10 @@ const workflowState: WorkflowState = {
5650
parallels: {},
5751
}
5852

59-
function pausedResult(billingAttribution: BillingAttributionSnapshot): ExecutionResult {
53+
function pausedResult(
54+
executionId: string,
55+
billingAttribution: BillingAttributionSnapshot
56+
): ExecutionResult {
6057
return {
6158
success: true,
6259
output: {},
@@ -77,7 +74,7 @@ function pausedResult(billingAttribution: BillingAttributionSnapshot): Execution
7774
metadata: {
7875
workflowId: ids.workflow,
7976
workspaceId: ids.workspace,
80-
executionId: ids.execution,
77+
executionId,
8178
userId: ids.owner,
8279
billingAttribution,
8380
},
@@ -87,17 +84,34 @@ function pausedResult(billingAttribution: BillingAttributionSnapshot): Execution
8784
}
8885
}
8986

90-
async function logStatus() {
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) {
91105
const [row] = await db
92106
.select({ status: workflowExecutionLogs.status })
93107
.from(workflowExecutionLogs)
94-
.where(eq(workflowExecutionLogs.executionId, ids.execution))
108+
.where(eq(workflowExecutionLogs.executionId, executionId))
95109
return row?.status
96110
}
97111

98-
function resume() {
112+
function resume(executionId: string) {
99113
return PauseResumeManager.enqueueOrStartResume({
100-
executionId: ids.execution,
114+
executionId,
101115
workflowId: ids.workflow,
102116
contextId: CONTEXT_ID,
103117
resumeInput: {},
@@ -132,8 +146,12 @@ beforeAll(async () => {
132146
})
133147
})
134148

149+
afterEach(() => {
150+
vi.restoreAllMocks()
151+
})
152+
135153
afterAll(async () => {
136-
await db.delete(resumeQueue).where(eq(resumeQueue.parentExecutionId, ids.execution))
154+
// Deleting a paused execution cascades to its resume queue entries.
137155
await db.delete(pausedExecutions).where(eq(pausedExecutions.workflowId, ids.workflow))
138156
await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow))
139157
await db
@@ -144,44 +162,53 @@ afterAll(async () => {
144162
})
145163

146164
describe('handlePostExecutionPauseState', () => {
147-
it('publishes a pause only after the run log is finalized, so an immediate resume finds a claimable log', async () => {
148-
const billingAttribution = await resolveBillingAttribution({
149-
actorUserId: ids.owner,
150-
workspaceId: ids.workspace,
151-
})
152-
const loggingSession = new LoggingSession(ids.workflow, ids.execution, 'api', 'pause-publish')
153-
await loggingSession.safeStart({
154-
userId: ids.owner,
155-
workspaceId: ids.workspace,
156-
billingAttribution,
157-
workflowState,
158-
})
159-
expect(await logStatus()).toBe('running')
165+
it('publishes a pause only after its log is finalized, so an immediate resume finds a claimable log', async () => {
166+
const { executionId, loggingSession, result } = await startRun()
160167

161168
/** Holds the core's background log finalizer open, as a slow trace projection would. */
162169
const finalizer = createDeferred<void>()
163170
loggingSession.setPostExecutionPromise(
164171
finalizer.promise.then(() => loggingSession.safeCompleteWithPause({ traceSpans: [] }))
165172
)
166173

174+
const persistPauseResult = PauseResumeManager.persistPauseResult
175+
const logFinalizedAtPublish: boolean[] = []
176+
vi.spyOn(PauseResumeManager, 'persistPauseResult').mockImplementation((args) => {
177+
logFinalizedAtPublish.push(loggingSession.hasCompleted())
178+
return persistPauseResult.call(PauseResumeManager, args)
179+
})
180+
167181
const publish = handlePostExecutionPauseState({
168-
result: pausedResult(billingAttribution),
182+
result,
169183
workflowId: ids.workflow,
170-
executionId: ids.execution,
184+
executionId,
185+
loggingSession,
186+
})
187+
finalizer.resolve()
188+
await publish
189+
190+
expect(logFinalizedAtPublish).toEqual([true])
191+
expect(await logStatus(executionId)).toBe('pending')
192+
await expect(resume(executionId)).resolves.toMatchObject({ status: 'starting' })
193+
})
194+
195+
it('fails the run instead of publishing a pause whose log was never finalized', async () => {
196+
const { executionId, loggingSession, result } = await startRun()
197+
198+
/** The core's finalizer swallows its own failures, so a lost pause write still settles. */
199+
loggingSession.setPostExecutionPromise(Promise.resolve())
200+
201+
await handlePostExecutionPauseState({
202+
result,
203+
workflowId: ids.workflow,
204+
executionId,
171205
loggingSession,
172206
})
173-
await Promise.race([publish, sleep(UNGUARDED_PUBLISH_WINDOW_MS)])
174207

175-
expect(await logStatus()).toBe('running')
176-
await expect(resume()).rejects.toMatchObject({
208+
expect(await logStatus(executionId)).toBe('failed')
209+
await expect(resume(executionId)).rejects.toMatchObject({
177210
name: 'ResumeAdmissionError',
178211
statusCode: 404,
179212
})
180-
181-
finalizer.resolve()
182-
await publish
183-
184-
expect(await logStatus()).toBe('pending')
185-
await expect(resume()).resolves.toMatchObject({ status: 'starting' })
186213
})
187214
})

‎apps/sim/lib/workflows/executor/pause-persistence.ts‎

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -21,14 +21,15 @@ interface HandlePostExecutionPauseStateArgs {
2121
* Every caller of `executeWorkflowCore` must call this after execution completes
2222
* to ensure HITL pause state is persisted to the database and queued resumes are drained.
2323
*
24-
* - If execution is paused with a valid snapshot: persists to `paused_executions` table
24+
* - If execution is paused but its log was never finalized: marks execution as failed
2525
* - If execution is paused without a snapshot: marks execution as failed
26+
* - If execution is paused with a valid snapshot: persists to `paused_executions` table
2627
* - If execution is not paused: processes any queued resume entries
2728
*
28-
* A pause is published only after the core's post-execution logging settles. The
29-
* core finalizes the run log in the background, and a resume claims that log only
30-
* once it has left `running`, so publishing first lets an immediate resume be
31-
* rejected as no longer resumable.
29+
* A pause is published only after the core's post-execution logging has persisted
30+
* the paused log. A resume claims that log only once it has left `running`, so a
31+
* pause published before then (or without it ever happening) is rejected as no
32+
* longer resumable.
3233
*/
3334
export async function handlePostExecutionPauseState({
3435
result,
@@ -39,7 +40,10 @@ export async function handlePostExecutionPauseState({
3940
}: HandlePostExecutionPauseStateArgs): Promise<void> {
4041
if (result.status === 'paused') {
4142
await loggingSession.waitForPostExecution()
42-
if (!result.snapshotSeed) {
43+
if (!loggingSession.hasCompleted()) {
44+
logger.error('Paused execution log was not finalized', { executionId })
45+
await loggingSession.markAsFailed('Failed to record paused execution')
46+
} else if (!result.snapshotSeed) {
4347
logger.error('Missing snapshot seed for paused execution', { executionId })
4448
await loggingSession.markAsFailed('Missing snapshot seed for paused execution')
4549
} else {

0 commit comments

Comments
 (0)