require suspendable state for checkpoints, fix snapshot processing queue

This commit is contained in:
nicktrn
2025-04-30 16:54:02 +01:00
parent bb9ac5026a
commit 3d978d337a
5 changed files with 526 additions and 227 deletions
@@ -401,46 +401,64 @@ export class ManagedRunController {
}) satisfies SupervisorSocket;
socket.on("run:notify", async ({ version, run }) => {
// Generate a unique ID for the notification
const notificationId = Math.random().toString(36).substring(2, 15);
// Use this to track the notification incl. any processing
const notification = {
id: notificationId,
runId: run.friendlyId,
version,
};
// Lock this to the current run and snapshot IDs
const controller = {
runFriendlyId: this.runFriendlyId,
snapshotFriendlyId: this.snapshotFriendlyId,
};
this.sendDebugLog({
runId: run.friendlyId,
message: "run:notify received by runner",
properties: { version, runId: run.friendlyId },
properties: {
notification,
controller,
},
});
if (!this.runFriendlyId) {
if (!controller.runFriendlyId) {
this.sendDebugLog({
runId: run.friendlyId,
message: "run:notify: ignoring notification, no local run ID",
properties: {
currentRunId: this.runFriendlyId,
currentSnapshotId: this.snapshotFriendlyId,
notification,
controller,
},
});
return;
}
if (run.friendlyId !== this.runFriendlyId) {
if (run.friendlyId !== controller.runFriendlyId) {
this.sendDebugLog({
runId: run.friendlyId,
message: "run:notify: ignoring notification for different run",
properties: {
currentRunId: this.runFriendlyId,
currentSnapshotId: this.snapshotFriendlyId,
notificationRunId: run.friendlyId,
notification,
controller,
},
});
return;
}
const latestSnapshot = await this.httpClient.getRunExecutionData(this.runFriendlyId);
const latestSnapshot = await this.httpClient.getRunExecutionData(controller.runFriendlyId);
if (!latestSnapshot.success) {
this.sendDebugLog({
runId: this.runFriendlyId,
runId: run.friendlyId,
message: "run:notify: failed to get latest snapshot data",
properties: {
currentRunId: this.runFriendlyId,
currentSnapshotId: this.snapshotFriendlyId,
notification,
controller,
error: latestSnapshot.error,
},
});
@@ -451,19 +469,29 @@ export class ManagedRunController {
if (!this.currentExecution) {
this.sendDebugLog({
runId: runExecutionData.run.friendlyId,
message: "handleSnapshotChange: no current execution",
runId: run.friendlyId,
message: "run:notify: no current execution",
properties: {
notification,
controller,
},
});
return;
}
const [error] = await tryCatch(this.currentExecution.handleSnapshotChange(runExecutionData));
const [error] = await tryCatch(
this.currentExecution.enqueueSnapshotChangeAndWait(runExecutionData)
);
if (error) {
this.sendDebugLog({
runId: runExecutionData.run.friendlyId,
message: "handleSnapshotChange: unexpected error",
properties: { error: error.message },
runId: run.friendlyId,
message: "run:notify: unexpected error",
properties: {
notification,
controller,
error: error.message,
},
});
}
});
@@ -32,7 +32,6 @@ const Env = z.object({
TRIGGER_MACHINE_MEMORY: z.string().default("0"),
TRIGGER_RUNNER_ID: z.string(),
TRIGGER_METADATA_URL: z.string().optional(),
TRIGGER_PRE_SUSPEND_WAIT_MS: z.coerce.number().default(200),
// Timeline metrics
TRIGGER_POD_SCHEDULED_AT_MS: DateEnv,
@@ -119,9 +118,6 @@ export class RunnerEnv {
get TRIGGER_METADATA_URL() {
return this.env.TRIGGER_METADATA_URL;
}
get TRIGGER_PRE_SUSPEND_WAIT_MS() {
return this.env.TRIGGER_PRE_SUSPEND_WAIT_MS;
}
get TRIGGER_POD_SCHEDULED_AT_MS() {
return this.env.TRIGGER_POD_SCHEDULED_AT_MS;
}
@@ -5,6 +5,7 @@ import {
type TaskRunExecutionMetrics,
type TaskRunExecutionResult,
TaskRunExecutionRetry,
TaskRunExecutionStatus,
type TaskRunFailedExecutionResult,
WorkerManifest,
} from "@trigger.dev/core/v3";
@@ -18,6 +19,7 @@ import { RunExecutionSnapshotPoller } from "./poller.js";
import { assertExhaustive, tryCatch } from "@trigger.dev/core/utils";
import { MetadataClient } from "./overrides.js";
import { randomBytes } from "node:crypto";
import { SnapshotManager, SnapshotState } from "./snapshotManager.js";
class ExecutionAbortError extends Error {
constructor(message: string) {
@@ -50,9 +52,9 @@ export class RunExecution {
private executionAbortController: AbortController;
private _runFriendlyId?: string;
private currentSnapshotId?: string;
private currentAttemptNumber?: number;
private currentTaskRunEnv?: Record<string, string>;
private snapshotManager?: SnapshotManager;
private dequeuedAt?: Date;
private podScheduledAt?: Date;
@@ -148,6 +150,10 @@ export class RunExecution {
this.sendRuntimeDebugLog(debugLog.message, debugLog.properties);
});
taskRunProcess.onSetSuspendable.attach(async ({ suspendable }) => {
this.suspendable = suspendable;
});
return taskRunProcess;
}
@@ -165,71 +171,22 @@ export class RunExecution {
/**
* Called by the RunController when it receives a websocket notification
* or when the snapshot poller detects a change
* or when the snapshot poller detects a change.
*
* This is the main entry point for snapshot changes, but processing is deferred to the snapshot manager.
*/
public async handleSnapshotChange(runData: RunExecutionData): Promise<void> {
public async enqueueSnapshotChangeAndWait(runData: RunExecutionData): Promise<void> {
if (this.isShuttingDown) {
this.sendDebugLog("handleSnapshotChange: shutting down, skipping");
this.sendDebugLog("enqueueSnapshotChangeAndWait: shutting down, skipping");
return;
}
const { run, snapshot, completedWaitpoints } = runData;
const snapshotMetadata = {
incomingRunId: run.friendlyId,
incomingSnapshotId: snapshot.friendlyId,
completedWaitpoints: completedWaitpoints.length,
};
// Ensure we have run details
if (!this.runFriendlyId || !this.currentSnapshotId) {
this.sendDebugLog(
"handleSnapshotChange: missing run or snapshot ID",
snapshotMetadata,
run.friendlyId
);
if (!this.snapshotManager) {
this.sendDebugLog("enqueueSnapshotChangeAndWait: missing snapshot manager");
return;
}
// Ensure the run ID matches
if (run.friendlyId !== this.runFriendlyId) {
// Send debug log to both runs
this.sendDebugLog("handleSnapshotChange: mismatched run IDs", snapshotMetadata);
this.sendDebugLog(
"handleSnapshotChange: mismatched run IDs",
snapshotMetadata,
run.friendlyId
);
return;
}
this.snapshotChangeQueue.push(runData);
await this.processSnapshotChangeQueue();
}
private snapshotChangeQueue: RunExecutionData[] = [];
private snapshotChangeQueueLock = false;
private async processSnapshotChangeQueue() {
if (this.snapshotChangeQueueLock) {
return;
}
this.snapshotChangeQueueLock = true;
while (this.snapshotChangeQueue.length > 0) {
const runData = this.snapshotChangeQueue.shift();
if (!runData) {
continue;
}
const [error] = await tryCatch(this.processSnapshotChange(runData));
if (error) {
this.sendDebugLog("Failed to process snapshot change", { error: error.message });
}
}
this.snapshotChangeQueueLock = false;
await this.snapshotManager.handleSnapshotChange(runData);
}
private async processSnapshotChange(runData: RunExecutionData): Promise<void> {
@@ -240,21 +197,13 @@ export class RunExecution {
completedWaitpoints: completedWaitpoints.length,
};
// Check if the incoming snapshot is newer than the current one
if (!this.currentSnapshotId || snapshot.friendlyId < this.currentSnapshotId) {
this.sendDebugLog(
"handleSnapshotChange: received older snapshot, skipping",
snapshotMetadata
);
return;
}
if (snapshot.friendlyId === this.currentSnapshotId) {
if (!this.snapshotManager) {
this.sendDebugLog("handleSnapshotChange: missing snapshot manager", snapshotMetadata);
return;
}
if (this.currentAttemptNumber && this.currentAttemptNumber !== run.attemptNumber) {
this.sendDebugLog("ERROR: attempt number mismatch", snapshotMetadata);
this.sendDebugLog("error: attempt number mismatch", snapshotMetadata);
await this.taskRunProcess?.suspend();
return;
}
@@ -264,12 +213,6 @@ export class RunExecution {
// Reset the snapshot poll interval so we don't do unnecessary work
this.snapshotPoller?.resetCurrentInterval();
// Update internal state
this.currentSnapshotId = snapshot.friendlyId;
// Update services
this.snapshotPoller?.updateSnapshotId(snapshot.friendlyId);
switch (snapshot.executionStatus) {
case "PENDING_CANCEL": {
const [error] = await tryCatch(this.cancel());
@@ -285,14 +228,14 @@ export class RunExecution {
return;
}
case "QUEUED": {
this.sendDebugLog("Run was re-queued", snapshotMetadata);
this.sendDebugLog("run was re-queued", snapshotMetadata);
// Pretend we've just suspended the run. This will kill the process without failing the run.
await this.taskRunProcess?.suspend();
return;
}
case "FINISHED": {
this.sendDebugLog("Run is finished", snapshotMetadata);
this.sendDebugLog("run is finished", snapshotMetadata);
// Pretend we've just suspended the run. This will kill the process without failing the run.
await this.taskRunProcess?.suspend();
@@ -300,80 +243,13 @@ export class RunExecution {
}
case "QUEUED_EXECUTING":
case "EXECUTING_WITH_WAITPOINTS": {
this.sendDebugLog("Run is executing with waitpoints", snapshotMetadata);
this.sendDebugLog("run is executing with waitpoints", snapshotMetadata);
const [error] = await tryCatch(this.taskRunProcess?.cleanup(false));
if (error) {
this.sendDebugLog("Failed to cleanup task run process, carrying on", {
...snapshotMetadata,
error: error.message,
});
}
if (snapshot.friendlyId !== this.currentSnapshotId) {
this.sendDebugLog("Snapshot changed after cleanup, abort", snapshotMetadata);
this.abortExecution();
return;
}
await sleep(this.env.TRIGGER_PRE_SUSPEND_WAIT_MS);
if (snapshot.friendlyId !== this.currentSnapshotId) {
this.sendDebugLog("Snapshot changed after suspend threshold, abort", snapshotMetadata);
this.abortExecution();
return;
}
if (!this.runFriendlyId || !this.currentSnapshotId) {
this.sendDebugLog(
"handleSnapshotChange: Missing run ID or snapshot ID after suspension, abort",
snapshotMetadata
);
this.abortExecution();
return;
}
const suspendResult = await this.httpClient.suspendRun(
this.runFriendlyId,
this.currentSnapshotId
);
if (!suspendResult.success) {
this.sendDebugLog("Failed to suspend run, staying alive 🎶", {
...snapshotMetadata,
error: suspendResult.error,
});
this.sendDebugLog("checkpoint: suspend request failed", {
...snapshotMetadata,
error: suspendResult.error,
});
// This is fine, we'll wait for the next status change
return;
}
if (!suspendResult.data.ok) {
this.sendDebugLog("checkpoint: failed to suspend run", {
snapshotId: this.currentSnapshotId,
error: suspendResult.data.error,
});
// This is fine, we'll wait for the next status change
return;
}
this.sendDebugLog("Suspending, any day now 🚬", snapshotMetadata);
// Wait for next status change
// Wait for next status change - suspension is handled by the snapshot manager
return;
}
case "SUSPENDED": {
this.sendDebugLog("Run was suspended, kill the process", snapshotMetadata);
this.sendDebugLog("run was suspended, kill the process", snapshotMetadata);
// This will kill the process and fail the execution with a SuspendedProcessError
await this.taskRunProcess?.suspend();
@@ -381,17 +257,17 @@ export class RunExecution {
return;
}
case "PENDING_EXECUTING": {
this.sendDebugLog("Run is pending execution", snapshotMetadata);
this.sendDebugLog("run is pending execution", snapshotMetadata);
if (completedWaitpoints.length === 0) {
this.sendDebugLog("No waitpoints to complete, nothing to do", snapshotMetadata);
this.sendDebugLog("no waitpoints to complete, nothing to do", snapshotMetadata);
return;
}
const [error] = await tryCatch(this.restore());
if (error) {
this.sendDebugLog("Failed to restore execution", {
this.sendDebugLog("failed to restore execution", {
...snapshotMetadata,
error: error.message,
});
@@ -403,16 +279,16 @@ export class RunExecution {
return;
}
case "EXECUTING": {
this.sendDebugLog("Run is now executing", snapshotMetadata);
this.sendDebugLog("run is now executing", snapshotMetadata);
if (completedWaitpoints.length === 0) {
return;
}
this.sendDebugLog("Processing completed waitpoints", snapshotMetadata);
this.sendDebugLog("processing completed waitpoints", snapshotMetadata);
if (!this.taskRunProcess) {
this.sendDebugLog("No task run process, ignoring completed waitpoints", snapshotMetadata);
this.sendDebugLog("no task run process, ignoring completed waitpoints", snapshotMetadata);
this.abortExecution();
return;
@@ -425,7 +301,7 @@ export class RunExecution {
return;
}
case "RUN_CREATED": {
this.sendDebugLog("Invalid status change", snapshotMetadata);
this.sendDebugLog("invalid status change", snapshotMetadata);
this.abortExecution();
return;
@@ -441,11 +317,11 @@ export class RunExecution {
}: {
isWarmStart?: boolean;
}): Promise<WorkloadRunAttemptStartResponseBody & { metrics: TaskRunExecutionMetrics }> {
if (!this.runFriendlyId || !this.currentSnapshotId) {
throw new Error("Cannot start attempt: missing run or snapshot ID");
if (!this.runFriendlyId || !this.snapshotManager) {
throw new Error("Cannot start attempt: missing run or snapshot manager");
}
this.sendDebugLog("Starting attempt");
this.sendDebugLog("starting attempt");
const attemptStartedAt = Date.now();
@@ -456,7 +332,7 @@ export class RunExecution {
const start = await this.httpClient.startRunAttempt(
this.runFriendlyId,
this.currentSnapshotId,
this.snapshotManager.snapshotId,
{ isWarmStart }
);
@@ -469,14 +345,17 @@ export class RunExecution {
}
// A snapshot was just created, so update the snapshot ID
this.currentSnapshotId = start.data.snapshot.friendlyId;
this.snapshotManager.updateSnapshot(
start.data.snapshot.friendlyId,
start.data.snapshot.executionStatus
);
// Also set or update the attempt number - we do this to detect illegal attempt number changes, e.g. from stalled runners coming back online
const attemptNumber = start.data.run.attemptNumber;
if (attemptNumber && attemptNumber > 0) {
this.currentAttemptNumber = attemptNumber;
} else {
this.sendDebugLog("ERROR: invalid attempt number returned from start attempt", {
this.sendDebugLog("error: invalid attempt number returned from start attempt", {
attemptNumber: String(attemptNumber),
});
}
@@ -487,7 +366,7 @@ export class RunExecution {
podScheduledAt: this.podScheduledAt?.getTime(),
});
this.sendDebugLog("Started attempt");
this.sendDebugLog("started attempt");
return { ...start.data, metrics };
}
@@ -499,18 +378,29 @@ export class RunExecution {
public async execute(runOpts: RunExecutionRunOptions): Promise<void> {
// Setup initial state
this.runFriendlyId = runOpts.runFriendlyId;
this.currentSnapshotId = runOpts.snapshotFriendlyId;
// Create snapshot manager
this.snapshotManager = new SnapshotManager({
runFriendlyId: runOpts.runFriendlyId,
initialSnapshotId: runOpts.snapshotFriendlyId,
// We're just guessing here, but "PENDING_EXECUTING" is probably fine
initialStatus: "PENDING_EXECUTING",
logger: this.logger,
onSnapshotChange: this.processSnapshotChange.bind(this),
onSuspendable: this.handleSuspendable.bind(this),
});
this.dequeuedAt = runOpts.dequeuedAt;
this.podScheduledAt = runOpts.podScheduledAt;
// Create and start services
this.snapshotPoller = new RunExecutionSnapshotPoller({
runFriendlyId: this.runFriendlyId,
snapshotFriendlyId: this.currentSnapshotId,
snapshotFriendlyId: this.snapshotManager.snapshotId,
httpClient: this.httpClient,
logger: this.logger,
snapshotPollIntervalSeconds: this.env.TRIGGER_SNAPSHOT_POLL_INTERVAL_SECONDS,
handleSnapshotChange: this.handleSnapshotChange.bind(this),
handleSnapshotChange: this.enqueueSnapshotChangeAndWait.bind(this),
});
this.snapshotPoller.start();
@@ -520,7 +410,7 @@ export class RunExecution {
);
if (startError) {
this.sendDebugLog("Failed to start attempt", { error: startError.message });
this.sendDebugLog("failed to start attempt", { error: startError.message });
this.stopServices();
return;
@@ -529,7 +419,7 @@ export class RunExecution {
const [executeError] = await tryCatch(this.executeRunWrapper(start));
if (executeError) {
this.sendDebugLog("Failed to execute run", { error: executeError.message });
this.sendDebugLog("failed to execute run", { error: executeError.message });
this.stopServices();
return;
@@ -562,7 +452,7 @@ export class RunExecution {
})
);
this.sendDebugLog("Run execution completed", { error: executeError?.message });
this.sendDebugLog("run execution completed", { error: executeError?.message });
if (!executeError) {
this.stopServices();
@@ -570,7 +460,7 @@ export class RunExecution {
}
if (executeError instanceof SuspendedProcessError) {
this.sendDebugLog("Run was suspended", {
this.sendDebugLog("run was suspended", {
run: run.friendlyId,
snapshot: snapshot.friendlyId,
error: executeError.message,
@@ -580,7 +470,7 @@ export class RunExecution {
}
if (executeError instanceof ExecutionAbortError) {
this.sendDebugLog("Run was interrupted", {
this.sendDebugLog("run was interrupted", {
run: run.friendlyId,
snapshot: snapshot.friendlyId,
error: executeError.message,
@@ -589,7 +479,7 @@ export class RunExecution {
return;
}
this.sendDebugLog("Error while executing attempt", {
this.sendDebugLog("error while executing attempt", {
error: executeError.message,
runId: run.friendlyId,
snapshotId: snapshot.friendlyId,
@@ -605,7 +495,7 @@ export class RunExecution {
const [completeError] = await tryCatch(this.complete({ completion }));
if (completeError) {
this.sendDebugLog("Failed to complete run", { error: completeError.message });
this.sendDebugLog("failed to complete run", { error: completeError.message });
}
this.stopServices();
@@ -641,7 +531,7 @@ export class RunExecution {
// Set up an abort handler that will cleanup the task run process
this.executionAbortController.signal.addEventListener("abort", async () => {
this.sendDebugLog("Execution aborted during task run, cleaning up process", {
this.sendDebugLog("execution aborted during task run, cleaning up process", {
runId: execution.run.id,
});
@@ -662,13 +552,13 @@ export class RunExecution {
);
// If we get here, the task completed normally
this.sendDebugLog("Completed run attempt", { attemptSuccess: completion.ok });
this.sendDebugLog("completed run attempt", { attemptSuccess: completion.ok });
// The execution has finished, so we can cleanup the task run process. Killing it should be safe.
const [error] = await tryCatch(this.taskRunProcess.cleanup(true));
if (error) {
this.sendDebugLog("Failed to cleanup task run process, submitting completion anyway", {
this.sendDebugLog("failed to cleanup task run process, submitting completion anyway", {
error: error.message,
});
}
@@ -676,7 +566,7 @@ export class RunExecution {
const [completionError] = await tryCatch(this.complete({ completion }));
if (completionError) {
this.sendDebugLog("Failed to complete run", { error: completionError.message });
this.sendDebugLog("failed to complete run", { error: completionError.message });
}
}
@@ -700,13 +590,13 @@ export class RunExecution {
}
private async complete({ completion }: { completion: TaskRunExecutionResult }): Promise<void> {
if (!this.runFriendlyId || !this.currentSnapshotId) {
throw new Error("Cannot complete run: missing run or snapshot ID");
if (!this.runFriendlyId || !this.snapshotManager) {
throw new Error("cannot complete run: missing run or snapshot manager");
}
const completionResult = await this.httpClient.completeRunAttempt(
this.runFriendlyId,
this.currentSnapshotId,
this.snapshotManager.snapshotId,
{ completion }
);
@@ -727,32 +617,33 @@ export class RunExecution {
completion: TaskRunExecutionResult;
result: CompleteRunAttemptResult;
}) {
this.sendDebugLog("Handling completion result", {
this.sendDebugLog(`completion result: ${result.attemptStatus}`, {
attemptSuccess: completion.ok,
attemptStatus: result.attemptStatus,
snapshotId: result.snapshot.friendlyId,
runId: result.run.friendlyId,
});
// Update our snapshot ID to match the completion result
// This ensures any subsequent API calls use the correct snapshot
this.currentSnapshotId = result.snapshot.friendlyId;
const snapshotStatus = this.convertAttemptStatusToSnapshotStatus(result.attemptStatus);
// Update our snapshot ID to match the completion result to ensure any subsequent API calls use the correct snapshot
this.updateSnapshot(result.snapshot.friendlyId, snapshotStatus);
const { attemptStatus } = result;
if (attemptStatus === "RUN_FINISHED") {
this.sendDebugLog("Run finished");
this.sendDebugLog("run finished");
return;
}
if (attemptStatus === "RUN_PENDING_CANCEL") {
this.sendDebugLog("Run pending cancel");
this.sendDebugLog("run pending cancel");
return;
}
if (attemptStatus === "RETRY_QUEUED") {
this.sendDebugLog("Retry queued");
this.sendDebugLog("retry queued");
return;
}
@@ -773,6 +664,28 @@ export class RunExecution {
assertExhaustive(attemptStatus);
}
private updateSnapshot(snapshotId: string, status: TaskRunExecutionStatus) {
this.snapshotManager?.updateSnapshot(snapshotId, status);
this.snapshotPoller?.updateSnapshotId(snapshotId);
}
private convertAttemptStatusToSnapshotStatus(
attemptStatus: CompleteRunAttemptResult["attemptStatus"]
): TaskRunExecutionStatus {
switch (attemptStatus) {
case "RUN_FINISHED":
return "FINISHED";
case "RUN_PENDING_CANCEL":
return "PENDING_CANCEL";
case "RETRY_QUEUED":
return "QUEUED";
case "RETRY_IMMEDIATELY":
return "EXECUTING";
default:
assertExhaustive(attemptStatus);
}
}
private measureExecutionMetrics({
attemptCreatedAt,
dequeuedAt,
@@ -813,7 +726,7 @@ export class RunExecution {
}
private async retryImmediately({ retryOpts }: { retryOpts: TaskRunExecutionRetry }) {
this.sendDebugLog("Retrying run immediately", {
this.sendDebugLog("retrying run immediately", {
timestamp: retryOpts.timestamp,
delay: retryOpts.delay,
});
@@ -829,7 +742,7 @@ export class RunExecution {
const [startError, start] = await tryCatch(this.startAttempt({ isWarmStart: true }));
if (startError) {
this.sendDebugLog("Failed to start attempt for retry", { error: startError.message });
this.sendDebugLog("failed to start attempt for retry", { error: startError.message });
this.stopServices();
return;
@@ -838,7 +751,7 @@ export class RunExecution {
const [executeError] = await tryCatch(this.executeRunWrapper({ ...start, isWarmStart: true }));
if (executeError) {
this.sendDebugLog("Failed to execute run for retry", { error: executeError.message });
this.sendDebugLog("failed to execute run for retry", { error: executeError.message });
this.stopServices();
return;
@@ -851,10 +764,10 @@ export class RunExecution {
* Restores a suspended execution from PENDING_EXECUTING
*/
private async restore(): Promise<void> {
this.sendDebugLog("Restoring execution");
this.sendDebugLog("restoring execution");
if (!this.runFriendlyId || !this.currentSnapshotId) {
throw new Error("Cannot restore: missing run or snapshot ID");
if (!this.runFriendlyId || !this.snapshotManager) {
throw new Error("Cannot restore: missing run or snapshot manager");
}
// Short delay to give websocket time to reconnect
@@ -865,7 +778,7 @@ export class RunExecution {
const continuationResult = await this.httpClient.continueRunExecution(
this.runFriendlyId,
this.currentSnapshotId
this.snapshotManager.snapshotId
);
if (!continuationResult.success) {
@@ -881,7 +794,7 @@ export class RunExecution {
*/
private async processEnvOverrides() {
if (!this.env.TRIGGER_METADATA_URL) {
this.sendDebugLog("No metadata URL, skipping env overrides");
this.sendDebugLog("no metadata url, skipping env overrides");
return;
}
@@ -889,11 +802,11 @@ export class RunExecution {
const overrides = await metadataClient.getEnvOverrides();
if (!overrides) {
this.sendDebugLog("No env overrides, skipping");
this.sendDebugLog("no env overrides, skipping");
return;
}
this.sendDebugLog("Processing env overrides", overrides);
this.sendDebugLog("processing env overrides", overrides);
// Override the env with the new values
this.env.override(overrides);
@@ -916,21 +829,24 @@ export class RunExecution {
private async onHeartbeat() {
if (!this.runFriendlyId) {
this.sendDebugLog("Heartbeat: missing run ID");
this.sendDebugLog("heartbeat: missing run ID");
return;
}
if (!this.currentSnapshotId) {
this.sendDebugLog("Heartbeat: missing snapshot ID");
if (!this.snapshotManager) {
this.sendDebugLog("heartbeat: missing snapshot manager");
return;
}
this.sendDebugLog("Heartbeat: started");
this.sendDebugLog("heartbeat");
const response = await this.httpClient.heartbeatRun(this.runFriendlyId, this.currentSnapshotId);
const response = await this.httpClient.heartbeatRun(
this.runFriendlyId,
this.snapshotManager.snapshotId
);
if (!response.success) {
this.sendDebugLog("Heartbeat: failed", { error: response.error });
this.sendDebugLog("heartbeat: failed", { error: response.error });
}
this.lastHeartbeat = new Date();
@@ -947,7 +863,7 @@ export class RunExecution {
properties: {
...properties,
runId: this.runFriendlyId,
snapshotId: this.currentSnapshotId,
snapshotId: this.currentSnapshotFriendlyId,
executionId: this.id,
executionRestoreCount: this.restoreCount,
lastHeartbeat: this.lastHeartbeat?.toISOString(),
@@ -967,7 +883,7 @@ export class RunExecution {
properties: {
...properties,
runId: this.runFriendlyId,
snapshotId: this.currentSnapshotId,
snapshotId: this.currentSnapshotFriendlyId,
executionId: this.id,
executionRestoreCount: this.restoreCount,
lastHeartbeat: this.lastHeartbeat?.toISOString(),
@@ -975,6 +891,10 @@ export class RunExecution {
});
}
private set suspendable(suspendable: boolean) {
this.snapshotManager?.setSuspendable(suspendable);
}
// Ensure we can only set this once
private set runFriendlyId(id: string) {
if (this._runFriendlyId) {
@@ -989,7 +909,7 @@ export class RunExecution {
}
public get currentSnapshotFriendlyId(): string | undefined {
return this.currentSnapshotId;
return this.snapshotManager?.snapshotId;
}
public get taskRunEnv(): Record<string, string> | undefined {
@@ -1008,7 +928,7 @@ export class RunExecution {
private abortExecution() {
if (this.isAborted) {
this.sendDebugLog("Execution already aborted");
this.sendDebugLog("execution already aborted");
return;
}
@@ -1024,5 +944,79 @@ export class RunExecution {
this.isShuttingDown = true;
this.snapshotPoller?.stop();
this.taskRunProcess?.onTaskRunHeartbeat.detach();
this.snapshotManager?.cleanup();
}
private async handleSuspendable(suspendableSnapshot: SnapshotState) {
this.sendDebugLog("handleSuspendable", { suspendableSnapshot });
if (!this.snapshotManager) {
this.sendDebugLog("handleSuspendable: missing snapshot manager");
return;
}
// Ensure this is the current snapshot
if (suspendableSnapshot.id !== this.currentSnapshotFriendlyId) {
this.sendDebugLog("snapshot changed before cleanup, abort", {
suspendableSnapshot,
currentSnapshotId: this.currentSnapshotFriendlyId,
});
this.abortExecution();
return;
}
// First cleanup the task run process
const [error] = await tryCatch(this.taskRunProcess?.cleanup(false));
if (error) {
this.sendDebugLog("failed to cleanup task run process, carrying on", {
suspendableSnapshot,
error: error.message,
});
}
// Double check snapshot hasn't changed after cleanup
if (suspendableSnapshot.id !== this.currentSnapshotFriendlyId) {
this.sendDebugLog("snapshot changed after cleanup, abort", {
suspendableSnapshot,
currentSnapshotId: this.currentSnapshotFriendlyId,
});
this.abortExecution();
return;
}
if (!this.runFriendlyId) {
this.sendDebugLog("missing run ID for suspension, abort", { suspendableSnapshot });
this.abortExecution();
return;
}
// Call the suspend API with the current snapshot ID
const suspendResult = await this.httpClient.suspendRun(
this.runFriendlyId,
suspendableSnapshot.id
);
if (!suspendResult.success) {
this.sendDebugLog("failed to suspend run, staying alive 🎶", {
suspendableSnapshot,
error: suspendResult.error,
});
// This is fine, we'll wait for the next status change
return;
}
if (!suspendResult.data.ok) {
this.sendDebugLog("checkpoint: failed to suspend run", {
suspendableSnapshot,
error: suspendResult.data.error,
});
// This is fine, we'll wait for the next status change
return;
}
this.sendDebugLog("suspending, any day now 🚬", { suspendableSnapshot });
}
}
@@ -84,6 +84,7 @@ export class RunExecutionSnapshotPoller {
this.poller.resetCurrentInterval();
}
// The snapshot ID is only used as an indicator of when a poller got stuck
updateSnapshotId(snapshotFriendlyId: string) {
this.snapshotFriendlyId = snapshotFriendlyId;
}
@@ -0,0 +1,280 @@
import { tryCatch } from "@trigger.dev/core/utils";
import { RunLogger, SendDebugLogOptions } from "./logger.js";
import { TaskRunExecutionStatus, type RunExecutionData } from "@trigger.dev/core/v3";
import { assertExhaustive } from "@trigger.dev/core/utils";
export type SnapshotState = {
id: string;
status: TaskRunExecutionStatus;
};
type SnapshotHandler = (runData: RunExecutionData) => Promise<void>;
type SuspendableHandler = (suspendableSnapshot: SnapshotState) => Promise<void>;
type SnapshotManagerOptions = {
runFriendlyId: string;
initialSnapshotId: string;
initialStatus: TaskRunExecutionStatus;
logger: RunLogger;
onSnapshotChange: SnapshotHandler;
onSuspendable: SuspendableHandler;
};
type QueuedChange =
| { type: "snapshot"; data: RunExecutionData }
| { type: "suspendable"; value: boolean };
type QueuedChangeItem = {
change: QueuedChange;
resolve: () => void;
reject: (error: Error) => void;
};
export class SnapshotManager {
private state: SnapshotState;
private runFriendlyId: string;
private logger: RunLogger;
private isSuspendable: boolean = false;
private readonly onSnapshotChange: SnapshotHandler;
private readonly onSuspendable: SuspendableHandler;
private changeQueue: QueuedChangeItem[] = [];
private isProcessingQueue = false;
constructor(opts: SnapshotManagerOptions) {
this.runFriendlyId = opts.runFriendlyId;
this.logger = opts.logger;
this.state = {
id: opts.initialSnapshotId,
status: opts.initialStatus,
};
this.onSnapshotChange = opts.onSnapshotChange;
this.onSuspendable = opts.onSuspendable;
}
public get snapshotId(): string {
return this.state.id;
}
public get status(): TaskRunExecutionStatus {
return this.state.status;
}
public get suspendable(): boolean {
return this.isSuspendable;
}
public async setSuspendable(suspendable: boolean): Promise<void> {
if (this.isSuspendable === suspendable) {
this.sendDebugLog(`skipping suspendable update, already ${suspendable}`);
return;
}
this.sendDebugLog(`setting suspendable to ${suspendable}`);
return this.enqueueSnapshotChange({ type: "suspendable", value: suspendable });
}
public updateSnapshot(snapshotId: string, status: TaskRunExecutionStatus) {
// Check if this is an old snapshot
if (snapshotId < this.state.id) {
this.sendDebugLog("skipping update for old snapshot", {
incomingId: snapshotId,
currentId: this.state.id,
});
return;
}
this.state = { id: snapshotId, status };
}
public async handleSnapshotChange(runData: RunExecutionData): Promise<void> {
if (!this.statusCheck(runData)) {
return;
}
return this.enqueueSnapshotChange({ type: "snapshot", data: runData });
}
private statusCheck(runData: RunExecutionData): boolean {
const { run, snapshot } = runData;
const statusCheckData = {
incomingId: snapshot.friendlyId,
incomingStatus: snapshot.executionStatus,
currentId: this.state.id,
currentStatus: this.state.status,
};
// Ensure run ID matches
if (run.friendlyId !== this.runFriendlyId) {
this.sendDebugLog("skipping update for mismatched run ID", {
statusCheckData,
});
return false;
}
// Skip if this is an old snapshot
if (snapshot.friendlyId < this.state.id) {
this.sendDebugLog("skipping update for old snapshot", {
statusCheckData,
});
return false;
}
// Skip if this is the current snapshot
if (snapshot.friendlyId === this.state.id) {
this.sendDebugLog("skipping update for duplicate snapshot", {
statusCheckData,
});
return false;
}
return true;
}
private async enqueueSnapshotChange(change: QueuedChange): Promise<void> {
return new Promise<void>((resolve, reject) => {
// For suspendable changes, resolve and remove any pending suspendable changes
// since only the last one matters
if (change.type === "suspendable") {
const pendingSuspendable = this.changeQueue.filter(
(item) => item.change.type === "suspendable"
);
// Resolve any pending suspendable changes - they're effectively done since we're superseding them
pendingSuspendable.forEach((item) => item.resolve());
// Remove them from the queue
this.changeQueue = this.changeQueue.filter((item) => item.change.type !== "suspendable");
}
this.changeQueue.push({ change, resolve, reject });
// Sort queue:
// 1. Suspendable changes always go to the back
// 2. Snapshot changes are ordered by ID
this.changeQueue.sort((a, b) => {
if (a.change.type === "suspendable" && b.change.type === "snapshot") {
return 1; // a goes after b
}
if (a.change.type === "snapshot" && b.change.type === "suspendable") {
return -1; // a goes before b
}
if (a.change.type === "snapshot" && b.change.type === "snapshot") {
// sort snapshot changes by creation time, CUIDs are sortable
return a.change.data.snapshot.friendlyId.localeCompare(b.change.data.snapshot.friendlyId);
}
return 0; // both suspendable, maintain insertion order
});
// Start processing if not already running
this.processQueue().catch((error) => {
this.sendDebugLog("error processing queue", { error: error.message });
});
});
}
private async processQueue() {
if (this.isProcessingQueue) {
return;
}
this.isProcessingQueue = true;
try {
while (this.changeQueue.length > 0) {
const item = this.changeQueue[0];
if (!item) {
break;
}
const [error] = await tryCatch(this.applyChange(item.change));
// Remove from queue and resolve/reject promise
this.changeQueue.shift();
if (error) {
item.reject(error);
} else {
item.resolve();
}
}
} finally {
this.isProcessingQueue = false;
}
}
private async applyChange(change: QueuedChange): Promise<void> {
switch (change.type) {
case "snapshot": {
// Double check we should process this snapshot
if (!this.statusCheck(change.data)) {
return;
}
const { snapshot } = change.data;
this.updateSnapshot(snapshot.friendlyId, snapshot.executionStatus);
const oldState = { ...this.state };
this.sendDebugLog(`status changed to ${snapshot.executionStatus}`, {
oldId: oldState.id,
newId: snapshot.friendlyId,
oldStatus: oldState.status,
newStatus: snapshot.executionStatus,
});
// Execute handler
await this.onSnapshotChange(change.data);
// Check suspendable state after snapshot change
await this.checkSuspendableState();
break;
}
case "suspendable": {
this.isSuspendable = change.value;
// Check suspendable state after suspendable change
await this.checkSuspendableState();
break;
}
default: {
assertExhaustive(change);
}
}
}
private async checkSuspendableState() {
if (
this.isSuspendable &&
(this.state.status === "EXECUTING_WITH_WAITPOINTS" ||
this.state.status === "QUEUED_EXECUTING")
) {
this.sendDebugLog("run is now suspendable, executing handler");
await this.onSuspendable(this.state);
}
}
public cleanup() {
// Clear any pending changes
this.changeQueue = [];
}
protected sendDebugLog(message: string, properties?: SendDebugLogOptions["properties"]) {
this.logger.sendDebugLog({
runId: this.runFriendlyId,
message: `[snapshot] ${message}`,
properties: {
...properties,
snapshotId: this.state.id,
status: this.state.status,
suspendable: this.isSuspendable,
queueLength: this.changeQueue.length,
isProcessingQueue: this.isProcessingQueue,
},
});
}
}