Skip to content

Commit a79bd57

Browse files
committed
fix(webapp): clear idempotency key and re-trigger dead runs in batchTrigger
1 parent d8c3530 commit a79bd57

7 files changed

Lines changed: 374 additions & 4 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
Re-trigger runs in batchTrigger when an idempotency key matches a failed or dead run.

‎apps/webapp/app/v3/services/batchTriggerV3.server.ts‎

Lines changed: 20 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ import type {
44
IOPacket,
55
} from "@trigger.dev/core/v3";
66
import { packetRequiresOffloading, parsePacket } from "@trigger.dev/core/v3";
7-
import type { BatchTaskRun, TaskRunAttempt } from "@trigger.dev/database";
7+
import type { BatchTaskRun, TaskRunAttempt, TaskRunStatus } from "@trigger.dev/database";
88
import { isUniqueConstraintError, Prisma } from "@trigger.dev/database";
99
import type { RunStore } from "@internal/run-store";
1010
import pMap from "p-map";
@@ -27,7 +27,11 @@ import { mintBatchFriendlyId } from "~/v3/runOpsMigration/mintBatchFriendlyId.se
2727
import { batchTriggerWorker } from "../batchTriggerWorker.server";
2828
import { guardQueueSizeLimitsForEnv } from "../queueSizeLimits.server";
2929
import { downloadPacketFromObjectStore, uploadPacketToObjectStore } from "../objectStore.server";
30-
import { isFinalAttemptStatus, isFinalRunStatus } from "../taskStatus";
30+
import {
31+
isFinalAttemptStatus,
32+
isFinalRunStatus,
33+
shouldIdempotencyKeyBeCleared,
34+
} from "../taskStatus";
3135
import { startActiveSpan } from "../tracer.server";
3236
import { BaseService, ServiceValidationError } from "./baseService.server";
3337
import { OutOfEntitlementError, TriggerTaskService } from "./triggerTask.server";
@@ -377,6 +381,14 @@ export class BatchTriggerV3Service extends BaseService {
377381
return mintFriendlyIdForKind({ ...target, region });
378382
}
379383

384+
private async prepareRunData(
385+
environment: AuthenticatedEnvironment,
386+
body: BatchTriggerTaskV2RequestBody,
387+
batchFriendlyId: string
388+
): Promise<Array<RunItemData>> {
389+
return this.#prepareRunData(environment, body, batchFriendlyId);
390+
}
391+
380392
async #prepareRunData(
381393
environment: AuthenticatedEnvironment,
382394
body: BatchTriggerTaskV2RequestBody,
@@ -450,7 +462,12 @@ export class BatchTriggerV3Service extends BaseService {
450462
);
451463

