From 51bb4c887a3e4f546cf9d2500fd1b59295bacfd0 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Fri, 31 May 2024 11:24:13 +0100 Subject: [PATCH] Support custom queue when triggering a task (#1138) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Created a v3-catalog test script for queues * SDK: Fix for calling trigger and passing a custom queue * Support custom queue in TriggerTaskService * Improved the script in the catalog so it’s clearer what’s going on * Remove the concurrencyLimit from a queue if the limit is null * Fix for the test code… stupid --- .changeset/itchy-chairs-itch.md | 5 ++ apps/webapp/app/v3/marqs/index.server.ts | 4 + .../services/createBackgroundWorker.server.ts | 2 + .../app/v3/services/triggerTask.server.ts | 48 +++++++++++- packages/trigger-sdk/src/v3/shared.ts | 4 +- references/v3-catalog/package.json | 5 +- references/v3-catalog/src/queues.ts | 73 +++++++++++++++++++ 7 files changed, 134 insertions(+), 7 deletions(-) create mode 100644 .changeset/itchy-chairs-itch.md create mode 100644 references/v3-catalog/src/queues.ts diff --git a/.changeset/itchy-chairs-itch.md b/.changeset/itchy-chairs-itch.md new file mode 100644 index 000000000..6d01a6dee --- /dev/null +++ b/.changeset/itchy-chairs-itch.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Fix for calling trigger and passing a custom queue diff --git a/apps/webapp/app/v3/marqs/index.server.ts b/apps/webapp/app/v3/marqs/index.server.ts index 1e369c685..e24077e01 100644 --- a/apps/webapp/app/v3/marqs/index.server.ts +++ b/apps/webapp/app/v3/marqs/index.server.ts @@ -82,6 +82,10 @@ export class MarQS { return this.redis.set(this.keys.queueConcurrencyLimitKey(env, queue), concurrency); } + public async removeQueueConcurrencyLimits(env: AuthenticatedEnvironment, queue: string) { + return this.redis.del(this.keys.queueConcurrencyLimitKey(env, queue)); + } + public async updateEnvConcurrencyLimits(env: AuthenticatedEnvironment) { await this.#callUpdateGlobalConcurrencyLimits({ envConcurrencyLimitKey: this.keys.envConcurrencyLimitKey(env), diff --git a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index 9045075e1..793980222 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -179,6 +179,8 @@ export async function createBackgroundTasks( taskQueue.name, taskQueue.concurrencyLimit ); + } else { + await marqs?.removeQueueConcurrencyLimits(environment, taskQueue.name); } } catch (error) { if (error instanceof Prisma.PrismaClientKnownRequestError) { diff --git a/apps/webapp/app/v3/services/triggerTask.server.ts b/apps/webapp/app/v3/services/triggerTask.server.ts index 1d6a59443..9c5c0dba6 100644 --- a/apps/webapp/app/v3/services/triggerTask.server.ts +++ b/apps/webapp/app/v3/services/triggerTask.server.ts @@ -5,11 +5,11 @@ import { packetRequiresOffloading, } from "@trigger.dev/core/v3"; import { createHash } from "node:crypto"; -import { $transaction } from "~/db.server"; +import { $transaction, prisma } from "~/db.server"; import { AuthenticatedEnvironment } from "~/services/apiAuth.server"; +import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server"; import { eventRepository } from "../eventRepository.server"; import { generateFriendlyId } from "../friendlyIdentifiers"; -import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server"; import { uploadToObjectStore } from "../r2.server"; import { BaseService } from "./baseService.server"; @@ -111,7 +111,12 @@ export class TriggerTaskService extends BaseService { select: { lastNumber: true }, }); - const queueName = sanitizeQueueName(body.options?.queue?.name ?? `task/${taskId}`); + let queueName = sanitizeQueueName(body.options?.queue?.name ?? `task/${taskId}`); + + // Check that the queuename is not an empty string + if (!queueName) { + queueName = sanitizeQueueName(`task/${taskId}`); + } event.setAttribute("queueName", queueName); span.setAttribute("queueName", queueName); @@ -182,6 +187,43 @@ export class TriggerTaskService extends BaseService { } } + if (body.options?.queue) { + const concurrencyLimit = body.options.queue.concurrencyLimit + ? Math.max(0, body.options.queue.concurrencyLimit) + : null; + const taskQueue = await prisma.taskQueue.upsert({ + where: { + runtimeEnvironmentId_name: { + runtimeEnvironmentId: environment.id, + name: queueName, + }, + }, + update: { + concurrencyLimit, + rateLimit: body.options.queue.rateLimit, + }, + create: { + friendlyId: generateFriendlyId("queue"), + name: queueName, + concurrencyLimit, + runtimeEnvironmentId: environment.id, + projectId: environment.projectId, + rateLimit: body.options.queue.rateLimit, + type: "NAMED", + }, + }); + + if (typeof taskQueue.concurrencyLimit === "number") { + await marqs?.updateQueueConcurrencyLimits( + environment, + taskQueue.name, + taskQueue.concurrencyLimit + ); + } else { + await marqs?.removeQueueConcurrencyLimits(environment, taskQueue.name); + } + } + return taskRun; }); diff --git a/packages/trigger-sdk/src/v3/shared.ts b/packages/trigger-sdk/src/v3/shared.ts index ff1121f6a..2499dce68 100644 --- a/packages/trigger-sdk/src/v3/shared.ts +++ b/packages/trigger-sdk/src/v3/shared.ts @@ -353,7 +353,7 @@ export function createTask