Skip to content

Commit a3bd4e2

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
feat: add runtime-aware restore routing and supervisor subscriptions
Adds opt-in versioned worker-queue assignment and routes checkpoint resumes to queues matching the saved runtime, preserving the run's dispatch class, region and channel. Legacy routing remains available, and task retries that start fresh return to the original assignment. Mono-RevId: ffc26bd998a3409f27f5011bb11a8bc2b2919db6
1 parent 35f380f commit a3bd4e2

26 files changed

Lines changed: 1091 additions & 81 deletions

‎apps/supervisor/src/env.test.ts‎

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,104 @@ const base = {
2121
OTEL_EXPORTER_OTLP_ENDPOINT: "http://localhost:4318",
2222
};
2323

24+
describe("worker queue selection", () => {
25+
const ondemand = {
26+
class: "ondemand",
27+
phase: "fresh",
28+
compat: "container",
29+
channel: "stable",
30+
};
31+
const restore = {
32+
class: "ondemand",
33+
phase: "restore",
34+
compat: "container",
35+
channel: "canary",
36+
};
37+
38+
it("keeps legacy default and scheduled selection when subscriptions are absent", () => {
39+
expect(Env.parse(base).TRIGGER_WORKER_QUEUE_CLASS).toBe("default");
40+
expect(
41+
Env.parse({ ...base, TRIGGER_WORKER_QUEUE_CLASS: "scheduled" }).TRIGGER_WORKER_QUEUE_CLASS
42+
).toBe("scheduled");
43+
});
44+
45+
it("parses multiple subscriptions without adding a legacy queue class", () => {
46+
const subscriptions = [{ ...ondemand, weight: 0.25 }, restore];
47+
const parsed = Env.parse({
48+
...base,
49+
TRIGGER_CHECKPOINT_URL: "http://localhost:8089",
50+
TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: JSON.stringify(subscriptions),
51+
});
52+
expect(parsed.TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS).toEqual(subscriptions);
53+
expect(parsed.TRIGGER_WORKER_QUEUE_CLASS).toBeUndefined();
54+
});
55+
56+
it.each(["default", "scheduled"])(
57+
"rejects subscriptions with explicit %s selection",
58+
(queueClass) => {
59+
expect(() =>
60+
Env.parse({
61+
...base,
62+
TRIGGER_WORKER_QUEUE_CLASS: queueClass,
63+
TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: JSON.stringify([ondemand]),
64+
})
65+
).toThrow("mutually exclusive");
66+
}
67+
);
68+
69+
it.each(["not-json", "[]", JSON.stringify([{ ...restore, phase: "unknown" }])])(
70+
"rejects invalid subscriptions: %s",
71+
(subscriptions) => {
72+
expect(() =>
73+
Env.parse({ ...base, TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: subscriptions })
74+
).toThrow();
75+
}
76+
);
77+
78+
it.each([
79+
{ config: {}, subscription: { ...ondemand, compat: "compute" }, error: "compatibility" },
80+
{
81+
config: { COMPUTE_GATEWAY_URL: "http://localhost:8080" },
82+
subscription: ondemand,
83+
error: "compatibility",
84+
},
85+
{ config: {}, subscription: restore, error: "TRIGGER_CHECKPOINT_URL" },
86+
{
87+
config: {
88+
KUBERNETES_FORCE_ENABLED: "true",
89+
KUBERNETES_RUN_CRD_ENABLED: "true",
90+
KUBERNETES_RUNNER_RUNTIME: "microvm",
91+
},
92+
subscription: { ...restore, compat: "compute" },
93+
error: "COMPUTE_GATEWAY_URL",
94+
},
95+
])(
96+
"rejects incompatible or unsupported subscriptions: $error",
97+
({ config, subscription, error }) => {
98+
expect(() =>
99+
Env.parse({
100+
...base,
101+
...config,
102+
TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: JSON.stringify([subscription]),
103+
})
104+
).toThrow(error);
105+
}
106+
);
107+
108+
it("accepts shared fresh work and compute restores on the compute backend", () => {
109+
expect(() =>
110+
Env.parse({
111+
...base,
112+
COMPUTE_GATEWAY_URL: "http://localhost:8080",
113+
TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: JSON.stringify([
114+
{ ...ondemand, compat: "any" },
115+
{ ...restore, compat: "compute" },
116+
]),
117+
})
118+
).not.toThrow();
119+
});
120+
});
121+
24122
describe("Env superRefine - backpressure source awareness", () => {
25123
it("pod-count source can be enabled without a Redis host", () => {
26124
expect(() =>

‎apps/supervisor/src/env.ts‎

Lines changed: 59 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { randomUUID } from "crypto";
22
import { env as stdEnv } from "std-env";
33
import { z } from "zod";
4+
import { WeightedWorkerQueueSubscriptions } from "@trigger.dev/core/v3/runEngineWorker";
45
import {
56
AdditionalEnvVars,
67
BoolEnv,
@@ -56,7 +57,23 @@ export const Env = z
5657
// Which worker-queue class this supervisor fleet serves. "default" pulls the
5758
// region queue (standard/agent runs); "scheduled" pulls the dedicated
5859
// scheduled-lineage queue. Run a separate fleet per class for isolation.
59-
TRIGGER_WORKER_QUEUE_CLASS: z.enum(["default", "scheduled"]).default("default"),
60+
TRIGGER_WORKER_QUEUE_CLASS: z.enum(["default", "scheduled"]).optional(),
61+
// JSON array of weighted v2 subscriptions; omit to retain legacy queue selection.
62+
TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS: z
63+
.string()
64+
.transform((value, ctx) => {
65+
try {
66+
return JSON.parse(value) as unknown;
67+
} catch {
68+
ctx.addIssue({
69+
code: z.ZodIssueCode.custom,
70+
message: "Worker queue subscriptions must be valid JSON",
71+
});
72+
return z.NEVER;
73+
}
74+
})
75+
.pipe(WeightedWorkerQueueSubscriptions)
76+
.optional(),
6077
TRIGGER_DEQUEUE_INTERVAL_MS: z.coerce.number().int().default(250),
6178
TRIGGER_DEQUEUE_IDLE_INTERVAL_MS: z.coerce.number().int().default(1000),
6279
TRIGGER_DEQUEUE_MAX_RUN_COUNT: z.coerce.number().int().default(1),
@@ -326,6 +343,44 @@ export const Env = z
326343
TRIGGER_WIDE_EVENTS_NOISY_ROUTES: BoolEnv.default(false),
327344
})
328345
.superRefine((data, ctx) => {
346+
if (data.TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS) {
347+
if (data.TRIGGER_WORKER_QUEUE_CLASS !== undefined) {
348+
ctx.addIssue({
349+
code: z.ZodIssueCode.custom,
350+
message:
351+
"TRIGGER_WORKER_QUEUE_CLASS and TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS are mutually exclusive",
352+
path: ["TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS"],
353+
});
354+
}
355+
356+
const compatibility =
357+
data.COMPUTE_GATEWAY_URL ||
358+
(data.KUBERNETES_FORCE_ENABLED &&
359+
data.KUBERNETES_RUN_CRD_ENABLED &&
360+
data.KUBERNETES_RUNNER_RUNTIME === "microvm")
361+
? "compute"
362+
: "container";
363+
364+
for (const [index, subscription] of data.TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS.entries()) {
365+
if (subscription.compat !== "any" && subscription.compat !== compatibility) {
366+
ctx.addIssue({
367+
code: z.ZodIssueCode.custom,
368+
message: `Subscription compatibility must be any or ${compatibility} for this backend`,
369+
path: ["TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS", index, "compat"],
370+
});
371+
} else if (
372+
subscription.phase === "restore" &&
373+
!(compatibility === "compute" ? data.COMPUTE_GATEWAY_URL : data.TRIGGER_CHECKPOINT_URL)
374+
) {
375+
ctx.addIssue({
376+
code: z.ZodIssueCode.custom,
377+
message: `Restore subscriptions require ${compatibility === "compute" ? "COMPUTE_GATEWAY_URL" : "TRIGGER_CHECKPOINT_URL"}`,
378+
path: ["TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS", index],
379+
});
380+
}
381+
}
382+
}
383+
329384
if (
330385
data.TRIGGER_DEQUEUE_BACKPRESSURE_POD_COUNT_ENABLED &&
331386
data.TRIGGER_DEQUEUE_BACKPRESSURE_POD_COUNT_RELEASE >=
@@ -404,6 +459,9 @@ export const Env = z
404459
})
405460
.transform((data) => ({
406461
...data,
462+
TRIGGER_WORKER_QUEUE_CLASS:
463+
data.TRIGGER_WORKER_QUEUE_CLASS ??
464+
(data.TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS === undefined ? "default" : undefined),
407465
COMPUTE_TRACE_OTLP_ENDPOINT: data.COMPUTE_TRACE_OTLP_ENDPOINT ?? `${data.TRIGGER_API_URL}/otel`,
408466
}));
409467

‎apps/supervisor/src/index.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -349,6 +349,7 @@ class ManagedSupervisor {
349349
queueConsumerEnabled: env.TRIGGER_DEQUEUE_ENABLED,
350350
maxRunCount: env.TRIGGER_DEQUEUE_MAX_RUN_COUNT,
351351
queueClass: env.TRIGGER_WORKER_QUEUE_CLASS,
352+
subscriptions: env.TRIGGER_WORKER_QUEUE_SUBSCRIPTIONS,
352353
metricsRegistry: register,
353354
scaling: {
354355
strategy: env.TRIGGER_DEQUEUE_SCALING_STRATEGY,

‎apps/webapp/app/env.server.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1562,6 +1562,8 @@ const EnvironmentSchema = z
15621562
// this is "1"; an org set to true splits even when this is "0"). Never
15631563
// applies to DEVELOPMENT environments.
15641564
TRIGGER_WORKER_QUEUE_SCHEDULED_SPLIT_ENABLED: z.string().default("0"),
1565+
// Controls new assignments only; existing v2 runs retain their routing.
1566+
TRIGGER_WORKER_QUEUE_V2_ENABLED: z.string().default("0"),
15651567

15661568
TRIGGER_MOLLIFIER_ENABLED: z.string().default("0"),
15671569
// Separate switch for the drainer (consumer side) so it can be split

‎apps/webapp/app/presenters/v3/SpanPresenter.server.ts‎

Lines changed: 16 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import {
1111
} from "@trigger.dev/core/v3";
1212

1313
import { AttemptId, getMaxDuration, parseTraceparent } from "@trigger.dev/core/v3/isomorphic";
14+
import { isV2WorkerQueue } from "@trigger.dev/core/v3/workers";
1415
import {
1516
runOpsLegacyReplica,
1617
runOpsNewReplica,
@@ -40,6 +41,7 @@ import { buildSyntheticSpanRun } from "~/v3/mollifier/syntheticSpanRun.server";
4041
import { engine } from "~/v3/runEngine.server";
4142
import { runStore } from "~/v3/runStore.server";
4243
import { runTriggeredAt } from "~/v3/runTimestamps";
44+
import { workerRegionRegistry } from "~/v3/workerRegions.server";
4345
import { getTaskEventStoreTableForRun, type TaskEventStoreTable } from "~/v3/taskEventStore.server";
4446
import { isFailedRunStatus, isFinalRunStatus } from "~/v3/taskStatus";
4547
import { BasePresenter } from "./basePresenter.server";
@@ -330,23 +332,30 @@ export class SpanPresenter extends BasePresenter {
330332
let region: { name: string; location: string | null } | null = null;
331333

332334
if (environment.type !== "DEVELOPMENT" && run.engine !== "V1") {
335+
const v2Region = isV2WorkerQueue(run.workerQueue)
336+
? (run.region ?? baseWorkerQueue(run.workerQueue))
337+
: undefined;
338+
// V2 identifies geography, not the executing group. Reuse the registry to
339+
// retain the indexed masterQueue lookup for the region's location metadata.
340+
const masterQueue = v2Region
341+
? (workerRegionRegistry.current()?.find((group) => group.region === v2Region)
342+
?.masterQueue ?? v2Region)
343+
: baseWorkerQueue(run.workerQueue);
333344
const workerGroup = await this._replica.workerInstanceGroup.findFirst({
334345
select: {
335346
name: true,
336347
location: true,
337348
},
338-
where: {
339-
// masterQueue is unique and IS the run's backing queue, so this finds
340-
// the group the run actually ran on.
341-
masterQueue: baseWorkerQueue(run.workerQueue),
342-
},
349+
where: { masterQueue },
343350
});
344351

345352
// Show the stamped geo region as the name so a migrated run never reveals
346353
// its compute backing; fall back to the group name for unstamped runs.
347354
region = workerGroup
348-
? { name: run.region ?? workerGroup.name, location: workerGroup.location }
349-
: null;
355+
? { name: v2Region ?? run.region ?? workerGroup.name, location: workerGroup.location }
356+
: v2Region
357+
? { name: v2Region, location: null }
358+
: null;
350359
}
351360

352361
// Only AGENT-tagged runs can be session-bound, so skip the SessionRun lookup
Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,85 @@
1+
import { describe, expect, it } from "vitest";
2+
import { workerQueueForBirth } from "~/runEngine/concerns/workerQueueAssignment.server";
3+
import { FEATURE_FLAG } from "~/v3/featureFlags";
4+
5+
const assignment = {
6+
workerQueue: "us-east-1-next",
7+
region: "us-east-1",
8+
envType: "PRODUCTION",
9+
orgFeatureFlags: null,
10+
globalDefault: false,
11+
workerGroups: [{ masterQueue: "us-east-1-next", workloadType: "MICROVM" as const }],
12+
rootTriggerSource: "api",
13+
splitEnabled: true,
14+
};
15+
16+
describe("workerQueueForBirth", () => {
17+
it("keeps legacy assignment and scheduled lineage when v2 is off", () => {
18+
expect(workerQueueForBirth(assignment)).toBe("us-east-1-next");
19+
expect(workerQueueForBirth({ ...assignment, rootTriggerSource: "schedule" })).toBe(
20+
"us-east-1-next:scheduled"
21+
);
22+
});
23+
24+
it("keeps legacy routing when the selected queue is absent from a loaded registry", () => {
25+
const input = {
26+
...assignment,
27+
region: assignment.workerQueue,
28+
globalDefault: true,
29+
workerGroups: [{ masterQueue: "us-west-2", workloadType: "CONTAINER" as const }],
30+
};
31+
expect(workerQueueForBirth(input)).toBe("us-east-1-next");
32+
expect(workerQueueForBirth({ ...input, rootTriggerSource: "schedule" })).toBe(
33+
"us-east-1-next:scheduled"
34+
);
35+
});
36+
37+
it("uses the geographic region and selected runtime when v2 is enabled", () => {
38+
expect(workerQueueForBirth({ ...assignment, globalDefault: true })).toBe(
39+
"us-east-1:v2:ondemand:fresh:compute:stable"
40+
);
41+
expect(
42+
workerQueueForBirth({
43+
...assignment,
44+
globalDefault: true,
45+
workerGroups: [{ masterQueue: "us-east-1-next", workloadType: "CONTAINER" }],
46+
})
47+
).toBe("us-east-1:v2:ondemand:fresh:container:stable");
48+
});
49+
50+
it("honors org opt-in, channel and explicit shared-runtime assignment", () => {
51+
expect(
52+
workerQueueForBirth({
53+
...assignment,
54+
orgFeatureFlags: {
55+
[FEATURE_FLAG.workerQueueV2Enabled]: true,
56+
[FEATURE_FLAG.workerQueueCompatibility]: "any",
57+
[FEATURE_FLAG.workerQueueChannel]: "canary",
58+
},
59+
rootTriggerSource: "schedule",
60+
})
61+
).toBe("us-east-1:v2:scheduled:fresh:any:canary");
62+
});
63+
64+
it("honors org opt-out even when the global default is enabled", () => {
65+
expect(
66+
workerQueueForBirth({
67+
...assignment,
68+
globalDefault: true,
69+
orgFeatureFlags: { [FEATURE_FLAG.workerQueueV2Enabled]: false },
70+
})
71+
).toBe("us-east-1-next");
72+
});
73+
74+
it("leaves development on the legacy path even with an org opt-in", () => {
75+
expect(
76+
workerQueueForBirth({
77+
...assignment,
78+
envType: "DEVELOPMENT",
79+
splitEnabled: false,
80+
globalDefault: true,
81+
orgFeatureFlags: { [FEATURE_FLAG.workerQueueV2Enabled]: true },
82+
})
83+
).toBe("us-east-1-next");
84+
});
85+
});

0 commit comments

Comments
 (0)