Frozen run fixes (#1286)

* When resuming a batch, only do marqs operations once

* Made TaskRunDependency clearer in the Prisma schema

* New ResumeDependentParentsService service, use it from checkpoints

* WIP on making resuming more robust

* Turn the declarative schedules off because they make debugging other runs painful

* Resuming batches when there’s an attempt is working

* If there’s no attempt then create one

* Added a log if there are no span events to complete

* If Graphile addJob doesn’t return a row, log and return undefined. No throw

* Pass prisma into the ResumeDependentParentsService

* Removed the todos

* Pass Prisma through to the checkpoint service

* Fix for not checking the batch item correctly

* Fix for when a log flush times out and the process is checkpointed

* Fix for when a log flush times out and the process is checkpointed

* Another test run that does batches with failed subtasks

* Don’t call ResumeTaskRunDependenciesService anymore (we have a new service)

* Only resume if the run is in a final state

* If an attempt doesn’t exist, fix for creating queue with sanitized name

* If DEV then don’t resume using marqs/batches. The CLI manages it

* We don’t need to check the run status again, it’s in the main function now

* Added TaskRunAttempt taskRunId index

* Only allow calling ResumeDependentParentsService with a run ID

* Put the flushing back to what it was
This commit is contained in:
Matt Aitken
2024-09-09 16:03:36 +01:00
committed by GitHub
parent e5f0aafb14
commit 392453e8a6
22 changed files with 624 additions and 231 deletions
+6
View File
@@ -0,0 +1,6 @@
---
"trigger.dev": patch
"@trigger.dev/core": patch
---
Fix for when a log flush times out and the process is checkpointed
+3 -5
View File
@@ -318,7 +318,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
identifier: K, identifier: K,
payload: z.infer<TMessageCatalog[K]>, payload: z.infer<TMessageCatalog[K]>,
options?: ZodWorkerEnqueueOptions options?: ZodWorkerEnqueueOptions
): Promise<GraphileJob> { ): Promise<GraphileJob | undefined> {
const task = this.#tasks[identifier]; const task = this.#tasks[identifier];
const optionsWithoutTx = removeUndefinedKeys(omit(options ?? {}, ["tx"])); const optionsWithoutTx = removeUndefinedKeys(omit(options ?? {}, ["tx"]));
@@ -439,11 +439,9 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
identifier, identifier,
payload, payload,
spec, spec,
error: JSON.stringify(rows.error),
}); });
return { job: undefined, durationInMs: Math.floor(durationInMs) };
throw new Error(
`Failed to add job to queue, zod parsing error: ${JSON.stringify(rows.error)}`
);
} }
const job = rows.data[0]; const job = rows.data[0];
@@ -796,6 +796,12 @@ function RunTimelineLine({ title, state }: RunTimelineLineProps) {
function RunError({ error }: { error: TaskRunError }) { function RunError({ error }: { error: TaskRunError }) {
switch (error.type) { switch (error.type) {
case "STRING_ERROR": case "STRING_ERROR":
return (
<div className="flex flex-col gap-2 rounded-sm border border-rose-500/50 px-3 pb-3 pt-2">
<Header3 className="text-rose-500">Error</Header3>
<Callout variant="error">{error.raw}</Callout>
</div>
);
case "CUSTOM_ERROR": { case "CUSTOM_ERROR": {
return ( return (
<div className="flex flex-col gap-2 rounded-sm border border-rose-500/50 px-3 pb-3 pt-2"> <div className="flex flex-col gap-2 rounded-sm border border-rose-500/50 px-3 pb-3 pt-2">
@@ -134,7 +134,7 @@ export class DeliverScheduledEventService {
id, id,
}, },
data: { data: {
workerJobId: workerJob.id, workerJobId: workerJob?.id,
nextEventTimestamp: runAt, nextEventTimestamp: runAt,
}, },
}); });
@@ -226,6 +226,7 @@ export class EventRepository {
const events = await this.queryIncompleteEvents({ spanId }); const events = await this.queryIncompleteEvents({ spanId });
if (events.length === 0) { if (events.length === 0) {
logger.warn("No incomplete events found for spanId", { spanId, options });
return; return;
} }
+2 -18
View File
@@ -42,26 +42,10 @@ export class FailedTaskRunService extends BaseService {
id: taskRun.id, id: taskRun.id,
status: "SYSTEM_FAILURE", status: "SYSTEM_FAILURE",
completedAt: new Date(), completedAt: new Date(),
attemptStatus: "FAILED",
error: sanitizeError(completion.error),
}); });
// Get the final attempt and add the error to it, if it's not already set
const finalAttempt = await this._prisma.taskRunAttempt.findFirst({
where: {
taskRunId: taskRun.id,
},
orderBy: { id: "desc" },
});
if (finalAttempt && !finalAttempt.error) {
// Haven't set the status because the attempt might still be running
await this._prisma.taskRunAttempt.update({
where: { id: finalAttempt.id },
data: {
error: sanitizeError(completion.error),
},
});
}
// Now we need to "complete" the task run event/span // Now we need to "complete" the task run event/span
await eventRepository.completeEvent(taskRun.spanId, { await eventRepository.completeEvent(taskRun.spanId, {
endTime: new Date(), endTime: new Date(),
@@ -5,7 +5,6 @@ import { eventRepository } from "../eventRepository.server";
import { isCancellableRunStatus } from "../taskStatus"; import { isCancellableRunStatus } from "../taskStatus";
import { BaseService } from "./baseService.server"; import { BaseService } from "./baseService.server";
import { FinalizeTaskRunService } from "./finalizeTaskRun.server"; import { FinalizeTaskRunService } from "./finalizeTaskRun.server";
import { ResumeTaskRunDependenciesService } from "./resumeTaskRunDependencies.server";
export class CancelAttemptService extends BaseService { export class CancelAttemptService extends BaseService {
public async call( public async call(
@@ -61,13 +60,15 @@ export class CancelAttemptService extends BaseService {
}, },
}); });
const isCancellable = isCancellableRunStatus(taskRunAttempt.taskRun.status);
const finalizeService = new FinalizeTaskRunService(tx); const finalizeService = new FinalizeTaskRunService(tx);
await finalizeService.call({ await finalizeService.call({
id: taskRunId, id: taskRunId,
status: isCancellableRunStatus(taskRunAttempt.taskRun.status) ? "INTERRUPTED" : undefined, status: isCancellable ? "INTERRUPTED" : undefined,
completedAt: isCancellableRunStatus(taskRunAttempt.taskRun.status) completedAt: isCancellable ? cancelledAt : undefined,
? cancelledAt attemptStatus: isCancellable ? "CANCELED" : undefined,
: undefined, error: isCancellable ? { type: "STRING_ERROR", raw: reason } : undefined,
}); });
}); });
@@ -84,10 +85,6 @@ export class CancelAttemptService extends BaseService {
return eventRepository.cancelEvent(event, cancelledAt, reason); return eventRepository.cancelEvent(event, cancelledAt, reason);
}) })
); );
if (environment?.type !== "DEVELOPMENT") {
await ResumeTaskRunDependenciesService.enqueue(taskRunAttempt.id, this._prisma);
}
}); });
} }
} }
@@ -76,6 +76,11 @@ export class CancelTaskRunService extends BaseService {
runtimeEnvironment: true, runtimeEnvironment: true,
lockedToVersion: true, lockedToVersion: true,
}, },
attemptStatus: "CANCELED",
error: {
type: "STRING_ERROR",
raw: opts.reason,
},
}); });
const inProgressEvents = await eventRepository.queryIncompleteEvents({ const inProgressEvents = await eventRepository.queryIncompleteEvents({
@@ -17,7 +17,6 @@ import { createExceptionPropertiesFromError, eventRepository } from "../eventRep
import { marqs } from "~/v3/marqs/index.server"; import { marqs } from "~/v3/marqs/index.server";
import { BaseService } from "./baseService.server"; import { BaseService } from "./baseService.server";
import { CancelAttemptService } from "./cancelAttempt.server"; import { CancelAttemptService } from "./cancelAttempt.server";
import { ResumeTaskRunDependenciesService } from "./resumeTaskRunDependencies.server";
import { MAX_TASK_RUN_ATTEMPTS } from "~/consts"; import { MAX_TASK_RUN_ATTEMPTS } from "~/consts";
import { CreateCheckpointService } from "./createCheckpoint.server"; import { CreateCheckpointService } from "./createCheckpoint.server";
import { TaskRun } from "@trigger.dev/database"; import { TaskRun } from "@trigger.dev/database";
@@ -76,6 +75,12 @@ export class CompleteAttemptService extends BaseService {
id: run.id, id: run.id,
status: "SYSTEM_FAILURE", status: "SYSTEM_FAILURE",
completedAt: new Date(), completedAt: new Date(),
attemptStatus: "FAILED",
error: {
type: "INTERNAL_ERROR",
code: "TASK_EXECUTION_FAILED",
message: "Tried to complete attempt but it doesn't exist",
},
}); });
// No attempt, so there's no message to ACK // No attempt, so there's no message to ACK
@@ -149,10 +154,6 @@ export class CompleteAttemptService extends BaseService {
}, },
}); });
if (!env || env.type !== "DEVELOPMENT") {
await ResumeTaskRunDependenciesService.enqueue(taskRunAttempt.id, this._prisma);
}
return "COMPLETED"; return "COMPLETED";
} }
@@ -355,10 +356,6 @@ export class CompleteAttemptService extends BaseService {
}); });
} }
if (!env || env.type !== "DEVELOPMENT") {
await ResumeTaskRunDependenciesService.enqueue(taskRunAttempt.id, this._prisma);
}
return "COMPLETED"; return "COMPLETED";
} }
} }
@@ -4,7 +4,6 @@ import { marqs } from "~/v3/marqs/index.server";
import { BaseService } from "./baseService.server"; import { BaseService } from "./baseService.server";
import { logger } from "~/services/logger.server"; import { logger } from "~/services/logger.server";
import { AuthenticatedEnvironment } from "~/services/apiAuth.server"; import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { ResumeTaskRunDependenciesService } from "./resumeTaskRunDependencies.server";
import { CRASHABLE_ATTEMPT_STATUSES, isCrashableRunStatus } from "../taskStatus"; import { CRASHABLE_ATTEMPT_STATUSES, isCrashableRunStatus } from "../taskStatus";
import { sanitizeError } from "@trigger.dev/core/v3"; import { sanitizeError } from "@trigger.dev/core/v3";
import { FinalizeTaskRunService } from "./finalizeTaskRun.server"; import { FinalizeTaskRunService } from "./finalizeTaskRun.server";
@@ -69,6 +68,13 @@ export class CrashTaskRunService extends BaseService {
}, },
}, },
}, },
attemptStatus: "FAILED",
error: {
type: "INTERNAL_ERROR",
code: "TASK_RUN_CRASHED",
message: opts.reason,
stackTrace: opts.logs,
},
}); });
const inProgressEvents = await eventRepository.queryIncompleteEvents( const inProgressEvents = await eventRepository.queryIncompleteEvents(
@@ -146,12 +152,6 @@ export class CrashTaskRunService extends BaseService {
}), }),
}, },
}); });
if (environment.type === "DEVELOPMENT") {
return;
}
await ResumeTaskRunDependenciesService.enqueue(attempt.id, this._prisma);
}); });
} }
} }
@@ -4,16 +4,11 @@ import type { Checkpoint, CheckpointRestoreEvent } from "@trigger.dev/database";
import { logger } from "~/services/logger.server"; import { logger } from "~/services/logger.server";
import { marqs } from "~/v3/marqs/index.server"; import { marqs } from "~/v3/marqs/index.server";
import { generateFriendlyId } from "../friendlyIdentifiers"; import { generateFriendlyId } from "../friendlyIdentifiers";
import { import { isFreezableAttemptStatus, isFreezableRunStatus } from "../taskStatus";
isFinalAttemptStatus,
isFinalRunStatus,
isFreezableAttemptStatus,
isFreezableRunStatus,
} from "../taskStatus";
import { BaseService } from "./baseService.server"; import { BaseService } from "./baseService.server";
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server"; import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
import { ResumeBatchRunService } from "./resumeBatchRun.server"; import { ResumeBatchRunService } from "./resumeBatchRun.server";
import { ResumeTaskDependencyService } from "./resumeTaskDependency.server"; import { ResumeDependentParentsService } from "./resumeDependentParents.server";
export class CreateCheckpointService extends BaseService { export class CreateCheckpointService extends BaseService {
public async call( public async call(
@@ -177,127 +172,15 @@ export class CreateCheckpointService extends BaseService {
}); });
await marqs?.cancelHeartbeat(attempt.taskRunId); await marqs?.cancelHeartbeat(attempt.taskRunId);
const dependency = await this._prisma.taskRunDependency.findFirst({ const resumeService = new ResumeDependentParentsService(this._prisma);
select: { const result = await resumeService.call({ id: attempt.taskRunId });
id: true,
taskRunId: true,
},
where: {
taskRun: {
friendlyId: reason.friendlyId,
},
},
});
logger.log("CreateCheckpointService: Created checkpoint WAIT_FOR_TASK", { if (result.success) {
checkpointId: checkpoint.id, logger.log("CreateCheckpointService: Resumed dependent parents", result);
runFriendlyId: reason.friendlyId, } else {
dependencyId: dependency?.id, logger.error("CreateCheckpointService: Failed to resume dependent parents", result);
});
if (!dependency) {
logger.error("CreateCheckpointService: Dependency not found", {
friendlyId: reason.friendlyId,
});
return {
success: true,
checkpoint,
event: checkpointEvent,
keepRunAlive: false,
};
} }
const childRun = await this._prisma.taskRun.findFirst({
select: {
id: true,
status: true,
},
where: {
id: dependency.taskRunId,
},
});
if (!childRun) {
logger.error("CreateCheckpointService: Dependency child run not found", {
taskRunId: dependency.taskRunId,
runFriendlyId: reason.friendlyId,
dependencyId: dependency.id,
});
return {
success: true,
checkpoint,
event: checkpointEvent,
keepRunAlive: false,
};
}
const isFinished = isFinalRunStatus(childRun.status);
if (!isFinished) {
logger.debug("CreateCheckpointService: Dependency child run not finished", {
taskRunId: dependency.taskRunId,
runFriendlyId: reason.friendlyId,
dependencyId: dependency.id,
childRunStatus: childRun.status,
childRunId: childRun.id,
});
return {
success: true,
checkpoint,
event: checkpointEvent,
keepRunAlive: false,
};
}
const lastAttempt = await this._prisma.taskRunAttempt.findFirst({
select: {
id: true,
status: true,
},
where: {
taskRunId: dependency.taskRunId,
},
orderBy: {
createdAt: "desc",
},
});
if (!lastAttempt) {
logger.debug("CreateCheckpointService: Dependency child attempt not found", {
taskRunId: dependency.taskRunId,
runFriendlyId: reason.friendlyId,
dependencyId: dependency?.id,
});
return {
success: true,
checkpoint,
event: checkpointEvent,
keepRunAlive: false,
};
}
if (!isFinalAttemptStatus(lastAttempt.status)) {
logger.debug("CreateCheckpointService: Dependency child attempt not final", {
taskRunId: dependency.taskRunId,
runFriendlyId: reason.friendlyId,
dependencyId: dependency.id,
lastAttemptId: lastAttempt.id,
lastAttemptStatus: lastAttempt.status,
});
return {
success: true,
checkpoint,
event: checkpointEvent,
keepRunAlive: false,
};
}
//resume the dependent task
await ResumeTaskDependencyService.enqueue(dependency.id, lastAttempt.id, this._prisma);
return { return {
success: true, success: true,
checkpoint, checkpoint,
@@ -45,6 +45,11 @@ export class ExpireEnqueuedRunService extends BaseService {
status: "EXPIRED", status: "EXPIRED",
expiredAt: new Date(), expiredAt: new Date(),
completedAt: new Date(), completedAt: new Date(),
attemptStatus: "FAILED",
error: {
type: "STRING_ERROR",
raw: `Run expired because the TTL (${run.ttl}) was reached`,
},
}); });
await eventRepository.completeEvent(run.spanId, { await eventRepository.completeEvent(run.spanId, {
@@ -1,16 +1,24 @@
import { sanitizeError, TaskRunError } from "@trigger.dev/core/v3";
import { type Prisma, type TaskRun } from "@trigger.dev/database"; import { type Prisma, type TaskRun } from "@trigger.dev/database";
import { logger } from "~/services/logger.server"; import { logger } from "~/services/logger.server";
import { marqs } from "~/v3/marqs/index.server"; import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { BaseService } from "./baseService.server"; import {
import { isFailedRunStatus, type FINAL_RUN_STATUSES } from "../taskStatus"; isFailedRunStatus,
import { PerformTaskAttemptAlertsService } from "./alerts/performTaskAttemptAlerts.server"; type FINAL_ATTEMPT_STATUSES,
type FINAL_RUN_STATUSES,
} from "../taskStatus";
import { PerformTaskRunAlertsService } from "./alerts/performTaskRunAlerts.server"; import { PerformTaskRunAlertsService } from "./alerts/performTaskRunAlerts.server";
import { BaseService } from "./baseService.server";
import { ResumeDependentParentsService } from "./resumeDependentParents.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
type BaseInput = { type BaseInput = {
id: string; id: string;
status?: FINAL_RUN_STATUSES; status?: FINAL_RUN_STATUSES;
expiredAt?: Date; expiredAt?: Date;
completedAt?: Date; completedAt?: Date;
attemptStatus?: FINAL_ATTEMPT_STATUSES;
error?: TaskRunError;
}; };
type InputWithInclude<T extends Prisma.TaskRunInclude> = BaseInput & { type InputWithInclude<T extends Prisma.TaskRunInclude> = BaseInput & {
@@ -32,6 +40,8 @@ export class FinalizeTaskRunService extends BaseService {
expiredAt, expiredAt,
completedAt, completedAt,
include, include,
attemptStatus,
error,
}: T extends Prisma.TaskRunInclude ? InputWithInclude<T> : InputWithoutInclude): Promise< }: T extends Prisma.TaskRunInclude ? InputWithInclude<T> : InputWithoutInclude): Promise<
Output<T> Output<T>
> { > {
@@ -56,6 +66,20 @@ export class FinalizeTaskRunService extends BaseService {
...(include ? { include } : {}), ...(include ? { include } : {}),
}); });
if (attemptStatus || error) {
await this.finalizeAttempt({ attemptStatus, error, run });
}
//resume any dependencies
const resumeService = new ResumeDependentParentsService(this._prisma);
const result = await resumeService.call({ id: run.id });
if (result.success) {
logger.log("FinalizeTaskRunService: Resumed dependent parents", { result });
} else {
logger.error("FinalizeTaskRunService: Failed to resume dependent parents", { result });
}
//enqueue alert //enqueue alert
if (isFailedRunStatus(run.status)) { if (isFailedRunStatus(run.status)) {
await PerformTaskRunAlertsService.enqueue(run.id, this._prisma); await PerformTaskRunAlertsService.enqueue(run.id, this._prisma);
@@ -63,4 +87,85 @@ export class FinalizeTaskRunService extends BaseService {
return run as Output<T>; return run as Output<T>;
} }
async finalizeAttempt({
attemptStatus,
error,
run,
}: {
attemptStatus?: FINAL_ATTEMPT_STATUSES;
error?: TaskRunError;
run: TaskRun;
}) {
if (attemptStatus || error) {
const latestAttempt = await this._prisma.taskRunAttempt.findFirst({
where: { taskRunId: run.id },
orderBy: { id: "desc" },
take: 1,
});
if (latestAttempt) {
logger.debug("Finalizing run attempt", {
id: latestAttempt.id,
status: attemptStatus,
error,
});
await this._prisma.taskRunAttempt.update({
where: { id: latestAttempt.id },
data: { status: attemptStatus, error: error ? sanitizeError(error) : undefined },
});
} else {
logger.debug("Finalizing run no attempt found", {
runId: run.id,
attemptStatus,
error,
});
const workerTask = await this._prisma.backgroundWorkerTask.findFirst({
select: {
id: true,
workerId: true,
runtimeEnvironmentId: true,
},
where: {
id: run.lockedById!,
},
});
if (!workerTask) {
logger.error("FinalizeTaskRunService: No worker task found", { runId: run.id });
return;
}
const queue = await this._prisma.taskQueue.findUnique({
where: {
runtimeEnvironmentId_name: {
runtimeEnvironmentId: workerTask.runtimeEnvironmentId,
name: sanitizeQueueName(run.queue),
},
},
});
if (!queue) {
logger.error("FinalizeTaskRunService: No queue found", { runId: run.id });
return;
}
await this._prisma.taskRunAttempt.create({
data: {
number: 1,
friendlyId: generateFriendlyId("attempt"),
taskRunId: run.id,
backgroundWorkerId: workerTask?.workerId,
backgroundWorkerTaskId: workerTask?.id,
queueId: queue.id,
runtimeEnvironmentId: workerTask.runtimeEnvironmentId,
status: attemptStatus,
error: error ? sanitizeError(error) : undefined,
},
});
}
}
}
} }
@@ -4,22 +4,13 @@ import { marqs } from "~/v3/marqs/index.server";
import { BaseService } from "./baseService.server"; import { BaseService } from "./baseService.server";
import { logger } from "~/services/logger.server"; import { logger } from "~/services/logger.server";
const finishedBatchRunStatuses = ["COMPLETED", "FAILED", "CANCELED"];
export class ResumeBatchRunService extends BaseService { export class ResumeBatchRunService extends BaseService {
public async call(batchRunId: string) { public async call(batchRunId: string) {
const batchRun = await this._prisma.batchTaskRun.findFirst({ const batchRun = await this._prisma.batchTaskRun.findFirst({
where: { where: {
id: batchRunId, id: batchRunId,
dependentTaskAttemptId: {
not: null,
},
status: "PENDING",
items: {
every: {
taskRunAttemptId: {
not: null,
},
},
},
}, },
include: { include: {
dependentTaskAttempt: { dependentTaskAttempt: {
@@ -38,6 +29,26 @@ export class ResumeBatchRunService extends BaseService {
}); });
if (!batchRun || !batchRun.dependentTaskAttempt) { if (!batchRun || !batchRun.dependentTaskAttempt) {
logger.error(
"ResumeBatchRunService: Batch run doesn't exist or doesn't have a dependent attempt",
{
batchRun,
}
);
return;
}
if (batchRun.status === "COMPLETED") {
logger.debug("ResumeBatchRunService: Batch run is already completed", {
batchRun: batchRun,
});
return;
}
if (batchRun.items.some((item) => !finishedBatchRunStatuses.includes(item.status))) {
logger.debug("ResumeBatchRunService: All items aren't yet completed", {
batchRun: batchRun,
});
return; return;
} }
@@ -61,34 +72,45 @@ export class ResumeBatchRunService extends BaseService {
const dependentRun = batchRun.dependentTaskAttempt.taskRun; const dependentRun = batchRun.dependentTaskAttempt.taskRun;
if (batchRun.dependentTaskAttempt.status === "PAUSED" && batchRun.checkpointEventId) { if (batchRun.dependentTaskAttempt.status === "PAUSED" && batchRun.checkpointEventId) {
// We need to update the batchRun status so we don't resume it again logger.debug("ResumeBatchRunService: Attempt is paused and has a checkpoint event", {
await this._prisma.batchTaskRun.update({ batchRunId: batchRun.id,
where: { dependentTaskAttempt: batchRun.dependentTaskAttempt,
id: batchRun.id, checkpointEventId: batchRun.checkpointEventId,
},
data: {
status: "COMPLETED",
},
}); });
await marqs?.enqueueMessage( // We need to update the batchRun status so we don't resume it again
environment, const wasUpdated = await this.#setBatchToCompletedOnce(batchRun.id);
dependentRun.queue, if (wasUpdated) {
dependentRun.id, logger.debug("ResumeBatchRunService: Resuming dependent run with checkpoint", {
{ batchRunId: batchRun.id,
type: "RESUME", dependentTaskAttemptId: batchRun.dependentTaskAttempt.id,
completedAttemptIds: [], });
resumableAttemptId: batchRun.dependentTaskAttempt.id, await marqs?.enqueueMessage(
environment,
dependentRun.queue,
dependentRun.id,
{
type: "RESUME",
completedAttemptIds: [],
resumableAttemptId: batchRun.dependentTaskAttempt.id,
checkpointEventId: batchRun.checkpointEventId,
taskIdentifier: batchRun.dependentTaskAttempt.taskRun.taskIdentifier,
projectId: batchRun.dependentTaskAttempt.runtimeEnvironment.projectId,
environmentId: batchRun.dependentTaskAttempt.runtimeEnvironment.id,
environmentType: batchRun.dependentTaskAttempt.runtimeEnvironment.type,
},
dependentRun.concurrencyKey ?? undefined
);
} else {
logger.debug("ResumeBatchRunService: with checkpoint was already completed", {
batchRunId: batchRun.id,
dependentTaskAttempt: batchRun.dependentTaskAttempt,
checkpointEventId: batchRun.checkpointEventId, checkpointEventId: batchRun.checkpointEventId,
taskIdentifier: batchRun.dependentTaskAttempt.taskRun.taskIdentifier, hasCheckpointEvent: !!batchRun.checkpointEventId,
projectId: batchRun.dependentTaskAttempt.runtimeEnvironment.projectId, });
environmentId: batchRun.dependentTaskAttempt.runtimeEnvironment.id, }
environmentType: batchRun.dependentTaskAttempt.runtimeEnvironment.type,
},
dependentRun.concurrencyKey ?? undefined
);
} else { } else {
logger.debug("Batch run resume: Attempt is not paused or there's no checkpoint event", { logger.debug("ResumeBatchRunService: attempt is not paused or there's no checkpoint event", {
batchRunId: batchRun.id, batchRunId: batchRun.id,
dependentTaskAttempt: batchRun.dependentTaskAttempt, dependentTaskAttempt: batchRun.dependentTaskAttempt,
checkpointEventId: batchRun.checkpointEventId, checkpointEventId: batchRun.checkpointEventId,
@@ -98,23 +120,60 @@ export class ResumeBatchRunService extends BaseService {
if (batchRun.dependentTaskAttempt.status === "PAUSED" && !batchRun.checkpointEventId) { if (batchRun.dependentTaskAttempt.status === "PAUSED" && !batchRun.checkpointEventId) {
// In case of race conditions the status can be PAUSED without a checkpoint event // In case of race conditions the status can be PAUSED without a checkpoint event
// When the checkpoint is created, it will continue the run // When the checkpoint is created, it will continue the run
logger.error("Batch run resume: Attempt is paused but there's no checkpoint event", { logger.error("ResumeBatchRunService: attempt is paused but there's no checkpoint event", {
batchRunId: batchRun.id, batchRunId: batchRun.id,
dependentTaskAttemptId: batchRun.dependentTaskAttempt.id, dependentTaskAttemptId: batchRun.dependentTaskAttempt.id,
}); });
return; return;
} }
await marqs?.replaceMessage(dependentRun.id, { // We need to update the batchRun status so we don't resume it again
type: "RESUME", const wasUpdated = await this.#setBatchToCompletedOnce(batchRun.id);
completedAttemptIds: batchRun.items.map((item) => item.taskRunAttemptId).filter(Boolean), if (wasUpdated) {
resumableAttemptId: batchRun.dependentTaskAttempt.id, logger.debug("ResumeBatchRunService: Resuming dependent run without checkpoint", {
checkpointEventId: batchRun.checkpointEventId ?? undefined, batchRunId: batchRun.id,
taskIdentifier: batchRun.dependentTaskAttempt.taskRun.taskIdentifier, dependentTaskAttemptId: batchRun.dependentTaskAttempt.id,
projectId: batchRun.dependentTaskAttempt.runtimeEnvironment.projectId, });
environmentId: batchRun.dependentTaskAttempt.runtimeEnvironment.id, await marqs?.replaceMessage(dependentRun.id, {
environmentType: batchRun.dependentTaskAttempt.runtimeEnvironment.type, type: "RESUME",
}); completedAttemptIds: batchRun.items.map((item) => item.taskRunAttemptId).filter(Boolean),
resumableAttemptId: batchRun.dependentTaskAttempt.id,
checkpointEventId: batchRun.checkpointEventId ?? undefined,
taskIdentifier: batchRun.dependentTaskAttempt.taskRun.taskIdentifier,
projectId: batchRun.dependentTaskAttempt.runtimeEnvironment.projectId,
environmentId: batchRun.dependentTaskAttempt.runtimeEnvironment.id,
environmentType: batchRun.dependentTaskAttempt.runtimeEnvironment.type,
});
} else {
logger.debug("ResumeBatchRunService: without checkpoint was already completed", {
batchRunId: batchRun.id,
dependentTaskAttempt: batchRun.dependentTaskAttempt,
checkpointEventId: batchRun.checkpointEventId,
hasCheckpointEvent: !!batchRun.checkpointEventId,
});
}
}
}
async #setBatchToCompletedOnce(batchRunId: string) {
const result = await this._prisma.batchTaskRun.updateMany({
where: {
id: batchRunId,
status: {
not: "COMPLETED", // Ensure the status is not already "COMPLETED"
},
},
data: {
status: "COMPLETED",
},
});
// Check if any records were updated
if (result.count > 0) {
// The status was changed, so we return true
return true;
} else {
return false;
} }
} }
@@ -0,0 +1,257 @@
import { Prisma } from "@trigger.dev/database";
import { logger } from "~/services/logger.server";
import { isFinalAttemptStatus, isFinalRunStatus } from "../taskStatus";
import { BaseService } from "./baseService.server";
import { ResumeBatchRunService } from "./resumeBatchRun.server";
import { ResumeTaskDependencyService } from "./resumeTaskDependency.server";
import { $transaction } from "~/db.server";
type Output =
| {
success: true;
action:
| "resume-scheduled"
| "batch-resume-scheduled"
| "no-dependencies"
| "not-finished"
| "dev";
}
| {
success: false;
error: string;
};
type Dependency = Prisma.TaskRunDependencyGetPayload<{
include: {
taskRun: true;
dependentAttempt: true;
dependentBatchRun: true;
};
}>;
/** This will resume a dependent (parent) run if there is one and it makes sense. */
export class ResumeDependentParentsService extends BaseService {
public async call({ id }: { id: string }): Promise<Output> {
try {
const dependency = await this._prisma.taskRunDependency.findFirst({
include: {
taskRun: {
include: {
runtimeEnvironment: true,
},
},
dependentAttempt: true,
dependentBatchRun: true,
},
where: {
taskRunId: id,
},
});
logger.log("ResumeDependentParentsService: tried to find dependency", {
runId: id,
dependency: dependency,
});
if (!dependency) {
logger.log("ResumeDependentParentsService: dependency not found", {
runId: id,
});
//no dependency, that's fine most runs won't have one.
return {
success: true,
action: "no-dependencies",
};
}
if (!isFinalRunStatus(dependency.taskRun.status)) {
logger.debug(
"ResumeDependentParentsService: run not finished yet, can't resume parent yet",
{
runId: id,
dependency,
}
);
// the child run isn't finished yet, so we can't resume the parent yet.
return {
success: true,
action: "not-finished",
};
}
if (dependency.taskRun.runtimeEnvironment.type === "DEVELOPMENT") {
logger.debug("ResumeDependentParentsService: runs are resumed on device for DEV runs.", {
runId: id,
dependency,
});
return {
success: true,
action: "dev",
};
}
if (dependency.dependentAttempt) {
return this.#singleRunDependency(dependency);
} else if (dependency.dependentBatchRun) {
return this.#batchRunDependency(dependency);
} else {
logger.error("ResumeDependentParentsService: dependency has no dependencies", {
runId: id,
dependency,
});
return {
success: false,
error: `Dependency has no dependencies (single or batch)`,
};
}
} catch (error) {
return {
success: false,
error: error instanceof Error ? error.message : JSON.stringify(error),
};
}
}
async #singleRunDependency(dependency: Dependency): Promise<Output> {
logger.debug(
`ResumeDependentParentsService.singleRunDependency(): Resuming dependent parent for run`,
{
dependency,
}
);
const lastAttempt = await this._prisma.taskRunAttempt.findFirst({
select: {
id: true,
status: true,
},
where: {
taskRunId: dependency.taskRunId,
},
orderBy: {
id: "desc",
},
});
if (!lastAttempt) {
logger.error(
"ResumeDependentParentsService.singleRunDependency(): dependency child attempt not found",
{
dependency,
}
);
return {
success: false,
error: `Dependency child attempt not found for run ${dependency.taskRunId}`,
};
}
if (!isFinalAttemptStatus(lastAttempt.status)) {
//We still want to continue if this happens because the run is final but log it
logger.error(
"ResumeDependentParentsService.singleRunDependency(): dependency child attempt not final, but the run is.",
{
dependency,
lastAttempt,
}
);
return {
success: false,
error: `Dependency child attempt not final, but the run is`,
};
}
//resume the dependent task
await ResumeTaskDependencyService.enqueue(dependency.id, lastAttempt.id, this._prisma);
return {
success: true,
action: "resume-scheduled",
};
}
async #batchRunDependency(dependency: Dependency): Promise<Output> {
logger.debug(
`ResumeDependentParentsService.batchRunDependency(): Resuming dependent batch for run`,
{
dependency,
}
);
if (!dependency.dependentBatchRun) {
logger.error(
"ResumeDependentParentsService.batchRunDependency(): dependency has no dependent batch",
{
dependency,
}
);
return {
success: false,
error: `Dependency has no dependent batch`,
};
}
const lastAttempt = await this._prisma.taskRunAttempt.findFirst({
select: {
id: true,
status: true,
},
where: {
taskRunId: dependency.taskRunId,
},
orderBy: {
id: "desc",
},
});
if (!lastAttempt) {
logger.error(
"ResumeDependentParentsService.singleRunDependency(): dependency child attempt not found",
{
dependency,
}
);
return {
success: false,
error: `Dependency child attempt not found for run ${dependency.taskRunId}`,
};
}
logger.log(
"ResumeDependentParentsService.batchRunDependency(): Setting the batchTaskRunItem to COMPLETED",
{
dependency,
lastAttempt,
}
);
await $transaction(this._prisma, async (tx) => {
await tx.batchTaskRunItem.update({
where: {
batchTaskRunId_taskRunId: {
batchTaskRunId: dependency.dependentBatchRun!.id,
taskRunId: dependency.taskRunId,
},
},
data: {
status: "COMPLETED",
taskRunAttemptId: lastAttempt.id,
},
});
await ResumeBatchRunService.enqueue(dependency.dependentBatchRun!.id, tx);
});
return {
success: true,
action: "batch-resume-scheduled",
};
}
}
+6 -1
View File
@@ -45,7 +45,12 @@ export const FINAL_RUN_STATUSES = [
export type FINAL_RUN_STATUSES = (typeof FINAL_RUN_STATUSES)[number]; export type FINAL_RUN_STATUSES = (typeof FINAL_RUN_STATUSES)[number];
export const FINAL_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["CANCELED", "COMPLETED", "FAILED"]; export const FINAL_ATTEMPT_STATUSES = [
"CANCELED",
"COMPLETED",
"FAILED",
] satisfies TaskRunAttemptStatus[];
export type FINAL_ATTEMPT_STATUSES = (typeof FINAL_ATTEMPT_STATUSES)[number];
export const FREEZABLE_RUN_STATUSES: TaskRunStatus[] = ["EXECUTING", "RETRYING_AFTER_FAILURE"]; export const FREEZABLE_RUN_STATUSES: TaskRunStatus[] = ["EXECUTING", "RETRYING_AFTER_FAILURE"];
export const FREEZABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["EXECUTING", "FAILED"]; export const FREEZABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["EXECUTING", "FAILED"];
@@ -399,7 +399,7 @@ class FlushingProcess {
private _flushPromise: Promise<void>; private _flushPromise: Promise<void>;
constructor(private readonly doFlush: () => Promise<void>) { constructor(private readonly doFlush: () => Promise<void>) {
this._flushPromise = this.doFlush(); this._flushPromise = this.doFlush().catch(() => {});
} }
waitForCompletion() { waitForCompletion() {
@@ -0,0 +1,2 @@
-- CreateIndex
CREATE INDEX CONCURRENTLY IF NOT EXISTS "TaskRunAttempt_taskRunId_idx" ON "TaskRunAttempt" ("taskRunId");
+3
View File
@@ -1796,9 +1796,11 @@ model TaskRunTag {
@@index([name, id]) @@index([name, id])
} }
/// This is used for triggerAndWait and batchTriggerAndWait. The taskRun is the child task, it points at a parent attempt or a batch
model TaskRunDependency { model TaskRunDependency {
id String @id @default(cuid()) id String @id @default(cuid())
/// The child run
taskRun TaskRun @relation(fields: [taskRunId], references: [id], onDelete: Cascade, onUpdate: Cascade) taskRun TaskRun @relation(fields: [taskRunId], references: [id], onDelete: Cascade, onUpdate: Cascade)
taskRunId String @unique taskRunId String @unique
@@ -1880,6 +1882,7 @@ model TaskRunAttempt {
alerts ProjectAlert[] alerts ProjectAlert[]
@@unique([taskRunId, number]) @@unique([taskRunId, number])
@@index([taskRunId])
} }
enum TaskRunAttemptStatus { enum TaskRunAttemptStatus {
@@ -26,8 +26,35 @@ export const batchParentTask = task({
}, },
}); });
export const batchParentWitFailsTask = task({
id: "batch-parent-with-fails-task",
retry: {
maxAttempts: 1,
},
run: async () => {
const response = await taskThatFails.batchTriggerAndWait([
{ payload: false },
{ payload: true },
{ payload: false },
]);
logger.info("Batch response", { response });
const respone2 = await taskThatFails.batchTriggerAndWait([
{ payload: true },
{ payload: false },
{ payload: true },
]);
logger.info("Batch response2", { respone2 });
},
});
export const batchChildTask = task({ export const batchChildTask = task({
id: "batch-child-task", id: "batch-child-task",
retry: {
maxAttempts: 2,
},
run: async (payload: string, { ctx }) => { run: async (payload: string, { ctx }) => {
logger.info("Processing child task", { payload }); logger.info("Processing child task", { payload });
@@ -36,3 +63,21 @@ export const batchChildTask = task({
return `${payload} - processed`; return `${payload} - processed`;
}, },
}); });
export const taskThatFails = task({
id: "task-that-fails",
retry: {
maxAttempts: 2,
},
run: async (fail: boolean) => {
logger.info(`Will fail ${fail}`);
if (fail) {
throw new Error("Task failed");
}
return {
foo: "bar",
};
},
});
@@ -0,0 +1,35 @@
import { logger, task } from "@trigger.dev/sdk/v3";
type Payload = {};
export const crashparent = task({
id: "crashparent",
run: async (payload: Payload, { ctx }) => {
logger.log("crashparent started");
const result = await crash.triggerAndWait({});
logger.log("crashparent done", { result });
const results = await crash.batchTriggerAndWait([
{ payload: {} },
{ payload: {} },
{ payload: {} },
{ payload: {} },
{ payload: {} },
]);
logger.log("crashparent batch done", { results });
},
});
export const crash = task({
id: "crash",
run: async (payload: Payload, { ctx }) => {
logger.log(`${ctx.run.version}`);
process.exit(1);
return {
foo: "bar",
};
},
});
@@ -3,7 +3,7 @@ import { logger, schedules, task } from "@trigger.dev/sdk/v3";
export const firstScheduledTask = schedules.task({ export const firstScheduledTask = schedules.task({
id: "first-scheduled-task", id: "first-scheduled-task",
//every other minute //every other minute
cron: "0 */2 * * *", // cron: "0 */2 * * *",
run: async (payload, { ctx }) => { run: async (payload, { ctx }) => {
const distanceInMs = const distanceInMs =
payload.timestamp.getTime() - (payload.lastTimestamp ?? new Date()).getTime(); payload.timestamp.getTime() - (payload.lastTimestamp ?? new Date()).getTime();
@@ -22,10 +22,10 @@ export const firstScheduledTask = schedules.task({
export const secondScheduledTask = schedules.task({ export const secondScheduledTask = schedules.task({
id: "second-scheduled-task", id: "second-scheduled-task",
cron: { // cron: {
pattern: "0 5 * * *", // pattern: "0 5 * * *",
timezone: "Asia/Tokyo", // timezone: "Asia/Tokyo",
}, // },
run: async (payload) => {}, run: async (payload) => {},
}); });