From 9eed4a360e7d476eaa4f5ad86eb919b6e536f00d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 19:58:31 -0700 Subject: [PATCH 1/4] fix(tables): write the cell state when a resumed run throws runResumeAndCellTerminal only wrote the cell terminal after a resume returned, so a resume that threw left the cell on its last partial running state, where it could not even be cancelled. The resume manager now records, on the error it rethrows, when the attempt kept its pause resumable (admission refused, run buffer unavailable). The resume job mirrors that onto the cell: back to paused when the pause was kept, failed otherwise. A failed cell write is logged without masking the resume error, which is still rethrown. --- apps/sim/background/resume-execution.ts | 49 +++++++++++--- .../resume-governed-subject.test.ts | 44 ++++++++++++ .../human-in-the-loop-manager.test.ts | 67 +++++++++++++++++++ .../executor/human-in-the-loop-manager.ts | 22 ++++++ .../mocks/human-in-the-loop-manager.mock.ts | 2 + 5 files changed, 174 insertions(+), 10 deletions(-) diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index 0d763e391ca..f1ea65fe977 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 { @@ -22,6 +23,7 @@ import type { CellResumeContext } from '@/lib/table/workflow-columns' import { createResumeAttemptTimeoutController, PauseResumeManager, + wasPausedExecutionRetained, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits' import { ExecutionSnapshot } from '@/executor/execution/snapshot' @@ -359,6 +361,27 @@ 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 the execution instead: a pause the resume kept + * resumable goes back to paused, and any other failure ended the run. + */ +async function writeFailedResumeCellTerminal(writers: CellWriters, error: unknown): Promise { + try { + if (wasPausedExecutionRetained(error)) { + await writers.writeCellTerminal('paused', null) + } else { + await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed')) + } + } catch (writeError) { + logger.error( + 'Failed to write the cell state after a failed resume', + projectResolvedSecretDiagnosticError(writeError, undefined) + ) + } +} + async function runResumeAndCellTerminal( payload: ResumeExecutionPayload, pausedExecution: Awaited>, @@ -367,16 +390,22 @@ async function runResumeAndCellTerminal( timeoutController: ReturnType ): Promise>> { if (!pausedExecution) throw new Error('Paused execution missing — already nulled by caller') - const result = await PauseResumeManager.startResumeExecution({ - resumeEntryId: payload.resumeEntryId, - resumeExecutionId: payload.resumeExecutionId, - pausedExecution, - contextId: payload.contextId, - resumeInput: payload.resumeInput, - userId: payload.userId, - onBlockComplete: writers.cellOnBlockComplete, - abortSignal: signal, - }) + let result: Awaited> + try { + result = await PauseResumeManager.startResumeExecution({ + resumeEntryId: payload.resumeEntryId, + resumeExecutionId: payload.resumeExecutionId, + pausedExecution, + contextId: payload.contextId, + resumeInput: payload.resumeInput, + userId: payload.userId, + onBlockComplete: writers.cellOnBlockComplete, + abortSignal: signal, + }) + } catch (error) { + await writeFailedResumeCellTerminal(writers, error) + throw error + } if (result.status === 'paused') { await writers.writeCellTerminal('paused', null) diff --git a/apps/sim/background/resume-governed-subject.test.ts b/apps/sim/background/resume-governed-subject.test.ts index 4f6a62778e4..e430c624d70 100644 --- a/apps/sim/background/resume-governed-subject.test.ts +++ b/apps/sim/background/resume-governed-subject.test.ts @@ -57,6 +57,7 @@ const mocks = { ...hoisted, getPausedExecutionById: humanInTheLoopManagerMockFns.mockGetPausedExecutionById, startResumeExecution: humanInTheLoopManagerMockFns.mockStartResumeExecution, + wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, findCellContextByExecutionId: tableWorkflowColumnsMockFns.mockFindCellContextByExecutionId, @@ -217,4 +218,47 @@ 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 + } + + it('marks the cell failed when the resumed run itself failed', async () => { + const runFailure = new Error('writeLedger: Unique constraint violation') + mocks.startResumeExecution.mockRejectedValueOnce(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') + mocks.startResumeExecution.mockRejectedValueOnce(admissionRefusal) + mocks.wasPausedExecutionRetained.mockReturnValueOnce(true) + + await expect(executeResumeJob(PAYLOAD)).rejects.toBe(admissionRefusal) + + expect(lastCellExecutionState()).toMatchObject({ + status: 'pending', + executionId: 'parent-execution-1', + jobId: 'paused-parent-execution-1', + }) + }, 20_000) + + it('still reports the resume failure when the cell write also fails', async () => { + const runFailure = new Error('Block failed') + mocks.startResumeExecution.mockRejectedValueOnce(runFailure) + mocks.writeWorkflowGroupState.mockRejectedValueOnce(new Error('Database unavailable')) + + await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) + }, 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..8d55aabb5f3 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 @@ -64,6 +64,7 @@ import { PauseResumeManager, requireResumeDeploymentVersion, updateResumeOutputInAggregationBuffers, + wasPausedExecutionRetained, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { getAutomaticResumeWaitingMetadata } from '@/lib/workflows/executor/paused-execution-metadata' import { AUTOMATIC_RESUME_WAITING_REASON_MAX_LENGTH } from '@/lib/workflows/executor/resume-policy' @@ -210,6 +211,72 @@ describe('queued resume attempt deadlines', () => { }) }) +describe('which failed resumes keep the pause resumable', () => { + 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', + } + + beforeEach(() => { + resetDbChainMock() + }) + + it('keeps the pause resumable when the paused log can no longer be claimed', async () => { + const markResumeAttemptFailedSpy = vi + .spyOn(PauseResumeManager, 'markResumeAttemptFailed') + .mockResolvedValueOnce() + + try { + const thrown = await PauseResumeManager.startResumeExecution(resumeArgs).catch( + (error: unknown) => error + ) + + expect(thrown).toMatchObject({ name: 'ResumeAdmissionError' }) + expect(wasPausedExecutionRetained(thrown)).toBe(true) + } finally { + markResumeAttemptFailedSpy.mockRestore() + } + }) + + it('does not keep the pause when the resumed run itself failed', async () => { + const rawError = new Error('Block failed') + const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals + const runResumeExecutionSpy = vi + .spyOn(managerInternals, 'runResumeExecution') + .mockRejectedValueOnce(rawError) + const markResumeFailedSpy = vi + .spyOn(managerInternals, 'markResumeFailed') + .mockResolvedValueOnce() + const processQueuedResumesSpy = vi + .spyOn(PauseResumeManager, 'processQueuedResumes') + .mockResolvedValueOnce() + + try { + await expect(PauseResumeManager.startResumeExecution(resumeArgs)).rejects.toBe(rawError) + expect(wasPausedExecutionRetained(rawError)).toBe(false) + } finally { + runResumeExecutionSpy.mockRestore() + markResumeFailedSpy.mockRestore() + processQueuedResumesSpy.mockRestore() + } + }) +}) + 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..c32ad66cbab 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 @@ -134,6 +134,26 @@ class ResumeAdmissionError extends Error { } } +/** + * Errors from resume attempts that failed without consuming their pause: the + * attempt was refused or could not start, so the paused execution stays + * resumable. Every other failed resume leaves the execution terminal. + */ +const pauseRetainingFailures = new WeakSet() + +/** + * Whether `error` was thrown by a resume attempt that left its paused execution + * resumable, so a caller mirroring the run's state (a table cell) keeps it paused + * rather than failed. + */ +export function wasPausedExecutionRetained(error: unknown): boolean { + return typeof error === 'object' && error !== null && pauseRetainingFailures.has(error) +} + +function retainPausedExecution(error: unknown): void { + if (typeof error === 'object' && error !== null) pauseRetainingFailures.add(error) +} + /** Matches the paused execution mode to the deployment recorded on its durable root log. */ export function requireResumeDeploymentVersion( useDraftState: unknown, @@ -987,6 +1007,7 @@ export class PauseResumeManager { const message = toError(error).message await releaseExecutionSlot(resumeEntryId) if (error instanceof ResumeAdmissionError) { + retainPausedExecution(error) await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, @@ -997,6 +1018,7 @@ export class PauseResumeManager { retryable: error.retryable, }) } else if (message === RUN_BUFFER_UNAVAILABLE_ERROR) { + retainPausedExecution(error) await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, diff --git a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts index c9011d2bc47..1e19a89ef3e 100644 --- a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts +++ b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts @@ -162,6 +162,7 @@ export const humanInTheLoopManagerMockFns = { } ), mockCreateResumeAttemptTimeoutController: vi.fn(), + mockWasPausedExecutionRetained: vi.fn((_error: unknown): boolean => false), mockExtractResumeBillingAttributionFromSnapshot: vi.fn(), mockComputeEarliestResumeAt: vi.fn( (points: Iterable, options: { after?: Date } = {}): Date | null => { @@ -195,6 +196,7 @@ export const humanInTheLoopManagerMock = { humanInTheLoopManagerMockFns.mockUpdateResumeOutputInAggregationBuffers, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, + wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained, extractResumeBillingAttributionFromSnapshot: humanInTheLoopManagerMockFns.mockExtractResumeBillingAttributionFromSnapshot, computeEarliestResumeAt: humanInTheLoopManagerMockFns.mockComputeEarliestResumeAt, From e19a8ab8387eddf5d7a0739e43e3598239d875d3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 20:11:53 -0700 Subject: [PATCH 2/4] fix(tables): mirror what a failed resume did instead of inferring it from the error An admission refusal can also mean the execution already finished, so treating every refusal as a kept pause wrote paused over a completed cell. markResumeAttemptFailed and markResumeFailed now return what their transaction actually did (pause still resumable / execution failed), the manager records that outcome on the rethrown error, and the resume job writes paused, error, or nothing. --- apps/sim/background/resume-execution.ts | 11 +- .../resume-governed-subject.test.ts | 18 ++- .../human-in-the-loop-manager.test.ts | 109 ++++++++++++------ .../executor/human-in-the-loop-manager.ts | 64 ++++++---- .../mocks/human-in-the-loop-manager.mock.ts | 6 +- 5 files changed, 139 insertions(+), 69 deletions(-) diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index f1ea65fe977..2e9b2a58d8a 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -22,8 +22,8 @@ import { classifyWorkflowCellTerminalResult } from '@/lib/table/workflow-cell-re import type { CellResumeContext } from '@/lib/table/workflow-columns' import { createResumeAttemptTimeoutController, + getFailedResumeOutcome, PauseResumeManager, - wasPausedExecutionRetained, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits' import { ExecutionSnapshot } from '@/executor/execution/snapshot' @@ -364,12 +364,15 @@ 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 the execution instead: a pause the resume kept - * resumable goes back to paused, and any other failure ended the run. + * partial `running` state. Mirror what the failed attempt did to the execution: + * a pause that stayed resumable goes back to paused, a failed execution fails + * the cell, and an attempt that changed nothing leaves the cell alone. */ async function writeFailedResumeCellTerminal(writers: CellWriters, error: unknown): Promise { + const outcome = getFailedResumeOutcome(error) + if (!outcome) return try { - if (wasPausedExecutionRetained(error)) { + if (outcome === 'pause_retained') { await writers.writeCellTerminal('paused', null) } else { await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed')) diff --git a/apps/sim/background/resume-governed-subject.test.ts b/apps/sim/background/resume-governed-subject.test.ts index e430c624d70..75f895ae9b4 100644 --- a/apps/sim/background/resume-governed-subject.test.ts +++ b/apps/sim/background/resume-governed-subject.test.ts @@ -57,7 +57,7 @@ const mocks = { ...hoisted, getPausedExecutionById: humanInTheLoopManagerMockFns.mockGetPausedExecutionById, startResumeExecution: humanInTheLoopManagerMockFns.mockStartResumeExecution, - wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained, + getFailedResumeOutcome: humanInTheLoopManagerMockFns.mockGetFailedResumeOutcome, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, findCellContextByExecutionId: tableWorkflowColumnsMockFns.mockFindCellContextByExecutionId, @@ -226,9 +226,10 @@ describe('resuming a paused table cell', () => { return payload?.executionState } - it('marks the cell failed when the resumed run itself failed', async () => { + it('marks the cell failed when the resume failed the execution', async () => { const runFailure = new Error('writeLedger: Unique constraint violation') mocks.startResumeExecution.mockRejectedValueOnce(runFailure) + mocks.getFailedResumeOutcome.mockReturnValueOnce('execution_failed') await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) @@ -242,7 +243,7 @@ describe('resuming a paused table cell', () => { it('puts the cell back to paused when the pause stayed resumable', async () => { const admissionRefusal = new Error('Execution can no longer be resumed') mocks.startResumeExecution.mockRejectedValueOnce(admissionRefusal) - mocks.wasPausedExecutionRetained.mockReturnValueOnce(true) + mocks.getFailedResumeOutcome.mockReturnValueOnce('pause_retained') await expect(executeResumeJob(PAYLOAD)).rejects.toBe(admissionRefusal) @@ -253,9 +254,20 @@ describe('resuming a paused table cell', () => { }) }, 20_000) + it('leaves the cell alone when the failed attempt changed nothing', async () => { + const staleRefusal = new Error('Execution can no longer be resumed') + mocks.startResumeExecution.mockRejectedValueOnce(staleRefusal) + mocks.getFailedResumeOutcome.mockReturnValueOnce(undefined) + + await expect(executeResumeJob(PAYLOAD)).rejects.toBe(staleRefusal) + + expect(lastCellExecutionState()).toBeUndefined() + }, 20_000) + it('still reports the resume failure when the cell write also fails', async () => { const runFailure = new Error('Block failed') mocks.startResumeExecution.mockRejectedValueOnce(runFailure) + mocks.getFailedResumeOutcome.mockReturnValueOnce('execution_failed') mocks.writeWorkflowGroupState.mockRejectedValueOnce(new Error('Database unavailable')) await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) 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 8d55aabb5f3..4f9159ae90f 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 @@ -61,10 +61,10 @@ vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadat import { createResumeAttemptTimeoutController, + getFailedResumeOutcome, PauseResumeManager, requireResumeDeploymentVersion, updateResumeOutputInAggregationBuffers, - wasPausedExecutionRetained, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { getAutomaticResumeWaitingMetadata } from '@/lib/workflows/executor/paused-execution-metadata' import { AUTOMATIC_RESUME_WAITING_REASON_MAX_LENGTH } from '@/lib/workflows/executor/resume-policy' @@ -85,7 +85,7 @@ if (!humanInTheLoopLogger) { } interface PauseResumeManagerInternals { - markResumeFailed: (...args: unknown[]) => Promise + markResumeFailed: (...args: unknown[]) => Promise runResumeExecution: (...args: unknown[]) => Promise } @@ -211,7 +211,7 @@ describe('queued resume attempt deadlines', () => { }) }) -describe('which failed resumes keep the pause resumable', () => { +describe('what a failed resume did to its paused execution', () => { type StartResumeArgs = Parameters[0] const pausedExecution = { id: 'paused-execution-1', @@ -231,50 +231,87 @@ describe('which failed resumes keep the pause resumable', () => { 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('keeps the pause resumable when the paused log can no longer be claimed', async () => { - const markResumeAttemptFailedSpy = vi - .spyOn(PauseResumeManager, 'markResumeAttemptFailed') - .mockResolvedValueOnce() + it('reports the pause still resumable when a refused attempt leaves it paused', async () => { + queueTableRows(workflowExecutionLogs, [{ status: 'paused' }]) + queueTableRows(pausedExecutions, [{ automaticResumeRetryCount: 0, status: 'paused' }]) - try { - const thrown = await PauseResumeManager.startResumeExecution(resumeArgs).catch( - (error: unknown) => error - ) + await expect(PauseResumeManager.markResumeAttemptFailed(attemptArgs)).resolves.toBe(true) + }) - expect(thrown).toMatchObject({ name: 'ResumeAdmissionError' }) - expect(wasPausedExecutionRetained(thrown)).toBe(true) - } finally { - markResumeAttemptFailedSpy.mockRestore() + 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('does not keep the pause when the resumed run itself failed', async () => { - const rawError = new Error('Block failed') - const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals - const runResumeExecutionSpy = vi - .spyOn(managerInternals, 'runResumeExecution') - .mockRejectedValueOnce(rawError) - const markResumeFailedSpy = vi - .spyOn(managerInternals, 'markResumeFailed') - .mockResolvedValueOnce() - const processQueuedResumesSpy = vi - .spyOn(PauseResumeManager, 'processQueuedResumes') - .mockResolvedValueOnce() + it.each([ + { stillResumable: true, outcome: 'pause_retained' }, + { stillResumable: false, outcome: undefined }, + ])( + 'records a refused attempt as $outcome when the pause is resumable: $stillResumable', + async ({ stillResumable, outcome }) => { + const markResumeAttemptFailedSpy = vi + .spyOn(PauseResumeManager, 'markResumeAttemptFailed') + .mockResolvedValueOnce(stillResumable) + + try { + const thrown = await PauseResumeManager.startResumeExecution(resumeArgs).catch( + (error: unknown) => error + ) - try { - await expect(PauseResumeManager.startResumeExecution(resumeArgs)).rejects.toBe(rawError) - expect(wasPausedExecutionRetained(rawError)).toBe(false) - } finally { - runResumeExecutionSpy.mockRestore() - markResumeFailedSpy.mockRestore() - processQueuedResumesSpy.mockRestore() + expect(thrown).toMatchObject({ name: 'ResumeAdmissionError' }) + expect(getFailedResumeOutcome(thrown)).toBe(outcome) + } finally { + markResumeAttemptFailedSpy.mockRestore() + } } - }) + ) + + it.each([ + { executionFailed: true, outcome: 'execution_failed' }, + { executionFailed: false, outcome: undefined }, + ])( + 'records a failed run as $outcome when it failed the execution: $executionFailed', + async ({ executionFailed, outcome }) => { + const rawError = new Error('Block failed') + const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals + const runResumeExecutionSpy = vi + .spyOn(managerInternals, 'runResumeExecution') + .mockRejectedValueOnce(rawError) + const markResumeFailedSpy = vi + .spyOn(managerInternals, 'markResumeFailed') + .mockResolvedValueOnce(executionFailed) + const processQueuedResumesSpy = vi + .spyOn(PauseResumeManager, 'processQueuedResumes') + .mockResolvedValueOnce() + + try { + await expect(PauseResumeManager.startResumeExecution(resumeArgs)).rejects.toBe(rawError) + expect(getFailedResumeOutcome(rawError)).toBe(outcome) + } finally { + runResumeExecutionSpy.mockRestore() + markResumeFailedSpy.mockRestore() + processQueuedResumesSpy.mockRestore() + } + } + ) }) describe('resume failure diagnostic projection', () => { 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 c32ad66cbab..b2a4c3d33f7 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 @@ -88,6 +88,8 @@ 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; a resume can never claim a log in one of them. */ +const TERMINAL_EXECUTION_LOG_STATUSES: readonly string[] = ['cancelled', 'failed', 'completed'] 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' @@ -135,23 +137,27 @@ class ResumeAdmissionError extends Error { } /** - * Errors from resume attempts that failed without consuming their pause: the - * attempt was refused or could not start, so the paused execution stays - * resumable. Every other failed resume leaves the execution terminal. + * 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. An attempt that changed neither (the + * execution had already finished or was cancelled) records nothing. */ -const pauseRetainingFailures = new WeakSet() +export type FailedResumeOutcome = 'pause_retained' | 'execution_failed' + +const failedResumeOutcomes = new WeakMap() /** - * Whether `error` was thrown by a resume attempt that left its paused execution - * resumable, so a caller mirroring the run's state (a table cell) keeps it paused - * rather than failed. + * The {@link FailedResumeOutcome} recorded on an error rethrown by + * `startResumeExecution`, so a caller mirroring the run's state (a table cell) + * follows the execution rather than guessing from the error. */ -export function wasPausedExecutionRetained(error: unknown): boolean { - return typeof error === 'object' && error !== null && pauseRetainingFailures.has(error) +export function getFailedResumeOutcome(error: unknown): FailedResumeOutcome | undefined { + return typeof error === 'object' && error !== null ? failedResumeOutcomes.get(error) : undefined } -function retainPausedExecution(error: unknown): void { - if (typeof error === 'object' && error !== null) pauseRetainingFailures.add(error) +function recordFailedResumeOutcome(error: unknown, outcome: FailedResumeOutcome | undefined): void { + if (outcome && typeof error === 'object' && error !== null) + failedResumeOutcomes.set(error, outcome) } /** Matches the paused execution mode to the deployment recorded on its durable root log. */ @@ -1007,8 +1013,7 @@ export class PauseResumeManager { const message = toError(error).message await releaseExecutionSlot(resumeEntryId) if (error instanceof ResumeAdmissionError) { - retainPausedExecution(error) - await PauseResumeManager.markResumeAttemptFailed({ + const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, @@ -1017,23 +1022,25 @@ export class PauseResumeManager { preserveForRetry: true, retryable: error.retryable, }) + recordFailedResumeOutcome(error, pauseResumable ? 'pause_retained' : undefined) } else if (message === RUN_BUFFER_UNAVAILABLE_ERROR) { - retainPausedExecution(error) - await PauseResumeManager.markResumeAttemptFailed({ + const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, contextId, failureReason: message, }) + recordFailedResumeOutcome(error, pauseResumable ? 'pause_retained' : undefined) } else { - await PauseResumeManager.markResumeFailed({ + const executionFailed = await PauseResumeManager.markResumeFailed({ resumeEntryId, pausedExecutionId: pausedExecution.id, parentExecutionId: pausedExecution.executionId, contextId, failureReason: message, }) + recordFailedResumeOutcome(error, executionFailed ? 'execution_failed' : undefined) } logger.error( 'Resume execution failed', @@ -2214,10 +2221,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) @@ -2246,7 +2253,7 @@ export class PauseResumeManager { .set({ status: 'cancelled', updatedAt: now, nextResumeAt: null }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) } - return + return false } await tx @@ -2257,7 +2264,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) @@ -2268,6 +2275,8 @@ export class PauseResumeManager { sql`${workflowExecutionLogs.status} != 'cancelled'` ) ) + + return true }) } @@ -2279,10 +2288,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) @@ -2339,7 +2348,7 @@ export class PauseResumeManager { .set({ status: 'cancelled', updatedAt: now, nextResumeAt: null }) .where(eq(pausedExecutions.id, args.pausedExecutionId)) } - return + return false } await tx @@ -2370,7 +2379,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) @@ -2391,6 +2400,13 @@ export class PauseResumeManager { )` ) ) + + return ( + pausedExecution !== undefined && + isResumablePausedStatus(pausedExecution.status) && + executionLog !== undefined && + !TERMINAL_EXECUTION_LOG_STATUSES.includes(executionLog.status) + ) }) } diff --git a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts index 1e19a89ef3e..a070fcaddf3 100644 --- a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts +++ b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts @@ -162,7 +162,9 @@ export const humanInTheLoopManagerMockFns = { } ), mockCreateResumeAttemptTimeoutController: vi.fn(), - mockWasPausedExecutionRetained: vi.fn((_error: unknown): boolean => false), + mockGetFailedResumeOutcome: vi.fn( + (_error: unknown): 'pause_retained' | 'execution_failed' | undefined => undefined + ), mockExtractResumeBillingAttributionFromSnapshot: vi.fn(), mockComputeEarliestResumeAt: vi.fn( (points: Iterable, options: { after?: Date } = {}): Date | null => { @@ -196,7 +198,7 @@ export const humanInTheLoopManagerMock = { humanInTheLoopManagerMockFns.mockUpdateResumeOutputInAggregationBuffers, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, - wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained, + getFailedResumeOutcome: humanInTheLoopManagerMockFns.mockGetFailedResumeOutcome, extractResumeBillingAttributionFromSnapshot: humanInTheLoopManagerMockFns.mockExtractResumeBillingAttributionFromSnapshot, computeEarliestResumeAt: humanInTheLoopManagerMockFns.mockComputeEarliestResumeAt, From d362cc24b035feb6128edcb2eb8faa0be375d343 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 20:38:13 -0700 Subject: [PATCH 3/4] fix(tables): type the terminal execution log statuses against the persisted vocabulary --- .../executor/human-in-the-loop-manager.ts | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) 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 b2a4c3d33f7..3922c640988 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,8 +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; a resume can never claim a log in one of them. */ -const TERMINAL_EXECUTION_LOG_STATUSES: readonly string[] = ['cancelled', 'failed', 'completed'] +/** + * 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' @@ -2405,7 +2417,7 @@ export class PauseResumeManager { pausedExecution !== undefined && isResumablePausedStatus(pausedExecution.status) && executionLog !== undefined && - !TERMINAL_EXECUTION_LOG_STATUSES.includes(executionLog.status) + !isTerminalExecutionLogStatus(executionLog.status) ) }) } From 8c419ed6069de76729131320cf0008a74132a2b6 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 21:44:38 -0700 Subject: [PATCH 4/4] fix(tables): report a failed resume's outcome through a hook instead of the error The outcome rode on the thrown error through a module WeakMap, so it was lost when draining queued resumes threw after the settle. startResumeExecution now takes onAttemptFailed, called right after the settle transaction; a hook failure is logged and never replaces the attempt's error. --- apps/sim/background/resume-execution.ts | 56 +++----- .../resume-governed-subject.test.ts | 41 +++--- .../human-in-the-loop-manager.test.ts | 131 +++++++++++++----- .../executor/human-in-the-loop-manager.ts | 42 +++--- .../mocks/human-in-the-loop-manager.mock.ts | 4 - 5 files changed, 155 insertions(+), 119 deletions(-) diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index 2e9b2a58d8a..db86c516b1a 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -22,7 +22,7 @@ import { classifyWorkflowCellTerminalResult } from '@/lib/table/workflow-cell-re import type { CellResumeContext } from '@/lib/table/workflow-columns' import { createResumeAttemptTimeoutController, - getFailedResumeOutcome, + type FailedResumeOutcome, PauseResumeManager, } from '@/lib/workflows/executor/human-in-the-loop-manager' import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits' @@ -365,23 +365,18 @@ 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, a failed execution fails - * the cell, and an attempt that changed nothing leaves the cell alone. + * a pause that stayed resumable goes back to paused, and a failed execution + * fails the cell. */ -async function writeFailedResumeCellTerminal(writers: CellWriters, error: unknown): Promise { - const outcome = getFailedResumeOutcome(error) - if (!outcome) return - try { - if (outcome === 'pause_retained') { - await writers.writeCellTerminal('paused', null) - } else { - await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed')) - } - } catch (writeError) { - logger.error( - 'Failed to write the cell state after a failed resume', - projectResolvedSecretDiagnosticError(writeError, undefined) - ) +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')) } } @@ -393,22 +388,17 @@ async function runResumeAndCellTerminal( timeoutController: ReturnType ): Promise>> { if (!pausedExecution) throw new Error('Paused execution missing — already nulled by caller') - let result: Awaited> - try { - result = await PauseResumeManager.startResumeExecution({ - resumeEntryId: payload.resumeEntryId, - resumeExecutionId: payload.resumeExecutionId, - pausedExecution, - contextId: payload.contextId, - resumeInput: payload.resumeInput, - userId: payload.userId, - onBlockComplete: writers.cellOnBlockComplete, - abortSignal: signal, - }) - } catch (error) { - await writeFailedResumeCellTerminal(writers, error) - throw error - } + const result = await PauseResumeManager.startResumeExecution({ + resumeEntryId: payload.resumeEntryId, + resumeExecutionId: payload.resumeExecutionId, + pausedExecution, + contextId: payload.contextId, + resumeInput: payload.resumeInput, + userId: payload.userId, + onBlockComplete: writers.cellOnBlockComplete, + onAttemptFailed: (outcome, error) => writeFailedResumeCellTerminal(writers, outcome, error), + abortSignal: signal, + }) if (result.status === 'paused') { await writers.writeCellTerminal('paused', null) diff --git a/apps/sim/background/resume-governed-subject.test.ts b/apps/sim/background/resume-governed-subject.test.ts index 75f895ae9b4..d7589d6eb45 100644 --- a/apps/sim/background/resume-governed-subject.test.ts +++ b/apps/sim/background/resume-governed-subject.test.ts @@ -51,13 +51,13 @@ 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 = { ...hoisted, getPausedExecutionById: humanInTheLoopManagerMockFns.mockGetPausedExecutionById, startResumeExecution: humanInTheLoopManagerMockFns.mockStartResumeExecution, - getFailedResumeOutcome: humanInTheLoopManagerMockFns.mockGetFailedResumeOutcome, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, findCellContextByExecutionId: tableWorkflowColumnsMockFns.mockFindCellContextByExecutionId, @@ -226,10 +226,23 @@ describe('resuming a paused table cell', () => { 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') - mocks.startResumeExecution.mockRejectedValueOnce(runFailure) - mocks.getFailedResumeOutcome.mockReturnValueOnce('execution_failed') + failResume('execution_failed', runFailure) await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) @@ -242,8 +255,7 @@ describe('resuming a paused table cell', () => { it('puts the cell back to paused when the pause stayed resumable', async () => { const admissionRefusal = new Error('Execution can no longer be resumed') - mocks.startResumeExecution.mockRejectedValueOnce(admissionRefusal) - mocks.getFailedResumeOutcome.mockReturnValueOnce('pause_retained') + failResume('pause_retained', admissionRefusal) await expect(executeResumeJob(PAYLOAD)).rejects.toBe(admissionRefusal) @@ -253,24 +265,5 @@ describe('resuming a paused table cell', () => { jobId: 'paused-parent-execution-1', }) }, 20_000) - - it('leaves the cell alone when the failed attempt changed nothing', async () => { - const staleRefusal = new Error('Execution can no longer be resumed') - mocks.startResumeExecution.mockRejectedValueOnce(staleRefusal) - mocks.getFailedResumeOutcome.mockReturnValueOnce(undefined) - - await expect(executeResumeJob(PAYLOAD)).rejects.toBe(staleRefusal) - - expect(lastCellExecutionState()).toBeUndefined() - }, 20_000) - - it('still reports the resume failure when the cell write also fails', async () => { - const runFailure = new Error('Block failed') - mocks.startResumeExecution.mockRejectedValueOnce(runFailure) - mocks.getFailedResumeOutcome.mockReturnValueOnce('execution_failed') - mocks.writeWorkflowGroupState.mockRejectedValueOnce(new Error('Database unavailable')) - - await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure) - }, 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 4f9159ae90f..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,7 +61,7 @@ vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadat import { createResumeAttemptTimeoutController, - getFailedResumeOutcome, + type FailedResumeOutcome, PauseResumeManager, requireResumeDeploymentVersion, updateResumeOutputInAggregationBuffers, @@ -262,56 +262,113 @@ describe('what a failed resume did to its paused execution', () => { ) it.each([ - { stillResumable: true, outcome: 'pause_retained' }, - { stillResumable: false, outcome: undefined }, + { logStatus: 'running', pauseStatus: 'paused', executionFailed: true }, + { logStatus: 'cancelled', pauseStatus: 'paused', executionFailed: false }, + { logStatus: 'running', pauseStatus: 'cancelling', executionFailed: false }, ])( - 'records a refused attempt as $outcome when the pause is resumable: $stillResumable', - async ({ stillResumable, outcome }) => { + '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 { - const thrown = await PauseResumeManager.startResumeExecution(resumeArgs).catch( - (error: unknown) => error - ) - - expect(thrown).toMatchObject({ name: 'ResumeAdmissionError' }) - expect(getFailedResumeOutcome(thrown)).toBe(outcome) + await expect(PauseResumeManager.startResumeExecution(args)).rejects.toMatchObject({ + name: 'ResumeAdmissionError', + }) + expect(outcomes).toEqual(reported) } finally { markResumeAttemptFailedSpy.mockRestore() } } ) - it.each([ - { executionFailed: true, outcome: 'execution_failed' }, - { executionFailed: false, outcome: undefined }, - ])( - 'records a failed run as $outcome when it failed the execution: $executionFailed', - async ({ executionFailed, outcome }) => { - const rawError = new Error('Block failed') + 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 - const runResumeExecutionSpy = vi - .spyOn(managerInternals, 'runResumeExecution') - .mockRejectedValueOnce(rawError) - const markResumeFailedSpy = vi - .spyOn(managerInternals, 'markResumeFailed') - .mockResolvedValueOnce(executionFailed) - const processQueuedResumesSpy = vi - .spyOn(PauseResumeManager, 'processQueuedResumes') - .mockResolvedValueOnce() + 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() + ) + } - try { - await expect(PauseResumeManager.startResumeExecution(resumeArgs)).rejects.toBe(rawError) - expect(getFailedResumeOutcome(rawError)).toBe(outcome) - } finally { - runResumeExecutionSpy.mockRestore() - markResumeFailedSpy.mockRestore() - processQueuedResumesSpy.mockRestore() + 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', () => { 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 3922c640988..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 @@ -151,27 +151,10 @@ 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. An attempt that changed neither (the - * execution had already finished or was cancelled) records nothing. + * resumed run failed the execution. */ export type FailedResumeOutcome = 'pause_retained' | 'execution_failed' -const failedResumeOutcomes = new WeakMap() - -/** - * The {@link FailedResumeOutcome} recorded on an error rethrown by - * `startResumeExecution`, so a caller mirroring the run's state (a table cell) - * follows the execution rather than guessing from the error. - */ -export function getFailedResumeOutcome(error: unknown): FailedResumeOutcome | undefined { - return typeof error === 'object' && error !== null ? failedResumeOutcomes.get(error) : undefined -} - -function recordFailedResumeOutcome(error: unknown, outcome: FailedResumeOutcome | undefined): void { - if (outcome && typeof error === 'object' && error !== null) - failedResumeOutcomes.set(error, outcome) -} - /** Matches the paused execution mode to the deployment recorded on its durable root log. */ export function requireResumeDeploymentVersion( useDraftState: unknown, @@ -457,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 } @@ -902,6 +892,7 @@ export class PauseResumeManager { sendEvent, onStream, onBlockComplete, + onAttemptFailed, abortSignal, } = args @@ -1024,6 +1015,7 @@ export class PauseResumeManager { } catch (error) { const message = toError(error).message await releaseExecutionSlot(resumeEntryId) + let outcome: FailedResumeOutcome | undefined if (error instanceof ResumeAdmissionError) { const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, @@ -1034,7 +1026,7 @@ export class PauseResumeManager { preserveForRetry: true, retryable: error.retryable, }) - recordFailedResumeOutcome(error, pauseResumable ? 'pause_retained' : undefined) + if (pauseResumable) outcome = 'pause_retained' } else if (message === RUN_BUFFER_UNAVAILABLE_ERROR) { const pauseResumable = await PauseResumeManager.markResumeAttemptFailed({ resumeEntryId, @@ -1043,7 +1035,7 @@ export class PauseResumeManager { contextId, failureReason: message, }) - recordFailedResumeOutcome(error, pauseResumable ? 'pause_retained' : undefined) + if (pauseResumable) outcome = 'pause_retained' } else { const executionFailed = await PauseResumeManager.markResumeFailed({ resumeEntryId, @@ -1052,7 +1044,15 @@ export class PauseResumeManager { contextId, failureReason: message, }) - recordFailedResumeOutcome(error, executionFailed ? 'execution_failed' : undefined) + 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', diff --git a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts index a070fcaddf3..c9011d2bc47 100644 --- a/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts +++ b/packages/testing/src/mocks/human-in-the-loop-manager.mock.ts @@ -162,9 +162,6 @@ export const humanInTheLoopManagerMockFns = { } ), mockCreateResumeAttemptTimeoutController: vi.fn(), - mockGetFailedResumeOutcome: vi.fn( - (_error: unknown): 'pause_retained' | 'execution_failed' | undefined => undefined - ), mockExtractResumeBillingAttributionFromSnapshot: vi.fn(), mockComputeEarliestResumeAt: vi.fn( (points: Iterable, options: { after?: Date } = {}): Date | null => { @@ -198,7 +195,6 @@ export const humanInTheLoopManagerMock = { humanInTheLoopManagerMockFns.mockUpdateResumeOutputInAggregationBuffers, createResumeAttemptTimeoutController: humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController, - getFailedResumeOutcome: humanInTheLoopManagerMockFns.mockGetFailedResumeOutcome, extractResumeBillingAttributionFromSnapshot: humanInTheLoopManagerMockFns.mockExtractResumeBillingAttributionFromSnapshot, computeEarliestResumeAt: humanInTheLoopManagerMockFns.mockComputeEarliestResumeAt,