Skip to content

Commit b84828d

Browse files
fix(async-jobs): retain accepted run IDs for immediate lookup (#8490)
* fix(async-jobs): retain accepted run IDs for immediate lookup * fix(async-jobs): finish cancellation scans after root errors
1 parent a1ee856 commit b84828d

6 files changed

Lines changed: 269 additions & 4 deletions

File tree

Lines changed: 130 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,130 @@
1+
import { db } from '@sim/db'
2+
import { idempotencyKey } from '@sim/db/schema'
3+
import { asyncJobsRegionMock } from '@sim/testing/mocks/async-jobs-region.mock'
4+
import { triggerSdkMock, triggerSdkMockFns } from '@sim/testing/mocks/trigger-sdk.mock'
5+
import { generateId } from '@sim/utils/id'
6+
import { inArray } from 'drizzle-orm'
7+
import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest'
8+
import { TriggerDevJobQueue } from '@/lib/core/async-jobs/backends/trigger-dev'
9+
10+
vi.mock('@trigger.dev/sdk', () => triggerSdkMock)
11+
vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock)
12+
13+
const receipts: string[] = []
14+
const runs = new Map<
15+
string,
16+
{ id: string; taskIdentifier: string; status: string; payload: unknown; createdAt: Date }
17+
>()
18+
19+
beforeEach(() => {
20+
triggerSdkMockFns.mockRunsList.mockReset()
21+
runs.clear()
22+
triggerSdkMockFns.mockTasksTrigger.mockImplementation(async (type, payload) => {
23+
const id = `run_${generateId()}`
24+
runs.set(id, { id, taskIdentifier: type, status: 'QUEUED', payload, createdAt: new Date() })
25+
return { id }
26+
})
27+
triggerSdkMockFns.mockRunsRetrieve.mockImplementation(async (id) => {
28+
const run = runs.get(id)
29+
if (!run) throw new Error('Run not found')
30+
return run
31+
})
32+
triggerSdkMockFns.mockRunsCancel.mockImplementation(async (id) => {
33+
const run = runs.get(id)
34+
if (!run) throw new Error('Run not found')
35+
run.status = 'CANCELED'
36+
return run
37+
})
38+
})
39+
40+
afterAll(async () => {
41+
if (receipts.length) await db.delete(idempotencyKey).where(inArray(idempotencyKey.key, receipts))
42+
})
43+
44+
async function enqueue() {
45+
const executionId = generateId()
46+
const workflowId = generateId()
47+
const rootJobId = `workflow-execution:${executionId}`
48+
receipts.push(`trigger-job:${rootJobId}`)
49+
const id = await new TriggerDevJobQueue().enqueue(
50+
'workflow-execution',
51+
{ executionId, workflowId },
52+
{ jobId: rootJobId }
53+
)
54+
return { id, binding: { workflowId, executionId, rootJobId } }
55+
}
56+
57+
describe('accepted Trigger runs before tag indexing', () => {
58+
it('persists a receipt and lets another queue instance read the accepted run', async () => {
59+
const { id, binding } = await enqueue()
60+
const reader = new TriggerDevJobQueue()
61+
expect(await reader.getJob(binding.rootJobId)).toMatchObject({
62+
id,
63+
status: 'pending',
64+
metadata: { workflowId: binding.workflowId },
65+
})
66+
const stored = await db
67+
.select()
68+
.from(idempotencyKey)
69+
.where(inArray(idempotencyKey.key, receipts))
70+
expect(
71+
stored.some(
72+
(row) =>
73+
row.key === `trigger-job:${binding.rootJobId}` &&
74+
(row.result as { runId: string }).runId === id
75+
)
76+
).toBe(true)
77+
})
78+
79+
it('cancels an accepted run while every tag search is still empty', async () => {
80+
const { id, binding } = await enqueue()
81+
expect(await new TriggerDevJobQueue().cancelByExecution(binding, 'standalone')).toBe(1)
82+
expect(runs.get(id)?.status).toBe('CANCELED')
83+
})
84+
85+
it('keeps a retry discoverable and cancels a run only once when its tags catch up', async () => {
86+
const { id, binding } = await enqueue()
87+
triggerSdkMockFns.mockTasksTrigger.mockResolvedValueOnce({ id })
88+
await new TriggerDevJobQueue().enqueue(
89+
'workflow-execution',
90+
{
91+
workflowId: binding.workflowId,
92+
executionId: binding.executionId,
93+
},
94+
{ jobId: binding.rootJobId }
95+
)
96+
triggerSdkMockFns.mockRunsList.mockImplementation(() => ({
97+
async *[Symbol.asyncIterator]() {
98+
yield {
99+
id,
100+
taskIdentifier: 'workflow-execution',
101+
tags: [`workflowId:${binding.workflowId}`, `executionId:${binding.executionId}`],
102+
}
103+
},
104+
}))
105+
const queue = new TriggerDevJobQueue()
106+
expect(await queue.getJob(binding.rootJobId)).toMatchObject({ id })
107+
expect(await queue.cancelByExecution(binding, 'standalone')).toBe(1)
108+
expect(runs.get(id)?.status).toBe('CANCELED')
109+
})
110+
111+
it('does not cancel a receipt belonging to another workflow or cancellation scope', async () => {
112+
const { id, binding } = await enqueue()
113+
const queue = new TriggerDevJobQueue()
114+
expect(
115+
await queue.cancelByExecution({ ...binding, workflowId: generateId() }, 'standalone')
116+
).toBe(0)
117+
expect(
118+
await queue.cancelByExecution({ ...binding, executionId: generateId() }, 'standalone')
119+
).toBe(0)
120+
expect(await queue.cancelByExecution(binding, 'resume')).toBe(0)
121+
expect(runs.get(id)?.status).toBe('QUEUED')
122+
})
123+
124+
it('does not report a completed receipt as a successful cancellation', async () => {
125+
const { id, binding } = await enqueue()
126+
runs.get(id)!.status = 'COMPLETED'
127+
expect(await new TriggerDevJobQueue().cancelByExecution(binding, 'standalone')).toBe(0)
128+
expect(runs.get(id)?.status).toBe('COMPLETED')
129+
})
130+
})

