Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions apps/sim/background/resume-execution.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand All @@ -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'
Expand Down Expand Up @@ -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<void> {
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<ReturnType<typeof PauseResumeManager.getPausedExecutionById>>,
Expand All @@ -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,
})

Expand Down
49 changes: 49 additions & 0 deletions apps/sim/background/resume-governed-subject.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down Expand Up @@ -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<void>
}) => {
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)
})
})
165 changes: 163 additions & 2 deletions apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -61,6 +61,7 @@ vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadat

import {
createResumeAttemptTimeoutController,
type FailedResumeOutcome,
PauseResumeManager,
requireResumeDeploymentVersion,
updateResumeOutputInAggregationBuffers,
Expand All @@ -84,7 +85,7 @@ if (!humanInTheLoopLogger) {
}

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

Expand Down Expand Up @@ -210,6 +211,166 @@ describe('queued resume attempt deadlines', () => {
})
})

describe('what a failed resume did to its paused execution', () => {
type StartResumeArgs = Parameters<typeof PauseResumeManager.startResumeExecution>[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<void>) {
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()
Expand Down
Loading
Loading