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
1 change: 1 addition & 0 deletions internal-packages/run-engine/src/engine/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,7 @@ export class RunEngine {
this.pendingVersionSystem = new PendingVersionSystem({
resources,
enqueueSystem: this.enqueueSystem,
executionSnapshotSystem: this.executionSnapshotSystem,
queueRunsPendingVersionBatchSize: options.queueRunsWaitingForWorkerBatchSize,
lagRetryDelayMs: options.pendingVersionLagRetryDelayMs,
lagMaxRetries: options.pendingVersionLagMaxRetries,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,8 @@ export class DelayedRunSystem {
throw new ServiceValidationError("Cannot reschedule a run that is not delayed");
}

const isParked = snapshot.runStatus === "PENDING_VERSION";
Comment thread
0ski marked this conversation as resolved.
Comment thread
0ski marked this conversation as resolved.

const updatedRun = await this.$.runStore.rescheduleRun(
runId,
{
Expand All @@ -57,6 +59,13 @@ export class DelayedRunSystem {
environmentType: snapshot.environmentType,
projectId: snapshot.projectId,
organizationId: snapshot.organizationId,
...(isParked
? {
Comment thread
0ski marked this conversation as resolved.
executionStatus: "RUN_CREATED" as const,
runStatus: "PENDING_VERSION" as const,
description: "Parked run was rescheduled to a future date",
}
: {}),
},
},
prisma
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import {
import type { TaskRunError } from "@trigger.dev/core/v3/schemas";
import type { MinimalAuthenticatedEnvironment } from "../../shared/index.js";
import type { EnqueueSystem } from "./enqueueSystem.js";
import type { ExecutionSnapshotSystem } from "./executionSnapshotSystem.js";
import type { SystemResources } from "./systems.js";

import { boundedIn } from "@trigger.dev/database";
Expand All @@ -28,6 +29,7 @@ export type PendingVersionSystemOptions = {
*/
lagMaxRetries?: number;
externalDeploymentParkDeadlineMs?: number;
executionSnapshotSystem: ExecutionSnapshotSystem;
};

const DEFAULT_LAG_RETRY_DELAY_MS = 5_000;
Expand Down Expand Up @@ -60,10 +62,12 @@ export function readExternalDeploymentIdAnnotation(annotations: unknown): string
export class PendingVersionSystem {
private readonly $: SystemResources;
private readonly enqueueSystem: EnqueueSystem;
private readonly executionSnapshotSystem: ExecutionSnapshotSystem;

constructor(private readonly options: PendingVersionSystemOptions) {
this.$ = options.resources;
this.enqueueSystem = options.enqueueSystem;
this.executionSnapshotSystem = options.executionSnapshotSystem;
}

async enqueueRunsForBackgroundWorker(backgroundWorkerId: string, attempt: number = 0) {
Expand Down Expand Up @@ -231,6 +235,20 @@ export class PendingVersionSystem {
}

if (stillDelayed) {
await this.executionSnapshotSystem.createExecutionSnapshot(
tx,
{
run: { id: run.id, status: "DELAYED" },
snapshot: { executionStatus: "DELAYED", description: "Run is delayed" },
batchId: run.batchId ?? undefined,
environmentId: backgroundWorker.runtimeEnvironment.id,
environmentType: backgroundWorker.runtimeEnvironment.type,
projectId: backgroundWorker.runtimeEnvironment.project.id,
organizationId: backgroundWorker.runtimeEnvironment.organization.id,
},
store
);

return true;
}

Expand Down Expand Up @@ -470,6 +488,22 @@ export class PendingVersionSystem {
);
}

if (run.delayUntil && run.delayUntil > new Date()) {
this.$.logger.info(
"expireParkedExternalDeploymentRun: run is not due yet, re-arming the park deadline",
{ runId, externalDeploymentId, delayUntil: run.delayUntil }
);

await this.scheduleExternalDeploymentParkDeadline({
runId,
externalDeploymentId,
ttl: run.ttl,
delayUntil: run.delayUntil,
});

return;
}
Comment thread
0ski marked this conversation as resolved.
Comment thread
0ski marked this conversation as resolved.

const error: TaskRunError = {
type: "STRING_ERROR",
raw: `Run expired because no deployment with external id '${externalDeploymentId}' became available`,
Expand Down Expand Up @@ -574,6 +608,20 @@ export class PendingVersionSystem {
}

if (stillDelayed) {
await this.executionSnapshotSystem.createExecutionSnapshot(
tx,
{
run: { id: run.id, status: "DELAYED" },
snapshot: { executionStatus: "DELAYED", description: "Run is delayed" },
batchId: run.batchId ?? undefined,
environmentId: env.id,
environmentType: env.type,
projectId: env.project.id,
organizationId: env.organization.id,
},
store
);

return true;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -566,6 +566,207 @@ describe("RunEngine external deployment parking", () => {
}
);

containerTest(
"a run released while still delayed records a DELAYED snapshot and no longer reports as parked",
async ({ prisma, redisOptions }) => {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = createEngine(prisma, redisOptions);

try {
const taskIdentifier = "test-task";

const run = await engine.trigger(
{
number: 1,
friendlyId: "run_1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t1234",
spanId: "s1234",
queue: `task/${taskIdentifier}`,
isTest: false,
tags: [],
delayUntil: new Date(Date.now() + 60 * 60 * 1000),
annotations: {
triggerSource: "sdk",
triggerAction: "trigger",
rootTriggerSource: "sdk",
externalDeploymentId: "commit-released",
},
parkedOnExternalDeploymentId: "commit-released",
},
prisma
);

const worker = await setupBackgroundWorker(
engine,
authenticatedEnvironment,
taskIdentifier
);
await nameDeploymentWithExternalId(prisma, worker.worker.id, "commit-released");

await engine.pendingVersionSystem.enqueueRunsForBackgroundWorker(worker.worker.id);

const released = await prisma.taskRun.findFirstOrThrow({ where: { id: run.id } });
expect(released.status).toBe("DELAYED");
expect(released.lockedToVersionId).toBe(worker.worker.id);

const afterRelease = await prisma.taskRunExecutionSnapshot.findMany({
where: { runId: run.id },
orderBy: { createdAt: "asc" },
select: { executionStatus: true, runStatus: true },
});

expect(afterRelease.at(-1)?.runStatus).toBe("DELAYED");
expect(afterRelease.at(-1)?.executionStatus).toBe("DELAYED");

// A later debounce push must not re-label an already-released run as parked.
await engine.delayedRunSystem.rescheduleDelayedRun({
runId: run.id,
delayUntil: new Date(Date.now() + 2 * 60 * 60 * 1000),
tx: prisma,
});

const afterPush = await prisma.taskRunExecutionSnapshot.findMany({
where: { runId: run.id },
orderBy: { createdAt: "asc" },
select: { executionStatus: true, runStatus: true },
});

expect(afterPush.at(-1)?.runStatus).toBe("DELAYED");
expect(afterPush.at(-1)?.executionStatus).toBe("DELAYED");
} finally {
await engine.quit();
}
}
);

containerTest(
"a debounce push on a parked run keeps the snapshot parked instead of reporting it delayed",
async ({ prisma, redisOptions }) => {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = createEngine(prisma, redisOptions);

try {
const taskIdentifier = "test-task";

const run = await engine.trigger(
{
number: 1,
friendlyId: "run_1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t1234",
spanId: "s1234",
queue: `task/${taskIdentifier}`,
isTest: false,
tags: [],
delayUntil: new Date(Date.now() + 60 * 1000),
annotations: {
triggerSource: "sdk",
triggerAction: "trigger",
rootTriggerSource: "sdk",
externalDeploymentId: "commit-snapshot",
},
parkedOnExternalDeploymentId: "commit-snapshot",
},
prisma
);

await engine.delayedRunSystem.rescheduleDelayedRun({
runId: run.id,
delayUntil: new Date(Date.now() + 10 * 60 * 1000),
tx: prisma,
});

const snapshots = await prisma.taskRunExecutionSnapshot.findMany({
where: { runId: run.id },
orderBy: { createdAt: "asc" },
select: { executionStatus: true, runStatus: true },
});

const latest = snapshots.at(-1);

assertNonNullable(latest);
expect(latest.runStatus).toBe("PENDING_VERSION");
expect(latest.executionStatus).toBe("RUN_CREATED");

const stillParked = await prisma.taskRun.findFirstOrThrow({ where: { id: run.id } });
expect(stillParked.status).toBe("PENDING_VERSION");
} finally {
await engine.quit();
}
}
);

containerTest(
"the parking deadline re-arms instead of expiring a run whose delay was pushed past it",
async ({ prisma, redisOptions }) => {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = createEngine(prisma, redisOptions);

try {
const taskIdentifier = "test-task";

const run = await engine.trigger(
{
number: 1,
friendlyId: "run_1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t1234",
spanId: "s1234",
queue: `task/${taskIdentifier}`,
isTest: false,
tags: [],
delayUntil: new Date(Date.now() + 60 * 1000),
annotations: {
triggerSource: "sdk",
triggerAction: "trigger",
rootTriggerSource: "sdk",
externalDeploymentId: "commit-pushed",
},
parkedOnExternalDeploymentId: "commit-pushed",
},
prisma
);

// Stand in for a debounce push: the run's delay moves out, but nothing re-arms the
// deadline that was computed when the run was first parked.
const pushedDelayUntil = new Date(Date.now() + 60 * 60 * 1000);
await prisma.taskRun.update({
where: { id: run.id },
data: { delayUntil: pushedDelayUntil },
});

await engine.pendingVersionSystem.expireParkedExternalDeploymentRun({
runId: run.id,
externalDeploymentId: "commit-pushed",
});

const stillParked = await prisma.taskRun.findFirstOrThrow({ where: { id: run.id } });

expect(stillParked.status).toBe("PENDING_VERSION");
expect(stillParked.statusReason).toBe("EXTERNAL_DEPLOYMENT_PENDING");
expect(stillParked.expiredAt).toBeNull();
} finally {
await engine.quit();
}
}
);

containerTest(
"the parking deadline expires a run whose deployment never arrived",
async ({ prisma, redisOptions }) => {
Expand Down
7 changes: 4 additions & 3 deletions internal-packages/run-store/src/PostgresRunStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1439,9 +1439,10 @@ export class PostgresRunStore implements RunStore {
executionSnapshots: {
create: {
engine: "V2",
executionStatus: "DELAYED",
description: "Delayed run was rescheduled to a future date",
runStatus: "DELAYED",
executionStatus: data.snapshot.executionStatus ?? "DELAYED",
description:
data.snapshot.description ?? "Delayed run was rescheduled to a future date",
runStatus: data.snapshot.runStatus ?? "DELAYED",
environmentId: data.snapshot.environmentId,
environmentType: data.snapshot.environmentType,
projectId: data.snapshot.projectId,
Expand Down
3 changes: 3 additions & 0 deletions internal-packages/run-store/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,9 @@ export type RescheduleSnapshotInput = {
environmentType: RuntimeEnvironmentType;
projectId: string;
organizationId: string;
executionStatus?: TaskRunExecutionStatus;
runStatus?: TaskRunStatus;
description?: string;
};

export type LockSnapshotInput = {
Expand Down