Skip to content
18 changes: 18 additions & 0 deletions .changeset/queue-concurrency-overrides.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
---
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
---

Adjust queue concurrency at runtime, per key and in total. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideTotalConcurrencyLimit` and `queues.resetTotalConcurrencyLimit` adjust the cap across all keys.

```ts
import { queues } from "@trigger.dev/sdk";

await queues.overrideConcurrencyLimit("my-queue", 20, { concurrencyKey: "tenant-123" });
await queues.resetConcurrencyLimit("my-queue", { concurrencyKey: "tenant-123" });

await queues.overrideTotalConcurrencyLimit("my-queue", 100);
await queues.resetTotalConcurrencyLimit("my-queue");
```

Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with total concurrency limits enabled.
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import { json } from "@remix-run/server-runtime";
import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3";
import { z } from "zod";
import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server";

const BodySchema = z.object({
type: RetrieveQueueType.default("id"),
concurrencyKey: z.string().min(1).max(128),
concurrencyLimit: z.number().int().min(0).max(100000),
});

const route = createActionApiRoute(
{
body: BodySchema,
params: z.object({
queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")),
}),
authorization: {
action: "write",
resource: () => ({ type: "queues" }),
},
},
async ({ params, body, authentication }) => {
const input: RetrieveQueueParam =
body.type === "id"
? params.queueParam
: {
type: body.type,
name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"),
};

return concurrencySystem.queues
.overrideConcurrencyKeyLimit(
authentication.environment,
input,
body.concurrencyKey,
body.concurrencyLimit
)
.match(
(queue) => {
return json(
toQueueItem({
friendlyId: queue.friendlyId,
name: queue.name,
type: queue.type,
running: queue.running,
queued: queue.queued,
concurrencyLimit: queue.concurrencyLimit,
concurrencyLimitBase: queue.concurrencyLimitBase,
concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt,
concurrencyLimitOverriddenBy: null,
paused: queue.paused,
}),
{ status: 200 }
);
Comment thread
matt-aitken marked this conversation as resolved.
},
(error) => {
switch (error.type) {
case "queue_not_found": {
return json({ error: "Queue not found" }, { status: 404 });
}
case "invalid_override":
case "concurrency_limit_exceeds_maximum":
case "too_many_key_overrides": {
return json({ error: error.message }, { status: 400 });
}
case "queue_update_failed": {
return json(
{ error: "Failed to update queue concurrency key limit" },
{ status: 500 }
);
}
case "sync_queue_concurrency_to_engine_failed": {
return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 });
}
case "get_queue_stats_failed": {
return json({ error: "Failed to read queue stats" }, { status: 500 });
}
case "other": {
return json(
{ error: "Failed to update queue concurrency key limit" },
{
status: 500,
}
);
}
default: {
return json(
{ error: "Failed to update queue concurrency key limit" },
{
status: 500,
}
);
}
}
}
);
}
);

export const { action } = route;
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import { json } from "@remix-run/server-runtime";
import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3";
import { z } from "zod";
import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server";

const BodySchema = z.object({
type: RetrieveQueueType.default("id"),
concurrencyKey: z.string().min(1).max(128),
});

const route = createActionApiRoute(
{
body: BodySchema,
params: z.object({
queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")),
}),
authorization: {
action: "write",
resource: () => ({ type: "queues" }),
},
},
async ({ params, body, authentication }) => {
const input: RetrieveQueueParam =
body.type === "id"
? params.queueParam
: {
type: body.type,
name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"),
};

return concurrencySystem.queues
.resetConcurrencyKeyLimit(authentication.environment, input, body.concurrencyKey)
.match(
(queue) => {
return json(
toQueueItem({
friendlyId: queue.friendlyId,
name: queue.name,
type: queue.type,
running: queue.running,
queued: queue.queued,
concurrencyLimit: queue.concurrencyLimit,
concurrencyLimitBase: queue.concurrencyLimitBase,
concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt,
concurrencyLimitOverriddenBy: null,
paused: queue.paused,
}),
{ status: 200 }
);
},
(error) => {
switch (error.type) {
case "queue_not_found": {
return json({ error: "Queue not found" }, { status: 404 });
}
case "queue_not_overridden": {
return json(
{ error: "This concurrency key does not have an override" },
{ status: 400 }
);
}
case "queue_update_failed": {
return json({ error: "Failed to reset the concurrency key limit" }, { status: 500 });
}
case "sync_queue_concurrency_to_engine_failed": {
return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 });
}
case "get_queue_stats_failed": {
return json({ error: "Failed to read queue stats" }, { status: 500 });
}
case "other": {
return json(
{ error: "Failed to reset the concurrency key limit" },
{
status: 500,
}
);
}
default: {
return json(
{ error: "Failed to reset the concurrency key limit" },
{
status: 500,
}
);
}
}
}
);
}
);

export const { action } = route;
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
import { json } from "@remix-run/server-runtime";
import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3";
import { z } from "zod";
import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server";
import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server";
import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server";

const BodySchema = z.object({
type: RetrieveQueueType.default("id"),
concurrencyLimit: z.number().int().min(0).max(100000),
});

const route = createActionApiRoute(
{
body: BodySchema,
params: z.object({
queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")),
}),
authorization: {
action: "write",
resource: () => ({ type: "queues" }),
},
},
async ({ params, body, authentication }) => {
const input: RetrieveQueueParam =
body.type === "id"
? params.queueParam
: {
type: body.type,
name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"),
};

return concurrencySystem.queues
.overrideTotalConcurrencyLimit(authentication.environment, input, body.concurrencyLimit)
.match(
(queue) => {
return json(
toQueueItem({
friendlyId: queue.friendlyId,
name: queue.name,
type: queue.type,
running: queue.running,
queued: queue.queued,
concurrencyLimit: queue.concurrencyLimit,
concurrencyLimitBase: queue.concurrencyLimitBase,
concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt,
concurrencyLimitOverriddenBy: null,
paused: queue.paused,
}),
{ status: 200 }
);
},
(error) => {
switch (error.type) {
case "queue_not_found": {
return json({ error: "Queue not found" }, { status: 404 });
}
case "invalid_override":
case "concurrency_limit_exceeds_maximum": {
return json({ error: error.message }, { status: 400 });
}
case "queue_update_failed": {
return json(
{ error: "Failed to update queue total concurrency limit" },
{ status: 500 }
);
}
case "sync_queue_concurrency_to_engine_failed": {
return json({ error: "Failed to sync the total concurrency limit" }, { status: 500 });
}
case "get_queue_stats_failed": {
return json({ error: "Failed to read queue stats" }, { status: 500 });
}
case "other": {
return json(
{ error: "Failed to update queue total concurrency limit" },
{
status: 500,
}
);
}
default: {
return json(
{ error: "Failed to update queue total concurrency limit" },
{
status: 500,
}
);
}
}
}
);
}
);

export const { action } = route;
Loading
Loading