Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
提交
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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 combined. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideCombinedConcurrencyLimit` and `queues.resetCombinedConcurrencyLimit` 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.overrideCombinedConcurrencyLimit("my-queue", 100);
await queues.resetCombinedConcurrencyLimit("my-queue");
```

Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with combined concurrency limits enabled.
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
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.action;
/** The builder's loader answers non-POST methods with a 405. */
export const loader = route.loader;
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
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"),
});

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
.resetTotalConcurrencyLimit(authentication.environment, input)
.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: "The queue total concurrency limit is not overridden" },
{ status: 400 }
);
}
case "queue_update_failed": {
return json(
{ error: "Failed to reset the 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 reset the queue total concurrency limit" },
{
status: 500,
}
);
}
default: {
return json(
{ error: "Failed to reset the queue total concurrency limit" },
{
status: 500,
}
);
}
}
}
);
}
);

export const action = route.action;
/** The builder's loader answers non-POST methods with a 405. */
export const loader = route.loader;
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
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.action;
/** The builder's loader answers non-POST methods with a 405. */
export const loader = route.loader;
Loading
Loading