diff --git a/.changeset/wild-swans-battle.md b/.changeset/wild-swans-battle.md
new file mode 100644
index 000000000..b1bf3ccda
--- /dev/null
+++ b/.changeset/wild-swans-battle.md
@@ -0,0 +1,6 @@
+---
+"@trigger.dev/sdk": patch
+"@trigger.dev/core": patch
+---
+
+Feature: Run execution concurrency limits
diff --git a/apps/webapp/app/components/JobsStatusTable.tsx b/apps/webapp/app/components/JobsStatusTable.tsx
index dbf66e2ab..88c0a8ff1 100644
--- a/apps/webapp/app/components/JobsStatusTable.tsx
+++ b/apps/webapp/app/components/JobsStatusTable.tsx
@@ -16,19 +16,23 @@ export type JobEnvironment = {
lastRun?: Date;
version: string;
enabled: boolean;
+ concurrencyLimit?: number | null;
+ concurrencyLimitGroup?: { name: string; concurrencyLimit: number } | null;
};
type JobStatusTableProps = {
environments: JobEnvironment[];
+ displayStyle?: "short" | "long";
};
-export function JobStatusTable({ environments }: JobStatusTableProps) {
+export function JobStatusTable({ environments, displayStyle = "short" }: JobStatusTableProps) {
return (
diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts
index e3f6d62ec..3e3b7ca89 100644
--- a/apps/webapp/app/env.server.ts
+++ b/apps/webapp/app/env.server.ts
@@ -18,14 +18,7 @@ const EnvironmentSchema = z.object({
REMIX_APP_PORT: z.string().optional(),
LOGIN_ORIGIN: z.string().default("http://localhost:3030"),
APP_ORIGIN: z.string().default("http://localhost:3030"),
- APP_ENV: z
- .union([
- z.literal("development"),
- z.literal("production"),
- z.literal("test"),
- z.literal("staging"),
- ])
- .default(process.env.NODE_ENV),
+ APP_ENV: z.string().default(process.env.NODE_ENV),
SECRET_STORE: SecretStoreOptionsSchema.default("DATABASE"),
POSTHOG_PROJECT_KEY: z.string().optional(),
TELEMETRY_TRIGGER_API_KEY: z.string().optional(),
@@ -59,6 +52,17 @@ const EnvironmentSchema = z.object({
AWS_SQS_QUEUE_URL: z.string().optional(),
AWS_SQS_BATCH_SIZE: z.coerce.number().int().optional().default(10),
DISABLE_SSE: z.string().optional(),
+
+ // Redis options
+ REDIS_HOST: z.string().optional(),
+ REDIS_READER_HOST: z.string().optional(),
+ REDIS_READER_PORT: z.coerce.number().optional(),
+ REDIS_PORT: z.coerce.number().optional(),
+ REDIS_USERNAME: z.string().optional(),
+ REDIS_PASSWORD: z.string().optional(),
+
+ DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT: z.coerce.number().int().default(10),
+ DEFAULT_DEV_ENV_EXECUTION_ATTEMPTS: z.coerce.number().int().positive().default(1),
});
export type Environment = z.infer
;
diff --git a/apps/webapp/app/models/jobRun.server.ts b/apps/webapp/app/models/jobRun.server.ts
new file mode 100644
index 000000000..537879ce2
--- /dev/null
+++ b/apps/webapp/app/models/jobRun.server.ts
@@ -0,0 +1,45 @@
+import type { JobRun, JobRunStatus } from "@trigger.dev/database";
+
+const COMPLETED_STATUSES: Array = [
+ "CANCELED",
+ "ABORTED",
+ "SUCCESS",
+ "TIMED_OUT",
+ "INVALID_PAYLOAD",
+ "FAILURE",
+ "UNRESOLVED_AUTH",
+];
+
+export function isRunCompleted(status: JobRunStatus) {
+ return COMPLETED_STATUSES.includes(status);
+}
+
+export type RunBasicStatus = "WAITING" | "PENDING" | "RUNNING" | "COMPLETED" | "FAILED";
+
+export function runBasicStatus(status: JobRunStatus): RunBasicStatus {
+ switch (status) {
+ case "WAITING_ON_CONNECTIONS":
+ case "QUEUED":
+ case "PREPROCESSING":
+ case "PENDING":
+ return "PENDING";
+ case "STARTED":
+ case "EXECUTING":
+ case "WAITING_TO_CONTINUE":
+ case "WAITING_TO_EXECUTE":
+ return "RUNNING";
+ case "FAILURE":
+ case "TIMED_OUT":
+ case "UNRESOLVED_AUTH":
+ case "CANCELED":
+ case "ABORTED":
+ case "INVALID_PAYLOAD":
+ return "FAILED";
+ case "SUCCESS":
+ return "COMPLETED";
+ default: {
+ const _exhaustiveCheck: never = status;
+ throw new Error(`Non-exhaustive match for value: ${status}`);
+ }
+ }
+}
diff --git a/apps/webapp/app/models/jobRunExecution.server.ts b/apps/webapp/app/models/jobRunExecution.server.ts
deleted file mode 100644
index c22cc5846..000000000
--- a/apps/webapp/app/models/jobRunExecution.server.ts
+++ /dev/null
@@ -1,47 +0,0 @@
-import { JobRun } from "@trigger.dev/database";
-import { PrismaClientOrTransaction } from "~/db.server";
-import { executionWorker } from "~/services/worker.server";
-
-export async function dequeueRunExecutionV2(run: JobRun, tx: PrismaClientOrTransaction) {
- return await executionWorker.dequeue(`job_run:${run.id}`, {
- tx,
- });
-}
-
-export type EnqueueRunExecutionV3Options = {
- runAt?: Date;
- skipRetrying?: boolean;
-};
-
-export async function enqueueRunExecutionV3(
- run: JobRun,
- tx: PrismaClientOrTransaction,
- options: EnqueueRunExecutionV3Options = {}
-) {
- const reason = run.status === "PREPROCESSING" ? "PREPROCESS" : "EXECUTE_JOB";
-
- return await executionWorker.enqueue(
- "performRunExecutionV3",
- {
- id: run.id,
- reason: reason,
- },
- {
- tx,
- runAt: options.runAt,
- queueName: `job_run:${run.id}`,
- jobKey: `job_run:${reason}:${run.id}`,
- maxAttempts: options.skipRetrying ? 1 : undefined,
- }
- );
-}
-
-export async function dequeueRunExecutionV3(run: JobRun, tx: PrismaClientOrTransaction) {
- await executionWorker.dequeue(`job_run:EXECUTE_JOB:${run.id}`, {
- tx,
- });
-
- await executionWorker.dequeue(`job_run:PREPROCESS:${run.id}`, {
- tx,
- });
-}
diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts
index 347dd5981..b88520a58 100644
--- a/apps/webapp/app/platform/zodWorker.server.ts
+++ b/apps/webapp/app/platform/zodWorker.server.ts
@@ -94,6 +94,11 @@ export type ZodWorkerCleanupOptions = {
type ZodWorkerReporter = (event: string, properties: Record) => Promise;
+export interface ZodWorkerRateLimiter {
+ forbiddenFlags(): Promise;
+ wrapTask(t: Task, rescheduler: Task): Task;
+}
+
export type ZodWorkerOptions = {
name: string;
runnerOptions: RunnerOptions;
@@ -104,6 +109,7 @@ export type ZodWorkerOptions = {
cleanup?: ZodWorkerCleanupOptions;
reporter?: ZodWorkerReporter;
shutdownTimeoutInMs?: number;
+ rateLimiter?: ZodWorkerRateLimiter;
};
export class ZodWorker {
@@ -116,6 +122,7 @@ export class ZodWorker {
#runner?: GraphileRunner;
#cleanup: ZodWorkerCleanupOptions | undefined;
#reporter?: ZodWorkerReporter;
+ #rateLimiter?: ZodWorkerRateLimiter;
#shutdownTimeoutInMs?: number;
#shuttingDown = false;
@@ -128,6 +135,7 @@ export class ZodWorker {
this.#recurringTasks = options.recurringTasks;
this.#cleanup = options.cleanup;
this.#reporter = options.reporter;
+ this.#rateLimiter = options.rateLimiter;
this.#shutdownTimeoutInMs = options.shutdownTimeoutInMs ?? 60000; // default to 60 seconds
}
@@ -151,6 +159,7 @@ export class ZodWorker {
noHandleSignals: true,
taskList: this.#createTaskListFromTasks(),
parsedCronItems,
+ forbiddenFlags: this.#rateLimiter?.forbiddenFlags.bind(this.#rateLimiter),
});
if (!this.#runner) {
@@ -395,7 +404,11 @@ export class ZodWorker {
return this.#handleMessage(key, payload, helpers);
};
- taskList[key] = task;
+ if (this.#rateLimiter) {
+ taskList[key] = this.#rateLimiter.wrapTask(task, this.#rescheduleTask.bind(this));
+ } else {
+ taskList[key] = task;
+ }
}
for (const [key] of Object.entries(this.#recurringTasks ?? {})) {
@@ -425,6 +438,19 @@ export class ZodWorker {
return taskList;
}
+ async #rescheduleTask(payload: unknown, helpers: JobHelpers) {
+ this.#logDebug("Rescheduling task", { payload, job: helpers.job });
+
+ await this.enqueue(helpers.job.task_identifier, payload, {
+ runAt: helpers.job.run_at,
+ queueName: helpers.job.queue_name ?? undefined,
+ priority: helpers.job.priority,
+ jobKey: helpers.job.key ?? undefined,
+ flags: Object.keys(helpers.job.flags ?? []),
+ maxAttempts: helpers.job.max_attempts,
+ });
+ }
+
#createCronItemsFromRecurringTasks() {
const cronItems: CronItem[] = [];
diff --git a/apps/webapp/app/presenters/JobPresenter.server.ts b/apps/webapp/app/presenters/JobPresenter.server.ts
index 957ed7200..023ee89aa 100644
--- a/apps/webapp/app/presenters/JobPresenter.server.ts
+++ b/apps/webapp/app/presenters/JobPresenter.server.ts
@@ -43,6 +43,13 @@ export class JobPresenter {
eventSpecification: true,
properties: true,
status: true,
+ concurrencyLimit: true,
+ concurrencyLimitGroup: {
+ select: {
+ name: true,
+ concurrencyLimit: true,
+ },
+ },
runs: {
select: {
createdAt: true,
@@ -186,6 +193,8 @@ export class JobPresenter {
enabled: alias.version.status === "ACTIVE",
lastRun: alias.version.runs.at(0)?.createdAt,
version: alias.version.version,
+ concurrencyLimit: alias.version.concurrencyLimit,
+ concurrencyLimitGroup: alias.version.concurrencyLimitGroup,
}));
const projectRootPath = projectPath({ slug: organizationSlug }, { slug: projectSlug });
diff --git a/apps/webapp/app/presenters/RunListPresenter.server.ts b/apps/webapp/app/presenters/RunListPresenter.server.ts
index 63a1df2b5..fabfcd896 100644
--- a/apps/webapp/app/presenters/RunListPresenter.server.ts
+++ b/apps/webapp/app/presenters/RunListPresenter.server.ts
@@ -6,14 +6,15 @@ export type Direction = z.infer;
type RunListOptions = {
userId: string;
- jobSlug: string;
+ jobSlug?: string;
organizationSlug: string;
projectSlug: string;
direction?: Direction;
cursor?: string;
+ pageSize?: number;
};
-const PAGE_SIZE = 20;
+const DEFAULT_PAGE_SIZE = 20;
export type RunList = Awaited>;
@@ -31,6 +32,7 @@ export class RunListPresenter {
projectSlug,
direction = "forward",
cursor,
+ pageSize = DEFAULT_PAGE_SIZE,
}: RunListOptions) {
const directionMultiplier = direction === "forward" ? 1 : -1;
@@ -41,6 +43,7 @@ export class RunListPresenter {
startedAt: true,
completedAt: true,
createdAt: true,
+ executionDuration: true,
isTest: true,
status: true,
environment: {
@@ -59,11 +62,19 @@ export class RunListPresenter {
version: true,
},
},
+ job: {
+ select: {
+ slug: true,
+ title: true,
+ },
+ },
},
where: {
- job: {
- slug: jobSlug,
- },
+ job: jobSlug
+ ? {
+ slug: jobSlug,
+ }
+ : undefined,
project: {
slug: projectSlug,
},
@@ -82,8 +93,8 @@ export class RunListPresenter {
},
},
orderBy: [{ id: "desc" }],
- //take an extra page to tell if there are more
- take: directionMultiplier * (PAGE_SIZE + 1),
+ //take an extra record to tell if there are more
+ take: directionMultiplier * (pageSize + 1),
//skip the cursor if there is one
skip: cursor ? 1 : 0,
cursor: cursor
@@ -93,7 +104,7 @@ export class RunListPresenter {
: undefined,
});
- const hasMore = runs.length > PAGE_SIZE;
+ const hasMore = runs.length > pageSize;
//get cursors for next and previous pages
let next: string | undefined;
@@ -102,19 +113,21 @@ export class RunListPresenter {
case "forward":
previous = cursor ? runs.at(0)?.id : undefined;
if (hasMore) {
- next = runs[PAGE_SIZE - 1]?.id;
+ next = runs[pageSize - 1]?.id;
}
break;
case "backward":
if (hasMore) {
previous = runs[1]?.id;
+ next = runs[pageSize]?.id;
+ } else {
+ next = runs[pageSize - 1]?.id;
}
- next = runs[PAGE_SIZE - 1]?.id;
break;
}
const runsToReturn =
- direction === "backward" && hasMore ? runs.slice(1, PAGE_SIZE + 1) : runs.slice(0, PAGE_SIZE);
+ direction === "backward" && hasMore ? runs.slice(1, pageSize + 1) : runs.slice(0, pageSize);
return {
runs: runsToReturn.map((run) => ({
@@ -123,6 +136,7 @@ export class RunListPresenter {
startedAt: run.startedAt,
completedAt: run.completedAt,
createdAt: run.createdAt,
+ executionDuration: run.executionDuration,
isTest: run.isTest,
status: run.status,
version: run.version?.version ?? "unknown",
@@ -131,6 +145,7 @@ export class RunListPresenter {
slug: run.environment.slug,
userId: run.environment.orgMember?.userId,
},
+ job: run.job,
})),
pagination: {
next,
diff --git a/apps/webapp/app/presenters/RunPresenter.server.ts b/apps/webapp/app/presenters/RunPresenter.server.ts
index 8b14e0e49..c658d5aa6 100644
--- a/apps/webapp/app/presenters/RunPresenter.server.ts
+++ b/apps/webapp/app/presenters/RunPresenter.server.ts
@@ -5,6 +5,7 @@ import {
StyleSchema,
} from "@trigger.dev/core";
import { PrismaClient, prisma } from "~/db.server";
+import { isRunCompleted, runBasicStatus } from "~/models/jobRun.server";
import { mergeProperties } from "~/utils/mergeProperties.server";
import { taskListToTree } from "~/utils/taskListToTree";
@@ -67,6 +68,8 @@ export class RunPresenter {
id: run.id,
number: run.number,
status: run.status,
+ basicStatus: runBasicStatus(run.status),
+ isFinished: isRunCompleted(run.status),
startedAt: run.startedAt,
completedAt: run.completedAt,
isTest: run.isTest,
@@ -82,6 +85,8 @@ export class RunPresenter {
runConnections: run.runConnections,
missingConnections: run.missingConnections,
error: runError,
+ executionDuration: run.executionDuration,
+ executionCount: run.executionCount,
};
}
@@ -112,6 +117,8 @@ export class RunPresenter {
isTest: true,
properties: true,
output: true,
+ executionCount: true,
+ executionDuration: true,
version: {
select: {
version: true,
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug._index/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug._index/route.tsx
index 4282a4dc8..f7aaed590 100644
--- a/apps/webapp/app/routes/_app.orgs.$organizationSlug._index/route.tsx
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug._index/route.tsx
@@ -10,6 +10,8 @@ import {
PageTitleRow,
PageTitle,
PageButtons,
+ PageInfoRow,
+ PageInfoGroup,
} from "~/components/primitives/PageHeader";
import { Paragraph } from "~/components/primitives/Paragraph";
import { useOrganization } from "~/hooks/useOrganizations";
@@ -38,6 +40,13 @@ export default function Page() {
+
+
+
+ UID: {organization.id}
+
+
+
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.environments/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.environments/route.tsx
index 412b9619b..d0a3ad7f9 100644
--- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.environments/route.tsx
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.environments/route.tsx
@@ -105,12 +105,14 @@ export default function Page() {
};
}, [selected, clients]);
- const isAnyClientFullyConfigured = useMemo(() => {
- return clients.some((client) => {
- const { DEVELOPMENT, PRODUCTION } = client.endpoints;
- return PRODUCTION.state === "configured" && DEVELOPMENT.state === PRODUCTION.state;
- });
- }, [clients]);
+ const isAnyClientFullyConfigured = clients.some((client) => {
+ const { DEVELOPMENT, PRODUCTION, STAGING } = client.endpoints;
+ return (
+ PRODUCTION.state === "configured" ||
+ DEVELOPMENT.state === "configured" ||
+ (STAGING && STAGING.state === "configured")
+ );
+ });
const organization = useOrganization();
const project = useProject();
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/ListPagination.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/ListPagination.tsx
index 0992803bc..514ba4db6 100644
--- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/ListPagination.tsx
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/ListPagination.tsx
@@ -22,31 +22,39 @@ export function ListPagination({
function NextButton({ cursor }: { cursor?: string }) {
const path = useCursorPath(cursor, "forward");
- return path ? (
+ return (
!path && e.preventDefault()}
>
Next
- ) : null;
+ );
}
function PreviousButton({ cursor }: { cursor?: string }) {
const path = useCursorPath(cursor, "backward");
- return path ? (
+ return (
!path && e.preventDefault()}
>
Prev
- ) : null;
+ );
}
function useCursorPath(cursor: string | undefined, direction: Direction) {
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route.tsx
index ceeb43196..4561f60b4 100644
--- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route.tsx
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route.tsx
@@ -72,8 +72,8 @@ export default function Page() {
-
+
+
{(open) => (
@@ -32,7 +32,7 @@ export default function Page() {
Environments
-
+
{job.status === "ACTIVE" && (
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx
index 1a76ae80a..e7638fb17 100644
--- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx
@@ -297,7 +297,9 @@ export default function Page() {
label={}
description={
<>
- Run #{run.number}{" "}
+ {typeof run.number === "number"
+ ? `Run #${run.number}`
+ : `Run ${run.id.slice(0, 8)}`}
{runStatusTitle(run.status).toLocaleLowerCase()}
diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.runs/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.runs/route.tsx
new file mode 100644
index 000000000..978263524
--- /dev/null
+++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.runs/route.tsx
@@ -0,0 +1,88 @@
+import { useNavigation } from "@remix-run/react";
+import { LoaderFunctionArgs } from "@remix-run/server-runtime";
+import { typedjson, useTypedLoaderData } from "remix-typedjson";
+import { PageBody, PageContainer } from "~/components/layout/AppLayout";
+import { LinkButton } from "~/components/primitives/Buttons";
+import {
+ PageButtons,
+ PageDescription,
+ PageHeader,
+ PageTitle,
+ PageTitleRow,
+} from "~/components/primitives/PageHeader";
+import { RunsTable } from "~/components/runs/RunsTable";
+import { useOrganization } from "~/hooks/useOrganizations";
+import { useProject } from "~/hooks/useProject";
+import { RunListPresenter } from "~/presenters/RunListPresenter.server";
+import { requireUserId } from "~/services/session.server";
+import { ProjectParamSchema, docsPath, projectPath } from "~/utils/pathBuilder";
+import { ListPagination } from "../_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/ListPagination";
+import { RunListSearchSchema } from "../_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route";
+
+export const loader = async ({ request, params }: LoaderFunctionArgs) => {
+ const userId = await requireUserId(request);
+ const { projectParam, organizationSlug } = ProjectParamSchema.parse(params);
+
+ const url = new URL(request.url);
+ const s = Object.fromEntries(url.searchParams.entries());
+ const searchParams = RunListSearchSchema.parse(s);
+
+ const presenter = new RunListPresenter();
+ const list = await presenter.call({
+ userId,
+ projectSlug: projectParam,
+ organizationSlug,
+ direction: searchParams.direction,
+ cursor: searchParams.cursor,
+ pageSize: 25,
+ });
+
+ return typedjson({
+ list,
+ });
+};
+
+export default function Page() {
+ const { list } = useTypedLoaderData();
+ const navigation = useNavigation();
+ const isLoading = navigation.state !== "idle";
+ const organization = useOrganization();
+ const project = useProject();
+
+ return (
+
+
+
+
+
+
+ Run documentation
+
+
+
+ All job runs in this project
+
+
+
+
+
+
+ );
+}
diff --git a/apps/webapp/app/routes/_app.orgs.new/route.tsx b/apps/webapp/app/routes/_app.orgs.new/route.tsx
index f6ba49474..fbb2f55a9 100644
--- a/apps/webapp/app/routes/_app.orgs.new/route.tsx
+++ b/apps/webapp/app/routes/_app.orgs.new/route.tsx
@@ -97,6 +97,7 @@ export default function NewOrganizationPage() {
{...conform.input(orgName, { type: "text" })}
placeholder="Your Organization name"
icon="organization"
+ autoFocus
/>
E.g. your company name or your workspace name.
{orgName.error}
diff --git a/apps/webapp/app/services/executions/createExecutionEvent.server.ts b/apps/webapp/app/services/executions/createExecutionEvent.server.ts
index 2d5f83791..971ac1088 100644
--- a/apps/webapp/app/services/executions/createExecutionEvent.server.ts
+++ b/apps/webapp/app/services/executions/createExecutionEvent.server.ts
@@ -10,6 +10,7 @@ export type CreateExecutionEventInput = {
eventTime: Date;
eventType: "start" | "finish";
drift?: number;
+ concurrencyLimitGroupId?: string | null;
};
export class CreateExecutionEventService {
@@ -25,7 +26,8 @@ export class CreateExecutionEventService {
"run_id",
"event_time",
"event_type",
- "drift_amount_in_ms"
+ "drift_amount_in_ms",
+ "concurrency_limit_group_id"
) VALUES (
${input.organizationId},
${input.projectId},
@@ -34,7 +36,8 @@ export class CreateExecutionEventService {
${input.runId},
${input.eventTime},
${input.eventType === "start" ? 1 : -1},
- ${input.drift}
+ ${input.drift},
+ ${input.concurrencyLimitGroupId}
)
`;
}
diff --git a/apps/webapp/app/services/jobs/registerJob.server.ts b/apps/webapp/app/services/jobs/registerJob.server.ts
index 63aeb5cb5..82d6367c8 100644
--- a/apps/webapp/app/services/jobs/registerJob.server.ts
+++ b/apps/webapp/app/services/jobs/registerJob.server.ts
@@ -14,6 +14,7 @@ import type { RuntimeEnvironment } from "~/models/runtimeEnvironment.server";
import type { AuthenticatedEnvironment } from "../apiAuth.server";
import { logger } from "../logger.server";
import { RegisterScheduleSourceService } from "../schedules/registerScheduleSource.server";
+import { executionRateLimiter } from "../runExecutionRateLimiter.server";
export class RegisterJobService {
#prismaClient: PrismaClient;
@@ -105,32 +106,28 @@ export class RegisterJobService {
},
});
- // Upsert the JobQueue
- const queueName = "default";
+ const { examples, ...eventSpecification } = metadata.event;
// Job Queues are going to be deprecated or used for something else, we're just doing this for now
- const jobQueue = await this.#prismaClient.jobQueue.upsert({
- where: {
- environmentId_name: {
- environmentId: environment.id,
- name: queueName,
- },
- },
- create: {
- environment: {
- connect: {
- id: environment.id,
- },
- },
- name: queueName,
- maxJobs: DEFAULT_MAX_CONCURRENT_RUNS,
- },
- update: {
- maxJobs: DEFAULT_MAX_CONCURRENT_RUNS,
- },
- });
-
- const { examples, ...eventSpecification } = metadata.event;
+ const concurrencyLimitGroup =
+ typeof metadata.concurrencyLimit === "object"
+ ? await this.#prismaClient.concurrencyLimitGroup.upsert({
+ where: {
+ environmentId_name: {
+ environmentId: environment.id,
+ name: metadata.concurrencyLimit.id,
+ },
+ },
+ create: {
+ environmentId: environment.id,
+ name: metadata.concurrencyLimit.id,
+ concurrencyLimit: metadata.concurrencyLimit.limit,
+ },
+ update: {
+ concurrencyLimit: metadata.concurrencyLimit.limit,
+ },
+ })
+ : null;
// Upsert the JobVersion
const jobVersion = await this.#prismaClient.jobVersion.upsert({
@@ -142,57 +139,29 @@ export class RegisterJobService {
},
},
create: {
- job: {
- connect: {
- id: job.id,
- },
- },
- endpoint: {
- connect: {
- id: endpoint.id,
- },
- },
- environment: {
- connect: {
- id: environment.id,
- },
- },
- organization: {
- connect: {
- id: environment.organizationId,
- },
- },
- project: {
- connect: {
- id: environment.projectId,
- },
- },
- queue: {
- connect: {
- id: jobQueue.id,
- },
- },
+ jobId: job.id,
+ endpointId: endpoint.id,
+ environmentId: environment.id,
+ organizationId: environment.organizationId,
+ projectId: environment.projectId,
version: metadata.version,
eventSpecification,
preprocessRuns: metadata.preprocessRuns,
startPosition: "LATEST",
status: "ACTIVE",
+ concurrencyLimitGroupId: concurrencyLimitGroup?.id ?? null,
+ concurrencyLimit:
+ typeof metadata.concurrencyLimit === "number" ? metadata.concurrencyLimit : null,
},
update: {
status: "ACTIVE",
startPosition: "LATEST",
eventSpecification,
preprocessRuns: metadata.preprocessRuns,
- queue: {
- connect: {
- id: jobQueue.id,
- },
- },
- endpoint: {
- connect: {
- id: endpoint.id,
- },
- },
+ endpointId: endpoint.id,
+ concurrencyLimitGroupId: concurrencyLimitGroup?.id ?? null,
+ concurrencyLimit:
+ typeof metadata.concurrencyLimit === "number" ? metadata.concurrencyLimit : null,
},
include: {
integrations: {
@@ -200,9 +169,28 @@ export class RegisterJobService {
integration: true,
},
},
+ concurrencyLimitGroup: true,
},
});
+ try {
+ if (jobVersion.concurrencyLimitGroup) {
+ // Upsert the maxSize for the concurrency limit group
+ await executionRateLimiter?.putConcurrencyLimitGroup(
+ jobVersion.concurrencyLimitGroup,
+ environment
+ );
+ }
+
+ await executionRateLimiter?.putJobVersionConcurrencyLimit(jobVersion, environment);
+ } catch (error) {
+ logger.error("Error setting concurrency limit", {
+ error,
+ jobVersionId: jobVersion.id,
+ environmentId: environment.id,
+ });
+ }
+
// Upsert the examples and delete any that are no longer in the metadata
const upsertedExamples = new Set();
if (examples) {
diff --git a/apps/webapp/app/services/runExecutionRateLimiter.server.ts b/apps/webapp/app/services/runExecutionRateLimiter.server.ts
new file mode 100644
index 000000000..b482c5c50
--- /dev/null
+++ b/apps/webapp/app/services/runExecutionRateLimiter.server.ts
@@ -0,0 +1,406 @@
+import { env } from "~/env.server";
+import {
+ Callback,
+ Cluster,
+ ClusterNode,
+ ClusterOptions,
+ Redis,
+ RedisOptions,
+ Result,
+} from "ioredis";
+import { JobHelpers, Task } from "graphile-worker";
+import { singleton } from "~/utils/singleton";
+import { logger } from "./logger.server";
+import { ZodWorkerRateLimiter } from "~/platform/zodWorker.server";
+import {
+ ConcurrencyLimitGroup,
+ JobRun,
+ JobVersion,
+ RuntimeEnvironment,
+} from "@trigger.dev/database";
+
+export interface RunExecutionRateLimiter {
+ putConcurrencyLimitGroup(
+ concurrencyLimitGroup: ConcurrencyLimitGroup,
+ env: RuntimeEnvironment
+ ): Promise;
+ putJobVersionConcurrencyLimit(jobVersion: JobVersion, env: RuntimeEnvironment): Promise;
+ setMaxSizeForFlag(flag: string, maxSize: number): Promise;
+ delMaxSizeForFlag(flag: string): Promise;
+ flagsForRun(
+ run: JobRun,
+ version: JobVersion & {
+ environment: RuntimeEnvironment;
+ concurrencyLimitGroup?: ConcurrencyLimitGroup;
+ }
+ ): string[];
+}
+
+declare module "ioredis" {
+ interface RedisCommander {
+ beforeTask(
+ setKey: string,
+ maxSizeKey: string,
+ forbiddenFlagsKey: string,
+ jobId: string,
+ timestamp: string,
+ windowSize: string,
+ forbiddenFlag: string,
+ maxSize: string,
+ callback?: Callback
+ ): Result;
+ rollbackBeforeTask(keys: number, ...args: string[]): Result;
+
+ afterTask(
+ setKey: string,
+ maxSizeKey: string,
+ forbiddenFlagsKey: string,
+ jobId: string,
+ timestamp: string,
+ windowSize: string,
+ forbiddenFlag: string,
+ maxSize: string,
+ callback?: Callback
+ ): Result;
+ }
+}
+
+type RedisRunExecutionRateLimiterOptions = {
+ redis?: RedisOptions;
+ cluster?: {
+ startupNodes: ClusterNode[];
+ options?: ClusterOptions;
+ };
+ defaultConcurrency?: number;
+ windowSize?: number;
+ prefix?: string;
+};
+
+const FORBIDDEN_FLAG_KEY = "forbiddenFlags";
+const KEY_PREFIX = "tr:exec:";
+
+class RedisRunExecutionRateLimiter implements RunExecutionRateLimiter, ZodWorkerRateLimiter {
+ private redis: Redis | Cluster;
+ private defaultMaxSize: number;
+ private windowSize: number;
+
+ constructor(options?: RedisRunExecutionRateLimiterOptions) {
+ this.redis = options?.cluster
+ ? new Redis.Cluster(options.cluster.startupNodes, options.cluster.options)
+ : new Redis(options?.redis ?? {});
+ this.defaultMaxSize = options?.defaultConcurrency ?? 10;
+ this.windowSize = options?.windowSize ?? 1000 * 15 * 60; // 2 minutes
+
+ this.redis.defineCommand("beforeTask", {
+ numberOfKeys: 3,
+ lua: `
+local setKey = KEYS[1]
+local maxSizeKey = KEYS[2]
+local forbiddenFlagsKey = KEYS[3]
+local jobId = ARGV[1]
+local timestamp = ARGV[2]
+local windowSize = ARGV[3]
+local forbiddenFlag = ARGV[4]
+local defaultMaxSize = ARGV[5]
+
+local maxSize = tonumber(redis.call('GET', maxSizeKey) or defaultMaxSize)
+local currentSize = redis.call('ZCOUNT', setKey, timestamp - windowSize, timestamp)
+
+if currentSize < maxSize then
+ redis.call('ZADD', setKey, timestamp, jobId)
+
+ return true
+else
+ redis.call('SADD', forbiddenFlagsKey, forbiddenFlag)
+
+ return false
+end
+ `,
+ });
+
+ // This will remove the job ID from the ZSET
+ this.redis.defineCommand("rollbackBeforeTask", {
+ lua: `
+for i, key in ipairs(KEYS) do
+ redis.call('ZREM', key, ARGV[1])
+end
+ `,
+ });
+
+ this.redis.defineCommand("afterTask", {
+ numberOfKeys: 3,
+ lua: `
+local setKey = KEYS[1]
+local maxSizeKey = KEYS[2]
+local forbiddenFlagsKey = KEYS[3]
+local jobId = ARGV[1]
+local timestamp = ARGV[2]
+local windowSize = ARGV[3]
+local forbiddenFlag = ARGV[4]
+local defaultMaxSize = ARGV[5]
+
+local maxSize = tonumber(redis.call('GET', maxSizeKey) or defaultMaxSize)
+
+-- Remove the job ID from the ZSET
+redis.call('ZREM', setKey, jobId)
+
+-- Count the current number of jobs in the window
+local currentSize = redis.call('ZCOUNT', setKey, timestamp - windowSize, timestamp)
+
+-- The cleanup of old job IDs is now an essential part of maintaining the ZSET's size
+redis.call('ZREMRANGEBYSCORE', setKey, '-inf', timestamp - windowSize)
+
+-- Update the forbidden flags based on the current size
+if currentSize < maxSize then
+ -- Only remove the forbidden flag if it's no longer needed
+ redis.call('SREM', forbiddenFlagsKey, forbiddenFlag)
+ return true
+else
+ -- No need to add the forbidden flag here as it should be handled in beforeTask
+ return false
+end
+
+ `,
+ });
+
+ if (this.redis instanceof Redis) {
+ logger.debug("⚡ RedisGraphileRateLimiter connected to Redis", {
+ host: this.redis.options.host,
+ port: this.redis.options.port,
+ });
+ } else {
+ logger.debug("⚡ RedisGraphileRateLimiter connected to Redis Cluster", {
+ nodes: this.redis.nodes,
+ });
+ }
+ }
+
+ async forbiddenFlags(): Promise {
+ return this.redis.smembers(FORBIDDEN_FLAG_KEY);
+ }
+
+ async putConcurrencyLimitGroup(
+ concurrencyLimitGroup: ConcurrencyLimitGroup,
+ env: RuntimeEnvironment
+ ): Promise {
+ await this.setMaxSizeForFlag(
+ this.flagForConcurrencyLimitGroup(concurrencyLimitGroup, env),
+ concurrencyLimitGroup.concurrencyLimit
+ );
+ }
+
+ async putJobVersionConcurrencyLimit(
+ jobVersion: JobVersion,
+ env: RuntimeEnvironment
+ ): Promise {
+ const flag = this.flagForJobVersion(jobVersion, env);
+
+ if (typeof jobVersion.concurrencyLimit === "number" && jobVersion.concurrencyLimit > 0) {
+ await this.setMaxSizeForFlag(flag, jobVersion.concurrencyLimit);
+ } else {
+ await this.delMaxSizeForFlag(flag);
+ }
+ }
+
+ flagsForRun(
+ run: JobRun,
+ version: JobVersion & {
+ environment: RuntimeEnvironment;
+ concurrencyLimitGroup?: ConcurrencyLimitGroup | null;
+ }
+ ): string[] {
+ const flags = [this.flagForOrganization(run)];
+
+ if (version.concurrencyLimitGroup) {
+ flags.push(
+ this.flagForConcurrencyLimitGroup(version.concurrencyLimitGroup, version.environment)
+ );
+ } else if (typeof version.concurrencyLimit === "number" && version.concurrencyLimit > 0) {
+ flags.push(this.flagForJobVersion(version, version.environment));
+ }
+
+ return flags;
+ }
+
+ flagForConcurrencyLimitGroup(
+ concurrencyLimitGroup: ConcurrencyLimitGroup,
+ env: RuntimeEnvironment
+ ): string {
+ return `rl:group:${env.id}:${env.slug}:${concurrencyLimitGroup.name}`;
+ }
+
+ flagForOrganization(run: JobRun): string {
+ return `rl:org:${run.organizationId}`;
+ }
+
+ flagForJobVersion(version: JobVersion, env: RuntimeEnvironment): string {
+ return `rl:job:${env.slug}:${version.id}`;
+ }
+
+ async setMaxSizeForFlag(flag: string, maxSize: number): Promise {
+ await this.redis.set(`${flag}:maxSize`, String(maxSize));
+ }
+
+ async delMaxSizeForFlag(flag: string): Promise {
+ await this.redis.del(`${flag}:maxSize`);
+ }
+
+ wrapTask(t: Task, rescheduler: Task): Task {
+ return async (payload: unknown, helpers: JobHelpers) => {
+ const flags = Object.keys(helpers.job.flags ?? {}).filter((flag) => flag.startsWith("rl:"));
+
+ if (flags.length === 0) {
+ return t(payload, helpers);
+ }
+
+ let passedFlags = [];
+
+ for (const flag of flags) {
+ const result = await this.#callBeforeTask(flag, String(helpers.job.id));
+
+ if (
+ (result.status === "fulfilled" && result.value === null) ||
+ result.status === "rejected"
+ ) {
+ logger.debug("Rolling back passed flags", {
+ flag,
+ passedFlags,
+ jobId: String(helpers.job.id),
+ result,
+ });
+ // If there are any passed flags, we need to roll them back
+ await this.#rollbackPassedFlags(passedFlags, String(helpers.job.id));
+
+ return await rescheduler(payload, helpers);
+ }
+
+ passedFlags.push(flag);
+ }
+
+ try {
+ await t(payload, helpers);
+ } finally {
+ const afterResults = await Promise.allSettled(
+ flags.map(async (flag) => this.#callAfterTask(flag, String(helpers.job.id)))
+ );
+ }
+ };
+ }
+
+ async #callBeforeTask(
+ flag: string,
+ jobId: string
+ ): Promise<
+ | { status: "fulfilled"; value: number | null; durationInMs: number }
+ | { status: "rejected"; error: any }
+ > {
+ try {
+ const now = performance.now();
+ const value = await this.redis.beforeTask(
+ flag,
+ `${flag}:maxSize`,
+ FORBIDDEN_FLAG_KEY,
+ jobId,
+ String(Date.now()),
+ String(this.windowSize),
+ flag,
+ String(this.defaultMaxSize)
+ );
+
+ const durationInMs = performance.now() - now;
+
+ return {
+ status: "fulfilled",
+ value,
+ durationInMs,
+ };
+ } catch (error) {
+ logger.error("Failed to call beforeTask", { error, flag, jobId });
+
+ return {
+ status: "rejected",
+ error,
+ };
+ }
+ }
+
+ // Method for rolling back passed flags using a single Lua script
+ async #rollbackPassedFlags(passedFlags: string[], jobId: string) {
+ if (passedFlags.length > 0) {
+ await this.redis.rollbackBeforeTask(passedFlags.length, ...passedFlags, jobId);
+ }
+ }
+
+ async #callAfterTask(flag: string, jobId: string) {
+ try {
+ const now = performance.now();
+
+ const results = await this.redis.afterTask(
+ flag,
+ `${flag}:maxSize`,
+ FORBIDDEN_FLAG_KEY,
+ jobId,
+ String(Date.now()),
+ String(this.windowSize),
+ flag,
+ String(this.defaultMaxSize)
+ );
+
+ const durationInMs = performance.now() - now;
+
+ return {
+ results,
+ durationInMs,
+ };
+ } catch (error) {
+ logger.error("Failed to call afterTask", { error, flag, jobId });
+ }
+ }
+}
+
+export const executionRateLimiter = singleton("execution-rate-limiter", getRateLimiter);
+
+function getRateLimiter() {
+ if (env.REDIS_HOST && env.REDIS_PORT) {
+ if (env.REDIS_READER_HOST) {
+ return new RedisRunExecutionRateLimiter({
+ cluster: {
+ startupNodes: [
+ { host: env.REDIS_HOST, port: env.REDIS_PORT },
+ { host: env.REDIS_READER_HOST, port: env.REDIS_READER_PORT ?? env.REDIS_PORT },
+ ],
+ options: {
+ keyPrefix: KEY_PREFIX,
+ scaleReads: "slave",
+ redisOptions: {
+ password: env.REDIS_PASSWORD,
+ tls: {
+ checkServerIdentity: () => {
+ // disable TLS verification
+ return undefined
+ }
+ },
+ enableAutoPipelining: true,
+ },
+ dnsLookup: (address, callback) => callback(null, address),
+ slotsRefreshTimeout: 10000,
+ },
+ },
+ defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
+ });
+ } else {
+ return new RedisRunExecutionRateLimiter({
+ redis: {
+ keyPrefix: KEY_PREFIX,
+ port: env.REDIS_PORT,
+ host: env.REDIS_HOST,
+ username: env.REDIS_USERNAME,
+ password: env.REDIS_PASSWORD,
+ enableAutoPipelining: true,
+ tls: {}
+ },
+ defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
+ });
+ }
+ }
+}
diff --git a/apps/webapp/app/services/runs/cancelRun.server.ts b/apps/webapp/app/services/runs/cancelRun.server.ts
index f33457a81..4b4d5926a 100644
--- a/apps/webapp/app/services/runs/cancelRun.server.ts
+++ b/apps/webapp/app/services/runs/cancelRun.server.ts
@@ -1,6 +1,6 @@
import { PrismaClient, prisma } from "~/db.server";
-import { executionWorker } from "../worker.server";
-import { dequeueRunExecutionV3 } from "~/models/jobRunExecution.server";
+import { PerformRunExecutionV3Service } from "./performRunExecutionV3.server";
+import { ResumeRunService } from "./resumeRun.server";
export class CancelRunService {
#prismaClient: PrismaClient;
@@ -39,7 +39,8 @@ export class CancelRunService {
},
});
- await dequeueRunExecutionV3(run, tx);
+ await PerformRunExecutionV3Service.dequeue(run, tx);
+ await ResumeRunService.dequeue(run, tx);
});
} catch (error) {
throw error;
diff --git a/apps/webapp/app/services/runs/continueRun.server.ts b/apps/webapp/app/services/runs/continueRun.server.ts
index c001b8a4f..982f0daac 100644
--- a/apps/webapp/app/services/runs/continueRun.server.ts
+++ b/apps/webapp/app/services/runs/continueRun.server.ts
@@ -1,6 +1,5 @@
-import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { $transaction, Prisma, PrismaClient, prisma } from "~/db.server";
-import { enqueueRunExecutionV3 } from "~/models/jobRunExecution.server";
+import { ResumeRunService } from "./resumeRun.server";
const RESUMABLE_STATUSES = ["FAILURE", "TIMED_OUT", "UNRESOLVED_AUTH", "ABORTED", "CANCELED"];
@@ -39,9 +38,7 @@ export class ContinueRunService {
},
});
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(run, tx);
},
{ timeout: 10000 }
);
diff --git a/apps/webapp/app/services/runs/createRun.server.ts b/apps/webapp/app/services/runs/createRun.server.ts
index 131835875..aab48b746 100644
--- a/apps/webapp/app/services/runs/createRun.server.ts
+++ b/apps/webapp/app/services/runs/createRun.server.ts
@@ -31,12 +31,6 @@ export class CreateRunService {
},
});
- const jobQueue = await this.#prismaClient.jobQueue.findUniqueOrThrow({
- where: {
- id: version.queueId,
- },
- });
-
const eventRecord = await this.#prismaClient.eventRecord.findUniqueOrThrow({
where: {
id: eventId,
@@ -44,22 +38,8 @@ export class CreateRunService {
});
return await $transaction(this.#prismaClient, async (tx) => {
- // Get the current max number for the given jobId
- const latestJob = await tx.jobRun.findFirst({
- where: { jobId: job.id },
- orderBy: { id: "desc" },
- select: {
- number: true,
- },
- });
-
- // Increment the number for the new execution
- const newNumber = (latestJob?.number ?? 0) + 1;
-
- // Create the new execution with the incremented number
const run = await tx.jobRun.create({
data: {
- number: newNumber,
preprocess: version.preprocessRuns,
jobId: job.id,
versionId: version.id,
@@ -68,7 +48,6 @@ export class CreateRunService {
organizationId: environment.organizationId,
projectId: environment.projectId,
endpointId: endpoint.id,
- queueId: jobQueue.id,
externalAccountId: eventRecord.externalAccountId
? eventRecord.externalAccountId
: undefined,
diff --git a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts
index a9f836058..52fc0ed78 100644
--- a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts
+++ b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts
@@ -16,7 +16,12 @@ import {
supportsFeature,
} from "@trigger.dev/core";
import { BloomFilter } from "@trigger.dev/core-backend";
-import { RuntimeEnvironmentType, type Task } from "@trigger.dev/database";
+import {
+ ConcurrencyLimitGroup,
+ JobRun,
+ JobVersion,
+ RuntimeEnvironment,
+} from "@trigger.dev/database";
import { generateErrorMessage } from "zod-error";
import { eventRecordToApiJson } from "~/api.server";
import {
@@ -26,7 +31,7 @@ import {
} from "~/consts";
import { $transaction, PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server";
import { detectResponseIsTimeout } from "~/models/endpoint.server";
-import { enqueueRunExecutionV3 } from "~/models/jobRunExecution.server";
+import { isRunCompleted } from "~/models/jobRun.server";
import { resolveRunConnections } from "~/models/runConnection.server";
import { prepareTasksForCaching, prepareTasksForCachingLegacy } from "~/models/task.server";
import { CompleteRunTaskService } from "~/routes/api.v1.runs.$runId.tasks.$id.complete";
@@ -36,8 +41,11 @@ import { EndpointApi } from "../endpointApi.server";
import { createExecutionEvent } from "../executions/createExecutionEvent.server";
import { logger } from "../logger.server";
import { ResumeTaskService } from "../tasks/resumeTask.server";
-import { workerQueue } from "../worker.server";
+import { executionWorker, workerQueue } from "../worker.server";
import { forceYieldCoordinator } from "./forceYieldCoordinator.server";
+import { ResumeRunService } from "./resumeRun.server";
+import { executionRateLimiter } from "../runExecutionRateLimiter.server";
+import { env } from "~/env.server";
type FoundRun = NonNullable>>;
type FoundTask = FoundRun["tasks"][number];
@@ -58,8 +66,15 @@ export type PerformRunExecutionV3Input = {
* @deprecated Resuming tasks now goes through ResumeTaskService, this is included here for backwards compatibility
*/
resumeTaskId?: string;
+
+ /**
+ * Specifies whether this should be the last attempt to execute the run. If so, we can't retry the run in case of a failure.
+ */
+ lastAttempt: boolean;
};
+export type RunExecutionPriority = "initial" | "resume";
+
export class PerformRunExecutionV3Service {
#prismaClient: PrismaClient;
@@ -74,206 +89,85 @@ export class PerformRunExecutionV3Service {
return;
}
- switch (input.reason) {
- case "PREPROCESS": {
- await this.#executePreprocessing(run);
- break;
- }
- case "EXECUTE_JOB": {
- await this.#executeJob(run, input, driftInMs);
- break;
- }
- }
+ await this.#executeJob(run, input, driftInMs);
}
- // Execute the preprocessing step of a run, which will send the payload to the endpoint and give the job
- // an opportunity to generate run properties based on the payload.
- // If the endpoint is not available, or the response is not ok,
- // the run execution will be marked as failed and the run will start
- async #executePreprocessing(run: FoundRun) {
- const client = new EndpointApi(run.environment.apiKey, run.endpoint.url);
- const event = eventRecordToApiJson(run.event);
-
- const { response, parser } = await client.preprocessRunRequest({
- event,
- job: {
- id: run.version.job.slug,
- version: run.version.version,
- },
- run: {
+ static async enqueue(
+ run: JobRun & {
+ version: JobVersion & {
+ environment: RuntimeEnvironment;
+ concurrencyLimitGroup?: ConcurrencyLimitGroup | null;
+ };
+ },
+ priority: RunExecutionPriority,
+ tx: PrismaClientOrTransaction,
+ options: {
+ runAt?: Date;
+ skipRetrying?: boolean;
+ } = {}
+ ) {
+ return await executionWorker.enqueue(
+ "performRunExecutionV3",
+ {
id: run.id,
- isTest: run.isTest,
+ reason: "EXECUTE_JOB",
},
- environment: {
- id: run.environment.id,
- slug: run.environment.slug,
- type: run.environment.type,
- },
- organization: {
- id: run.organization.id,
- slug: run.organization.slug,
- title: run.organization.title,
- },
- account: run.externalAccount
- ? {
- id: run.externalAccount.identifier,
- metadata: run.externalAccount.metadata,
- }
- : undefined,
- });
-
- if (!response) {
- return await this.#failRunExecution(this.#prismaClient, "PREPROCESS", run, {
- message: "Could not connect to the endpoint",
- });
- }
-
- if (!response.ok) {
- return await this.#failRunExecution(this.#prismaClient, "PREPROCESS", run, {
- message: `Endpoint responded with ${response.status} status code`,
- });
- }
-
- const rawBody = await response.text();
- const safeBody = safeJsonZodParse(parser, rawBody);
-
- if (!safeBody) {
- return await this.#failRunExecution(this.#prismaClient, "PREPROCESS", run, {
- message: "Endpoint responded with invalid JSON",
- });
- }
-
- if (!safeBody.success) {
- return await this.#failRunExecution(this.#prismaClient, "PREPROCESS", run, {
- message: generateErrorMessage(safeBody.error.issues),
- });
- }
-
- if (safeBody.data.abort) {
- return this.#failRunExecution(
- this.#prismaClient,
- "PREPROCESS",
- run,
- { message: "Endpoint aborted the run" },
- "ABORTED"
- );
- } else {
- await $transaction(this.#prismaClient, async (tx) => {
- await tx.jobRun.update({
- where: {
- id: run.id,
- },
- data: {
- status: "STARTED",
- startedAt: new Date(),
- properties: safeBody.data.properties,
- forceYieldImmediately: false,
- },
- });
-
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
- });
- }
+ {
+ tx,
+ runAt: options.runAt,
+ jobKey: `job_run:EXECUTE_JOB:${run.id}`,
+ maxAttempts: options.skipRetrying ? env.DEFAULT_DEV_ENV_EXECUTION_ATTEMPTS : undefined,
+ flags: executionRateLimiter?.flagsForRun(run, run.version) ?? [],
+ priority: priority === "initial" ? 0 : -1,
+ }
+ );
}
+
+ static async dequeue(run: JobRun, tx: PrismaClientOrTransaction) {
+ await executionWorker.dequeue(`job_run:EXECUTE_JOB:${run.id}`, {
+ tx,
+ });
+ }
+
async #executeJob(run: FoundRun, input: PerformRunExecutionV3Input, driftInMs: number = 0) {
try {
- const { isRetry, resumeTaskId } = input;
-
- if (run.status === "CANCELED") {
- await this.#cancelExecution(run);
+ if (isRunCompleted(run.status)) {
return;
}
- try {
- if (
- typeof process.env.BLOCKED_ORGS === "string" &&
- process.env.BLOCKED_ORGS.includes(run.organizationId)
- ) {
- logger.debug("Skipping execution for blocked org", {
- orgId: run.organizationId,
- });
-
- await this.#prismaClient.jobRun.update({
- where: {
- id: run.id,
- },
- data: {
- status: "CANCELED",
- completedAt: new Date(),
- },
- });
-
- return;
- }
- } catch (e) {}
-
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url);
const event = eventRecordToApiJson(run.event);
const startedAt = new Date();
- const { executionCount } = await this.#prismaClient.jobRun.update({
- where: {
- id: run.id,
- },
- data: {
- status: run.status === "QUEUED" ? "STARTED" : run.status,
- startedAt: run.startedAt ?? new Date(),
- executionCount: {
- increment: 1,
- },
- },
- select: {
- executionCount: true,
- },
- });
-
const connections = await resolveRunConnections(run.runConnections);
if (!connections.success) {
- return this.#failRunExecution(this.#prismaClient, "EXECUTE_JOB", run, {
+ return this.#failRunExecution(this.#prismaClient, run, {
message: `Could not resolve all connections for run ${run.id}. This should not happen`,
});
}
- let resumedTask: Task | undefined;
-
- if (resumeTaskId) {
- resumedTask =
- (await this.#prismaClient.task.findUnique({
- where: {
- id: resumeTaskId,
- },
- })) ?? undefined;
-
- if (resumedTask) {
- resumedTask = await this.#prismaClient.task.update({
- where: {
- id: resumeTaskId,
- },
- data: {
- status: resumedTask.noop ? "COMPLETED" : "RUNNING",
- completedAt: resumedTask.noop ? new Date() : undefined,
- },
- });
- }
- }
-
const sourceContext = RunSourceContextSchema.safeParse(run.event.sourceContext);
const executionBody = await this.#createExecutionBody(
run,
- [run.tasks, resumedTask].flat().filter(Boolean),
+ run.tasks,
startedAt,
- isRetry,
+ false,
connections.auth,
event,
sourceContext.success ? sourceContext.data : undefined
);
- forceYieldCoordinator.registerRun(run.id);
+ await this.#prismaClient.jobRun.update({
+ where: {
+ id: run.id,
+ },
+ data: {
+ status: "EXECUTING",
+ },
+ });
await createExecutionEvent({
eventType: "start",
@@ -284,8 +178,12 @@ export class PerformRunExecutionV3Service {
projectId: run.projectId,
jobId: run.jobId,
runId: run.id,
+ concurrencyLimitGroupId: run.version.concurrencyLimitGroupId,
});
+ forceYieldCoordinator.registerRun(run.id);
+
+ // TODO: add the ability to abort the execution from any server using Redis pub/sub
const { response, parser, errorParser, headersParser, durationInMs } =
await client.executeJobRequest(executionBody);
@@ -298,12 +196,13 @@ export class PerformRunExecutionV3Service {
projectId: run.projectId,
jobId: run.jobId,
runId: run.id,
+ concurrencyLimitGroupId: run.version.concurrencyLimitGroupId,
});
forceYieldCoordinator.deregisterRun(run.id);
if (!response) {
- return await this.#failRunExecutionWithRetry({
+ return await this.#failRunExecutionWithRetry(run, input.lastAttempt, {
message: `Connection could not be established to the endpoint (${run.endpoint.url})`,
});
}
@@ -393,14 +292,9 @@ export class PerformRunExecutionV3Service {
if (errorBody && errorBody.success) {
// Only retry if the error isn't a 4xx
if (response.status >= 400 && response.status <= 499) {
- return await this.#failRunExecution(
- this.#prismaClient,
- "EXECUTE_JOB",
- run,
- errorBody.data
- );
+ return await this.#failRunExecution(this.#prismaClient, run, errorBody.data);
} else {
- return await this.#failRunExecutionWithRetry(errorBody.data);
+ return await this.#failRunExecutionWithRetry(run, input.lastAttempt, errorBody.data);
}
}
@@ -408,7 +302,6 @@ export class PerformRunExecutionV3Service {
if (response.status >= 400 && response.status <= 499 && response.status !== 408) {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
{
message: `Endpoint responded with ${response.status} status code`,
@@ -423,11 +316,10 @@ export class PerformRunExecutionV3Service {
this.#prismaClient,
run,
input,
- durationInMs,
- executionCount
+ durationInMs
);
} else {
- return await this.#failRunExecutionWithRetry({
+ return await this.#failRunExecutionWithRetry(run, input.lastAttempt, {
message: `Endpoint responded with ${response.status} status code`,
});
}
@@ -439,7 +331,6 @@ export class PerformRunExecutionV3Service {
if (!safeBody) {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
{
message: "Endpoint responded with invalid JSON",
@@ -452,7 +343,6 @@ export class PerformRunExecutionV3Service {
if (!safeBody.success) {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
{
message: generateErrorMessage(safeBody.error.issues),
@@ -491,7 +381,6 @@ export class PerformRunExecutionV3Service {
break;
}
case "CANCELED": {
- await this.#cancelExecution(run);
break;
}
case "UNRESOLVED_AUTH_ERROR": {
@@ -644,6 +533,9 @@ export class PerformRunExecutionV3Service {
executionDuration: {
increment: durationInMs,
},
+ executionCount: {
+ increment: 1,
+ },
},
});
@@ -661,17 +553,18 @@ export class PerformRunExecutionV3Service {
run: FoundRun,
data: RunJobResumeWithTask,
durationInMs: number,
- executionCount: number = 1
+ executionCountIncrement: number = 1
) {
return await $transaction(this.#prismaClient, async (tx) => {
await tx.jobRun.update({
where: { id: run.id },
data: {
+ status: "WAITING_TO_CONTINUE",
executionDuration: {
increment: durationInMs,
},
executionCount: {
- increment: executionCount,
+ increment: executionCountIncrement,
},
},
});
@@ -744,7 +637,6 @@ export class PerformRunExecutionV3Service {
case "ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
childError.error ?? undefined,
"FAILURE",
@@ -754,7 +646,6 @@ export class PerformRunExecutionV3Service {
case "INVALID_PAYLOAD": {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
childError.errors,
"INVALID_PAYLOAD",
@@ -774,7 +665,6 @@ export class PerformRunExecutionV3Service {
case "UNRESOLVED_AUTH_ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
- "EXECUTE_JOB",
run,
childError.issues,
"UNRESOLVED_AUTH",
@@ -805,14 +695,7 @@ export class PerformRunExecutionV3Service {
});
}
- await this.#failRunExecution(
- tx,
- "EXECUTE_JOB",
- execution,
- data.error ?? undefined,
- "FAILURE",
- durationInMs
- );
+ await this.#failRunExecution(tx, execution, data.error ?? undefined, "FAILURE", durationInMs);
});
}
@@ -822,14 +705,7 @@ export class PerformRunExecutionV3Service {
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
- await this.#failRunExecution(
- tx,
- "EXECUTE_JOB",
- execution,
- data.issues,
- "UNRESOLVED_AUTH",
- durationInMs
- );
+ await this.#failRunExecution(tx, execution, data.issues, "UNRESOLVED_AUTH", durationInMs);
});
}
@@ -839,14 +715,7 @@ export class PerformRunExecutionV3Service {
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
- await this.#failRunExecution(
- tx,
- "EXECUTE_JOB",
- execution,
- data.errors,
- "INVALID_PAYLOAD",
- durationInMs
- );
+ await this.#failRunExecution(tx, execution, data.errors, "INVALID_PAYLOAD", durationInMs);
});
}
@@ -860,7 +729,6 @@ export class PerformRunExecutionV3Service {
if (run.yieldedExecutions.length + 1 > MAX_RUN_YIELDED_EXECUTIONS) {
return await this.#failRunExecution(
tx,
- "EXECUTE_JOB",
run,
{
message: `Run has yielded too many times, the maximum is ${MAX_RUN_YIELDED_EXECUTIONS}`,
@@ -875,6 +743,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
+ status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -892,9 +761,7 @@ export class PerformRunExecutionV3Service {
},
});
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(run, tx);
});
}
@@ -910,6 +777,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
+ status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -933,9 +801,7 @@ export class PerformRunExecutionV3Service {
},
});
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(run, tx);
});
}
@@ -981,9 +847,7 @@ export class PerformRunExecutionV3Service {
output: data.output ? (JSON.parse(data.output) as any) : undefined,
});
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(run, tx);
});
}
@@ -1035,6 +899,7 @@ export class PerformRunExecutionV3Service {
status: "WAITING",
run: {
update: {
+ status: "WAITING_TO_CONTINUE",
executionDuration: {
increment: durationInMs,
},
@@ -1054,8 +919,7 @@ export class PerformRunExecutionV3Service {
prisma: PrismaClientOrTransaction,
run: FoundRun,
input: PerformRunExecutionV3Input,
- durationInMs: number,
- executionCount: number
+ durationInMs: number
) {
await $transaction(prisma, async (tx) => {
const executionDuration = run.executionDuration + durationInMs;
@@ -1064,7 +928,6 @@ export class PerformRunExecutionV3Service {
if (executionDuration >= run.organization.maximumExecutionTimePerRunInMs) {
await this.#failRunExecution(
tx,
- "EXECUTE_JOB",
run,
{
message: `Execution timed out after ${
@@ -1112,7 +975,6 @@ export class PerformRunExecutionV3Service {
await this.#failRunExecution(
tx,
- "EXECUTE_JOB",
run,
{
message: `Function timeout detected in ${
@@ -1147,102 +1009,73 @@ export class PerformRunExecutionV3Service {
});
// The run has timed out, so we need to enqueue a new execution
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(run, tx);
});
}
- async #failRunExecutionWithRetry(output: Record): Promise {
+ async #failRunExecutionWithRetry(
+ run: FoundRun,
+ lastAttempt: boolean,
+ output: Record
+ ): Promise {
+ if (lastAttempt) {
+ return await this.#failRunExecution(this.#prismaClient, run, output);
+ }
+
+ await this.#prismaClient.jobRun.update({
+ where: { id: run.id },
+ data: {
+ status: "WAITING_TO_EXECUTE",
+ },
+ });
+
throw new Error(JSON.stringify(output));
}
async #failRunExecution(
prisma: PrismaClientOrTransaction,
- reason: "EXECUTE_JOB" | "PREPROCESS",
run: FoundRun,
output: Record,
status: "FAILURE" | "ABORTED" | "TIMED_OUT" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE",
durationInMs: number = 0
): Promise {
await $transaction(prisma, async (tx) => {
- switch (reason) {
- case "EXECUTE_JOB": {
- // If the execution is an EXECUTE_JOB reason, we need to fail the run
- await tx.jobRun.update({
- where: { id: run.id },
- data: {
- completedAt: new Date(),
- status,
- output,
- executionDuration: {
- increment: durationInMs,
- },
- tasks: {
- updateMany: {
- where: {
- status: {
- in: ["WAITING", "RUNNING", "PENDING"],
- },
- },
- data: {
- status: status === "TIMED_OUT" ? "CANCELED" : "ERRORED",
- completedAt: new Date(),
- },
+ // If the execution is an EXECUTE_JOB reason, we need to fail the run
+ await tx.jobRun.update({
+ where: { id: run.id },
+ data: {
+ completedAt: new Date(),
+ status,
+ output,
+ executionDuration: {
+ increment: durationInMs,
+ },
+ tasks: {
+ updateMany: {
+ where: {
+ status: {
+ in: ["WAITING", "RUNNING", "PENDING"],
},
},
- forceYieldImmediately: false,
- },
- });
-
- await workerQueue.enqueue(
- "deliverRunSubscriptions",
- {
- id: run.id,
- },
- { tx }
- );
-
- break;
- }
- case "PREPROCESS": {
- // If the status is ABORTED, we need to fail the run
- if (status === "ABORTED") {
- await tx.jobRun.update({
- where: { id: run.id },
data: {
+ status: status === "TIMED_OUT" ? "CANCELED" : "ERRORED",
completedAt: new Date(),
- status,
- output,
},
- });
-
- break;
- }
-
- await tx.jobRun.update({
- where: {
- id: run.id,
},
- data: {
- status: "STARTED",
- startedAt: new Date(),
- },
- });
+ },
+ forceYieldImmediately: false,
+ },
+ });
- await enqueueRunExecutionV3(run, tx, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
-
- break;
- }
- }
+ await workerQueue.enqueue(
+ "deliverRunSubscriptions",
+ {
+ id: run.id,
+ },
+ { tx }
+ );
});
}
-
- async #cancelExecution(run: FoundRun) {
- return;
- }
}
function prepareNoOpTasksBloomFilter(possibleTasks: FoundTask[]): string {
diff --git a/apps/webapp/app/services/runs/resumeRun.server.ts b/apps/webapp/app/services/runs/resumeRun.server.ts
new file mode 100644
index 000000000..6b99035ff
--- /dev/null
+++ b/apps/webapp/app/services/runs/resumeRun.server.ts
@@ -0,0 +1,158 @@
+import { JobRun, RuntimeEnvironmentType } from "@trigger.dev/database";
+import { PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server";
+import { workerQueue } from "../worker.server";
+import { PerformRunExecutionV3Service, RunExecutionPriority } from "./performRunExecutionV3.server";
+
+type FoundRun = NonNullable>>;
+
+export class ResumeRunService {
+ #prismaClient: PrismaClient;
+
+ constructor(prismaClient: PrismaClient = prisma) {
+ this.#prismaClient = prismaClient;
+ }
+
+ public async call(id: string) {
+ const run = await findRun(this.#prismaClient, id);
+
+ if (!run) {
+ return;
+ }
+
+ switch (run.status) {
+ case "ABORTED":
+ case "CANCELED":
+ case "FAILURE":
+ case "INVALID_PAYLOAD":
+ case "SUCCESS":
+ case "TIMED_OUT":
+ case "UNRESOLVED_AUTH": {
+ return;
+ }
+ case "QUEUED": {
+ await this.#resumeQueuedRun(run);
+ break;
+ }
+ case "WAITING_TO_EXECUTE": {
+ await this.#executeRun(run, "resume");
+ break;
+ }
+ case "WAITING_TO_CONTINUE": {
+ await this.#resumeWaitingToContinueRun(run);
+ break;
+ }
+ case "STARTED": {
+ await this.#resumeStartedRun(run);
+ break;
+ }
+ case "PENDING":
+ case "PREPROCESSING": {
+ await this.#resumePendingRun(run);
+ break;
+ }
+ case "EXECUTING": {
+ throw new Error("Cannot resume a run that is currently executing");
+ }
+ case "WAITING_ON_CONNECTIONS": {
+ throw new Error("Cannot resume a run that is waiting on connections");
+ }
+ default: {
+ const _exhaustiveCheck: never = run.status;
+ throw new Error(`Non-exhaustive match for value: ${run.status}`);
+ }
+ }
+ }
+
+ async #resumeQueuedRun(run: FoundRun) {
+ await this.#prismaClient.jobRun.update({
+ where: {
+ id: run.id,
+ },
+ data: {
+ startedAt: run.startedAt ?? new Date(),
+ },
+ });
+
+ await this.#executeRun(run, "initial");
+ }
+
+ async #resumeStartedRun(run: FoundRun) {
+ await this.#prismaClient.jobRun.update({
+ where: {
+ id: run.id,
+ },
+ data: {
+ status: "WAITING_TO_EXECUTE",
+ },
+ });
+
+ await this.#executeRun(run, "initial");
+ }
+
+ async #resumeWaitingToContinueRun(run: FoundRun) {
+ await this.#prismaClient.jobRun.update({
+ where: {
+ id: run.id,
+ },
+ data: {
+ status: "WAITING_TO_EXECUTE",
+ },
+ });
+
+ await this.#executeRun(run, "resume");
+ }
+
+ async #resumePendingRun(run: FoundRun) {
+ await this.#prismaClient.jobRun.update({
+ where: {
+ id: run.id,
+ },
+ data: {
+ status: "QUEUED",
+ startedAt: new Date(),
+ },
+ });
+
+ await this.#executeRun(run, "initial");
+ }
+
+ async #executeRun(run: FoundRun, priority: RunExecutionPriority) {
+ await PerformRunExecutionV3Service.enqueue(run, priority, this.#prismaClient, {
+ skipRetrying: run.version.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
+ });
+ }
+
+ static async enqueue(run: JobRun, tx: PrismaClientOrTransaction, runAt?: Date) {
+ return await workerQueue.enqueue(
+ "resumeRun",
+ {
+ id: run.id,
+ },
+ {
+ tx,
+ runAt: runAt ?? run.createdAt,
+ jobKey: `run_resume:${run.id}`,
+ }
+ );
+ }
+
+ static async dequeue(run: JobRun, tx: PrismaClientOrTransaction) {
+ await workerQueue.dequeue(`run_resume:${run.id}`, {
+ tx,
+ });
+ }
+}
+
+async function findRun(prisma: PrismaClientOrTransaction, id: string) {
+ return await prisma.jobRun.findUnique({
+ where: { id },
+ include: {
+ version: {
+ include: {
+ environment: true,
+ concurrencyLimitGroup: true,
+ },
+ },
+ },
+ });
+}
diff --git a/apps/webapp/app/services/runs/startRun.server.ts b/apps/webapp/app/services/runs/startRun.server.ts
index 4ad39970b..62f9f0ed1 100644
--- a/apps/webapp/app/services/runs/startRun.server.ts
+++ b/apps/webapp/app/services/runs/startRun.server.ts
@@ -1,13 +1,13 @@
import {
- RuntimeEnvironmentType,
type ConnectionType,
type Integration,
type IntegrationConnection,
} from "@trigger.dev/database";
import type { PrismaClient, PrismaClientOrTransaction } from "~/db.server";
-import { prisma } from "~/db.server";
-import { enqueueRunExecutionV3 } from "~/models/jobRunExecution.server";
+import { $transaction, prisma } from "~/db.server";
import { workerQueue } from "../worker.server";
+import { ResumeRunService } from "./resumeRun.server";
+import { createHash } from "node:crypto";
type FoundRun = NonNullable>>;
type RunConnectionsByKey = Awaited>;
@@ -59,23 +59,24 @@ export class StartRunService {
: undefined
)
.filter(Boolean);
+ const lockId = jobIdToLockId(run.jobId);
- const updateRun = async () => {
- if (run.preprocess) {
- // Start the jobRun and increment the jobCount
- return await this.#prismaClient.jobRun.update({
- where: { id },
- data: {
- status: "PREPROCESSING",
- runConnections: {
- create: createRunConnections,
- },
- },
+ await $transaction(
+ this.#prismaClient,
+ async (tx) => {
+ await tx.$executeRaw`SELECT pg_advisory_xact_lock(${lockId})`;
+
+ const counter = await tx.jobCounter.upsert({
+ where: { jobId: run.jobId },
+ update: { lastNumber: { increment: 1 } },
+ create: { jobId: run.jobId, lastNumber: 1 },
+ select: { lastNumber: true },
});
- } else {
- return await this.#prismaClient.jobRun.update({
+
+ const updatedRun = await this.#prismaClient.jobRun.update({
where: { id },
data: {
+ number: counter.lastNumber,
status: "QUEUED",
queuedAt: new Date(),
runConnections: {
@@ -83,14 +84,11 @@ export class StartRunService {
},
},
});
- }
- };
- const updatedRun = await updateRun();
-
- await enqueueRunExecutionV3(updatedRun, this.#prismaClient, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(updatedRun, tx);
+ },
+ { timeout: 60000 }
+ );
}
async #handleMissingConnections(id: string, runConnectionsByKey: RunConnectionsByKey) {
@@ -237,3 +235,8 @@ async function createRunConnections(tx: PrismaClientOrTransaction, run: FoundRun
function hasMissingConnections(runConnectionsByKey: RunConnectionsByKey) {
return Object.values(runConnectionsByKey).some((connection) => connection.result === "missing");
}
+
+function jobIdToLockId(jobId: string): number {
+ // Convert jobId to a unique lock identifier
+ return parseInt(createHash("sha256").update(jobId).digest("hex").slice(0, 8), 16);
+}
diff --git a/apps/webapp/app/services/tasks/resumeTask.server.ts b/apps/webapp/app/services/tasks/resumeTask.server.ts
index c47b04791..03a59b473 100644
--- a/apps/webapp/app/services/tasks/resumeTask.server.ts
+++ b/apps/webapp/app/services/tasks/resumeTask.server.ts
@@ -1,8 +1,7 @@
import { PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server";
-import { workerQueue } from "../worker.server";
-import { enqueueRunExecutionV3 } from "~/models/jobRunExecution.server";
-import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { logger } from "../logger.server";
+import { ResumeRunService } from "../runs/resumeRun.server";
+import { workerQueue } from "../worker.server";
type FoundTask = Awaited>;
@@ -81,9 +80,7 @@ export class ResumeTaskService {
}
}
- await enqueueRunExecutionV3(task.run, this.#prismaClient, {
- skipRetrying: task.run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
- });
+ await ResumeRunService.enqueue(task.run, this.#prismaClient);
}
public static async enqueue(id: string, runAt?: Date, tx?: PrismaClientOrTransaction) {
diff --git a/apps/webapp/app/services/tasks/runTask.server.ts b/apps/webapp/app/services/tasks/runTask.server.ts
index eebce56cd..f0a8edd28 100644
--- a/apps/webapp/app/services/tasks/runTask.server.ts
+++ b/apps/webapp/app/services/tasks/runTask.server.ts
@@ -71,7 +71,7 @@ export class RunTaskService {
status = "CANCELED";
} else {
status =
- delayUntilInFuture || callbackEnabled || taskBody.trigger
+ delayUntilInFuture || callbackEnabled
? "WAITING"
: taskBody.noop
? "COMPLETED"
@@ -180,7 +180,7 @@ export class RunTaskService {
if (existingTask) {
if (existingTask.status === "CANCELED") {
const existingTaskStatus =
- delayUntilInFuture || callbackEnabled || taskBody.trigger
+ delayUntilInFuture || callbackEnabled
? "WAITING"
: taskBody.noop
? "COMPLETED"
diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts
index eae2ba59e..92ac8a1ca 100644
--- a/apps/webapp/app/services/worker.server.ts
+++ b/apps/webapp/app/services/worker.server.ts
@@ -26,6 +26,8 @@ import { DeliverRunSubscriptionsService } from "./runs/deliverRunSubscriptions.s
import { ResumeTaskService } from "./tasks/resumeTask.server";
import { ExpireDispatcherService } from "./dispatchers/expireDispatcher.server";
import { InvokeEphemeralDispatcherService } from "./dispatchers/invokeEphemeralEventDispatcher.server";
+import { ResumeRunService } from "./runs/resumeRun.server";
+import { executionRateLimiter } from "./runExecutionRateLimiter.server";
import { DeliverWebhookRequestService } from "./sources/deliverWebhookRequest.server";
const workerCatalog = {
@@ -97,6 +99,9 @@ const workerCatalog = {
expireDispatcher: z.object({
id: z.string(),
}),
+ resumeRun: z.object({
+ id: z.string(),
+ }),
};
const executionWorkerCatalog = {
@@ -225,7 +230,6 @@ function getWorkerQueue() {
"events.invokeDispatcher": {
priority: 0, // smaller number = higher priority
maxAttempts: 6,
- queueName: (payload) => `dispatcher:${payload.id}`, // use a queue for a dispatcher so runs are created sequentially
handler: async (payload, job) => {
const service = new InvokeDispatcherService();
@@ -412,6 +416,15 @@ function getWorkerQueue() {
handler: async (payload) => {
const service = new ExpireDispatcherService();
+ return await service.call(payload.id);
+ },
+ },
+ resumeRun: {
+ priority: 0,
+ maxAttempts: 10,
+ handler: async (payload, job) => {
+ const service = new ResumeRunService();
+
return await service.call(payload.id);
},
},
@@ -433,6 +446,7 @@ function getExecutionWorkerQueue() {
},
shutdownTimeoutInMs: env.GRACEFUL_SHUTDOWN_TIMEOUT,
schema: executionWorkerCatalog,
+ rateLimiter: executionRateLimiter,
tasks: {
performRunExecutionV2: {
priority: 0, // smaller number = higher priority
@@ -445,6 +459,7 @@ function getExecutionWorkerQueue() {
reason: payload.reason,
resumeTaskId: payload.resumeTaskId,
isRetry: payload.isRetry,
+ lastAttempt: job.max_attempts === job.attempts,
});
},
},
@@ -461,6 +476,7 @@ function getExecutionWorkerQueue() {
id: payload.id,
reason: payload.reason,
isRetry: false,
+ lastAttempt: job.max_attempts === job.attempts,
},
driftInMs
);
diff --git a/apps/webapp/app/utils.ts b/apps/webapp/app/utils.ts
index 5a50c6d84..bf8bbb621 100644
--- a/apps/webapp/app/utils.ts
+++ b/apps/webapp/app/utils.ts
@@ -158,19 +158,10 @@ export const obfuscateApiKey = (apiKey: string) => {
return `${prefix}_${slug}_${"*".repeat(secretPart.length)}`;
};
-export function appEnvTitleTag(appEnv?: "test" | "production" | "development" | "staging"): string {
- if (!appEnv) {
+export function appEnvTitleTag(appEnv?: string): string {
+ if (!appEnv || appEnv === "production") {
return "";
}
- switch (appEnv) {
- case "test":
- return " (test)";
- case "production":
- return "";
- case "development":
- return " (dev)";
- case "staging":
- return " (staging)";
- }
+ return ` (${appEnv})`
}
diff --git a/apps/webapp/app/utils/pathBuilder.ts b/apps/webapp/app/utils/pathBuilder.ts
index 8683d4ff5..a9417fb59 100644
--- a/apps/webapp/app/utils/pathBuilder.ts
+++ b/apps/webapp/app/utils/pathBuilder.ts
@@ -126,6 +126,10 @@ export function projectPath(organization: OrgForPath, project: ProjectForPath) {
return `/orgs/${organizationParam(organization)}/projects/${projectParam(project)}`;
}
+export function projectRunsPath(organization: OrgForPath, project: ProjectForPath) {
+ return `${projectPath(organization, project)}/runs`;
+}
+
export function projectSetupPath(organization: OrgForPath, project: ProjectForPath) {
return `${projectPath(organization, project)}/setup`;
}
diff --git a/apps/webapp/package.json b/apps/webapp/package.json
index 82c60d8f0..03952397d 100644
--- a/apps/webapp/package.json
+++ b/apps/webapp/package.json
@@ -84,6 +84,7 @@
"highlight.run": "^7.3.4",
"humanize-duration": "^3.27.3",
"intl-parse-accept-language": "^1.0.0",
+ "ioredis": "^5.3.2",
"isbot": "^3.6.5",
"jsonpointer": "^5.0.1",
"lodash.omit": "^4.5.0",
diff --git a/docker/dev-compose.yml b/docker/dev-compose.yml
index 015139b68..642510e28 100644
--- a/docker/dev-compose.yml
+++ b/docker/dev-compose.yml
@@ -2,6 +2,7 @@ version: "3"
volumes:
database-data:
+ redis-data:
networks:
app_network:
@@ -42,3 +43,21 @@ services:
PORT: 3030
networks:
- app_network
+
+ redis:
+ container_name: redis
+ image: redis:7
+ restart: always
+ volumes:
+ - redis-data:/data
+ networks:
+ - app_network
+ ports:
+ - 6379:6379
+
+ redisinsight:
+ image: redislabs/redisinsight:latest
+ ports:
+ - "8001:8001"
+ volumes:
+ - redis-data:/redisinsight
diff --git a/docker/docker-compose.yml b/docker/docker-compose.yml
index 41e935745..7c6f6cc58 100644
--- a/docker/docker-compose.yml
+++ b/docker/docker-compose.yml
@@ -3,6 +3,7 @@ version: "3"
volumes:
database-data:
pgadmin-data:
+ redis-data:
networks:
app_network:
@@ -41,3 +42,21 @@ services:
- 5480:80
depends_on:
- database
+
+ redis:
+ container_name: redis
+ image: redis:7
+ restart: always
+ volumes:
+ - redis-data:/data
+ networks:
+ - app_network
+ ports:
+ - 6379:6379
+
+ redisinsight:
+ image: redislabs/redisinsight:latest
+ ports:
+ - "8001:8001"
+ volumes:
+ - redis-data:/redisinsight
diff --git a/docs/_snippets/jobs/options.mdx b/docs/_snippets/jobs/options.mdx
index 49298360c..31a60f3aa 100644
--- a/docs/_snippets/jobs/options.mdx
+++ b/docs/_snippets/jobs/options.mdx
@@ -36,6 +36,9 @@
The `enabled` property is an optional property that specifies whether the Job is enabled or not. The Job will be enabled by default if you omit this property. When a job is disabled, no new runs will be triggered or resumed. In progress runs will continue to run until they are finished or delayed by using `io.wait`.
+
+ The `concurrencyLimit` property is an optional property that specifies the maximum number of concurrent run executions. If this property is omitted, the job can potentially use up the full concurrency of an environment. You can also create a limit on a group of jobs by defining a [ConcurrencyLimit](/sdk/triggerclient/instancemethods/concurrency-limit) object.
+
The `onSuccess` property is an optional property that specifies a callback function to run when the Job finishes successfully. The callback function receives a [Run Notification](/sdk/run-notification) object as it's only parameter.
diff --git a/docs/documentation/concepts/limits.mdx b/docs/documentation/concepts/limits.mdx
index cb06a5ccd..c7b644564 100644
--- a/docs/documentation/concepts/limits.mdx
+++ b/docs/documentation/concepts/limits.mdx
@@ -16,7 +16,7 @@ The following limits apply to the Trigger.dev Cloud service and users of the sel
| Connected Integrations | Up to 50 | Up to 1000 | Custom |
| Task Output Size | 3MB | 3MB | 3MB |
| [Tasks per Run](#tasks-per-runs) | Up to 250 | Up to 1000 | Custom |
-| [Concurrent Run Executions](#concurrent-run-executions) | Up to 10 | Up to 10 | Custom |
+| [Concurrent Run Executions per Environment](#concurrent-run-executions) | Up to 10 | Up to 100 | Custom |
| [Maximum Task Duration](#maximum-task-duration) | < 2m | < 2m | < Deployment Grace Period |
| [Maximum Run Execution Duration](#maximum-total-run-execution-duration) | up to 15m | up to 2 hrs | Custom |
| [Yielded Executions per Run](#yielded-executions-per-run) | Up to 100 | Up to 100 | Custom |
@@ -88,6 +88,51 @@ This does not include runs that are waiting for a [io.wait()](/sdk/io/wait) to c
Going over this limit does not abort or cancel runs, but it will prevent new run executions until the number of concurrent executions drops below the limit.
+You can limit the execution concurrency of a specific job like so:
+
+```ts
+client.defineJob({
+ id: `test-job-1`,
+ name: `Test Job 1`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "test",
+ }),
+ concurrencyLimit: 5, // Limit this job to 5 concurrent executions
+});
+```
+
+Alternatively, you can limit a group of jobs concurrency limit by defining a concurrency limit and passing it to the `defineJobs()` method:
+
+```ts
+const concurrencyLimit = client.defineConcurrencyLimit({
+ id: `test-shared`,
+ limit: 5, // Limit all jobs in this group to 5 concurrent executions
+});
+
+client.defineJob({
+ id: `test-job-1`,
+ name: `Test Job 1`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "test",
+ }),
+ concurrencyLimit,
+});
+
+client.defineJob({
+ id: `test-job-2`,
+ name: `Test Job 2`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "test",
+ }),
+ concurrencyLimit,
+});
+```
+
+The two jobs above will share the same concurrency limit, so between them they can only have 5 concurrent executions.
+
### Maximum Task Duration
The Maximum Task Duration is the maximum amount of time a single Task can run for. This limit is partly enforced by the Trigger.dev server, but also by the execution runtime of your deployed serverless function.
diff --git a/docs/documentation/concepts/runs.mdx b/docs/documentation/concepts/runs.mdx
index e65ebb870..95c90c001 100644
--- a/docs/documentation/concepts/runs.mdx
+++ b/docs/documentation/concepts/runs.mdx
@@ -24,7 +24,7 @@ client.defineJob({
run: async (payload, io, ctx) => {
// 2. Regular code and Tasks
// 3. Optionally return data from run execution
- return { status: 'success' }
+ return { status: "success" };
},
});
```
@@ -64,6 +64,52 @@ A few things you can do with `io`:
The `context` object gives you access to information about the current Run, Job, Environment, Organization and Event. [View the full reference](/sdk/context) for `context`.
+## Run Statuses
+
+### Pending
+
+The run has been created but has not started yet. This is the initial status of a run.
+
+### Queued
+
+The run is waiting to be executed. Runs can be queued because of [Run Execution Concurrency Limits](/documentation/concepts/limits#concurrent-run-executions)
+
+### Waiting on Connections
+
+If a run depends on a hosted integration, it will be in this status until the integration is ready.
+
+### Executing
+
+The run is currently executing. This means that the run function is running.
+
+### Waiting
+
+The run is waiting, either because of a call to `io.wait()` or because a task failed and will be retried at some point in the future. Runs in this state don't count towards concurrency limits.
+
+### Failed
+
+The run failed. This can happen if the run function throws an error or if a task fails and the run is not configured to retry.
+
+### Completed
+
+The run completed successfully. This means that the run function finished executing and all tasks completed successfully.
+
+### Cancelled
+
+The run was cancelled. This can happen if the run is cancelled manually.
+
+### Timed Out
+
+The run timed out. This can happen if the run exceeds the maximum run duration, or if we receive a serverless function execution timed out response when hitting your endpoint repeatedly with no new task creation.
+
+### Invalid Payload
+
+The run failed because the payload was invalid.
+
+### Unresolved Auth
+
+The run failed because the auth data could not be resolved when using a custom Auth Resolver.
+
## References
diff --git a/docs/mint.json b/docs/mint.json
index cc99e1ade..b48bd8080 100644
--- a/docs/mint.json
+++ b/docs/mint.json
@@ -1,7 +1,9 @@
{
"$schema": "https://mintlify.com/schema.json",
"name": "Trigger.dev",
- "openapi": ["/openapi.yml"],
+ "openapi": [
+ "/openapi.yml"
+ ],
"logo": {
"dark": "/logo/dark.png",
"light": "/logo/light.png",
@@ -253,7 +255,10 @@
"pages": [
{
"group": "Airtable",
- "pages": ["integrations/apis/airtable", "integrations/apis/airtable-tasks"]
+ "pages": [
+ "integrations/apis/airtable",
+ "integrations/apis/airtable-tasks"
+ ]
},
{
"group": "GitHub",
@@ -279,16 +284,25 @@
},
{
"group": "Plain",
- "pages": ["integrations/apis/plain", "integrations/apis/plain-tasks"]
+ "pages": [
+ "integrations/apis/plain",
+ "integrations/apis/plain-tasks"
+ ]
},
"integrations/apis/replicate",
{
"group": "SendGrid",
- "pages": ["integrations/apis/sendgrid", "integrations/apis/sendgrid-tasks"]
+ "pages": [
+ "integrations/apis/sendgrid",
+ "integrations/apis/sendgrid-tasks"
+ ]
},
{
"group": "Resend",
- "pages": ["integrations/apis/resend", "integrations/apis/resend-tasks"]
+ "pages": [
+ "integrations/apis/resend",
+ "integrations/apis/resend-tasks"
+ ]
},
{
"group": "Shopify",
@@ -300,7 +314,10 @@
},
{
"group": "Slack",
- "pages": ["integrations/apis/slack", "integrations/apis/slack-tasks"]
+ "pages": [
+ "integrations/apis/slack",
+ "integrations/apis/slack-tasks"
+ ]
},
"integrations/apis/stripe",
{
@@ -344,6 +361,7 @@
"sdk/triggerclient/instancemethods/define-dynamic-trigger",
"sdk/triggerclient/instancemethods/define-dynamic-schedule",
"sdk/triggerclient/instancemethods/define-auth-resolver",
+ "sdk/triggerclient/instancemethods/concurrency-limit",
"sdk/triggerclient/instancemethods/on"
]
}
@@ -388,7 +406,10 @@
"sdk/dynamictrigger/constructor",
{
"group": "Instance methods",
- "pages": ["sdk/dynamictrigger/register", "sdk/dynamictrigger/unregister"]
+ "pages": [
+ "sdk/dynamictrigger/register",
+ "sdk/dynamictrigger/unregister"
+ ]
}
]
},
@@ -399,7 +420,10 @@
"sdk/dynamicschedule/constructor",
{
"group": "Instance methods",
- "pages": ["sdk/dynamicschedule/register", "sdk/dynamicschedule/unregister"]
+ "pages": [
+ "sdk/dynamicschedule/register",
+ "sdk/dynamicschedule/unregister"
+ ]
}
]
},
@@ -411,7 +435,9 @@
},
{
"group": "HTTP Reference",
- "pages": ["sdk/api-reference/events/create-an-event"]
+ "pages": [
+ "sdk/api-reference/events/create-an-event"
+ ]
},
{
"group": "React SDK",
@@ -425,7 +451,9 @@
},
{
"group": "Overview",
- "pages": ["examples/introduction"]
+ "pages": [
+ "examples/introduction"
+ ]
}
],
"footerSocials": {
@@ -438,4 +466,4 @@
"apiKey": "phc_hwYmedO564b3Ik8nhA4Csrb5SueY0EwFJWCbseGwWW"
}
}
-}
+}
\ No newline at end of file
diff --git a/docs/sdk/triggerclient/instancemethods/concurrency-limit.mdx b/docs/sdk/triggerclient/instancemethods/concurrency-limit.mdx
new file mode 100644
index 000000000..48c8c1a17
--- /dev/null
+++ b/docs/sdk/triggerclient/instancemethods/concurrency-limit.mdx
@@ -0,0 +1,47 @@
+---
+title: "defineConcurrencyLimit()"
+description: "Define a concurrency limit group to control the concurrency of your jobs."
+---
+
+You can control the concurrency of run executions for a group of jobs using a concurrency limit group.
+
+
+
+```ts example
+const concurrencyLimit = client.defineConcurrencyLimit({
+ id: `test-shared`,
+ limit: 5, // Limit all jobs in this group to 5 concurrent executions
+});
+
+client.defineJob({
+ id: `test-job-1`,
+ name: `Test Job 1`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "test",
+ }),
+ concurrencyLimit,
+});
+
+client.defineJob({
+ id: `test-job-2`,
+ name: `Test Job 2`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "test",
+ }),
+ concurrencyLimit,
+});
+```
+
+
+
+## Parameters
+
+
+ The ID of the concurrency limit group.
+
+
+
+ The maximum number of concurrent executions allowed for this group.
+
diff --git a/packages/core/src/schemas/api.ts b/packages/core/src/schemas/api.ts
index c09410b33..9a7882eae 100644
--- a/packages/core/src/schemas/api.ts
+++ b/packages/core/src/schemas/api.ts
@@ -273,6 +273,11 @@ export const QueueOptionsSchema = z.object({
export type QueueOptions = z.infer;
+export const ConcurrencyLimitOptionsSchema = z.object({
+ id: z.string(),
+ limit: z.number(),
+});
+
export const JobMetadataSchema = z.object({
id: z.string(),
name: z.string(),
@@ -284,6 +289,7 @@ export const JobMetadataSchema = z.object({
enabled: z.boolean(),
startPosition: z.enum(["initial", "latest"]),
preprocessRuns: z.boolean(),
+ concurrencyLimit: ConcurrencyLimitOptionsSchema.or(z.number().int().positive()).optional(),
});
export type JobMetadata = z.infer;
@@ -879,7 +885,6 @@ export const RunTaskOptionsSchema = z.object({
/** A No Operation means that the code won't be executed. This is used internally to implement features like [io.wait()](https://trigger.dev/docs/sdk/io/wait). */
noop: z.boolean().default(false),
redact: RedactSchema.optional(),
- trigger: TriggerMetadataSchema.optional(),
parallel: z.boolean().optional(),
});
diff --git a/packages/core/src/schemas/runs.ts b/packages/core/src/schemas/runs.ts
index e62c1331b..8fc448b68 100644
--- a/packages/core/src/schemas/runs.ts
+++ b/packages/core/src/schemas/runs.ts
@@ -18,6 +18,9 @@ export const RunStatusSchema = z.union([
z.literal("CANCELED"),
z.literal("UNRESOLVED_AUTH"),
z.literal("INVALID_PAYLOAD"),
+ z.literal("EXECUTING"),
+ z.literal("WAITING_TO_CONTINUE"),
+ z.literal("WAITING_TO_EXECUTE"),
]);
export const RunTaskSchema = z.object({
diff --git a/packages/database/prisma/migrations/20231117145312_add_additional_run_statuses/migration.sql b/packages/database/prisma/migrations/20231117145312_add_additional_run_statuses/migration.sql
new file mode 100644
index 000000000..98fe1903c
--- /dev/null
+++ b/packages/database/prisma/migrations/20231117145312_add_additional_run_statuses/migration.sql
@@ -0,0 +1,11 @@
+-- AlterEnum
+-- This migration adds more than one value to an enum.
+-- With PostgreSQL versions 11 and earlier, this is not possible
+-- in a single migration. This can be worked around by creating
+-- multiple migrations, each migration adding only one value to
+-- the enum.
+
+
+ALTER TYPE "JobRunStatus" ADD VALUE 'EXECUTING';
+ALTER TYPE "JobRunStatus" ADD VALUE 'WAITING_TO_CONTINUE';
+ALTER TYPE "JobRunStatus" ADD VALUE 'WAITING_TO_EXECUTE';
diff --git a/packages/database/prisma/migrations/20231121144353_make_job_run_number_optional/migration.sql b/packages/database/prisma/migrations/20231121144353_make_job_run_number_optional/migration.sql
new file mode 100644
index 000000000..f8979aaa8
--- /dev/null
+++ b/packages/database/prisma/migrations/20231121144353_make_job_run_number_optional/migration.sql
@@ -0,0 +1,2 @@
+-- AlterTable
+ALTER TABLE "JobRun" ALTER COLUMN "number" DROP NOT NULL;
diff --git a/packages/database/prisma/migrations/20231121154359_add_job_counter_table/migration.sql b/packages/database/prisma/migrations/20231121154359_add_job_counter_table/migration.sql
new file mode 100644
index 000000000..6f7b36264
--- /dev/null
+++ b/packages/database/prisma/migrations/20231121154359_add_job_counter_table/migration.sql
@@ -0,0 +1,7 @@
+-- CreateTable
+CREATE TABLE "JobCounter" (
+ "jobId" TEXT NOT NULL,
+ "lastNumber" INTEGER NOT NULL DEFAULT 0,
+
+ CONSTRAINT "JobCounter_pkey" PRIMARY KEY ("jobId")
+);
diff --git a/packages/database/prisma/migrations/20231121154545_seed_job_counter_tables/migration.sql b/packages/database/prisma/migrations/20231121154545_seed_job_counter_tables/migration.sql
new file mode 100644
index 000000000..8bc544ef5
--- /dev/null
+++ b/packages/database/prisma/migrations/20231121154545_seed_job_counter_tables/migration.sql
@@ -0,0 +1,10 @@
+-- This is an empty migration.
+INSERT INTO
+ "JobCounter" ("jobId", "lastNumber")
+SELECT
+ "jobId",
+ MAX(number)
+FROM
+ "JobRun"
+GROUP BY
+ "jobId";
\ No newline at end of file
diff --git a/packages/database/prisma/migrations/20231122210707_add_concurrency_limit_tables_and_columns/migration.sql b/packages/database/prisma/migrations/20231122210707_add_concurrency_limit_tables_and_columns/migration.sql
new file mode 100644
index 000000000..b04d0eefe
--- /dev/null
+++ b/packages/database/prisma/migrations/20231122210707_add_concurrency_limit_tables_and_columns/migration.sql
@@ -0,0 +1,30 @@
+-- AlterTable
+ALTER TABLE "JobRun" ADD COLUMN "concurrencyLimitGroupId" TEXT;
+
+-- AlterTable
+ALTER TABLE "JobVersion" ADD COLUMN "concurrencyLimit" INTEGER,
+ADD COLUMN "concurrencyLimitGroupId" TEXT;
+
+-- CreateTable
+CREATE TABLE "ConcurrencyLimitGroup" (
+ "id" TEXT NOT NULL,
+ "name" TEXT NOT NULL,
+ "concurrencyLimit" INTEGER NOT NULL,
+ "environmentId" TEXT NOT NULL,
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "updatedAt" TIMESTAMP(3) NOT NULL,
+
+ CONSTRAINT "ConcurrencyLimitGroup_pkey" PRIMARY KEY ("id")
+);
+
+-- CreateIndex
+CREATE UNIQUE INDEX "ConcurrencyLimitGroup_environmentId_name_key" ON "ConcurrencyLimitGroup"("environmentId", "name");
+
+-- AddForeignKey
+ALTER TABLE "JobVersion" ADD CONSTRAINT "JobVersion_concurrencyLimitGroupId_fkey" FOREIGN KEY ("concurrencyLimitGroupId") REFERENCES "ConcurrencyLimitGroup"("id") ON DELETE SET NULL ON UPDATE CASCADE;
+
+-- AddForeignKey
+ALTER TABLE "ConcurrencyLimitGroup" ADD CONSTRAINT "ConcurrencyLimitGroup_environmentId_fkey" FOREIGN KEY ("environmentId") REFERENCES "RuntimeEnvironment"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+
+-- AddForeignKey
+ALTER TABLE "JobRun" ADD CONSTRAINT "JobRun_concurrencyLimitGroupId_fkey" FOREIGN KEY ("concurrencyLimitGroupId") REFERENCES "ConcurrencyLimitGroup"("id") ON DELETE SET NULL ON UPDATE CASCADE;
diff --git a/packages/database/prisma/migrations/20231122212600_make_job_queues_optional/migration.sql b/packages/database/prisma/migrations/20231122212600_make_job_queues_optional/migration.sql
new file mode 100644
index 000000000..b59219304
--- /dev/null
+++ b/packages/database/prisma/migrations/20231122212600_make_job_queues_optional/migration.sql
@@ -0,0 +1,17 @@
+-- DropForeignKey
+ALTER TABLE "JobRun" DROP CONSTRAINT "JobRun_queueId_fkey";
+
+-- DropForeignKey
+ALTER TABLE "JobVersion" DROP CONSTRAINT "JobVersion_queueId_fkey";
+
+-- AlterTable
+ALTER TABLE "JobRun" ALTER COLUMN "queueId" DROP NOT NULL;
+
+-- AlterTable
+ALTER TABLE "JobVersion" ALTER COLUMN "queueId" DROP NOT NULL;
+
+-- AddForeignKey
+ALTER TABLE "JobVersion" ADD CONSTRAINT "JobVersion_queueId_fkey" FOREIGN KEY ("queueId") REFERENCES "JobQueue"("id") ON DELETE SET NULL ON UPDATE CASCADE;
+
+-- AddForeignKey
+ALTER TABLE "JobRun" ADD CONSTRAINT "JobRun_queueId_fkey" FOREIGN KEY ("queueId") REFERENCES "JobQueue"("id") ON DELETE SET NULL ON UPDATE CASCADE;
diff --git a/packages/database/prisma/migrations/20231123113308_remove_concurrency_group_from_run/migration.sql b/packages/database/prisma/migrations/20231123113308_remove_concurrency_group_from_run/migration.sql
new file mode 100644
index 000000000..9bb2b98b7
--- /dev/null
+++ b/packages/database/prisma/migrations/20231123113308_remove_concurrency_group_from_run/migration.sql
@@ -0,0 +1,11 @@
+/*
+ Warnings:
+
+ - You are about to drop the column `concurrencyLimitGroupId` on the `JobRun` table. All the data in the column will be lost.
+
+*/
+-- DropForeignKey
+ALTER TABLE "JobRun" DROP CONSTRAINT "JobRun_concurrencyLimitGroupId_fkey";
+
+-- AlterTable
+ALTER TABLE "JobRun" DROP COLUMN "concurrencyLimitGroupId";
diff --git a/packages/database/prisma/migrations/20231123115015_add_concurrency_limit_group_id_to_run_executions/migration.sql b/packages/database/prisma/migrations/20231123115015_add_concurrency_limit_group_id_to_run_executions/migration.sql
new file mode 100644
index 000000000..4529052e9
--- /dev/null
+++ b/packages/database/prisma/migrations/20231123115015_add_concurrency_limit_group_id_to_run_executions/migration.sql
@@ -0,0 +1,4 @@
+ALTER TABLE
+ "triggerdotdev_events"."run_executions"
+ADD
+ COLUMN "concurrency_limit_group_id" text;
\ No newline at end of file
diff --git a/packages/database/prisma/schema.prisma b/packages/database/prisma/schema.prisma
index 810e76e14..a3ce71e17 100644
--- a/packages/database/prisma/schema.prisma
+++ b/packages/database/prisma/schema.prisma
@@ -328,6 +328,7 @@ model RuntimeEnvironment {
scheduleSources ScheduleSource[]
ExternalAccount ExternalAccount[]
httpEndpointEnvironments TriggerHttpEndpointEnvironment[]
+ concurrencyLimitGroups ConcurrencyLimitGroup[]
keyValueItems KeyValueItem[]
webhookEnvironments WebhookEnvironment[]
webhookRequestDeliveries WebhookRequestDelivery[]
@@ -488,12 +489,16 @@ model JobVersion {
project Project @relation(fields: [projectId], references: [id], onDelete: Cascade, onUpdate: Cascade)
projectId String
- queue JobQueue @relation(fields: [queueId], references: [id])
- queueId String
+ queue JobQueue? @relation(fields: [queueId], references: [id])
+ queueId String?
startPosition JobStartPosition @default(INITIAL)
preprocessRuns Boolean @default(false)
+ concurrencyLimit Int?
+ concurrencyLimitGroup ConcurrencyLimitGroup? @relation(fields: [concurrencyLimitGroupId], references: [id])
+ concurrencyLimitGroupId String?
+
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
@@ -532,6 +537,23 @@ model EventExample {
@@unique([slug, jobVersionId])
}
+model ConcurrencyLimitGroup {
+ id String @id @default(cuid())
+ name String
+
+ concurrencyLimit Int
+
+ environment RuntimeEnvironment @relation(fields: [environmentId], references: [id], onDelete: Cascade, onUpdate: Cascade)
+ environmentId String
+
+ createdAt DateTime @default(now())
+ updatedAt DateTime @updatedAt
+
+ jobVersion JobVersion[]
+
+ @@unique([environmentId, name])
+}
+
model JobQueue {
id String @id @default(cuid())
name String
@@ -718,7 +740,7 @@ enum PayloadType {
model JobRun {
id String @id @default(cuid())
- number Int
+ number Int?
internal Boolean @default(false)
job Job @relation(fields: [jobId], references: [id], onDelete: Cascade, onUpdate: Cascade)
@@ -742,8 +764,8 @@ model JobRun {
project Project @relation(fields: [projectId], references: [id], onDelete: Cascade, onUpdate: Cascade)
projectId String
- queue JobQueue @relation(fields: [queueId], references: [id])
- queueId String
+ queue JobQueue? @relation(fields: [queueId], references: [id])
+ queueId String?
externalAccount ExternalAccount? @relation(fields: [externalAccountId], references: [id], onDelete: Cascade, onUpdate: Cascade)
externalAccountId String?
@@ -787,6 +809,9 @@ enum JobRunStatus {
WAITING_ON_CONNECTIONS
PREPROCESSING
STARTED
+ EXECUTING
+ WAITING_TO_CONTINUE
+ WAITING_TO_EXECUTE
SUCCESS
FAILURE
TIMED_OUT
@@ -796,6 +821,11 @@ enum JobRunStatus {
INVALID_PAYLOAD
}
+model JobCounter {
+ jobId String @id
+ lastNumber Int @default(0)
+}
+
model JobRunAutoYieldExecution {
id String @id @default(cuid())
diff --git a/packages/trigger-sdk/src/concurrencyLimit.ts b/packages/trigger-sdk/src/concurrencyLimit.ts
new file mode 100644
index 000000000..98fdba209
--- /dev/null
+++ b/packages/trigger-sdk/src/concurrencyLimit.ts
@@ -0,0 +1,16 @@
+export type ConcurrencyLimitOptions = {
+ id: string;
+ limit: number;
+};
+
+export class ConcurrencyLimit {
+ constructor(private options: ConcurrencyLimitOptions) {}
+
+ get id() {
+ return this.options.id;
+ }
+
+ get limit() {
+ return this.options.limit;
+ }
+}
diff --git a/packages/trigger-sdk/src/job.ts b/packages/trigger-sdk/src/job.ts
index b6336a376..da5f7423b 100644
--- a/packages/trigger-sdk/src/job.ts
+++ b/packages/trigger-sdk/src/job.ts
@@ -20,6 +20,7 @@ import type {
import { slugifyId } from "./utils";
import { runLocalStorage } from "./runLocalStorage";
import { Prettify } from "@trigger.dev/core";
+import { ConcurrencyLimit } from "./concurrencyLimit";
export type JobOptions<
TTrigger extends Trigger>,
@@ -60,9 +61,16 @@ export type JobOptions<
});
``` */
integrations?: TIntegrations;
- /** @deprecated This property is deprecated and no longer effects the execution of the Job
- * */
- queue?: QueueOptions | string;
+
+ /**
+ * The `concurrencyLimit` property is used to limit the number of concurrent run executions of a job.
+ * Can be a number which represents the limit or a `ConcurrencyLimit` instance which can be used to
+ * group together multiple jobs to share the same concurrency limit.
+ *
+ * If undefined the job will be limited only by the server's global concurrency limit, or if you are using the
+ * Trigger.dev Cloud service, the concurrency limit of your plan.
+ */
+ concurrencyLimit?: number | ConcurrencyLimit;
/** The `enabled` property is used to enable or disable the Job. If you disable a Job, it will not run. */
enabled?: boolean;
/** This function gets called automatically when a Run is Triggered.
@@ -174,6 +182,12 @@ export class Job<
enabled: this.enabled,
preprocessRuns: this.trigger.preprocessRuns,
internal,
+ concurrencyLimit:
+ typeof this.options.concurrencyLimit === "number"
+ ? this.options.concurrencyLimit
+ : typeof this.options.concurrencyLimit === "object"
+ ? { id: this.options.concurrencyLimit.id, limit: this.options.concurrencyLimit.limit }
+ : undefined,
};
}
diff --git a/packages/trigger-sdk/src/triggerClient.ts b/packages/trigger-sdk/src/triggerClient.ts
index 307cb8539..a64897e75 100644
--- a/packages/trigger-sdk/src/triggerClient.ts
+++ b/packages/trigger-sdk/src/triggerClient.ts
@@ -118,6 +118,7 @@ const registerSourceEvent: EventSpecification = {
import EventEmitter from "node:events";
import * as packageJson from "../package.json";
+import { ConcurrencyLimit, ConcurrencyLimitOptions } from "./concurrencyLimit";
import { formatSchemaErrors } from "./utils/formatSchemaErrors";
import { WebhookDeliveryContext, WebhookSource } from "./triggers/webhook";
import { KeyValueStore } from "./store/keyValueStore";
@@ -742,6 +743,10 @@ export class TriggerClient {
return endpoint;
}
+ defineConcurrencyLimit(options: ConcurrencyLimitOptions) {
+ return new ConcurrencyLimit(options);
+ }
+
attach(job: Job, any>): void {
this.#registeredJobs[job.id] = job;
job.trigger.attachToJob(this, job);
@@ -1788,6 +1793,12 @@ export class TriggerClient {
enabled: job.enabled,
preprocessRuns: job.trigger.preprocessRuns,
internal,
+ concurrencyLimit:
+ typeof job.options.concurrencyLimit === "number"
+ ? job.options.concurrencyLimit
+ : typeof job.options.concurrencyLimit === "object"
+ ? { id: job.options.concurrencyLimit.id, limit: job.options.concurrencyLimit.limit }
+ : undefined,
};
}
diff --git a/perf/src/index.ts b/perf/src/index.ts
index c3edeacad..4083f9d94 100644
--- a/perf/src/index.ts
+++ b/perf/src/index.ts
@@ -108,8 +108,8 @@ async function mainParallel() {
async function mainParallelBulk() {
const batches = 1;
- const concurrency = 50;
- const eventsPer = 20;
+ const concurrency = 10;
+ const eventsPer = 10;
console.log("Preparing perf tests...");
@@ -169,7 +169,34 @@ async function mainSerial() {
}
}
-mainParallelBulk().catch((err) => {
+async function mainConcurrency() {
+ const batches = 1;
+ const concurrency = 10;
+ const eventsPer = 5;
+
+ console.log("Preparing perf tests...");
+
+ await new Promise((resolve) => setTimeout(resolve, 5000));
+
+ console.log("Starting perf tests in 1 second...");
+
+ // wait for 1 seconds
+ await new Promise((resolve) => setTimeout(resolve, 1000));
+
+ // Send 5 events per second for 30 seconds (1 event == 10 runs)
+ for (let i = 0; i < batches; i++) {
+ console.log(`Sending ${concurrency} x ${eventsPer} events... batch ${i + 1}/${batches}`);
+ await Promise.all(new Array(concurrency).fill(0).map(() => sendEvents(eventsPer)));
+
+ await new Promise((resolve) => setTimeout(resolve, 250));
+ }
+}
+
+async function mainSingle() {
+ await sendEvent();
+}
+
+mainConcurrency().catch((err) => {
console.error(err);
process.exit(1);
});
diff --git a/perf/src/trigger.ts b/perf/src/trigger.ts
index 6f8006fa5..3d79eb07f 100644
--- a/perf/src/trigger.ts
+++ b/perf/src/trigger.ts
@@ -6,6 +6,11 @@ export const triggerClient = new TriggerClient({
apiUrl: process.env.TRIGGER_API_URL!,
});
+const concurrencyLimit = triggerClient.defineConcurrencyLimit({
+ id: `perf-test-shared`,
+ limit: 5,
+});
+
triggerClient.defineJob({
id: `perf-test-1`,
name: `Perf Test 1`,
@@ -13,11 +18,12 @@ triggerClient.defineJob({
trigger: eventTrigger({
name: "perf.test",
}),
+ concurrencyLimit,
run: async (payload, io, ctx) => {
await io.runTask(
"task-1",
async (task) => {
- await new Promise((resolve) => setTimeout(resolve, 2000));
+ await new Promise((resolve) => setTimeout(resolve, 5000));
return {
value: Math.random(),
@@ -26,6 +32,8 @@ triggerClient.defineJob({
{ name: "task 1" }
);
+ await io.wait("wait", 10);
+
await io.runTask(
"task-2",
async (task) => {
@@ -35,5 +43,111 @@ triggerClient.defineJob({
},
{ name: "task 2" }
);
+
+ await io.runTask(
+ "task-3",
+ async (task) => {
+ await new Promise((resolve) => setTimeout(resolve, 2000));
+
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 3" }
+ );
+ },
+});
+
+triggerClient.defineJob({
+ id: `perf-test-2`,
+ name: `Perf Test 2`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "perf.test",
+ }),
+ concurrencyLimit: 5,
+ run: async (payload, io, ctx) => {
+ await io.runTask(
+ "task-1",
+ async (task) => {
+ await new Promise((resolve) => setTimeout(resolve, 5000));
+
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 1" }
+ );
+
+ await io.wait("wait", 10);
+
+ await io.runTask(
+ "task-2",
+ async (task) => {
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 2" }
+ );
+
+ await io.runTask(
+ "task-3",
+ async (task) => {
+ await new Promise((resolve) => setTimeout(resolve, 2000));
+
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 3" }
+ );
+ },
+});
+
+triggerClient.defineJob({
+ id: `perf-test-3`,
+ name: `Perf Test 3`,
+ version: "1.0.0",
+ trigger: eventTrigger({
+ name: "perf.test",
+ }),
+ concurrencyLimit,
+ run: async (payload, io, ctx) => {
+ await io.runTask(
+ "task-1",
+ async (task) => {
+ await new Promise((resolve) => setTimeout(resolve, 5000));
+
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 1" }
+ );
+
+ await io.wait("wait", 10);
+
+ await io.runTask(
+ "task-2",
+ async (task) => {
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 2" }
+ );
+
+ await io.runTask(
+ "task-3",
+ async (task) => {
+ await new Promise((resolve) => setTimeout(resolve, 2000));
+
+ return {
+ value: Math.random(),
+ };
+ },
+ { name: "task 3" }
+ );
},
});
diff --git a/perf/tsconfig.json b/perf/tsconfig.json
index 67f0df985..d0d1f2687 100644
--- a/perf/tsconfig.json
+++ b/perf/tsconfig.json
@@ -12,6 +12,8 @@
"@trigger.dev/express/*": ["../packages/express/src/*"],
"@trigger.dev/core": ["../packages/core/src/index"],
"@trigger.dev/core/*": ["../packages/core/src/*"],
+ "@trigger.dev/core-backend": ["../packages/core-backend/src/index"],
+ "@trigger.dev/core-backend/*": ["../packages/core-backend/src/*"],
"@trigger.dev/integration-kit": ["../packages/integration-kit/src/index"],
"@trigger.dev/integration-kit/*": ["../packages/integration-kit/src/*"],
"@trigger.dev/github": ["../integrations/github/src/index"],
diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml
index a736f021a..f37df3fd5 100644
--- a/pnpm-lock.yaml
+++ b/pnpm-lock.yaml
@@ -175,6 +175,7 @@ importers:
highlight.run: ^7.3.4
humanize-duration: ^3.27.3
intl-parse-accept-language: ^1.0.0
+ ioredis: ^5.3.2
isbot: ^3.6.5
jsonpointer: ^5.0.1
lodash.omit: ^4.5.0
@@ -283,6 +284,7 @@ importers:
highlight.run: 7.3.4
humanize-duration: 3.27.3
intl-parse-accept-language: 1.0.0
+ ioredis: 5.3.2
isbot: 3.6.5
jsonpointer: 5.0.1
lodash.omit: 4.5.0
@@ -7922,6 +7924,10 @@ packages:
/@humanwhocodes/object-schema/1.2.1:
resolution: {integrity: sha512-ZnQMnLV4e7hDlUvw8H+U8ASL02SS2Gn6+9Ac3wGGLIe7+je2AeAOxPY+izIPJDfFDb7eDjev0Us8MO1iFRN8hA==}
+ /@ioredis/commands/1.2.0:
+ resolution: {integrity: sha512-Sx1pU8EM64o2BrqNpEO1CNLtKQwyhuXuqyfH7oGKCk+1a33d2r5saW8zNwm3j6BTExtjrv2BxTgzzkMwts6vGg==}
+ dev: false
+
/@isaacs/cliui/8.0.2:
resolution: {integrity: sha512-O8jcjabXaleOG9DQ0+ARXWZBTfnP4WNAqzuiJK7ll44AmxGKv/J2M4TPjxjY3znBCfvBXFzucm1twdyFybFqEA==}
engines: {node: '>=12'}
@@ -18013,6 +18019,11 @@ packages:
resolution: {integrity: sha512-rQ1+kcj+ttHG0MKVGBUXwayCCF1oh39BF5COIpRzuCEv8Mwjv0XucrI2ExNTOn9IlLifGClWQcU9BrZORvtw6Q==}
engines: {node: '>=6'}
+ /cluster-key-slot/1.1.2:
+ resolution: {integrity: sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==}
+ engines: {node: '>=0.10.0'}
+ dev: false
+
/co/4.6.0:
resolution: {integrity: sha512-QVb0dM5HvG+uaxitm8wONl7jltx8dqhfU33DcqtOZcLSVIKSDDLDi7+0LbAKiyI8hD9u42m2YxXSkMGWThaecQ==}
engines: {iojs: '>= 1.0.0', node: '>= 0.12.0'}
@@ -18953,6 +18964,11 @@ packages:
/delegates/1.0.0:
resolution: {integrity: sha512-bd2L678uiWATM6m5Z1VzNCErI3jiGzt6HGY8OVICs40JQq/HALfbyNJmp0UDakEY4pMMaN0Ly5om/B1VI/+xfQ==}
+ /denque/2.1.0:
+ resolution: {integrity: sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw==}
+ engines: {node: '>=0.10'}
+ dev: false
+
/depd/2.0.0:
resolution: {integrity: sha512-g7nH6P6dyDioJogAAGprGpCtVImJhpPk/roCzdb3fIh61/s/nPsfR6onyMwkCAR/OlC3yBC0lESvUoQEAssIrw==}
engines: {node: '>= 0.8'}
@@ -23265,6 +23281,23 @@ packages:
loose-envify: 1.4.0
dev: false
+ /ioredis/5.3.2:
+ resolution: {integrity: sha512-1DKMMzlIHM02eBBVOFQ1+AolGjs6+xEcM4PDL7NqOS6szq7H9jSaEkIUH6/a5Hl241LzW6JLSiAbNvTQjUupUA==}
+ engines: {node: '>=12.22.0'}
+ dependencies:
+ '@ioredis/commands': 1.2.0
+ cluster-key-slot: 1.1.2
+ debug: 4.3.4
+ denque: 2.1.0
+ lodash.defaults: 4.2.0
+ lodash.isarguments: 3.1.0
+ redis-errors: 1.2.0
+ redis-parser: 3.0.0
+ standard-as-callback: 2.1.0
+ transitivePeerDependencies:
+ - supports-color
+ dev: false
+
/ip/1.1.8:
resolution: {integrity: sha512-PuExPYUiu6qMBQb4l06ecm6T6ujzhmh+MeJcW9wa89PoAz5pvd4zPgN5WJV104mb6S2T1AwNIAaB70JNrLQWhg==}
@@ -25035,6 +25068,14 @@ packages:
resolution: {integrity: sha512-FT1yDzDYEoYWhnSGnpE/4Kj1fLZkDFyqRb7fNt6FdYOSxlUWAtp42Eh6Wb0rGIv/m9Bgo7x4GhQbm5Ys4SG5ow==}
dev: true
+ /lodash.defaults/4.2.0:
+ resolution: {integrity: sha512-qjxPLHd3r5DnsdGacqOMU6pb/avJzdh9tFX2ymgoZE27BmjXrNy/y4LoaiTeAb+O3gL8AfpJGtqfX/ae2leYYQ==}
+ dev: false
+
+ /lodash.isarguments/3.1.0:
+ resolution: {integrity: sha512-chi4NHZlZqZD18a0imDHnZPrDeBbTtVN7GXMwuGdRH9qotxAjYs3aVLKc7zNOG9eddR5Ksd8rvFEBc9SsggPpg==}
+ dev: false
+
/lodash.isplainobject/4.0.6:
resolution: {integrity: sha512-oSXzaWypCMHkPC3NvBEaPHf0KsA5mvPrOPgQWDsbg8n7orZ290M0BmC/jgRZ4vcJ6DTAhjrsSYgdsW/F+MFOBA==}
dev: true
@@ -29249,6 +29290,18 @@ packages:
strip-indent: 3.0.0
dev: false
+ /redis-errors/1.2.0:
+ resolution: {integrity: sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==}
+ engines: {node: '>=4'}
+ dev: false
+
+ /redis-parser/3.0.0:
+ resolution: {integrity: sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==}
+ engines: {node: '>=4'}
+ dependencies:
+ redis-errors: 1.2.0
+ dev: false
+
/reduce-css-calc/2.1.8:
resolution: {integrity: sha512-8liAVezDmUcH+tdzoEGrhfbGcP7nOV4NkGE3a74+qqvE7nt9i4sKLGBuZNOnpI4WiGksiNPklZxva80061QiPg==}
dependencies:
@@ -30773,6 +30826,10 @@ packages:
get-source: 2.0.12
dev: true
+ /standard-as-callback/2.1.0:
+ resolution: {integrity: sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A==}
+ dev: false
+
/static-extend/0.1.2:
resolution: {integrity: sha512-72E9+uLc27Mt718pMHt9VMNiAL4LMsmDbBva8mxWUCkT07fSzEGMYUCk0XWY6lp0j6RBAG4cJ3mWuZv2OE3s0g==}
engines: {node: '>=0.10.0'}