Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
15 changes: 15 additions & 0 deletions .changeset/queue-concurrency-overrides.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
---
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
---

Adjust a queue's combined concurrency limit at runtime. `queues.overrideCombinedConcurrencyLimit` raises or lowers the cap on concurrent runs across all of a queue's `concurrencyKey` values, and `queues.resetCombinedConcurrencyLimit` reverts to the declared configuration.

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

await queues.overrideCombinedConcurrencyLimit("my-queue", 100);
await queues.resetCombinedConcurrencyLimit("my-queue");
```

Overrides survive deploys. 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;
152 changes: 151 additions & 1 deletion apps/webapp/app/v3/services/concurrencySystem.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,12 @@ import { errAsync, fromPromise, okAsync } from "neverthrow";
import type { PrismaClientOrTransaction } from "~/db.server";
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { removeQueueConcurrencyLimits, updateQueueConcurrencyLimits } from "../runQueue.server";
import {
removeQueueConcurrencyLimits,
removeQueueTotalConcurrencyLimits,
updateQueueConcurrencyLimits,
updateQueueTotalConcurrencyLimits,
} from "../runQueue.server";
import { engine } from "../runEngine.server";

export type ConcurrencySystemOptions = {
Expand Down Expand Up @@ -77,6 +82,32 @@ export class ConcurrencySystem {
.andThen((queue) => syncQueueConcurrencyToEngine(environment, queue))
.andThen((queue) => getQueueStats(environment, queue));
},
overrideTotalConcurrencyLimit: (
environment: AuthenticatedEnvironment,
queue: QueueInput,
totalConcurrencyLimit: number,
overriddenBy?: User
) => {
return findQueueFromInput(this.db, environment, queue)
.andThen((queue) =>
overrideQueueTotalConcurrencyLimit(
this.db,
environment,
queue,
totalConcurrencyLimit,
overriddenBy
)
)
.andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue))
.andThen((queue) => getQueueStats(environment, queue));
},
resetTotalConcurrencyLimit: (environment: AuthenticatedEnvironment, queue: QueueInput) => {
return findQueueFromInput(this.db, environment, queue)
.andThen((queue) => syncQueueTotalConcurrencyResetToEngine(environment, queue))
.andThen((queue) => resetQueueTotalConcurrencyLimit(this.db, queue))
.andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue))
.andThen((queue) => getQueueStats(environment, queue));
},
/**
* Recalculates the materialized limit of every percent-based override in the environment
* against its CURRENT maximumConcurrencyLimit and syncs changed queues to the run engine.
Expand Down Expand Up @@ -316,6 +347,125 @@ function syncQueueConcurrencyToEngine(environment: AuthenticatedEnvironment, que
}
}

function overrideQueueTotalConcurrencyLimit(
db: PrismaClientOrTransaction,
environment: AuthenticatedEnvironment,
queue: TaskQueue,
totalConcurrencyLimit: number,
overriddenBy?: User
) {
const maximum = environment.maximumConcurrencyLimit;

if (!Number.isFinite(totalConcurrencyLimit) || totalConcurrencyLimit < 0) {
return errAsync({
type: "invalid_override" as const,
message: "Combined concurrency limit must be a non-negative number",
});
}

if (totalConcurrencyLimit > maximum) {
return errAsync({
type: "concurrency_limit_exceeds_maximum" as const,
message: `Combined concurrency limit (${totalConcurrencyLimit}) cannot exceed the environment limit (${maximum})`,
});
}

const totalConcurrencyLimitBase = queue.totalConcurrencyLimitOverriddenAt
? queue.totalConcurrencyLimitBase
: queue.totalConcurrencyLimit;

return fromPromise(
db.taskQueue.update({
where: { id: queue.id },
data: {
totalConcurrencyLimit,
totalConcurrencyLimitBase: totalConcurrencyLimitBase ?? null,
totalConcurrencyLimitOverriddenAt: new Date(),
totalConcurrencyLimitOverriddenBy: overriddenBy?.id ?? null,
},
}),
(error) => ({
type: "queue_update_failed" as const,
cause: error,
})
);
}

/**
* Enforce first, then persist: syncs the engine to the declared base BEFORE clearing
* the override marker, so an engine failure leaves the marker set and a retry
* converges instead of being rejected while the overridden limit stays enforced.
*/
function syncQueueTotalConcurrencyResetToEngine(
environment: AuthenticatedEnvironment,
queue: TaskQueue
) {
if (queue.totalConcurrencyLimitOverriddenAt === null) {
return errAsync({ type: "queue_not_overridden" as const });
}

if (typeof queue.totalConcurrencyLimitBase === "number") {
return fromPromise(
updateQueueTotalConcurrencyLimits(environment, queue.name, queue.totalConcurrencyLimitBase),
(error) => ({
type: "sync_queue_concurrency_to_engine_failed" as const,
cause: error,
})
).andThen(() => okAsync(queue));
}

return fromPromise(removeQueueTotalConcurrencyLimits(environment, queue.name), (error) => ({
type: "sync_queue_concurrency_to_engine_failed" as const,
cause: error,
})).andThen(() => okAsync(queue));
}

function resetQueueTotalConcurrencyLimit(db: PrismaClientOrTransaction, queue: TaskQueue) {
if (queue.totalConcurrencyLimitOverriddenAt === null) {
return errAsync({ type: "queue_not_overridden" as const });
}

return fromPromise(
db.taskQueue.update({
where: { id: queue.id },
data: {
totalConcurrencyLimit: queue.totalConcurrencyLimitBase,
totalConcurrencyLimitBase: null,
totalConcurrencyLimitOverriddenAt: null,
totalConcurrencyLimitOverriddenBy: null,
},
}),
(error) => ({
type: "queue_update_failed" as const,
cause: error,
})
);
}

/**
* The total limit key is separate from the per-queue limit key that pause zeroes,
* so it syncs regardless of the paused state.
*/
function syncQueueTotalConcurrencyToEngine(
environment: AuthenticatedEnvironment,
queue: TaskQueue
) {
if (typeof queue.totalConcurrencyLimit === "number") {
return fromPromise(
updateQueueTotalConcurrencyLimits(environment, queue.name, queue.totalConcurrencyLimit),
(error) => ({
type: "sync_queue_concurrency_to_engine_failed" as const,
cause: error,
})
).andThen(() => okAsync(queue));
}

return fromPromise(removeQueueTotalConcurrencyLimits(environment, queue.name), (error) => ({
type: "sync_queue_concurrency_to_engine_failed" as const,
cause: error,
})).andThen(() => okAsync(queue));
}

function getQueueStats(environment: AuthenticatedEnvironment, queue: TaskQueue) {
return fromPromise(
Promise.all([
Expand Down
Loading
Loading