Skip to content

Commit 3c737a5

Browse files
committed
fix(async-jobs): finish cancellation scans after root errors
1 parent ab1c14e commit 3c737a5

3 files changed

Lines changed: 61 additions & 12 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ const runs = new Map<
1717
>()
1818

1919
beforeEach(() => {
20+
triggerSdkMockFns.mockRunsList.mockReset()
2021
runs.clear()
2122
triggerSdkMockFns.mockTasksTrigger.mockImplementation(async (type, payload) => {
2223
const id = `run_${generateId()}`

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

Lines changed: 46 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,9 @@
1+
import { idempotencyKey } from '@sim/db/schema'
12
import {
23
asyncJobsRegionMock,
34
asyncJobsRegionMockFns,
45
} from '@sim/testing/mocks/async-jobs-region.mock'
5-
import { dbChainMockFns } from '@sim/testing/mocks/database.mock'
6+
import { dbChainMockFns, queueTableRows } from '@sim/testing/mocks/database.mock'
67
import { getMockLogger } from '@sim/testing/mocks/logger.mock'
78
import {
89
MockTriggerApiError as MockApiError,
@@ -329,7 +330,7 @@ describe('TriggerDevJobQueue status mapping', () => {
329330

330331
describe('TriggerDevJobQueue cancellation', () => {
331332
beforeEach(() => {
332-
mockList.mockReturnValue(
333+
mockList.mockReset().mockReturnValue(
333334
createListPage([
334335
{
335336
id: 'run-1',
@@ -476,6 +477,49 @@ describe('TriggerDevJobQueue cancellation', () => {
476477
})
477478
})
478479

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+
479523
it('cancels legacy workflow-tagged runs only after payload verification', async () => {
480524
mockList.mockReturnValueOnce(createListPage([])).mockReturnValueOnce(
481525
createListPage([

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

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -551,16 +551,20 @@ export class TriggerDevJobQueue implements JobQueueBackend {
551551

552552
try {
553553
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
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)
564568
}
565569
}
566570

0 commit comments

Comments
 (0)