From 11cd349a5d825620006d3a3f4513b2a192b75e7b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 10:48:30 -0700 Subject: [PATCH] fix(execution): complete async and resume jobs on marked workflow user failures The workflow-execution and resume-execution Trigger.dev tasks re-threw every execution error, so a user's own workflow failure (their function code raising, a condition expression that does not parse) failed the run and paged #eng-errors. The existing finalized-by-core guard in workflow-execution also ran before core had recorded the failure, so it never applied. Failures stay faults by default. A job completes with success: false only when core recorded the failure and it was positively marked as the workflow's own at its source (markWorkflowUserFailure): - the Function route's 422 for a runtime or compile error in user/custom-tool code, carried across executeTool's flattening as ToolResponse.workflowUserFailure - a condition expression that threw or that the sandbox could not parse - the HTTP block receiving a 4xx from the URL the workflow called - a missing required field, now a WorkflowValidationError attributed to its block Child workflows carry the mark through their cause chain. Anything unmarked, including database and other platform errors inside blocks, still faults. /api/jobs and the queue-job status fallback project the completed failure result back to failed. Webhook and schedule execution are unchanged. --- apps/sim/app/api/jobs/[jobId]/route.ts | 6 +- .../async-preprocessing-correlation.test.ts | 11 +- apps/sim/background/resume-execution.test.ts | 40 +++++ apps/sim/background/resume-execution.ts | 13 ++ .../sim/background/workflow-execution.test.ts | 157 ++++++++++++++++++ apps/sim/background/workflow-execution.ts | 11 ++ apps/sim/executor/constants.ts | 1 + .../executor/handlers/api/api-handler.test.ts | 19 +++ apps/sim/executor/handlers/api/api-handler.ts | 11 ++ .../condition/condition-handler.test.ts | 47 ++++++ .../handlers/condition/condition-handler.ts | 8 +- .../function/function-handler.test.ts | 23 ++- .../handlers/function/function-handler.ts | 3 +- apps/sim/executor/utils/errors.ts | 29 +++- .../workflows/executor/execution-status.ts | 10 +- .../workflows/executor/job-failure.test.ts | 95 +++++++++++ .../sim/lib/workflows/executor/job-failure.ts | 73 ++++++++ .../workflows/executor/job-outcome.test.ts | 37 +++++ .../sim/lib/workflows/executor/job-outcome.ts | 30 ++++ apps/sim/serializer/index.test.ts | 43 +++++ apps/sim/serializer/index.ts | 10 +- apps/sim/tools/index.test.ts | 25 +++ apps/sim/tools/index.ts | 13 +- apps/sim/tools/types.ts | 6 + 24 files changed, 700 insertions(+), 21 deletions(-) create mode 100644 apps/sim/background/workflow-execution.test.ts create mode 100644 apps/sim/lib/workflows/executor/job-failure.test.ts create mode 100644 apps/sim/lib/workflows/executor/job-failure.ts create mode 100644 apps/sim/lib/workflows/executor/job-outcome.test.ts create mode 100644 apps/sim/lib/workflows/executor/job-outcome.ts diff --git a/apps/sim/app/api/jobs/[jobId]/route.ts b/apps/sim/app/api/jobs/[jobId]/route.ts index 4cb066efd12..bcaa5caaf0a 100644 --- a/apps/sim/app/api/jobs/[jobId]/route.ts +++ b/apps/sim/app/api/jobs/[jobId]/route.ts @@ -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') @@ -66,15 +67,16 @@ export const GET = withRouteHandler( return createErrorResponse('Access denied', 403) } + const outcome = projectWorkflowJobOutcome(job) const response: Record = { 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) { diff --git a/apps/sim/background/async-preprocessing-correlation.test.ts b/apps/sim/background/async-preprocessing-correlation.test.ts index bdc774bc027..a09e26b469a 100644 --- a/apps/sim/background/async-preprocessing-correlation.test.ts +++ b/apps/sim/background/async-preprocessing-correlation.test.ts @@ -26,7 +26,6 @@ const { mockExecuteWorkflowCore, mockExecutionSnapshot, mockWasExecutionFinalizedByCore, - mockHasExecutionResult, mockIsWorkflowTimedOut, mockGetScheduleTimeValues, mockGetSubBlockValue, @@ -34,7 +33,6 @@ const { mockExecuteWorkflowCore: vi.fn(), mockExecutionSnapshot: vi.fn(), mockWasExecutionFinalizedByCore: vi.fn(), - mockHasExecutionResult: vi.fn(), mockIsWorkflowTimedOut: vi.fn(() => false), mockGetScheduleTimeValues: vi.fn(), mockGetSubBlockValue: vi.fn(), @@ -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' @@ -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([ @@ -379,7 +372,6 @@ describe('async preprocessing correlation threading', () => { executionTimeout: {}, }) mockExecuteWorkflowCore.mockRejectedValueOnce(rawError) - mockHasExecutionResult.mockImplementation((error) => error === rawError) mockWasExecutionFinalizedByCore.mockReturnValue(true) await expect( @@ -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() }) diff --git a/apps/sim/background/resume-execution.test.ts b/apps/sim/background/resume-execution.test.ts index fb2e2cc79b0..20bdc9af3bf 100644 --- a/apps/sim/background/resume-execution.test.ts +++ b/apps/sim/background/resume-execution.test.ts @@ -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 { @@ -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, diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index 0d763e391ca..99147ee8617 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -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' @@ -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') { + return { + ...buildWorkflowJobFailureResult({ error, workflowId, executionId: resumeExecutionId }), + parentExecutionId, + status: 'failed' as const, + } + } throw error } finally { timeoutController?.cleanup() diff --git a/apps/sim/background/workflow-execution.test.ts b/apps/sim/background/workflow-execution.test.ts new file mode 100644 index 00000000000..b7c0e8699f2 --- /dev/null +++ b/apps/sim/background/workflow-execution.test.ts @@ -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() + }) +}) diff --git a/apps/sim/background/workflow-execution.ts b/apps/sim/background/workflow-execution.ts index b04e5a0610c..c1b8692e7fc 100644 --- a/apps/sim/background/workflow-execution.ts +++ b/apps/sim/background/workflow-execution.ts @@ -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' @@ -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 } diff --git a/apps/sim/executor/constants.ts b/apps/sim/executor/constants.ts index e0b4b08c283..2bde2267dcf 100644 --- a/apps/sim/executor/constants.ts +++ b/apps/sim/executor/constants.ts @@ -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, diff --git a/apps/sim/executor/handlers/api/api-handler.test.ts b/apps/sim/executor/handlers/api/api-handler.test.ts index f7096fb5e61..b6301a66f54 100644 --- a/apps/sim/executor/handlers/api/api-handler.test.ts +++ b/apps/sim/executor/handlers/api/api-handler.test.ts @@ -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' @@ -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) + } + ) }) diff --git a/apps/sim/executor/handlers/api/api-handler.ts b/apps/sim/executor/handlers/api/api-handler.ts index 1db45b81246..1d10c32d587 100644 --- a/apps/sim/executor/handlers/api/api-handler.ts +++ b/apps/sim/executor/handlers/api/api-handler.ts @@ -1,6 +1,7 @@ import { createLogger } from '@sim/logger' import { BlockType, HTTP } from '@/executor/constants' import type { BlockHandler, ExecutionContext } from '@/executor/types' +import { markWorkflowUserFailure } from '@/executor/utils/errors' import type { SerializedBlock } from '@/serializer/types' import { executeTool } from '@/tools' import { getTool } from '@/tools/utils' @@ -127,6 +128,16 @@ export class ApiBlockHandler implements BlockHandler { timestamp: new Date().toISOString(), }) + // The URL the workflow called rejected the request it sent. + const status = result.output?.status + if ( + typeof status === 'number' && + status >= HTTP.STATUS.BAD_REQUEST && + status < HTTP.STATUS.SERVER_ERROR + ) { + markWorkflowUserFailure(error) + } + throw error } diff --git a/apps/sim/executor/handlers/condition/condition-handler.test.ts b/apps/sim/executor/handlers/condition/condition-handler.test.ts index 4b347469487..0023490d14e 100644 --- a/apps/sim/executor/handlers/condition/condition-handler.test.ts +++ b/apps/sim/executor/handlers/condition/condition-handler.test.ts @@ -25,6 +25,7 @@ vi.mock('@/executor/utils/block-data', () => ({ })) import { collectBlockData } from '@/executor/utils/block-data' +import { isWorkflowUserFailure } from '@/executor/utils/errors' import { executeTool } from '@/tools' const mockExecuteTool = executeTool as ReturnType @@ -367,6 +368,52 @@ describe('ConditionBlockHandler', () => { expect(JSON.stringify(mockConditionLogger.error.mock.calls)).not.toContain(secret) }) + describe('workflow user failures', () => { + const conditions = [ + { id: 'cond1', title: 'if', value: ' === true' }, + { id: 'else1', title: 'else', value: '' }, + ] + + it('marks an expression that threw as the workflow user failure', async () => { + mockExecuteTool.mockResolvedValueOnce(threwAt(0, 'Cannot read properties of undefined')) + + const thrown = await handler + .execute(mockContext, mockBlock, { conditions: JSON.stringify(conditions) }) + .catch((error: unknown) => error) + + expect(isWorkflowUserFailure(thrown)).toBe(true) + }) + + it('marks an expression the sandbox could not parse as the workflow user failure', async () => { + const syntaxError = { + success: false, + error: + "Syntax Error: Line 3: ` === true` - Unexpected token '<'", + workflowUserFailure: true as const, + } + mockExecuteTool.mockResolvedValueOnce(syntaxError) + mockExecuteTool.mockResolvedValueOnce(syntaxError) + + const thrown = await handler + .execute(mockContext, mockBlock, { conditions: JSON.stringify(conditions) }) + .catch((error: unknown) => error) + + expect(thrown).toBeInstanceOf(Error) + expect(isWorkflowUserFailure(thrown)).toBe(true) + }) + + it('does not mark an evaluation that timed out', async () => { + mockExecuteTool.mockResolvedValue({ success: false, error: 'Request timed out after 5000ms' }) + + const thrown = await handler + .execute(mockContext, mockBlock, { conditions: JSON.stringify(conditions) }) + .catch((error: unknown) => error) + + expect(thrown).toBeInstanceOf(Error) + expect(isWorkflowUserFailure(thrown)).toBe(false) + }) + }) + it('preserves routing metadata when the target block is disabled', async () => { mockExecuteTool.mockResolvedValueOnce(matchedAt(0)) diff --git a/apps/sim/executor/handlers/condition/condition-handler.ts b/apps/sim/executor/handlers/condition/condition-handler.ts index 5086cdd81a6..49cdd415166 100644 --- a/apps/sim/executor/handlers/condition/condition-handler.ts +++ b/apps/sim/executor/handlers/condition/condition-handler.ts @@ -10,6 +10,7 @@ import type { BlockOutput } from '@/blocks/types' import { BlockType, DEFAULTS, EDGE } from '@/executor/constants' import type { BlockHandler, ExecutionContext } from '@/executor/types' import { collectBlockData } from '@/executor/utils/block-data' +import { markWorkflowUserFailure } from '@/executor/utils/errors' import { createEnvVarPattern } from '@/executor/utils/reference-validation' import { buildBranchNodeId, @@ -295,7 +296,8 @@ async function evaluateSingleCondition( if (result.retryable === false) { throw new NonRetryableExecutionError(result.error ?? 'Condition evaluation is indeterminate') } - throw new Error(result.error ?? 'Condition evaluation failed') + const error = new Error(result.error ?? 'Condition evaluation failed') + throw result.workflowUserFailure ? markWorkflowUserFailure(error) : error } return Boolean(result.output?.result) @@ -505,7 +507,9 @@ export class ConditionBlockHandler implements BlockHandler { return null case 'expression-threw': logger.error('Failed to evaluate condition', { conditionCount: conditions.length }) - throw conditionError(conditions[evaluation.index], evaluation.message) + throw markWorkflowUserFailure( + conditionError(conditions[evaluation.index], evaluation.message) + ) case 'no-verdict': if (!evaluation.retryable) { throw new NonRetryableExecutionError( diff --git a/apps/sim/executor/handlers/function/function-handler.test.ts b/apps/sim/executor/handlers/function/function-handler.test.ts index 4dfdf6b849a..999ed69cc5c 100644 --- a/apps/sim/executor/handlers/function/function-handler.test.ts +++ b/apps/sim/executor/handlers/function/function-handler.test.ts @@ -5,7 +5,7 @@ import { NonRetryableExecutionError } from '@/lib/execution/non-retryable-error' import { BlockType } from '@/executor/constants' import { FunctionBlockHandler } from '@/executor/handlers/function/function-handler' import type { ExecutionContext } from '@/executor/types' -import { readTrustedExecutionCost } from '@/executor/utils/errors' +import { isWorkflowUserFailure, readTrustedExecutionCost } from '@/executor/utils/errors' import type { SerializedBlock } from '@/serializer/types' import { executeTool } from '@/tools' @@ -109,4 +109,25 @@ describe('FunctionBlockHandler', () => { expect(readTrustedExecutionCost(thrown)).toEqual(cost) } ) + + it.each([ + { workflowUserFailure: true as const, expected: true }, + { workflowUserFailure: undefined, expected: false }, + ])( + 'marks the failure as the workflow user failure only when the tool reports one ($expected)', + async ({ workflowUserFailure, expected }) => { + mockExecuteTool.mockResolvedValueOnce({ + success: false, + output: { result: null, stdout: '' }, + error: "ValueError: kind ''", + ...(workflowUserFailure ? { workflowUserFailure } : {}), + }) + + const thrown = await handler + .execute(mockContext, mockBlock, { code: 'raise ValueError()' }) + .catch((error: unknown) => error) + + expect(isWorkflowUserFailure(thrown)).toBe(expected) + } + ) }) diff --git a/apps/sim/executor/handlers/function/function-handler.ts b/apps/sim/executor/handlers/function/function-handler.ts index 2d66a019e1f..366b714b47f 100644 --- a/apps/sim/executor/handlers/function/function-handler.ts +++ b/apps/sim/executor/handlers/function/function-handler.ts @@ -12,7 +12,7 @@ import { normalizeSecretMountPolicy } from '@/lib/mothership/secret-mount-policy import { BlockType } from '@/executor/constants' import type { BlockHandler, ExecutionContext } from '@/executor/types' import { collectBlockData } from '@/executor/utils/block-data' -import { attachTrustedExecutionCost } from '@/executor/utils/errors' +import { attachTrustedExecutionCost, markWorkflowUserFailure } from '@/executor/utils/errors' import { FUNCTION_BLOCK_CONTEXT_VARS_KEY, FUNCTION_BLOCK_DISPLAY_CODE_KEY, @@ -116,6 +116,7 @@ export class FunctionBlockHandler implements BlockHandler { result.retryable === false ? new NonRetryableExecutionError(result.error || 'Function execution is indeterminate') : new Error(result.error || 'Function execution failed') + if (result.workflowUserFailure) markWorkflowUserFailure(error) attachTrustedExecutionCost(error, result.output?.cost) throw error } diff --git a/apps/sim/executor/utils/errors.ts b/apps/sim/executor/utils/errors.ts index e3f0a9ee9b8..fd29c9c3fe6 100644 --- a/apps/sim/executor/utils/errors.ts +++ b/apps/sim/executor/utils/errors.ts @@ -1,4 +1,4 @@ -import { getErrorMessage } from '@sim/utils/errors' +import { findCause, getErrorMessage } from '@sim/utils/errors' import { HttpError } from '@/lib/core/utils/http-error' import type { ExecutionContext, ExecutionResult } from '@/executor/types' import type { SerializedBlock } from '@/serializer/types' @@ -141,6 +141,33 @@ function isRecordedThrown(value: unknown): value is object { return (typeof value === 'object' || typeof value === 'function') && value !== null } +/** + * Marks `error` as caused by the workflow itself: its user code, an expression, + * its configuration, or a request it sent that the target rejected. Mark only a + * failure positively known to be the workflow's; an unmarked failure is treated + * as Sim failing to run the workflow. + */ +export function markWorkflowUserFailure(error: T): T { + return Object.assign(error, { workflowUserFailure: true as const }) +} + +/** + * Whether `error`, or an error in its `.cause` chain, was marked with + * {@link markWorkflowUserFailure}. The chain is walked because block and child + * workflow boundaries wrap the failure they report. + */ +export function isWorkflowUserFailure(error: unknown): boolean { + return ( + findCause( + error, + (value): value is Error => + value instanceof Error && + 'workflowUserFailure' in value && + value.workflowUserFailure === true + ) !== undefined + ) +} + export interface BlockExecutionErrorDetails { block: SerializedBlock error: Error | string diff --git a/apps/sim/lib/workflows/executor/execution-status.ts b/apps/sim/lib/workflows/executor/execution-status.ts index edac4c8649f..7776e3dcffb 100644 --- a/apps/sim/lib/workflows/executor/execution-status.ts +++ b/apps/sim/lib/workflows/executor/execution-status.ts @@ -14,6 +14,7 @@ import { RESUME_EXECUTION_JOB_ID_PREFIX, WORKFLOW_EXECUTION_JOB_ID_PREFIX, } from '@/lib/workflows/executor/execution-job-ids' +import { projectWorkflowJobOutcome } from '@/lib/workflows/executor/job-outcome' import { getAutomaticResumeWaitingMetadata } from '@/lib/workflows/executor/paused-execution-metadata' import type { PausePoint } from '@/executor/types' @@ -97,8 +98,13 @@ function projectQueueJob( job: Job, input: Pick ): WorkflowExecutionStatusResponse { + const outcome = projectWorkflowJobOutcome(job) const status: WorkflowExecutionStatusResponse['status'] = - job.status === 'pending' ? 'queued' : job.status === 'processing' ? 'running' : job.status + outcome.status === 'pending' + ? 'queued' + : outcome.status === 'processing' + ? 'running' + : outcome.status const startedAt = job.startedAt ?? job.createdAt const endedAt = job.completedAt ?? null @@ -113,7 +119,7 @@ function projectQueueJob( totalDurationMs: endedAt ? endedAt.getTime() - startedAt.getTime() : null, paused: null, cost: null, - error: status === 'failed' ? (job.error ?? 'Execution failed') : null, + error: status === 'failed' ? (outcome.error ?? 'Execution failed') : null, finalOutput: input.includeOutput && status === 'completed' ? extractJobFinalOutput(job.output) : null, blockOutputs: null, diff --git a/apps/sim/lib/workflows/executor/job-failure.test.ts b/apps/sim/lib/workflows/executor/job-failure.test.ts new file mode 100644 index 00000000000..e0eec4fa31b --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-failure.test.ts @@ -0,0 +1,95 @@ +/** + * @vitest-environment node + */ +import { describe, expect, it } from 'vitest' +import { classifyWorkflowJobFailure } from '@/lib/workflows/executor/job-failure' +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 writeLedger = { + id: 'write-ledger', + metadata: { id: 'table', name: 'writeLedger' }, +} as SerializedBlock + +/** + * Core throws first and finalizes the execution log from a post-execution + * promise; this session resolves that promise the way core does, by flagging + * the thrown error once the log write lands. + */ +function sessionFinalizing(error: Error, finalized: boolean) { + return { + waitForPostExecution: async () => { + if (finalized) Object.assign(error, { executionFinalizedByCore: true }) + }, + } +} + +describe('classifyWorkflowJobFailure', () => { + it('treats a marked failure that core records only after throwing as the workflow outcome', async () => { + const error = buildBlockExecutionError({ + block: planPanels, + error: markWorkflowUserFailure(new Error("ValueError: kind ''")), + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-1', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('workflow_failure') + }) + + it('treats a marked failure inside a child workflow as the workflow outcome', async () => { + const childFailure = buildBlockExecutionError({ + block: planPanels, + error: markWorkflowUserFailure(new Error("KeyError: 'id'")), + }) + const error = new Error(`callSheets: "sheets" failed: ${childFailure.message}`, { + cause: childFailure, + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-2', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('workflow_failure') + }) + + it('faults a marked failure core could not record', async () => { + const error = buildBlockExecutionError({ + block: planPanels, + error: markWorkflowUserFailure(new Error("ValueError: kind ''")), + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-3', + loggingSession: sessionFinalizing(error, false), + }) + ).resolves.toBe('job_fault') + }) + + it('faults an unmarked failure inside a block even when core recorded it', async () => { + const error = buildBlockExecutionError({ + block: writeLedger, + error: new Error('An internal error occurred while running this block'), + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-4', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('job_fault') + }) +}) diff --git a/apps/sim/lib/workflows/executor/job-failure.ts b/apps/sim/lib/workflows/executor/job-failure.ts new file mode 100644 index 00000000000..e16076f14e0 --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-failure.ts @@ -0,0 +1,73 @@ +import { getErrorMessage } from '@sim/utils/errors' +import type { LoggingSession } from '@/lib/logs/execution/logging-session' +import { wasExecutionFinalizedByCore } from '@/lib/workflows/executor/execution-core' +import { hasExecutionResult, isWorkflowUserFailure } from '@/executor/utils/errors' + +/** + * How a background job that ran a workflow ends when the execution threw. + * + * - `workflow_failure`: the workflow itself failed (see `markWorkflowUserFailure`) + * and core recorded it in the execution log. That is the run's outcome, so the + * job completes and reports it instead of paging engineering. + * - `job_fault`: everything else, including any failure not positively known to + * be the workflow's. The job re-throws so the queue fails the run and alerts. + */ +export type WorkflowJobFailure = 'workflow_failure' | 'job_fault' + +/** + * What a workflow job returns when its workflow failed. The job itself + * completed, so readers of job status project it back to a failed run (see + * `projectWorkflowJobOutcome`). + */ +export interface WorkflowJobFailureResult { + success: false + workflowId: string + executionId: string + output: unknown + error: string + executedAt: string +} + +/** Builds the {@link WorkflowJobFailureResult} for a workflow that failed. */ +export function buildWorkflowJobFailureResult(params: { + error: unknown + workflowId: string + executionId: string +}): WorkflowJobFailureResult { + return { + success: false, + workflowId: params.workflowId, + executionId: params.executionId, + output: hasExecutionResult(params.error) ? params.error.executionResult.output : {}, + error: getErrorMessage(params.error, 'Execution failed'), + executedAt: new Date().toISOString(), + } +} + +/** + * Decides whether an execution failure is the workflow's outcome or a fault in + * the job running it. Only call once the run's post-execution work has settled + * (see {@link classifyWorkflowJobFailure}). + */ +export function classifySettledWorkflowJobFailure( + error: unknown, + executionId: string +): WorkflowJobFailure { + return wasExecutionFinalizedByCore(error, executionId) && isWorkflowUserFailure(error) + ? 'workflow_failure' + : 'job_fault' +} + +/** + * {@link classifySettledWorkflowJobFailure} for a caller holding the run's + * logging session. Core throws before its post-execution work records the + * failure, so the finalized signal is only reliable once that work settles. + */ +export async function classifyWorkflowJobFailure(params: { + error: unknown + executionId: string + loggingSession: Pick +}): Promise { + await params.loggingSession.waitForPostExecution() + return classifySettledWorkflowJobFailure(params.error, params.executionId) +} diff --git a/apps/sim/lib/workflows/executor/job-outcome.test.ts b/apps/sim/lib/workflows/executor/job-outcome.test.ts new file mode 100644 index 00000000000..ea135d9a256 --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-outcome.test.ts @@ -0,0 +1,37 @@ +/** + * @vitest-environment node + */ +import { describe, expect, it } from 'vitest' +import { projectWorkflowJobOutcome } from '@/lib/workflows/executor/job-outcome' + +describe('projectWorkflowJobOutcome', () => { + it('reports a workflow job that completed with a failed workflow as failed', () => { + expect( + projectWorkflowJobOutcome({ + type: 'workflow-execution', + status: 'completed', + output: { success: false, error: 'planPanels: ValueError' }, + }) + ).toEqual({ status: 'failed', error: 'planPanels: ValueError' }) + }) + + it('leaves a resume that reported its cancellation without an error as the queue reported it', () => { + expect( + projectWorkflowJobOutcome({ + type: 'resume-execution', + status: 'completed', + output: { success: false, status: 'cancelled' }, + }) + ).toEqual({ status: 'completed' }) + }) + + it('leaves a job type that never returns a failure result as the queue reported it', () => { + expect( + projectWorkflowJobOutcome({ + type: 'webhook-execution', + status: 'completed', + output: { success: false, error: 'Gmail 2 is missing required fields: Label' }, + }) + ).toEqual({ status: 'completed' }) + }) +}) diff --git a/apps/sim/lib/workflows/executor/job-outcome.ts b/apps/sim/lib/workflows/executor/job-outcome.ts new file mode 100644 index 00000000000..a729e9a843f --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-outcome.ts @@ -0,0 +1,30 @@ +import { toStringOrNull } from '@sim/utils/coerce' +import { toRecordOrNull } from '@sim/utils/object' +import { JOB_STATUS, type Job, type JobStatus, type JobType } from '@/lib/core/async-jobs/types' + +/** Job types that complete with a `WorkflowJobFailureResult` when their workflow fails. */ +const WORKFLOW_FAILURE_RESULT_JOB_TYPES: ReadonlySet = new Set([ + 'workflow-execution', + 'resume-execution', +]) + +/** + * The status a workflow job stands for. The queue only knows whether the job + * faulted; a job that completed with a `WorkflowJobFailureResult` ran a + * workflow that failed, and callers polling the job must see that failed run. + */ +export function projectWorkflowJobOutcome(job: Pick): { + status: JobStatus + error?: string +} { + const reported = + job.error === undefined ? { status: job.status } : { status: job.status, error: job.error } + if (job.status !== JOB_STATUS.COMPLETED || !WORKFLOW_FAILURE_RESULT_JOB_TYPES.has(job.type)) { + return reported + } + + const output = toRecordOrNull(job.output) + const error = toStringOrNull(output?.error) + if (output?.success !== false || error === null) return reported + return { status: JOB_STATUS.FAILED, error } +} diff --git a/apps/sim/serializer/index.test.ts b/apps/sim/serializer/index.test.ts index e85f9325295..132b3391a67 100644 --- a/apps/sim/serializer/index.test.ts +++ b/apps/sim/serializer/index.test.ts @@ -24,7 +24,9 @@ import { } from '@sim/testing/mocks' import { describe, expect, it, vi } from 'vitest' import { DAGBuilder } from '@/executor/dag/builder' +import { isWorkflowUserFailure } from '@/executor/utils/errors' import { Serializer } from '@/serializer/index' +import type { BlockState } from '@/stores/workflows/workflow/types' import { getToolMetadata, getToolParams } from '@/tools/metadata' vi.mocked(getToolMetadata).mockImplementation(toolsMetadataMock.getToolMetadata) @@ -452,6 +454,47 @@ describe('Serializer', () => { } ) + it.concurrent( + 'attributes a missing required field to its block as a workflow user failure', + () => { + const serializer = new Serializer() + const waitBlockMissingRequired: BlockState = { + id: 'wait-block', + type: 'wait', + name: 'Wait Block', + position: { x: 0, y: 0 }, + subBlocks: { + timeValue: { id: 'timeValue', type: 'short-input', value: '' }, + timeUnit: { id: 'timeUnit', type: 'dropdown', value: 'seconds' }, + }, + outputs: {}, + enabled: true, + } + + const thrown = (() => { + try { + serializer.serializeWorkflow( + { 'wait-block': waitBlockMissingRequired }, + [], + {}, + undefined, + true + ) + } catch (error) { + return error + } + })() + + expect(thrown).toMatchObject({ + name: 'WorkflowValidationError', + blockId: 'wait-block', + blockType: 'wait', + blockName: 'Wait Block', + }) + expect(isWorkflowUserFailure(thrown)).toBe(true) + } + ) + it.concurrent('should handle empty string values as missing', () => { const serializer = new Serializer() diff --git a/apps/sim/serializer/index.ts b/apps/sim/serializer/index.ts index c90ba4758ac..1ef826c73b0 100644 --- a/apps/sim/serializer/index.ts +++ b/apps/sim/serializer/index.ts @@ -19,6 +19,7 @@ import { import { getBlock } from '@/blocks' import { isCustomBlockType, RESERVED_PARAMS } from '@/blocks/custom/build-config' import type { SubBlockConfig } from '@/blocks/types' +import { markWorkflowUserFailure } from '@/executor/utils/errors' import type { SerializedBlock, SerializedWorkflow } from '@/serializer/types' import type { BlockState, Loop, Parallel } from '@/stores/workflows/workflow/types' import { generateLoopBlocks, generateParallelBlocks } from '@/stores/workflows/workflow/utils' @@ -272,8 +273,13 @@ export class Serializer { const { missingRequiredFields } = collectBlockFieldIssues(block, blockConfig, params) if (missingRequiredFields.length > 0) { const blockName = block.name || blockConfig.name || 'Block' - throw new Error( - `${blockName} is missing required fields: ${missingRequiredFields.join(', ')}` + throw markWorkflowUserFailure( + new WorkflowValidationError( + `${blockName} is missing required fields: ${missingRequiredFields.join(', ')}`, + block.id, + block.type, + blockName + ) ) } } diff --git a/apps/sim/tools/index.test.ts b/apps/sim/tools/index.test.ts index b92a94d5d18..1636647693e 100644 --- a/apps/sim/tools/index.test.ts +++ b/apps/sim/tools/index.test.ts @@ -1852,6 +1852,31 @@ describe('executeTool Function', () => { expect(result.output).not.toHaveProperty('cost') }) + it.each([ + { status: 422, expected: true }, + { status: 500, expected: undefined }, + { status: 503, expected: undefined }, + ])( + 'marks a Function failure as the workflow user failure only for the user-code status ($status)', + async ({ status, expected }) => { + mockExecuteFunction.mockResolvedValueOnce( + Response.json( + { success: false, error: 'ValueError: boom', output: { result: null, stdout: '' } }, + { status } + ) + ) + + const result = await executeTool( + 'function_execute', + { code: 'raise ValueError("boom")' }, + { executionContext: createToolExecutionContext({ userId: 'user-1' }) } + ) + + expect(result.success).toBe(false) + expect(result.workflowUserFailure).toBe(expected) + } + ) + it('does not log a secret-bearing non-OK response stream error', async () => { const secret = 'function-body-stream-secret-value' const streamError = `${secret} __var_API_KEY __sim_code_0_binding_0` diff --git a/apps/sim/tools/index.ts b/apps/sim/tools/index.ts index b5710fe348b..d47d6b5af89 100644 --- a/apps/sim/tools/index.ts +++ b/apps/sim/tools/index.ts @@ -91,6 +91,7 @@ import { assertPermissionsAllowed } from '@/ee/access-control/utils/permission-c import { isCustomTool, isMcpTool } from '@/executor/constants' import { resolveSkillContent } from '@/executor/handlers/agent/skills-resolver' import type { ExecutionContext, UserFile } from '@/executor/types' +import { isWorkflowUserFailure, markWorkflowUserFailure } from '@/executor/utils/errors' import { resolveEnvVarReferences } from '@/executor/utils/reference-validation' import { projectResolvedSecretDiagnosticContent } from '@/executor/utils/resolved-secret-content-projection' import { @@ -2424,6 +2425,7 @@ async function executeToolImplementation( // thrown error into a result object; an upstream provider's status stays // on `output` where it cannot be mistaken for ours. ...(error instanceof HttpError ? { statusCode: error.statusCode } : {}), + ...(isWorkflowUserFailure(error) ? { workflowUserFailure: true as const } : {}), timing: { startTime: startTimeISO, endTime: endTimeISO, @@ -2609,6 +2611,12 @@ function isToolResponse(value: unknown): value is ToolResponse { return isRecordLike(value) && typeof value.success === 'boolean' && isRecordLike(value.output) } +/** + * Status the Function route answers a failure in the user's own code with (a + * runtime or compile error). The route keeps 5xx for its own faults. + */ +const FUNCTION_USER_CODE_FAILURE_STATUS = 422 + async function executeDeclaredInternalOperation({ toolId, tool, @@ -2790,10 +2798,13 @@ async function executeDeclaredInternalOperation({ } catch { errorData = errorText } - throw createTransformedErrorFromErrorInfo( + const error = createTransformedErrorFromErrorInfo( { status: response.status, statusText: response.statusText, data: errorData }, tool.errorExtractor ) + throw isFunctionOperation && response.status === FUNCTION_USER_CODE_FAILURE_STATUS + ? markWorkflowUserFailure(error) + : error } if (tool.transformResponse) return tool.transformResponse(response, params, { signal }) diff --git a/apps/sim/tools/types.ts b/apps/sim/tools/types.ts index c26205d140f..b2cac8d2974 100644 --- a/apps/sim/tools/types.ts +++ b/apps/sim/tools/types.ts @@ -113,6 +113,12 @@ export interface ToolResponse { * status — a provider's 404 must never become the workflow API's status. */ statusCode?: number + /** + * True when the tool failed because of the workflow itself (see + * `markWorkflowUserFailure`), carried across the same flattening as + * `statusCode` so the block can report the failure as the workflow's. + */ + workflowUserFailure?: true resources?: MothershipResource[] // Resources to auto-open/show in UI largeValueKeys?: string[] fileKeys?: string[]