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
81 changes: 72 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,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'

Expand Down Expand Up @@ -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([
Expand Down Expand Up @@ -379,7 +374,6 @@ describe('async preprocessing correlation threading', () => {
executionTimeout: {},
})
mockExecuteWorkflowCore.mockRejectedValueOnce(rawError)
mockHasExecutionResult.mockImplementation((error) => error === rawError)
mockWasExecutionFinalizedByCore.mockReturnValue(true)

await expect(
Expand All @@ -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()
})
Expand Down Expand Up @@ -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) })
)
})
})
})
25 changes: 25 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,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 {
Expand Down Expand Up @@ -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,
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') {
Comment thread
waleedlatif1 marked this conversation as resolved.
return {
...buildWorkflowJobFailureResult({ error, workflowId, executionId: resumeExecutionId }),
parentExecutionId,
status: 'failed' as const,
}
}
throw error
} finally {
timeoutController?.cleanup()
Expand Down
14 changes: 14 additions & 0 deletions apps/sim/background/schedule-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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 {
Expand All @@ -1211,6 +1222,7 @@ export async function executeScheduleJob(
return
}

jobFault = { error }
logger.error(`[${requestId}] Error processing schedule ${payload.scheduleId}`, error, {
cause: describeError(error),
})
Expand All @@ -1230,6 +1242,8 @@ export async function executeScheduleJob(
trace.getActiveSpan()?.recordException(toError(recoveryError))
}
}

if (jobFault) throw jobFault.error
})
} finally {
timeoutController.cleanup()
Expand Down
Loading
Loading