Files
triggerdotdev--trigger.dev/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts
T
Eric Allam dac9c83bdc chore(webapp,run-engine): downgrade boundary log noise to warn (#3462)
## Summary

Several boundary catches and customer-input validation paths were
logging at `error` level for failures the system already handles
gracefully — disconnect on auth failure, return undefined, skip retries,
etc. This batch routes them to `warn` (which stays in stdout) or counts
them as OTel metrics, so visibility is preserved without surfacing them
as alerts.

## Changes

**New helper / pattern:**
- `apiBuilder.server.ts` — `logBoundaryError(message, error, url)`
inspects the inner error type at loader/action boundary catches;
downgrades to `warn` for `AbortError`, `ServiceValidationError`, and
`EngineServiceValidationError`.
- `platform.v3.server.ts` — `platform_client.failures_total` OTel
counter with `{function, kind}` labels; helper
`recordPlatformFailure(fn, kind)` replaces the previous error-level
logging across all `BillingClient` wrappers.

**Log-level downgrades:**
- `handleSocketIo.server.ts` — `Worker authentication failed` → warn
(system disconnects on failure; refs TRI-8863)
- `waitpointSystem.ts` — when `runStatus === "CANCELED"` in the
suspended-without-checkpoint branch, skip the throw and warn instead
(benign cancel-vs-resume race, nothing to resume)
- `runAttemptSystem.ts` — `flushedMetadata` parse/validate failures →
warn (customer-side data shape, system returns gracefully)
- `batch-queue/index.ts` — final-attempt failures with
`result.skipRetries` → warn (callbacks already opted out of retry, e.g.
queue size limit hit)
- `queryPerformanceMonitor.server.ts` — slow queries → warn
(observability signal, not an application error)
- `timeoutDeployment.server.ts` — deployment-state mismatch in the
timeout job → warn (timeout-vs-completion race)

**Inner error preservation:**
- `waitpointCompletionPacket.server.ts` — `logger.error(uploadError)`
before throwing the `ServiceValidationError` wrapper, so the underlying
upload error stays visible.

## Why

The pattern across all of these is the same: a boundary log treated any
thrown/returned error as `error` regardless of cause, even when the
cause was an expected, system-handled condition (client disconnect,
customer quota, race condition, schema validation of customer data).
That made the logs noisy and made it harder to spot real bugs.

Where the underlying signal is still useful operationally (slow queries,
billing call failures), we route it to OTel metrics with low-cardinality
labels so dashboards and alerts can be tuned independently of error
logs.

## Test plan

- [ ] `pnpm run typecheck --filter webapp`
- [ ] `pnpm run build --filter @internal/run-engine`
- [ ] Trigger a run on hello-world and verify task lifecycle is
unaffected
- [ ] Cancel a suspended run and verify the cancel-while-suspended
branch in `waitpointSystem.ts` returns `{status: "skipped"}` instead of
throwing
- [ ] Confirm `platform_client.failures_total` counter shows up in
metrics with `{function, kind}` labels when the billing client errors
2026-04-29 10:00:22 +01:00

2104 lines
66 KiB
TypeScript

import {
createCache,
createLRUMemoryStore,
DefaultStatefulContext,
Namespace,
RedisCacheStore,
UnkeyCache,
} from "@internal/cache";
import { RedisOptions } from "@internal/redis";
import { startSpan } from "@internal/tracing";
import { tryCatch } from "@trigger.dev/core/utils";
import {
CompleteRunAttemptResult,
ExecutionResult,
FlushedRunMetadata,
GitMeta,
MachinePreset,
MachinePresetName,
StartRunAttemptResult,
TaskRunContext,
TaskRunError,
TaskRunExecution,
TaskRunExecutionDeployment,
TaskRunExecutionOrganization,
TaskRunExecutionProject,
TaskRunExecutionQueue,
TaskRunExecutionResult,
TaskRunFailedExecutionResult,
TaskRunInternalError,
TaskRunSuccessfulExecutionResult,
} from "@trigger.dev/core/v3/schemas";
import {
extractIdempotencyKeyScope,
getUserProvidedIdempotencyKey,
} from "@trigger.dev/core/v3/serverOnly";
import { parsePacket } from "@trigger.dev/core/v3/utils/ioSerialization";
import {
$transaction,
PrismaClientOrTransaction,
RuntimeEnvironmentType,
TaskRun,
} from "@trigger.dev/database";
import { MAX_TASK_RUN_ATTEMPTS } from "../consts.js";
import { runStatusFromError, ServiceValidationError } from "../errors.js";
import { sendNotificationToWorker } from "../eventBus.js";
import { getMachinePreset, machinePresetFromName } from "../machinePresets.js";
import { retryOutcomeFromCompletion } from "../retrying.js";
import {
isExecuting,
isFinishedOrPendingFinished,
isInitialState,
isPendingExecuting,
} from "../statuses.js";
import { RunEngineOptions } from "../types.js";
import { BatchSystem } from "./batchSystem.js";
import { DelayedRunSystem } from "./delayedRunSystem.js";
import {
EnhancedExecutionSnapshot,
executionResultFromSnapshot,
ExecutionSnapshotSystem,
getLatestExecutionSnapshot,
} from "./executionSnapshotSystem.js";
import { SystemResources } from "./systems.js";
import { WaitpointSystem } from "./waitpointSystem.js";
import { BatchId, RunId } from "@trigger.dev/core/v3/isomorphic";
export type RunAttemptSystemOptions = {
resources: SystemResources;
executionSnapshotSystem: ExecutionSnapshotSystem;
batchSystem: BatchSystem;
waitpointSystem: WaitpointSystem;
delayedRunSystem: DelayedRunSystem;
retryWarmStartThresholdMs?: number;
machines: RunEngineOptions["machines"];
redisOptions: RedisOptions;
};
type BackwardsCompatibleTaskRunExecution = Omit<TaskRunExecution, "task" | "attempt" | "run"> & {
task: TaskRunExecution["task"] & {
exportName: string | undefined;
};
attempt: TaskRunExecution["attempt"] & {
id: string;
backgroundWorkerId: string;
backgroundWorkerTaskId: string;
status: string;
};
run: TaskRunExecution["run"] & {
context: undefined;
durationMs: number;
costInCents: number;
baseCostInCents: number;
};
};
const ORG_FRESH_TTL = 60000 * 60 * 24; // 1 day
const ORG_STALE_TTL = 60000 * 60 * 24 * 2; // 2 days
const PROJECT_FRESH_TTL = 60000 * 60 * 24; // 1 day
const PROJECT_STALE_TTL = 60000 * 60 * 24 * 2; // 2 days
const TASK_FRESH_TTL = 60000 * 60 * 24; // 1 day
const TASK_STALE_TTL = 60000 * 60 * 24 * 2; // 2 days
const MACHINE_PRESET_FRESH_TTL = 60000 * 60 * 24; // 1 day
const MACHINE_PRESET_STALE_TTL = 60000 * 60 * 24 * 2; // 2 days
const DEPLOYMENT_FRESH_TTL = 60000 * 60 * 24; // 1 day
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
export class RunAttemptSystem {
private readonly $: SystemResources;
private readonly executionSnapshotSystem: ExecutionSnapshotSystem;
private readonly batchSystem: BatchSystem;
private readonly waitpointSystem: WaitpointSystem;
private readonly delayedRunSystem: DelayedRunSystem;
private readonly cache: UnkeyCache<{
tasks: BackwardsCompatibleTaskRunExecution["task"];
machinePresets: MachinePreset;
deployments: TaskRunExecutionDeployment;
queues: TaskRunExecutionQueue;
projects: TaskRunExecutionProject;
orgs: TaskRunExecutionOrganization;
}>;
constructor(private readonly options: RunAttemptSystemOptions) {
this.$ = options.resources;
this.executionSnapshotSystem = options.executionSnapshotSystem;
this.batchSystem = options.batchSystem;
this.waitpointSystem = options.waitpointSystem;
this.delayedRunSystem = options.delayedRunSystem;
const ctx = new DefaultStatefulContext();
const memory = createLRUMemoryStore(5000);
const redisCacheStore = new RedisCacheStore({
name: "run-attempt-system",
connection: {
...options.redisOptions,
keyPrefix: "engine:run-attempt-system:cache:",
},
useModernCacheKeyBuilder: true,
});
this.cache = createCache({
orgs: new Namespace<TaskRunExecutionOrganization>(ctx, {
stores: [memory, redisCacheStore],
fresh: ORG_FRESH_TTL,
stale: ORG_STALE_TTL,
}),
projects: new Namespace<TaskRunExecutionProject>(ctx, {
stores: [memory, redisCacheStore],
fresh: PROJECT_FRESH_TTL,
stale: PROJECT_STALE_TTL,
}),
tasks: new Namespace<BackwardsCompatibleTaskRunExecution["task"]>(ctx, {
stores: [memory, redisCacheStore],
fresh: TASK_FRESH_TTL,
stale: TASK_STALE_TTL,
}),
machinePresets: new Namespace<MachinePreset>(ctx, {
stores: [memory, redisCacheStore],
fresh: MACHINE_PRESET_FRESH_TTL,
stale: MACHINE_PRESET_STALE_TTL,
}),
deployments: new Namespace<TaskRunExecutionDeployment>(ctx, {
stores: [memory, redisCacheStore],
fresh: DEPLOYMENT_FRESH_TTL,
stale: DEPLOYMENT_STALE_TTL,
}),
queues: new Namespace<TaskRunExecutionQueue>(ctx, {
stores: [memory, redisCacheStore],
fresh: QUEUE_FRESH_TTL,
stale: QUEUE_STALE_TTL,
}),
});
}
public async resolveTaskRunContext(runId: string): Promise<TaskRunContext> {
const run = await this.$.readOnlyPrisma.taskRun.findFirst({
where: {
id: runId,
},
select: {
id: true,
createdAt: true,
updatedAt: true,
executedAt: true,
baseCostInCents: true,
projectId: true,
organizationId: true,
friendlyId: true,
lockedById: true,
lockedQueueId: true,
queue: true,
attemptNumber: true,
status: true,
ttl: true,
machinePreset: true,
runTags: true,
isTest: true,
replayedFromTaskRunFriendlyId: true,
idempotencyKey: true,
idempotencyKeyOptions: true,
startedAt: true,
maxAttempts: true,
taskVersion: true,
maxDurationInSeconds: true,
usageDurationMs: true,
costInCents: true,
traceContext: true,
priorityMs: true,
taskIdentifier: true,
runtimeEnvironment: {
select: {
id: true,
slug: true,
type: true,
branchName: true,
git: true,
organizationId: true,
},
},
parentTaskRunId: true,
rootTaskRunId: true,
batchId: true,
workerQueue: true,
},
});
if (!run) {
throw new ServiceValidationError("Task run not found", 404);
}
const [task, queue, organization, project, machinePreset, deployment] = await Promise.all([
run.lockedById
? this.#resolveTaskRunExecutionTask(run.lockedById)
: Promise.resolve({
id: run.taskIdentifier,
filePath: "unknown",
}),
this.#resolveTaskRunExecutionQueue({
lockedQueueId: run.lockedQueueId ?? undefined,
queueName: run.queue,
runtimeEnvironmentId: run.runtimeEnvironment.id,
}),
this.#resolveTaskRunExecutionOrganization(run.runtimeEnvironment.organizationId),
this.#resolveTaskRunExecutionProjectByRuntimeEnvironmentId(run.runtimeEnvironment.id),
run.lockedById
? this.#resolveTaskRunExecutionMachinePreset(run.lockedById, run.machinePreset)
: Promise.resolve(
getMachinePreset({
defaultMachine: this.options.machines.defaultMachine,
machines: this.options.machines.machines,
config: undefined,
run,
})
),
run.lockedById
? this.#resolveTaskRunExecutionDeployment(run.lockedById)
: Promise.resolve(undefined),
]);
return {
run: {
id: run.friendlyId,
tags: run.runTags,
isTest: run.isTest,
isReplay: !!run.replayedFromTaskRunFriendlyId,
createdAt: run.createdAt,
startedAt: run.startedAt ?? run.createdAt,
idempotencyKey: getUserProvidedIdempotencyKey(run) ?? undefined,
idempotencyKeyScope: extractIdempotencyKeyScope(run),
maxAttempts: run.maxAttempts ?? undefined,
version: run.taskVersion ?? "unknown",
maxDuration: run.maxDurationInSeconds ?? undefined,
priority: run.priorityMs === 0 ? undefined : run.priorityMs / 1_000,
parentTaskRunId: run.parentTaskRunId ? RunId.toFriendlyId(run.parentTaskRunId) : undefined,
rootTaskRunId: run.rootTaskRunId ? RunId.toFriendlyId(run.rootTaskRunId) : undefined,
region: run.runtimeEnvironment.type !== "DEVELOPMENT" ? run.workerQueue : undefined,
},
attempt: {
number: run.attemptNumber ?? 1,
startedAt: run.startedAt ?? new Date(),
},
task,
queue,
organization,
project,
machine: machinePreset,
deployment,
environment: {
id: run.runtimeEnvironment.id,
slug: run.runtimeEnvironment.slug,
type: run.runtimeEnvironment.type,
branchName: run.runtimeEnvironment.branchName ?? undefined,
git: safeParseGitMeta(run.runtimeEnvironment.git),
},
batch: run.batchId ? { id: BatchId.toFriendlyId(run.batchId) } : undefined,
};
}
public async startRunAttempt({
runId,
snapshotId,
workerId,
runnerId,
isWarmStart,
tx,
}: {
runId: string;
snapshotId: string;
workerId?: string;
runnerId?: string;
isWarmStart?: boolean;
tx?: PrismaClientOrTransaction;
}): Promise<StartRunAttemptResult> {
const prisma = tx ?? this.$.prisma;
return startSpan(
this.$.tracer,
"startRunAttempt",
async (span) => {
return this.$.runLock.lock("startRunAttempt", [runId], async () => {
const latestSnapshot = await getLatestExecutionSnapshot(prisma, runId);
if (latestSnapshot.id !== snapshotId) {
//if there is a big delay between the snapshot and the attempt, the snapshot might have changed
//we just want to log because elsewhere it should have been put back into a state where it can be attempted
this.$.logger.warn(
"RunEngine.createRunAttempt(): snapshot has changed since the attempt was created, ignoring.",
{
snapshotId,
latestSnapshotId: latestSnapshot.id,
}
);
throw new ServiceValidationError("Snapshot changed inside startRunAttempt", 409, {
snapshotId,
latestSnapshotId: latestSnapshot.id,
});
}
const taskRun = await this.$.readOnlyPrisma.taskRun.findFirst({
where: {
id: runId,
},
select: {
id: true,
friendlyId: true,
attemptNumber: true,
projectId: true,
runtimeEnvironmentId: true,
status: true,
lockedById: true,
ttl: true,
},
});
this.$.logger.debug("Creating a task run attempt", { taskRun });
if (!taskRun) {
throw new ServiceValidationError("Task run not found", 404);
}
span.setAttribute("projectId", taskRun.projectId);
span.setAttribute("environmentId", taskRun.runtimeEnvironmentId);
span.setAttribute("taskRunId", taskRun.id);
span.setAttribute("taskRunFriendlyId", taskRun.friendlyId);
if (isFinishedOrPendingFinished(latestSnapshot.executionStatus)) {
throw new ServiceValidationError("Task run is already finished", 400);
}
if (!taskRun.lockedById) {
throw new ServiceValidationError("Task run is not locked", 400);
}
//increment the attempt number (start at 1)
const nextAttemptNumber = (taskRun.attemptNumber ?? 0) + 1;
if (nextAttemptNumber > MAX_TASK_RUN_ATTEMPTS) {
await this.attemptFailed({
runId: taskRun.id,
snapshotId,
completion: {
ok: false,
id: taskRun.id,
error: {
type: "INTERNAL_ERROR",
code: "TASK_RUN_CRASHED",
message: "Max attempts reached.",
},
},
tx: prisma,
});
throw new ServiceValidationError("Max attempts reached", 400);
}
const result = await $transaction(
prisma,
async (tx) => {
const run = await tx.taskRun.update({
where: {
id: taskRun.id,
},
data: {
status: "EXECUTING",
attemptNumber: nextAttemptNumber,
executedAt: taskRun.attemptNumber === null ? new Date() : undefined,
isWarmStart: isWarmStart ?? false,
},
select: {
id: true,
createdAt: true,
updatedAt: true,
executedAt: true,
baseCostInCents: true,
projectId: true,
organizationId: true,
friendlyId: true,
lockedById: true,
lockedQueueId: true,
queue: true,
attemptNumber: true,
status: true,
ttl: true,
metadata: true,
metadataType: true,
machinePreset: true,
payload: true,
payloadType: true,
runTags: true,
isTest: true,
replayedFromTaskRunFriendlyId: true,
idempotencyKey: true,
idempotencyKeyOptions: true,
startedAt: true,
maxAttempts: true,
taskVersion: true,
maxDurationInSeconds: true,
usageDurationMs: true,
costInCents: true,
traceContext: true,
priorityMs: true,
batchId: true,
realtimeStreamsVersion: true,
runtimeEnvironment: {
select: {
id: true,
slug: true,
type: true,
branchName: true,
git: true,
organizationId: true,
},
},
parentTaskRunId: true,
rootTaskRunId: true,
workerQueue: true,
taskEventStore: true,
},
});
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(tx, {
run,
snapshot: {
executionStatus: "EXECUTING",
description: `Attempt created, starting execution${
isWarmStart ? " (warm start)" : ""
}`,
},
previousSnapshotId: latestSnapshot.id,
environmentId: latestSnapshot.environmentId,
environmentType: latestSnapshot.environmentType,
projectId: latestSnapshot.projectId,
organizationId: latestSnapshot.organizationId,
batchId: latestSnapshot.batchId ?? undefined,
completedWaitpoints: latestSnapshot.completedWaitpoints,
workerId,
runnerId,
});
if (taskRun.ttl) {
//don't expire the run, it's going to execute
await this.$.worker.ack(`expireRun:${taskRun.id}`);
}
return { updatedRun: run, snapshot: newSnapshot };
},
(error) => {
this.$.logger.error("RunEngine.createRunAttempt(): prisma.$transaction error", {
code: error.code,
meta: error.meta,
stack: error.stack,
message: error.message,
name: error.name,
});
throw new ServiceValidationError(
"Failed to update task run and execution snapshot",
500
);
}
);
if (!result) {
this.$.logger.error("RunEngine.createRunAttempt(): failed to create task run attempt", {
runId: taskRun.id,
nextAttemptNumber,
});
throw new ServiceValidationError("Failed to create task run attempt", 500);
}
const { updatedRun, snapshot } = result;
this.$.eventBus.emit("runAttemptStarted", {
time: new Date(),
run: {
id: updatedRun.id,
status: updatedRun.status,
createdAt: updatedRun.createdAt,
updatedAt: updatedRun.updatedAt,
attemptNumber: nextAttemptNumber,
baseCostInCents: updatedRun.baseCostInCents,
executedAt: updatedRun.executedAt ?? undefined,
},
organization: {
id: updatedRun.runtimeEnvironment.organizationId,
},
project: {
id: updatedRun.projectId,
},
environment: {
id: updatedRun.runtimeEnvironment.id,
},
});
const environmentGit = safeParseGitMeta(updatedRun.runtimeEnvironment.git);
const [metadata, task, queue, organization, project, machinePreset, deployment] =
await Promise.all([
parsePacket({
data: updatedRun.metadata ?? undefined,
dataType: updatedRun.metadataType,
}),
this.#resolveTaskRunExecutionTask(taskRun.lockedById),
this.#resolveTaskRunExecutionQueue({
lockedQueueId: updatedRun.lockedQueueId ?? undefined,
queueName: updatedRun.queue,
runtimeEnvironmentId: updatedRun.runtimeEnvironment.id,
}),
this.#resolveTaskRunExecutionOrganization(
updatedRun.runtimeEnvironment.organizationId
),
this.#resolveTaskRunExecutionProjectByRuntimeEnvironmentId(
updatedRun.runtimeEnvironment.id
),
this.#resolveTaskRunExecutionMachinePreset(
taskRun.lockedById,
updatedRun.machinePreset
),
this.#resolveTaskRunExecutionDeployment(taskRun.lockedById),
]);
const execution: BackwardsCompatibleTaskRunExecution = {
attempt: {
number: nextAttemptNumber,
startedAt: latestSnapshot.updatedAt,
/** @deprecated */
id: "deprecated",
/** @deprecated */
backgroundWorkerId: "deprecated",
/** @deprecated */
backgroundWorkerTaskId: "deprecated",
/** @deprecated */
status: "deprecated",
},
run: {
id: updatedRun.friendlyId,
payload: updatedRun.payload,
payloadType: updatedRun.payloadType,
createdAt: updatedRun.createdAt,
tags: updatedRun.runTags,
isTest: updatedRun.isTest,
isReplay: !!updatedRun.replayedFromTaskRunFriendlyId,
idempotencyKey: getUserProvidedIdempotencyKey(updatedRun) ?? undefined,
idempotencyKeyScope: extractIdempotencyKeyScope(updatedRun),
startedAt: updatedRun.startedAt ?? updatedRun.createdAt,
maxAttempts: updatedRun.maxAttempts ?? undefined,
version: updatedRun.taskVersion ?? "unknown",
metadata,
maxDuration: updatedRun.maxDurationInSeconds ?? undefined,
/** @deprecated */
context: undefined,
/** @deprecated */
durationMs: updatedRun.usageDurationMs,
/** @deprecated */
costInCents: updatedRun.costInCents,
/** @deprecated */
baseCostInCents: updatedRun.baseCostInCents,
traceContext: updatedRun.traceContext as Record<string, string | undefined>,
priority: updatedRun.priorityMs === 0 ? undefined : updatedRun.priorityMs / 1_000,
parentTaskRunId: updatedRun.parentTaskRunId
? RunId.toFriendlyId(updatedRun.parentTaskRunId)
: undefined,
rootTaskRunId: updatedRun.rootTaskRunId
? RunId.toFriendlyId(updatedRun.rootTaskRunId)
: undefined,
region:
updatedRun.runtimeEnvironment.type !== "DEVELOPMENT"
? updatedRun.workerQueue
: undefined,
realtimeStreamsVersion: updatedRun.realtimeStreamsVersion ?? undefined,
},
task,
queue,
environment: {
id: updatedRun.runtimeEnvironment.id,
slug: updatedRun.runtimeEnvironment.slug,
type: updatedRun.runtimeEnvironment.type,
branchName: updatedRun.runtimeEnvironment.branchName ?? undefined,
git: environmentGit,
},
organization,
project,
machine: machinePreset,
deployment,
batch: updatedRun.batchId
? {
id: BatchId.toFriendlyId(updatedRun.batchId),
}
: undefined,
};
return { run: updatedRun, snapshot, execution };
});
},
{
attributes: { runId, snapshotId },
}
);
}
public async completeRunAttempt({
runId,
snapshotId,
completion,
workerId,
runnerId,
}: {
runId: string;
snapshotId: string;
completion: TaskRunExecutionResult;
workerId?: string;
runnerId?: string;
}): Promise<CompleteRunAttemptResult> {
await this.#notifyMetadataUpdated(runId, completion);
switch (completion.ok) {
case true: {
return this.attemptSucceeded({
runId,
snapshotId,
completion,
tx: this.$.prisma,
workerId,
runnerId,
});
}
case false: {
return this.attemptFailed({
runId,
snapshotId,
completion,
tx: this.$.prisma,
workerId,
runnerId,
});
}
}
}
public async attemptSucceeded({
runId,
snapshotId,
completion,
tx,
workerId,
runnerId,
}: {
runId: string;
snapshotId: string;
completion: TaskRunSuccessfulExecutionResult;
tx: PrismaClientOrTransaction;
workerId?: string;
runnerId?: string;
}): Promise<CompleteRunAttemptResult> {
const prisma = tx ?? this.$.prisma;
return startSpan(
this.$.tracer,
"#completeRunAttemptSuccess",
async (span) => {
return this.$.runLock.lock("attemptSucceeded", [runId], async () => {
const latestSnapshot = await getLatestExecutionSnapshot(prisma, runId);
if (latestSnapshot.id !== snapshotId) {
throw new ServiceValidationError("Snapshot ID doesn't match the latest snapshot", 400);
}
if (latestSnapshot.executionStatus === "FINISHED") {
throw new ServiceValidationError("Run is already finished", 400);
}
span.setAttribute("completionStatus", completion.ok);
span.setAttribute("runId", runId);
const completedAt = new Date();
// Read current usage values to calculate new totals (safe under runLock)
const currentRun = await this.$.readOnlyPrisma.taskRun.findFirst({
where: { id: runId },
select: {
usageDurationMs: true,
costInCents: true,
machinePreset: true,
},
});
if (!currentRun) {
throw new ServiceValidationError("Run not found", 404);
}
// Calculate new usage totals
const updatedUsage = this.#calculateUpdatedUsage({
runId,
currentUsageDurationMs: currentRun.usageDurationMs,
currentCostInCents: currentRun.costInCents,
attemptDurationMs: completion.usage?.durationMs ?? 0,
machinePresetName: currentRun.machinePreset,
environmentType: latestSnapshot.environmentType,
});
const run = await prisma.taskRun.update({
where: { id: runId },
data: {
status: "COMPLETED_SUCCESSFULLY",
completedAt,
output: completion.output,
outputType: completion.outputType,
usageDurationMs: updatedUsage.usageDurationMs,
costInCents: updatedUsage.costInCents,
executionSnapshots: {
create: {
executionStatus: "FINISHED",
description: "Task completed successfully",
runStatus: "COMPLETED_SUCCESSFULLY",
attemptNumber: latestSnapshot.attemptNumber,
environmentId: latestSnapshot.environmentId,
environmentType: latestSnapshot.environmentType,
projectId: latestSnapshot.projectId,
organizationId: latestSnapshot.organizationId,
workerId,
runnerId,
},
},
},
select: {
id: true,
friendlyId: true,
status: true,
attemptNumber: true,
spanId: true,
updatedAt: true,
associatedWaitpoint: {
select: {
id: true,
},
},
project: {
select: {
organizationId: true,
},
},
batchId: true,
createdAt: true,
completedAt: true,
taskEventStore: true,
parentTaskRunId: true,
usageDurationMs: true,
costInCents: true,
runtimeEnvironmentId: true,
projectId: true,
},
});
const newSnapshot = await getLatestExecutionSnapshot(prisma, runId);
await this.$.runQueue.acknowledgeMessage(run.project.organizationId, runId);
// We need to manually emit this as we created the final snapshot as part of the task run update
this.$.eventBus.emit("executionSnapshotCreated", {
time: newSnapshot.createdAt,
run: {
id: newSnapshot.runId,
},
snapshot: {
...newSnapshot,
completedWaitpointIds: newSnapshot.completedWaitpoints.map((wp) => wp.id),
},
});
// Complete the waitpoint if it exists (runs without waiting parents have no waitpoint)
if (run.associatedWaitpoint) {
await this.waitpointSystem.completeWaitpoint({
id: run.associatedWaitpoint.id,
output: completion.output
? { value: completion.output, type: completion.outputType, isError: false }
: undefined,
});
}
this.$.eventBus.emit("runSucceeded", {
time: completedAt,
run: {
id: runId,
status: run.status,
spanId: run.spanId,
output: completion.output,
outputType: completion.outputType,
createdAt: run.createdAt,
completedAt: run.completedAt,
taskEventStore: run.taskEventStore,
usageDurationMs: run.usageDurationMs,
costInCents: run.costInCents,
updatedAt: run.updatedAt,
attemptNumber: run.attemptNumber ?? 1,
},
organization: {
id: run.project.organizationId,
},
project: {
id: run.projectId,
},
environment: {
id: run.runtimeEnvironmentId,
},
});
await this.#finalizeRun(run);
return {
attemptStatus: "RUN_FINISHED",
snapshot: newSnapshot,
run,
};
});
},
{
attributes: { runId, snapshotId },
}
);
}
public async attemptFailed({
runId,
snapshotId,
workerId,
runnerId,
completion,
forceRequeue,
tx,
}: {
runId: string;
snapshotId: string;
workerId?: string;
runnerId?: string;
completion: TaskRunFailedExecutionResult;
forceRequeue?: boolean;
tx: PrismaClientOrTransaction;
}): Promise<CompleteRunAttemptResult> {
const prisma = this.$.prisma;
return startSpan(
this.$.tracer,
"completeRunAttemptFailure",
async (span) => {
return this.$.runLock.lock("attemptFailed", [runId], async () => {
const latestSnapshot = await getLatestExecutionSnapshot(prisma, runId);
if (latestSnapshot.id !== snapshotId) {
throw new ServiceValidationError("Snapshot ID doesn't match the latest snapshot", 400);
}
if (latestSnapshot.executionStatus === "FINISHED") {
throw new ServiceValidationError("Run is already finished", 400);
}
span.setAttribute("completionStatus", completion.ok);
//remove waitpoints blocking the run
const deletedCount = await this.waitpointSystem.clearBlockingWaitpoints({ runId, tx });
if (deletedCount > 0) {
this.$.logger.debug("Cleared blocking waitpoints", { runId, deletedCount });
}
const failedAt = new Date();
const retryResult = await retryOutcomeFromCompletion(this.$.readOnlyPrisma, {
runId,
error: completion.error,
retryUsingQueue: forceRequeue ?? false,
retrySettings: completion.retry,
attemptNumber: latestSnapshot.attemptNumber,
});
// Force requeue means it was crashed so the attempt span needs to be closed
if (forceRequeue) {
const minimalRun = await this.$.readOnlyPrisma.taskRun.findFirst({
where: {
id: runId,
},
select: {
status: true,
spanId: true,
maxAttempts: true,
runtimeEnvironment: {
select: {
organizationId: true,
},
},
taskEventStore: true,
createdAt: true,
completedAt: true,
updatedAt: true,
},
});
if (!minimalRun) {
throw new ServiceValidationError("Run not found", 404);
}
this.$.eventBus.emit("runAttemptFailed", {
time: failedAt,
run: {
id: runId,
status: minimalRun.status,
spanId: minimalRun.spanId,
error: completion.error,
attemptNumber: latestSnapshot.attemptNumber ?? 0,
createdAt: minimalRun.createdAt,
completedAt: minimalRun.completedAt,
taskEventStore: minimalRun.taskEventStore,
updatedAt: minimalRun.updatedAt,
},
});
}
switch (retryResult.outcome) {
case "cancel_run": {
const result = await this.cancelRun({
runId,
completedAt: failedAt,
reason: retryResult.reason,
finalizeRun: true,
attemptDurationMs: completion.usage?.durationMs,
tx: prisma,
});
return {
attemptStatus:
result.snapshot.executionStatus === "PENDING_CANCEL"
? "RUN_PENDING_CANCEL"
: "RUN_FINISHED",
...result,
};
}
case "fail_run": {
return await this.#permanentlyFailRun({
runId,
latestSnapshot,
failedAt,
error: retryResult.sanitizedError,
workerId,
runnerId,
attemptDurationMs: completion.usage?.durationMs,
});
}
case "retry": {
const retryAt = new Date(retryResult.settings.timestamp);
// Calculate new usage totals using the current machine's rate
// (retryResult includes current usage values from the read in retryOutcomeFromCompletion)
const updatedUsage = this.#calculateUpdatedUsage({
runId,
currentUsageDurationMs: retryResult.usageDurationMs,
currentCostInCents: retryResult.costInCents,
attemptDurationMs: completion.usage?.durationMs ?? 0,
machinePresetName: retryResult.machinePreset,
environmentType: latestSnapshot.environmentType,
});
const run = await prisma.taskRun.update({
where: {
id: runId,
},
data: {
machinePreset: retryResult.machine,
usageDurationMs: updatedUsage.usageDurationMs,
costInCents: updatedUsage.costInCents,
},
include: {
runtimeEnvironment: {
include: {
project: true,
organization: true,
orgMember: true,
},
},
},
});
const nextAttemptNumber =
latestSnapshot.attemptNumber === null ? 1 : latestSnapshot.attemptNumber + 1;
if (retryResult.wasOOMError) {
this.$.eventBus.emit("runAttemptFailed", {
time: failedAt,
run: {
id: runId,
status: run.status,
spanId: run.spanId,
error: completion.error,
attemptNumber: latestSnapshot.attemptNumber ?? 0,
createdAt: run.createdAt,
completedAt: run.completedAt,
taskEventStore: run.taskEventStore,
updatedAt: run.updatedAt,
},
});
}
this.$.eventBus.emit("runRetryScheduled", {
time: failedAt,
run: {
id: run.id,
status: run.status,
friendlyId: run.friendlyId,
attemptNumber: nextAttemptNumber,
queue: run.queue,
taskIdentifier: run.taskIdentifier,
traceContext: run.traceContext as Record<string, string | undefined>,
baseCostInCents: run.baseCostInCents,
spanId: run.spanId,
nextMachineAfterOOM: retryResult.machine,
updatedAt: run.updatedAt,
error: completion.error,
createdAt: run.createdAt,
taskEventStore: run.taskEventStore,
},
organization: {
id: run.runtimeEnvironment.organizationId,
},
environment: run.runtimeEnvironment,
retryAt,
});
//if it's a long delay and we support checkpointing, put it back in the queue
if (
forceRequeue ||
retryResult.method === "queue" ||
(this.options.retryWarmStartThresholdMs !== undefined &&
retryResult.settings.delay >= this.options.retryWarmStartThresholdMs)
) {
//we nack the message, requeuing it for later
const nackResult = await this.tryNackAndRequeue({
run,
environment: run.runtimeEnvironment,
orgId: run.runtimeEnvironment.organizationId,
projectId: run.runtimeEnvironment.project.id,
timestamp: retryAt.getTime(),
error: {
type: "INTERNAL_ERROR",
code: "TASK_RUN_DEQUEUED_MAX_RETRIES",
message: `We tried to dequeue the run the maximum number of times but it wouldn't start executing`,
},
tx: prisma,
});
if (!nackResult.wasRequeued) {
return {
attemptStatus: "RUN_FINISHED",
...nackResult,
};
} else {
return { attemptStatus: "RETRY_QUEUED", ...nackResult };
}
}
//it will continue running because the retry delay is short
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(
prisma,
{
run,
snapshot: {
executionStatus: "EXECUTING",
description: "Attempt failed with a short delay, starting a new attempt",
},
previousSnapshotId: latestSnapshot.id,
environmentId: latestSnapshot.environmentId,
environmentType: latestSnapshot.environmentType,
projectId: latestSnapshot.projectId,
organizationId: latestSnapshot.organizationId,
workerId,
runnerId,
}
);
//the worker can fetch the latest snapshot and should create a new attempt
await sendNotificationToWorker({
runId,
snapshot: newSnapshot,
eventBus: this.$.eventBus,
});
return {
attemptStatus: "RETRY_IMMEDIATELY",
...executionResultFromSnapshot(newSnapshot),
};
}
}
});
},
{
attributes: { runId, snapshotId },
}
);
}
public async systemFailure({
runId,
error,
tx,
}: {
runId: string;
error: TaskRunInternalError;
tx?: PrismaClientOrTransaction;
}): Promise<CompleteRunAttemptResult> {
const prisma = tx ?? this.$.prisma;
return startSpan(
this.$.tracer,
"systemFailure",
async (span) => {
const latestSnapshot = await getLatestExecutionSnapshot(prisma, runId);
//already finished
if (latestSnapshot.executionStatus === "FINISHED") {
//todo check run is in the correct state
return {
attemptStatus: "RUN_FINISHED",
snapshot: latestSnapshot,
run: {
id: runId,
friendlyId: latestSnapshot.runFriendlyId,
status: latestSnapshot.runStatus,
attemptNumber: latestSnapshot.attemptNumber,
},
};
}
const result = await this.attemptFailed({
runId,
snapshotId: latestSnapshot.id,
completion: {
ok: false,
id: runId,
error,
},
tx: prisma,
});
return result;
},
{
attributes: {
runId,
},
}
);
}
public async tryNackAndRequeue({
run,
environment,
orgId,
projectId,
timestamp,
error,
workerId,
runnerId,
checkpointId,
completedWaitpoints,
batchId,
tx,
}: {
run: TaskRun;
environment: {
id: string;
type: RuntimeEnvironmentType;
};
orgId: string;
projectId: string;
timestamp?: number;
error: TaskRunInternalError;
workerId?: string;
runnerId?: string;
checkpointId?: string;
tx?: PrismaClientOrTransaction;
completedWaitpoints?: {
id: string;
index?: number;
}[];
batchId?: string;
}): Promise<{ wasRequeued: boolean } & ExecutionResult> {
const prisma = tx ?? this.$.prisma;
return await this.$.runLock.lock("tryNackAndRequeue", [run.id], async () => {
//we nack the message, this allows another worker to pick up the run
const gotRequeued = await this.$.runQueue.nackMessage({
orgId,
messageId: run.id,
retryAt: timestamp,
});
if (!gotRequeued) {
const result = await this.systemFailure({
runId: run.id,
error,
tx: prisma,
});
return { wasRequeued: false, ...result };
}
const requeuedRun = await prisma.taskRun.update({
where: {
id: run.id,
},
data: {
status: "PENDING",
},
select: {
id: true,
status: true,
attemptNumber: true,
},
});
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(prisma, {
run: requeuedRun,
snapshot: {
executionStatus: "QUEUED",
description: "Requeued the run after a failure",
},
environmentId: environment.id,
environmentType: environment.type,
projectId: projectId,
organizationId: orgId,
workerId,
runnerId,
checkpointId,
completedWaitpoints,
batchId,
});
return {
wasRequeued: true,
snapshot: {
id: newSnapshot.id,
friendlyId: newSnapshot.friendlyId,
executionStatus: newSnapshot.executionStatus,
description: newSnapshot.description,
createdAt: newSnapshot.createdAt,
},
run: {
id: newSnapshot.runId,
friendlyId: newSnapshot.runFriendlyId,
status: newSnapshot.runStatus,
attemptNumber: newSnapshot.attemptNumber,
},
};
});
}
/**
Call this to cancel a run.
If the run is in-progress it will change it's state to PENDING_CANCEL and notify the worker.
If the run is not in-progress it will finish it.
You can pass `finalizeRun` in if you know it's no longer running, e.g. the worker has messaged to say it's done.
You can optionally pass `attemptDurationMs` if you have completion data with usage info.
*/
async cancelRun({
runId,
workerId,
runnerId,
completedAt,
reason,
finalizeRun,
bulkActionId,
attemptDurationMs,
tx,
}: {
runId: string;
workerId?: string;
runnerId?: string;
completedAt?: Date;
reason?: string;
finalizeRun?: boolean;
bulkActionId?: string;
attemptDurationMs?: number;
tx?: PrismaClientOrTransaction;
}): Promise<ExecutionResult & { alreadyFinished: boolean }> {
const prisma = tx ?? this.$.prisma;
reason = reason ?? "Canceled by user";
return startSpan(this.$.tracer, "cancelRun", async (span) => {
return this.$.runLock.lock("cancelRun", [runId], async () => {
const latestSnapshot = await getLatestExecutionSnapshot(prisma, runId);
//already finished, do nothing
if (latestSnapshot.executionStatus === "FINISHED") {
if (bulkActionId) {
await prisma.taskRun.update({
where: { id: runId },
data: {
bulkActionGroupIds: {
push: bulkActionId,
},
},
});
}
return {
alreadyFinished: true,
...executionResultFromSnapshot(latestSnapshot),
};
}
//is pending cancellation and we're not finalizing, alert the worker again
if (latestSnapshot.executionStatus === "PENDING_CANCEL" && !finalizeRun) {
await sendNotificationToWorker({
runId,
snapshot: latestSnapshot,
eventBus: this.$.eventBus,
});
return {
alreadyFinished: false,
...executionResultFromSnapshot(latestSnapshot),
};
}
//set the run to cancelled immediately
const error: TaskRunError = {
type: "STRING_ERROR",
raw: reason,
};
// Calculate updated usage if we have attempt duration data
let usageUpdate: { usageDurationMs: number; costInCents: number } | undefined;
if (attemptDurationMs !== undefined) {
const currentRun = await this.$.readOnlyPrisma.taskRun.findFirst({
where: { id: runId },
select: {
usageDurationMs: true,
costInCents: true,
machinePreset: true,
},
});
if (!currentRun) {
throw new ServiceValidationError("Run not found", 404);
}
usageUpdate = this.#calculateUpdatedUsage({
runId,
currentUsageDurationMs: currentRun.usageDurationMs,
currentCostInCents: currentRun.costInCents,
attemptDurationMs,
machinePresetName: currentRun.machinePreset,
environmentType: latestSnapshot.environmentType,
});
}
const run = await prisma.taskRun.update({
where: { id: runId },
data: {
status: "CANCELED",
completedAt: finalizeRun ? completedAt ?? new Date() : completedAt,
error,
bulkActionGroupIds: bulkActionId
? {
push: bulkActionId,
}
: undefined,
...(usageUpdate && {
usageDurationMs: usageUpdate.usageDurationMs,
costInCents: usageUpdate.costInCents,
}),
},
select: {
id: true,
friendlyId: true,
status: true,
attemptNumber: true,
spanId: true,
batchId: true,
createdAt: true,
completedAt: true,
taskEventStore: true,
parentTaskRunId: true,
delayUntil: true,
updatedAt: true,
runtimeEnvironment: {
select: {
organizationId: true,
},
},
associatedWaitpoint: {
select: {
id: true,
},
},
childRuns: {
select: {
id: true,
},
},
},
});
//if the run is delayed and hasn't started yet, we need to prevent it being added to the queue in future
if (isInitialState(latestSnapshot.executionStatus) && run.delayUntil) {
await this.delayedRunSystem.preventDelayedRunFromBeingEnqueued({ runId });
}
//remove it from the queue and release concurrency
await this.$.runQueue.acknowledgeMessage(run.runtimeEnvironment.organizationId, runId, {
removeFromWorkerQueue: true,
});
//if executing, we need to message the worker to cancel the run and put it into `PENDING_CANCEL` status
//unless finalizeRun is true (worker is known to be dead), in which case skip straight to FINISHED
if (
isExecuting(latestSnapshot.executionStatus) ||
isPendingExecuting(latestSnapshot.executionStatus)
) {
if (!finalizeRun) {
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(prisma, {
run,
snapshot: {
executionStatus: "PENDING_CANCEL",
description: "Run was cancelled",
},
previousSnapshotId: latestSnapshot.id,
environmentId: latestSnapshot.environmentId,
environmentType: latestSnapshot.environmentType,
projectId: latestSnapshot.projectId,
organizationId: latestSnapshot.organizationId,
workerId,
runnerId,
});
//the worker needs to be notified so it can kill the run and complete the attempt
await sendNotificationToWorker({
runId,
snapshot: newSnapshot,
eventBus: this.$.eventBus,
});
return {
alreadyFinished: false,
...executionResultFromSnapshot(newSnapshot),
};
}
// finalizeRun is true — fall through to finish the run immediately
}
//not executing, so we will actually finish the run
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(prisma, {
run,
snapshot: {
executionStatus: "FINISHED",
description: "Run was cancelled, not finished",
},
previousSnapshotId: latestSnapshot.id,
environmentId: latestSnapshot.environmentId,
environmentType: latestSnapshot.environmentType,
projectId: latestSnapshot.projectId,
organizationId: latestSnapshot.organizationId,
workerId,
runnerId,
});
// Complete the waitpoint if it exists (runs without waiting parents have no waitpoint)
if (run.associatedWaitpoint) {
await this.waitpointSystem.completeWaitpoint({
id: run.associatedWaitpoint.id,
output: { value: JSON.stringify(error), isError: true },
});
}
await this.#finalizeRun(run);
this.$.eventBus.emit("runCancelled", {
time: new Date(),
run: {
id: run.id,
status: run.status,
friendlyId: run.friendlyId,
spanId: run.spanId,
taskEventStore: run.taskEventStore,
createdAt: run.createdAt,
completedAt: run.completedAt,
error,
updatedAt: run.updatedAt,
attemptNumber: run.attemptNumber ?? 1,
},
organization: {
id: latestSnapshot.organizationId,
},
project: {
id: latestSnapshot.projectId,
},
environment: {
id: latestSnapshot.environmentId,
},
});
//schedule the cancellation of all the child runs
//it will call this function for each child,
//which will recursively cancel all children if they need to be
if (run.childRuns.length > 0) {
for (const childRun of run.childRuns) {
await this.$.worker.enqueue({
id: `cancelRun:${childRun.id}`,
job: "cancelRun",
payload: { runId: childRun.id, completedAt: run.completedAt ?? new Date(), reason },
});
}
}
return {
alreadyFinished: false,
...executionResultFromSnapshot(newSnapshot),
};
});
});
}
async #permanentlyFailRun({
runId,
latestSnapshot,
failedAt,
error,
workerId,
runnerId,
attemptDurationMs,
}: {
runId: string;
latestSnapshot: EnhancedExecutionSnapshot;
failedAt: Date;
error: TaskRunError;
workerId?: string;
runnerId?: string;
attemptDurationMs?: number;
}): Promise<CompleteRunAttemptResult> {
const prisma = this.$.prisma;
return startSpan(this.$.tracer, "permanentlyFailRun", async (span) => {
const status = runStatusFromError(error, latestSnapshot.environmentType);
const truncatedError = this.#truncateTaskRunError(error);
// Read current usage values to calculate new totals
const currentRun = await this.$.readOnlyPrisma.taskRun.findFirst({
where: { id: runId },
select: {
usageDurationMs: true,
costInCents: true,
machinePreset: true,
},
});
if (!currentRun) {
throw new ServiceValidationError("Run not found", 404);
}
// Calculate new usage totals
const updatedUsage = this.#calculateUpdatedUsage({
runId,
currentUsageDurationMs: currentRun.usageDurationMs,
currentCostInCents: currentRun.costInCents,
attemptDurationMs: attemptDurationMs ?? 0,
machinePresetName: currentRun.machinePreset,
environmentType: latestSnapshot.environmentType,
});
//run permanently failed
const run = await prisma.taskRun.update({
where: {
id: runId,
},
data: {
status,
completedAt: failedAt,
error: truncatedError,
usageDurationMs: updatedUsage.usageDurationMs,
costInCents: updatedUsage.costInCents,
},
select: {
id: true,
friendlyId: true,
status: true,
attemptNumber: true,
spanId: true,
batchId: true,
parentTaskRunId: true,
updatedAt: true,
usageDurationMs: true,
costInCents: true,
associatedWaitpoint: {
select: {
id: true,
},
},
runtimeEnvironment: {
select: {
id: true,
type: true,
organizationId: true,
project: {
select: {
id: true,
organizationId: true,
},
},
},
},
taskEventStore: true,
createdAt: true,
completedAt: true,
},
});
const newSnapshot = await this.executionSnapshotSystem.createExecutionSnapshot(prisma, {
run,
snapshot: {
executionStatus: "FINISHED",
description: "Run failed",
},
previousSnapshotId: latestSnapshot.id,
environmentId: run.runtimeEnvironment.id,
environmentType: run.runtimeEnvironment.type,
projectId: run.runtimeEnvironment.project.id,
organizationId: run.runtimeEnvironment.project.organizationId,
workerId,
runnerId,
});
await this.$.runQueue.acknowledgeMessage(run.runtimeEnvironment.organizationId, runId, {
removeFromWorkerQueue: true,
});
// Complete the waitpoint if it exists (runs without waiting parents have no waitpoint)
if (run.associatedWaitpoint) {
await this.waitpointSystem.completeWaitpoint({
id: run.associatedWaitpoint.id,
output: { value: JSON.stringify(truncatedError), isError: true },
});
}
this.$.eventBus.emit("runFailed", {
time: failedAt,
run: {
id: runId,
status: run.status,
spanId: run.spanId,
error,
taskEventStore: run.taskEventStore,
createdAt: run.createdAt,
completedAt: run.completedAt,
updatedAt: run.updatedAt,
attemptNumber: run.attemptNumber ?? 1,
usageDurationMs: run.usageDurationMs,
costInCents: run.costInCents,
},
organization: {
id: run.runtimeEnvironment.project.organizationId,
},
project: {
id: run.runtimeEnvironment.project.id,
},
environment: {
id: run.runtimeEnvironment.id,
},
});
await this.#finalizeRun(run);
return {
attemptStatus: "RUN_FINISHED",
snapshot: newSnapshot,
run,
};
});
}
/*
* Whether the run succeeds, fails, is cancelled… we need to run these operations
*/
async #finalizeRun({ id, batchId }: { id: string; batchId: string | null }) {
if (batchId) {
await this.batchSystem.scheduleCompleteBatch({ batchId });
}
//cancel the heartbeats
await this.$.worker.ack(`heartbeatSnapshot.${id}`);
}
async #resolveTaskRunExecutionTask(
backgroundWorkerTaskId: string
): Promise<BackwardsCompatibleTaskRunExecution["task"]> {
const result = await this.cache.tasks.swr(backgroundWorkerTaskId, async () => {
const task = await this.$.readOnlyPrisma.backgroundWorkerTask.findFirstOrThrow({
where: {
id: backgroundWorkerTaskId,
},
select: {
id: true,
slug: true,
filePath: true,
exportName: true,
},
});
return {
id: task.slug,
filePath: task.filePath,
exportName: task.exportName ?? undefined,
};
});
if (result.err) {
throw result.err;
}
if (!result.val) {
throw new ServiceValidationError(
`Could not resolve task execution data for task ${backgroundWorkerTaskId}`
);
}
return result.val;
}
async #resolveTaskRunExecutionOrganization(
organizationId: string
): Promise<TaskRunExecutionOrganization> {
const result = await this.cache.orgs.swr(organizationId, async () => {
const organization = await this.$.readOnlyPrisma.organization.findFirstOrThrow({
where: { id: organizationId },
select: {
id: true,
title: true,
slug: true,
},
});
return {
id: organization.id,
name: organization.title,
slug: organization.slug,
};
});
if (result.err) {
throw result.err;
}
if (!result.val) {
throw new ServiceValidationError(
`Could not resolve organization data for organization ${organizationId}`
);
}
return result.val;
}
async #resolveTaskRunExecutionProjectByRuntimeEnvironmentId(
runtimeEnvironmentId: string
): Promise<TaskRunExecutionProject> {
const result = await this.cache.projects.swr(runtimeEnvironmentId, async () => {
const { project } = await this.$.readOnlyPrisma.runtimeEnvironment.findFirstOrThrow({
where: { id: runtimeEnvironmentId },
select: {
id: true,
project: {
select: {
id: true,
name: true,
slug: true,
externalRef: true,
},
},
},
});
return {
id: project.id,
name: project.name,
slug: project.slug,
ref: project.externalRef,
};
});
if (result.err) {
throw result.err;
}
if (!result.val) {
throw new ServiceValidationError(
`Could not resolve project data for project ${runtimeEnvironmentId}`
);
}
return result.val;
}
async #resolveTaskRunExecutionMachinePreset(
backgroundWorkerTaskId: string,
runMachinePreset: string | null
): Promise<MachinePreset> {
if (runMachinePreset) {
return machinePresetFromName(
this.options.machines.machines,
runMachinePreset as MachinePresetName
);
}
const result = await this.cache.machinePresets.swr(backgroundWorkerTaskId, async () => {
const { machineConfig } = await this.$.readOnlyPrisma.backgroundWorkerTask.findFirstOrThrow({
where: {
id: backgroundWorkerTaskId,
},
select: {
machineConfig: true,
},
});
return getMachinePreset({
machines: this.options.machines.machines,
defaultMachine: this.options.machines.defaultMachine,
config: machineConfig,
run: { machinePreset: null },
});
});
if (result.err) {
throw result.err;
}
if (!result.val) {
throw new ServiceValidationError(
`Could not resolve machine preset for task ${backgroundWorkerTaskId}`
);
}
return result.val;
}
async #resolveTaskRunExecutionQueue(params: {
lockedQueueId?: string;
queueName: string;
runtimeEnvironmentId: string;
}): Promise<TaskRunExecutionQueue> {
// Cache key should be based on queue identity, not run ID
// Using lockedQueueId if available, otherwise environment + queue name
const cacheKey = params.lockedQueueId ?? `${params.runtimeEnvironmentId}:${params.queueName}`;
const result = await this.cache.queues.swr(cacheKey, async () => {
const queue = params.lockedQueueId
? await this.$.readOnlyPrisma.taskQueue.findFirst({
where: {
id: params.lockedQueueId,
},
select: {
id: true,
friendlyId: true,
name: true,
},
})
: await this.$.readOnlyPrisma.taskQueue.findFirst({
where: {
runtimeEnvironmentId: params.runtimeEnvironmentId,
name: params.queueName,
},
select: {
id: true,
friendlyId: true,
name: true,
},
});
if (!queue) {
// Return synthetic queue so run/span view still loads (e.g. createFailedTaskRun with fallback queue)
return {
id: params.queueName,
name: params.queueName,
};
}
return {
id: queue.friendlyId,
name: queue.name,
};
});
if (result.err) {
throw result.err;
}
if (!result.val) {
throw new ServiceValidationError(
`Could not resolve queue data for queue ${params.queueName}`,
404
);
}
return result.val;
}
async #resolveTaskRunExecutionDeployment(
backgroundWorkerTaskId: string
): Promise<TaskRunExecutionDeployment | undefined> {
const result = await this.cache.deployments.swr(backgroundWorkerTaskId, async () => {
const { worker } = await this.$.readOnlyPrisma.backgroundWorkerTask.findFirstOrThrow({
where: { id: backgroundWorkerTaskId },
select: {
worker: {
select: {
deployment: true,
},
},
},
});
if (!worker.deployment) {
return undefined;
}
return {
id: worker.deployment.friendlyId,
shortCode: worker.deployment.shortCode,
version: worker.deployment.version,
runtime: worker.deployment.runtime ?? "unknown",
runtimeVersion: worker.deployment.runtimeVersion ?? "unknown",
git: safeParseGitMeta(worker.deployment.git),
};
});
if (result.err) {
throw result.err;
}
return result.val;
}
async #notifyMetadataUpdated(runId: string, completion: TaskRunExecutionResult) {
if (completion.metadata) {
this.$.eventBus.emit("runMetadataUpdated", {
time: new Date(),
run: {
id: runId,
metadata: completion.metadata,
},
});
return;
}
if (completion.flushedMetadata) {
const [, packet] = await tryCatch(parsePacket(completion.flushedMetadata));
if (!packet) {
return;
}
const metadata = FlushedRunMetadata.safeParse(packet);
if (!metadata.success) {
// Customer's metadata operations don't match the schema (typically
// non-JSON values in `operations[].value`). System ignores it.
this.$.logger.warn(
"RunEngine.completeRunAttempt(): failed to validate flushed metadata",
{
runId,
flushedMetadata: completion.flushedMetadata,
error: metadata.error,
}
);
return;
}
this.$.eventBus.emit("runMetadataUpdated", {
time: new Date(),
run: {
id: runId,
metadata: metadata.data,
},
});
}
}
#truncateTaskRunError(error: TaskRunError): TaskRunError {
if (error.type !== "BUILT_IN_ERROR") {
return error;
}
return {
type: "BUILT_IN_ERROR",
name: truncateString(error.name, 1024),
message: truncateString(error.message, 1024 * 16), // 16kb
stackTrace: truncateString(error.stackTrace, 1024 * 16), // 16kb
};
}
// PostgreSQL int4 max value (~24.85 days in milliseconds)
static readonly MAX_INT4 = 2_147_483_647;
/**
* Calculates the updated usage values for a run by adding the attempt's usage to the current totals.
* This should be called under the runLock to ensure safe read-modify-write.
* Cost is only calculated for non-dev environments.
*/
#calculateUpdatedUsage({
runId,
currentUsageDurationMs,
currentCostInCents,
attemptDurationMs,
machinePresetName,
environmentType,
}: {
runId: string;
currentUsageDurationMs: number;
currentCostInCents: number;
attemptDurationMs: number;
machinePresetName: string | null;
environmentType: RuntimeEnvironmentType;
}): { usageDurationMs: number; costInCents: number } {
let usageDurationMs = currentUsageDurationMs + attemptDurationMs;
// Overflow protection: cap at PostgreSQL int4 max value
if (usageDurationMs > RunAttemptSystem.MAX_INT4) {
this.$.logger.error("usageDurationMs overflow detected, capping at max int4 value", {
runId,
currentUsageDurationMs,
attemptDurationMs,
calculatedTotal: usageDurationMs,
cappedAt: RunAttemptSystem.MAX_INT4,
});
usageDurationMs = RunAttemptSystem.MAX_INT4;
}
// Only calculate cost for non-dev environments
let costInCents = currentCostInCents;
if (environmentType !== "DEVELOPMENT") {
const machinePreset = machinePresetName
? machinePresetFromName(
this.options.machines.machines,
machinePresetName as MachinePresetName
)
: machinePresetFromName(
this.options.machines.machines,
this.options.machines.defaultMachine
);
costInCents = currentCostInCents + attemptDurationMs * machinePreset.centsPerMs;
}
return {
usageDurationMs,
costInCents,
};
}
}
export function safeParseGitMeta(git: unknown): GitMeta | undefined {
const parsed = GitMeta.safeParse(git);
if (parsed.success) {
return parsed.data;
}
return undefined;
}
function truncateString(str: string | undefined, maxLength: number): string {
if (!str) {
return "";
}
return str.slice(0, maxLength);
}