452464
if (cachedRun) {
453-
if (cachedRun.idempotencyKeyExpiresAt && cachedRun.idempotencyKeyExpiresAt < new Date()) {
465+
const isExpired =
466+
cachedRun.idempotencyKeyExpiresAt && cachedRun.idempotencyKeyExpiresAt < new Date();
467+
const shouldClear =
468+
isExpired || shouldIdempotencyKeyBeCleared(cachedRun.status as TaskRunStatus);
469+
470+
if (shouldClear) {
454471
expiredRunIds.add(cachedRun.friendlyId);
455472

456473
return {
Lines changed: 339 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,339 @@
1+
import { describe, expect, it, vi } from "vitest";
2+
3+
// Mock DB layer singletons
4+
vi.mock("~/db.server", () => ({
5+
prisma: {},
6+
$replica: {},
7+
runOpsNewPrisma: {},
8+
runOpsLegacyPrisma: {},
9+
runOpsNewReplica: {},
10+
runOpsLegacyReplica: {},
11+
}));
12+
13+
import type { TaskRunStatus } from "@trigger.dev/database";
14+
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
15+
import { BatchTriggerV3Service } from "~/v3/services/batchTriggerV3.server";
16+
import { shouldIdempotencyKeyBeCleared } from "~/v3/taskStatus";
17+
18+
vi.setConfig({ testTimeout: 60_000 });
19+
20+
function fakeEnv(): AuthenticatedEnvironment {
21+
return {
22+
id: "env_test_123",
23+
organizationId: "org_test_123",
24+
organization: { featureFlags: {} },
25+
type: "DEVELOPMENT",
26+
} as unknown as AuthenticatedEnvironment;
27+
}
28+
29+
describe("shouldIdempotencyKeyBeCleared (unit)", () => {
30+
const CLEARABLE_STATUSES: TaskRunStatus[] = [
31+
"CRASHED",
32+
"SYSTEM_FAILURE",
33+
"TIMED_OUT",
34+
"EXPIRED",
35+
"COMPLETED_WITH_ERRORS",
36+
"INTERRUPTED",
37+
];
38+
39+
const NON_CLEARABLE_STATUSES: TaskRunStatus[] = [
40+
"COMPLETED_SUCCESSFULLY",
41+
"EXECUTING",
42+
"PENDING",
43+
"WAITING_FOR_DEPLOY",
44+
"PAUSED",
45+
"DELAYED",
46+
"CANCELED",
47+
];
48+
49+
for (const status of CLEARABLE_STATUSES) {
50+
it(`returns true for clearable failure status: ${status}`, () => {
51+
expect(shouldIdempotencyKeyBeCleared(status)).toBe(true);
52+
});
53+
}
54+
55+
for (const status of NON_CLEARABLE_STATUSES) {
56+
it(`returns false for non-clearable status: ${status}`, () => {
57+
expect(shouldIdempotencyKeyBeCleared(status)).toBe(false);
58+
});
59+
}
60+
});
61+
62+
describe("BatchTriggerV3Service #prepareRunData idempotency key status check", () => {
63+
const CLEARABLE_FAILURE_STATUSES: TaskRunStatus[] = [
64+
"CRASHED",
65+
"SYSTEM_FAILURE",
66+
"TIMED_OUT",
67+
"EXPIRED",
68+
"COMPLETED_WITH_ERRORS",
69+
"INTERRUPTED",
70+
];
71+
72+
const NON_CLEARABLE_STATUSES: TaskRunStatus[] = [
73+
"COMPLETED_SUCCESSFULLY",
74+
"EXECUTING",
75+
"PENDING",
76+
"WAITING_FOR_DEPLOY",
77+
];
78+
79+
for (const status of CLEARABLE_FAILURE_STATUSES) {
80+
it(`clears idempotency key and mints fresh run for clearable status: ${status}`, async () => {
81+
const clearIdempotencyKeyMock = vi.fn().mockResolvedValue({ count: 1 });
82+
const mockRunStore = {
83+
findRunsByIdempotencyKeys: vi.fn().mockResolvedValue([
84+
{
85+
id: "run_internal_dead",
86+
createdAt: new Date(),
87+
friendlyId: "run_dead_123",
88+
idempotencyKey: "key_dead",
89+
idempotencyKeyExpiresAt: null,
90+
status,
91+
},
92+
]),
93+
clearIdempotencyKey: clearIdempotencyKeyMock,
94+
};
95+
96+
const service = new BatchTriggerV3Service(
97+
undefined,
98+
undefined,
99+
{} as any,
100+
mockRunStore as any,
101+
async () => "cuid"
102+
);
103+
104+
const body = {
105+
items: [
106+
{
107+
task: "test-task",
108+
payload: "{}",
109+
options: {
110+
idempotencyKey: "key_dead",
111+
},
112+
},
113+
],
114+
};
115+
116+
const runs = await (service as any).prepareRunData(fakeEnv(), body, "batch_123");
117+
118+
// Verify clearIdempotencyKey was called for the dead run
119+
expect(clearIdempotencyKeyMock).toHaveBeenCalledTimes(1);
120+
expect(clearIdempotencyKeyMock).toHaveBeenCalledWith(
121+
{ byFriendlyIds: ["run_dead_123"] },
122+
expect.anything()
123+
);
124+
125+
// Verify the run returned is NOT cached and has a freshly minted ID
126+
expect(runs).toHaveLength(1);
127+
expect(runs[0].isCached).toBe(false);
128+
expect(runs[0].id).not.toBe("run_dead_123");
129+
expect(runs[0].taskIdentifier).toBe("test-task");
130+
expect(runs[0].idempotencyKey).toBe("key_dead");
131+
});
132+
}
133+
134+
for (const status of NON_CLEARABLE_STATUSES) {
135+
it(`reuses cached run and does NOT clear key for status: ${status}`, async () => {
136+
const clearIdempotencyKeyMock = vi.fn().mockResolvedValue({ count: 0 });
137+
const mockRunStore = {
138+
findRunsByIdempotencyKeys: vi.fn().mockResolvedValue([
139+
{
140+
id: "run_internal_live",
141+
createdAt: new Date(),
142+
friendlyId: "run_live_123",
143+
idempotencyKey: "key_live",
144+
idempotencyKeyExpiresAt: null,
145+
status,
146+
},
147+
]),
148+
clearIdempotencyKey: clearIdempotencyKeyMock,
149+
};
150+
151+
const service = new BatchTriggerV3Service(
152+
undefined,
153+
undefined,
154+
{} as any,
155+
mockRunStore as any,
156+
async () => "cuid"
157+
);
158+
159+
const body = {
160+
items: [
161+
{
162+
task: "test-task",
163+
payload: "{}",
164+
options: {
165+
idempotencyKey: "key_live",
166+
},
167+
},
168+
],
169+
};
170+
171+
const runs = await (service as any).prepareRunData(fakeEnv(), body, "batch_123");
172+
173+
// Verify clearIdempotencyKey was NOT called
174+
expect(clearIdempotencyKeyMock).not.toHaveBeenCalled();
175+
176+
// Verify the run returned IS cached and preserves existing run ID
177+
expect(runs).toHaveLength(1);
178+
expect(runs[0].isCached).toBe(true);
179+
expect(runs[0].id).toBe("run_live_123");
180+
expect(runs[0].taskIdentifier).toBe("test-task");
181+
expect(runs[0].idempotencyKey).toBe("key_live");
182+
});
183+
}
184+
185+
it("handles a mixed batch with fresh keys, live cached runs, and dead runs correctly", async () => {
186+
const clearIdempotencyKeyMock = vi.fn().mockResolvedValue({ count: 7 });
187+
const mockRunStore = {
188+
findRunsByIdempotencyKeys: vi.fn().mockImplementation(async ({ idempotencyKeys }) => {
189+
const matches = [
190+
{
191+
id: "r1",
192+
createdAt: new Date(),
193+
friendlyId: "run_crashed",
194+
idempotencyKey: "k_crashed",
195+
idempotencyKeyExpiresAt: null,
196+
status: "CRASHED",
197+
},
198+
{
199+
id: "r2",
200+
createdAt: new Date(),
201+
friendlyId: "run_success",
202+
idempotencyKey: "k_success",
203+
idempotencyKeyExpiresAt: null,
204+
status: "COMPLETED_SUCCESSFULLY",
205+
},
206+
{
207+
id: "r3",
208+
createdAt: new Date(),
209+
friendlyId: "run_sys_fail",
210+
idempotencyKey: "k_sys_fail",
211+
idempotencyKeyExpiresAt: null,
212+
status: "SYSTEM_FAILURE",
213+
},
214+
{
215+
id: "r4",
216+
createdAt: new Date(),
217+
friendlyId: "run_executing",
218+
idempotencyKey: "k_executing",
219+
idempotencyKeyExpiresAt: null,
220+
status: "EXECUTING",
221+
},
222+
{
223+
id: "r5",
224+
createdAt: new Date(),
225+
friendlyId: "run_timeout",
226+
idempotencyKey: "k_timeout",
227+
idempotencyKeyExpiresAt: null,
228+
status: "TIMED_OUT",
229+
},
230+
{
231+
id: "r6",
232+
createdAt: new Date(),
233+
friendlyId: "run_errors",
234+
idempotencyKey: "k_errors",
235+
idempotencyKeyExpiresAt: null,
236+
status: "COMPLETED_WITH_ERRORS",
237+
},
238+
{
239+
id: "r7",
240+
createdAt: new Date(),
241+
friendlyId: "run_interrupted",
242+
idempotencyKey: "k_interrupted",
243+
idempotencyKeyExpiresAt: null,
244+
status: "INTERRUPTED",
245+
},
246+
{
247+
id: "r8",
248+
createdAt: new Date(),
249+
friendlyId: "run_expired_status",
250+
idempotencyKey: "k_expired_status",
251+
idempotencyKeyExpiresAt: null,
252+
status: "EXPIRED",
253+
},
254+
{
255+
id: "r9",
256+
createdAt: new Date(),
257+
friendlyId: "run_expired_ttl",
258+
idempotencyKey: "k_expired_ttl",
259+
idempotencyKeyExpiresAt: new Date(Date.now() - 10_000),
260+
status: "PENDING",
261+
},
262+
];
263+
return matches.filter((m) => idempotencyKeys.includes(m.idempotencyKey));
264+
}),
265+
clearIdempotencyKey: clearIdempotencyKeyMock,
266+
};
267+
268+
const service = new BatchTriggerV3Service(
269+
undefined,
270+
undefined,
271+
{} as any,
272+
mockRunStore as any,
273+
async () => "cuid"
274+
);
275+
276+
const body = {
277+
items: [
278+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_crashed" } },
279+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_success" } },
280+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_sys_fail" } },
281+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_executing" } },
282+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_timeout" } },
283+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_errors" } },
284+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_interrupted" } },
285+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_expired_status" } },
286+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_expired_ttl" } },
287+
{ task: "t1", payload: "{}", options: { idempotencyKey: "k_brand_new" } },
288+
],
289+
};
290+
291+
const runs = await (service as any).prepareRunData(fakeEnv(), body, "batch_mixed_123");
292+
293+
// All clearable failure statuses + expired TTL run must be cleared
294+
expect(clearIdempotencyKeyMock).toHaveBeenCalledTimes(1);
295+
const clearedIds = clearIdempotencyKeyMock.mock.calls[0][0].byFriendlyIds.sort();
296+
expect(clearedIds).toEqual(
297+
[
298+
"run_crashed",
299+
"run_sys_fail",
300+
"run_timeout",
301+
"run_errors",
302+
"run_interrupted",
303+
"run_expired_status",
304+
"run_expired_ttl",
305+
].sort()
306+
);
307+
308+
expect(runs).toHaveLength(10);
309+
// k_crashed: cleared, minted new
310+
expect(runs[0].isCached).toBe(false);
311+
expect(runs[0].id).not.toBe("run_crashed");
312+
// k_success: live, cached
313+
expect(runs[1].isCached).toBe(true);
314+
expect(runs[1].id).toBe("run_success");
315+
// k_sys_fail: cleared, minted new
316+
expect(runs[2].isCached).toBe(false);
317+
expect(runs[2].id).not.toBe("run_sys_fail");
318+
// k_executing: live, cached
319+
expect(runs[3].isCached).toBe(true);
320+
expect(runs[3].id).toBe("run_executing");
321+
// k_timeout: cleared, minted new
322+
expect(runs[4].isCached).toBe(false);
323+
expect(runs[4].id).not.toBe("run_timeout");
324+
// k_errors: cleared, minted new
325+
expect(runs[5].isCached).toBe(false);
326+
expect(runs[5].id).not.toBe("run_errors");
327+
// k_interrupted: cleared, minted new
328+
expect(runs[6].isCached).toBe(false);
329+
expect(runs[6].id).not.toBe("run_interrupted");
330+
// k_expired_status: cleared, minted new
331+
expect(runs[7].isCached).toBe(false);
332+
expect(runs[7].id).not.toBe("run_expired_status");
333+
// k_expired_ttl: cleared, minted new
334+
expect(runs[8].isCached).toBe(false);
335+
expect(runs[8].id).not.toBe("run_expired_ttl");
336+
// k_brand_new: fresh, minted new
337+
expect(runs[9].isCached).toBe(false);
338+
});
339+
});

0 commit comments

Comments
 (0)