From ab1c14ee4a903e01f86bae284a9c2219e1f64785 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Wed, 30 Sep 2026 16:37:22 -0700 Subject: [PATCH 1/2] fix(async-jobs): retain accepted run IDs for immediate lookup --- .../backends/trigger-dev.integration.ts | 129 ++++++++++++++++++ .../async-jobs/backends/trigger-dev.test.ts | 8 ++ .../core/async-jobs/backends/trigger-dev.ts | 67 ++++++++- apps/sim/lib/core/async-jobs/types.ts | 2 + .../cancel-workflow-execution.test.ts | 4 + .../execution/cancel-workflow-execution.ts | 12 +- 6 files changed, 219 insertions(+), 3 deletions(-) create mode 100644 apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts new file mode 100644 index 00000000000..cc83875fe0b --- /dev/null +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts @@ -0,0 +1,129 @@ +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(() => { + 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(() => ({ + 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') + }) +}) diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts index 70247fdefd8..236cc5b2cf9 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts @@ -2,6 +2,7 @@ import { asyncJobsRegionMock, asyncJobsRegionMockFns, } from '@sim/testing/mocks/async-jobs-region.mock' +import { dbChainMockFns } from '@sim/testing/mocks/database.mock' import { getMockLogger } from '@sim/testing/mocks/logger.mock' import { MockTriggerApiError as MockApiError, @@ -170,6 +171,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() diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts index e65f04f9efb..28ab6709c3a 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts @@ -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, @@ -78,7 +81,36 @@ async function retrieveRunById(jobId: string): Promise { } } -/** 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 { + 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 { + 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 { for await (const candidate of runs.list({ tag: `jobId:${jobId}`, limit: 1 })) { return runs.retrieve(candidate.id) @@ -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 } @@ -417,7 +461,10 @@ export class TriggerDevJobQueue implements JobQueueBackend { async getJob(jobId: string): Promise { 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 @@ -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 = { @@ -501,6 +550,20 @@ export class TriggerDevJobQueue implements JobQueueBackend { } try { + if (binding.rootJobId) { + 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 + } + } + await scanAndCancelTriggerRuns({ binding, cancelJob: (jobId) => this.cancelJob(jobId), diff --git a/apps/sim/lib/core/async-jobs/types.ts b/apps/sim/lib/core/async-jobs/types.ts index 1cbecd49de0..d3c893cd54a 100644 --- a/apps/sim/lib/core/async-jobs/types.ts +++ b/apps/sim/lib/core/async-jobs/types.ts @@ -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' diff --git a/apps/sim/lib/execution/cancel-workflow-execution.test.ts b/apps/sim/lib/execution/cancel-workflow-execution.test.ts index d05fb780570..07ec60b3c8a 100644 --- a/apps/sim/lib/execution/cancel-workflow-execution.test.ts +++ b/apps/sim/lib/execution/cancel-workflow-execution.test.ts @@ -197,6 +197,7 @@ describe('cancelWorkflowExecution', () => { { workflowId: 'wf-1', executionId: 'ex-1', + rootJobId: 'workflow-execution:ex-1', }, 'standalone' ) @@ -775,6 +776,7 @@ describe('cancelWorkflowExecution', () => { { workflowId: 'wf-1', executionId: 'ex-1', + rootJobId: 'workflow-execution:ex-1', }, 'standalone' ) @@ -795,6 +797,7 @@ describe('cancelWorkflowExecution', () => { { workflowId: 'wf-1', executionId: 'ex-1', + rootJobId: 'workflow-execution:ex-1', }, 'standalone' ) @@ -1268,6 +1271,7 @@ describe('cancelWorkflowExecution', () => { { workflowId: 'wf-1', executionId: 'ex-1', + rootJobId: 'workflow-execution:ex-1', }, 'standalone' ) diff --git a/apps/sim/lib/execution/cancel-workflow-execution.ts b/apps/sim/lib/execution/cancel-workflow-execution.ts index 89565495b5b..c5b0ddc09d9 100644 --- a/apps/sim/lib/execution/cancel-workflow-execution.ts +++ b/apps/sim/lib/execution/cancel-workflow-execution.ts @@ -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') @@ -64,7 +65,16 @@ async function cancelQueuedExecutionJobs( ): Promise { 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, From 3c737a501da73c0ccf40219a5f774c260599714f Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Wed, 30 Sep 2026 16:52:53 -0700 Subject: [PATCH 2/2] fix(async-jobs): finish cancellation scans after root errors --- .../backends/trigger-dev.integration.ts | 1 + .../async-jobs/backends/trigger-dev.test.ts | 48 ++++++++++++++++++- .../core/async-jobs/backends/trigger-dev.ts | 24 ++++++---- 3 files changed, 61 insertions(+), 12 deletions(-) diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts index cc83875fe0b..8ba419f040b 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.integration.ts @@ -17,6 +17,7 @@ const runs = new Map< >() beforeEach(() => { + triggerSdkMockFns.mockRunsList.mockReset() runs.clear() triggerSdkMockFns.mockTasksTrigger.mockImplementation(async (type, payload) => { const id = `run_${generateId()}` diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts index 236cc5b2cf9..e4cc6089868 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts @@ -1,8 +1,9 @@ +import { idempotencyKey } from '@sim/db/schema' import { asyncJobsRegionMock, asyncJobsRegionMockFns, } from '@sim/testing/mocks/async-jobs-region.mock' -import { dbChainMockFns } from '@sim/testing/mocks/database.mock' +import { dbChainMockFns, queueTableRows } from '@sim/testing/mocks/database.mock' import { getMockLogger } from '@sim/testing/mocks/logger.mock' import { MockTriggerApiError as MockApiError, @@ -329,7 +330,7 @@ describe('TriggerDevJobQueue status mapping', () => { describe('TriggerDevJobQueue cancellation', () => { beforeEach(() => { - mockList.mockReturnValue( + mockList.mockReset().mockReturnValue( createListPage([ { id: 'run-1', @@ -476,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() + 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([ diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts index 28ab6709c3a..e2f13aba5c9 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts @@ -551,16 +551,20 @@ export class TriggerDevJobQueue implements JobQueueBackend { try { if (binding.rootJobId) { - 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 + 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) } }