Skip to content

Commit d93f848

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
perf(webapp): avoid task registration conflict errors
Mono-RevId: fc08d78116292c16f3f9fc2145a2f6ce74bb11d0
1 parent fdbc2e3 commit d93f848

2 files changed

Lines changed: 148 additions & 59 deletions

File tree

‎apps/webapp/app/v3/services/createBackgroundWorker.server.test.ts‎

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,13 @@
11
import { ScheduleEngine } from "@internal/schedule-engine";
22
import { containerTest } from "@internal/testcontainers";
3+
import type { BackgroundWorkerMetadata } from "@trigger.dev/core/v3";
34
import type { PrismaClient } from "@trigger.dev/database";
45
import { describe, expect, vi } from "vitest";
56
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
67
import { FEATURE_FLAG } from "~/v3/featureFlags";
78
import {
89
CreateBackgroundWorkerService,
10+
createWorkerResources,
911
syncDeclarativeSchedules,
1012
} from "~/v3/services/createBackgroundWorker.server";
1113

@@ -123,6 +125,22 @@ function declarativeTasks(schedule: { cron: string; timezone: string; window?: s
123125
return [{ id: "my-task", schedule }] as TasksArg;
124126
}
125127

128+
function workerMetadata(description: string) {
129+
return {
130+
contentHash: "duplicate-task-content",
131+
tasks: [
132+
{
133+
id: "duplicate-task",
134+
description,
135+
filePath: "src/trigger/duplicate-task.ts",
136+
exportName: "duplicateTask",
137+
queue: { name: "duplicate-task-queue" },
138+
},
139+
],
140+
queues: [{ name: "duplicate-task-queue" }],
141+
} as unknown as BackgroundWorkerMetadata;
142+
}
143+
126144
async function seedScheduledTask(
127145
prisma: PrismaClient,
128146
projectId: string,
@@ -152,6 +170,82 @@ async function seedScheduledTask(
152170
});
153171
}
154172

