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..f3e5aca2e4e 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,8 @@ vi.mock('@/executor/execution/snapshot', () => ({ ExecutionSnapshot: mockExecutionSnapshot, })) -vi.mock('@/executor/utils/errors', () => ({ - hasExecutionResult: mockHasExecutionResult, -})) - +import { buildBlockExecutionError } from '@/executor/utils/errors' +import type { SerializedBlock } from '@/serializer/types' import { executeScheduleJob } from './schedule-execution' import { executeWorkflowJob } from './workflow-execution' @@ -125,7 +121,6 @@ const principal = { describe('async preprocessing correlation threading', () => { beforeEach(() => { mockWasExecutionFinalizedByCore.mockReturnValue(false) - mockHasExecutionResult.mockReturnValue(false) mockIsWorkflowTimedOut.mockReturnValue(false) resetDbChainMock() dbChainMockFns.limit.mockResolvedValue([ @@ -379,7 +374,6 @@ describe('async preprocessing correlation threading', () => { executionTimeout: {}, }) mockExecuteWorkflowCore.mockRejectedValueOnce(rawError) - mockHasExecutionResult.mockImplementation((error) => error === rawError) mockWasExecutionFinalizedByCore.mockReturnValue(true) await expect( @@ -395,7 +389,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() }) @@ -625,4 +620,72 @@ describe('async preprocessing correlation threading', () => { }) ) }) + + describe('scheduled run failures', () => { + const schedulePayload = { + scheduleId: 'schedule-1', + workflowId: 'workflow-1', + workspaceId: 'workspace-1', + billingAttribution, + now: '2025-01-01T00:00:00.000Z', + scheduledFor: '2025-01-01T00:00:00.000Z', + } + + beforeEach(() => { + mockPreprocessExecution.mockResolvedValueOnce({ + success: true, + actorUserId: 'actor-1', + workflowRecord: { + id: 'workflow-1', + userId: 'owner-1', + workspaceId: 'workspace-1', + variables: {}, + }, + billingAttribution, + executionTimeout: {}, + }) + }) + + it('faults the job on a failure core never recorded, after recording the schedule failure', async () => { + const engineError = new Error('Workflow state not found') + mockExecuteWorkflowCore.mockRejectedValueOnce(engineError) + + await expect( + executeScheduleJob({ + ...schedulePayload, + executionId: 'execution-schedule-fault', + requestId: 'request-schedule-fault', + }) + ).rejects.toBe(engineError) + + expect(dbChainMockFns.set).toHaveBeenCalledWith( + expect.objectContaining({ lastQueuedAt: null, lastFailedAt: expect.any(Date) }) + ) + }) + + it('completes the job when the workflow failed in a block core recorded', async () => { + mockExecuteWorkflowCore.mockRejectedValueOnce( + buildBlockExecutionError({ + block: { + id: 'plan-panels', + metadata: { id: 'function', name: 'planPanels' }, + } as SerializedBlock, + error: new Error("ValueError: kind ''"), + }) + ) + mockWasExecutionFinalizedByCore.mockReturnValue(true) + + await expect( + executeScheduleJob({ + ...schedulePayload, + executionId: 'execution-schedule-failure', + requestId: 'request-schedule-failure', + }) + ).resolves.toBeUndefined() + + expect(dbChainMockFns.set).toHaveBeenCalledWith( + expect.objectContaining({ lastQueuedAt: null, lastFailedAt: expect.any(Date) }) + ) + }) + }) }) diff --git a/apps/sim/background/resume-execution.test.ts b/apps/sim/background/resume-execution.test.ts index fb2e2cc79b0..6ec846c3fdb 100644 --- a/apps/sim/background/resume-execution.test.ts +++ b/apps/sim/background/resume-execution.test.ts @@ -29,6 +29,8 @@ vi.mock('@/executor/execution/snapshot', () => ({ })) import { executeResumeJob, type ResumeExecutionPayload } from '@/background/resume-execution' +import { buildBlockExecutionError } from '@/executor/utils/errors' +import type { SerializedBlock } from '@/serializer/types' const { mockFindCellContextByExecutionId } = tableWorkflowColumnsMockFns const { @@ -104,6 +106,29 @@ describe('executeResumeJob terminal errors', () => { expect(rawError.message).toContain(secret) }) + it('completes the job when the resumed workflow failed in a block core recorded', async () => { + const blockError = Object.assign( + buildBlockExecutionError({ + block: { + id: 'plan-panels', + metadata: { id: 'function', name: 'planPanels' }, + } as SerializedBlock, + error: 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('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/schedule-execution.ts b/apps/sim/background/schedule-execution.ts index 0a34f19d5fd..d7176e0aa18 100644 --- a/apps/sim/background/schedule-execution.ts +++ b/apps/sim/background/schedule-execution.ts @@ -38,6 +38,7 @@ import { executeWorkflowCore, wasExecutionFinalizedByCore, } from '@/lib/workflows/executor/execution-core' +import { classifyWorkflowJobFailure } from '@/lib/workflows/executor/job-failure' import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence' import { loadDeployedWorkflowState } from '@/lib/workflows/persistence/utils' import { notifyScheduleAutoDisabled } from '@/lib/workflows/schedules/disable-notifications' @@ -855,6 +856,13 @@ export async function executeScheduleJob( disableReason, }) + /** + * A platform fault, re-thrown only after the schedule's own bookkeeping + * (failure count, next run, claim) has run, so faulting the job to alert + * on it never leaves the schedule claimed or its cadence stalled. + */ + let jobFault: { error: unknown } | undefined + try { const [scheduleRecord] = await db .select({ @@ -1203,6 +1211,9 @@ export async function executeScheduleJob( `Error updating schedule ${payload.scheduleId} after execution error`, 'consecutive_failures' ) + + const failure = await classifyWorkflowJobFailure({ error, executionId, loggingSession }) + if (failure === 'job_fault') jobFault = { error } } } catch (error: unknown) { try { @@ -1211,6 +1222,7 @@ export async function executeScheduleJob( return } + jobFault = { error } logger.error(`[${requestId}] Error processing schedule ${payload.scheduleId}`, error, { cause: describeError(error), }) @@ -1230,6 +1242,8 @@ export async function executeScheduleJob( trace.getActiveSpan()?.recordException(toError(recoveryError)) } } + + if (jobFault) throw jobFault.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..1d364ae44a9 --- /dev/null +++ b/apps/sim/background/workflow-execution.test.ts @@ -0,0 +1,143 @@ +/** + * @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 } 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 block failure is recorded only after core throws', async () => { + const blockError = buildBlockExecutionError({ + block: planPanels, + error: 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 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/handlers/generic/generic-handler.ts b/apps/sim/executor/handlers/generic/generic-handler.ts index f57f8ad825f..e5c5020df8c 100644 --- a/apps/sim/executor/handlers/generic/generic-handler.ts +++ b/apps/sim/executor/handlers/generic/generic-handler.ts @@ -358,6 +358,7 @@ export class GenericBlockHandler implements BlockHandler { // status (hosted-key 429/503) would be lost here. Carry it onto the // error so `getExecutionErrorStatus` can still reach the API caller. ...(typeof result.statusCode === 'number' ? { statusCode: result.statusCode } : {}), + ...(result.isSystemError ? { isSystemError: true } : {}), }) throw error 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..bfb165a4700 --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-failure.test.ts @@ -0,0 +1,115 @@ +/** + * @vitest-environment node + */ +import { describe, expect, it } from 'vitest' +import { classifyWorkflowJobFailure } from '@/lib/workflows/executor/job-failure' +import { buildBlockExecutionError } from '@/executor/utils/errors' +import { WorkflowValidationError } from '@/serializer' +import type { SerializedBlock } from '@/serializer/types' + +const block = { + id: 'block-1', + metadata: { id: 'function', name: 'planPanels' }, +} 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 block failure that core finalizes after throwing as the workflow outcome', async () => { + const error = buildBlockExecutionError({ block, error: new Error("ValueError: kind ''") }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-1', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('workflow_failure') + }) + + it('treats a pre-execution validation failure on a block as the workflow outcome', async () => { + const error = new WorkflowValidationError( + 'Gmail 2 is missing required fields: Label', + 'block-2', + 'gmail', + 'Gmail 2' + ) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-2', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('workflow_failure') + }) + + it('faults a block failure core could not record', async () => { + const error = buildBlockExecutionError({ block, error: new Error('boom') }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-3', + loggingSession: sessionFinalizing(error, false), + }) + ).resolves.toBe('job_fault') + }) + + it('faults an engine failure that no block owns even when core recorded it', async () => { + const error = new Error('Cannot read properties of undefined (reading "edges")') + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-4', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('job_fault') + }) +}) + +describe('classifyWorkflowJobFailure for failures in Sim code', () => { + it('faults a block failure caused by a programming error in Sim code', async () => { + const error = buildBlockExecutionError({ + block, + error: new TypeError('rows.flatMap is not a function'), + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-5', + loggingSession: sessionFinalizing(error, true), + }) + ).resolves.toBe('job_fault') + }) + + it('faults a block failure a tool flattened from a system error', async () => { + const error = Object.assign(new Error('writeLibraryRows: rows.flatMap is not a function'), { + blockId: 'write-library-rows', + blockName: 'writeLibraryRows', + blockType: 'table', + isSystemError: true, + }) + + await expect( + classifyWorkflowJobFailure({ + error, + executionId: 'execution-6', + 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..5702f484b0e --- /dev/null +++ b/apps/sim/lib/workflows/executor/job-failure.ts @@ -0,0 +1,104 @@ +import { findCause, getErrorMessage, isSystemError } from '@sim/utils/errors' +import type { LoggingSession } from '@/lib/logs/execution/logging-session' +import { wasExecutionFinalizedByCore } from '@/lib/workflows/executor/execution-core' +import { + classifyExecutionError, + hasExecutionResult, + type WorkflowExecutionErrorCode, +} from '@/executor/utils/errors' + +/** + * How a background job that ran a workflow ends when the execution threw. + * + * - `workflow_failure`: the workflow itself failed and core recorded it in the + * execution log. That is the run's outcome, so the job completes and reports + * it; faulting the job would page engineering for a user's workflow. + * - `job_fault`: a failure no workflow part owns (engine, setup, a failure core + * could not record). The job re-throws so the queue marks the run failed and + * alerts on it. + */ +export type WorkflowJobFailure = 'workflow_failure' | 'job_fault' + +/** + * What a workflow job returns when the workflow failed. The job itself + * completed, so readers of job status project `success: false` 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(), + } +} + +/** + * Codes that describe the workflow's own outcome even without a block to blame. + * `CANCELLED` is absent: a cancelled run returns rather than throws, and the + * code also matches any engine message containing "cancelled". + */ +const UNATTRIBUTED_WORKFLOW_FAILURE_CODES: ReadonlySet = new Set([ + 'INVALID_INPUT', + 'USAGE_LIMIT_EXCEEDED', +]) + +/** + * Whether Sim's own code, not the workflow, caused the failure anywhere in the + * `.cause` chain. User code never produces a system error, since sandboxes + * return its errors as data. + */ +function isSystemFailure(error: unknown): boolean { + return findCause(error, isSystemError) !== undefined +} + +/** + * 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}). Attribution comes from + * {@link classifyExecutionError}, the single place raw execution errors are + * interpreted, so the job queue, the v2 API and parent workflows agree on what + * a failure was. + */ +export function classifySettledWorkflowJobFailure( + error: unknown, + executionId: string +): WorkflowJobFailure { + if (!wasExecutionFinalizedByCore(error, executionId)) return 'job_fault' + if (isSystemFailure(error)) return 'job_fault' + + const { code, blockId } = classifyExecutionError(error) + return blockId !== undefined || UNATTRIBUTED_WORKFLOW_FAILURE_CODES.has(code) + ? '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..8b5499bbe33 100644 --- a/apps/sim/serializer/index.test.ts +++ b/apps/sim/serializer/index.test.ts @@ -25,6 +25,7 @@ import { import { describe, expect, it, vi } from 'vitest' import { DAGBuilder } from '@/executor/dag/builder' 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 +453,39 @@ describe('Serializer', () => { } ) + it.concurrent('attributes a missing required field to the block that lacks it', () => { + 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, + } + + expect(() => + serializer.serializeWorkflow( + { 'wait-block': waitBlockMissingRequired }, + [], + {}, + undefined, + true + ) + ).toThrow( + expect.objectContaining({ + name: 'WorkflowValidationError', + blockId: 'wait-block', + blockType: 'wait', + blockName: 'Wait Block', + }) + ) + }) + 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..8163ee8e742 100644 --- a/apps/sim/serializer/index.ts +++ b/apps/sim/serializer/index.ts @@ -272,8 +272,11 @@ 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 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..07a4034bf52 100644 --- a/apps/sim/tools/index.test.ts +++ b/apps/sim/tools/index.test.ts @@ -1147,6 +1147,85 @@ describe('executeTool Function', () => { } }) + it.each([ + { thrown: new TypeError('rows.flatMap is not a function'), isSystemError: true }, + { + thrown: new Error('Workspace ID is required in execution context'), + isSystemError: undefined, + }, + ])( + 'marks a failure as a system error only when Sim code raised a programming error ($thrown.name)', + async ({ thrown, isSystemError }) => { + const mockTool = { + id: 'test_operation_input_throws', + name: 'Test Operation Input Throws', + description: 'Throws while building its operation input', + version: '1.0.0', + params: {}, + operation: { + input: () => { + throw thrown + }, + }, + } satisfies InternalToolConfig> + ;(tools as Record).test_operation_input_throws = mockTool + + try { + const result = await executeTool( + 'test_operation_input_throws', + {}, + { + executionContext: createToolExecutionContext({ + userId: 'user-1', + workspaceId: 'workspace-1', + workflowId: 'workflow-1', + }), + } + ) + + expect(result.success).toBe(false) + expect(result.isSystemError).toBe(isSystemError) + } finally { + Reflect.deleteProperty(tools, 'test_operation_input_throws') + } + } + ) + + it('marks a declared operation with no registered handler as a system error', async () => { + const mockTool = { + id: 'test_unregistered_operation', + name: 'Test Unregistered Operation', + description: 'Declares an operation no handler implements', + version: '1.0.0', + params: {}, + operation: { input: createInternalToolOperationInput }, + } satisfies InternalToolConfig> + ;(tools as Record).test_unregistered_operation = mockTool + mockGetInternalToolOperationHandler.mockResolvedValueOnce(undefined) + + try { + const result = await executeTool( + 'test_unregistered_operation', + {}, + { + executionContext: createToolExecutionContext({ + userId: 'user-1', + workspaceId: 'workspace-1', + workflowId: 'workflow-1', + }), + } + ) + + expect(result).toMatchObject({ + success: false, + error: 'No internal operation registered for test_unregistered_operation', + isSystemError: true, + }) + } finally { + Reflect.deleteProperty(tools, 'test_unregistered_operation') + } + }) + it('preserves actorless schedule authority for registered operations', async () => { const mockTool = { id: 'test_actorless_registered_operation', diff --git a/apps/sim/tools/index.ts b/apps/sim/tools/index.ts index b5710fe348b..cbb1da11a65 100644 --- a/apps/sim/tools/index.ts +++ b/apps/sim/tools/index.ts @@ -1,6 +1,13 @@ import { createLogger } from '@sim/logger' import { isLoopbackIp, unwrapIpv6Brackets } from '@sim/security/ssrf' -import { describeError, getErrorMessage, toError } from '@sim/utils/errors' +import { + describeError, + findCause, + getErrorMessage, + isSystemError, + SystemError, + toError, +} from '@sim/utils/errors' import { sleep } from '@sim/utils/helpers' import { isPlainRecord, isRecordLike } from '@sim/utils/object' import { backoffWithJitter, parseRetryAfter } from '@sim/utils/retry' @@ -2424,6 +2431,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 } : {}), + ...(findCause(error, isSystemError) ? { isSystemError: true as const } : {}), timing: { startTime: startTimeISO, endTime: endTimeISO, @@ -2731,7 +2739,7 @@ async function executeDeclaredInternalOperation({ }) } else { const handler = await getInternalToolOperationHandler(toolId) - if (!handler) throw new Error(`No internal operation registered for ${toolId}`) + if (!handler) throw new SystemError(`No internal operation registered for ${toolId}`) const requestedTimeout = Number(params.timeout) const operationTimeout = Number.isFinite(requestedTimeout) && requestedTimeout > 0 @@ -3286,7 +3294,7 @@ async function executeMcpTool( logger.info(`[${actualRequestId}] Executing MCP tool: ${toolId}`) validateRequestBodySize(JSON.stringify(params), actualRequestId, `mcp:${toolId}`) const handler = await getInternalToolOperationHandler(toolId) - if (!handler) throw new Error(`No internal operation registered for ${toolId}`) + if (!handler) throw new SystemError(`No internal operation registered for ${toolId}`) const resultResponse = await handler({ toolId, input: params, diff --git a/apps/sim/tools/types.ts b/apps/sim/tools/types.ts index c26205d140f..1da65f3bc89 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 failure was a programming error in Sim's own code rather than + * anything the workflow did, carried across the same flattening as + * `statusCode` so the job running the workflow can still fault on it. + */ + isSystemError?: true resources?: MothershipResource[] // Resources to auto-open/show in UI largeValueKeys?: string[] fileKeys?: string[] diff --git a/packages/utils/src/errors.test.ts b/packages/utils/src/errors.test.ts index 168afff1818..e3d5dc2df3c 100644 --- a/packages/utils/src/errors.test.ts +++ b/packages/utils/src/errors.test.ts @@ -4,6 +4,7 @@ import { getPostgresCancellationReason, getPostgresErrorCode, getTransientDatabaseFailure, + isProgrammingError, } from '@sim/utils/errors' import { describe, expect, it } from 'vitest' @@ -236,3 +237,27 @@ describe('describeError', () => { expect(described?.causeChain?.length).toBeLessThanOrEqual(10) }) }) + +describe('isProgrammingError', () => { + it('recognizes a defect the running code raised itself', () => { + const thrown = (() => { + try { + JSON.parse('null').flatMap() + } catch (error) { + return error + } + })() + + expect(isProgrammingError(thrown)).toBe(true) + }) + + it('does not treat a runtime API reporting an external failure as a defect', () => { + const networkFailure = new TypeError('fetch failed', { + cause: Object.assign(new Error('connect ECONNREFUSED 127.0.0.1:443'), { + code: 'ECONNREFUSED', + }), + }) + + expect(isProgrammingError(networkFailure)).toBe(false) + }) +}) diff --git a/packages/utils/src/errors.ts b/packages/utils/src/errors.ts index 14248fb4e25..696d16d227a 100644 --- a/packages/utils/src/errors.ts +++ b/packages/utils/src/errors.ts @@ -317,6 +317,50 @@ export function findCause( return undefined } +/** + * Whether `value` is a JavaScript runtime programming error: a type, reference + * or range fault raised by the running code itself. A runtime API reporting an + * external failure wraps what it observed as `cause` (`fetch` rejects a network + * failure as `TypeError('fetch failed', { cause })`), while a defect such as + * calling a missing method never carries one. + */ +export function isProgrammingError( + value: unknown +): value is TypeError | ReferenceError | RangeError { + return ( + (value instanceof TypeError || + value instanceof ReferenceError || + value instanceof RangeError) && + value.cause === undefined + ) +} + +/** + * A broken invariant in Sim's own code, never a rejection of the caller's + * input. Throw it where an invariant is checked so the failure is recognised as + * a platform fault even after it is flattened into a plain error. + */ +export class SystemError extends Error { + readonly isSystemError = true + + constructor(message: string, options?: ErrorOptions) { + super(message, options) + this.name = 'SystemError' + } +} + +/** + * Whether Sim's own code caused `value`: a {@link SystemError}, an error that + * carries its `isSystemError` flag across a flattening boundary, or a + * programming error (see {@link isProgrammingError}). + */ +export function isSystemError(value: unknown): value is Error { + return ( + isProgrammingError(value) || + (value instanceof Error && 'isSystemError' in value && value.isSystemError === true) + ) +} + function readPgErrorField(error: unknown, field: string): string | undefined { const seen = new Set() let current: unknown = error