diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index db86c516b1a..f6a8fa9b7cd 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -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 @@ -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 { - 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 + case 'execution_failed': + await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed')) + return } } @@ -385,7 +415,8 @@ async function runResumeAndCellTerminal( pausedExecution: Awaited>, writers: CellWriters, signal: AbortSignal | undefined, - timeoutController: ReturnType + timeoutController: ReturnType, + onCompletedBeforeFailure?: () => void ): Promise>> { if (!pausedExecution) throw new Error('Paused execution missing — already nulled by caller') const result = await PauseResumeManager.startResumeExecution({ @@ -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, }) diff --git a/apps/sim/background/resume-governed-subject.test.ts b/apps/sim/background/resume-governed-subject.test.ts index d7589d6eb45..a689b1e78b0 100644 --- a/apps/sim/background/resume-governed-subject.test.ts +++ b/apps/sim/background/resume-governed-subject.test.ts @@ -220,13 +220,26 @@ 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 ({ @@ -234,7 +247,7 @@ describe('resuming a paused table cell', () => { }: { onAttemptFailed?: (outcome: FailedResumeOutcome, error: unknown) => Promise }) => { - await onAttemptFailed?.(outcome, error) + await onAttemptFailed?.(outcome, error).catch(() => undefined) throw error } ) @@ -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) diff --git a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts index 95205271c98..b5b182b2ef4 100644 --- a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts +++ b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts @@ -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, @@ -85,7 +86,7 @@ if (!humanInTheLoopLogger) { } interface PauseResumeManagerInternals { - markResumeFailed: (...args: unknown[]) => Promise + markResumeFailed: (...args: unknown[]) => Promise runResumeExecution: (...args: unknown[]) => Promise } @@ -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) } ) @@ -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') @@ -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) @@ -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( @@ -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') }) diff --git a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts index 732240b12f9..82f64214ac3 100644 --- a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts +++ b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts @@ -149,11 +149,12 @@ class ResumeAdmissionError extends Error { } /** - * What a failed resume attempt did to its paused execution, as reported by the - * transaction that settled the attempt: the pause stayed resumable, or the - * resumed run failed the execution. + * What a failed resume attempt left its paused execution as, read from the + * transaction that settled the attempt: the pause stayed resumable, the + * resumed run failed the execution, or the run completed before a later step + * of the attempt threw. */ -export type FailedResumeOutcome = 'pause_retained' | 'execution_failed' +export type FailedResumeOutcome = 'pause_retained' | 'execution_failed' | 'execution_completed' /** Matches the paused execution mode to the deployment recorded on its durable root log. */ export function requireResumeDeploymentVersion( @@ -924,18 +925,29 @@ export class PauseResumeManager { }) if (result.status === 'paused') { + /** + * A pause that cannot be saved fails the execution. Fail the log with the + * reason, then throw so the attempt settles as failed below. The thrown + * message stays stable; the underlying error rides on `cause`. + */ const effectiveExecutionId = result.metadata?.executionId ?? resumeExecutionId - if (!result.snapshotSeed) { - logger.error('Missing snapshot seed for paused resume execution', { - resumeExecutionId, - }) + const failPause = async (reason: string, cause?: unknown): Promise => { + logger.error( + reason, + cause === undefined + ? { resumeExecutionId } + : projectResolvedSecretDiagnosticError(cause, undefined, { resumeExecutionId }) + ) await LoggingSession.markExecutionAsFailed( effectiveExecutionId, - 'Missing snapshot seed for paused execution', + cause === undefined ? reason : `${reason}: ${toError(cause).message}`, undefined, pausedExecution.workflowId ) - await releaseExecutionSlot(resumeEntryId) + throw new Error(reason, { cause }) + } + if (!result.snapshotSeed) { + await failPause('Missing snapshot seed for paused execution') } else { try { await PauseResumeManager.persistPauseResult({ @@ -947,19 +959,7 @@ export class PauseResumeManager { executorUserId: result.metadata?.userId, }) } catch (pauseError) { - logger.error( - 'Failed to persist pause result for resumed execution', - projectResolvedSecretDiagnosticError(pauseError, undefined, { - resumeExecutionId, - }) - ) - await LoggingSession.markExecutionAsFailed( - effectiveExecutionId, - `Failed to persist pause state: ${toError(pauseError).message}`, - undefined, - pausedExecution.workflowId - ) - await releaseExecutionSlot(resumeEntryId) + await failPause('Failed to persist pause state', pauseError) } } } else { @@ -1037,14 +1037,13 @@ export class PauseResumeManager { }) if (pauseResumable) outcome = 'pause_retained' } else { - const executionFailed = await PauseResumeManager.markResumeFailed({ + outcome = await PauseResumeManager.markResumeFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, contextId, failureReason: message, }) - if (executionFailed) outcome = 'execution_failed' } if (outcome && onAttemptFailed) { await onAttemptFailed(outcome, error).catch((hookError: unknown) => { @@ -2233,7 +2232,7 @@ export class PauseResumeManager { parentExecutionId: string contextId: string failureReason: string - }): Promise { + }): Promise<'execution_failed' | 'execution_completed' | undefined> { const now = new Date() return execDb.transaction(async (tx) => { @@ -2265,7 +2264,7 @@ export class PauseResumeManager { .set({ status: 'cancelled', updatedAt: now, nextResumeAt: null }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) } - return false + return undefined } await tx @@ -2276,19 +2275,24 @@ export class PauseResumeManager { }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) - if (pausedExecution?.status === 'cancelling') return false + if (pausedExecution?.status === 'cancelling') return undefined - await tx - .update(workflowExecutionLogs) - .set(terminalExecutionLogFields('failed', now)) - .where( - and( - eq(workflowExecutionLogs.executionId, args.parentExecutionId), - sql`${workflowExecutionLogs.status} != 'cancelled'` + /** The run completed before a later step of the attempt threw; its outcome stands. */ + if (executionLog?.status === 'completed') return 'execution_completed' + + if (executionLog?.status !== 'failed') { + await tx + .update(workflowExecutionLogs) + .set(terminalExecutionLogFields('failed', now)) + .where( + and( + eq(workflowExecutionLogs.executionId, args.parentExecutionId), + sql`${workflowExecutionLogs.status} != 'cancelled'` + ) ) - ) + } - return true + return 'execution_failed' }) }