‎apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts‎

Lines changed: 53 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
1+
import { idempotencyKey } from '@sim/db/schema'
12
import {
23
asyncJobsRegionMock,
34
asyncJobsRegionMockFns,
45
} from '@sim/testing/mocks/async-jobs-region.mock'
6+
import { dbChainMockFns, queueTableRows } from '@sim/testing/mocks/database.mock'
57
import { getMockLogger } from '@sim/testing/mocks/logger.mock'
68
import {
79
MockTriggerApiError as MockApiError,
@@ -170,6 +172,13 @@ describe('TriggerDevJobQueue enqueue', () => {
170172
expect(mockTrigger).not.toHaveBeenCalled()
171173
})
172174

175+
it('preserves ambiguous acceptance when the run receipt cannot be persisted', async () => {
176+
dbChainMockFns.onConflictDoUpdate.mockRejectedValueOnce(new Error('database unavailable'))
177+
await expect(
178+
new TriggerDevJobQueue().enqueue('workflow-execution', {}, { jobId: 'workflow:1' })
179+
).rejects.toMatchObject({ acceptance: 'unknown', retryable: true })
180+
})
181+
173182
it('classifies a client response as proven non-acceptance', async () => {
174183
mockTrigger.mockRejectedValueOnce(new MockApiError(422, 'invalid payload'))
175184
const queue = new TriggerDevJobQueue()
@@ -321,7 +330,7 @@ describe('TriggerDevJobQueue status mapping', () => {
321330

322331
describe('TriggerDevJobQueue cancellation', () => {
323332
beforeEach(() => {
324-
mockList.mockReturnValue(
333+
mockList.mockReset().mockReturnValue(
325334
createListPage([
326335
{
327336
id: 'run-1',
@@ -468,6 +477,49 @@ describe('TriggerDevJobQueue cancellation', () => {
468477
})
469478
})
470479

480+
it.each(['receipt', 'retrieve', 'cancel'] as const)(
481+
'continues every discovery phase after a root %s failure',
482+
async (failurePhase) => {
483+
const failure = new Error(`root ${failurePhase} unavailable`)
484+
const payload = { workflowId: 'workflow-1', executionId: 'execution-1' }
485+
const cancelled = new Set<string>()
486+
if (failurePhase === 'receipt') {
487+
dbChainMockFns.limit.mockRejectedValueOnce(failure)
488+
} else {
489+
queueTableRows(idempotencyKey, [{ result: { runId: 'run_root' } }])
490+
}
491+
mockRetrieve.mockImplementation(async (id: string) => {
492+
if (id === 'run_root' && failurePhase === 'retrieve') throw failure
493+
return { id, taskIdentifier: 'workflow-execution', status: 'QUEUED', payload }
494+
})
495+
mockCancel.mockImplementation(async (id: string) => {
496+
if (id === 'run_root') throw failure
497+
cancelled.add(id)
498+
})
499+
mockList
500+
.mockReturnValueOnce(
501+
createListPage([
502+
{ id: 'tagged', tags: ['workflowId:workflow-1', 'executionId:execution-1'] },
503+
])
504+
)
505+
.mockReturnValueOnce(
506+
createListPage([{ id: 'legacy-tagged', tags: ['workflowId:workflow-1'] }])
507+
)
508+
.mockReturnValueOnce(createListPage([{ id: 'legacy-untagged', tags: [] }]))
509+
510+
await expect(
511+
new TriggerDevJobQueue().cancelByExecution(
512+
{
513+
...payload,
514+
rootJobId: 'workflow-execution:execution-1',
515+
},
516+
'standalone'
517+
)
518+
).rejects.toBe(failure)
519+
expect(cancelled).toEqual(new Set(['tagged', 'legacy-tagged', 'legacy-untagged']))
520+
}
521+
)
522+
471523
it('cancels legacy workflow-tagged runs only after payload verification', async () => {
472524
mockList.mockReturnValueOnce(createListPage([])).mockReturnValueOnce(
473525
createListPage([

‎apps/sim/lib/core/async-jobs/backends/trigger-dev.ts‎

Lines changed: 69 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,12 @@
1+
import { db } from '@sim/db'
2+
import { idempotencyKey } from '@sim/db/schema'
13
import { createLogger } from '@sim/logger'
24
import { sha256Hex } from '@sim/security/hash'
35
import { toError } from '@sim/utils/errors'
46
import { isRecordLike } from '@sim/utils/object'
57
import { taskContext } from '@trigger.dev/core/v3'
68
import { ApiError, runs, type TriggerOptions, tasks } from '@trigger.dev/sdk'
9+
import { eq } from 'drizzle-orm'
710
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
811
import {
912
AsyncJobEnqueueError,
@@ -78,7 +81,36 @@ async function retrieveRunById(jobId: string): Promise<TriggerRun | null> {
7881
}
7982
}
8083

81-
/** Resolves a caller-chosen job id through the `jobId:` tag set at enqueue. */
84+
/**
85+
* Retains Trigger's accepted run ID before enqueue can return success. The
86+
* idempotency result table shares receipts across app processes; its ordinary
87+
* retention bounds storage, with tag lookup retained for older jobs.
88+
*/
89+
async function storeRunReceipt(jobId: string, runId: string): Promise<void> {
90+
const result = { runId }
91+
await db
92+
.insert(idempotencyKey)
93+
.values({ key: `trigger-job:${jobId}`, result })
94+
.onConflictDoUpdate({
95+
target: idempotencyKey.key,
96+
set: { result, createdAt: new Date() },
97+
})
98+
}
99+
100+
/** Reads accepted run IDs without depending on Trigger's asynchronous tag index. */
101+
async function retrieveRunByReceipt(jobId: string): Promise<TriggerRun | null> {
102+
const [receipt] = await db
103+
.select({ result: idempotencyKey.result })
104+
.from(idempotencyKey)
105+
.where(eq(idempotencyKey.key, `trigger-job:${jobId}`))
106+
.limit(1)
107+
const result = receipt?.result
108+
return isRecordLike(result) && typeof result.runId === 'string'
109+
? retrieveRunById(result.runId)
110+
: null
111+
}
112+
113+
/** Resolves legacy or expired receipts through the `jobId:` tag set at enqueue. */
82114
async function retrieveRunByJobIdTag(jobId: string): Promise<TriggerRun | null> {
83115
for await (const candidate of runs.list({ tag: `jobId:${jobId}`, limit: 1 })) {
84116
return runs.retrieve(candidate.id)
@@ -331,6 +363,18 @@ export class TriggerDevJobQueue implements JobQueueBackend {
331363
throw classifyTriggerEnqueueError(error)
332364
}
333365

366+
if (options?.jobId) {
367+
try {
368+
await storeRunReceipt(options.jobId, handle.id)
369+
} catch (error) {
370+
throw new AsyncJobEnqueueError('Trigger run accepted but its receipt could not be stored', {
371+
acceptance: 'unknown',
372+
retryable: true,
373+
cause: error,
374+
})
375+
}
376+
}
377+
334378
logger.debug('Enqueued job via trigger.dev', { jobId: handle.id, type, taskId, tags })
335379
return handle.id
336380
}
@@ -417,7 +461,10 @@ export class TriggerDevJobQueue implements JobQueueBackend {
417461

418462
async getJob(jobId: string): Promise<Job | null> {
419463
try {
420-
const run = (await retrieveRunById(jobId)) ?? (await retrieveRunByJobIdTag(jobId))
464+
const run =
465+
(await retrieveRunById(jobId)) ??
466+
(jobId.startsWith(TRIGGER_RUN_ID_PREFIX) ? null : await retrieveRunByReceipt(jobId)) ??
467+
(await retrieveRunByJobIdTag(jobId))
421468
if (!run) {
422469
logger.debug('Job not found in trigger.dev', { jobId })
423470
return null
@@ -492,7 +539,9 @@ export class TriggerDevJobQueue implements JobQueueBackend {
492539
const executionTag = buildExecutionTag(binding.executionId)
493540
const workflowTag = buildWorkflowTag(binding.workflowId)
494541
const allowedTaskIdentifiers = EXECUTION_JOB_TYPES_BY_CANCELLATION_SCOPE[scope]
542+
let cancelledRootRunId: string | undefined
495543
const isAllowedTask = (run: CancellationListRun) =>
544+
run.id !== cancelledRootRunId &&
496545
(allowedTaskIdentifiers as readonly string[]).includes(run.taskIdentifier)
497546
const cutoff = new Date(Date.now() - JOB_PENDING_RETENTION_HOURS * 60 * 60 * 1000)
498547
const state: CancellationScanState = {
@@ -501,6 +550,24 @@ export class TriggerDevJobQueue implements JobQueueBackend {
501550
}
502551

503552
try {
553+
if (binding.rootJobId) {
554+
try {
555+
const root = await this.getJob(binding.rootJobId)
556+
if (
557+
root &&
558+
(allowedTaskIdentifiers as readonly string[]).includes(root.type) &&
559+
(root.status === JOB_STATUS.PENDING || root.status === JOB_STATUS.PROCESSING) &&
560+
payloadMatchesExecution(root.payload, binding)
561+
) {
562+
await this.cancelJob(root.id)
563+
cancelledRootRunId = root.id
564+
state.cancelledJobs += 1
565+
}
566+
} catch (error) {
567+
recordCancellationCandidateFailure(state, error)
568+
}
569+
}
570+
504571
await scanAndCancelTriggerRuns({
505572
binding,
506573
cancelJob: (jobId) => this.cancelJob(jobId),

‎apps/sim/lib/core/async-jobs/types.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -163,6 +163,8 @@ export interface EnqueueOptions {
163163
export interface ExecutionJobBinding {
164164
workflowId: string
165165
executionId: string
166+
/** Known root job identity; cancellation must still verify workflow, execution, and scope. */
167+
rootJobId?: string
166168
}
167169

168170
export type ExecutionJobCancellationScope = 'standalone' | 'resume'

‎apps/sim/lib/execution/cancel-workflow-execution.test.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -197,6 +197,7 @@ describe('cancelWorkflowExecution', () => {
197197
{
198198
workflowId: 'wf-1',
199199
executionId: 'ex-1',
200+
rootJobId: 'workflow-execution:ex-1',
200201
},
201202
'standalone'
202203
)
@@ -775,6 +776,7 @@ describe('cancelWorkflowExecution', () => {
775776
{
776777
workflowId: 'wf-1',
777778
executionId: 'ex-1',
779+
rootJobId: 'workflow-execution:ex-1',
778780
},
779781
'standalone'
780782
)
@@ -795,6 +797,7 @@ describe('cancelWorkflowExecution', () => {
795797
{
796798
workflowId: 'wf-1',
797799
executionId: 'ex-1',
800+
rootJobId: 'workflow-execution:ex-1',
798801
},
799802
'standalone'
800803
)
@@ -1268,6 +1271,7 @@ describe('cancelWorkflowExecution', () => {
12681271
{
12691272
workflowId: 'wf-1',
12701273
executionId: 'ex-1',
1274+
rootJobId: 'workflow-execution:ex-1',
12711275
},
12721276
'standalone'
12731277
)

‎apps/sim/lib/execution/cancel-workflow-execution.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import {
2626
type PublishableWorkflowGroupCancellation,
2727
publishWorkflowGroupCancellationEvent,
2828
} from '@/lib/table/workflow-group-cancellation'
29+
import { WORKFLOW_EXECUTION_JOB_ID_PREFIX } from '@/lib/workflows/executor/execution-job-ids'
2930
import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager'
3031

3132
const logger = createLogger('CancelWorkflowExecution')
@@ -64,7 +65,16 @@ async function cancelQueuedExecutionJobs(
6465
): Promise<number> {
6566
try {
6667
const queue = await getJobQueue()
67-
return await queue.cancelByExecution({ workflowId, executionId }, scope)
68+
return await queue.cancelByExecution(
69+
{
70+
workflowId,
71+
executionId,
72+
...(scope === 'standalone'
73+
? { rootJobId: `${WORKFLOW_EXECUTION_JOB_ID_PREFIX}${executionId}` }
74+
: {}),
75+
},
76+
scope
77+
)
6878
} catch (error) {
6979
logger.warn('Failed to cancel queued execution jobs', {
7080
workflowId,

0 commit comments

Comments
 (0)