ae46e3f7c8
This PR implements a new run TTL system and queue size limits to prevent unbounded queue growth which should help prevent situations where queues enter a "death spiral" where the queue will never be able to catch up. The main/correct way to battle this situation is to enforce a maximum TTL on all runs (e.g. up to 14 days) where runs that have been queued for that maximum TTL will get auto-expired, making room for newer runs to execute. This required creating a new TTL system that can handle higher workloads and is now deeply integrated into the RunQueue. When runs are enqueued with a TTL, they are added to their normal queue as well as to the TTL queue. When runs are dequeued, they are removed from both their normal queue and the TTL queue. If runs are dequeued by the TTL system, they are removed from their normal queue. Both these dequeues happen automatically so there is no race condition. The TTL expiration system is also made reliable by expiring runs via a Redis worker, which is enqueued to atomically inside the TTL dequeue lua script. ### Optional associated waitpoints Additionally, this PR implements an optimization where runs that aren't triggered with a dependent parent run will no longer create an associated waitpoint. Associated waitpoints are then lazily created if a dependent run wants to wait for the child run post-facto (via debounce or idempotency), which is a rare situation but is possible. This means fewer waitpoint creations but also fewer waitpoint completions for runs with no dependencies. ### Environment Queue Limits Prevents any single queue growing too large by enforcing queue size limits at trigger time. - Queue size checks happen at trigger time - runs are rejected if queue would exceed limit - Dashboard UI shows queue limits on both the Queues page and a new Limits page - In-memory caching for queue size checks to reduce Redis load ### Batch trigger fixes Currently when a batch item cannot be created for whatever reason (e.g. queue limits) the run will never get created, which means a stalled run if using `batchTriggerAndWait`. We've updated the system to handle this differently: now when a batch item cannot be triggered and converted into a run, we will eventually (after retrying 8 times up to 30s) we will create a "pre-failed" run with the error details, correctly resolving the batchTriggerAndWait.
499 lines
17 KiB
TypeScript
499 lines
17 KiB
TypeScript
import {
|
|
RunDuplicateIdempotencyKeyError,
|
|
RunEngine,
|
|
RunOneTimeUseTokenError,
|
|
} from "@internal/run-engine";
|
|
import { Tracer } from "@opentelemetry/api";
|
|
import { tryCatch } from "@trigger.dev/core/utils";
|
|
import {
|
|
TaskRunError,
|
|
taskRunErrorEnhancer,
|
|
taskRunErrorToString,
|
|
TriggerTaskRequestBody,
|
|
TriggerTraceContext,
|
|
} from "@trigger.dev/core/v3";
|
|
import {
|
|
parseTraceparent,
|
|
RunId,
|
|
serializeTraceparent,
|
|
stringifyDuration,
|
|
} from "@trigger.dev/core/v3/isomorphic";
|
|
import type { PrismaClientOrTransaction } from "@trigger.dev/database";
|
|
import { createTags } from "~/models/taskRunTag.server";
|
|
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
|
|
import { logger } from "~/services/logger.server";
|
|
import { parseDelay } from "~/utils/delays";
|
|
import { handleMetadataPacket } from "~/utils/packets";
|
|
import { startSpan } from "~/v3/tracing.server";
|
|
import type {
|
|
TriggerTaskServiceOptions,
|
|
TriggerTaskServiceResult,
|
|
} from "../../v3/services/triggerTask.server";
|
|
import { clampMaxDuration } from "../../v3/utils/maxDuration";
|
|
import { IdempotencyKeyConcern } from "../concerns/idempotencyKeys.server";
|
|
import type {
|
|
PayloadProcessor,
|
|
QueueManager,
|
|
TraceEventConcern,
|
|
TriggerRacepoints,
|
|
TriggerRacepointSystem,
|
|
TriggerTaskRequest,
|
|
TriggerTaskValidator,
|
|
} from "../types";
|
|
import { ServiceValidationError } from "~/v3/services/common.server";
|
|
|
|
class NoopTriggerRacepointSystem implements TriggerRacepointSystem {
|
|
async waitForRacepoint(options: { racepoint: TriggerRacepoints; id: string }): Promise<void> {
|
|
return;
|
|
}
|
|
}
|
|
|
|
export class RunEngineTriggerTaskService {
|
|
private readonly queueConcern: QueueManager;
|
|
private readonly validator: TriggerTaskValidator;
|
|
private readonly payloadProcessor: PayloadProcessor;
|
|
private readonly idempotencyKeyConcern: IdempotencyKeyConcern;
|
|
private readonly prisma: PrismaClientOrTransaction;
|
|
private readonly engine: RunEngine;
|
|
private readonly tracer: Tracer;
|
|
private readonly traceEventConcern: TraceEventConcern;
|
|
private readonly triggerRacepointSystem: TriggerRacepointSystem;
|
|
private readonly metadataMaximumSize: number;
|
|
|
|
constructor(opts: {
|
|
prisma: PrismaClientOrTransaction;
|
|
engine: RunEngine;
|
|
queueConcern: QueueManager;
|
|
validator: TriggerTaskValidator;
|
|
payloadProcessor: PayloadProcessor;
|
|
idempotencyKeyConcern: IdempotencyKeyConcern;
|
|
traceEventConcern: TraceEventConcern;
|
|
tracer: Tracer;
|
|
metadataMaximumSize: number;
|
|
triggerRacepointSystem?: TriggerRacepointSystem;
|
|
}) {
|
|
this.prisma = opts.prisma;
|
|
this.engine = opts.engine;
|
|
this.queueConcern = opts.queueConcern;
|
|
this.validator = opts.validator;
|
|
this.payloadProcessor = opts.payloadProcessor;
|
|
this.idempotencyKeyConcern = opts.idempotencyKeyConcern;
|
|
this.tracer = opts.tracer;
|
|
this.traceEventConcern = opts.traceEventConcern;
|
|
this.metadataMaximumSize = opts.metadataMaximumSize;
|
|
this.triggerRacepointSystem = opts.triggerRacepointSystem ?? new NoopTriggerRacepointSystem();
|
|
}
|
|
|
|
public async call({
|
|
taskId,
|
|
environment,
|
|
body,
|
|
options = {},
|
|
attempt = 0,
|
|
}: {
|
|
taskId: string;
|
|
environment: AuthenticatedEnvironment;
|
|
body: TriggerTaskRequestBody;
|
|
options?: TriggerTaskServiceOptions;
|
|
attempt?: number;
|
|
}): Promise<TriggerTaskServiceResult | undefined> {
|
|
return await startSpan(this.tracer, "RunEngineTriggerTaskService.call()", async (span) => {
|
|
span.setAttribute("taskId", taskId);
|
|
span.setAttribute("attempt", attempt);
|
|
|
|
const runFriendlyId = options?.runFriendlyId ?? RunId.generate().friendlyId;
|
|
const triggerRequest = {
|
|
taskId,
|
|
friendlyId: runFriendlyId,
|
|
environment,
|
|
body,
|
|
options,
|
|
} satisfies TriggerTaskRequest;
|
|
|
|
// Validate max attempts
|
|
const maxAttemptsValidation = this.validator.validateMaxAttempts({
|
|
taskId,
|
|
attempt,
|
|
});
|
|
|
|
if (!maxAttemptsValidation.ok) {
|
|
throw maxAttemptsValidation.error;
|
|
}
|
|
|
|
// Validate tags
|
|
const tagValidation = this.validator.validateTags({
|
|
tags: body.options?.tags,
|
|
});
|
|
|
|
if (!tagValidation.ok) {
|
|
throw tagValidation.error;
|
|
}
|
|
|
|
// Validate entitlement (unless skipChecks is enabled)
|
|
let planType: string | undefined;
|
|
|
|
if (!options.skipChecks) {
|
|
const entitlementValidation = await this.validator.validateEntitlement({
|
|
environment,
|
|
});
|
|
|
|
if (!entitlementValidation.ok) {
|
|
throw entitlementValidation.error;
|
|
}
|
|
|
|
// Extract plan type from entitlement response
|
|
planType = entitlementValidation.plan?.type;
|
|
} else {
|
|
// When skipChecks is enabled, planType should be passed via options
|
|
planType = options.planType;
|
|
|
|
if (!planType) {
|
|
logger.warn("Plan type not set but skipChecks is enabled", {
|
|
taskId,
|
|
environment: {
|
|
id: environment.id,
|
|
type: environment.type,
|
|
projectId: environment.projectId,
|
|
organizationId: environment.organizationId,
|
|
},
|
|
});
|
|
}
|
|
}
|
|
|
|
// Parse delay from either explicit delay option or debounce.delay
|
|
const delaySource = body.options?.delay ?? body.options?.debounce?.delay;
|
|
const [parseDelayError, delayUntil] = await tryCatch(parseDelay(delaySource));
|
|
|
|
if (parseDelayError) {
|
|
throw new ServiceValidationError(`Invalid delay ${delaySource}`);
|
|
}
|
|
|
|
// Validate debounce options
|
|
if (body.options?.debounce) {
|
|
if (!delayUntil) {
|
|
throw new ServiceValidationError(
|
|
`Debounce requires a valid delay duration. Provided: ${body.options.debounce.delay}`
|
|
);
|
|
}
|
|
|
|
// Always validate debounce.delay separately since it's used for rescheduling
|
|
// This catches the case where options.delay is valid but debounce.delay is invalid
|
|
const [debounceDelayError, debounceDelayUntil] = await tryCatch(
|
|
parseDelay(body.options.debounce.delay)
|
|
);
|
|
|
|
if (debounceDelayError || !debounceDelayUntil) {
|
|
throw new ServiceValidationError(
|
|
`Invalid debounce delay: ${body.options.debounce.delay}. ` +
|
|
`Supported formats: {number}s, {number}m, {number}h, {number}d, {number}w`
|
|
);
|
|
}
|
|
}
|
|
|
|
const ttl =
|
|
typeof body.options?.ttl === "number"
|
|
? stringifyDuration(body.options?.ttl)
|
|
: body.options?.ttl ?? (environment.type === "DEVELOPMENT" ? "10m" : undefined);
|
|
|
|
// Get parent run if specified
|
|
const parentRun = body.options?.parentRunId
|
|
? await this.prisma.taskRun.findFirst({
|
|
where: {
|
|
id: RunId.fromFriendlyId(body.options.parentRunId),
|
|
runtimeEnvironmentId: environment.id,
|
|
},
|
|
})
|
|
: undefined;
|
|
|
|
// Validate parent run
|
|
const parentRunValidation = this.validator.validateParentRun({
|
|
taskId,
|
|
parentRun: parentRun ?? undefined,
|
|
resumeParentOnCompletion: body.options?.resumeParentOnCompletion,
|
|
});
|
|
|
|
if (!parentRunValidation.ok) {
|
|
throw parentRunValidation.error;
|
|
}
|
|
|
|
const idempotencyKeyConcernResult = await this.idempotencyKeyConcern.handleTriggerRequest(
|
|
triggerRequest,
|
|
parentRun?.taskEventStore
|
|
);
|
|
|
|
if (idempotencyKeyConcernResult.isCached) {
|
|
return idempotencyKeyConcernResult;
|
|
}
|
|
|
|
const { idempotencyKey, idempotencyKeyExpiresAt } = idempotencyKeyConcernResult;
|
|
|
|
if (idempotencyKey) {
|
|
await this.triggerRacepointSystem.waitForRacepoint({
|
|
racepoint: "idempotencyKey",
|
|
id: idempotencyKey,
|
|
});
|
|
}
|
|
|
|
const lockedToBackgroundWorker = body.options?.lockToVersion
|
|
? await this.prisma.backgroundWorker.findFirst({
|
|
where: {
|
|
projectId: environment.projectId,
|
|
runtimeEnvironmentId: environment.id,
|
|
version: body.options?.lockToVersion,
|
|
},
|
|
select: {
|
|
id: true,
|
|
version: true,
|
|
sdkVersion: true,
|
|
cliVersion: true,
|
|
},
|
|
})
|
|
: undefined;
|
|
|
|
const { queueName, lockedQueueId } = await this.queueConcern.resolveQueueProperties(
|
|
triggerRequest,
|
|
lockedToBackgroundWorker ?? undefined
|
|
);
|
|
|
|
if (!options.skipChecks) {
|
|
const queueSizeGuard = await this.queueConcern.validateQueueLimits(
|
|
environment,
|
|
queueName
|
|
);
|
|
|
|
if (!queueSizeGuard.ok) {
|
|
throw new ServiceValidationError(
|
|
`Cannot trigger ${taskId} as the queue size limit for this environment has been reached. The maximum size is ${queueSizeGuard.maximumSize}`
|
|
);
|
|
}
|
|
}
|
|
|
|
const metadataPacket = body.options?.metadata
|
|
? handleMetadataPacket(
|
|
body.options?.metadata,
|
|
body.options?.metadataType ?? "application/json",
|
|
this.metadataMaximumSize
|
|
)
|
|
: undefined;
|
|
|
|
//upsert tags
|
|
const tags = await createTags(
|
|
{
|
|
tags: body.options?.tags,
|
|
projectId: environment.projectId,
|
|
},
|
|
this.prisma
|
|
);
|
|
|
|
const depth = parentRun ? parentRun.depth + 1 : 0;
|
|
|
|
const workerQueue = await this.queueConcern.getWorkerQueue(environment, body.options?.region);
|
|
|
|
try {
|
|
return await this.traceEventConcern.traceRun(
|
|
triggerRequest,
|
|
parentRun?.taskEventStore,
|
|
async (event, store) => {
|
|
event.setAttribute("queueName", queueName);
|
|
span.setAttribute("queueName", queueName);
|
|
event.setAttribute("runId", runFriendlyId);
|
|
span.setAttribute("runId", runFriendlyId);
|
|
|
|
const payloadPacket = await this.payloadProcessor.process(triggerRequest);
|
|
|
|
const taskRun = await this.engine.trigger(
|
|
{
|
|
friendlyId: runFriendlyId,
|
|
environment: environment,
|
|
idempotencyKey,
|
|
idempotencyKeyExpiresAt: idempotencyKey ? idempotencyKeyExpiresAt : undefined,
|
|
idempotencyKeyOptions: body.options?.idempotencyKeyOptions,
|
|
taskIdentifier: taskId,
|
|
payload: payloadPacket.data ?? "",
|
|
payloadType: payloadPacket.dataType,
|
|
context: body.context,
|
|
traceContext: this.#propagateExternalTraceContext(
|
|
event.traceContext,
|
|
parentRun?.traceContext,
|
|
event.traceparent?.spanId
|
|
),
|
|
traceId: event.traceId,
|
|
spanId: event.spanId,
|
|
parentSpanId:
|
|
options.parentAsLinkType === "replay" ? undefined : event.traceparent?.spanId,
|
|
replayedFromTaskRunFriendlyId: options.replayedFromTaskRunFriendlyId,
|
|
lockedToVersionId: lockedToBackgroundWorker?.id,
|
|
taskVersion: lockedToBackgroundWorker?.version,
|
|
sdkVersion: lockedToBackgroundWorker?.sdkVersion,
|
|
cliVersion: lockedToBackgroundWorker?.cliVersion,
|
|
concurrencyKey: body.options?.concurrencyKey,
|
|
queue: queueName,
|
|
lockedQueueId,
|
|
workerQueue,
|
|
isTest: body.options?.test ?? false,
|
|
delayUntil,
|
|
queuedAt: delayUntil ? undefined : new Date(),
|
|
maxAttempts: body.options?.maxAttempts,
|
|
taskEventStore: store,
|
|
ttl,
|
|
tags,
|
|
oneTimeUseToken: options.oneTimeUseToken,
|
|
parentTaskRunId: parentRun?.id,
|
|
rootTaskRunId: parentRun?.rootTaskRunId ?? parentRun?.id,
|
|
batch: options?.batchId
|
|
? {
|
|
id: options.batchId,
|
|
index: options.batchIndex ?? 0,
|
|
}
|
|
: undefined,
|
|
resumeParentOnCompletion: body.options?.resumeParentOnCompletion,
|
|
depth,
|
|
metadata: metadataPacket?.data,
|
|
metadataType: metadataPacket?.dataType,
|
|
seedMetadata: metadataPacket?.data,
|
|
seedMetadataType: metadataPacket?.dataType,
|
|
maxDurationInSeconds: body.options?.maxDuration
|
|
? clampMaxDuration(body.options.maxDuration)
|
|
: undefined,
|
|
machine: body.options?.machine,
|
|
priorityMs: body.options?.priority ? body.options.priority * 1_000 : undefined,
|
|
queueTimestamp:
|
|
options.queueTimestamp ??
|
|
(parentRun && body.options?.resumeParentOnCompletion
|
|
? parentRun.queueTimestamp ?? undefined
|
|
: undefined),
|
|
scheduleId: options.scheduleId,
|
|
scheduleInstanceId: options.scheduleInstanceId,
|
|
createdAt: options.overrideCreatedAt,
|
|
bulkActionId: body.options?.bulkActionId,
|
|
planType,
|
|
realtimeStreamsVersion: options.realtimeStreamsVersion,
|
|
debounce: body.options?.debounce,
|
|
// When debouncing with triggerAndWait, create a span for the debounced trigger
|
|
onDebounced:
|
|
body.options?.debounce && body.options?.resumeParentOnCompletion
|
|
? async ({ existingRun, waitpoint, debounceKey }) => {
|
|
return await this.traceEventConcern.traceDebouncedRun(
|
|
triggerRequest,
|
|
parentRun?.taskEventStore,
|
|
{
|
|
existingRun,
|
|
debounceKey,
|
|
incomplete: waitpoint.status === "PENDING",
|
|
isError: waitpoint.outputIsError,
|
|
},
|
|
async (spanEvent) => {
|
|
const spanId =
|
|
options?.parentAsLinkType === "replay"
|
|
? spanEvent.spanId
|
|
: spanEvent.traceparent?.spanId
|
|
? `${spanEvent.traceparent.spanId}:${spanEvent.spanId}`
|
|
: spanEvent.spanId;
|
|
return spanId;
|
|
}
|
|
);
|
|
}
|
|
: undefined,
|
|
},
|
|
this.prisma
|
|
);
|
|
|
|
// If the returned run has a different friendlyId, it was debounced.
|
|
// For triggerAndWait: stop the outer span since a replacement debounced span was created via onDebounced.
|
|
// For regular trigger: let the span complete normally - no replacement span needed since the
|
|
// original run already has its span from when it was first created.
|
|
if (
|
|
taskRun.friendlyId !== runFriendlyId &&
|
|
body.options?.debounce &&
|
|
body.options?.resumeParentOnCompletion
|
|
) {
|
|
event.stop();
|
|
}
|
|
|
|
const error = taskRun.error ? TaskRunError.parse(taskRun.error) : undefined;
|
|
|
|
if (error) {
|
|
event.failWithError(error);
|
|
}
|
|
|
|
const result = { run: taskRun, error, isCached: false };
|
|
|
|
if (result?.error) {
|
|
throw new ServiceValidationError(
|
|
taskRunErrorToString(taskRunErrorEnhancer(result.error))
|
|
);
|
|
}
|
|
|
|
return result;
|
|
}
|
|
);
|
|
} catch (error) {
|
|
if (error instanceof RunDuplicateIdempotencyKeyError) {
|
|
//retry calling this function, because this time it will return the idempotent run
|
|
return await this.call({
|
|
taskId,
|
|
environment,
|
|
body,
|
|
options: { ...options, runFriendlyId },
|
|
attempt: attempt + 1,
|
|
});
|
|
}
|
|
|
|
if (error instanceof RunOneTimeUseTokenError) {
|
|
throw new ServiceValidationError(
|
|
`Cannot trigger ${taskId} with a one-time use token as it has already been used.`
|
|
);
|
|
}
|
|
|
|
throw error;
|
|
}
|
|
});
|
|
}
|
|
|
|
#propagateExternalTraceContext(
|
|
traceContext: Record<string, unknown>,
|
|
parentRunTraceContext: unknown,
|
|
parentSpanId: string | undefined
|
|
): TriggerTraceContext {
|
|
if (!parentRunTraceContext) {
|
|
return traceContext;
|
|
}
|
|
|
|
const parsedParentRunTraceContext = TriggerTraceContext.safeParse(parentRunTraceContext);
|
|
|
|
if (!parsedParentRunTraceContext.success) {
|
|
return traceContext;
|
|
}
|
|
|
|
const { external } = parsedParentRunTraceContext.data;
|
|
|
|
if (!external) {
|
|
return traceContext;
|
|
}
|
|
|
|
if (!external.traceparent) {
|
|
return traceContext;
|
|
}
|
|
|
|
const parsedTraceparent = parseTraceparent(external.traceparent);
|
|
|
|
if (!parsedTraceparent) {
|
|
return traceContext;
|
|
}
|
|
|
|
const newExternalTraceparent = serializeTraceparent(
|
|
parsedTraceparent.traceId,
|
|
parentSpanId ?? parsedTraceparent.spanId,
|
|
parsedTraceparent.traceFlags
|
|
);
|
|
|
|
return {
|
|
...traceContext,
|
|
external: {
|
|
...external,
|
|
traceparent: newExternalTraceparent,
|
|
},
|
|
};
|
|
}
|
|
}
|