diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index 0d763e391ca..db86c516b1a 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -1,4 +1,5 @@ import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { task, timeout } from '@trigger.dev/sdk' import { @@ -21,6 +22,7 @@ import { classifyWorkflowCellTerminalResult } from '@/lib/table/workflow-cell-re import type { CellResumeContext } from '@/lib/table/workflow-columns' import { createResumeAttemptTimeoutController, + type FailedResumeOutcome, PauseResumeManager, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits' @@ -359,6 +361,25 @@ async function buildResumeCellWriters( return { cellOnBlockComplete, writeCellTerminal } } +/** + * 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. + */ +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')) + } +} + async function runResumeAndCellTerminal( payload: ResumeExecutionPayload, pausedExecution: Awaited>, @@ -375,6 +396,7 @@ async function runResumeAndCellTerminal( resumeInput: payload.resumeInput, userId: payload.userId, onBlockComplete: writers.cellOnBlockComplete, + onAttemptFailed: (outcome, error) => writeFailedResumeCellTerminal(writers, outcome, error), abortSignal: signal, }) diff --git a/apps/sim/background/resume-governed-subject.test.ts b/apps/sim/background/resume-governed-subject.test.ts index 4f6a62778e4..d7589d6eb45 100644 --- a/apps/sim/background/resume-governed-subject.test.ts +++ b/apps/sim/background/resume-governed-subject.test.ts @@ -51,6 +51,7 @@ vi.mock('@/executor/execution/snapshot', () => ({ ExecutionSnapshot: { fromJSON: hoisted.snapshotFromJson }, })) +import type { FailedResumeOutcome } from '@/lib/workflows/executor/human-in-the-loop-manager' import { executeResumeJob, type ResumeExecutionPayload } from '@/background/resume-execution' const mocks = { @@ -217,4 +218,52 @@ describe('resuming a paused table cell', () => { const [cascadePayload] = mocks.runRowCascadeLoop.mock.calls[0] expect(cascadePayload.capabilityGovernedUserId).toBe('requesting-member') }, 20_000) + + describe('when the resume throws', () => { + /** 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. */ + function failResume(outcome: FailedResumeOutcome, error: Error) { + mocks.startResumeExecution.mockImplementationOnce( + async ({ + onAttemptFailed, + }: { + onAttemptFailed?: (outcome: FailedResumeOutcome, error: unknown) => Promise + }) => { + await onAttemptFailed?.(outcome, error) + throw error + } + ) + } + + it('marks the cell failed when the resume failed the execution', async () => { + const runFailure = new Error('writeLedger: Unique constraint violation') + failResume('execution_failed', runFailure) + + await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) + + expect(lastCellExecutionState()).toMatchObject({ + status: 'error', + executionId: 'parent-execution-1', + error: 'writeLedger: Unique constraint violation', + }) + }, 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) + + await expect(executeResumeJob(PAYLOAD)).rejects.toBe(admissionRefusal) + + expect(lastCellExecutionState()).toMatchObject({ + status: 'pending', + executionId: 'parent-execution-1', + jobId: 'paused-parent-execution-1', + }) + }, 20_000) + }) }) 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 d407d5c9525..95205271c98 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 @@ -18,7 +18,7 @@ import { largeValueMetadataMock, largeValueMetadataMockFns, } from '@sim/testing/mocks/large-value-metadata.mock' -import { beforeEach, describe, expect, it, vi } from 'vitest' +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' @@ -61,6 +61,7 @@ vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadat import { createResumeAttemptTimeoutController, + type FailedResumeOutcome, PauseResumeManager, requireResumeDeploymentVersion, updateResumeOutputInAggregationBuffers, @@ -84,7 +85,7 @@ if (!humanInTheLoopLogger) { } interface PauseResumeManagerInternals { - markResumeFailed: (...args: unknown[]) => Promise + markResumeFailed: (...args: unknown[]) => Promise runResumeExecution: (...args: unknown[]) => Promise } @@ -210,6 +211,166 @@ describe('queued resume attempt deadlines', () => { }) }) +describe('what a failed resume did to its paused execution', () => { + type StartResumeArgs = Parameters[0] + const pausedExecution = { + id: 'paused-execution-1', + workflowId: 'workflow-1', + executionId: 'parent-execution-1', + pausePoints: { + 'context-1': { contextId: 'context-1', blockId: 'hitl-1' }, + }, + executionSnapshot: createSnapshotSeed(), + metadata: {}, + } as StartResumeArgs['pausedExecution'] + const resumeArgs: StartResumeArgs = { + resumeEntryId: 'resume-entry-1', + resumeExecutionId: 'resume-execution-1', + pausedExecution, + contextId: 'context-1', + resumeInput: { approved: true }, + userId: 'user-1', + } + const attemptArgs = { + resumeEntryId: 'resume-entry-1', + pausedExecutionId: 'paused-execution-1', + parentExecutionId: 'parent-execution-1', + contextId: 'context-1', + failureReason: 'Execution can no longer be resumed', + preserveForRetry: true, + } + + beforeEach(() => { + resetDbChainMock() + }) + + it('reports the pause still resumable when a refused attempt leaves it paused', async () => { + queueTableRows(workflowExecutionLogs, [{ status: 'paused' }]) + queueTableRows(pausedExecutions, [{ automaticResumeRetryCount: 0, status: 'paused' }]) + + await expect(PauseResumeManager.markResumeAttemptFailed(attemptArgs)).resolves.toBe(true) + }) + + it.each(['completed', 'failed', 'cancelled'])( + 'reports the pause not resumable when the execution is already %s', + async (logStatus) => { + queueTableRows(workflowExecutionLogs, [{ status: logStatus }]) + queueTableRows(pausedExecutions, [{ automaticResumeRetryCount: 0, status: 'paused' }]) + + await expect(PauseResumeManager.markResumeAttemptFailed(attemptArgs)).resolves.toBe(false) + } + ) + + it.each([ + { logStatus: 'running', pauseStatus: 'paused', executionFailed: true }, + { logStatus: 'cancelled', pauseStatus: 'paused', executionFailed: false }, + { logStatus: 'running', pauseStatus: 'cancelling', executionFailed: false }, + ])( + 'reports the execution failed: $executionFailed for a $logStatus log and $pauseStatus pause', + async ({ logStatus, pauseStatus, executionFailed }) => { + queueTableRows(workflowExecutionLogs, [{ status: logStatus }]) + queueTableRows(pausedExecutions, [{ status: pauseStatus }]) + const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals + + await expect(managerInternals.markResumeFailed(attemptArgs)).resolves.toBe(executionFailed) + } + ) + + /** Resume args that collect every outcome the manager reports. */ + function argsReportingOutcomes(onAttemptFailed?: () => Promise) { + const outcomes: FailedResumeOutcome[] = [] + return { + outcomes, + args: { + ...resumeArgs, + onAttemptFailed: async (outcome: FailedResumeOutcome) => { + outcomes.push(outcome) + await onAttemptFailed?.() + }, + }, + } + } + + it.each([ + { stillResumable: true, reported: ['pause_retained'] }, + { stillResumable: false, reported: [] }, + ])( + 'reports a refused attempt as $reported when the pause is resumable: $stillResumable', + async ({ stillResumable, reported }) => { + const markResumeAttemptFailedSpy = vi + .spyOn(PauseResumeManager, 'markResumeAttemptFailed') + .mockResolvedValueOnce(stillResumable) + const { outcomes, args } = argsReportingOutcomes() + + try { + await expect(PauseResumeManager.startResumeExecution(args)).rejects.toMatchObject({ + name: 'ResumeAdmissionError', + }) + expect(outcomes).toEqual(reported) + } finally { + markResumeAttemptFailedSpy.mockRestore() + } + } + ) + + describe('when the resumed run fails', () => { + const rawError = new Error('Block failed') + const spies: { mockRestore: () => void }[] = [] + + function failRun(options: { executionFailed: boolean; queuedResumesError?: Error }) { + const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals + spies.push( + vi.spyOn(managerInternals, 'runResumeExecution').mockRejectedValueOnce(rawError), + vi + .spyOn(managerInternals, 'markResumeFailed') + .mockResolvedValueOnce(options.executionFailed), + options.queuedResumesError + ? vi + .spyOn(PauseResumeManager, 'processQueuedResumes') + .mockRejectedValueOnce(options.queuedResumesError) + : vi.spyOn(PauseResumeManager, 'processQueuedResumes').mockResolvedValueOnce() + ) + } + + afterEach(() => { + for (const spy of spies.splice(0)) spy.mockRestore() + }) + + it.each([ + { executionFailed: true, reported: ['execution_failed'] }, + { executionFailed: false, reported: [] }, + ])( + 'reports $reported when it failed the execution: $executionFailed', + async ({ executionFailed, reported }) => { + failRun({ executionFailed }) + const { outcomes, args } = argsReportingOutcomes() + + await expect(PauseResumeManager.startResumeExecution(args)).rejects.toBe(rawError) + expect(outcomes).toEqual(reported) + } + ) + + it('reports the outcome before draining queued resumes, which may throw', async () => { + failRun({ executionFailed: true, queuedResumesError: new Error('queue drain failed') }) + const { outcomes, args } = argsReportingOutcomes() + + await expect(PauseResumeManager.startResumeExecution(args)).rejects.toThrow( + 'queue drain failed' + ) + expect(outcomes).toEqual(['execution_failed']) + }) + + it('rethrows the run failure when the outcome handler fails', async () => { + failRun({ executionFailed: true }) + const { args } = argsReportingOutcomes(async () => { + throw new Error('Database unavailable') + }) + + await expect(PauseResumeManager.startResumeExecution(args)).rejects.toBe(rawError) + }) + }) +}) + describe('resume failure diagnostic projection', () => { beforeEach(() => { resetDbChainMock() 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 8125ead3911..732240b12f9 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 @@ -38,6 +38,7 @@ import { terminalExecutionLogFields, } from '@/lib/logs/execution/cancellation' import { LoggingSession } from '@/lib/logs/execution/logging-session' +import type { PersistedWorkflowExecutionStatus } from '@/lib/logs/types' import { cleanupExecutionBase64Cache } from '@/lib/uploads/utils/user-file-base64.server' import { executeWorkflowCore } from '@/lib/workflows/executor/execution-core' import { @@ -88,6 +89,19 @@ const execDb = dbFor('exec') const logger = createLogger('HumanInTheLoopManager') const RUN_BUFFER_UNAVAILABLE_ERROR = 'Run buffer temporarily unavailable' const RESUMABLE_PAUSED_STATUSES = ['paused', 'partially_resumed'] as const +/** + * Statuses of a finished execution's log, the same set the attempt-failure + * `CASE` preserves; a resume can never claim a log in one of them. + */ +const TERMINAL_EXECUTION_LOG_STATUSES = [ + 'cancelled', + 'failed', + 'completed', +] as const satisfies readonly PersistedWorkflowExecutionStatus[] + +function isTerminalExecutionLogStatus(status: string): boolean { + return (TERMINAL_EXECUTION_LOG_STATUSES as readonly string[]).includes(status) +} const CANCELLABLE_PAUSED_STATUSES = ['paused', 'partially_resumed'] as const const AUTOMATIC_RESUME_INTERVENTION_PREFIX = 'Automatic resume requires manual intervention: ' const PAUSED_CANCELLATION_QUEUE_FAILURE_REASON = 'Paused execution cancellation requested' @@ -134,6 +148,13 @@ 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. + */ +export type FailedResumeOutcome = 'pause_retained' | 'execution_failed' + /** Matches the paused execution mode to the deployment recorded on its durable root log. */ export function requireResumeDeploymentVersion( useDraftState: unknown, @@ -419,6 +440,13 @@ interface StartResumeExecutionArgs { sendEvent?: (event: ExecutionEvent) => void onStream?: (streamingExec: StreamingExecution) => Promise onBlockComplete?: (blockId: string, data: BlockCompletionCallbackData) => Promise + /** + * Called once a failed attempt is settled, so a caller mirroring the run's + * state (a table cell) follows the execution. Not called when the attempt + * changed neither (the execution had already finished or was cancelled). A + * throw is logged and never replaces the attempt's error. + */ + onAttemptFailed?: (outcome: FailedResumeOutcome, error: unknown) => Promise abortSignal?: AbortSignal } @@ -864,6 +892,7 @@ export class PauseResumeManager { sendEvent, onStream, onBlockComplete, + onAttemptFailed, abortSignal, } = args @@ -986,8 +1015,9 @@ export class PauseResumeManager { } catch (error) { const message = toError(error).message await releaseExecutionSlot(resumeEntryId) + let outcome: FailedResumeOutcome | undefined if (error instanceof ResumeAdmissionError) { - await PauseResumeManager.markResumeAttemptFailed({ + const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, @@ -996,22 +1026,33 @@ export class PauseResumeManager { preserveForRetry: true, retryable: error.retryable, }) + if (pauseResumable) outcome = 'pause_retained' } else if (message === RUN_BUFFER_UNAVAILABLE_ERROR) { - await PauseResumeManager.markResumeAttemptFailed({ + const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, contextId, failureReason: message, }) + if (pauseResumable) outcome = 'pause_retained' } else { - await PauseResumeManager.markResumeFailed({ + const executionFailed = 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) => { + logger.error( + 'Failed to report a failed resume attempt', + projectResolvedSecretDiagnosticError(hookError, undefined, { resumeExecutionId }) + ) + }) } logger.error( 'Resume execution failed', @@ -2192,10 +2233,10 @@ export class PauseResumeManager { parentExecutionId: string contextId: string failureReason: string - }): Promise { + }): Promise { const now = new Date() - await execDb.transaction(async (tx) => { + return execDb.transaction(async (tx) => { const executionLog = await tx .select({ status: workflowExecutionLogs.status }) .from(workflowExecutionLogs) @@ -2224,7 +2265,7 @@ export class PauseResumeManager { .set({ status: 'cancelled', updatedAt: now, nextResumeAt: null }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) } - return + return false } await tx @@ -2235,7 +2276,7 @@ export class PauseResumeManager { }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) - if (pausedExecution?.status === 'cancelling') return + if (pausedExecution?.status === 'cancelling') return false await tx .update(workflowExecutionLogs) @@ -2246,6 +2287,8 @@ export class PauseResumeManager { sql`${workflowExecutionLogs.status} != 'cancelled'` ) ) + + return true }) } @@ -2257,10 +2300,10 @@ export class PauseResumeManager { failureReason: string preserveForRetry?: boolean retryable?: boolean - }): Promise { + }): Promise { const now = new Date() - await execDb.transaction(async (tx) => { + return execDb.transaction(async (tx) => { const executionLog = await tx .select({ status: workflowExecutionLogs.status }) .from(workflowExecutionLogs) @@ -2317,7 +2360,7 @@ export class PauseResumeManager { .set({ status: 'cancelled', updatedAt: now, nextResumeAt: null }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) } - return + return false } await tx @@ -2348,7 +2391,7 @@ export class PauseResumeManager { }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) - if (pausedExecution?.status === 'cancelling') return + if (pausedExecution?.status === 'cancelling') return false await tx .update(workflowExecutionLogs) @@ -2369,6 +2412,13 @@ export class PauseResumeManager { )` ) ) + + return ( + pausedExecution !== undefined && + isResumablePausedStatus(pausedExecution.status) && + executionLog !== undefined && + !isTerminalExecutionLogStatus(executionLog.status) + ) }) }