Skip to content

Commit 9272d01

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
feat: add versioned queues and authorized multi-queue dequeue
Adds authenticated, weighted multi-queue consumption with separate dispatch classes and execution phases. Existing consumers retain their legacy queue selection; requested subscriptions are validated and authorized before dequeue. Mono-RevId: 1b2b009af1ec70bc9bc23eb7d1a9a78cb60862e7
1 parent 25317cf commit 9272d01

31 files changed

Lines changed: 1496 additions & 246 deletions

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { isValidDatabaseUrl } from "./utils/db";
66
import { parseRunOpsShards, validateShardListAgainstNewUrl } from "~/v3/runOpsShards.server";
77
import { isValidRegex } from "./utils/regex";
88
import { isValidDuration } from "./services/realtime/duration.server";
9+
import { WorkerQueueSubscriptionPolicyEnv } from "./runEngine/concerns/workerQueueSubscriptions.server";
910

1011
// `z.string()` constrained to a `parseDuration`-parseable string (e.g.
1112
// `7d`, `1h`). Validated at boot so a typo'd duration fails fast.
@@ -1226,6 +1227,8 @@ const EnvironmentSchema = z
12261227
RUN_ENGINE_PROCESS_WORKER_QUEUE_DEBOUNCE_MS: z.coerce.number().int().default(200),
12271228
RUN_ENGINE_DEQUEUE_BLOCKING_TIMEOUT_SECONDS: z.coerce.number().int().default(10),
12281229
RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES: z.string().optional(),
1230+
// JSON map of worker-group IDs to allowed v2 subscriptions. Unlisted groups are denied.
1231+
RUN_ENGINE_WORKER_QUEUE_SUBSCRIPTIONS: WorkerQueueSubscriptionPolicyEnv,
12291232
RUN_ENGINE_MASTER_QUEUE_CONSUMERS_INTERVAL_MS: z.coerce.number().int().default(1000),
12301233
// Off by default. Enable on a single service (e.g. the engine worker) so only one
12311234
// instance reports worker queue length, rather than every replica.
Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,100 @@
1+
import { startTestServer, type TestServer } from "@internal/testcontainers/webapp";
2+
import type { WorkerQueueSubscription } from "@trigger.dev/core/v3/workers";
3+
import { afterAll, beforeAll, describe, expect, it } from "vitest";
4+
import {
5+
seedEngineFixtures,
6+
workerHeaders,
7+
type EngineFixtures,
8+
} from "../../test/bench/lib/engineFixtures";
9+
10+
const WORKER_GROUP_ID = "worker-group-legacy-auth";
11+
const LEGACY_MASTER_QUEUE = "my workers";
12+
const subscription: WorkerQueueSubscription = {
13+
class: "ondemand",
14+
phase: "fresh",
15+
compat: "container",
16+
channel: "stable",
17+
};
18+
19+
let server: TestServer;
20+
let fixtures: EngineFixtures;
21+
22+
beforeAll(async () => {
23+
server = await startTestServer({
24+
extraEnv: {
25+
RUN_ENGINE_WORKER_QUEUE_SUBSCRIPTIONS: JSON.stringify({
26+
[WORKER_GROUP_ID]: [subscription],
27+
}),
28+
},
29+
overrideEnv: { RUN_ENGINE_WORKER_ENABLED: "1" },
30+
});
31+
fixtures = await seedEngineFixtures(server.prisma, {
32+
taskCount: 1,
33+
workerGroupId: WORKER_GROUP_ID,
34+
masterQueue: LEGACY_MASTER_QUEUE,
35+
enableFastPath: true,
36+
});
37+
}, 180_000);
38+
39+
afterAll(async () => {
40+
await server?.stop();
41+
}, 120_000);
42+
43+
function workerAction(path: string, body: unknown) {
44+
return server.webapp.fetch(path, {
45+
method: "POST",
46+
headers: workerHeaders(fixtures, "legacy-worker-instance"),
47+
body: JSON.stringify(body),
48+
});
49+
}
50+
51+
describe("dequeue request validation", () => {
52+
it("authenticates legacy workers and rejects invalid requests without popping their queue", async () => {
53+
const heartbeat = await workerAction("/engine/v1/worker-actions/heartbeat", {
54+
cpu: { used: 0, available: 1 },
55+
memory: { used: 0, available: 1 },
56+
tasks: [],
57+
});
58+
expect(heartbeat.status).toBe(200);
59+
60+
const trigger = await server.webapp.fetch(
61+
`/api/v1/tasks/${fixtures.taskIdentifiers[0]!}/trigger`,
62+
{
63+
method: "POST",
64+
headers: {
65+
Authorization: `Bearer ${fixtures.environmentApiKey}`,
66+
"content-type": "application/json",
67+
},
68+
body: JSON.stringify({ payload: { message: "dequeue route validation" } }),
69+
}
70+
);
71+
expect(trigger.ok).toBe(true);
72+
73+
const invalidRequests: Array<[body: unknown, status: number]> = [
74+
[{ queueClass: "restore" }, 400],
75+
[{ maxRunCount: "invalid" }, 400],
76+
[{ maxResources: { cpu: "invalid", memory: 1 } }, 400],
77+
[{ subscriptions: [] }, 422],
78+
[{ subscriptions: [{ ...subscription, weight: 2 }] }, 422],
79+
[{ queueClass: "default", subscriptions: [subscription] }, 422],
80+
[{ maxRunCount: "invalid", subscriptions: [] }, 400],
81+
[{ queueClass: "invalid", subscriptions: [subscription] }, 400],
82+
];
83+
for (const [body, status] of invalidRequests) {
84+
expect((await workerAction("/engine/v1/worker-actions/dequeue", body)).status).toBe(status);
85+
}
86+
87+
const invalidV2Region = await workerAction("/engine/v1/worker-actions/dequeue", {
88+
subscriptions: [subscription],
89+
});
90+
expect(invalidV2Region.status).toBe(422);
91+
92+
const legacyDequeue = await workerAction("/engine/v1/worker-actions/dequeue", {});
93+
expect(legacyDequeue.status).toBe(200);
94+
expect(await legacyDequeue.json()).toHaveLength(1);
95+
96+
const emptyDequeue = await workerAction("/engine/v1/worker-actions/dequeue", {});
97+
expect(emptyDequeue.status).toBe(200);
98+
expect(await emptyDequeue.json()).toEqual([]);
99+
}, 180_000);
100+
});

