Compare commits

...

1 Commits

Author SHA1 Message Date
Eric Allam 79393deacf feat(webapp): gate worker dequeues by worker queue via env var
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 10m57s
🚀 Publish Trigger.dev Docker / publish-worker-v4 (push) Has been skipped
🚀 Publish Trigger.dev Docker / publish-webapp (push) Has been skipped
🚀 Publish Trigger.dev Docker / publish-worker (push) Has been skipped
🚀 Publish Trigger.dev Docker / scan-webapp (push) Has been skipped
🚀 Publish Trigger.dev Docker / units (push) Failing after 13m27s
🚀 Publish Trigger.dev Docker / 📣 Dispatch main image (push) Has been cancelled
Add RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES: a comma-separated list of
worker queues (or base regions, which also cover the :scheduled split) for
which the engine API refuses worker dequeue requests and returns no work, so
those runs stay queued instead of being handed to workers that cannot run
them. Unset means no gating. Blocked dequeues increment the
run_engine.dequeue.blocked otel counter, tagged by worker_queue and region.
2026-06-23 10:17:29 +01:00
6 changed files with 114 additions and 3 deletions
+6
View File
@@ -0,0 +1,6 @@
---
area: webapp
type: feature
---
Add a `RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES` setting that refuses worker dequeue requests for the listed worker queues (or base regions), so their runs stay queued instead of being handed to workers that can't run them.
+1
View File
@@ -809,6 +809,7 @@ const EnvironmentSchema = z
RUN_ENGINE_RETRY_WARM_START_THRESHOLD_MS: z.coerce.number().int().default(30_000),
RUN_ENGINE_PROCESS_WORKER_QUEUE_DEBOUNCE_MS: z.coerce.number().int().default(200),
RUN_ENGINE_DEQUEUE_BLOCKING_TIMEOUT_SECONDS: z.coerce.number().int().default(10),
RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES: z.string().optional(),
RUN_ENGINE_MASTER_QUEUE_CONSUMERS_INTERVAL_MS: z.coerce.number().int().default(1000),
RUN_ENGINE_MASTER_QUEUE_COOLOFF_PERIOD_MS: z.coerce.number().int().default(10_000),
RUN_ENGINE_MASTER_QUEUE_COOLOFF_COUNT_THRESHOLD: z.coerce.number().int().default(10),
@@ -0,0 +1,29 @@
import { getMeter } from "@internal/tracing";
import { env } from "~/env.server";
import {
baseWorkerQueue,
matchesDisabledWorkerQueue,
parseDisabledWorkerQueues,
} from "./workerQueueSplit.server";
const meter = getMeter("run-engine-dequeue-gate");
const blockedDequeueCounter = meter.createCounter("run_engine.dequeue.blocked", {
description:
"Count of worker dequeue requests refused because the worker queue is gated off via RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES",
});
const disabledWorkerQueues = parseDisabledWorkerQueues(
env.RUN_ENGINE_DEQUEUE_DISABLED_WORKER_QUEUES
);
export function isWorkerQueueDequeueDisabled(workerQueue: string): boolean {
return matchesDisabledWorkerQueue(workerQueue, disabledWorkerQueues);
}
export function recordBlockedDequeue(workerQueue: string): void {
blockedDequeueCounter.add(1, {
worker_queue: workerQueue,
region: baseWorkerQueue(workerQueue),
});
}
@@ -109,3 +109,25 @@ export function workerQueueForClass(
return masterQueue;
}
export function parseDisabledWorkerQueues(raw: string | undefined): Set<string> {
return new Set(
(raw ?? "")
.split(",")
.map((entry) => entry.trim())
.filter(Boolean)
);
}
export function matchesDisabledWorkerQueue(
workerQueue: string,
disabledWorkerQueues: ReadonlySet<string>
): boolean {
if (disabledWorkerQueues.size === 0) {
return false;
}
return (
disabledWorkerQueues.has(workerQueue) || disabledWorkerQueues.has(baseWorkerQueue(workerQueue))
);
}
@@ -28,6 +28,10 @@ import { singleton } from "~/utils/singleton";
import { resolveVariablesForEnvironment } from "~/v3/environmentVariables/environmentVariablesRepository.server";
import { machinePresetFromName } from "~/v3/machinePresets.server";
import { workerQueueForClass } from "~/runEngine/concerns/workerQueueSplit.server";
import {
isWorkerQueueDequeueDisabled,
recordBlockedDequeue,
} from "~/runEngine/concerns/dequeueGate.server";
import { WithRunEngine, WithRunEngineOptions } from "../baseService.server";
const authenticatedWorkerInstanceCache = singleton(
@@ -377,11 +381,16 @@ export class AuthenticatedWorkerInstance extends WithRunEngine {
runnerId?: string;
queueClass?: WorkerQueueClass;
}): Promise<DequeuedMessage[]> {
// Derive the actual queue from this worker's own masterQueue + class, so a
// token can only ever reach its own region's queues (default or :scheduled).
const workerQueue = workerQueueForClass(this.masterQueue, queueClass);
if (isWorkerQueueDequeueDisabled(workerQueue)) {
recordBlockedDequeue(workerQueue);
return [];
}
return await this._engine.dequeueFromWorkerQueue({
consumerId: this.workerInstanceId,
workerQueue: workerQueueForClass(this.masterQueue, queueClass),
workerQueue,
workerId: this.workerInstanceId,
runnerId,
});
@@ -0,0 +1,44 @@
import { describe, expect, it } from "vitest";
import {
matchesDisabledWorkerQueue,
parseDisabledWorkerQueues,
} from "~/runEngine/concerns/workerQueueSplit.server";
describe("parseDisabledWorkerQueues", () => {
it("returns an empty set for undefined or empty input", () => {
expect(parseDisabledWorkerQueues(undefined).size).toBe(0);
expect(parseDisabledWorkerQueues("").size).toBe(0);
expect(parseDisabledWorkerQueues(" , ,").size).toBe(0);
});
it("splits, trims, and drops empties", () => {
const parsed = parseDisabledWorkerQueues(" eu-central-1 , us-east-1:scheduled ,, ");
expect([...parsed]).toEqual(["eu-central-1", "us-east-1:scheduled"]);
});
});
describe("matchesDisabledWorkerQueue", () => {
it("never matches when the disabled set is empty", () => {
const empty = parseDisabledWorkerQueues(undefined);
expect(matchesDisabledWorkerQueue("eu-central-1", empty)).toBe(false);
expect(matchesDisabledWorkerQueue("eu-central-1:scheduled", empty)).toBe(false);
});
it("gates the base region and its scheduled split when the base region is listed", () => {
const disabled = parseDisabledWorkerQueues("eu-central-1");
expect(matchesDisabledWorkerQueue("eu-central-1", disabled)).toBe(true);
expect(matchesDisabledWorkerQueue("eu-central-1:scheduled", disabled)).toBe(true);
});
it("leaves other regions alone", () => {
const disabled = parseDisabledWorkerQueues("eu-central-1");
expect(matchesDisabledWorkerQueue("us-east-1", disabled)).toBe(false);
expect(matchesDisabledWorkerQueue("us-east-1:scheduled", disabled)).toBe(false);
});
it("gates only the scheduled split when a full worker queue is listed", () => {
const disabled = parseDisabledWorkerQueues("eu-central-1:scheduled");
expect(matchesDisabledWorkerQueue("eu-central-1:scheduled", disabled)).toBe(true);
expect(matchesDisabledWorkerQueue("eu-central-1", disabled)).toBe(false);
});
});