Skip to content
Merged
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
130 changes: 130 additions & 0 deletions apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
import { db } from '@sim/db'
import { idempotencyKey } from '@sim/db/schema'
import { asyncJobsRegionMock } from '@sim/testing/mocks/async-jobs-region.mock'
import { triggerSdkMock, triggerSdkMockFns } from '@sim/testing/mocks/trigger-sdk.mock'
import { generateId } from '@sim/utils/id'
import { inArray } from 'drizzle-orm'
import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest'
import { TriggerDevJobQueue } from '@/lib/core/async-jobs/backends/trigger-dev'

vi.mock('@trigger.dev/sdk', () => triggerSdkMock)
vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock)

const receipts: string[] = []
const runs = new Map<
string,
{ id: string; taskIdentifier: string; status: string; payload: unknown; createdAt: Date }
>()

beforeEach(() => {
triggerSdkMockFns.mockRunsList.mockReset()
runs.clear()
triggerSdkMockFns.mockTasksTrigger.mockImplementation(async (type, payload) => {
const id = `run_${generateId()}`
runs.set(id, { id, taskIdentifier: type, status: 'QUEUED', payload, createdAt: new Date() })
return { id }
})
triggerSdkMockFns.mockRunsRetrieve.mockImplementation(async (id) => {
const run = runs.get(id)
if (!run) throw new Error('Run not found')
return run
})
triggerSdkMockFns.mockRunsCancel.mockImplementation(async (id) => {
const run = runs.get(id)
if (!run) throw new Error('Run not found')
run.status = 'CANCELED'
return run
})
})

afterAll(async () => {
if (receipts.length) await db.delete(idempotencyKey).where(inArray(idempotencyKey.key, receipts))
})

async function enqueue() {
const executionId = generateId()
const workflowId = generateId()
const rootJobId = `workflow-execution:${executionId}`
receipts.push(`trigger-job:${rootJobId}`)
const id = await new TriggerDevJobQueue().enqueue(
'workflow-execution',
{ executionId, workflowId },
{ jobId: rootJobId }
)
return { id, binding: { workflowId, executionId, rootJobId } }
}

