Merge branch 'main' into v3/worker-attempt-creation

This commit is contained in:
nicktrn
2024-05-21 10:31:34 +01:00
3 changed files with 87 additions and 14 deletions
@@ -413,6 +413,7 @@ export class DevQueueConsumer {
data: {
lockedAt: new Date(),
lockedById: backgroundTask.id,
status: "EXECUTING",
lockedToVersionId: backgroundWorker.id,
},
include: {
@@ -26,7 +26,11 @@ import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { EnvironmentVariablesRepository } from "../environmentVariables/environmentVariablesRepository.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { socketIo } from "../handleSocketIo.server";
import { findCurrentWorkerDeployment } from "../models/workerDeployment.server";
import {
findCurrentWorkerDeployment,
getWorkerDeploymentFromWorker,
getWorkerDeploymentFromWorkerTask,
} from "../models/workerDeployment.server";
import { RestoreCheckpointService } from "../services/restoreCheckpoint.server";
import { SEMINTATTRS_FORCE_RECORDING, tracer } from "../tracer.server";
import { CrashTaskRunService } from "../services/crashTaskRun.server";
@@ -317,11 +321,11 @@ export class SharedQueueConsumer {
return;
}
const deployment = existingTaskRun.lockedToVersion?.deployment
? {
...existingTaskRun.lockedToVersion.deployment,
worker: existingTaskRun.lockedToVersion,
}
// Check if the task run is locked to a specific worker, if not, use the current worker deployment
const deployment = existingTaskRun.lockedById
? await getWorkerDeploymentFromWorkerTask(existingTaskRun.lockedById)
: existingTaskRun.lockedToVersionId
? await getWorkerDeploymentFromWorker(existingTaskRun.lockedToVersionId)
: await findCurrentWorkerDeployment(existingTaskRun.runtimeEnvironmentId);
if (!deployment || !deployment.worker) {
@@ -1,16 +1,30 @@
import type { Prettify } from "@trigger.dev/core";
import { CURRENT_DEPLOYMENT_LABEL } from "~/consts";
import { prisma } from "~/db.server";
import { Prisma, prisma } from "~/db.server";
export type CurrentWorkerDeployment = Prettify<NonNullable<Awaited<ReturnType<typeof findCurrentWorkerDeployment>>>>;
export type CurrentWorkerDeployment = Prettify<
NonNullable<Awaited<ReturnType<typeof findCurrentWorkerDeployment>>>
>;
export async function findCurrentWorkerDeployment(environmentId: string) {
type WorkerDeploymentWithWorkerTasks = Prisma.WorkerDeploymentGetPayload<{
include: {
worker: {
include: {
tasks: true;
};
};
};
}>;
export async function findCurrentWorkerDeployment(
environmentId: string
): Promise<WorkerDeploymentWithWorkerTasks | undefined> {
const promotion = await prisma.workerDeploymentPromotion.findUnique({
where: {
environmentId_label: {
environmentId,
label: CURRENT_DEPLOYMENT_LABEL,
}
},
},
include: {
deployment: {
@@ -20,10 +34,64 @@ export async function findCurrentWorkerDeployment(environmentId: string) {
tasks: true,
},
},
}
}
}
},
},
},
});
return promotion?.deployment;
}
}
export async function getWorkerDeploymentFromWorker(
workerId: string
): Promise<WorkerDeploymentWithWorkerTasks | undefined> {
const worker = await prisma.backgroundWorker.findUnique({
where: {
id: workerId,
},
include: {
deployment: true,
tasks: true,
},
});
if (!worker?.deployment) {
return;
}
const { deployment, ...workerWithoutDeployment } = worker;
return {
...deployment,
worker: workerWithoutDeployment,
};
}
export async function getWorkerDeploymentFromWorkerTask(
workerTaskId: string
): Promise<WorkerDeploymentWithWorkerTasks | undefined> {
const workerTask = await prisma.backgroundWorkerTask.findUnique({
where: {
id: workerTaskId,
},
include: {
worker: {
include: {
deployment: true,
tasks: true,
},
},
},
});
if (!workerTask?.worker.deployment) {
return;
}
const { deployment, ...workerWithoutDeployment } = workerTask.worker;
return {
...deployment,
worker: workerWithoutDeployment,
};
}