Skip to content

Commit 4e9273e

Browse files
committed
fix(run-engine,webapp): resolve dequeue worker version fresh per task
1 parent 8dc8e1b commit 4e9273e

9 files changed

Lines changed: 627 additions & 61 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+
Fixed a brief window after promoting or rolling back a deployment where newly triggered runs could still execute on the previous version. New runs now pick up the current version immediately.

apps/webapp/app/env.server.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,7 @@ const EnvironmentSchema = z
179179
// Explicit positive opt-in. Split behavior is unreachable unless this is true
180180
// AND the distinct-DB sentinel confirms the two URLs are physically distinct DBs.
181181
RUN_OPS_SPLIT_ENABLED: BoolEnv.default(false),
182+
RUN_OPS_WORKER_VERSION_FRESH_READ_ENABLED: BoolEnv.default(true),
182183
// Canonical connection URL for the dedicated NEW run-ops DB — drives the runtime pool, the split
183184
// decision, replication, and migrations. Optional so single-DB installs never set it.
184185
RUN_OPS_DATABASE_URL: z

apps/webapp/app/v3/runOpsMigration/controlPlaneCache.server.ts

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -228,7 +228,6 @@ export class ControlPlaneCache {
228228
this.#bump(`env:${id}`);
229229
}
230230

