Skip to content
Merged
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
70 changes: 59 additions & 11 deletions apps/webapp/app/v3/runOpsMigration/controlPlaneCache.server.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,4 @@
import type {
BackgroundWorker,
BackgroundWorkerTask,
Prisma,
RuntimeEnvironmentType,
TaskQueue,
WorkerDeployment,
} from "@trigger.dev/database";
import type { BackgroundWorker, Prisma, RuntimeEnvironmentType } from "@trigger.dev/database";
import { BoundedTtlCache } from "~/services/realtime/boundedTtlCache";
import type { AuthenticatedEnvironment } from "@trigger.dev/core/v3/auth/environment";

Expand Down Expand Up @@ -49,12 +42,67 @@ export type ResolvedEnv = {
concurrencyLimitBurstFactor: Prisma.Decimal;
};

/**
* The BackgroundWorkerTask columns the dequeue resolve path reads. Mirrors run-engine's
* `ResolvedWorkerTask` exactly. The unread heavy JSON columns (`payloadSchema`, `config`,
* `queueConfig`, `description`) are dropped so this hot control-plane read stops shipping
* ~62KB/query (and each cached entry stays small); `machineConfig`/`retryConfig` are read
* at dequeue and stay.
*/
export type ResolvedWorkerTask = {
id: string;
slug: string;
machineConfig: Prisma.JsonValue | null;
retryConfig: Prisma.JsonValue | null;
maxDurationInSeconds: number | null;
};

/** The `select` that yields a `ResolvedWorkerTask`. */
export const resolvedWorkerTaskSelect = {
id: true,
slug: true,
machineConfig: true,
retryConfig: true,
maxDurationInSeconds: true,
} satisfies Prisma.BackgroundWorkerTaskSelect;

/** Mirrors run-engine's `ResolvedTaskQueue` exactly. `id` + `name` (the matcher keys on both). */
export type ResolvedTaskQueue = {
id: string;
name: string;
};

/** The `select` that yields a `ResolvedTaskQueue`. */
export const resolvedTaskQueueSelect = {
id: true,
name: true,
} satisfies Prisma.TaskQueueSelect;

/**
* Mirrors run-engine's `ResolvedWorkerDeployment` exactly. Drops the unread heavy JSON columns
* (`externalBuildData`, `buildServerMetadata`, `errorData`, `git`) from this single-row read.
*/
export type ResolvedWorkerDeployment = {
id: string;
friendlyId: string;
imageReference: string | null;
imagePlatform: string;
};

/** The `select` that yields a `ResolvedWorkerDeployment`. */
export const resolvedWorkerDeploymentSelect = {
id: true,
friendlyId: true,
imageReference: true,
imagePlatform: true,
} satisfies Prisma.WorkerDeploymentSelect;

/** Mirrors `WorkerDeploymentWithWorkerTasks` in `dequeueSystem.ts` exactly. */
export type ResolvedWorkerVersion = {
worker: BackgroundWorker;
tasks: BackgroundWorkerTask[];
queues: TaskQueue[];
deployment: WorkerDeployment | null;
tasks: ResolvedWorkerTask[];
queues: ResolvedTaskQueue[];
deployment: ResolvedWorkerDeployment | null;
};