describe('accepted Trigger runs before tag indexing', () => {
it('persists a receipt and lets another queue instance read the accepted run', async () => {
const { id, binding } = await enqueue()
const reader = new TriggerDevJobQueue()
expect(await reader.getJob(binding.rootJobId)).toMatchObject({
id,
status: 'pending',
metadata: { workflowId: binding.workflowId },
})
const stored = await db
.select()
.from(idempotencyKey)
.where(inArray(idempotencyKey.key, receipts))
expect(
stored.some(
(row) =>
row.key === `trigger-job:${binding.rootJobId}` &&
(row.result as { runId: string }).runId === id
)
).toBe(true)
})

it('cancels an accepted run while every tag search is still empty', async () => {
const { id, binding } = await enqueue()
expect(await new TriggerDevJobQueue().cancelByExecution(binding, 'standalone')).toBe(1)
expect(runs.get(id)?.status).toBe('CANCELED')
})

it('keeps a retry discoverable and cancels a run only once when its tags catch up', async () => {
const { id, binding } = await enqueue()
triggerSdkMockFns.mockTasksTrigger.mockResolvedValueOnce({ id })
await new TriggerDevJobQueue().enqueue(
'workflow-execution',
{
workflowId: binding.workflowId,
executionId: binding.executionId,
},
{ jobId: binding.rootJobId }
)
triggerSdkMockFns.mockRunsList.mockImplementation(() => ({
Comment thread
icecrasher321 marked this conversation as resolved.
async *[Symbol.asyncIterator]() {
yield {
id,
taskIdentifier: 'workflow-execution',
tags: [`workflowId:${binding.workflowId}`, `executionId:${binding.executionId}`],
}
},
}))
const queue = new TriggerDevJobQueue()
expect(await queue.getJob(binding.rootJobId)).toMatchObject({ id })
expect(await queue.cancelByExecution(binding, 'standalone')).toBe(1)
expect(runs.get(id)?.status).toBe('CANCELED')
})

it('does not cancel a receipt belonging to another workflow or cancellation scope', async () => {
const { id, binding } = await enqueue()
const queue = new TriggerDevJobQueue()
expect(
await queue.cancelByExecution({ ...binding, workflowId: generateId() }, 'standalone')
).toBe(0)
expect(
await queue.cancelByExecution({ ...binding, executionId: generateId() }, 'standalone')
).toBe(0)
expect(await queue.cancelByExecution(binding, 'resume')).toBe(0)
expect(runs.get(id)?.status).toBe('QUEUED')
})

it('does not report a completed receipt as a successful cancellation', async () => {
const { id, binding } = await enqueue()
runs.get(id)!.status = 'COMPLETED'
expect(await new TriggerDevJobQueue().cancelByExecution(binding, 'standalone')).toBe(0)
expect(runs.get(id)?.status).toBe('COMPLETED')
})
})
54 changes: 53 additions & 1 deletion apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import { idempotencyKey } from '@sim/db/schema'
import {
asyncJobsRegionMock,
asyncJobsRegionMockFns,
} from '@sim/testing/mocks/async-jobs-region.mock'
import { dbChainMockFns, queueTableRows } from '@sim/testing/mocks/database.mock'
import { getMockLogger } from '@sim/testing/mocks/logger.mock'
import {
MockTriggerApiError as MockApiError,
Expand Down Expand Up @@ -170,6 +172,13 @@ describe('TriggerDevJobQueue enqueue', () => {
expect(mockTrigger).not.toHaveBeenCalled()
})

it('preserves ambiguous acceptance when the run receipt cannot be persisted', async () => {
dbChainMockFns.onConflictDoUpdate.mockRejectedValueOnce(new Error('database unavailable'))
await expect(
new TriggerDevJobQueue().enqueue('workflow-execution', {}, { jobId: 'workflow:1' })
).rejects.toMatchObject({ acceptance: 'unknown', retryable: true })
})

it('classifies a client response as proven non-acceptance', async () => {
mockTrigger.mockRejectedValueOnce(new MockApiError(422, 'invalid payload'))
const queue = new TriggerDevJobQueue()
Expand Down Expand Up @@ -321,7 +330,7 @@ describe('TriggerDevJobQueue status mapping', () => {

describe('TriggerDevJobQueue cancellation', () => {
beforeEach(() => {
mockList.mockReturnValue(
mockList.mockReset().mockReturnValue(
createListPage([
{
id: 'run-1',
Expand Down Expand Up @@ -468,6 +477,49 @@ describe('TriggerDevJobQueue cancellation', () => {
})
})

it.each(['receipt', 'retrieve', 'cancel'] as const)(
'continues every discovery phase after a root %s failure',
async (failurePhase) => {
const failure = new Error(`root ${failurePhase} unavailable`)
const payload = { workflowId: 'workflow-1', executionId: 'execution-1' }
const cancelled = new Set<string>()
if (failurePhase === 'receipt') {
dbChainMockFns.limit.mockRejectedValueOnce(failure)
} else {
queueTableRows(idempotencyKey, [{ result: { runId: 'run_root' } }])
}
mockRetrieve.mockImplementation(async (id: string) => {
if (id === 'run_root' && failurePhase === 'retrieve') throw failure
return { id, taskIdentifier: 'workflow-execution', status: 'QUEUED', payload }
})
mockCancel.mockImplementation(async (id: string) => {
if (id === 'run_root') throw failure
cancelled.add(id)
})
mockList
.mockReturnValueOnce(
createListPage([
{ id: 'tagged', tags: ['workflowId:workflow-1', 'executionId:execution-1'] },
])
)
.mockReturnValueOnce(
createListPage([{ id: 'legacy-tagged', tags: ['workflowId:workflow-1'] }])
)
.mockReturnValueOnce(createListPage([{ id: 'legacy-untagged', tags: [] }]))

await expect(
new TriggerDevJobQueue().cancelByExecution(
{
...payload,
rootJobId: 'workflow-execution:execution-1',
},
'standalone'
)
).rejects.toBe(failure)
expect(cancelled).toEqual(new Set(['tagged', 'legacy-tagged', 'legacy-untagged']))
}
)

it('cancels legacy workflow-tagged runs only after payload verification', async () => {
mockList.mockReturnValueOnce(createListPage([])).mockReturnValueOnce(
createListPage([
Expand Down
71 changes: 69 additions & 2 deletions apps/sim/lib/core/async-jobs/backends/trigger-dev.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,12 @@
import { db } from '@sim/db'
import { idempotencyKey } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { sha256Hex } from '@sim/security/hash'
import { toError } from '@sim/utils/errors'
import { isRecordLike } from '@sim/utils/object'
import { taskContext } from '@trigger.dev/core/v3'
import { ApiError, runs, type TriggerOptions, tasks } from '@trigger.dev/sdk'
import { eq } from 'drizzle-orm'
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
import {
AsyncJobEnqueueError,
Expand Down Expand Up @@ -78,7 +81,36 @@ async function retrieveRunById(jobId: string): Promise<TriggerRun | null> {
}
}

/** Resolves a caller-chosen job id through the `jobId:` tag set at enqueue. */
/**
* Retains Trigger's accepted run ID before enqueue can return success. The
* idempotency result table shares receipts across app processes; its ordinary
* retention bounds storage, with tag lookup retained for older jobs.
*/
async function storeRunReceipt(jobId: string, runId: string): Promise<void> {
const result = { runId }
await db
.insert(idempotencyKey)
.values({ key: `trigger-job:${jobId}`, result })
.onConflictDoUpdate({
target: idempotencyKey.key,
set: { result, createdAt: new Date() },
})
}

/** Reads accepted run IDs without depending on Trigger's asynchronous tag index. */
async function retrieveRunByReceipt(jobId: string): Promise<TriggerRun | null> {
const [receipt] = await db
.select({ result: idempotencyKey.result })
.from(idempotencyKey)
.where(eq(idempotencyKey.key, `trigger-job:${jobId}`))
.limit(1)
const result = receipt?.result
return isRecordLike(result) && typeof result.runId === 'string'
? retrieveRunById(result.runId)
: null
}

/** Resolves legacy or expired receipts through the `jobId:` tag set at enqueue. */
async function retrieveRunByJobIdTag(jobId: string): Promise<TriggerRun | null> {
for await (const candidate of runs.list({ tag: `jobId:${jobId}`, limit: 1 })) {
return runs.retrieve(candidate.id)
Expand Down Expand Up @@ -331,6 +363,18 @@ export class TriggerDevJobQueue implements JobQueueBackend {
throw classifyTriggerEnqueueError(error)
}

if (options?.jobId) {
try {
await storeRunReceipt(options.jobId, handle.id)
} catch (error) {
throw new AsyncJobEnqueueError('Trigger run accepted but its receipt could not be stored', {
acceptance: 'unknown',
retryable: true,
cause: error,
})
}
}

logger.debug('Enqueued job via trigger.dev', { jobId: handle.id, type, taskId, tags })
return handle.id
}
Expand Down Expand Up @@ -417,7 +461,10 @@ export class TriggerDevJobQueue implements JobQueueBackend {

async getJob(jobId: string): Promise<Job | null> {
try {
const run = (await retrieveRunById(jobId)) ?? (await retrieveRunByJobIdTag(jobId))
const run =
(await retrieveRunById(jobId)) ??
(jobId.startsWith(TRIGGER_RUN_ID_PREFIX) ? null : await retrieveRunByReceipt(jobId)) ??
(await retrieveRunByJobIdTag(jobId))
if (!run) {
logger.debug('Job not found in trigger.dev', { jobId })
return null
Expand Down Expand Up @@ -492,7 +539,9 @@ export class TriggerDevJobQueue implements JobQueueBackend {
const executionTag = buildExecutionTag(binding.executionId)
const workflowTag = buildWorkflowTag(binding.workflowId)
const allowedTaskIdentifiers = EXECUTION_JOB_TYPES_BY_CANCELLATION_SCOPE[scope]
let cancelledRootRunId: string | undefined
const isAllowedTask = (run: CancellationListRun) =>
run.id !== cancelledRootRunId &&
(allowedTaskIdentifiers as readonly string[]).includes(run.taskIdentifier)
const cutoff = new Date(Date.now() - JOB_PENDING_RETENTION_HOURS * 60 * 60 * 1000)
const state: CancellationScanState = {
Expand All @@ -501,6 +550,24 @@ export class TriggerDevJobQueue implements JobQueueBackend {
}

try {
if (binding.rootJobId) {
try {
const root = await this.getJob(binding.rootJobId)
if (
root &&
(allowedTaskIdentifiers as readonly string[]).includes(root.type) &&
(root.status === JOB_STATUS.PENDING || root.status === JOB_STATUS.PROCESSING) &&
payloadMatchesExecution(root.payload, binding)
) {
await this.cancelJob(root.id)
cancelledRootRunId = root.id
state.cancelledJobs += 1
}
} catch (error) {
recordCancellationCandidateFailure(state, error)
}
}

await scanAndCancelTriggerRuns({
binding,
cancelJob: (jobId) => this.cancelJob(jobId),
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/lib/core/async-jobs/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,8 @@ export interface EnqueueOptions {
export interface ExecutionJobBinding {
workflowId: string
executionId: string
/** Known root job identity; cancellation must still verify workflow, execution, and scope. */
rootJobId?: string
}

export type ExecutionJobCancellationScope = 'standalone' | 'resume'
Expand Down
4 changes: 4 additions & 0 deletions apps/sim/lib/execution/cancel-workflow-execution.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ describe('cancelWorkflowExecution', () => {
{
workflowId: 'wf-1',
executionId: 'ex-1',
rootJobId: 'workflow-execution:ex-1',
},
'standalone'
)
Expand Down Expand Up @@ -775,6 +776,7 @@ describe('cancelWorkflowExecution', () => {
{
workflowId: 'wf-1',
executionId: 'ex-1',
rootJobId: 'workflow-execution:ex-1',
},
'standalone'
)
Expand All @@ -795,6 +797,7 @@ describe('cancelWorkflowExecution', () => {
{
workflowId: 'wf-1',
executionId: 'ex-1',
rootJobId: 'workflow-execution:ex-1',
},
'standalone'
)
Expand Down Expand Up @@ -1268,6 +1271,7 @@ describe('cancelWorkflowExecution', () => {
{
workflowId: 'wf-1',
executionId: 'ex-1',
rootJobId: 'workflow-execution:ex-1',
},
'standalone'
)
Expand Down
12 changes: 11 additions & 1 deletion apps/sim/lib/execution/cancel-workflow-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import {
type PublishableWorkflowGroupCancellation,
publishWorkflowGroupCancellationEvent,
} from '@/lib/table/workflow-group-cancellation'
import { WORKFLOW_EXECUTION_JOB_ID_PREFIX } from '@/lib/workflows/executor/execution-job-ids'
import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager'

const logger = createLogger('CancelWorkflowExecution')
Expand Down Expand Up @@ -64,7 +65,16 @@ async function cancelQueuedExecutionJobs(
): Promise<number> {
try {
const queue = await getJobQueue()
return await queue.cancelByExecution({ workflowId, executionId }, scope)
return await queue.cancelByExecution(
{
workflowId,
executionId,
...(scope === 'standalone'
? { rootJobId: `${WORKFLOW_EXECUTION_JOB_ID_PREFIX}${executionId}` }
: {}),
},
scope
)
} catch (error) {
logger.warn('Failed to cancel queued execution jobs', {
workflowId,
Expand Down
Loading