Skip to content
57 changes: 46 additions & 11 deletions apps/sim/background/resume-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -178,13 +178,36 @@ export async function executeResumeJob(payload: ResumeExecutionPayload, signal?:
cellContext.rowId,
parentExecutionId,
async () => {
let completedBeforeFailure = false
const result = await runResumeAndCellTerminal(
payload,
pausedExecution,
writers,
attemptSignal,
attemptTimeoutController
)
attemptTimeoutController,
() => {
completedBeforeFailure = true
}
).catch(async (error: unknown) => {
/**
* The run completed and only a later step of the attempt threw, so its
* cell is completed: continue the cascade as a completed run would, and
* still surface the failure.
*/
if (completedBeforeFailure) {
await continueCascadeAfterResume(cellContext, billingAttribution, attemptSignal).catch(
(cascadeError: unknown) => {
logger.error(
'Failed to continue the cascade after a completed resume',
projectResolvedSecretDiagnosticError(cascadeError, undefined, {
resumeExecutionId,
})
)
}
)
}
throw error
})
if (result.status === 'paused' || result.status === 'cancelled') return result
await continueCascadeAfterResume(cellContext, billingAttribution, attemptSignal)
return result
Expand Down Expand Up @@ -364,19 +387,26 @@ async function buildResumeCellWriters(
/**
* A resume that throws never reaches the terminal write in
* {@link runResumeAndCellTerminal}, which would leave the cell showing its last
* partial `running` state. Mirror what the failed attempt did to the execution:
* a pause that stayed resumable goes back to paused, and a failed execution
* fails the cell.
* partial `running` state. Mirror what the failed attempt left the execution
* as: a pause that stayed resumable goes back to paused, a failed execution
* fails the cell, and a run that completed before a later step threw
* completes it.
*/
async function writeFailedResumeCellTerminal(
writers: CellWriters,
outcome: FailedResumeOutcome,
error: unknown
): Promise<void> {
if (outcome === 'pause_retained') {
await writers.writeCellTerminal('paused', null)
} else {
await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed'))
switch (outcome) {
case 'pause_retained':
await writers.writeCellTerminal('paused', null)
return
case 'execution_completed':
await writers.writeCellTerminal('completed', null)
return
Comment thread
waleedlatif1 marked this conversation as resolved.
case 'execution_failed':
await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed'))
return
}
}

Expand All @@ -385,7 +415,8 @@ async function runResumeAndCellTerminal(
pausedExecution: Awaited<ReturnType<typeof PauseResumeManager.getPausedExecutionById>>,
writers: CellWriters,
signal: AbortSignal | undefined,
timeoutController: ReturnType<typeof createTimeoutAbortController>
timeoutController: ReturnType<typeof createTimeoutAbortController>,
onCompletedBeforeFailure?: () => void
): Promise<Awaited<ReturnType<typeof PauseResumeManager.startResumeExecution>>> {
if (!pausedExecution) throw new Error('Paused execution missing — already nulled by caller')
const result = await PauseResumeManager.startResumeExecution({
Expand All @@ -396,7 +427,11 @@ async function runResumeAndCellTerminal(
resumeInput: payload.resumeInput,
userId: payload.userId,
onBlockComplete: writers.cellOnBlockComplete,
onAttemptFailed: (outcome, error) => writeFailedResumeCellTerminal(writers, outcome, error),
onAttemptFailed: async (outcome, error) => {
await writeFailedResumeCellTerminal(writers, outcome, error)
/** Only a cell saved as completed may start its downstream groups. */
if (outcome === 'execution_completed') onCompletedBeforeFailure?.()
},
abortSignal: signal,
})

Expand Down
50 changes: 48 additions & 2 deletions apps/sim/background/resume-governed-subject.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -220,21 +220,34 @@ describe('resuming a paused table cell', () => {
}, 20_000)

describe('when the resume throws', () => {
/** Downstream groups the row's cascade started after the resume. */
let startedGroups: string[]

beforeEach(() => {
startedGroups = []
mocks.runRowCascadeLoop.mockImplementation(async (payload: { groupId: string }) => {
startedGroups.push(payload.groupId)
})
})

/** The execution state the last cell write persisted. */
function lastCellExecutionState() {
const [, payload] = mocks.writeWorkflowGroupState.mock.calls.at(-1) ?? []
return payload?.executionState
}

/** Fails the resume the way the manager does: settle, report the outcome, rethrow. */
/**
* Fails the resume the way the manager does: settle, report the outcome (a
* failing handler is logged, never rethrown), rethrow the attempt's error.
*/
function failResume(outcome: FailedResumeOutcome, error: Error) {
mocks.startResumeExecution.mockImplementationOnce(
async ({
onAttemptFailed,
}: {
onAttemptFailed?: (outcome: FailedResumeOutcome, error: unknown) => Promise<void>
}) => {
await onAttemptFailed?.(outcome, error)
await onAttemptFailed?.(outcome, error).catch(() => undefined)
throw error
}
)
Expand All @@ -253,6 +266,39 @@ describe('resuming a paused table cell', () => {
})
}, 20_000)

it('marks the cell completed when the run completed before a later step failed', async () => {
const bookkeepingFailure = new Error('Database unavailable')
failResume('execution_completed', bookkeepingFailure)

await expect(executeResumeJob(PAYLOAD)).rejects.toBe(bookkeepingFailure)

expect(lastCellExecutionState()).toMatchObject({
status: 'completed',
executionId: 'parent-execution-1',
error: null,
})
expect(startedGroups).toEqual([NEXT_GROUP.id])
}, 20_000)

it('does not continue the cascade when the completed cell could not be saved', async () => {
const bookkeepingFailure = new Error('Database unavailable')
failResume('execution_completed', bookkeepingFailure)
mocks.writeWorkflowGroupState.mockRejectedValueOnce(new Error('Cell write failed'))

await expect(executeResumeJob(PAYLOAD)).rejects.toBe(bookkeepingFailure)

expect(startedGroups).toEqual([])
}, 20_000)

it('does not continue the cascade when the resume failed the execution', async () => {
const runFailure = new Error('Block failed')
failResume('execution_failed', runFailure)

await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure)

expect(startedGroups).toEqual([])
}, 20_000)

it('puts the cell back to paused when the pause stayed resumable', async () => {
const admissionRefusal = new Error('Execution can no longer be resumed')
failResume('pause_retained', admissionRefusal)
Expand Down
112 changes: 94 additions & 18 deletions apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { createTimeoutAbortController, getExecutionDeadlineAt } from '@/lib/core/execution-limits'
import { abortManualExecution } from '@/lib/execution/manual-cancellation'
import { terminalExecutionLogFields } from '@/lib/logs/execution/cancellation'
import { LoggingSession } from '@/lib/logs/execution/logging-session'

const {
mockExecuteWorkflowCore,
Expand Down Expand Up @@ -85,7 +86,7 @@ if (!humanInTheLoopLogger) {
}

interface PauseResumeManagerInternals {
markResumeFailed: (...args: unknown[]) => Promise<boolean>
markResumeFailed: (...args: unknown[]) => Promise<FailedResumeOutcome | undefined>
runResumeExecution: (...args: unknown[]) => Promise<unknown>
}

Expand Down Expand Up @@ -262,17 +263,36 @@ describe('what a failed resume did to its paused execution', () => {
)

it.each([
{ logStatus: 'running', pauseStatus: 'paused', executionFailed: true },
{ logStatus: 'cancelled', pauseStatus: 'paused', executionFailed: false },
{ logStatus: 'running', pauseStatus: 'cancelling', executionFailed: false },
{ logStatus: 'running', pauseStatus: 'paused', outcome: 'execution_failed' },
{ logStatus: 'cancelled', pauseStatus: 'paused', outcome: undefined },
{ logStatus: 'running', pauseStatus: 'cancelling', outcome: undefined },
{ logStatus: 'failed', pauseStatus: 'paused', outcome: 'execution_failed' },
{ logStatus: 'completed', pauseStatus: 'paused', outcome: 'execution_completed' },
])(
'reports the execution failed: $executionFailed for a $logStatus log and $pauseStatus pause',
async ({ logStatus, pauseStatus, executionFailed }) => {
'reports $outcome for a $logStatus log and $pauseStatus pause',
async ({ logStatus, pauseStatus, outcome }) => {
queueTableRows(workflowExecutionLogs, [{ status: logStatus }])
queueTableRows(pausedExecutions, [{ status: pauseStatus }])
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals

await expect(managerInternals.markResumeFailed(attemptArgs)).resolves.toBe(executionFailed)
await expect(managerInternals.markResumeFailed(attemptArgs)).resolves.toBe(outcome)
}
)

it.each([
{ logStatus: 'completed', updated: [resumeQueue, pausedExecutions] },
{ logStatus: 'failed', updated: [resumeQueue, pausedExecutions] },
{ logStatus: 'running', updated: [resumeQueue, pausedExecutions, workflowExecutionLogs] },
])(
'leaves a $logStatus execution log as it is when a resume fails late',
async ({ logStatus, updated }) => {
queueTableRows(workflowExecutionLogs, [{ status: logStatus }])
queueTableRows(pausedExecutions, [{ status: 'paused' }])
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals

await managerInternals.markResumeFailed(attemptArgs)

expect(dbChainMockFns.update.mock.calls.map(([table]) => table)).toEqual(updated)
}
)

Expand Down Expand Up @@ -313,17 +333,72 @@ describe('what a failed resume did to its paused execution', () => {
}
)

describe('when the resumed run pauses but its pause cannot be saved', () => {
const spies: { mockRestore: () => void }[] = []

function pauseRun(options: { snapshotSeed?: unknown; persistError?: Error }) {
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals
spies.push(
vi.spyOn(managerInternals, 'runResumeExecution').mockResolvedValueOnce({
success: true,
status: 'paused',
output: {},
logs: [],
pausePoints: [],
snapshotSeed: options.snapshotSeed,
metadata: { executionId: 'parent-execution-1', duration: 1, startTime: 'start' },
}),
vi.spyOn(managerInternals, 'markResumeFailed').mockResolvedValueOnce('execution_failed'),
vi.spyOn(LoggingSession, 'markExecutionAsFailed').mockResolvedValueOnce(),
vi.spyOn(PauseResumeManager, 'processQueuedResumes').mockResolvedValueOnce()
)
if (options.persistError) {
spies.push(
vi
.spyOn(PauseResumeManager, 'persistPauseResult')
.mockRejectedValueOnce(options.persistError)
)
}
}

afterEach(() => {
for (const spy of spies.splice(0)) spy.mockRestore()
})

it('fails the attempt when the pause state cannot be persisted', async () => {
pauseRun({ snapshotSeed: createSnapshotSeed(), persistError: new Error('lock timeout') })
const { outcomes, args } = argsReportingOutcomes()

await expect(PauseResumeManager.startResumeExecution(args)).rejects.toMatchObject({
message: 'Failed to persist pause state',
cause: new Error('lock timeout'),
})
expect(outcomes).toEqual(['execution_failed'])
})

it('fails the attempt when the paused run has no snapshot seed', async () => {
pauseRun({})
const { outcomes, args } = argsReportingOutcomes()

await expect(PauseResumeManager.startResumeExecution(args)).rejects.toThrow(
'Missing snapshot seed for paused execution'
)
expect(outcomes).toEqual(['execution_failed'])
})
})

describe('when the resumed run fails', () => {
const rawError = new Error('Block failed')
const spies: { mockRestore: () => void }[] = []

function failRun(options: { executionFailed: boolean; queuedResumesError?: Error }) {
function failRun(options: {
outcome: FailedResumeOutcome | undefined
queuedResumesError?: Error
}) {
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals
spies.push(
vi.spyOn(managerInternals, 'runResumeExecution').mockRejectedValueOnce(rawError),
vi
.spyOn(managerInternals, 'markResumeFailed')
.mockResolvedValueOnce(options.executionFailed),
vi.spyOn(managerInternals, 'markResumeFailed').mockResolvedValueOnce(options.outcome),
options.queuedResumesError
? vi
.spyOn(PauseResumeManager, 'processQueuedResumes')
Expand All @@ -337,12 +412,13 @@ describe('what a failed resume did to its paused execution', () => {
})

it.each([
{ executionFailed: true, reported: ['execution_failed'] },
{ executionFailed: false, reported: [] },
{ outcome: 'execution_failed' as const, reported: ['execution_failed'] },
{ outcome: 'execution_completed' as const, reported: ['execution_completed'] },
{ outcome: undefined, reported: [] },
])(
'reports $reported when it failed the execution: $executionFailed',
async ({ executionFailed, reported }) => {
failRun({ executionFailed })
'reports $reported when the attempt left the execution $outcome',
async ({ outcome, reported }) => {
failRun({ outcome })
const { outcomes, args } = argsReportingOutcomes()

await expect(PauseResumeManager.startResumeExecution(args)).rejects.toBe(rawError)
Expand All @@ -351,7 +427,7 @@ describe('what a failed resume did to its paused execution', () => {
)

it('reports the outcome before draining queued resumes, which may throw', async () => {
failRun({ executionFailed: true, queuedResumesError: new Error('queue drain failed') })
failRun({ outcome: 'execution_failed', queuedResumesError: new Error('queue drain failed') })
const { outcomes, args } = argsReportingOutcomes()

await expect(PauseResumeManager.startResumeExecution(args)).rejects.toThrow(
Expand All @@ -361,7 +437,7 @@ describe('what a failed resume did to its paused execution', () => {
})

it('rethrows the run failure when the outcome handler fails', async () => {
failRun({ executionFailed: true })
failRun({ outcome: 'execution_failed' })
const { args } = argsReportingOutcomes(async () => {
throw new Error('Database unavailable')
})
Expand Down
Loading
Loading