From f10b382bdf527b19738015bfe65acd7039759f83 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 10:31:23 +0100 Subject: [PATCH 1/3] feat(sdk,core,webapp,run-engine): declare queue gates on tasks and triggers The queue option on task() and on trigger now accepts an array: the first entry is the queue the run waits in, the rest (at most two) are gates the run must also hold a concurrency slot in while executing. A gate without a concurrencyKey inherits the run's own key; a literal key pins the gate to one shared slot pool. Task-level gates apply to every trigger; a trigger array replaces them for that run. Gates persist on BackgroundWorkerTask (task defaults, carried through the task metadata cache with no extra trigger-time queries) and on TaskRun (replay fidelity and payload rebuilds), and flow into the run queue message where the engine enforces them behind its flag. --- .changeset/queue-gates.md | 22 +++ .../app/runEngine/concerns/queues.server.ts | 48 ++++- .../runEngine/services/triggerTask.server.ts | 6 +- apps/webapp/app/runEngine/types.ts | 2 + .../app/services/taskMetadataCache.server.ts | 32 ++++ .../changeCurrentDeployment.server.ts | 3 + .../services/createBackgroundWorker.server.ts | 3 + .../app/v3/services/replayTaskRun.server.ts | 3 + .../migration.sql | 5 + .../database/prisma/schema.prisma | 7 + .../run-engine/src/engine/index.ts | 2 + .../src/engine/systems/enqueueSystem.ts | 28 +++ .../src/engine/tests/queueGates.test.ts | 137 ++++++++++++++ .../run-engine/src/engine/types.ts | 3 + .../run-engine/src/run-queue/types.ts | 4 +- internal-packages/run-store/src/types.ts | 2 + packages/core/src/v3/schemas/api.ts | 18 ++ packages/core/src/v3/schemas/resources.ts | 3 +- packages/core/src/v3/schemas/schemas.ts | 12 ++ packages/core/src/v3/types/tasks.ts | 35 +++- packages/trigger-sdk/src/v3/shared.ts | 173 ++++++++++++++---- 21 files changed, 499 insertions(+), 49 deletions(-) create mode 100644 .changeset/queue-gates.md create mode 100644 internal-packages/database/prisma/migrations/20260829120000_add_queue_gates/migration.sql create mode 100644 internal-packages/run-engine/src/engine/tests/queueGates.test.ts diff --git a/.changeset/queue-gates.md b/.changeset/queue-gates.md new file mode 100644 index 00000000000..b4934bc92b3 --- /dev/null +++ b/.changeset/queue-gates.md @@ -0,0 +1,22 @@ +--- +"@trigger.dev/sdk": patch +"@trigger.dev/core": patch +--- + +Hold a run's concurrency slot in more than one queue with queue gates. Pass an array as `queue`: the first entry is the queue the run waits in, and up to two more name gates, other queues the run must also have capacity in and occupies while it executes. A gate without a `concurrencyKey` uses the run's own key, so a shared `tenant` queue caps a tenant across every task; a literal key pins the gate to one slot pool, capping, say, all traffic to one external provider. + +```ts +import { queue, task } from "@trigger.dev/sdk"; + +export const tenant = queue({ name: "tenant", concurrencyLimit: 10 }); + +export const processWebhook = task({ + id: "process-webhook", + queue: [{ name: "webhooks", concurrencyLimit: 2 }, "tenant"], + run: async (payload) => {}, +}); + +await processWebhook.trigger(payload, { concurrencyKey: tenantId }); +``` + +The same array form works on `queue` when triggering, replacing the task's gates for that run. Enforcement happens server-side; servers without queue gates enabled accept the option but run without it. diff --git a/apps/webapp/app/runEngine/concerns/queues.server.ts b/apps/webapp/app/runEngine/concerns/queues.server.ts index 1374a34d288..9214e91b9b3 100644 --- a/apps/webapp/app/runEngine/concerns/queues.server.ts +++ b/apps/webapp/app/runEngine/concerns/queues.server.ts @@ -23,7 +23,12 @@ import { Namespace, } from "@internal/cache"; import { singleton } from "~/utils/singleton"; -import type { TaskMetadataCache, TaskMetadataEntry } from "~/services/taskMetadataCache.server"; +import { + parseTaskGates, + type TaskMetadataCache, + type TaskMetadataEntry, + type TaskMetadataGate, +} from "~/services/taskMetadataCache.server"; import { taskMetadataCacheInstance } from "~/services/taskMetadataCacheInstance.server"; import { recordTaskMetaResolve, @@ -95,6 +100,7 @@ export class DefaultQueueManager implements QueueManager { let lockedQueueId: string | undefined; let taskTtl: string | null | undefined; let taskKind: string | undefined; + let taskGates: TaskMetadataGate[] | null | undefined; // Determine queue name based on lockToVersion and provided options if (lockedBackgroundWorker) { @@ -146,6 +152,7 @@ export class DefaultQueueManager implements QueueManager { taskTtl = lockedMeta?.ttl ?? undefined; } taskKind = lockedMeta?.triggerSource; + taskGates = lockedMeta?.gates; } else { // No queue override - resolve default queue + TTL + triggerSource via cache, // falling back to a single BackgroundWorkerTask lookup on miss. @@ -184,6 +191,7 @@ export class DefaultQueueManager implements QueueManager { queueName = lockedMeta.queueName; lockedQueueId = lockedMeta.queueId ?? undefined; taskKind = lockedMeta.triggerSource; + taskGates = lockedMeta.gates; } } else { // Task is not locked to a specific version, use regular logic @@ -199,6 +207,7 @@ export class DefaultQueueManager implements QueueManager { queueName = taskInfo.queueName; taskTtl = taskInfo.taskTtl; taskKind = taskInfo.taskKind; + taskGates = taskInfo.taskGates; } // Sanitize the final determined queue name once @@ -211,17 +220,29 @@ export class DefaultQueueManager implements QueueManager { queueName = sanitizedQueueName; } + const requestedGates = request.body.options?.gates ?? taskGates ?? undefined; + const gates = requestedGates + ?.flatMap((gate) => { + const sanitized = sanitizeQueueName(gate.queue); + return sanitized ? [{ queue: sanitized, concurrencyKey: gate.concurrencyKey }] : []; + }) + .slice(0, 2); + return { queueName, lockedQueueId, taskTtl, taskKind, + gates: gates && gates.length > 0 ? gates : undefined, }; } - private async getTaskQueueInfo( - request: TriggerTaskRequest - ): Promise<{ queueName: string; taskTtl?: string | null; taskKind?: string | undefined }> { + private async getTaskQueueInfo(request: TriggerTaskRequest): Promise<{ + queueName: string; + taskTtl?: string | null; + taskKind?: string | undefined; + taskGates?: TaskMetadataGate[] | null; + }> { const { taskId, environment, body } = request; const { queue } = body.options ?? {}; @@ -243,6 +264,7 @@ export class DefaultQueueManager implements QueueManager { queueName: overriddenQueueName, taskTtl: meta?.ttl ?? undefined, taskKind: meta?.triggerSource, + taskGates: meta?.gates, }; } @@ -259,10 +281,20 @@ export class DefaultQueueManager implements QueueManager { taskId, environmentId: environment.id, }); - return { queueName: defaultQueueName, taskTtl: meta.ttl, taskKind: meta.triggerSource }; + return { + queueName: defaultQueueName, + taskTtl: meta.ttl, + taskKind: meta.triggerSource, + taskGates: meta.gates, + }; } - return { queueName: meta.queueName, taskTtl: meta.ttl, taskKind: meta.triggerSource }; + return { + queueName: meta.queueName, + taskTtl: meta.ttl, + taskKind: meta.triggerSource, + taskGates: meta.gates, + }; } /** @@ -320,6 +352,7 @@ export class DefaultQueueManager implements QueueManager { triggerSource: row.triggerSource, queueId: row.queue?.id ?? null, queueName: row.queue?.name ?? "", + gates: parseTaskGates(row.gates), }; // Fire-and-forget back-fill — `setByWorker` upserts the single field and @@ -340,6 +373,7 @@ export class DefaultQueueManager implements QueueManager { select: { ttl: true, triggerSource: true, + gates: true, queue: { select: { id: true, name: true } }, }, }); @@ -378,6 +412,7 @@ export class DefaultQueueManager implements QueueManager { select: { ttl: true, triggerSource: true, + gates: true, queue: { select: { id: true, name: true } }, }, }); @@ -395,6 +430,7 @@ export class DefaultQueueManager implements QueueManager { triggerSource: row.triggerSource, queueId: row.queue?.id ?? null, queueName: row.queue?.name ?? "", + gates: parseTaskGates(row.gates), }; // Fire-and-forget back-fill — atomically upserts the slug into both diff --git a/apps/webapp/app/runEngine/services/triggerTask.server.ts b/apps/webapp/app/runEngine/services/triggerTask.server.ts index d3320dbc219..2994f3ce7af 100644 --- a/apps/webapp/app/runEngine/services/triggerTask.server.ts +++ b/apps/webapp/app/runEngine/services/triggerTask.server.ts @@ -445,7 +445,7 @@ export class RunEngineTriggerTaskService { const parkedOnExternalDeploymentId = externalDeploymentResolution?.outcome === "park" ? externalDeploymentId : undefined; - const { queueName, lockedQueueId, taskTtl, taskKind } = + const { queueName, lockedQueueId, taskTtl, taskKind, gates } = await this.queueConcern.resolveQueueProperties( triggerRequest, lockedToBackgroundWorker ?? undefined @@ -663,6 +663,7 @@ export class RunEngineTriggerTaskService { options, queueName, lockedQueueId, + gates, workerQueue, region: migrated.region, enableFastPath: migrated.enableFastPath, @@ -743,6 +744,7 @@ export class RunEngineTriggerTaskService { options, queueName, lockedQueueId, + gates, workerQueue, region: migrated.region, enableFastPath: migrated.enableFastPath, @@ -905,6 +907,7 @@ export class RunEngineTriggerTaskService { options: TriggerTaskServiceOptions; queueName: string; lockedQueueId?: string; + gates?: Array<{ queue: string; concurrencyKey?: string }>; workerQueue?: string; region?: string; enableFastPath: boolean; @@ -971,6 +974,7 @@ export class RunEngineTriggerTaskService { : args.body.options?.concurrencyKey, queue: args.queueName, lockedQueueId: args.lockedQueueId, + gates: args.gates, workerQueue: args.workerQueue, region: args.region, enableFastPath: args.enableFastPath, diff --git a/apps/webapp/app/runEngine/types.ts b/apps/webapp/app/runEngine/types.ts index 4e415483120..a30606e92a6 100644 --- a/apps/webapp/app/runEngine/types.ts +++ b/apps/webapp/app/runEngine/types.ts @@ -46,6 +46,8 @@ export type QueueProperties = { lockedQueueId?: string; taskTtl?: string | null; taskKind?: string; + /** Other queues the run must also hold a concurrency slot in while executing. */ + gates?: Array<{ queue: string; concurrencyKey?: string }>; }; export type LockedBackgroundWorker = Pick< diff --git a/apps/webapp/app/services/taskMetadataCache.server.ts b/apps/webapp/app/services/taskMetadataCache.server.ts index 6130295a73f..e420329d74d 100644 --- a/apps/webapp/app/services/taskMetadataCache.server.ts +++ b/apps/webapp/app/services/taskMetadataCache.server.ts @@ -2,12 +2,16 @@ import type { Redis, Result, Callback } from "ioredis"; import type { TaskTriggerSource } from "@trigger.dev/database"; import { logger } from "./logger.server"; +export type TaskMetadataGate = { queue: string; concurrencyKey?: string }; + export type TaskMetadataEntry = { slug: string; ttl: string | null; triggerSource: TaskTriggerSource; queueId: string | null; queueName: string; + /** Task-declared gates, applied to every trigger that does not override them. */ + gates: TaskMetadataGate[] | null; }; export interface TaskMetadataCache { @@ -52,11 +56,37 @@ export type RedisTaskMetadataCacheOptions = { byWorkerTtlSeconds?: number; }; +/** + * BackgroundWorkerTask.gates is an untyped Json column; keep only well-shaped + * entries so a malformed value can never fail a trigger. + */ +export function parseTaskGates(gates: unknown): TaskMetadataGate[] | null { + if (!Array.isArray(gates) || gates.length === 0) { + return null; + } + + const parsed = gates.flatMap((gate) => { + if (!gate || typeof gate !== "object" || typeof (gate as any).queue !== "string") { + return []; + } + const concurrencyKey = (gate as any).concurrencyKey; + return [ + { + queue: (gate as any).queue, + concurrencyKey: typeof concurrencyKey === "string" ? concurrencyKey : undefined, + }, + ]; + }); + + return parsed.length > 0 ? parsed.slice(0, 2) : null; +} + type EncodedEntry = { t: string | null; k: TaskTriggerSource; q: string | null; n: string; + g?: TaskMetadataGate[] | null; }; function encode(entry: TaskMetadataEntry): string { @@ -65,6 +95,7 @@ function encode(entry: TaskMetadataEntry): string { k: entry.triggerSource, q: entry.queueId, n: entry.queueName, + g: entry.gates, }; return JSON.stringify(payload); } @@ -78,6 +109,7 @@ function decode(slug: string, raw: string): TaskMetadataEntry | null { triggerSource: parsed.k, queueId: parsed.q, queueName: parsed.n, + gates: parsed.g ?? null, }; } catch (error) { logger.error("Failed to decode task metadata cache entry", { slug, error }); diff --git a/apps/webapp/app/v3/services/changeCurrentDeployment.server.ts b/apps/webapp/app/v3/services/changeCurrentDeployment.server.ts index 5f5697dd7cb..021cac19174 100644 --- a/apps/webapp/app/v3/services/changeCurrentDeployment.server.ts +++ b/apps/webapp/app/v3/services/changeCurrentDeployment.server.ts @@ -6,6 +6,7 @@ import { logger } from "~/services/logger.server"; import { syncTaskIdentifiers } from "~/services/taskIdentifierRegistry.server"; import { type TaskMetadataCache, + parseTaskGates, type TaskMetadataEntry, } from "~/services/taskMetadataCache.server"; import { taskMetadataCacheInstance } from "~/services/taskMetadataCacheInstance.server"; @@ -119,6 +120,7 @@ export class ChangeCurrentDeploymentService extends BaseService { slug: true, triggerSource: true, ttl: true, + gates: true, queue: { select: { id: true, name: true } }, }, }) @@ -157,6 +159,7 @@ export class ChangeCurrentDeploymentService extends BaseService { triggerSource: t.triggerSource, queueId: t.queue?.id ?? null, queueName: t.queue?.name ?? "", + gates: parseTaskGates(t.gates), })); // Cache calls log+swallow internally. diff --git a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index f1cebe640d5..87777379744 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -437,6 +437,7 @@ async function createWorkerTask( exportName: task.exportName, retryConfig: task.retry, queueConfig: task.queue, + gates: task.gates, machineConfig: task.machine, triggerSource: resolvedTriggerSource, config: task.agentConfig ? (task.agentConfig as any) : undefined, @@ -454,6 +455,7 @@ async function createWorkerTask( triggerSource: resolvedTriggerSource, queueId: queue.id, queueName: queue.name, + gates: task.gates ?? null, }; } catch (error) { if (error instanceof Prisma.PrismaClientKnownRequestError) { @@ -477,6 +479,7 @@ async function createWorkerTask( triggerSource: resolvedTriggerSource, queueId: queue.id, queueName: queue.name, + gates: task.gates ?? null, }; } } else { diff --git a/apps/webapp/app/v3/services/replayTaskRun.server.ts b/apps/webapp/app/v3/services/replayTaskRun.server.ts index 750427fb327..09b0439a4ea 100644 --- a/apps/webapp/app/v3/services/replayTaskRun.server.ts +++ b/apps/webapp/app/v3/services/replayTaskRun.server.ts @@ -142,6 +142,9 @@ export class ReplayTaskRunService extends BaseService { : undefined, concurrencyKey: overrideOptions.concurrencyKey ?? existingTaskRun.concurrencyKey ?? undefined, + gates: Array.isArray(existingTaskRun.gates) + ? (existingTaskRun.gates as Array<{ queue: string; concurrencyKey?: string }>) + : undefined, maxAttempts: overrideOptions.maxAttempts, maxDuration: overrideOptions.maxDurationSeconds, machine: diff --git a/internal-packages/database/prisma/migrations/20260829120000_add_queue_gates/migration.sql b/internal-packages/database/prisma/migrations/20260829120000_add_queue_gates/migration.sql new file mode 100644 index 00000000000..d4dfce733d2 --- /dev/null +++ b/internal-packages/database/prisma/migrations/20260829120000_add_queue_gates/migration.sql @@ -0,0 +1,5 @@ +-- AlterTable +ALTER TABLE "BackgroundWorkerTask" ADD COLUMN "gates" JSONB; + +-- AlterTable +ALTER TABLE "TaskRun" ADD COLUMN "gates" JSONB; diff --git a/internal-packages/database/prisma/schema.prisma b/internal-packages/database/prisma/schema.prisma index ac0ac09b2f0..c88ba5b5887 100644 --- a/internal-packages/database/prisma/schema.prisma +++ b/internal-packages/database/prisma/schema.prisma @@ -741,6 +741,9 @@ model BackgroundWorkerTask { queueConfig Json? retryConfig Json? machineConfig Json? + /// Gates declared on the task: other queues its runs must also hold a concurrency + /// slot in while executing. Shape: [{ queue: string, concurrencyKey?: string }] + gates Json? queueId String? queue TaskQueue? @relation(fields: [queueId], references: [id], onDelete: SetNull, onUpdate: Cascade) @@ -1107,6 +1110,10 @@ model TaskRun { concurrencyKey String? + /// Gates for this run: other queues it must also hold a concurrency slot in while + /// executing. Shape: [{ queue: string, concurrencyKey?: string }] + gates Json? + delayUntil DateTime? queuedAt DateTime? ttl String? diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index 478786dd230..9d7bb8ff947 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -833,6 +833,7 @@ export class RunEngine { sdkVersion, cliVersion, concurrencyKey, + gates, workerQueue, region, enableFastPath, @@ -1011,6 +1012,7 @@ export class RunEngine { sdkVersion, cliVersion, concurrencyKey, + gates, queue, lockedQueueId, workerQueue, diff --git a/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts b/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts index 38c681c511e..33a79d0c61f 100644 --- a/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/enqueueSystem.ts @@ -16,6 +16,33 @@ export type EnqueueSystemOptions = { executionSnapshotSystem: ExecutionSnapshotSystem; }; +/** + * TaskRun.gates is an untyped Json column; admit correctness only needs well-shaped + * entries, so anything malformed is dropped rather than failing the enqueue. + */ +function parseRunGates( + gates: unknown +): Array<{ queue: string; concurrencyKey?: string }> | undefined { + if (!Array.isArray(gates) || gates.length === 0) { + return undefined; + } + + const parsed = gates.flatMap((gate) => { + if (!gate || typeof gate !== "object" || typeof (gate as any).queue !== "string") { + return []; + } + const concurrencyKey = (gate as any).concurrencyKey; + return [ + { + queue: (gate as any).queue, + concurrencyKey: typeof concurrencyKey === "string" ? concurrencyKey : undefined, + }, + ]; + }); + + return parsed.length > 0 ? parsed.slice(0, 2) : undefined; +} + export class EnqueueSystem { private readonly $: SystemResources; private readonly executionSnapshotSystem: ExecutionSnapshotSystem; @@ -177,6 +204,7 @@ export class EnqueueSystem { environmentType: env.type, queue: run.queue, concurrencyKey: run.concurrencyKey ?? undefined, + gates: parseRunGates(run.gates), timestamp, eligibleAtMs, attempt: 0, diff --git a/internal-packages/run-engine/src/engine/tests/queueGates.test.ts b/internal-packages/run-engine/src/engine/tests/queueGates.test.ts new file mode 100644 index 00000000000..8c06e9d8563 --- /dev/null +++ b/internal-packages/run-engine/src/engine/tests/queueGates.test.ts @@ -0,0 +1,137 @@ +import { containerTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { setTimeout } from "timers/promises"; +import { expect } from "vitest"; +import { RunEngine } from "../index.js"; +import { setupAuthenticatedEnvironment, setupBackgroundWorker } from "./setup.js"; + +vi.setConfig({ testTimeout: 60_000 }); + +async function waitFor(condition: () => Promise, timeoutMs = 20_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await condition()) { + return true; + } + await setTimeout(250); + } + return condition(); +} + +describe("RunEngine queue gates", () => { + containerTest( + "trigger persists gates and the queue enforces them", + async ({ prisma, redisOptions }) => { + const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); + + const engine = new RunEngine({ + prisma, + worker: { + redis: redisOptions, + workers: 1, + tasksPerWorker: 10, + pollIntervalMs: 100, + }, + queue: { + redis: redisOptions, + processWorkerQueueDebounceMs: 50, + gatesEnabled: true, + }, + runLock: { + redis: redisOptions, + }, + machines: { + defaultMachine: "small-1x", + machines: { + "small-1x": { + name: "small-1x" as const, + cpu: 0.5, + memory: 0.5, + centsPerMs: 0.0001, + }, + }, + baseCostInCents: 0.0001, + }, + tracer: trace.getTracer("test", "0.0.0"), + }); + + try { + const taskIdentifier = "test-task"; + + await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier); + + await engine.runQueue.updateQueueConcurrencyLimits( + authenticatedEnvironment, + "shared-gate", + 1 + ); + + const run1 = await engine.trigger( + { + number: 1, + friendlyId: "run_g1", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t1", + spanId: "s1", + workerQueue: "main", + queue: `task/${taskIdentifier}`, + gates: [{ queue: "shared-gate" }], + isTest: false, + tags: [], + }, + prisma + ); + + const run2 = await engine.trigger( + { + number: 2, + friendlyId: "run_g2", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t2", + spanId: "s2", + workerQueue: "main", + queue: `task/${taskIdentifier}`, + gates: [{ queue: "shared-gate" }], + isTest: false, + tags: [], + }, + prisma + ); + + const storedRun = await prisma.taskRun.findFirst({ where: { id: run1.id } }); + expect(storedRun?.gates).toEqual([{ queue: "shared-gate" }]); + + const oneAdmitted = await waitFor( + async () => + (await engine.runQueue.currentConcurrencyOfQueue( + authenticatedEnvironment, + "shared-gate" + )) === 1 + ); + expect(oneAdmitted).toBe(true); + + /** The second run must stay queued behind the full gate. */ + await setTimeout(2000); + expect( + await engine.runQueue.currentConcurrencyOfQueue(authenticatedEnvironment, "shared-gate") + ).toBe(1); + expect( + await engine.runQueue.lengthOfQueue(authenticatedEnvironment, `task/${taskIdentifier}`) + ).toBe(1); + expect(run2.id).toBeDefined(); + } finally { + await engine.quit(); + } + } + ); +}); diff --git a/internal-packages/run-engine/src/engine/types.ts b/internal-packages/run-engine/src/engine/types.ts index d8df88c4d2a..38905c22677 100644 --- a/internal-packages/run-engine/src/engine/types.ts +++ b/internal-packages/run-engine/src/engine/types.ts @@ -23,6 +23,7 @@ import type { workerCatalog } from "./workerCatalog.js"; import { type BillingPlan } from "./billingCache.js"; import type { DRRConfig } from "../batch-queue/types.js"; import type { PendingVersionRunIdLookup } from "./services/pendingVersionLookup.js"; +import type { QueueGate } from "../run-queue/types.js"; /** * Structural mirror of the webapp's CrossSeamGuardDecision @@ -325,6 +326,8 @@ export type TriggerParams = { sdkVersion?: string; cliVersion?: string; concurrencyKey?: string; + /** Other queues this run must also hold a concurrency slot in while executing. At most two. */ + gates?: QueueGate[]; workerQueue?: string; region?: string; /** When true, the run queue may push directly to the worker queue if concurrency is available. diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index 97fb99417fc..2cbfe40c775 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -9,11 +9,11 @@ import type { MinimalAuthenticatedEnvironment } from "../shared/index.js"; * the entry is keyed) and an extra slot held until release. `queue` is the bare * queue name; the org/project/env scope comes from the run's own payload. */ -const QueueGate = z.object({ +export const QueueGate = z.object({ queue: z.string().min(1).max(128), concurrencyKey: z.string().min(1).max(128).optional(), }); -type QueueGate = z.infer; +export type QueueGate = z.infer; export const InputPayload = z.object({ runId: z.string(), diff --git a/internal-packages/run-store/src/types.ts b/internal-packages/run-store/src/types.ts index 41fecf90bb5..35670d1b92a 100644 --- a/internal-packages/run-store/src/types.ts +++ b/internal-packages/run-store/src/types.ts @@ -174,6 +174,8 @@ export type CreateRunData = { sdkVersion?: string; cliVersion?: string; concurrencyKey?: string; + /** Other queues this run must also hold a concurrency slot in while executing. */ + gates?: Array<{ queue: string; concurrencyKey?: string }>; queue: string; lockedQueueId?: string; workerQueue?: string; diff --git a/packages/core/src/v3/schemas/api.ts b/packages/core/src/v3/schemas/api.ts index 5175f94cb19..7266b8125e6 100644 --- a/packages/core/src/v3/schemas/api.ts +++ b/packages/core/src/v3/schemas/api.ts @@ -309,6 +309,15 @@ export const TriggerTaskRequestBody = z concurrencyLimit: z.number().int().optional(), }) .optional(), + gates: z + .array( + z.object({ + queue: z.string().max(128), + concurrencyKey: z.string().max(128).optional(), + }) + ) + .max(2) + .optional(), concurrencyKey: ConcurrencyKeySchema.optional(), delay: z.string().or(z.coerce.date()).optional(), idempotencyKey: z @@ -415,6 +424,15 @@ export const BatchTriggerTaskItem = z.object({ name: z.string(), }) .optional(), + gates: z + .array( + z.object({ + queue: z.string().max(128), + concurrencyKey: z.string().max(128).optional(), + }) + ) + .max(2) + .optional(), tags: RunTags.optional(), test: z.boolean().optional(), ttl: z.string().or(z.number().nonnegative().int()).optional(), diff --git a/packages/core/src/v3/schemas/resources.ts b/packages/core/src/v3/schemas/resources.ts index bc0b7a74d8c..14f4774895a 100644 --- a/packages/core/src/v3/schemas/resources.ts +++ b/packages/core/src/v3/schemas/resources.ts @@ -1,5 +1,5 @@ import { z } from "zod"; -import { QueueManifest, RetryOptions, ScheduleMetadata } from "./schemas.js"; +import { QueueGateManifest, QueueManifest, RetryOptions, ScheduleMetadata } from "./schemas.js"; import { MachineConfig } from "./common.js"; import { WebhookVerifierArtifact, @@ -19,6 +19,7 @@ export const TaskResource = z.object({ filePath: z.string(), exportName: z.string().optional(), queue: QueueManifest.extend({ name: z.string().optional() }).optional(), + gates: QueueGateManifest.array().max(2).optional(), retry: RetryOptions.optional(), machine: MachineConfig.optional(), triggerSource: z.string().optional(), diff --git a/packages/core/src/v3/schemas/schemas.ts b/packages/core/src/v3/schemas/schemas.ts index ff669193c42..0f36cb5632f 100644 --- a/packages/core/src/v3/schemas/schemas.ts +++ b/packages/core/src/v3/schemas/schemas.ts @@ -188,6 +188,17 @@ export const QueueManifest = z.object({ export type QueueManifest = z.infer; +/** A gate a task's runs must also hold a concurrency slot in while executing. + * `queue` is the gate queue's name. When `concurrencyKey` is omitted the run's own + * `concurrencyKey` is used, so the gate is keyed per tenant; a literal value pins + * the gate to one shared slot pool. */ +export const QueueGateManifest = z.object({ + queue: z.string().max(128), + concurrencyKey: z.string().max(128).optional(), +}); + +export type QueueGateManifest = z.infer; + export const ScheduleMetadata = z.object({ cron: z.string(), timezone: z.string(), @@ -203,6 +214,7 @@ const taskMetadata = { id: z.string(), description: z.string().optional(), queue: QueueManifest.extend({ name: z.string().optional() }).optional(), + gates: QueueGateManifest.array().max(2).optional(), retry: RetryOptions.optional(), machine: MachineConfig.optional(), triggerSource: z.string().optional(), diff --git a/packages/core/src/v3/types/tasks.ts b/packages/core/src/v3/types/tasks.ts index 070cc96f3df..bbf9e524d3e 100644 --- a/packages/core/src/v3/types/tasks.ts +++ b/packages/core/src/v3/types/tasks.ts @@ -222,11 +222,23 @@ type CommonTaskOptions< }); * ``` */ - queue?: { - name?: string; - concurrencyLimit?: number; - totalConcurrencyLimit?: number; - }; + queue?: + | { + name?: string; + concurrencyLimit?: number; + totalConcurrencyLimit?: number; + } + | [ + ( + | string + | { + name?: string; + concurrencyLimit?: number; + totalConcurrencyLimit?: number; + } + ), + ...QueueGateRef[], + ]; /** Configure the spec of the [machine](https://trigger.dev/docs/machines) you want your task to run on. * * @example @@ -396,6 +408,13 @@ type CommonTaskOptions< agentConfig?: { type: string }; }; +/** + * A reference to a gate: another queue a run must also hold a concurrency slot in while + * it executes. A plain string names the gate queue; the object form pins the gate to a + * literal `concurrencyKey` instead of inheriting the run's own key. + */ +export type QueueGateRef = string | { name: string; concurrencyKey?: string }; + export type TaskOptions< TIdentifier extends string, TPayload = void, @@ -809,8 +828,12 @@ export type TriggerOptions = { /** * You can override the queue for the task. If a queue doesn't exist for the given name, the run will be in the PENDING_VERSION state until the queue is created.. + * + * An array names the queue to wait in first, then up to two gates: other queues this + * run must also hold a concurrency slot in while it executes. A gate without a + * `concurrencyKey` uses the run's own `concurrencyKey`. */ - queue?: string; + queue?: string | [string, ...QueueGateRef[]]; /** * The `concurrencyKey` creates a copy of the queue for every unique value of the key. diff --git a/packages/trigger-sdk/src/v3/shared.ts b/packages/trigger-sdk/src/v3/shared.ts index b962697d559..1483bd932ce 100644 --- a/packages/trigger-sdk/src/v3/shared.ts +++ b/packages/trigger-sdk/src/v3/shared.ts @@ -75,6 +75,7 @@ import { type SchemaParseFn, type Task, type TaskIdentifier, + type QueueGateRef, type TaskOptions, type TaskOptionsWithSchema, type TaskOutput, @@ -134,6 +135,69 @@ function resolveTriggerExternalDeploymentId(explicit?: string): string | undefin }); } +type NormalizedTaskQueue = { + queue?: { name?: string; concurrencyLimit?: number; totalConcurrencyLimit?: number }; + gates?: Array<{ queue: string; concurrencyKey?: string }>; +}; + +/** + * A task's `queue` option accepts a single queue or an array of `[queue, ...gates]`. + * Normalizes both forms into the queue the run waits in plus the gate list that is + * registered on the task and carried on every trigger. + */ +function normalizeTaskQueue(queue: TaskOptions["queue"]): NormalizedTaskQueue { + if (!queue) { + return {}; + } + + if (!Array.isArray(queue)) { + return { queue }; + } + + const [home, ...gates] = queue; + + return { + queue: typeof home === "string" ? { name: home } : home, + gates: gates.map((gate) => + typeof gate === "string" + ? { queue: gate } + : { queue: gate.name, concurrencyKey: gate.concurrencyKey } + ), + }; +} + +type NormalizedTriggerQueue = { + queueName?: string; + gates?: Array<{ queue: string; concurrencyKey?: string }>; +}; + +/** + * Trigger options accept `queue` as a name or as `[name, ...gates]`; an array's gates + * replace any task-level gates for that run. + */ +function normalizeTriggerQueue( + queue: string | [string, ...QueueGateRef[]] | undefined +): NormalizedTriggerQueue { + if (!queue) { + return {}; + } + + if (typeof queue === "string") { + return { queueName: queue }; + } + + const [home, ...gates] = queue; + + return { + queueName: home, + gates: gates.map((gate) => + typeof gate === "string" + ? { queue: gate } + : { queue: gate.name, concurrencyKey: gate.concurrencyKey } + ), + }; +} + export function queue(options: QueueOptions): Queue { resourceCatalog.registerQueueMetadata(options); @@ -172,6 +236,8 @@ export function createTask< | TaskOptions | TaskOptionsWithSchema ): Task | Task { + const normalizedQueue = normalizeTaskQueue(params.queue); + const task: Task = { id: params.id, description: params.description, @@ -183,7 +249,7 @@ export function createTask< payload, undefined, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, } ); @@ -196,7 +262,7 @@ export function createTask< options, undefined, undefined, - params.queue?.name + normalizedQueue.queue?.name ); }, triggerAndWait: (payload, options, requestOptions) => { @@ -207,7 +273,7 @@ export function createTask< payload, undefined, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, }, requestOptions @@ -228,7 +294,7 @@ export function createTask< payload, undefined, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, } ) @@ -248,7 +314,7 @@ export function createTask< undefined, options, undefined, - params.queue?.name + normalizedQueue.queue?.name ); }, }; @@ -258,7 +324,8 @@ export function createTask< resourceCatalog.registerTaskMetadata({ id: params.id, description: params.description, - queue: params.queue, + queue: normalizedQueue.queue, + gates: normalizedQueue.gates, retry: params.retry ? { ...defaultRetryOptions, ...params.retry } : undefined, machine: typeof params.machine === "string" ? { preset: params.machine } : params.machine, triggerSource: params.triggerSource, @@ -271,7 +338,7 @@ export function createTask< }, }); - const queue = params.queue; + const queue = normalizedQueue.queue; if (queue && typeof queue.name === "string") { resourceCatalog.registerQueueMetadata({ @@ -327,6 +394,8 @@ export function createSchemaTask< ? getSchemaParseFn>(params.schema) : undefined; + const normalizedQueue = normalizeTaskQueue(params.queue); + const task: TaskWithSchema = { id: params.id, description: params.description, @@ -338,7 +407,7 @@ export function createSchemaTask< payload, parsePayload, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, }, requestOptions @@ -352,7 +421,7 @@ export function createSchemaTask< options, parsePayload, requestOptions, - params.queue?.name + normalizedQueue.queue?.name ); }, triggerAndWait: (payload, options) => { @@ -363,7 +432,7 @@ export function createSchemaTask< payload, parsePayload, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, } ) @@ -383,7 +452,7 @@ export function createSchemaTask< payload, parsePayload, { - queue: params.queue?.name, + queue: normalizedQueue.queue?.name, ...options, } ) @@ -403,7 +472,7 @@ export function createSchemaTask< parsePayload, options, undefined, - params.queue?.name + normalizedQueue.queue?.name ); }, }; @@ -413,7 +482,8 @@ export function createSchemaTask< resourceCatalog.registerTaskMetadata({ id: params.id, description: params.description, - queue: params.queue, + queue: normalizedQueue.queue, + gates: normalizedQueue.gates, retry: params.retry ? { ...defaultRetryOptions, ...params.retry } : undefined, machine: typeof params.machine === "string" ? { preset: params.machine } : params.machine, triggerSource: params.triggerSource, @@ -427,7 +497,7 @@ export function createSchemaTask< schema: params.schema, }); - const queue = params.queue; + const queue = normalizedQueue.queue; if (queue && typeof queue.name === "string") { resourceCatalog.registerQueueMetadata({ @@ -731,7 +801,10 @@ export async function batchTriggerById( task: item.id, payload: payloadPacket.data, options: { - queue: item.options?.queue ? { name: item.options.queue } : undefined, + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -991,7 +1064,10 @@ export async function batchTriggerByIdAndWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: item.options?.queue ? { name: item.options.queue } : undefined, + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -1256,7 +1332,10 @@ export async function batchTriggerTasks( task: item.task.id, payload: payloadPacket.data, options: { - queue: item.options?.queue ? { name: item.options.queue } : undefined, + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -1521,7 +1600,10 @@ export async function batchTriggerAndWaitTasks( task: item.id, payload: payloadPacket.data, options: { - queue: item.options?.queue ? { name: item.options.queue } : undefined, + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2063,7 +2148,10 @@ async function* transformBatchItemsStreamForWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: item.options?.queue ? { name: item.options.queue } : undefined, + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2113,7 +2201,10 @@ async function* transformBatchByTaskItemsStream( task: taskIdentifier, payload: payloadPacket.data, options: { - queue: item.options?.queue - ? { name: item.options.queue } + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } : queue ? { name: queue } : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2280,11 +2375,12 @@ async function* transformSingleTaskBatchItemsStreamForWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: item.options?.queue - ? { name: item.options.queue } + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } : queue ? { name: queue } : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2334,7 +2430,10 @@ async function trigger_internal( { payload: triggerPayloadPacket.data, options: { - queue: options?.queue ? { name: options.queue } : undefined, + queue: normalizeTriggerQueue(options?.queue).queueName + ? { name: normalizeTriggerQueue(options?.queue).queueName! } + : undefined, + gates: normalizeTriggerQueue(options?.queue).gates, concurrencyKey: options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: triggerPayloadPacket.dataType, @@ -2419,11 +2518,12 @@ async function batchTrigger_internal( task: taskIdentifier, payload: payloadPacket.data, options: { - queue: item.options?.queue - ? { name: item.options.queue } + queue: normalizeTriggerQueue(item.options?.queue).queueName + ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } : queue ? { name: queue } : undefined, + gates: normalizeTriggerQueue(item.options?.queue).gates, concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2603,7 +2703,10 @@ async function triggerAndWait_internal Date: Sat, 29 Aug 2026 10:42:31 +0100 Subject: [PATCH 2/3] fix(core,sdk,webapp): bound gate tuples at two and normalize queues once The queue tuple types now reject more than two gate entries at compile time, matching the wire schemas. Trigger body construction normalizes the queue option once per call site through a shared helper instead of three times, and cached task metadata runs gates through the same parser as database reads. --- .../app/services/taskMetadataCache.server.ts | 2 +- packages/core/src/v3/types/tasks.ts | 28 ++--- packages/trigger-sdk/src/v3/shared.ts | 103 ++++++------------ 3 files changed, 47 insertions(+), 86 deletions(-) diff --git a/apps/webapp/app/services/taskMetadataCache.server.ts b/apps/webapp/app/services/taskMetadataCache.server.ts index e420329d74d..dcf742aa2fe 100644 --- a/apps/webapp/app/services/taskMetadataCache.server.ts +++ b/apps/webapp/app/services/taskMetadataCache.server.ts @@ -109,7 +109,7 @@ function decode(slug: string, raw: string): TaskMetadataEntry | null { triggerSource: parsed.k, queueId: parsed.q, queueName: parsed.n, - gates: parsed.g ?? null, + gates: parseTaskGates(parsed.g ?? null), }; } catch (error) { logger.error("Failed to decode task metadata cache entry", { slug, error }); diff --git a/packages/core/src/v3/types/tasks.ts b/packages/core/src/v3/types/tasks.ts index bbf9e524d3e..28afb2c5314 100644 --- a/packages/core/src/v3/types/tasks.ts +++ b/packages/core/src/v3/types/tasks.ts @@ -223,22 +223,10 @@ type CommonTaskOptions< * ``` */ queue?: - | { - name?: string; - concurrencyLimit?: number; - totalConcurrencyLimit?: number; - } - | [ - ( - | string - | { - name?: string; - concurrencyLimit?: number; - totalConcurrencyLimit?: number; - } - ), - ...QueueGateRef[], - ]; + | TaskQueueIn + | [TaskQueueIn | string] + | [TaskQueueIn | string, QueueGateRef] + | [TaskQueueIn | string, QueueGateRef, QueueGateRef]; /** Configure the spec of the [machine](https://trigger.dev/docs/machines) you want your task to run on. * * @example @@ -415,6 +403,12 @@ type CommonTaskOptions< */ export type QueueGateRef = string | { name: string; concurrencyKey?: string }; +type TaskQueueIn = { + name?: string; + concurrencyLimit?: number; + totalConcurrencyLimit?: number; +}; + export type TaskOptions< TIdentifier extends string, TPayload = void, @@ -833,7 +827,7 @@ export type TriggerOptions = { * run must also hold a concurrency slot in while it executes. A gate without a * `concurrencyKey` uses the run's own `concurrencyKey`. */ - queue?: string | [string, ...QueueGateRef[]]; + queue?: string | [string] | [string, QueueGateRef] | [string, QueueGateRef, QueueGateRef]; /** * The `concurrencyKey` creates a copy of the queue for every unique value of the key. diff --git a/packages/trigger-sdk/src/v3/shared.ts b/packages/trigger-sdk/src/v3/shared.ts index 1483bd932ce..42f0f465f4f 100644 --- a/packages/trigger-sdk/src/v3/shared.ts +++ b/packages/trigger-sdk/src/v3/shared.ts @@ -198,6 +198,26 @@ function normalizeTriggerQueue( }; } +/** + * Builds the queue and gates fields of a trigger request body from the `queue` + * option, normalizing once per call site. + */ +function triggerQueueBody( + queue: Parameters[0], + fallbackQueueName?: string +): { + queue?: { name: string }; + gates?: Array<{ queue: string; concurrencyKey?: string }>; +} { + const normalized = normalizeTriggerQueue(queue); + const name = normalized.queueName ?? fallbackQueueName; + + return { + queue: name ? { name } : undefined, + gates: normalized.gates, + }; +} + export function queue(options: QueueOptions): Queue { resourceCatalog.registerQueueMetadata(options); @@ -801,10 +821,7 @@ export async function batchTriggerById( task: item.id, payload: payloadPacket.data, options: { - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -1064,10 +1081,7 @@ export async function batchTriggerByIdAndWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -1332,10 +1346,7 @@ export async function batchTriggerTasks( task: item.task.id, payload: payloadPacket.data, options: { - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -1600,10 +1611,7 @@ export async function batchTriggerAndWaitTasks( task: item.id, payload: payloadPacket.data, options: { - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2148,10 +2153,7 @@ async function* transformBatchItemsStreamForWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2201,10 +2203,7 @@ async function* transformBatchByTaskItemsStream( task: taskIdentifier, payload: payloadPacket.data, options: { - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : queue - ? { name: queue } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue, queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2375,12 +2366,7 @@ async function* transformSingleTaskBatchItemsStreamForWait( payload: payloadPacket.data, options: { lockToVersion: taskContext.worker?.version, - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : queue - ? { name: queue } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue, queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2430,10 +2416,7 @@ async function trigger_internal( { payload: triggerPayloadPacket.data, options: { - queue: normalizeTriggerQueue(options?.queue).queueName - ? { name: normalizeTriggerQueue(options?.queue).queueName! } - : undefined, - gates: normalizeTriggerQueue(options?.queue).gates, + ...triggerQueueBody(options?.queue), concurrencyKey: options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: triggerPayloadPacket.dataType, @@ -2518,12 +2501,7 @@ async function batchTrigger_internal( task: taskIdentifier, payload: payloadPacket.data, options: { - queue: normalizeTriggerQueue(item.options?.queue).queueName - ? { name: normalizeTriggerQueue(item.options?.queue).queueName! } - : queue - ? { name: queue } - : undefined, - gates: normalizeTriggerQueue(item.options?.queue).gates, + ...triggerQueueBody(item.options?.queue, queue), concurrencyKey: item.options?.concurrencyKey, test: taskContext.ctx?.run.isTest, payloadType: payloadPacket.dataType, @@ -2703,10 +2681,7 @@ async function triggerAndWait_internal Date: Sat, 29 Aug 2026 10:57:15 +0100 Subject: [PATCH 3/3] fix(core,database): mirror gates into the run-ops schema and reject empty gate names The dedicated run-ops TaskRun schema gains the same nullable gates column and migration so run creation keeps working with run-operations splitting enabled (the schema parity test covers it). Gate names and keys are also required to be non-empty everywhere they enter, so a configured gate can never be silently dropped by sanitization. --- .../20260829120000_add_task_run_gates/migration.sql | 2 ++ internal-packages/run-ops-database/prisma/schema.prisma | 4 ++++ packages/core/src/v3/schemas/api.ts | 8 ++++---- packages/core/src/v3/schemas/schemas.ts | 4 ++-- 4 files changed, 12 insertions(+), 6 deletions(-) create mode 100644 internal-packages/run-ops-database/prisma/migrations/20260829120000_add_task_run_gates/migration.sql diff --git a/internal-packages/run-ops-database/prisma/migrations/20260829120000_add_task_run_gates/migration.sql b/internal-packages/run-ops-database/prisma/migrations/20260829120000_add_task_run_gates/migration.sql new file mode 100644 index 00000000000..72f07a937f3 --- /dev/null +++ b/internal-packages/run-ops-database/prisma/migrations/20260829120000_add_task_run_gates/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "TaskRun" ADD COLUMN "gates" JSONB; diff --git a/internal-packages/run-ops-database/prisma/schema.prisma b/internal-packages/run-ops-database/prisma/schema.prisma index 0a6c10dd034..b7a1fc77e4c 100644 --- a/internal-packages/run-ops-database/prisma/schema.prisma +++ b/internal-packages/run-ops-database/prisma/schema.prisma @@ -153,6 +153,10 @@ model TaskRun { concurrencyKey String? + /// Gates for this run: other queues it must also hold a concurrency slot in while + /// executing. Shape: [{ queue: string, concurrencyKey?: string }] + gates Json? + delayUntil DateTime? queuedAt DateTime? ttl String? diff --git a/packages/core/src/v3/schemas/api.ts b/packages/core/src/v3/schemas/api.ts index 7266b8125e6..df089a60126 100644 --- a/packages/core/src/v3/schemas/api.ts +++ b/packages/core/src/v3/schemas/api.ts @@ -312,8 +312,8 @@ export const TriggerTaskRequestBody = z gates: z .array( z.object({ - queue: z.string().max(128), - concurrencyKey: z.string().max(128).optional(), + queue: z.string().min(1).max(128), + concurrencyKey: z.string().min(1).max(128).optional(), }) ) .max(2) @@ -427,8 +427,8 @@ export const BatchTriggerTaskItem = z.object({ gates: z .array( z.object({ - queue: z.string().max(128), - concurrencyKey: z.string().max(128).optional(), + queue: z.string().min(1).max(128), + concurrencyKey: z.string().min(1).max(128).optional(), }) ) .max(2) diff --git a/packages/core/src/v3/schemas/schemas.ts b/packages/core/src/v3/schemas/schemas.ts index 0f36cb5632f..8993b2125ea 100644 --- a/packages/core/src/v3/schemas/schemas.ts +++ b/packages/core/src/v3/schemas/schemas.ts @@ -193,8 +193,8 @@ export type QueueManifest = z.infer; * `concurrencyKey` is used, so the gate is keyed per tenant; a literal value pins * the gate to one shared slot pool. */ export const QueueGateManifest = z.object({ - queue: z.string().max(128), - concurrencyKey: z.string().max(128).optional(), + queue: z.string().min(1).max(128), + concurrencyKey: z.string().min(1).max(128).optional(), }); export type QueueGateManifest = z.infer;