173+
describe("worker task creation", () => {
174+
containerTest(
175+
"preserves one task when concurrent registration transactions target the same worker and slug",
176+
async ({ prisma }) => {
177+
const { project, devEnv } = await seedProjectWithEnvs(prisma);
178+
const worker = await prisma.backgroundWorker.create({
179+
data: {
180+
friendlyId: `worker_${devEnv.id}`,
181+
contentHash: "duplicate-task-content",
182+
version: "20260811.1",
183+
metadata: {},
184+
projectId: project.id,
185+
runtimeEnvironmentId: devEnv.id,
186+
},
187+
});
188+
const queue = await prisma.taskQueue.create({
189+
data: {
190+
friendlyId: `queue_${devEnv.id}`,
191+
name: "duplicate-task-queue",
192+
type: "NAMED",
193+
version: "V2",
194+
paused: true,
195+
projectId: project.id,
196+
runtimeEnvironmentId: devEnv.id,
197+
},
198+
});
199+
const environment = { ...asEnv(devEnv), project } as AuthenticatedEnvironment;
200+
201+
const entries = await Promise.all(
202+
["first registration", "second registration"].map((description) =>
203+
prisma.$transaction((tx) =>
204+
createWorkerResources(workerMetadata(description), worker, environment, tx)
205+
)
206+
)
207+
);
208+
209+
expect(entries).toEqual([
210+
[
211+
{
212+
slug: "duplicate-task",
213+
ttl: null,
214+
triggerSource: "STANDARD",
215+
queueId: queue.id,
216+
queueName: queue.name,
217+
},
218+
],
219+
[
220+
{
221+
slug: "duplicate-task",
222+
ttl: null,
223+
triggerSource: "STANDARD",
224+
queueId: queue.id,
225+
queueName: queue.name,
226+
},
227+
],
228+
]);
229+
230+
const created = await prisma.backgroundWorkerTask.findUniqueOrThrow({
231+
where: { workerId_slug: { workerId: worker.id, slug: "duplicate-task" } },
232+
});
233+
expect(["first registration", "second registration"]).toContain(created.description);
234+
expect(await prisma.backgroundWorkerTask.count({ where: { workerId: worker.id } })).toBe(1);
235+
236+
await prisma.$transaction((tx) =>
237+
createWorkerResources(workerMetadata("replacement"), worker, environment, tx)
238+
);
239+
240+
const afterRetry = await prisma.backgroundWorkerTask.findUniqueOrThrow({
241+
where: { workerId_slug: { workerId: worker.id, slug: "duplicate-task" } },
242+
});
243+
expect(afterRetry.id).toBe(created.id);
244+
expect(afterRetry.description).toBe(created.description);
245+
}
246+
);
247+
});
248+
155249
describe("declarative schedule preflight", () => {
156250
containerTest("rejects identical retries before persisting a worker", async ({ prisma }) => {
157251
const { project, devEnv } = await seedProjectWithEnvs(prisma);

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

Lines changed: 54 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -410,7 +410,6 @@ async function createWorkerTask(
410410
prisma: PrismaClientOrTransaction,
411411
tasksToBackgroundFiles?: Map<string, string>
412412
): Promise<TaskMetadataEntry | null> {
413-
// Hoisted so the P2002 catch branch can return the same entry shape.
414413
let queue: TaskQueue | undefined;
415414
let resolvedTriggerSource: "SCHEDULED" | "AGENT" | "WEBHOOK" | "STANDARD" | undefined;
416415
let resolvedTtl: string | null | undefined;
@@ -445,28 +444,51 @@ async function createWorkerTask(
445444
resolvedTtl =
446445
typeof task.ttl === "number" ? (stringifyDuration(task.ttl) ?? null) : (task.ttl ?? null);
447446

448-
await prisma.backgroundWorkerTask.create({
449-
data: {
450-
friendlyId: generateFriendlyId("task"),
451-
projectId: worker.projectId,
452-
runtimeEnvironmentId: worker.runtimeEnvironmentId,
453-
workerId: worker.id,
454-
slug: task.id,
455-
description: task.description,
456-
filePath: task.filePath,
457-
exportName: task.exportName,
458-
retryConfig: task.retry,
459-
queueConfig: task.queue,
460-
machineConfig: task.machine,
461-
triggerSource: resolvedTriggerSource,
462-
config: task.agentConfig ? (task.agentConfig as any) : undefined,
463-
fileId: tasksToBackgroundFiles?.get(task.id) ?? null,
464-
maxDurationInSeconds: task.maxDuration ? clampMaxDuration(task.maxDuration) : null,
465-
ttl: resolvedTtl,
466-
queueId: queue.id,
467-
payloadSchema: task.payloadSchema as any,
468-
},
469-
});
447+
let taskPersisted = false;
448+
for (let attempt = 0; attempt < 3; attempt++) {
449+
const result = await prisma.backgroundWorkerTask.createMany({
450+
data: {
451+
friendlyId: generateFriendlyId("task"),
452+
projectId: worker.projectId,
453+
runtimeEnvironmentId: worker.runtimeEnvironmentId,
454+
workerId: worker.id,
455+
slug: task.id,
456+
description: task.description,
457+
filePath: task.filePath,
458+
exportName: task.exportName,
459+
retryConfig: task.retry,
460+
queueConfig: task.queue,
461+
machineConfig: task.machine,
462+
triggerSource: resolvedTriggerSource,
463+
config: task.agentConfig ? (task.agentConfig as any) : undefined,
464+
fileId: tasksToBackgroundFiles?.get(task.id) ?? null,
465+
maxDurationInSeconds: task.maxDuration ? clampMaxDuration(task.maxDuration) : null,
466+
ttl: resolvedTtl,
467+
queueId: queue.id,
468+
payloadSchema: task.payloadSchema as any,
469+
},
470+
skipDuplicates: true,
471+
});
472+
473+
if (result.count === 1) {
474+
taskPersisted = true;
475+
break;
476+
}
477+
478+
const existing = await prisma.backgroundWorkerTask.findFirst({
479+
where: { workerId: worker.id, slug: task.id },
480+
select: { id: true },
481+
});
482+
483+
if (existing) {
484+
taskPersisted = true;
485+
break;
486+
}
487+
}
488+
489+
if (!taskPersisted) {
490+
throw new Error("Failed to create background worker task after unique constraint conflicts");
491+
}
470492

471493
return {
472494
slug: task.id,
@@ -477,53 +499,26 @@ async function createWorkerTask(
477499
};
478500
} catch (error) {
479501
if (error instanceof Prisma.PrismaClientKnownRequestError) {
480-
// The error code for unique constraint violation in Prisma is P2002
481-
if (error.code === "P2002") {
482-
// Retry landing after the first attempt's row was already written.
483-
const existing = await prisma.backgroundWorkerTask.findFirst({
484-
where: { workerId: worker.id, slug: task.id },
485-
select: { id: true },
486-
});
487-
488-
logger.warn("Attempted to recreate background worker task", {
489-
task,
490-
worker,
491-
});
492-
493-
if (existing && queue && resolvedTriggerSource && resolvedTtl !== undefined) {
494-
return {
495-
slug: task.id,
496-
ttl: resolvedTtl,
497-
triggerSource: resolvedTriggerSource,
498-
queueId: queue.id,
499-
queueName: queue.name,
500-
};
501-
}
502-
} else {
503-
logger.error("Prisma Error creating background worker task", {
504-
error: {
505-
code: error.code,
506-
message: error.message,
507-
},
508-
task,
509-
worker,
510-
});
511-
}
502+
logger.error("Prisma Error creating background worker task", {
503+
error: {
504+
code: error.code,
505+
message: error.message,
506+
},
507+
workerId: worker.id,
508+
});
512509
} else if (error instanceof Error) {
513510
logger.error("Error creating background worker task", {
514511
error: {
515512
name: error.name,
516513
message: error.message,
517514
stack: error.stack,
518515
},
519-
task,
520-
worker,
516+
workerId: worker.id,
521517
});
522518
} else {
523519
logger.error("Unknown error creating background worker task", {
524520
error,
525-
task,
526-
worker,
521+
workerId: worker.id,
527522
});
528523
}
529524
return null;

0 commit comments

Comments
 (0)