231-
// worker version: key = `${environmentId}:${backgroundWorkerId ?? "current"}`
232231
getWorkerVersion(key: string): (ResolvedWorkerVersion | null) | undefined {
233232
return this.#read(this.#version, `version:${key}`);
234233
}
Lines changed: 320 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,320 @@
1+
/**
2+
* TRI-13291: the dequeue worker-version resolve reads the currently-promoted worker fresh on every
3+
* call. The previous env-keyed TTL cache served the superseded worker for up to the TTL after a
4+
* deployment promotion / dev re-register (nothing invalidated it); it has been removed, so a
5+
* promotion or re-register is reflected on the very next resolve. The DB is never mocked: every
6+
* query runs against the real Postgres container.
7+
*/
8+
import { postgresTest } from "@internal/testcontainers";
9+
import { describe, expect } from "vitest";
10+
import type { PrismaClient, PrismaReplicaClient } from "@trigger.dev/database";
11+
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
12+
import { ControlPlaneCache } from "./controlPlaneCache.server";
13+
import { ControlPlaneResolver } from "./controlPlaneResolver.server";
14+
15+
let n = 0;
16+
17+
async function seedEnv(prisma: PrismaClient, type: "PRODUCTION" | "DEVELOPMENT" = "PRODUCTION") {
18+
const s = n++;
19+
const organization = await prisma.organization.create({
20+
data: { title: `Org ${s}`, slug: `org-${s}` },
21+
});
22+
const project = await prisma.project.create({
23+
data: {
24+
name: `P ${s}`,
25+
slug: `p-${s}`,
26+
externalRef: `proj_${s}`,
27+
organizationId: organization.id,
28+
},
29+
});
30+
const environment = await prisma.runtimeEnvironment.create({
31+
data: {
32+
type,
33+
slug: `env-${s}`,
34+
projectId: project.id,
35+
organizationId: organization.id,
36+
apiKey: `tr_${s}`,
37+
pkApiKey: `pk_${s}`,
38+
shortcode: `sc_${s}`,
39+
},
40+
});
41+
return { organization, project, environment };
42+
}
43+
44+
async function seedWorkerWithTask(
45+
prisma: PrismaClient,
46+
ctx: { projectId: string; runtimeEnvironmentId: string },
47+
version: string,
48+
taskSlug: string
49+
) {
50+
const s = n++;
51+
const worker = await prisma.backgroundWorker.create({
52+
data: {
53+
friendlyId: `worker_${s}`,
54+
version,
55+
contentHash: `hash_${s}`,
56+
projectId: ctx.projectId,
57+
runtimeEnvironmentId: ctx.runtimeEnvironmentId,
58+
metadata: {},
59+
},
60+
});
61+
await prisma.backgroundWorkerTask.create({
62+
data: {
63+
friendlyId: `task_${s}`,
64+
slug: taskSlug,
65+
filePath: `src/${taskSlug}.ts`,
66+
workerId: worker.id,
67+
projectId: ctx.projectId,
68+
runtimeEnvironmentId: ctx.runtimeEnvironmentId,
69+
},
70+
});
71+
return worker;
72+
}
73+
74+
async function seedManagedDeployment(
75+
prisma: PrismaClient,
76+
ctx: { projectId: string; environmentId: string; workerId: string },
77+
version: string
78+
) {
79+
const s = n++;
80+
return prisma.workerDeployment.create({
81+
data: {
82+
friendlyId: `deploy_${s}`,
83+
shortCode: `dep${s}`,
84+
contentHash: `dhash_${s}`,
85+
version,
86+
type: "MANAGED",
87+
status: "DEPLOYED",
88+
projectId: ctx.projectId,
89+
environmentId: ctx.environmentId,
90+
workerId: ctx.workerId,
91+
},
92+
});
93+
}
94+
95+
async function promote(prisma: PrismaClient, environmentId: string, deploymentId: string) {
96+
await prisma.workerDeploymentPromotion.upsert({
97+
where: { environmentId_label: { environmentId, label: CURRENT_DEPLOYMENT_LABEL } },
98+
create: { deploymentId, environmentId, label: CURRENT_DEPLOYMENT_LABEL },
99+
update: { deploymentId },
100+
});
101+
}
102+
103+
function makeResolver(prisma: PrismaClient, freshRead = true) {
104+
return new ControlPlaneResolver({
105+
controlPlanePrimary: prisma,
106+
controlPlaneReplica: prisma as unknown as PrismaReplicaClient,
107+
cache: new ControlPlaneCache(),
108+
splitEnabled: () => true,
109+
workerVersionFreshReadEnabled: () => freshRead,
110+
});
111+
}
112+
113+
describe("ControlPlaneResolver worker-version dispatch freshness (TRI-13291)", () => {
114+
postgresTest(
115+
"deployed :current: resolves the newly-promoted worker on the next call (no stale cache)",
116+
async ({ prisma }) => {
117+
const { project, environment } = await seedEnv(prisma);
118+
const ctx = { projectId: project.id, runtimeEnvironmentId: environment.id };
119+
120+
const workerV1 = await seedWorkerWithTask(prisma, ctx, "20240101.1", "task-a");
121+
const workerV2 = await seedWorkerWithTask(prisma, ctx, "20240101.2", "task-a");
122+
const depV1 = await seedManagedDeployment(
123+
prisma,
124+
{ projectId: project.id, environmentId: environment.id, workerId: workerV1.id },
125+
"20240101.1"
126+
);
127+
const depV2 = await seedManagedDeployment(
128+
prisma,
129+
{ projectId: project.id, environmentId: environment.id, workerId: workerV2.id },
130+
"20240101.2"
131+
);
132+
133+
await promote(prisma, environment.id, depV1.id);
134+
135+
const resolver = makeResolver(prisma);
136+
137+
const first = await resolver.resolveWorkerVersion({
138+
environmentId: environment.id,
139+
type: "PRODUCTION",
140+
taskIdentifier: "task-a",
141+
});
142+
expect(first?.worker.id).toBe(workerV1.id);
143+
expect(first?.worker.version).toBe("20240101.1");
144+
145+
await promote(prisma, environment.id, depV2.id);
146+
147+
const afterPromotion = await resolver.resolveWorkerVersion({
148+
environmentId: environment.id,
149+
type: "PRODUCTION",
150+
taskIdentifier: "task-a",
151+
});
152+
expect(afterPromotion?.worker.id).toBe(workerV2.id);
153+
expect(afterPromotion?.worker.version).toBe("20240101.2");
154+
},
155+
30_000
156+
);
157+
158+
postgresTest(
159+
"deployed :current: resolves the rolled-back worker on the next call",
160+
async ({ prisma }) => {
161+
const { project, environment } = await seedEnv(prisma);
162+
const ctx = { projectId: project.id, runtimeEnvironmentId: environment.id };
163+
164+
const workerV1 = await seedWorkerWithTask(prisma, ctx, "20240101.1", "task-a");
165+
const workerV2 = await seedWorkerWithTask(prisma, ctx, "20240101.2", "task-a");
166+
const depV1 = await seedManagedDeployment(
167+
prisma,
168+
{ projectId: project.id, environmentId: environment.id, workerId: workerV1.id },
169+
"20240101.1"
170+
);
171+
const depV2 = await seedManagedDeployment(
172+
prisma,
173+
{ projectId: project.id, environmentId: environment.id, workerId: workerV2.id },
174+
"20240101.2"
175+
);
176+
177+
await promote(prisma, environment.id, depV2.id);
178+
const resolver = makeResolver(prisma);
179+
const onV2 = await resolver.resolveWorkerVersion({
180+
environmentId: environment.id,
181+
type: "PRODUCTION",
182+
taskIdentifier: "task-a",
183+
});
184+
expect(onV2?.worker.id).toBe(workerV2.id);
185+
186+
await promote(prisma, environment.id, depV1.id);
187+
const rolledBack = await resolver.resolveWorkerVersion({
188+
environmentId: environment.id,
189+
type: "PRODUCTION",
190+
taskIdentifier: "task-a",
191+
});
192+
expect(rolledBack?.worker.id).toBe(workerV1.id);
193+
},
194+
30_000
195+
);
196+
197+
postgresTest(
198+
"dev :current: resolves the re-registered worker on the next call",
199+
async ({ prisma }) => {
200+
const { project, environment } = await seedEnv(prisma, "DEVELOPMENT");
201+
const ctx = { projectId: project.id, runtimeEnvironmentId: environment.id };
202+
203+
const workerV1 = await seedWorkerWithTask(prisma, ctx, "20240101.1", "task-a");
204+
205+
const resolver = makeResolver(prisma);
206+
const first = await resolver.resolveWorkerVersion({
207+
environmentId: environment.id,
208+
type: "DEVELOPMENT",
209+
taskIdentifier: "task-a",
210+
});
211+
expect(first?.worker.id).toBe(workerV1.id);
212+
213+
const workerV2 = await seedWorkerWithTask(prisma, ctx, "20240101.2", "task-a");
214+
215+
const afterReregister = await resolver.resolveWorkerVersion({
216+
environmentId: environment.id,
217+
type: "DEVELOPMENT",
218+
taskIdentifier: "task-a",
219+
});
220+
expect(afterReregister?.worker.id).toBe(workerV2.id);
221+
},
222+
30_000
223+
);
224+
225+
postgresTest(
226+
"resolves only the matched task and queue (per-slug/per-queue dispatch)",
227+
async ({ prisma }) => {
228+
const { project, environment } = await seedEnv(prisma);
229+
const ctx = { projectId: project.id, runtimeEnvironmentId: environment.id };
230+
231+
const worker = await seedWorkerWithTask(prisma, ctx, "20240101.1", "task-a");
232+
await prisma.backgroundWorkerTask.create({
233+
data: {
234+
friendlyId: `task_extra_${n++}`,
235+
slug: "task-b",
236+
filePath: "src/task-b.ts",
237+
workerId: worker.id,
238+
projectId: project.id,
239+
runtimeEnvironmentId: environment.id,
240+
},
241+
});
242+
const qA = await prisma.taskQueue.create({
243+
data: {
244+
friendlyId: `q_${n++}`,
245+
name: "queue-a",
246+
runtimeEnvironmentId: environment.id,
247+
projectId: project.id,
248+
workers: { connect: { id: worker.id } },
249+
},
250+
});
251+
await prisma.taskQueue.create({
252+
data: {
253+
friendlyId: `q_${n++}`,
254+
name: "queue-b",
255+
runtimeEnvironmentId: environment.id,
256+
projectId: project.id,
257+
workers: { connect: { id: worker.id } },
258+
},
259+
});
260+
const dep = await seedManagedDeployment(
261+
prisma,
262+
{ projectId: project.id, environmentId: environment.id, workerId: worker.id },
263+
"20240101.1"
264+
);
265+
await promote(prisma, environment.id, dep.id);
266+
267+
const resolved = await makeResolver(prisma).resolveWorkerVersion({
268+
environmentId: environment.id,
269+
type: "PRODUCTION",
270+
taskIdentifier: "task-a",
271+
queue: { name: "queue-a" },
272+
});
273+
274+
expect(resolved?.tasks.map((t) => t.slug)).toEqual(["task-a"]);
275+
expect(resolved?.queues.map((q) => q.id)).toEqual([qA.id]);
276+
},
277+
30_000
278+
);
279+
280+
postgresTest(
281+
"kill-switch off falls back to the legacy env-keyed cache (serves the pre-promotion worker)",
282+
async ({ prisma }) => {
283+
const { project, environment } = await seedEnv(prisma);
284+
const ctx = { projectId: project.id, runtimeEnvironmentId: environment.id };
285+
286+
const workerV1 = await seedWorkerWithTask(prisma, ctx, "20240101.1", "task-a");
287+
const workerV2 = await seedWorkerWithTask(prisma, ctx, "20240101.2", "task-a");
288+
const depV1 = await seedManagedDeployment(
289+
prisma,
290+
{ projectId: project.id, environmentId: environment.id, workerId: workerV1.id },
291+
"20240101.1"
292+
);
293+
const depV2 = await seedManagedDeployment(
294+
prisma,
295+
{ projectId: project.id, environmentId: environment.id, workerId: workerV2.id },
296+
"20240101.2"
297+
);
298+
299+
await promote(prisma, environment.id, depV1.id);
300+
301+
const resolver = makeResolver(prisma, false);
302+
const first = await resolver.resolveWorkerVersion({
303+
environmentId: environment.id,
304+
type: "PRODUCTION",
305+
taskIdentifier: "task-a",
306+
});
307+
expect(first?.worker.id).toBe(workerV1.id);
308+
309+
await promote(prisma, environment.id, depV2.id);
310+
311+
const afterPromotion = await resolver.resolveWorkerVersion({
312+
environmentId: environment.id,
313+
type: "PRODUCTION",
314+
taskIdentifier: "task-a",
315+
});
316+
expect(afterPromotion?.worker.id).toBe(workerV1.id);
317+
},
318+
30_000
319+
);
320+
});

0 commit comments

Comments
 (0)