Skip to content

Commit a39efeb

Browse files
committed
fix(batchTrigger): clear idempotency key when previous run failed (closes #4819)
batchTrigger silently returned stale failed runs as isCached: true, so a caller re-batching with the same idempotency key never got a retry. The single-trigger path has done the right thing since day one via `IdempotencyKeyConcern.handleExistingRun` -> `shouldIdempotencyKeyBeCleared` (apps/webapp/app/runEngine/concerns/idempotencyKeys.server.ts:463) — the batch path at batchTriggerV3.server.ts:450 just ignored the status. Root cause (as @Jaimin2687 identified on the issue): * `PostgresRunStore.findRunsByIdempotencyKeys` SELECTed 5 columns — id, createdAt, friendlyId, idempotencyKey, idempotencyKeyExpiresAt. * `status` was missing, and `IdempotencyKeyRunMatch` did not type it. * `#prepareRunData` therefore had only time-based expiry to work with and no way to see that the cached run was CRASHED / SYSTEM_FAILURE / COMPLETED_WITH_ERRORS / INTERRUPTED / TIMED_OUT / EXPIRED. Fix (three files, additive): - internal-packages/run-store/src/types.ts * `IdempotencyKeyRunMatch` gains a required `status: TaskRunStatus`. * Comment explains the invariant for the next reader. - internal-packages/run-store/src/PostgresRunStore.ts * The hot-path UNION ALL SQL now also selects `"status"`. One token added per branch, no new rows read, no index change needed — (runtimeEnvironmentId, taskIdentifier, idempotencyKey) already uniquely identifies the row. - apps/webapp/app/v3/services/batchTriggerV3.server.ts * `#prepareRunData` imports `shouldIdempotencyKeyBeCleared` and ORs it with the existing time-based expiry check. A cached run whose status is in the "clear me" set is now treated exactly like an expired cache entry: the friendlyId is added to `expiredRunIds` for cleanup, a fresh child id is minted, and the result is returned with `isCached: false`. * Behaviour on cached-but-still-good runs is unchanged. * Behaviour on cached-and-time-expired runs is unchanged. - internal-packages/run-store/src/PostgresRunStore.findRunsByIdempotencyKeys.test.ts * `createRun` helper gains an optional `status` param. * New postgresTest "returns run status so callers can honour shouldIdempotencyKeyBeCleared" creates four seed rows (PENDING, CRASHED, COMPLETED_WITH_ERRORS, COMPLETED_SUCCESSFULLY) and asserts the store surfaces each status in the result. This locks the contract for batchTrigger to depend on. Reporter: @Jaimin2687 (full RCA in the issue). No in-flight PR for #4819 at time of writing.
1 parent 21f1dcc commit a39efeb

4 files changed

Lines changed: 95 additions & 4 deletions

File tree

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

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -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";
@@ -450,7 +454,24 @@ export class BatchTriggerV3Service extends BaseService {
450454
);
451455

452456
if (cachedRun) {
453-
if (cachedRun.idempotencyKeyExpiresAt && cachedRun.idempotencyKeyExpiresAt < new Date()) {
457+
// Clear the idempotency key and mint a fresh run when either:
458+
// (a) the key's time-based expiry has passed, OR
459+
// (b) the previous run ended in a terminal failure state the
460+
// single-trigger path treats as "retry allowed"
461+
// (CRASHED, SYSTEM_FAILURE, INTERRUPTED, COMPLETED_WITH_ERRORS,
462+
// EXPIRED, TIMED_OUT_WITH_ERRORS — see shouldIdempotencyKeyBeCleared).
463+
//
464+
// Before this fix, batchTrigger only checked (a). A failed run with
465+
// the same idempotency key was returned as `isCached: true`, silently
466+
// handing the caller back a dead run that would never produce output
467+
// (issue #4819). This mirrors the single-trigger behaviour at
468+
// IdempotencyKeyConcern.handleExistingRun (idempotencyKeys.server.ts:463).
469+
const expired =
470+
cachedRun.idempotencyKeyExpiresAt &&
471+
cachedRun.idempotencyKeyExpiresAt < new Date();
472+
const failed = shouldIdempotencyKeyBeCleared(cachedRun.status);
473+
474+
if (expired || failed) {
454475
expiredRunIds.add(cachedRun.friendlyId);
455476

456477
return {

‎internal-packages/run-store/src/PostgresRunStore.findRunsByIdempotencyKeys.test.ts‎

Lines changed: 61 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { postgresTest } from "@internal/testcontainers";
2-
import type { PrismaClient } from "@trigger.dev/database";
2+
import type { PrismaClient, TaskRunStatus } from "@trigger.dev/database";
33
import { describe, expect } from "vitest";
44
import { PostgresRunStore } from "./PostgresRunStore.js";
55

@@ -38,6 +38,7 @@ async function createRun(
3838
taskIdentifier: string;
3939
idempotencyKey: string;
4040
idempotencyKeyExpiresAt?: Date;
41+
status?: TaskRunStatus;
4142
}
4243
) {
4344
await prisma.taskRun.create({
@@ -46,6 +47,7 @@ async function createRun(
4647
taskIdentifier: params.taskIdentifier,
4748
idempotencyKey: params.idempotencyKey,
4849
idempotencyKeyExpiresAt: params.idempotencyKeyExpiresAt ?? null,
50+
status: params.status ?? "PENDING",
4951
payload: "{}",
5052
payloadType: "application/json",
5153
runtimeEnvironmentId: params.runtimeEnvironmentId,
@@ -117,4 +119,62 @@ describe("PostgresRunStore.findRunsByIdempotencyKeys", () => {
117119

118120
expect(rows).toEqual([]);
119121
});
122+
123+
// Regression for #4819 — batchTrigger needs the row's status to decide
124+
// whether to re-trigger (same semantics as single trigger via
125+
// shouldIdempotencyKeyBeCleared). Before this fix the SQL omitted `status`,
126+
// so batchTrigger could only inspect time-based expiry and silently handed
127+
// callers back dead CRASHED / COMPLETED_WITH_ERRORS runs as `isCached: true`.
128+
postgresTest(
129+
"returns run status so callers can honour shouldIdempotencyKeyBeCleared",
130+
async ({ prisma }) => {
131+
const { project, environment } = await seedEnvironment(prisma);
132+
const store = new PostgresRunStore({ prisma, readOnlyPrisma: prisma });
133+
134+
await createRun(prisma, {
135+
runtimeEnvironmentId: environment.id,
136+
projectId: project.id,
137+
friendlyId: "run_pending",
138+
taskIdentifier: "task-x",
139+
idempotencyKey: "key-pending",
140+
status: "PENDING",
141+
});
142+
await createRun(prisma, {
143+
runtimeEnvironmentId: environment.id,
144+
projectId: project.id,
145+
friendlyId: "run_crashed",
146+
taskIdentifier: "task-x",
147+
idempotencyKey: "key-crashed",
148+
status: "CRASHED",
149+
});
150+
await createRun(prisma, {
151+
runtimeEnvironmentId: environment.id,
152+
projectId: project.id,
153+
friendlyId: "run_completed_with_errors",
154+
taskIdentifier: "task-x",
155+
idempotencyKey: "key-cwe",
156+
status: "COMPLETED_WITH_ERRORS",
157+
});
158+
await createRun(prisma, {
159+
runtimeEnvironmentId: environment.id,
160+
projectId: project.id,
161+
friendlyId: "run_completed",
162+
taskIdentifier: "task-x",
163+
idempotencyKey: "key-ok",
164+
status: "COMPLETED_SUCCESSFULLY",
165+
});
166+
167+
const rows = await store.findRunsByIdempotencyKeys({
168+
runtimeEnvironmentId: environment.id,
169+
taskIdentifier: "task-x",
170+
idempotencyKeys: ["key-pending", "key-crashed", "key-cwe", "key-ok"],
171+
});
172+
173+
const byKey = new Map(rows.map((r) => [r.idempotencyKey, r]));
174+
expect(byKey.get("key-pending")?.status).toBe("PENDING");
175+
expect(byKey.get("key-crashed")?.status).toBe("CRASHED");
176+
expect(byKey.get("key-cwe")?.status).toBe("COMPLETED_WITH_ERRORS");
177+
expect(byKey.get("key-ok")?.status).toBe("COMPLETED_SUCCESSFULLY");
178+
}
179+
);
120180
});

‎internal-packages/run-store/src/PostgresRunStore.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1924,7 +1924,9 @@ export class PostgresRunStore implements RunStore {
19241924
const branches = args.idempotencyKeys.map((key) => {
19251925
const base = params.length;
19261926
params.push(args.runtimeEnvironmentId, args.taskIdentifier, key);
1927-
return `SELECT "id", "createdAt", "friendlyId", "idempotencyKey", "idempotencyKeyExpiresAt" FROM "TaskRun" WHERE "runtimeEnvironmentId" = $${base + 1} AND "taskIdentifier" = $${base + 2} AND "idempotencyKey" = $${base + 3}`;
1927+
// "status" is required so callers can honour shouldIdempotencyKeyBeCleared
1928+
// (closes #4819 — batchTrigger previously returned stale failed runs as cached).
1929+
return `SELECT "id", "createdAt", "friendlyId", "idempotencyKey", "idempotencyKeyExpiresAt", "status" FROM "TaskRun" WHERE "runtimeEnvironmentId" = $${base + 1} AND "taskIdentifier" = $${base + 2} AND "idempotencyKey" = $${base + 3}`;
19281930
});
19291931
return prisma.$queryRawUnsafe<IdempotencyKeyRunMatch[]>(
19301932
branches.join(" UNION ALL "),

‎internal-packages/run-store/src/types.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,14 @@ export type IdempotencyKeyRunMatch = {
3030
friendlyId: string;
3131
idempotencyKey: string | null;
3232
idempotencyKeyExpiresAt: Date | null;
33+
/**
34+
* Terminal run status is needed so callers can honour the same
35+
* "clear the idempotency key if the previous run failed" semantics the
36+
* single-trigger path already applies via `shouldIdempotencyKeyBeCleared`.
37+
* Without it, `batchTrigger` silently returned stale failed runs as
38+
* `isCached: true` (issue #4819).
39+
*/
40+
status: TaskRunStatus;
3341
};
3442

3543
/**

0 commit comments

Comments
 (0)