// The canonical authenticated-environment shape (slug/type/project/organization/orgMember/…)
Expand Down
62 changes: 52 additions & 10 deletions apps/webapp/app/v3/runOpsMigration/controlPlaneResolver.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@ import {
ControlPlaneCache,
DEFAULT_CP_CACHE_MAX_ENTRIES,
DEFAULT_CP_CACHE_TTL_MS,
resolvedTaskQueueSelect,
resolvedWorkerDeploymentSelect,
resolvedWorkerTaskSelect,
type ResolvedAuthenticatedEnv,
type ResolvedEnv,
type ResolvedWorkerVersion,
Expand Down Expand Up @@ -389,7 +392,11 @@ export class ControlPlaneResolver {
if (backgroundWorkerId) {
const worker = await client.backgroundWorker.findFirst({
where: { id: backgroundWorkerId },
include: { deployment: true, tasks: true, queues: true },
include: {
deployment: { select: resolvedWorkerDeploymentSelect },
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
});

if (!worker) {
Expand All @@ -411,7 +418,16 @@ export class ControlPlaneResolver {
where: { environmentId, label: CURRENT_DEPLOYMENT_LABEL },
include: {
deployment: {
include: { worker: { include: { tasks: true, queues: true } } },
select: {
...resolvedWorkerDeploymentSelect,
type: true,
worker: {
include: {
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
},
},
},
},
});
Expand All @@ -421,11 +437,17 @@ export class ControlPlaneResolver {
}

if (type === undefined || promotion.deployment.type === "MANAGED") {
const { worker } = promotion.deployment;
return {
worker: promotion.deployment.worker,
tasks: promotion.deployment.worker.tasks,
queues: promotion.deployment.worker.queues,
deployment: promotion.deployment,
worker,
tasks: worker.tasks,
queues: worker.queues,
deployment: {
id: promotion.deployment.id,
friendlyId: promotion.deployment.friendlyId,
imageReference: promotion.deployment.imageReference,
imagePlatform: promotion.deployment.imagePlatform,
},
};
}

Expand All @@ -434,7 +456,15 @@ export class ControlPlaneResolver {
const latestV2Deployment = await client.workerDeployment.findFirst({
where: { environmentId, type: "MANAGED" },
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
include: { worker: { include: { tasks: true, queues: true } } },
select: {
...resolvedWorkerDeploymentSelect,
worker: {
include: {
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
},
},
});

if (!latestV2Deployment?.worker) {
Expand All @@ -445,7 +475,12 @@ export class ControlPlaneResolver {
worker: latestV2Deployment.worker,
tasks: latestV2Deployment.worker.tasks,
queues: latestV2Deployment.worker.queues,
deployment: latestV2Deployment,
deployment: {
id: latestV2Deployment.id,
friendlyId: latestV2Deployment.friendlyId,
imageReference: latestV2Deployment.imageReference,
imagePlatform: latestV2Deployment.imagePlatform,
},
};
}

Expand All @@ -455,7 +490,11 @@ export class ControlPlaneResolver {
): Promise<ResolvedWorkerVersion | null> {
const worker = await client.backgroundWorker.findFirst({
where: { id: workerId },
include: { deployment: true, tasks: true, queues: true },
include: {
deployment: { select: resolvedWorkerDeploymentSelect },
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
});

if (!worker) {
Expand All @@ -471,7 +510,10 @@ export class ControlPlaneResolver {
): Promise<ResolvedWorkerVersion | null> {
const worker = await client.backgroundWorker.findFirst({
where: { runtimeEnvironmentId: environmentId },
include: { tasks: true, queues: true },
include: {
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
});

Expand Down
119 changes: 94 additions & 25 deletions internal-packages/run-engine/src/engine/controlPlaneResolver.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,8 @@
import type {
BackgroundWorker,
BackgroundWorkerTask,
Prisma,
PrismaClient,
RuntimeEnvironmentType,
TaskQueue,
WorkerDeployment,
} from "@trigger.dev/database";
import { CURRENT_DEPLOYMENT_LABEL } from "@trigger.dev/core/v3/isomorphic";
import type { AuthenticatedEnvironment } from "@trigger.dev/core/v3/auth/environment";
Expand Down Expand Up @@ -51,12 +48,70 @@ export type ResolvedEngineEnv = {
*/
export type ResolvedAuthenticatedEnv = AuthenticatedEnvironment & { git: Prisma.JsonValue | null };

/**
* The BackgroundWorkerTask columns the dequeue resolve path actually reads. The worker's
* whole task set (~73 rows) is fetched for one matched task, so pulling the unread heavy
* JSON columns (`payloadSchema`, `config`, `queueConfig`, `description`) shipped ~62KB/query
* on the hottest control-plane read. `machineConfig`/`retryConfig` are read at dequeue and stay.
*/
export type ResolvedWorkerTask = {
id: string;
slug: string;
machineConfig: Prisma.JsonValue | null;
retryConfig: Prisma.JsonValue | null;
maxDurationInSeconds: number | null;
};

/** The `select` that yields a `ResolvedWorkerTask`. */
export const resolvedWorkerTaskSelect = {
id: true,
slug: true,
machineConfig: true,
retryConfig: true,
maxDurationInSeconds: true,
} satisfies Prisma.BackgroundWorkerTaskSelect;
Comment thread
ericallam marked this conversation as resolved.

/**
* The TaskQueue columns the dequeue resolve path uses: `id` and `name` (the matcher keys on
* `lockedQueueId`/`name`). Drops the unread `rateLimit` JSON and the concurrency/type scalars.
*/
export type ResolvedTaskQueue = {
id: string;
name: string;
};

/** The `select` that yields a `ResolvedTaskQueue`. */
export const resolvedTaskQueueSelect = {
id: true,
name: true,
} satisfies Prisma.TaskQueueSelect;

/**
* The WorkerDeployment columns the dequeue resolve path reads: `id`, `friendlyId`,
* `imageReference`, `imagePlatform`. Drops the unread heavy JSON columns (`externalBuildData`,
* `buildServerMetadata`, `errorData`, `git`) that this single-row read otherwise ships.
*/
export type ResolvedWorkerDeployment = {
id: string;
friendlyId: string;
imageReference: string | null;
imagePlatform: string;
};

/** The `select` that yields a `ResolvedWorkerDeployment`. */
export const resolvedWorkerDeploymentSelect = {
id: true,
friendlyId: true,
imageReference: true,
imagePlatform: true,
} satisfies Prisma.WorkerDeploymentSelect;

/** Identical to dequeue's `WorkerDeploymentWithWorkerTasks`. */
export type ResolvedWorkerVersion = {
worker: BackgroundWorker;
tasks: BackgroundWorkerTask[];
queues: TaskQueue[];
deployment: WorkerDeployment | null;
tasks: ResolvedWorkerTask[];
queues: ResolvedTaskQueue[];
deployment: ResolvedWorkerDeployment | null;
};

export interface ControlPlaneResolver {
Expand Down Expand Up @@ -207,9 +262,9 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
id: workerId,
},
include: {
deployment: true,
tasks: true,
queues: true,
deployment: { select: resolvedWorkerDeploymentSelect },
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
});

Expand All @@ -231,8 +286,8 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
runtimeEnvironmentId: environmentId,
},
include: {
tasks: true,
queues: true,
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
});
Expand All @@ -250,9 +305,9 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
id: workerId,
},
include: {
deployment: true,
tasks: true,
queues: true,
deployment: { select: resolvedWorkerDeploymentSelect },
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
});

Expand All @@ -278,11 +333,13 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
},
include: {
deployment: {
include: {
select: {
...resolvedWorkerDeploymentSelect,
type: true,
worker: {
include: {
tasks: true,
queues: true,
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
},
},
Expand All @@ -296,11 +353,17 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {

if (promotion.deployment.type === "MANAGED") {
// This is a run engine v2 deployment, so return it
const { worker } = promotion.deployment;
return {
worker: promotion.deployment.worker,
tasks: promotion.deployment.worker.tasks,
queues: promotion.deployment.worker.queues,
deployment: promotion.deployment,
worker,
tasks: worker.tasks,
queues: worker.queues,
deployment: {
id: promotion.deployment.id,
friendlyId: promotion.deployment.friendlyId,
imageReference: promotion.deployment.imageReference,
imagePlatform: promotion.deployment.imagePlatform,
},
};
}

Expand All @@ -311,11 +374,12 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
type: "MANAGED",
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
include: {
select: {
...resolvedWorkerDeploymentSelect,
worker: {
include: {
tasks: true,
queues: true,
tasks: { select: resolvedWorkerTaskSelect },
queues: { select: resolvedTaskQueueSelect },
},
},
},
Expand All @@ -329,7 +393,12 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
worker: latestV2Deployment.worker,
tasks: latestV2Deployment.worker.tasks,
queues: latestV2Deployment.worker.queues,
deployment: latestV2Deployment,
deployment: {
id: latestV2Deployment.id,
friendlyId: latestV2Deployment.friendlyId,
imageReference: latestV2Deployment.imageReference,
imagePlatform: latestV2Deployment.imagePlatform,
},
};
}
}
Loading
Loading