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
6 changes: 6 additions & 0 deletions .server-changes/worker-deployment-lookup-ordering.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
area: webapp
type: improvement
---

Speeds up resolving the latest worker version and deployment for an environment, removing an occasional stall when triggering runs in projects that have accumulated many deployed versions.
Original file line number Diff line number Diff line change
Expand Up @@ -271,9 +271,6 @@ export class ClickHouseRunsRepository implements IRunsRepository {
in: ids,
},
},
orderBy: {
id: "desc",
},
select: {
id: true,
friendlyId: true,
Expand Down
6 changes: 2 additions & 4 deletions apps/webapp/app/v3/models/workerDeployment.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,14 +128,12 @@ export async function findCurrentWorkerDeployment({
}

// We need to get the latest deployment of the given type
const latestDeployment = await prisma.workerDeployment.findFirst({
const latestDeployment = await $prisma.workerDeployment.findFirst({
Comment thread
ericallam marked this conversation as resolved.
where: {
environmentId,
type,
},
orderBy: {
id: "desc",
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
Comment thread
ericallam marked this conversation as resolved.
select: {
id: true,
imageReference: true,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -432,7 +432,7 @@ export class ControlPlaneResolver {
// MANAGED deployment.
const latestV2Deployment = await client.workerDeployment.findFirst({
where: { environmentId, type: "MANAGED" },
orderBy: { id: "desc" },
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
include: { worker: { include: { tasks: true, queues: true } } },
});

Expand All @@ -455,7 +455,6 @@ export class ControlPlaneResolver {
const worker = await client.backgroundWorker.findFirst({
where: { id: workerId },
include: { deployment: true, tasks: true, queues: true },
orderBy: { id: "desc" },
});

if (!worker) {
Expand All @@ -472,7 +471,7 @@ export class ControlPlaneResolver {
const worker = await client.backgroundWorker.findFirst({
where: { runtimeEnvironmentId: environmentId },
include: { tasks: true, queues: true },
orderBy: { id: "desc" },
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
});

if (!worker) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,12 @@ async function seedControlPlane(prisma: PrismaClient) {
async function seedWorker(
prisma: PrismaClient,
ctx: { projectId: string; environmentId: string },
opts?: { promote?: boolean }
opts?: {
promote?: boolean;
createdAt?: Date;
deploymentType?: "MANAGED" | "UNMANAGED" | "V1";
deploymentCreatedAt?: Date;
}
) {
const n = seedCounter++;
const worker = await prisma.backgroundWorker.create({
Expand All @@ -74,6 +79,7 @@ async function seedWorker(
version: `2024.1.${n}`,
metadata: {},
engine: "V2",
...(opts?.createdAt ? { createdAt: opts.createdAt } : {}),
},
});
const task = await prisma.backgroundWorkerTask.create({
Expand Down Expand Up @@ -104,11 +110,12 @@ async function seedWorker(
contentHash: `hash_${n}`,
version: worker.version,
shortCode: `dep_${n}`,
type: "MANAGED",
type: opts?.deploymentType ?? "MANAGED",
status: "DEPLOYED",
projectId: ctx.projectId,
environmentId: ctx.environmentId,
workerId: worker.id,
...(opts?.deploymentCreatedAt ? { createdAt: opts.deploymentCreatedAt } : {}),
},
});
await prisma.workerDeploymentPromotion.create({
Expand Down Expand Up @@ -748,3 +755,103 @@ heteroPostgresTest(
expect(reads()).toBe(readsAfterFirst * 2);
}
);

heteroPostgresTest(
"resolveWorkerVersion (DEVELOPMENT) resolves the newest worker by createdAt, not by id",
async ({ prisma14 }) => {
const { environment, project } = await seedControlPlane(prisma14);
const ctx = { projectId: project.id, environmentId: environment.id };

const newest = await seedWorker(prisma14, ctx, {
createdAt: new Date("2026-07-31T12:00:00.000Z"),
});
const oldest = await seedWorker(prisma14, ctx, {
createdAt: new Date("2026-07-30T12:00:00.000Z"),
});

expect(oldest.worker.id > newest.worker.id).toBe(true);
expect(oldest.worker.createdAt < newest.worker.createdAt).toBe(true);

const resolver = new ControlPlaneResolver({
controlPlaneReplica: prisma14,
controlPlanePrimary: prisma14,
cache: new ControlPlaneCache(),
splitEnabled: () => true,
});

const resolved = await resolver.resolveWorkerVersion({
environmentId: environment.id,
type: "DEVELOPMENT",
});

expect(resolved).not.toBeNull();
expect(resolved!.worker.id).toBe(newest.worker.id);
}
);

heteroPostgresTest(
"resolveWorkerVersion latest-MANAGED fallback resolves by createdAt, not by id",
async ({ prisma14 }) => {
const { environment, project } = await seedControlPlane(prisma14);
const ctx = { projectId: project.id, environmentId: environment.id };

await seedWorker(prisma14, ctx, { promote: true, deploymentType: "V1" });

const newest = await seedWorker(prisma14, ctx, {
promote: false,
deploymentCreatedAt: new Date("2026-07-31T12:00:00.000Z"),
});
const oldest = await seedWorker(prisma14, ctx, {
promote: false,
deploymentCreatedAt: new Date("2026-07-30T12:00:00.000Z"),
});

const newestDeployment = await prisma14.workerDeployment.create({
data: {
friendlyId: `deployment_newest_${environment.id}`,
contentHash: "hash_newest",
version: newest.worker.version,
shortCode: "dep_newest",
type: "MANAGED",
status: "DEPLOYED",
projectId: project.id,
environmentId: environment.id,
workerId: newest.worker.id,
createdAt: new Date("2026-07-31T12:00:00.000Z"),
},
});
const oldestDeployment = await prisma14.workerDeployment.create({
data: {
friendlyId: `deployment_oldest_${environment.id}`,
contentHash: "hash_oldest",
version: oldest.worker.version,
shortCode: "dep_oldest",
type: "MANAGED",
status: "DEPLOYED",
projectId: project.id,
environmentId: environment.id,
workerId: oldest.worker.id,
createdAt: new Date("2026-07-30T12:00:00.000Z"),
},
});

expect(oldestDeployment.id > newestDeployment.id).toBe(true);
expect(oldestDeployment.createdAt < newestDeployment.createdAt).toBe(true);

const resolver = new ControlPlaneResolver({
controlPlaneReplica: prisma14,
controlPlanePrimary: prisma14,
cache: new ControlPlaneCache(),
splitEnabled: () => true,
});

const resolved = await resolver.resolveWorkerVersion({
environmentId: environment.id,
type: "PRODUCTION",
});

expect(resolved).not.toBeNull();
expect(resolved!.deployment!.id).toBe(newestDeployment.id);
expect(resolved!.worker.id).toBe(newest.worker.id);
}
);
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
CREATE INDEX CONCURRENTLY IF NOT EXISTS "WorkerDeployment_environmentId_createdAt_idx" ON "public"."WorkerDeployment"("environmentId", "createdAt");
1 change: 1 addition & 0 deletions internal-packages/database/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -2208,6 +2208,7 @@ model WorkerDeployment {
@@unique([projectId, shortCode])
@@unique([environmentId, version])
@@index([commitSHA])
@@index([environmentId, createdAt])
Comment thread
ericallam marked this conversation as resolved.
}

enum WorkerDeploymentStatus {
Expand Down
11 changes: 2 additions & 9 deletions internal-packages/run-engine/src/engine/controlPlaneResolver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -234,9 +234,7 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
tasks: true,
queues: true,
},
orderBy: {
id: "desc",
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
});

if (!worker) {
Expand All @@ -256,9 +254,6 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
tasks: true,
queues: true,
},
orderBy: {
id: "desc",
},
});

if (!worker) {
Expand Down Expand Up @@ -315,9 +310,7 @@ export class PassthroughControlPlaneResolver implements ControlPlaneResolver {
environmentId,
type: "MANAGED",
},
orderBy: {
id: "desc",
},
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
Comment thread
ericallam marked this conversation as resolved.
include: {
worker: {
include: {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -538,3 +538,142 @@ describe("DequeueSystem controlPlaneResolver (single-DB passthrough)", () => {
}
);
});

heteroPostgresTest(
"resolveWorkerVersion (DEVELOPMENT) resolves the newest worker by createdAt, not by id",
async ({ prisma14 }) => {
const cp = await seedControlPlane(
prisma14 as unknown as PrismaClient,
"cpord",
"ordering-task"
);

await prisma14.backgroundWorker.update({
where: { id: cp.worker.id },
data: { createdAt: new Date("2026-07-01T12:00:00.000Z") },
});

const newest = await prisma14.backgroundWorker.create({
data: {
friendlyId: generateFriendlyId("worker"),
contentHash: "hash_newest",
projectId: cp.project.id,
runtimeEnvironmentId: cp.environment.id,
version: "20260731.1",
metadata: {},
engine: "V2",
createdAt: new Date("2026-07-31T12:00:00.000Z"),
},
});
const oldest = await prisma14.backgroundWorker.create({
data: {
friendlyId: generateFriendlyId("worker"),
contentHash: "hash_oldest",
projectId: cp.project.id,
runtimeEnvironmentId: cp.environment.id,
version: "20260730.1",
metadata: {},
engine: "V2",
createdAt: new Date("2026-07-30T12:00:00.000Z"),
},
});

expect(oldest.id > newest.id).toBe(true);
expect(oldest.createdAt < newest.createdAt).toBe(true);

const resolver = new PassthroughControlPlaneResolver({
prisma: prisma14 as unknown as PrismaClient,
});

const resolved = await resolver.resolveWorkerVersion({
environmentId: cp.environment.id,
type: "DEVELOPMENT",
});

assertNonNullable(resolved);
expect(resolved.worker.id).toBe(newest.id);
}
);

heteroPostgresTest(
"resolveWorkerVersion latest-MANAGED fallback resolves by createdAt, not by id",
async ({ prisma14 }) => {
const cp = await seedControlPlane(
prisma14 as unknown as PrismaClient,
"cpfall",
"fallback-task"
);

await prisma14.workerDeployment.update({
where: { id: cp.deployment.id },
data: { type: "V1" },
});

const newestWorker = await prisma14.backgroundWorker.create({
data: {
friendlyId: generateFriendlyId("worker"),
contentHash: "hash_newest",
projectId: cp.project.id,
runtimeEnvironmentId: cp.environment.id,
version: "20260731.1",
metadata: {},
engine: "V2",
},
});
const oldestWorker = await prisma14.backgroundWorker.create({
data: {
friendlyId: generateFriendlyId("worker"),
contentHash: "hash_oldest",
projectId: cp.project.id,
runtimeEnvironmentId: cp.environment.id,
version: "20260730.1",
metadata: {},
engine: "V2",
},
});

const newestDeployment = await prisma14.workerDeployment.create({
data: {
friendlyId: generateFriendlyId("deployment"),
contentHash: "hash_newest",
version: "20260731.1",
shortCode: "short_code_newest",
status: "DEPLOYED",
projectId: cp.project.id,
environmentId: cp.environment.id,
workerId: newestWorker.id,
type: "MANAGED",
createdAt: new Date("2026-07-31T12:00:00.000Z"),
},
});
const oldestDeployment = await prisma14.workerDeployment.create({
data: {
friendlyId: generateFriendlyId("deployment"),
contentHash: "hash_oldest",
version: "20260730.1",
shortCode: "short_code_oldest",
status: "DEPLOYED",
projectId: cp.project.id,
environmentId: cp.environment.id,
workerId: oldestWorker.id,
type: "MANAGED",
createdAt: new Date("2026-07-30T12:00:00.000Z"),
},
});

expect(oldestDeployment.id > newestDeployment.id).toBe(true);
expect(oldestDeployment.createdAt < newestDeployment.createdAt).toBe(true);

const resolver = new PassthroughControlPlaneResolver({
prisma: prisma14 as unknown as PrismaClient,
});

const resolved = await resolver.resolveWorkerVersion({
environmentId: cp.environment.id,
type: "PRODUCTION",
});

assertNonNullable(resolved);
expect(resolved.deployment?.id).toBe(newestDeployment.id);
}
);
Loading