Skip to content

Commit 5f56ea9

Browse files
committed
fix(execution): complete background jobs on workflow failures, fault only on platform errors
Async API and resume jobs re-threw every execution error, so Trigger.dev marked the run failed and alerted on user workflow failures (a block's own error, a missing required field). The workflow-execution task also checked core's finalized signal before core had set it, so its existing guard never applied. - Add one classifier (lib/workflows/executor/job-failure) built on classifyExecutionError: a failure core recorded and attributed to a block is the workflow's outcome and the job completes with success: false; anything else (engine, setup, unrecorded, or a programming error in Sim code) faults. - Await core's post-execution work before classifying. - Carry programming errors thrown inside tools across the executeTool flattening as ToolResponse.isSystemError so they still fault the job. - Throw WorkflowValidationError with block attribution for missing required fields so pre-execution validation is attributed to its block. - Apply the rule to workflow-execution, resume-execution and schedule-execution; schedule re-throws platform faults only after its own failure bookkeeping. - Project a completed failure result back to failed in /api/jobs and the queue-job status fallback so pollers still see a failed run.
1 parent a657493 commit 5f56ea9

19 files changed

Lines changed: 693 additions & 9 deletions

‎apps/sim/app/api/jobs/[jobId]/route.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { checkHybridAuth } from '@/lib/auth/hybrid'
88
import { getJobQueue } from '@/lib/core/async-jobs'
99
import { generateRequestId } from '@/lib/core/utils/request'
1010
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
11+
import { projectWorkflowJobOutcome } from '@/lib/workflows/executor/job-outcome'
1112
import { createErrorResponse } from '@/app/api/workflows/utils'
1213

1314
const logger = createLogger('TaskStatusAPI')
@@ -66,15 +67,16 @@ export const GET = withRouteHandler(
6667
return createErrorResponse('Access denied', 403)
6768
}
6869

70+
const outcome = projectWorkflowJobOutcome(job)
6971
const response: Record<string, unknown> = {
7072
success: true,
7173
taskId,
72-
status: job.status,
74+
status: outcome.status,
7375
metadata: job.metadata,
7476
}
7577

7678
if (job.output !== undefined) response.output = job.output
77-
if (job.error !== undefined) response.error = job.error
79+
if (outcome.error !== undefined) response.error = outcome.error
7880

7981
return NextResponse.json(response)
8082
} catch (error: unknown) {

‎apps/sim/background/async-preprocessing-correlation.test.ts‎

Lines changed: 74 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -74,10 +74,13 @@ vi.mock('@/executor/execution/snapshot', () => ({
7474
ExecutionSnapshot: mockExecutionSnapshot,
7575
}))
7676

77-
vi.mock('@/executor/utils/errors', () => ({
77+
vi.mock('@/executor/utils/errors', async (importOriginal) => ({
78+
...(await importOriginal<typeof import('@/executor/utils/errors')>()),
7879
hasExecutionResult: mockHasExecutionResult,
7980
}))
8081

82+
import { buildBlockExecutionError } from '@/executor/utils/errors'
83+
import type { SerializedBlock } from '@/serializer/types'
8184
import { executeScheduleJob } from './schedule-execution'
8285
import { executeWorkflowJob } from './workflow-execution'
8386

@@ -395,7 +398,8 @@ describe('async preprocessing correlation threading', () => {
395398
})
396399
).rejects.toBe(rawError)
397400

398-
expect(loggingSessionMockFns.mockWaitForPostExecution).not.toHaveBeenCalled()
401+
// Core finalizes after throwing, so the task must settle that work before deciding.
402+
expect(loggingSessionMockFns.mockWaitForPostExecution).toHaveBeenCalled()
399403
expect(mockWasExecutionFinalizedByCore).toHaveBeenCalledWith(rawError, 'execution-finalized')
400404
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
401405
})
@@ -625,4 +629,72 @@ describe('async preprocessing correlation threading', () => {
625629
})
626630
)
627631
})
632+
633+
describe('scheduled run failures', () => {
634+
const schedulePayload = {
635+
scheduleId: 'schedule-1',
636+
workflowId: 'workflow-1',
637+
workspaceId: 'workspace-1',
638+
billingAttribution,
639+
now: '2025-01-01T00:00:00.000Z',
640+
scheduledFor: '2025-01-01T00:00:00.000Z',
641+
}
642+
643+
beforeEach(() => {
644+
mockPreprocessExecution.mockResolvedValueOnce({
645+
success: true,
646+
actorUserId: 'actor-1',
647+
workflowRecord: {
648+
id: 'workflow-1',
649+
userId: 'owner-1',
650+
workspaceId: 'workspace-1',
651+
variables: {},
652+
},
653+
billingAttribution,
654+
executionTimeout: {},
655+
})
656+
})
657+
658+
it('faults the job on a failure core never recorded, after recording the schedule failure', async () => {
659+
const engineError = new Error('Workflow state not found')
660+
mockExecuteWorkflowCore.mockRejectedValueOnce(engineError)
661+
662+
await expect(
663+
executeScheduleJob({
664+
...schedulePayload,
665+
executionId: 'execution-schedule-fault',
666+
requestId: 'request-schedule-fault',
667+
})
668+
).rejects.toBe(engineError)
669+
670+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
671+
expect.objectContaining({ lastQueuedAt: null, lastFailedAt: expect.any(Date) })
672+
)
673+
})
674+
675+
it('completes the job when the workflow failed in a block core recorded', async () => {
676+
mockExecuteWorkflowCore.mockRejectedValueOnce(
677+
buildBlockExecutionError({
678+
block: {
679+
id: 'plan-panels',
680+
metadata: { id: 'function', name: 'planPanels' },
681+
} as SerializedBlock,
682+
error: new Error("ValueError: kind ''"),
683+
})
684+
)
685+
mockWasExecutionFinalizedByCore.mockReturnValue(true)
686+
687+
await expect(
688+
executeScheduleJob({
689+
...schedulePayload,
690+
executionId: 'execution-schedule-failure',
691+
requestId: 'request-schedule-failure',
692+
})
693+
).resolves.toBeUndefined()
694+
695+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
696+
expect.objectContaining({ lastQueuedAt: null, lastFailedAt: expect.any(Date) })
697+
)
698+
})
699+
})
628700
})

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

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@ vi.mock('@/executor/execution/snapshot', () => ({
2929
}))
3030

