Skip to content

Commit 1918976

Browse files
committed
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.
1 parent 66e5705 commit 1918976

5 files changed

Lines changed: 174 additions & 10 deletions

File tree

‎apps/sim/background/resume-execution.ts‎

Lines changed: 39 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { createLogger } from '@sim/logger'
2+
import { getErrorMessage } from '@sim/utils/errors'
23
import { generateId } from '@sim/utils/id'
34
import { task, timeout } from '@trigger.dev/sdk'
45
import {
@@ -22,6 +23,7 @@ import type { CellResumeContext } from '@/lib/table/workflow-columns'
2223
import {
2324
createResumeAttemptTimeoutController,
2425
PauseResumeManager,
26+
wasPausedExecutionRetained,
2527
} from '@/lib/workflows/executor/human-in-the-loop-manager'
2628
import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits'
2729
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
@@ -359,6 +361,27 @@ async function buildResumeCellWriters(
359361
return { cellOnBlockComplete, writeCellTerminal }
360362
}
361363

364+
/**
365+
* A resume that throws never reaches the terminal write in
366+
* {@link runResumeAndCellTerminal}, which would leave the cell showing its last
367+
* partial `running` state. Mirror the execution instead: a pause the resume kept
368+
* resumable goes back to paused, and any other failure ended the run.
369+
*/
370+
async function writeFailedResumeCellTerminal(writers: CellWriters, error: unknown): Promise<void> {
371+
try {
372+
if (wasPausedExecutionRetained(error)) {
373+
await writers.writeCellTerminal('paused', null)
374+
} else {
375+
await writers.writeCellTerminal('error', getErrorMessage(error, 'Resume execution failed'))
376+
}
377+
} catch (writeError) {
378+
logger.error(
379+
'Failed to write the cell state after a failed resume',
380+
projectResolvedSecretDiagnosticError(writeError, undefined)
381+
)
382+
}
383+
}
384+
362385
async function runResumeAndCellTerminal(
363386
payload: ResumeExecutionPayload,
364387
pausedExecution: Awaited<ReturnType<typeof PauseResumeManager.getPausedExecutionById>>,
@@ -367,16 +390,22 @@ async function runResumeAndCellTerminal(
367390
timeoutController: ReturnType<typeof createTimeoutAbortController>
368391
): Promise<Awaited<ReturnType<typeof PauseResumeManager.startResumeExecution>>> {
369392
if (!pausedExecution) throw new Error('Paused execution missing — already nulled by caller')
370-
const result = await PauseResumeManager.startResumeExecution({
371-
resumeEntryId: payload.resumeEntryId,
372-
resumeExecutionId: payload.resumeExecutionId,
373-
pausedExecution,
374-
contextId: payload.contextId,
375-
resumeInput: payload.resumeInput,
376-
userId: payload.userId,
377-
onBlockComplete: writers.cellOnBlockComplete,
378-
abortSignal: signal,
379-
})
393+
let result: Awaited<ReturnType<typeof PauseResumeManager.startResumeExecution>>
394+
try {
395+
result = await PauseResumeManager.startResumeExecution({
396+
resumeEntryId: payload.resumeEntryId,
397+
resumeExecutionId: payload.resumeExecutionId,
398+
pausedExecution,
399+
contextId: payload.contextId,
400+
resumeInput: payload.resumeInput,
401+
userId: payload.userId,
402+
onBlockComplete: writers.cellOnBlockComplete,
403+
abortSignal: signal,
404+
})
405+
} catch (error) {
406+
await writeFailedResumeCellTerminal(writers, error)
407+
throw error
408+
}
380409

381410
if (result.status === 'paused') {
382411
await writers.writeCellTerminal('paused', null)

‎apps/sim/background/resume-governed-subject.test.ts‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@ const mocks = {
5757
...hoisted,
5858
getPausedExecutionById: humanInTheLoopManagerMockFns.mockGetPausedExecutionById,
5959
startResumeExecution: humanInTheLoopManagerMockFns.mockStartResumeExecution,
60+
wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained,
6061
createResumeAttemptTimeoutController:
6162
humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController,
6263
findCellContextByExecutionId: tableWorkflowColumnsMockFns.mockFindCellContextByExecutionId,
@@ -217,4 +218,47 @@ describe('resuming a paused table cell', () => {
217218
const [cascadePayload] = mocks.runRowCascadeLoop.mock.calls[0]
218219
expect(cascadePayload.capabilityGovernedUserId).toBe('requesting-member')
219220
}, 20_000)
221+
222+
describe('when the resume throws', () => {
223+
/** The execution state the last cell write persisted. */
224+
function lastCellExecutionState() {
225+
const [, payload] = mocks.writeWorkflowGroupState.mock.calls.at(-1) ?? []
226+
return payload?.executionState
227+
}
228+
229+
it('marks the cell failed when the resumed run itself failed', async () => {
230+
const runFailure = new Error('writeLedger: Unique constraint violation')
231+
mocks.startResumeExecution.mockRejectedValueOnce(runFailure)
232+
233+
await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure)
234+
235+
expect(lastCellExecutionState()).toMatchObject({
236+
status: 'error',
237+
executionId: 'parent-execution-1',
238+
error: 'writeLedger: Unique constraint violation',
239+
})
240+
}, 20_000)
241+
242+
it('puts the cell back to paused when the pause stayed resumable', async () => {
243+
const admissionRefusal = new Error('Execution can no longer be resumed')
244+
mocks.startResumeExecution.mockRejectedValueOnce(admissionRefusal)
245+
mocks.wasPausedExecutionRetained.mockReturnValueOnce(true)
246+
247+
await expect(executeResumeJob(PAYLOAD)).rejects.toBe(admissionRefusal)
248+
249+
expect(lastCellExecutionState()).toMatchObject({
250+
status: 'pending',
251+
executionId: 'parent-execution-1',
252+
jobId: 'paused-parent-execution-1',
253+
})
254+
}, 20_000)
255+
256+
it('still reports the resume failure when the cell write also fails', async () => {
257+
const runFailure = new Error('Block failed')
258+
mocks.startResumeExecution.mockRejectedValueOnce(runFailure)
259+
mocks.writeWorkflowGroupState.mockRejectedValueOnce(new Error('Database unavailable'))
260+
261+
await expect(executeResumeJob(PAYLOAD)).rejects.toBe(runFailure)
262+
}, 20_000)
263+
})
220264
})

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

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ import {
6464
PauseResumeManager,
6565
requireResumeDeploymentVersion,
6666
updateResumeOutputInAggregationBuffers,
67+
wasPausedExecutionRetained,
6768
} from '@/lib/workflows/executor/human-in-the-loop-manager'
6869
import { getAutomaticResumeWaitingMetadata } from '@/lib/workflows/executor/paused-execution-metadata'
6970
import { AUTOMATIC_RESUME_WAITING_REASON_MAX_LENGTH } from '@/lib/workflows/executor/resume-policy'
@@ -210,6 +211,72 @@ describe('queued resume attempt deadlines', () => {
210211
})
211212
})
212213