‎apps/webapp/app/routes/engine.v1.worker-actions.dequeue.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,12 +8,20 @@ import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.
88
export const action = createActionWorkerApiRoute(
99
{
1010
body: z.compile(WorkerApiDequeueRequestBody),
11+
bodyValidationErrorStatus: (error) =>
12+
error.issues.every((issue) => issue.path[0] === "subscriptions") ? 422 : 400,
1113
},
1214
async ({
1315
authenticatedWorker,
1416
runnerId,
1517
body,
1618
}): Promise<TypedResponse<WorkerApiDequeueResponseBody>> => {
17-
return json(await authenticatedWorker.dequeue({ runnerId, queueClass: body.queueClass }));
19+
return json(
20+
await authenticatedWorker.dequeue({
21+
runnerId,
22+
queueClass: body.queueClass,
23+
subscriptions: body.subscriptions,
24+
})
25+
);
1826
}
1927
);

‎apps/webapp/app/runEngine/concerns/dequeueGate.server.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,11 @@ const disabledWorkerQueues = parseDisabledWorkerQueues(
1717
env.RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES
1818
);
1919

20-
export function isWorkerQueueDequeueDisabled(workerQueue: string): boolean {
21-
return matchesDisabledWorkerQueue(workerQueue, disabledWorkerQueues);
20+
export function isWorkerQueueDequeueDisabled(
21+
workerQueue: string,
22+
version: "legacy" | "v2"
23+
): boolean {
24+
return matchesDisabledWorkerQueue(workerQueue, disabledWorkerQueues, version);
2225
}
2326

2427
export function recordBlockedDequeue(workerQueue: string): void {
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
import { RunEngine } from "@internal/run-engine";
2+
import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "@internal/run-engine/tests";
3+
import { assertNonNullable, containerTest } from "@internal/testcontainers";
4+
import { trace } from "@internal/tracing";
5+
import { generateFriendlyId } from "@trigger.dev/core/v3/isomorphic";
6+
import type { WorkerQueueSubscription } from "@trigger.dev/core/v3/workers";
7+
import { expect } from "vitest";
8+
import { dequeueWorkerQueues } from "~/runEngine/concerns/workerQueueDequeue.server";
9+
import { createWorkerQueueConsumer } from "~/runEngine/concerns/workerQueueSubscriptions.server";
10+
11+
containerTest(
12+
"authorizes the entire subscription set before pop and retains legacy dequeue",
13+
async ({ prisma, redisOptions }) => {
14+
const engine = new RunEngine({
15+
prisma,
16+
worker: { redis: redisOptions, workers: 1, tasksPerWorker: 10, pollIntervalMs: 100 },
17+
queue: { redis: redisOptions, masterQueueConsumersDisabled: true },
18+
runLock: { redis: redisOptions },
19+
machines: {
20+
defaultMachine: "small-1x",
21+
machines: {
22+
"small-1x": { name: "small-1x", cpu: 0.5, memory: 0.5, centsPerMs: 0.0001 },
23+
},
24+
baseCostInCents: 0.0005,
25+
},
26+
tracer: trace.getTracer("multi-queue-test"),
27+
});
28+
29+
try {
30+
const environment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
31+
const taskIdentifier = "test-task";
32+
await setupBackgroundWorker(engine, environment, taskIdentifier);
33+
const stable: WorkerQueueSubscription = {
34+
class: "ondemand",
35+
phase: "fresh",
36+
compat: "container",
37+
channel: "stable",
38+
};
39+
const canary: WorkerQueueSubscription = { ...stable, channel: "canary" };
40+
const scheduled: WorkerQueueSubscription = { ...stable, class: "scheduled" };
41+
const worker = createWorkerQueueConsumer({
42+
workerGroupId: environment.project.defaultWorkerGroupId!,
43+
workerInstanceId: "test-instance",
44+
masterQueue: "legacy-region",
45+
region: "us-east-1",
46+
workloadType: "CONTAINER" as const,
47+
allowedSubscriptions: [stable, scheduled],
48+
});
49+
50+
const lanes = [
51+
"us-east-1:v2:ondemand:fresh:container:stable",
52+
"us-east-1:v2:scheduled:fresh:container:stable",
53+
"legacy-region",
54+
"legacy-region:scheduled",
55+
];
56+
const runs = [];
57+
for (const [index, workerQueue] of lanes.entries()) {
58+
runs.push(
59+
await engine.trigger(
60+
{
61+
number: index + 1,
62+
friendlyId: generateFriendlyId("run"),
63+
environment,
64+
taskIdentifier,
65+
payload: "{}",
66+
payloadType: "application/json",
67+
context: {},
68+
traceContext: {},
69+
traceId: "test-trace",
70+
spanId: "test-span",
71+
workerQueue,
72+
enableFastPath: true,
73+
queue: `task/${taskIdentifier}`,
74+
isTest: false,
75+
tags: [],
76+
},
77+
prisma
78+
)
79+
);
80+
}
81+
82+
expect(() =>
83+
dequeueWorkerQueues({ engine, worker, subscriptions: [stable, canary] })
84+
).toThrow("not authorized");
85+
const unchanged = await engine.getRunExecutionData({ runId: runs[0]!.id });
86+
expect(unchanged?.snapshot.executionStatus).toBe("QUEUED");
87+
88+
const deliveredRunIds = new Set<string>();
89+
for (const remaining of [1, 0]) {
90+
const delivery = await dequeueWorkerQueues({
91+
engine,
92+
worker,
93+
subscriptions: [stable, scheduled],
94+
});
95+
expect(delivery).toHaveLength(1);
96+
assertNonNullable(delivery[0]);
97+
deliveredRunIds.add(delivery[0].run.id);
98+
expect(delivery[0].snapshot.executionStatus).toBe("PENDING_EXECUTING");
99+
expect(delivery[0].workerQueueLength).toBe(remaining);
100+
}
101+
expect(deliveredRunIds).toEqual(new Set([runs[0]!.id, runs[1]!.id]));
102+
103+
expect((await dequeueWorkerQueues({ engine, worker }))[0]?.run.id).toBe(runs[2]!.id);
104+
expect(
105+
(await dequeueWorkerQueues({ engine, worker, queueClass: "scheduled" }))[0]?.run.id
106+
).toBe(runs[3]!.id);
107+
} finally {
108+
await engine.quit();
109+
}
110+
}
111+
);
Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
import type { RunEngine } from "@internal/run-engine";
2+
import type {
3+
WeightedWorkerQueueSubscription,
4+
WorkerQueueClass,
5+
} from "@trigger.dev/core/v3/workers";
6+
import { isWorkerQueueDequeueDisabled, recordBlockedDequeue } from "./dequeueGate.server";
7+
import { workerQueueForClass } from "./workerQueueSplit.server";
8+
import {
9+
resolveWorkerQueueSubscriptions,
10+
type WorkerQueueConsumer,
11+
} from "./workerQueueSubscriptions.server";
12+
13+
export function dequeueWorkerQueues({
14+
engine,
15+
worker,
16+
runnerId,
17+
queueClass,
18+
subscriptions,
19+
}: {
20+
engine: RunEngine;
21+
worker: WorkerQueueConsumer;
22+
runnerId?: string;
23+
queueClass?: WorkerQueueClass;
24+
subscriptions?: WeightedWorkerQueueSubscription[];
25+
}) {
26+
const workerQueues = subscriptions
27+
? resolveWorkerQueueSubscriptions(worker, subscriptions)
28+
: [{ queue: workerQueueForClass(worker.masterQueue, queueClass), weight: 1 }];
29+
30+
const enabledQueues = workerQueues.filter(({ queue }) => {
31+
if (isWorkerQueueDequeueDisabled(queue, subscriptions ? "v2" : "legacy")) {
32+
recordBlockedDequeue(queue);
33+
return false;
34+
}
35+
return true;
36+
});
37+
if (enabledQueues.length === 0) {
38+
return Promise.resolve([]);
39+
}
40+
41+
return engine.dequeueFromWorkerQueues({
42+
consumerId: worker.workerInstanceId,
43+
workerQueues: enabledQueues,
44+
workerId: worker.workerInstanceId,
45+
runnerId,
46+
});
47+
}

0 commit comments

Comments
 (0)