Skip to content
Open
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/run-finalization-guard.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
area: webapp
type: fix
---

A durable guard improves reliability for runs waiting on triggerAndWait or batchTriggerAndWait if there's a database error that interrupts a child run finishing.
6 changes: 6 additions & 0 deletions apps/webapp/app/v3/runOpsMigration/unblockRouteCatalog.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,12 @@ export const UNBLOCK_ROUTES: readonly UnblockRoute[] = [
site: RUN_ATTEMPT_SYSTEM,
symbol: "#permanentlyFailRun",
},
{
id: "runAttempt.ensureFinalized",
kind: "RUN",
site: RUN_ATTEMPT_SYSTEM,
symbol: "ensureRunFinalized",
},
];

export function expectedCompleteWaitpointCallSites(): { site: string; symbol: string }[] {
Expand Down
8 changes: 8 additions & 0 deletions internal-packages/run-engine/src/engine/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,12 @@ export class RunEngine {
runId: payload.runId,
});
},
ensureRunFinalized: async ({ payload }) => {
await this.runAttemptSystem.ensureRunFinalized({
runId: payload.runId,
deferCount: payload.deferCount,
});
},
enqueueDelayedRun: async ({ payload }) => {
await this.delayedRunSystem.enqueueDelayedRun({ runId: payload.runId });
},
Expand Down Expand Up @@ -422,6 +428,7 @@ export class RunEngine {
this.ttlSystem = new TtlSystem({
resources,
waitpointSystem: this.waitpointSystem,
finalizationGuardDelayMs: this.options.finalizationGuardDelayMs,
});

const ttlWorkerCatalog = createTtlWorkerCatalog({
Expand Down Expand Up @@ -505,6 +512,7 @@ export class RunEngine {
delayedRunSystem: this.delayedRunSystem,
machines: this.options.machines,
retryWarmStartThresholdMs: this.options.retryWarmStartThresholdMs,
finalizationGuardDelayMs: this.options.finalizationGuardDelayMs,
redisOptions: this.options.cache?.redis ?? this.options.runLock.redis,
});

Expand Down
190 changes: 189 additions & 1 deletion internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import {
RedisCacheStore,
} from "@internal/cache";
import type { RedisOptions } from "@internal/redis";
import { startSpan } from "@internal/tracing";
import { startSpan, type Counter } from "@internal/tracing";
import { tryCatch } from "@trigger.dev/core/utils";
import type {
CompleteRunAttemptResult,
Expand Down Expand Up @@ -41,6 +41,7 @@ import { getMachinePreset, machinePresetFromName } from "../machinePresets.js";
import { retryOutcomeFromCompletion } from "../retrying.js";
import {
isExecuting,
isFinalRunStatus,
isFinishedOrPendingFinished,
isInitialState,
isPendingExecuting,
Expand All @@ -67,6 +68,7 @@ export type RunAttemptSystemOptions = {
waitpointSystem: WaitpointSystem;
delayedRunSystem: DelayedRunSystem;
retryWarmStartThresholdMs?: number;
finalizationGuardDelayMs?: number;
machines: RunEngineOptions["machines"];
redisOptions: RedisOptions;
};
Expand Down Expand Up @@ -102,12 +104,22 @@ const DEPLOYMENT_STALE_TTL = 60000 * 60 * 24 * 2; // 2 days
const QUEUE_FRESH_TTL = 60000 * 60; // 1 hour
const QUEUE_STALE_TTL = 60000 * 60 * 2; // 2 hours

/**
* How many times the finalization guard defers to an in-flight cancellation before
* delivering anyway. The worker-owned states normally exit within a heartbeat cycle
* or two, so exhausting this budget means the heartbeat itself was lost; resuming the
* parent with the identical cancel error then beats watching forever.
*/
const MAX_FINALIZATION_GUARD_DEFERRALS = 10;

export class RunAttemptSystem {
private readonly $: SystemResources;
private readonly executionSnapshotSystem: ExecutionSnapshotSystem;
private readonly batchSystem: BatchSystem;
private readonly waitpointSystem: WaitpointSystem;
private readonly delayedRunSystem: DelayedRunSystem;
private readonly finalizationGuardDelayMs: number;
private readonly rederivationsCounter: Counter;
private readonly cache: UnkeyCache<{
tasks: BackwardsCompatibleTaskRunExecution["task"];
machinePresets: MachinePreset;
Expand All @@ -123,6 +135,15 @@ export class RunAttemptSystem {
this.batchSystem = options.batchSystem;
this.waitpointSystem = options.waitpointSystem;
this.delayedRunSystem = options.delayedRunSystem;
this.finalizationGuardDelayMs = options.finalizationGuardDelayMs ?? 60_000;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · HIGH

The finalization guard delay is duplicated as a magic default of 60_000 in both RunAttemptSystem and TtlSystem constructors, with no shared constant.

Impact: The finalization guard delay is duplicated as a magic default of 60_000 in both RunAttemptSystem and TtlSystem constructors, with no shared constant. A future change to one default will silently diverge the two systems and produce inconsistent guard timing.

Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright.

this.rederivationsCounter = this.$.meter.createCounter(
"run_attempt_system.finalization_rederivations",
{
description:
"Lost run-finalization side effects re-delivered by the ensureRunFinalized guard",
unit: "runs",
}
);

const ctx = new DefaultStatefulContext();
const memory = createLRUMemoryStore(5000);
Expand Down Expand Up @@ -757,6 +778,8 @@ export class RunAttemptSystem {
environmentType: latestSnapshot.environmentType,
});

await this.#scheduleFinalizationGuard(runId);

const run = await this.$.runStore.completeAttemptSuccess(
runId,
{
Expand Down Expand Up @@ -1439,6 +1462,8 @@ export class RunAttemptSystem {
});
}

await this.#scheduleFinalizationGuard(runId);

const run = await this.$.runStore.cancelRun(
runId,
{
Expand Down Expand Up @@ -1654,6 +1679,8 @@ export class RunAttemptSystem {
environmentType: latestSnapshot.environmentType,
});

await this.#scheduleFinalizationGuard(runId);

//run permanently failed
const run = await this.$.runStore.failRunPermanently(
runId,
Expand Down Expand Up @@ -1769,6 +1796,167 @@ export class RunAttemptSystem {

//cancel the heartbeats
await this.$.worker.ack(`heartbeatSnapshot.${id}`);

await this.$.worker.ack(`ensureRunFinalized:${id}`);
}

/**
* Write-ahead guard for run finalization. Enqueued BEFORE the finish commit (so no
* finish write can exist without a durable watcher) and acked at the end of
* {@link #finalizeRun} once every inline side effect succeeded. It only ever
* executes when the inline path died in between.
*/
async #scheduleFinalizationGuard(runId: string, deferCount?: number): Promise<void> {
await this.$.worker.enqueue({
id: `ensureRunFinalized:${runId}`,
job: "ensureRunFinalized",
payload: { runId, deferCount },
availableAt: new Date(Date.now() + this.finalizationGuardDelayMs),
});
}

/**
* Re-delivers a finished run's finalization side effects: the queue ack, the
* associated waitpoint's completion, the parent unblock fan-out, and the batch
* completion nudge. Safe to run at-least-once and to race the inline path — every leg
* is idempotent, and completing an already-completed waitpoint still re-runs the
* blocked-run fan-out (which covers a lost `continueRunIfUnblocked` enqueue). A
* non-final run means the finish commit itself never landed; the caller's retry
* re-runs the whole completion, so there is nothing to re-deliver. A canceled run
* whose worker is still winding down re-arms the guard and waits: the cancellation
* finalize path owns that window, and completing early would resume the parent while
* the child is still running.
*/
public async ensureRunFinalized({
runId,
deferCount,
}: {
runId: string;
deferCount?: number;
}): Promise<void> {
return startSpan(this.$.tracer, "ensureRunFinalized", async (span) => {
span.setAttribute("runId", runId);

const run = await this.$.runStore.findRun(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · HIGH

The ensureRunFinalized method is public and directly callable from tests and other systems, but its cancellation deferral logic depends on deferCount being passed correctly.

Impact: The ensureRunFinalized method is public and directly callable from tests and other systems, but its cancellation deferral logic depends on deferCount being passed correctly. A caller that omits deferCount resets the budget to zero on every invocation, allowing an unbounded number of deferrals if the worker never reaches a finished execution state.

Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright.

{ id: runId },
{
select: {
id: true,
status: true,
output: true,
outputType: true,
error: true,
batchId: true,
runtimeEnvironmentId: true,
associatedWaitpoint: {
select: {
id: true,
status: true,
},
},
},
},
this.$.prisma
);

if (!run) {
this.$.logger.error("ensureRunFinalized: run not found", { runId });
return;
}

if (!isFinalRunStatus(run.status)) {
this.$.logger.debug("ensureRunFinalized: run is not final, nothing to re-deliver", {
runId,
status: run.status,
});
return;
}

if (run.status === "CANCELED") {
const latestSnapshot = await getLatestExecutionSnapshot(
this.$.prisma,
runId,
this.$.runStore
);

/**
* Defer only while a worker still owns the execution: those states carry
* heartbeats that force the cancellation finalize path, which completes the
* waitpoint with the run's actual wind-down and acks this guard, so the watch
* always terminates. In any other snapshot state (queued, delayed, suspended,
* created) nobody is left to produce a FINISHED snapshot for a canceled run,
* so the guard must deliver or the parent is stranded.
*/
const workerOwnsExecution =
isExecuting(latestSnapshot.executionStatus) ||
isPendingExecuting(latestSnapshot.executionStatus) ||
latestSnapshot.executionStatus === "PENDING_CANCEL";

if (latestSnapshot.executionStatus !== "FINISHED" && workerOwnsExecution) {
const currentDeferCount = deferCount ?? 0;

if (currentDeferCount < MAX_FINALIZATION_GUARD_DEFERRALS) {
this.$.logger.info(
"ensureRunFinalized: run is canceled but the worker is still winding down, keeping watch until the cancellation finalize path completes it",
{
runId,
executionStatus: latestSnapshot.executionStatus,
deferCount: currentDeferCount,
}
);
await this.#scheduleFinalizationGuard(runId, currentDeferCount + 1);
return;
}

this.rederivationsCounter.add(1, { leg: "cancel_deferral_budget" });
this.$.logger.warn(
"ensureRunFinalized: canceled run never reached a finished execution within the deferral budget, delivering anyway",
{
runId,
executionStatus: latestSnapshot.executionStatus,
deferCount: currentDeferCount,
}
);
}
}

span.setAttribute("runStatus", run.status);

const env = await this.$.controlPlaneResolver.resolveEnv(run.runtimeEnvironmentId);

if (env) {
await this.$.runQueue.acknowledgeMessage(env.organizationId, runId, {
removeFromWorkerQueue: true,
});
} else {
this.$.logger.error("ensureRunFinalized: environment not found, skipping queue ack", {
runId,
runtimeEnvironmentId: run.runtimeEnvironmentId,
});
}

if (run.associatedWaitpoint) {
const wasPending = run.associatedWaitpoint.status === "PENDING";

if (wasPending) {
this.rederivationsCounter.add(1, { leg: "waitpoint" });
this.$.logger.warn("ensureRunFinalized: re-deriving lost waitpoint completion", {
runId,
runStatus: run.status,
waitpointId: run.associatedWaitpoint.id,
});
}

await this.waitpointSystem.completeWaitpoint({
id: run.associatedWaitpoint.id,
output: this.waitpointSystem.buildWaitpointOutputFromRun(run),
});
}

if (run.batchId) {
await this.batchSystem.scheduleCompleteBatch({ batchId: run.batchId });
}
});
}

async #resolveTaskRunExecutionTask(
Expand Down
36 changes: 36 additions & 0 deletions internal-packages/run-engine/src/engine/systems/ttlSystem.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,33 @@ import { boundedIn } from "@trigger.dev/database";
export type TtlSystemOptions = {
resources: SystemResources;
waitpointSystem: WaitpointSystem;
finalizationGuardDelayMs?: number;
};

export class TtlSystem {
private readonly $: SystemResources;
private readonly waitpointSystem: WaitpointSystem;
private readonly finalizationGuardDelayMs: number;

constructor(private readonly options: TtlSystemOptions) {
this.$ = options.resources;
this.waitpointSystem = options.waitpointSystem;
this.finalizationGuardDelayMs = options.finalizationGuardDelayMs ?? 60_000;
}

/**
* Write-ahead guard for TTL expiry, mirroring the run attempt system's: enqueued
* before the EXPIRED commit so a crash or error between that commit and the
* waitpoint completion cannot strand a waiting parent, and acked once the inline
* side effects succeed.
*/
async #scheduleFinalizationGuard(runId: string): Promise<void> {
await this.$.worker.enqueue({
id: `ensureRunFinalized:${runId}`,
job: "ensureRunFinalized",
payload: { runId },
availableAt: new Date(Date.now() + this.finalizationGuardDelayMs),
});
}

async expireRun({ runId, tx }: { runId: string; tx?: PrismaClientOrTransaction }) {
Expand Down Expand Up @@ -62,6 +80,8 @@ export class TtlSystem {
raw: `Run expired because the TTL (${run.ttl}) was reached`,
};

await this.#scheduleFinalizationGuard(runId);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · CRITICAL

In the single-run TTL path, the guard is scheduled before expireRun commits, but the ack only runs after expireRun succeeds.

Impact: In the single-run TTL path, the guard is scheduled before expireRun commits, but the ack only runs after expireRun succeeds. If expireRun throws, the guard remains armed and later fires ensureRunFinalized for a run that was never expired. ensureRunFinalized only checks isFinalRunStatus, so a non-final run returns early, but the stale guard can still race a later legitimate finalization and re-deliver side effects fo…

Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · MEDIUM

The finalization guard is scheduled before the run is expired, but the ack in the single-run path is only executed inside the transaction callback after expireRun succeeds.

Impact: The finalization guard is scheduled before the run is expired, but the ack in the single-run path is only executed inside the transaction callback after expireRun succeeds. If expireRun throws, the guard remains armed and will later fire ensureRunFinalized for a run that was never expired, potentially causing incorrect finalization side effects.

Suggested fix: Fix the review finding before release.


const updatedRun = await this.$.runStore.expireRun(
runId,
{
Expand Down Expand Up @@ -121,6 +141,8 @@ export class TtlSystem {
project: { id: snapshot.projectId },
environment: { id: snapshot.environmentId },
});

await this.$.worker.ack(`ensureRunFinalized:${runId}`);
});
}

Expand Down Expand Up @@ -220,6 +242,10 @@ export class TtlSystem {
raw: "Run expired because the TTL was reached",
};

await pMap(runsToExpire, (run) => this.#scheduleFinalizationGuard(run.id), {
concurrency: 10,
});

await this.$.runStore.expireRunsBatch(runIdsToExpire, { error, now }, this.$.prisma);

// Process each run: enqueue waitpoint completion jobs and emit events
Expand Down Expand Up @@ -262,6 +288,16 @@ export class TtlSystem {
environment: { id: run.runtimeEnvironmentId },
});

/**

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · CRITICAL

The batch TTL path acks the finalization guard for runs without an associatedWaitpoint immediately after enqueueing finishWaitpoint, but the comment claims the guard stays armed fo

Impact: The batch TTL path acks the finalization guard for runs without an associatedWaitpoint immediately after enqueueing finishWaitpoint, but the comment claims the guard stays armed for runs with a waiting parent. The condition checks associatedWaitpoint, not whether a parent is actually waiting. A run with a waiting parent but no associatedWaitpoint will have its guard released prematurely, re-introducing the exact str…

Suggested fix: Review the cited evidence, fix the risk if confirmed, and rerun Shipwright.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shipwright · HIGH

In the batch path, the guard is acked for runs without an associatedWaitpoint immediately after enqueueing the finishWaitpoint job.

Impact: In the batch path, the guard is acked for runs without an associatedWaitpoint immediately after enqueueing the finishWaitpoint job. However, the comment says the guard stays armed for runs with a waiting parent and verifies the completion landed. The condition checks associatedWaitpoint, not whether the run has a waiting parent, so a run with an associatedWaitpoint that is not actually waiting on a parent will keep…

Suggested fix: Fix the review finding before release.

* Waitpoint completion in this path is delegated to the finishWaitpoint job,
* which can still exhaust its retries, so the guard stays armed for runs with
* a waiting parent and verifies the completion landed. Runs with no waitpoint
* have nothing left to re-deliver, so release their guard now.
*/
if (!run.associatedWaitpoint) {
await this.$.worker.ack(`ensureRunFinalized:${run.id}`);
}

expired.push(run.id);
} catch (e) {
this.$.logger.error("Failed to process expired run", {
Expand Down
Loading