Files
triggerdotdev--trigger.dev/apps/webapp/app/v3/services/executeTasksWaitingForDeploy.ts
Matt Aitken da6ce3c8d5 Concurrency page and more accurate tracking (#1252)
* Initial TaskRunConcurrencyTracker implementation

* MARQS calls a subscriber to events

* When enqueuing add the extra required metadata

* Track concurrency per environment for tasks too

* Admin page for global concurrency

* Use the new concurrency tracker on the tasks page

* Useful performance test task

* getAllTaskIdentifiers()

* New page for concurrency

* BackgroundWorkerTask index for quick lookup of task identifiers

* Added a way to get concurrency for environments

* Added upgrade/request more concurrency button

* Queued task column working

* Use defer and suspense

* Added queue column to the concurrency environments table

* Some comments added for clarity

* Fixed bad log message

* Sidemenu: move lower and rename to “Concurrency limits”

* Only show the environments, not tasks. Renamed to “Concurrency limits”
2024-08-13 11:43:46 +01:00

126 lines
3.3 KiB
TypeScript

import { PrismaClientOrTransaction } from "~/db.server";
import { workerQueue } from "~/services/worker.server";
import { marqs } from "~/v3/marqs/index.server";
import { BaseService } from "./baseService.server";
import { logger } from "~/services/logger.server";
export class ExecuteTasksWaitingForDeployService extends BaseService {
public async call(backgroundWorkerId: string) {
const backgroundWorker = await this._prisma.backgroundWorker.findFirst({
where: {
id: backgroundWorkerId,
},
include: {
runtimeEnvironment: {
include: {
project: true,
organization: true,
},
},
tasks: true,
},
});
if (!backgroundWorker) {
logger.error("Background worker not found", { id: backgroundWorkerId });
return;
}
const runsWaitingForDeploy = await this._prisma.taskRun.findMany({
where: {
runtimeEnvironmentId: backgroundWorker.runtimeEnvironmentId,
projectId: backgroundWorker.projectId,
status: "WAITING_FOR_DEPLOY",
taskIdentifier: {
in: backgroundWorker.tasks.map((task) => task.slug),
},
},
orderBy: {
number: "asc",
},
});
if (!runsWaitingForDeploy.length) {
return;
}
// Clear any runs awaiting deployment for execution
const pendingRuns = await this._prisma.taskRun.updateMany({
where: {
id: {
in: runsWaitingForDeploy.map((run) => run.id),
},
},
data: {
status: "PENDING",
},
});
if (pendingRuns.count) {
logger.debug("Task runs waiting for deploy are now ready for execution", {
tasks: runsWaitingForDeploy.map((run) => run.id),
total: pendingRuns.count,
});
}
if (!marqs) {
return;
}
const enqueues: Promise<any>[] = [];
let i = 0;
for (const run of runsWaitingForDeploy) {
enqueues.push(
marqs.enqueueMessage(
backgroundWorker.runtimeEnvironment,
run.queue,
run.id,
{
type: "EXECUTE",
taskIdentifier: run.taskIdentifier,
projectId: backgroundWorker.runtimeEnvironment.projectId,
environmentId: backgroundWorker.runtimeEnvironment.id,
environmentType: backgroundWorker.runtimeEnvironment.type,
},
run.concurrencyKey ?? undefined,
Date.now() + i * 5 // slight delay to help preserve order
)
);
i++;
}
const settled = await Promise.allSettled(enqueues);
if (settled.some((s) => s.status === "rejected")) {
const rejectedRuns: { id: string; reason: any }[] = [];
runsWaitingForDeploy.forEach((run, i) => {
if (settled[i].status === "rejected") {
const rejected = settled[i] as PromiseRejectedResult;
rejectedRuns.push({ id: run.id, reason: rejected.reason });
}
});
logger.error("Failed to requeue task runs for immediate execution", {
rejectedRuns,
});
}
}
static async enqueue(backgroundWorkerId: string, tx: PrismaClientOrTransaction, runAt?: Date) {
return await workerQueue.enqueue(
"v3.executeTasksWaitingForDeploy",
{
backgroundWorkerId,
},
{
tx,
runAt,
}
);
}
}