Skip to content

Commit ab1c14e

Browse files
committed
fix(async-jobs): retain accepted run IDs for immediate lookup
1 parent a1ee856 commit ab1c14e

6 files changed

Lines changed: 219 additions & 3 deletions

File tree

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

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import {
22
asyncJobsRegionMock,
33
asyncJobsRegionMockFns,
44
} from '@sim/testing/mocks/async-jobs-region.mock'
5+
import { dbChainMockFns } from '@sim/testing/mocks/database.mock'
56
import { getMockLogger } from '@sim/testing/mocks/logger.mock'
67
import {
78
MockTriggerApiError as MockApiError,
@@ -170,6 +171,13 @@ describe('TriggerDevJobQueue enqueue', () => {
170171
expect(mockTrigger).not.toHaveBeenCalled()
171172
})
172173

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

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

Lines changed: 65 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,20 @@ export class TriggerDevJobQueue implements JobQueueBackend {
501550
}
502551

503552
try {
553+
if (binding.rootJobId) {
554+
const root = await this.getJob(binding.rootJobId)
555+
if (
556+
root &&
557+
(allowedTaskIdentifiers as readonly string[]).includes(root.type) &&
558+
(root.status === JOB_STATUS.PENDING || root.status === JOB_STATUS.PROCESSING) &&
559+
payloadMatchesExecution(root.payload, binding)
560+
) {
561+
await this.cancelJob(root.id)
562+
cancelledRootRunId = root.id
563+
state.cancelledJobs += 1
564+
}
565+
}
566+
504567
await scanAndCancelTriggerRuns({
505568
binding,
506569
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)