From f62294a119f0bdca4e1e564df3768a622391a886 Mon Sep 17 00:00:00 2001 From: Oskar Otwinowski Date: Wed, 19 Aug 2026 16:38:30 +0200 Subject: [PATCH] fix(run-engine): correct park deadline and snapshot state for debounced parked runs Two defects surface when a run parked on an external deployment id is pushed by a debounce key. Both were reproduced against a local instance before fixing. 1. The run is expired before it is due. The park deadline is armed once, when the run is first parked, from max(now, delayUntil) + deadline. Debounce pushes delayUntil out afterwards: rescheduleDelayedRun reschedules enqueueDelayedRun:, and the redis-worker reschedule is an update-only ZADD, so expireParkedExternalDeploymentRun: is never re-armed. Repeat triggers on one key walk delayUntil past a deadline that no longer moves, and the run is expired with EXTERNAL_DEPLOYMENT_NOT_FOUND before it was ever due to start. Observed: a run due at 14:01:37 expired at 13:57:02. The expiry job already loads delayUntil, so it now re-arms from the current value and returns instead of expiring a run that is not due. Putting the guard there rather than in the debounce path covers every caller that moves delayUntil, and it stays bounded by the debounce max-duration contract. 2. The run reports itself as delayed while it is parked. rescheduleRun hardcoded a DELAYED/DELAYED execution snapshot, so a debounce push left the run row on PENDING_VERSION while its latest snapshot claimed DELAYED, and the run page described a parked run as delayed. The snapshot statuses are now supplied by the caller and default to DELAYED, so the ordinary delayed path is unchanged, and rescheduleDelayedRun passes the parked statuses through when the run is parked. --- .../run-engine/src/engine/index.ts | 1 + .../src/engine/systems/delayedRunSystem.ts | 9 + .../engine/systems/pendingVersionSystem.ts | 48 +++++ .../tests/externalDeploymentParking.test.ts | 201 ++++++++++++++++++ .../run-store/src/PostgresRunStore.ts | 7 +- internal-packages/run-store/src/types.ts | 3 + 6 files changed, 266 insertions(+), 3 deletions(-) diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index a9c6ea5b78b..e07d668f399 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -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, diff --git a/internal-packages/run-engine/src/engine/systems/delayedRunSystem.ts b/internal-packages/run-engine/src/engine/systems/delayedRunSystem.ts index d0d603f9f02..00867f891a9 100644 --- a/internal-packages/run-engine/src/engine/systems/delayedRunSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/delayedRunSystem.ts @@ -48,6 +48,8 @@ export class DelayedRunSystem { throw new ServiceValidationError("Cannot reschedule a run that is not delayed"); } + const isParked = snapshot.runStatus === "PENDING_VERSION"; + const updatedRun = await this.$.runStore.rescheduleRun( runId, { @@ -57,6 +59,13 @@ export class DelayedRunSystem { environmentType: snapshot.environmentType, projectId: snapshot.projectId, organizationId: snapshot.organizationId, + ...(isParked + ? { + executionStatus: "RUN_CREATED" as const, + runStatus: "PENDING_VERSION" as const, + description: "Parked run was rescheduled to a future date", + } + : {}), }, }, prisma diff --git a/internal-packages/run-engine/src/engine/systems/pendingVersionSystem.ts b/internal-packages/run-engine/src/engine/systems/pendingVersionSystem.ts index fddd7565458..000cf4c8e70 100644 --- a/internal-packages/run-engine/src/engine/systems/pendingVersionSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/pendingVersionSystem.ts @@ -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"; @@ -28,6 +29,7 @@ export type PendingVersionSystemOptions = { */ lagMaxRetries?: number; externalDeploymentParkDeadlineMs?: number; + executionSnapshotSystem: ExecutionSnapshotSystem; }; const DEFAULT_LAG_RETRY_DELAY_MS = 5_000; @@ -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) { @@ -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; } @@ -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; + } + const error: TaskRunError = { type: "STRING_ERROR", raw: `Run expired because no deployment with external id '${externalDeploymentId}' became available`, @@ -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; } diff --git a/internal-packages/run-engine/src/engine/tests/externalDeploymentParking.test.ts b/internal-packages/run-engine/src/engine/tests/externalDeploymentParking.test.ts index f815b1b8638..6f32c84aec5 100644 --- a/internal-packages/run-engine/src/engine/tests/externalDeploymentParking.test.ts +++ b/internal-packages/run-engine/src/engine/tests/externalDeploymentParking.test.ts @@ -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 }) => { diff --git a/internal-packages/run-store/src/PostgresRunStore.ts b/internal-packages/run-store/src/PostgresRunStore.ts index b3b5c565f03..b7d5086431b 100644 --- a/internal-packages/run-store/src/PostgresRunStore.ts +++ b/internal-packages/run-store/src/PostgresRunStore.ts @@ -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, diff --git a/internal-packages/run-store/src/types.ts b/internal-packages/run-store/src/types.ts index c09522ccd02..7c7f9566893 100644 --- a/internal-packages/run-store/src/types.ts +++ b/internal-packages/run-store/src/types.ts @@ -79,6 +79,9 @@ export type RescheduleSnapshotInput = { environmentType: RuntimeEnvironmentType; projectId: string; organizationId: string; + executionStatus?: TaskRunExecutionStatus; + runStatus?: TaskRunStatus; + description?: string; }; export type LockSnapshotInput = {