214+
describe('which failed resumes keep the pause resumable', () => {
215+
type StartResumeArgs = Parameters<typeof PauseResumeManager.startResumeExecution>[0]
216+
const pausedExecution = {
217+
id: 'paused-execution-1',
218+
workflowId: 'workflow-1',
219+
executionId: 'parent-execution-1',
220+
pausePoints: {
221+
'context-1': { contextId: 'context-1', blockId: 'hitl-1' },
222+
},
223+
executionSnapshot: createSnapshotSeed(),
224+
metadata: {},
225+
} as StartResumeArgs['pausedExecution']
226+
const resumeArgs: StartResumeArgs = {
227+
resumeEntryId: 'resume-entry-1',
228+
resumeExecutionId: 'resume-execution-1',
229+
pausedExecution,
230+
contextId: 'context-1',
231+
resumeInput: { approved: true },
232+
userId: 'user-1',
233+
}
234+
235+
beforeEach(() => {
236+
resetDbChainMock()
237+
})
238+
239+
it('keeps the pause resumable when the paused log can no longer be claimed', async () => {
240+
const markResumeAttemptFailedSpy = vi
241+
.spyOn(PauseResumeManager, 'markResumeAttemptFailed')
242+
.mockResolvedValueOnce()
243+
244+
try {
245+
const thrown = await PauseResumeManager.startResumeExecution(resumeArgs).catch(
246+
(error: unknown) => error
247+
)
248+
249+
expect(thrown).toMatchObject({ name: 'ResumeAdmissionError' })
250+
expect(wasPausedExecutionRetained(thrown)).toBe(true)
251+
} finally {
252+
markResumeAttemptFailedSpy.mockRestore()
253+
}
254+
})
255+
256+
it('does not keep the pause when the resumed run itself failed', async () => {
257+
const rawError = new Error('Block failed')
258+
const managerInternals = PauseResumeManager as unknown as PauseResumeManagerInternals
259+
const runResumeExecutionSpy = vi
260+
.spyOn(managerInternals, 'runResumeExecution')
261+
.mockRejectedValueOnce(rawError)
262+
const markResumeFailedSpy = vi
263+
.spyOn(managerInternals, 'markResumeFailed')
264+
.mockResolvedValueOnce()
265+
const processQueuedResumesSpy = vi
266+
.spyOn(PauseResumeManager, 'processQueuedResumes')
267+
.mockResolvedValueOnce()
268+
269+
try {
270+
await expect(PauseResumeManager.startResumeExecution(resumeArgs)).rejects.toBe(rawError)
271+
expect(wasPausedExecutionRetained(rawError)).toBe(false)
272+
} finally {
273+
runResumeExecutionSpy.mockRestore()
274+
markResumeFailedSpy.mockRestore()
275+
processQueuedResumesSpy.mockRestore()
276+
}
277+
})
278+
})
279+
213280
describe('resume failure diagnostic projection', () => {
214281
beforeEach(() => {
215282
resetDbChainMock()

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

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,26 @@ class ResumeAdmissionError extends Error {
134134
}
135135
}
136136

137+
/**
138+
* Errors from resume attempts that failed without consuming their pause: the
139+
* attempt was refused or could not start, so the paused execution stays
140+
* resumable. Every other failed resume leaves the execution terminal.
141+
*/
142+
const pauseRetainingFailures = new WeakSet<object>()
143+
144+
/**
145+
* Whether `error` was thrown by a resume attempt that left its paused execution
146+
* resumable, so a caller mirroring the run's state (a table cell) keeps it paused
147+
* rather than failed.
148+
*/
149+
export function wasPausedExecutionRetained(error: unknown): boolean {
150+
return typeof error === 'object' && error !== null && pauseRetainingFailures.has(error)
151+
}
152+
153+
function retainPausedExecution(error: unknown): void {
154+
if (typeof error === 'object' && error !== null) pauseRetainingFailures.add(error)
155+
}
156+
137157
/** Matches the paused execution mode to the deployment recorded on its durable root log. */
138158
export function requireResumeDeploymentVersion(
139159
useDraftState: unknown,
@@ -987,6 +1007,7 @@ export class PauseResumeManager {
9871007
const message = toError(error).message
9881008
await releaseExecutionSlot(resumeEntryId)
9891009
if (error instanceof ResumeAdmissionError) {
1010+
retainPausedExecution(error)
9901011
await PauseResumeManager.markResumeAttemptFailed({
9911012
resumeEntryId,
9921013
pausedExecutionId: pausedExecution.id,
@@ -997,6 +1018,7 @@ export class PauseResumeManager {
9971018
retryable: error.retryable,
9981019
})
9991020
} else if (message === RUN_BUFFER_UNAVAILABLE_ERROR) {
1021+
retainPausedExecution(error)
10001022
await PauseResumeManager.markResumeAttemptFailed({
10011023
resumeEntryId,
10021024
pausedExecutionId: pausedExecution.id,

‎packages/testing/src/mocks/human-in-the-loop-manager.mock.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,7 @@ export const humanInTheLoopManagerMockFns = {
162162
}
163163
),
164164
mockCreateResumeAttemptTimeoutController: vi.fn(),
165+
mockWasPausedExecutionRetained: vi.fn((_error: unknown): boolean => false),
165166
mockExtractResumeBillingAttributionFromSnapshot: vi.fn(),
166167
mockComputeEarliestResumeAt: vi.fn(
167168
(points: Iterable<MockPausePoint>, options: { after?: Date } = {}): Date | null => {
@@ -195,6 +196,7 @@ export const humanInTheLoopManagerMock = {
195196
humanInTheLoopManagerMockFns.mockUpdateResumeOutputInAggregationBuffers,
196197
createResumeAttemptTimeoutController:
197198
humanInTheLoopManagerMockFns.mockCreateResumeAttemptTimeoutController,
199+
wasPausedExecutionRetained: humanInTheLoopManagerMockFns.mockWasPausedExecutionRetained,
198200
extractResumeBillingAttributionFromSnapshot:
199201
humanInTheLoopManagerMockFns.mockExtractResumeBillingAttributionFromSnapshot,
200202
computeEarliestResumeAt: humanInTheLoopManagerMockFns.mockComputeEarliestResumeAt,

0 commit comments

Comments
 (0)