3131
import { executeResumeJob, type ResumeExecutionPayload } from '@/background/resume-execution'
32+
import { buildBlockExecutionError } from '@/executor/utils/errors'
33+
import type { SerializedBlock } from '@/serializer/types'
3234

3335
const { mockFindCellContextByExecutionId } = tableWorkflowColumnsMockFns
3436
const {
@@ -104,6 +106,29 @@ describe('executeResumeJob terminal errors', () => {
104106
expect(rawError.message).toContain(secret)
105107
})
106108

109+
it('completes the job when the resumed workflow failed in a block core recorded', async () => {
110+
const blockError = Object.assign(
111+
buildBlockExecutionError({
112+
block: {
113+
id: 'plan-panels',
114+
metadata: { id: 'function', name: 'planPanels' },
115+
} as SerializedBlock,
116+
error: new Error("ValueError: kind ''"),
117+
}),
118+
{ executionFinalizedByCore: true }
119+
)
120+
mockStartResumeExecution.mockRejectedValue(blockError)
121+
122+
await expect(executeResumeJob(payload)).resolves.toMatchObject({
123+
success: false,
124+
workflowId: 'workflow-1',
125+
executionId: 'resume-execution-1',
126+
parentExecutionId: 'parent-execution-1',
127+
status: 'failed',
128+
error: blockError.message,
129+
})
130+
})
131+
107132
it('starts a legacy attempt deadline before deserializing the full snapshot', async () => {
108133
mockStartResumeExecution.mockResolvedValue({
109134
success: true,

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@ import {
2323
createResumeAttemptTimeoutController,
2424
PauseResumeManager,
2525
} from '@/lib/workflows/executor/human-in-the-loop-manager'
26+
import {
27+
buildWorkflowJobFailureResult,
28+
classifySettledWorkflowJobFailure,
29+
} from '@/lib/workflows/executor/job-failure'
2630
import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits'
2731
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
2832
import type { SerializedSnapshot } from '@/executor/types'
@@ -230,6 +234,15 @@ export async function executeResumeJob(payload: ResumeExecutionPayload, signal?:
230234
workflowId,
231235
})
232236
)
237+
// The resumed run executes under its parent's id, and the manager settles
238+
// its post-execution work before re-throwing.
239+
if (classifySettledWorkflowJobFailure(error, parentExecutionId) === 'workflow_failure') {
240+
return {
241+
...buildWorkflowJobFailureResult({ error, workflowId, executionId: resumeExecutionId }),
242+
parentExecutionId,
243+
status: 'failed' as const,
244+
}
245+
}
233246
throw error
234247
} finally {
235248
timeoutController?.cleanup()

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

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@ import {
3838
executeWorkflowCore,
3939
wasExecutionFinalizedByCore,
4040
} from '@/lib/workflows/executor/execution-core'
41+
import { classifyWorkflowJobFailure } from '@/lib/workflows/executor/job-failure'
4142
import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence'
4243
import { loadDeployedWorkflowState } from '@/lib/workflows/persistence/utils'
4344
import { notifyScheduleAutoDisabled } from '@/lib/workflows/schedules/disable-notifications'
@@ -855,6 +856,13 @@ export async function executeScheduleJob(
855856
disableReason,
856857
})
857858

859+
/**
860+
* A platform fault, re-thrown only after the schedule's own bookkeeping
861+
* (failure count, next run, claim) has run, so faulting the job to alert
862+
* on it never leaves the schedule claimed or its cadence stalled.
863+
*/
864+
let jobFault: unknown
865+
858866
try {
859867
const [scheduleRecord] = await db
860868
.select({
@@ -1203,6 +1211,9 @@ export async function executeScheduleJob(
12031211
`Error updating schedule ${payload.scheduleId} after execution error`,
12041212
'consecutive_failures'
12051213
)
1214+
1215+
const failure = await classifyWorkflowJobFailure({ error, executionId, loggingSession })
1216+
if (failure === 'job_fault') jobFault = error
12061217
}
12071218
} catch (error: unknown) {
12081219
try {
@@ -1211,6 +1222,7 @@ export async function executeScheduleJob(
12111222
return
12121223
}
12131224

1225+
jobFault = error
12141226
logger.error(`[${requestId}] Error processing schedule ${payload.scheduleId}`, error, {
12151227
cause: describeError(error),
12161228
})
@@ -1230,6 +1242,8 @@ export async function executeScheduleJob(
12301242
trace.getActiveSpan()?.recordException(toError(recoveryError))
12311243
}
12321244
}
1245+
1246+
if (jobFault !== undefined) throw jobFault
12331247
})
12341248
} finally {
12351249
timeoutController.cleanup()
Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,143 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import {
5+
executionPreprocessingMock,
6+
executionPreprocessingMockFns,
7+
LoggingSessionMock,
8+
loggingSessionMock,
9+
loggingSessionMockFns,
10+
} from '@sim/testing'
11+
import {
12+
executionLimitsMock,
13+
executionLimitsMockFns,
14+
} from '@sim/testing/mocks/execution-limits.mock'
15+
import { beforeEach, describe, expect, it, vi } from 'vitest'
16+
17+
const { mockExecuteWorkflowCore, mockWasExecutionFinalizedByCore } = vi.hoisted(() => ({
18+
mockExecuteWorkflowCore: vi.fn(),
19+
mockWasExecutionFinalizedByCore: vi.fn(),
20+
}))
21+
22+
vi.mock('@/lib/execution/preprocessing', () => executionPreprocessingMock)
23+
vi.mock('@/lib/logs/execution/logging-session', () => loggingSessionMock)
24+
vi.mock('@/lib/core/execution-limits', () => executionLimitsMock)
25+
vi.mock('@/lib/workflows/executor/execution-core', () => ({
26+
executeWorkflowCore: mockExecuteWorkflowCore,
27+
wasExecutionFinalizedByCore: mockWasExecutionFinalizedByCore,
28+
}))
29+
vi.mock('@/lib/workflows/executor/pause-persistence', () => ({
30+
handlePostExecutionPauseState: vi.fn(),
31+
}))
32+
vi.mock('@/lib/logs/execution/trace-spans/trace-spans', () => ({
33+
buildTraceSpans: vi.fn(() => ({ traceSpans: [] })),
34+
}))
35+
vi.mock('@/executor/execution/snapshot', () => ({ ExecutionSnapshot: vi.fn() }))
36+
vi.mock('@/lib/uploads/utils/user-file-base64.server', () => ({
37+
cleanupExecutionBase64Cache: vi.fn(async () => {}),
38+
}))
39+
40+
import * as usageReservation from '@/lib/billing/calculations/usage-reservation'
41+
import { executeWorkflowJob, type WorkflowExecutionPayload } from '@/background/workflow-execution'
42+
import { buildBlockExecutionError } from '@/executor/utils/errors'
43+
import type { SerializedBlock } from '@/serializer/types'
44+
45+
const billingAttribution = {
46+
actorUserId: 'user-1',
47+
workspaceId: 'workspace-1',
48+
organizationId: null,
49+
billedAccountUserId: 'user-1',
50+
billingEntity: { type: 'user' as const, id: 'user-1' },
51+
billingPeriod: { start: '2026-09-01T00:00:00.000Z', end: '2026-10-01T00:00:00.000Z' },
52+
payerSubscription: null,
53+
}
54+
55+
const payload: WorkflowExecutionPayload = {
56+
workflowId: 'workflow-1',
57+
principal: {
58+
version: 1,
59+
principal: {
60+
kind: 'system',
61+
serviceId: 'internal',
62+
workspaceId: 'workspace-1',
63+
workflowId: 'workflow-1',
64+
},
65+
},
66+
userId: 'user-1',
67+
billingAttribution,
68+
workspaceId: 'workspace-1',
69+
executionId: 'execution-1',
70+
requestId: 'request-1',
71+
triggerType: 'api',
72+
}
73+
74+
const planPanels = {
75+
id: 'plan-panels',
76+
metadata: { id: 'function', name: 'planPanels' },
77+
} as SerializedBlock
78+
79+
describe('executeWorkflowJob fault vs workflow failure', () => {
80+
beforeEach(() => {
81+
vi.spyOn(usageReservation, 'refreshExecutionSlotExpiry').mockResolvedValue(true)
82+
vi.spyOn(usageReservation, 'releaseExecutionSlot').mockResolvedValue(undefined)
83+
executionLimitsMockFns.mockCreateTimeoutAbortController.mockImplementation(() => ({
84+
signal: new AbortController().signal,
85+
cleanup: vi.fn(),
86+
abort: vi.fn(),
87+
isTimedOut: () => false,
88+
timeoutMs: 120_000,
89+
}))
90+
LoggingSessionMock.mockImplementation(function LoggingSession() {
91+
return {
92+
safeCompleteWithError: loggingSessionMockFns.mockSafeCompleteWithError,
93+
waitForPostExecution: loggingSessionMockFns.mockWaitForPostExecution,
94+
markAsFailed: loggingSessionMockFns.mockMarkAsFailed,
95+
setExecutionDeadlineAt: loggingSessionMockFns.mockSetExecutionDeadlineAt,
96+
projectDiagnosticError: loggingSessionMockFns.mockProjectDiagnosticError,
97+
}
98+
})
99+
executionPreprocessingMockFns.mockPreprocessExecution.mockResolvedValue({
100+
success: true,
101+
actorUserId: 'user-1',
102+
billingAttribution,
103+
workflowRecord: {
104+
id: 'workflow-1',
105+
workspaceId: 'workspace-1',
106+
userId: 'user-1',
107+
variables: {},
108+
},
109+
})
110+
})
111+
112+
it('completes the job when a block failure is recorded only after core throws', async () => {
113+
const blockError = buildBlockExecutionError({
114+
block: planPanels,
115+
error: new Error("ValueError: Doctrine has no sheet layout for kind ''"),
116+
})
117+
let recorded = false
118+
mockExecuteWorkflowCore.mockRejectedValue(blockError)
119+
loggingSessionMockFns.mockWaitForPostExecution.mockImplementation(async () => {
120+
recorded = true
121+
})
122+
mockWasExecutionFinalizedByCore.mockImplementation(() => recorded)
123+
124+
const result = await executeWorkflowJob(payload)
125+
126+
expect(result).toMatchObject({
127+
success: false,
128+
workflowId: 'workflow-1',
129+
executionId: 'execution-1',
130+
error: blockError.message,
131+
})
132+
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
133+
})
134+
135+
it('faults the job when core never recorded the failure', async () => {
136+
const setupError = new Error('Workflow state not found')
137+
mockExecuteWorkflowCore.mockRejectedValue(setupError)
138+
mockWasExecutionFinalizedByCore.mockReturnValue(false)
139+
140+
await expect(executeWorkflowJob(payload)).rejects.toBe(setupError)
141+
expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalled()
142+
})
143+
})

0 commit comments

Comments
 (0)