From 60268108fedef69de5c1b9244f1cfd234382af3e Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Fri, 10 Jan 2025 17:01:23 +0000 Subject: [PATCH] cut resume payload queries in half --- .../v3/marqs/sharedQueueConsumer.server.ts | 302 ++++++++++++------ .../app/v3/services/resumeAttempt.server.ts | 24 +- 2 files changed, 217 insertions(+), 109 deletions(-) diff --git a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts index 80374055e..62f721a5b 100644 --- a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts @@ -17,6 +17,7 @@ import { ZodMessageSender } from "@trigger.dev/core/v3/zodMessageHandler"; import { BackgroundWorker, BackgroundWorkerTask, + Prisma, RuntimeEnvironment, TaskRun, TaskRunStatus, @@ -998,7 +999,157 @@ export class SharedQueueConsumer { } } +type AttemptForCompletion = Prisma.TaskRunAttemptGetPayload<{ + include: { + backgroundWorker: true; + backgroundWorkerTask: true; + taskRun: { + include: { + runtimeEnvironment: { + include: { + organization: true; + project: true; + }; + }; + tags: true; + }; + }; + queue: true; + }; +}>; + +type AttemptForExecution = Prisma.TaskRunAttemptGetPayload<{ + include: { + backgroundWorker: true; + backgroundWorkerTask: true; + runtimeEnvironment: { + include: { + organization: true; + project: true; + }; + }; + taskRun: { + include: { + tags: true; + batchItems: { + include: { + batchTaskRun: { + select: { + friendlyId: true; + }; + }; + }; + }; + }; + }; + queue: true; + }; +}>; + class SharedQueueTasks { + private _completionPayloadFromAttempt(attempt: AttemptForCompletion): TaskRunExecutionResult { + const ok = attempt.status === "COMPLETED"; + + if (ok) { + const success: TaskRunSuccessfulExecutionResult = { + ok, + id: attempt.taskRun.friendlyId, + output: attempt.output ?? undefined, + outputType: attempt.outputType, + taskIdentifier: attempt.taskRun.taskIdentifier, + }; + return success; + } else { + const failure: TaskRunFailedExecutionResult = { + ok, + id: attempt.taskRun.friendlyId, + error: attempt.error as TaskRunError, + taskIdentifier: attempt.taskRun.taskIdentifier, + }; + return failure; + } + } + + private async _executionFromAttempt( + attempt: AttemptForExecution, + machinePreset?: MachinePreset + ): Promise { + const { backgroundWorkerTask, taskRun, queue } = attempt; + + if (!machinePreset) { + machinePreset = machinePresetFromConfig(backgroundWorkerTask.machineConfig ?? {}); + } + + const metadata = await parsePacket({ + data: taskRun.metadata ?? undefined, + dataType: taskRun.metadataType, + }); + + const execution: ProdTaskRunExecution = { + task: { + id: backgroundWorkerTask.slug, + filePath: backgroundWorkerTask.filePath, + exportName: backgroundWorkerTask.exportName, + }, + attempt: { + id: attempt.friendlyId, + number: attempt.number, + startedAt: attempt.startedAt ?? attempt.createdAt, + backgroundWorkerId: attempt.backgroundWorkerId, + backgroundWorkerTaskId: attempt.backgroundWorkerTaskId, + status: "EXECUTING" as const, + }, + run: { + id: taskRun.friendlyId, + payload: taskRun.payload, + payloadType: taskRun.payloadType, + context: taskRun.context, + createdAt: taskRun.createdAt, + startedAt: taskRun.startedAt ?? taskRun.createdAt, + tags: taskRun.tags.map((tag) => tag.name), + isTest: taskRun.isTest, + idempotencyKey: taskRun.idempotencyKey ?? undefined, + durationMs: taskRun.usageDurationMs, + costInCents: taskRun.costInCents, + baseCostInCents: taskRun.baseCostInCents, + metadata, + maxDuration: taskRun.maxDurationInSeconds ?? undefined, + }, + queue: { + id: queue.friendlyId, + name: queue.name, + }, + environment: { + id: attempt.runtimeEnvironment.id, + slug: attempt.runtimeEnvironment.slug, + type: attempt.runtimeEnvironment.type, + }, + organization: { + id: attempt.runtimeEnvironment.organization.id, + slug: attempt.runtimeEnvironment.organization.slug, + name: attempt.runtimeEnvironment.organization.title, + }, + project: { + id: attempt.runtimeEnvironment.project.id, + ref: attempt.runtimeEnvironment.project.externalRef, + slug: attempt.runtimeEnvironment.project.slug, + name: attempt.runtimeEnvironment.project.name, + }, + batch: + taskRun.batchItems[0] && taskRun.batchItems[0].batchTaskRun + ? { id: taskRun.batchItems[0].batchTaskRun.friendlyId } + : undefined, + worker: { + id: attempt.backgroundWorkerId, + contentHash: attempt.backgroundWorker.contentHash, + version: attempt.backgroundWorker.version, + }, + machine: machinePreset, + }; + + return execution; + } + async getCompletionPayloadFromAttempt(id: string): Promise { const attempt = await prisma.taskRunAttempt.findFirst({ where: { @@ -1030,26 +1181,7 @@ class SharedQueueTasks { return; } - const ok = attempt.status === "COMPLETED"; - - if (ok) { - const success: TaskRunSuccessfulExecutionResult = { - ok, - id: attempt.taskRun.friendlyId, - output: attempt.output ?? undefined, - outputType: attempt.outputType, - taskIdentifier: attempt.taskRun.taskIdentifier, - }; - return success; - } else { - const failure: TaskRunFailedExecutionResult = { - ok, - id: attempt.taskRun.friendlyId, - error: attempt.error as TaskRunError, - taskIdentifier: attempt.taskRun.taskIdentifier, - }; - return failure; - } + return this._completionPayloadFromAttempt(attempt); } async getExecutionPayloadFromAttempt({ @@ -1162,78 +1294,10 @@ class SharedQueueTasks { }, }); } - - const { backgroundWorkerTask, taskRun, queue } = attempt; + const { backgroundWorkerTask, taskRun } = attempt; const machinePreset = machinePresetFromConfig(backgroundWorkerTask.machineConfig ?? {}); - - const metadata = await parsePacket({ - data: taskRun.metadata ?? undefined, - dataType: taskRun.metadataType, - }); - - const execution: ProdTaskRunExecution = { - task: { - id: backgroundWorkerTask.slug, - filePath: backgroundWorkerTask.filePath, - exportName: backgroundWorkerTask.exportName, - }, - attempt: { - id: attempt.friendlyId, - number: attempt.number, - startedAt: attempt.startedAt ?? attempt.createdAt, - backgroundWorkerId: attempt.backgroundWorkerId, - backgroundWorkerTaskId: attempt.backgroundWorkerTaskId, - status: "EXECUTING" as const, - }, - run: { - id: taskRun.friendlyId, - payload: taskRun.payload, - payloadType: taskRun.payloadType, - context: taskRun.context, - createdAt: taskRun.createdAt, - startedAt: taskRun.startedAt ?? taskRun.createdAt, - tags: taskRun.tags.map((tag) => tag.name), - isTest: taskRun.isTest, - idempotencyKey: taskRun.idempotencyKey ?? undefined, - durationMs: taskRun.usageDurationMs, - costInCents: taskRun.costInCents, - baseCostInCents: taskRun.baseCostInCents, - metadata, - maxDuration: taskRun.maxDurationInSeconds ?? undefined, - }, - queue: { - id: queue.friendlyId, - name: queue.name, - }, - environment: { - id: attempt.runtimeEnvironment.id, - slug: attempt.runtimeEnvironment.slug, - type: attempt.runtimeEnvironment.type, - }, - organization: { - id: attempt.runtimeEnvironment.organization.id, - slug: attempt.runtimeEnvironment.organization.slug, - name: attempt.runtimeEnvironment.organization.title, - }, - project: { - id: attempt.runtimeEnvironment.project.id, - ref: attempt.runtimeEnvironment.project.externalRef, - slug: attempt.runtimeEnvironment.project.slug, - name: attempt.runtimeEnvironment.project.name, - }, - batch: - taskRun.batchItems[0] && taskRun.batchItems[0].batchTaskRun - ? { id: taskRun.batchItems[0].batchTaskRun.friendlyId } - : undefined, - worker: { - id: attempt.backgroundWorkerId, - contentHash: attempt.backgroundWorker.contentHash, - version: attempt.backgroundWorker.version, - }, - machine: machinePreset, - }; - + const execution = await this._executionFromAttempt(attempt, machinePreset); const variables = await this.#buildEnvironmentVariables( attempt.runtimeEnvironment, taskRun.id, @@ -1252,6 +1316,64 @@ class SharedQueueTasks { return payload; } + async getResumePayload(attemptId: string): Promise< + | { + execution: ProdTaskRunExecution; + completion: TaskRunExecutionResult; + } + | undefined + > { + const attempt = await prisma.taskRunAttempt.findFirst({ + where: { + id: attemptId, + }, + include: { + backgroundWorker: true, + backgroundWorkerTask: true, + runtimeEnvironment: { + include: { + organization: true, + project: true, + }, + }, + taskRun: { + include: { + runtimeEnvironment: { + include: { + organization: true, + project: true, + }, + }, + tags: true, + batchItems: { + include: { + batchTaskRun: { + select: { + friendlyId: true, + }, + }, + }, + }, + }, + }, + queue: true, + }, + }); + + if (!attempt) { + logger.error("getExecutionPayloadFromAttempt: No attempt found", { id: attemptId }); + return; + } + + const execution = await this._executionFromAttempt(attempt); + const completion = this._completionPayloadFromAttempt(attempt); + + return { + execution, + completion, + }; + } + async getLatestExecutionPayloadFromRun( id: string, setToExecuting?: boolean, diff --git a/apps/webapp/app/v3/services/resumeAttempt.server.ts b/apps/webapp/app/v3/services/resumeAttempt.server.ts index fd3b70095..a2bf8db04 100644 --- a/apps/webapp/app/v3/services/resumeAttempt.server.ts +++ b/apps/webapp/app/v3/services/resumeAttempt.server.ts @@ -182,30 +182,16 @@ export class ResumeAttemptService extends BaseService { completedRunId: completedAttempt.taskRunId, }); - const completion = await sharedQueueTasks.getCompletionPayloadFromAttempt( - completedAttempt.id - ); + const resumePayload = await sharedQueueTasks.getResumePayload(completedAttempt.id); - if (!completion) { - logger.error("Failed to get completion payload"); + if (!resumePayload) { + logger.error("Failed to get resume payload"); await marqs?.acknowledgeMessage(attempt.taskRunId); return; } - completions.push(completion); - - const executionPayload = await sharedQueueTasks.getExecutionPayloadFromAttempt({ - id: completedAttempt.id, - skipStatusChecks: true, // already checked when getting the completion - }); - - if (!executionPayload) { - logger.error("Failed to get execution payload"); - await marqs?.acknowledgeMessage(attempt.taskRunId); - return; - } - - executions.push(executionPayload.execution); + completions.push(resumePayload.completion); + executions.push(resumePayload.execution); } await this.#setPostResumeStatuses(attempt);