Skip to content
Closed
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
6 changes: 4 additions & 2 deletions apps/sim/app/api/jobs/[jobId]/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { checkHybridAuth } from '@/lib/auth/hybrid'
import { getJobQueue } from '@/lib/core/async-jobs'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { projectWorkflowJobOutcome } from '@/lib/workflows/executor/job-outcome'
import { createErrorResponse } from '@/app/api/workflows/utils'

const logger = createLogger('TaskStatusAPI')
Expand Down Expand Up @@ -66,15 +67,16 @@ export const GET = withRouteHandler(
return createErrorResponse('Access denied', 403)
}

const outcome = projectWorkflowJobOutcome(job)
const response: Record<string, unknown> = {
success: true,
taskId,
status: job.status,
status: outcome.status,
metadata: job.metadata,
}

if (job.output !== undefined) response.output = job.output
if (job.error !== undefined) response.error = job.error
if (outcome.error !== undefined) response.error = outcome.error

return NextResponse.json(response)
} catch (error: unknown) {
Expand Down
11 changes: 2 additions & 9 deletions apps/sim/background/async-preprocessing-correlation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,15 +26,13 @@ const {
mockExecuteWorkflowCore,
mockExecutionSnapshot,
mockWasExecutionFinalizedByCore,
mockHasExecutionResult,
mockIsWorkflowTimedOut,
mockGetScheduleTimeValues,
mockGetSubBlockValue,
} = vi.hoisted(() => ({
mockExecuteWorkflowCore: vi.fn(),
mockExecutionSnapshot: vi.fn(),
mockWasExecutionFinalizedByCore: vi.fn(),
mockHasExecutionResult: vi.fn(),
mockIsWorkflowTimedOut: vi.fn(() => false),
mockGetScheduleTimeValues: vi.fn(),
mockGetSubBlockValue: vi.fn(),
Expand Down Expand Up @@ -74,10 +72,6 @@ vi.mock('@/executor/execution/snapshot', () => ({
ExecutionSnapshot: mockExecutionSnapshot,
}))

vi.mock('@/executor/utils/errors', () => ({
hasExecutionResult: mockHasExecutionResult,
}))

import { executeScheduleJob } from './schedule-execution'
import { executeWorkflowJob } from './workflow-execution'

Expand Down Expand Up @@ -125,7 +119,6 @@ const principal = {
describe('async preprocessing correlation threading', () => {
beforeEach(() => {
mockWasExecutionFinalizedByCore.mockReturnValue(false)
mockHasExecutionResult.mockReturnValue(false)
mockIsWorkflowTimedOut.mockReturnValue(false)
resetDbChainMock()
dbChainMockFns.limit.mockResolvedValue([
Expand Down Expand Up @@ -379,7 +372,6 @@ describe('async preprocessing correlation threading', () => {
executionTimeout: {},
})
mockExecuteWorkflowCore.mockRejectedValueOnce(rawError)
mockHasExecutionResult.mockImplementation((error) => error === rawError)
mockWasExecutionFinalizedByCore.mockReturnValue(true)

await expect(
Expand All @@ -395,7 +387,8 @@ describe('async preprocessing correlation threading', () => {
})
).rejects.toBe(rawError)

expect(loggingSessionMockFns.mockWaitForPostExecution).not.toHaveBeenCalled()
// Core finalizes after throwing, so the task must settle that work before deciding.
expect(loggingSessionMockFns.mockWaitForPostExecution).toHaveBeenCalled()
expect(mockWasExecutionFinalizedByCore).toHaveBeenCalledWith(rawError, 'execution-finalized')
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
Comment on lines +391 to 393

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Tests assert mock calls

These assertions check whether mocked functions were called. The repository’s testing directive prohibits mock-call assertions; the new workflow-execution test also uses them. Please verify the resulting behavior instead. This requirement must be satisfied before merging.

Context Used: CLAUDE.md (source)

Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!

})
Expand Down
40 changes: 40 additions & 0 deletions apps/sim/background/resume-execution.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,13 @@ vi.mock('@/executor/execution/snapshot', () => ({
}))

import { executeResumeJob, type ResumeExecutionPayload } from '@/background/resume-execution'
import { buildBlockExecutionError, markWorkflowUserFailure } from '@/executor/utils/errors'
import type { SerializedBlock } from '@/serializer/types'

const planPanels = {
id: 'plan-panels',
metadata: { id: 'function', name: 'planPanels' },
} as SerializedBlock

const { mockFindCellContextByExecutionId } = tableWorkflowColumnsMockFns
const {
Expand Down Expand Up @@ -104,6 +111,39 @@ describe('executeResumeJob terminal errors', () => {
expect(rawError.message).toContain(secret)
})

it('completes the job when the resumed workflow failed with a recorded workflow user failure', async () => {
const blockError = Object.assign(
buildBlockExecutionError({
block: planPanels,
error: markWorkflowUserFailure(new Error("ValueError: kind ''")),
}),
{ executionFinalizedByCore: true }
)
mockStartResumeExecution.mockRejectedValue(blockError)

await expect(executeResumeJob(payload)).resolves.toMatchObject({
success: false,
workflowId: 'workflow-1',
executionId: 'resume-execution-1',
parentExecutionId: 'parent-execution-1',
status: 'failed',
error: blockError.message,
})
})

it('faults the job on an unmarked block failure even when core recorded it', async () => {
const blockError = Object.assign(
buildBlockExecutionError({
block: planPanels,
error: new TypeError('rows.flatMap is not a function'),
}),
{ executionFinalizedByCore: true }
)
mockStartResumeExecution.mockRejectedValue(blockError)

await expect(executeResumeJob(payload)).rejects.toBe(blockError)
})

it('starts a legacy attempt deadline before deserializing the full snapshot', async () => {
mockStartResumeExecution.mockResolvedValue({
success: true,
Expand Down
13 changes: 13 additions & 0 deletions apps/sim/background/resume-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,10 @@ import {
createResumeAttemptTimeoutController,
PauseResumeManager,
} from '@/lib/workflows/executor/human-in-the-loop-manager'
import {
buildWorkflowJobFailureResult,
classifySettledWorkflowJobFailure,
} from '@/lib/workflows/executor/job-failure'
import { RESUME_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits'
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
import type { SerializedSnapshot } from '@/executor/types'
Expand Down Expand Up @@ -230,6 +234,15 @@ export async function executeResumeJob(payload: ResumeExecutionPayload, signal?:
workflowId,
})
)
// The resumed run executes under its parent's id, and the manager settles
// its post-execution work before re-throwing.
if (classifySettledWorkflowJobFailure(error, parentExecutionId) === 'workflow_failure') {

@cubic-dev-ai cubic-dev-ai Bot Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: This handled failure returns before the cell-context terminal writer runs. Marked workflow failures in table cells will therefore leave the cell/group stuck in its pre-resume state; write the cell terminal error state before returning this workflow-failure result.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At apps/sim/background/resume-execution.ts, line 239:

<comment>This handled failure returns before the cell-context terminal writer runs. Marked workflow failures in table cells will therefore leave the cell/group stuck in its pre-resume state; write the cell terminal error state before returning this workflow-failure result.</comment>

<file context>
@@ -230,6 +234,15 @@ export async function executeResumeJob(payload: ResumeExecutionPayload, signal?:
     )
+    // The resumed run executes under its parent's id, and the manager settles
+    // its post-execution work before re-throwing.
+    if (classifySettledWorkflowJobFailure(error, parentExecutionId) === 'workflow_failure') {
+      return {
+        ...buildWorkflowJobFailureResult({ error, workflowId, executionId: resumeExecutionId }),
</file context>
Fix with cubic

return {
...buildWorkflowJobFailureResult({ error, workflowId, executionId: resumeExecutionId }),
parentExecutionId,
status: 'failed' as const,
}
}
throw error
} finally {
timeoutController?.cleanup()
Expand Down
157 changes: 157 additions & 0 deletions apps/sim/background/workflow-execution.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
/**
* @vitest-environment node
*/
import {
executionPreprocessingMock,
executionPreprocessingMockFns,
LoggingSessionMock,
loggingSessionMock,
loggingSessionMockFns,
} from '@sim/testing'
import {
executionLimitsMock,
executionLimitsMockFns,
} from '@sim/testing/mocks/execution-limits.mock'
import { beforeEach, describe, expect, it, vi } from 'vitest'

const { mockExecuteWorkflowCore, mockWasExecutionFinalizedByCore } = vi.hoisted(() => ({
mockExecuteWorkflowCore: vi.fn(),
mockWasExecutionFinalizedByCore: vi.fn(),
}))

vi.mock('@/lib/execution/preprocessing', () => executionPreprocessingMock)
vi.mock('@/lib/logs/execution/logging-session', () => loggingSessionMock)
vi.mock('@/lib/core/execution-limits', () => executionLimitsMock)
vi.mock('@/lib/workflows/executor/execution-core', () => ({
executeWorkflowCore: mockExecuteWorkflowCore,
wasExecutionFinalizedByCore: mockWasExecutionFinalizedByCore,
}))
vi.mock('@/lib/workflows/executor/pause-persistence', () => ({
handlePostExecutionPauseState: vi.fn(),
}))
vi.mock('@/lib/logs/execution/trace-spans/trace-spans', () => ({
buildTraceSpans: vi.fn(() => ({ traceSpans: [] })),
}))
vi.mock('@/executor/execution/snapshot', () => ({ ExecutionSnapshot: vi.fn() }))
vi.mock('@/lib/uploads/utils/user-file-base64.server', () => ({
cleanupExecutionBase64Cache: vi.fn(async () => {}),
}))

import * as usageReservation from '@/lib/billing/calculations/usage-reservation'
import { executeWorkflowJob, type WorkflowExecutionPayload } from '@/background/workflow-execution'
import { buildBlockExecutionError, markWorkflowUserFailure } from '@/executor/utils/errors'
import type { SerializedBlock } from '@/serializer/types'

const billingAttribution = {
actorUserId: 'user-1',
workspaceId: 'workspace-1',
organizationId: null,
billedAccountUserId: 'user-1',
billingEntity: { type: 'user' as const, id: 'user-1' },
billingPeriod: { start: '2026-09-01T00:00:00.000Z', end: '2026-10-01T00:00:00.000Z' },
payerSubscription: null,
}

const payload: WorkflowExecutionPayload = {
workflowId: 'workflow-1',
principal: {
version: 1,
principal: {
kind: 'system',
serviceId: 'internal',
workspaceId: 'workspace-1',
workflowId: 'workflow-1',
},
},
userId: 'user-1',
billingAttribution,
workspaceId: 'workspace-1',
executionId: 'execution-1',
requestId: 'request-1',
triggerType: 'api',
}

const planPanels = {
id: 'plan-panels',
metadata: { id: 'function', name: 'planPanels' },
} as SerializedBlock

describe('executeWorkflowJob fault vs workflow failure', () => {
beforeEach(() => {
vi.spyOn(usageReservation, 'refreshExecutionSlotExpiry').mockResolvedValue(true)
vi.spyOn(usageReservation, 'releaseExecutionSlot').mockResolvedValue(undefined)
executionLimitsMockFns.mockCreateTimeoutAbortController.mockImplementation(() => ({
signal: new AbortController().signal,
cleanup: vi.fn(),
abort: vi.fn(),
isTimedOut: () => false,
timeoutMs: 120_000,
}))
LoggingSessionMock.mockImplementation(function LoggingSession() {
return {
safeCompleteWithError: loggingSessionMockFns.mockSafeCompleteWithError,
waitForPostExecution: loggingSessionMockFns.mockWaitForPostExecution,
markAsFailed: loggingSessionMockFns.mockMarkAsFailed,
setExecutionDeadlineAt: loggingSessionMockFns.mockSetExecutionDeadlineAt,
projectDiagnosticError: loggingSessionMockFns.mockProjectDiagnosticError,
}
})
executionPreprocessingMockFns.mockPreprocessExecution.mockResolvedValue({
success: true,
actorUserId: 'user-1',
billingAttribution,
workflowRecord: {
id: 'workflow-1',
workspaceId: 'workspace-1',
userId: 'user-1',
variables: {},
},
})
})

it('completes the job when a workflow user failure is recorded only after core throws', async () => {
const blockError = buildBlockExecutionError({
block: planPanels,
error: markWorkflowUserFailure(
new Error("ValueError: Doctrine has no sheet layout for kind ''")
),
})
let recorded = false
mockExecuteWorkflowCore.mockRejectedValue(blockError)
loggingSessionMockFns.mockWaitForPostExecution.mockImplementation(async () => {
recorded = true
})
mockWasExecutionFinalizedByCore.mockImplementation(() => recorded)

const result = await executeWorkflowJob(payload)

expect(result).toMatchObject({
success: false,
workflowId: 'workflow-1',
executionId: 'execution-1',
error: blockError.message,
})
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
})

it('faults the job on an unmarked block failure even when core recorded it', async () => {
const internalError = buildBlockExecutionError({
block: planPanels,
error: new Error('An internal error occurred while running this block'),
})
mockExecuteWorkflowCore.mockRejectedValue(internalError)
mockWasExecutionFinalizedByCore.mockReturnValue(true)

await expect(executeWorkflowJob(payload)).rejects.toBe(internalError)
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
})

it('faults the job when core never recorded the failure', async () => {
const setupError = new Error('Workflow state not found')
mockExecuteWorkflowCore.mockRejectedValue(setupError)
mockWasExecutionFinalizedByCore.mockReturnValue(false)

await expect(executeWorkflowJob(payload)).rejects.toBe(setupError)
expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalled()
})
})
11 changes: 11 additions & 0 deletions apps/sim/background/workflow-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,10 @@ import {
executeWorkflowCore,
wasExecutionFinalizedByCore,
} from '@/lib/workflows/executor/execution-core'
import {
buildWorkflowJobFailureResult,
classifyWorkflowJobFailure,
} from '@/lib/workflows/executor/job-failure'
import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence'
import { WORKFLOW_EXECUTION_CONCURRENCY_LIMIT } from '@/background/concurrency-limits'
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
Expand Down Expand Up @@ -313,6 +317,13 @@ export async function executeWorkflowJob(

if (error instanceof ExecutionTimeoutError) throw error

const failure = await classifyWorkflowJobFailure({ error, executionId, loggingSession })
if (failure === 'workflow_failure') {
return {
...buildWorkflowJobFailureResult({ error, workflowId, executionId }),
metadata: payload.metadata,
}
}
if (wasExecutionFinalizedByCore(error, executionId)) {
throw error
}
Expand Down
1 change: 1 addition & 0 deletions apps/sim/executor/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,7 @@ export const DEFAULTS = {
export const HTTP = {
STATUS: {
OK: 200,
BAD_REQUEST: 400,
FORBIDDEN: 403,
NOT_FOUND: 404,
TOO_MANY_REQUESTS: 429,
Expand Down
19 changes: 19 additions & 0 deletions apps/sim/executor/handlers/api/api-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { beforeEach, describe, expect, it, type Mock, vi } from 'vitest'
import { BlockType } from '@/executor/constants'
import { ApiBlockHandler } from '@/executor/handlers/api/api-handler'
import type { ExecutionContext } from '@/executor/types'
import { isWorkflowUserFailure } from '@/executor/utils/errors'
import type { SerializedBlock } from '@/serializer/types'
import { executeTool } from '@/tools'
import type { ToolConfig } from '@/tools/types'
Expand Down Expand Up @@ -125,4 +126,22 @@ describe('ApiBlockHandler', () => {
)
expect(mockExecuteTool).toHaveBeenCalled()
})

it.each([
{ output: { status: 400, statusText: 'Bad Request' }, expected: true },
{ output: { status: 503, statusText: 'Service Unavailable' }, expected: false },
{ output: {}, expected: false },
])(
'marks the failure as the workflow user failure only for a 4xx from the requested URL ($output.status)',
async ({ output, expected }) => {
mockExecuteTool.mockResolvedValue({ success: false, output, error: 'Request failed' })

const thrown = await handler
.execute(mockContext, mockBlock, { url: 'https://example.com/hook', method: 'POST' })
.catch((error: unknown) => error)

expect(thrown).toBeInstanceOf(Error)
expect(isWorkflowUserFailure(thrown)).toBe(expected)
}
)
})
Loading
Loading