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:
@@ -0,0 +1,6 @@
|
|||||||
|
---
|
||||||
|
"trigger.dev": patch
|
||||||
|
"@trigger.dev/core": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
Fix for when a log flush times out and the process is checkpointed
|
||||||
@@ -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];
|
||||||
|
|||||||
+6
@@ -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;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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",
|
||||||
|
};
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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() {
|
||||||
|
|||||||
+2
@@ -0,0 +1,2 @@
|
|||||||
|
-- CreateIndex
|
||||||
|
CREATE INDEX CONCURRENTLY IF NOT EXISTS "TaskRunAttempt_taskRunId_idx" ON "TaskRunAttempt" ("taskRunId");
|
||||||
@@ -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) => {},
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user