Skip to content

Commit aa14823

Browse files
committed
fix(tables): settle a resume whose pause cannot be saved as failed, and never fail a completed run
A resumed run that paused but whose pause state could not be persisted failed its log yet returned a paused result, so the cell showed paused and the resume entry was marked completed. It now throws after failing the log, so the attempt settles as failed and reports execution_failed. markResumeFailed also rewrote a completed log as failed when a step after a completed run threw. A completed run's outcome now stands, and an already failed log keeps its original end time.
1 parent 7e11748 commit aa14823

2 files changed

Lines changed: 98 additions & 25 deletions

File tree

‎apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts‎

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
2222
import { createTimeoutAbortController, getExecutionDeadlineAt } from '@/lib/core/execution-limits'
2323
import { abortManualExecution } from '@/lib/execution/manual-cancellation'
2424
import { terminalExecutionLogFields } from '@/lib/logs/execution/cancellation'
25+
import { LoggingSession } from '@/lib/logs/execution/logging-session'
2526

2627
const {
2728
mockExecuteWorkflowCore,
@@ -265,6 +266,8 @@ describe('what a failed resume did to its paused execution', () => {
265266
{ logStatus: 'running', pauseStatus: 'paused', executionFailed: true },
266267
{ logStatus: 'cancelled', pauseStatus: 'paused', executionFailed: false },
267268
{ logStatus: 'running', pauseStatus: 'cancelling', executionFailed: false },
269+
{ logStatus: 'failed', pauseStatus: 'paused', executionFailed: true },
270+
{ logStatus: 'completed', pauseStatus: 'paused', executionFailed: false },
268271
])(
269272
'reports the execution failed: $executionFailed for a $logStatus log and $pauseStatus pause',
270273
async ({ logStatus, pauseStatus, executionFailed }) => {
@@ -276,6 +279,23 @@ describe('what a failed resume did to its paused execution', () => {
276279
}
277280
)
278281

282+
it.each([
283+
{ logStatus: 'completed', updated: [resumeQueue] },
284+
{ logStatus: 'failed', updated: [resumeQueue, pausedExecutions] },
285+
{ logStatus: 'running', updated: [resumeQueue, pausedExecutions, workflowExecutionLogs] },
286+
])(
287+
'leaves a $logStatus execution log as it is when a resume fails late',
288+
async ({ logStatus, updated }) => {
289+
queueTableRows(workflowExecutionLogs, [{ status: logStatus }])
290+
queueTableRows(pausedExecutions, [{ status: 'paused' }])
291+
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals
292+
293+
await managerInternals.markResumeFailed(attemptArgs)
294+
295+
expect(dbChainMockFns.update.mock.calls.map(([table]) => table)).toEqual(updated)
296+
}
297+
)
298+
279299
/** Resume args that collect every outcome the manager reports. */
280300
function argsReportingOutcomes(onAttemptFailed?: () => Promise<void>) {
281301
const outcomes: FailedResumeOutcome[] = []
@@ -313,6 +333,59 @@ describe('what a failed resume did to its paused execution', () => {
313333
}
314334
)
315335

336+
describe('when the resumed run pauses but its pause cannot be saved', () => {
337+
const spies: { mockRestore: () => void }[] = []
338+
339+
function pauseRun(options: { snapshotSeed?: unknown; persistError?: Error }) {
340+
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals
341+
spies.push(
342+
vi.spyOn(managerInternals, 'runResumeExecution').mockResolvedValueOnce({
343+
success: true,
344+
status: 'paused',
345+
output: {},
346+
logs: [],
347+
pausePoints: [],
348+
snapshotSeed: options.snapshotSeed,
349+
metadata: { executionId: 'parent-execution-1', duration: 1, startTime: 'start' },
350+
}),
351+
vi.spyOn(managerInternals, 'markResumeFailed').mockResolvedValueOnce(true),
352+
vi.spyOn(LoggingSession, 'markExecutionAsFailed').mockResolvedValueOnce(),
353+
vi.spyOn(PauseResumeManager, 'processQueuedResumes').mockResolvedValueOnce()
354+
)
355+
if (options.persistError) {
356+
spies.push(
357+
vi
358+
.spyOn(PauseResumeManager, 'persistPauseResult')
359+
.mockRejectedValueOnce(options.persistError)
360+
)
361+
}
362+
}
363+
364+
afterEach(() => {
365+
for (const spy of spies.splice(0)) spy.mockRestore()
366+
})
367+
368+
it('fails the attempt when the pause state cannot be persisted', async () => {
369+
pauseRun({ snapshotSeed: createSnapshotSeed(), persistError: new Error('lock timeout') })
370+
const { outcomes, args } = argsReportingOutcomes()
371+
372+
await expect(PauseResumeManager.startResumeExecution(args)).rejects.toThrow(
373+
'Failed to persist pause state: lock timeout'
374+
)
375+
expect(outcomes).toEqual(['execution_failed'])
376+
})
377+
378+
it('fails the attempt when the paused run has no snapshot seed', async () => {
379+
pauseRun({})
380+
const { outcomes, args } = argsReportingOutcomes()
381+
382+
await expect(PauseResumeManager.startResumeExecution(args)).rejects.toThrow(
383+
'Missing snapshot seed for paused execution'
384+
)
385+
expect(outcomes).toEqual(['execution_failed'])
386+
})
387+
})
388+
316389
describe('when the resumed run fails', () => {
317390
const rawError = new Error('Block failed')
318391
const spies: { mockRestore: () => void }[] = []

‎apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts‎

Lines changed: 25 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -924,18 +924,22 @@ export class PauseResumeManager {
924924
})
925925

926926
if (result.status === 'paused') {
927+
/**
928+
* A pause that cannot be saved fails the execution. Fail the log with the
929+
* reason, then throw so the attempt settles as failed below.
930+
*/
927931
const effectiveExecutionId = result.metadata?.executionId ?? resumeExecutionId
928-
if (!result.snapshotSeed) {
929-
logger.error('Missing snapshot seed for paused resume execution', {
930-
resumeExecutionId,
931-
})
932+
const failPause = async (message: string, cause?: unknown): Promise<never> => {
932933
await LoggingSession.markExecutionAsFailed(
933934
effectiveExecutionId,
934-
'Missing snapshot seed for paused execution',
935+
message,
935936
undefined,
936937
pausedExecution.workflowId
937938
)
938-
await releaseExecutionSlot(resumeEntryId)
939+
throw new Error(message, { cause })
940+
}
941+
if (!result.snapshotSeed) {
942+
await failPause('Missing snapshot seed for paused execution')
939943
} else {
940944
try {
941945
await PauseResumeManager.persistPauseResult({
@@ -947,19 +951,10 @@ export class PauseResumeManager {
947951
executorUserId: result.metadata?.userId,
948952
})
949953
} catch (pauseError) {
950-
logger.error(
951-
'Failed to persist pause result for resumed execution',
952-
projectResolvedSecretDiagnosticError(pauseError, undefined, {
953-
resumeExecutionId,
954-
})
955-
)
956-
await LoggingSession.markExecutionAsFailed(
957-
effectiveExecutionId,
954+
await failPause(
958955
`Failed to persist pause state: ${toError(pauseError).message}`,
959-
undefined,
960-
pausedExecution.workflowId
956+
pauseError
961957
)
962-
await releaseExecutionSlot(resumeEntryId)
963958
}
964959
}
965960
} else {
@@ -2268,6 +2263,9 @@ export class PauseResumeManager {
22682263
return false
22692264
}
22702265

2266+
/** The run completed before a later step threw; its outcome stands. */
2267+
if (executionLog?.status === 'completed') return false
2268+
22712269
await tx
22722270
.update(pausedExecutions)
22732271
.set({
@@ -2278,15 +2276,17 @@ export class PauseResumeManager {
22782276

22792277
if (pausedExecution?.status === 'cancelling') return false
22802278

2281-
await tx
2282-
.update(workflowExecutionLogs)
2283-
.set(terminalExecutionLogFields('failed', now))
2284-
.where(
2285-
and(
2286-
eq(workflowExecutionLogs.executionId, args.parentExecutionId),
2287-
sql`${workflowExecutionLogs.status} != 'cancelled'`
2279+
if (executionLog?.status !== 'failed') {
2280+
await tx
2281+
.update(workflowExecutionLogs)
2282+
.set(terminalExecutionLogFields('failed', now))
2283+
.where(
2284+
and(
2285+
eq(workflowExecutionLogs.executionId, args.parentExecutionId),
2286+
sql`${workflowExecutionLogs.status} != 'cancelled'`
2287+
)
22882288
)
2289-
)
2289+
}
22902290

22912291
return true
22922292
})

0 commit comments

Comments
 (0)