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 (
Env
Last Run
+ {displayStyle === "long" && Concurrency}
Version
Status
@@ -42,6 +46,23 @@ export function JobStatusTable({ environments }: JobStatusTableProps) {
{environment.lastRun ? : "Never Run"}
+ {displayStyle === "long" && (
+
+ {environment.concurrencyLimitGroup ? (
+
+ {environment.concurrencyLimitGroup.name}
+
+ ({environment.concurrencyLimitGroup.concurrencyLimit})
+
+
+ ) : typeof environment.concurrencyLimit === "number" ? (
+ {environment.concurrencyLimit}
+ ) : (
+ Not specified
+ )}
+
+ )}
+
{environment.version}
diff --git a/apps/webapp/app/components/navigation/SideMenu.tsx b/apps/webapp/app/components/navigation/SideMenu.tsx
index 75b7a2ada..3591f4c34 100644
--- a/apps/webapp/app/components/navigation/SideMenu.tsx
+++ b/apps/webapp/app/components/navigation/SideMenu.tsx
@@ -26,6 +26,7 @@ import {
projectEnvironmentsPath,
projectHttpEndpointsPath,
projectPath,
+ projectRunsPath,
projectSetupPath,
projectTriggersPath,
} from "~/utils/pathBuilder";
@@ -120,6 +121,12 @@ export function SideMenu({ user, project, organization, organizations }: SideMen
to={projectPath(organization, project)}
data-action="jobs"
/>
+
(
}
);
-type LinkPropsType = Pick & React.ComponentProps;
-export const LinkButton = ({ to, ...props }: LinkPropsType) => {
+type LinkPropsType = Pick &
+ React.ComponentProps;
+export const LinkButton = ({ to, onClick, ...props }: LinkPropsType) => {
const innerRef = useRef(null);
if (props.shortcut) {
useShortcutKeys({
@@ -297,6 +298,7 @@ export const LinkButton = ({ to, ...props }: LinkPropsType) => {
href={to.toString()}
ref={innerRef}
className={cn("group outline-none", props.fullWidth ? "w-full" : "")}
+ onClick={onClick}
>
@@ -307,6 +309,7 @@ export const LinkButton = ({ to, ...props }: LinkPropsType) => {
to={to}
ref={innerRef}
className={cn("group outline-none", props.fullWidth ? "w-full" : "")}
+ onClick={onClick}
>
diff --git a/apps/webapp/app/components/runs/RunStatuses.tsx b/apps/webapp/app/components/runs/RunStatuses.tsx
index 087006ba2..f993e963b 100644
--- a/apps/webapp/app/components/runs/RunStatuses.tsx
+++ b/apps/webapp/app/components/runs/RunStatuses.tsx
@@ -3,6 +3,7 @@ import {
CheckCircleIcon,
ClockIcon,
ExclamationTriangleIcon,
+ PauseCircleIcon,
WrenchIcon,
XCircleIcon,
} from "@heroicons/react/24/solid";
@@ -28,12 +29,13 @@ export function RunStatusIcon({ status, className }: { status: JobRunStatus; cla
case "SUCCESS":
return ;
case "PENDING":
- case "QUEUED":
+ case "WAITING_TO_CONTINUE":
return ;
+ case "QUEUED":
+ case "WAITING_TO_EXECUTE":
+ return ;
case "PREPROCESSING":
case "STARTED":
- case "WAITING_TO_CONTINUE":
- case "WAITING_TO_EXECUTE":
case "EXECUTING":
return ;
case "TIMED_OUT":
@@ -63,13 +65,12 @@ export function runStatusTitle(status: JobRunStatus): string {
case "STARTED":
return "In progress";
case "QUEUED":
+ case "WAITING_TO_EXECUTE":
return "Queued";
case "EXECUTING":
return "Executing";
case "WAITING_TO_CONTINUE":
return "Waiting";
- case "WAITING_TO_EXECUTE":
- return "Queued";
case "FAILURE":
return "Failed";
case "TIMED_OUT":
@@ -105,7 +106,7 @@ export function runStatusClassNameColor(status: JobRunStatus): string {
case "WAITING_TO_EXECUTE":
return "text-blue-500";
case "QUEUED":
- return "text-amber-300";
+ return "text-slate-500";
case "FAILURE":
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
diff --git a/apps/webapp/app/components/runs/RunsTable.tsx b/apps/webapp/app/components/runs/RunsTable.tsx
index 1494e27b6..0cb184c89 100644
--- a/apps/webapp/app/components/runs/RunsTable.tsx
+++ b/apps/webapp/app/components/runs/RunsTable.tsx
@@ -24,6 +24,7 @@ type RunTableItem = {
environment: {
type: RuntimeEnvironmentType;
};
+ job: { title: string; slug: string };
status: JobRunStatus;
startedAt: Date | null;
completedAt: Date | null;
@@ -36,6 +37,7 @@ type RunTableItem = {
type RunsTableProps = {
total: number;
hasFilters: boolean;
+ showJob?: boolean;
runs: RunTableItem[];
isLoading?: boolean;
runsParentPath: string;
@@ -46,6 +48,7 @@ export function RunsTable({
hasFilters,
runs,
isLoading = false,
+ showJob = false,
runsParentPath,
}: RunsTableProps) {
return (
@@ -53,6 +56,7 @@ export function RunsTable({
Run
+ {showJob && Job}
Env
Status
Started
@@ -68,21 +72,24 @@ export function RunsTable({
{total === 0 && !hasFilters ? (
-
-
+
+
) : runs.length === 0 ? (
-
+
) : (
runs.map((run) => {
- const path = `${runsParentPath}/${run.id}/trigger`;
+ const path = showJob
+ ? `${runsParentPath}/jobs/${run.job.slug}/runs/${run.id}/trigger`
+ : `${runsParentPath}/${run.id}/trigger`;
return (
{typeof run.number === "number" ? `#${run.number}` : "-"}
+ {showJob && {run.job.slug}}
@@ -130,6 +137,7 @@ export function RunsTable({
);
}
+
function NoRuns({ title }: { title: string }) {
return (
diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts
index a495b7707..39373e920 100644
--- a/apps/webapp/app/env.server.ts
+++ b/apps/webapp/app/env.server.ts
@@ -69,6 +69,7 @@ const EnvironmentSchema = z.object({
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/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts
index d5f5831df..a32bd3020 100644
--- a/apps/webapp/app/platform/zodWorker.server.ts
+++ b/apps/webapp/app/platform/zodWorker.server.ts
@@ -11,15 +11,6 @@ import type {
} from "graphile-worker";
import { run as graphileRun, parseCronItems } from "graphile-worker";
-import {
- Callback,
- Cluster,
- ClusterNode,
- ClusterOptions,
- Redis,
- RedisOptions,
- Result,
-} from "ioredis";
import omit from "lodash.omit";
import { z } from "zod";
import { PrismaClient, PrismaClientOrTransaction } from "~/db.server";
@@ -102,6 +93,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;
@@ -112,7 +108,7 @@ export type ZodWorkerOptions = {
cleanup?: ZodWorkerCleanupOptions;
reporter?: ZodWorkerReporter;
shutdownTimeoutInMs?: number;
- rateLimiter?: GraphileRateLimiter;
+ rateLimiter?: ZodWorkerRateLimiter;
};
export class ZodWorker {
@@ -125,7 +121,7 @@ export class ZodWorker {
#runner?: GraphileRunner;
#cleanup: ZodWorkerCleanupOptions | undefined;
#reporter?: ZodWorkerReporter;
- #rateLimiter?: GraphileRateLimiter;
+ #rateLimiter?: ZodWorkerRateLimiter;
#shutdownTimeoutInMs?: number;
#shuttingDown = false;
@@ -432,9 +428,10 @@ export class ZodWorker {
await this.enqueue(helpers.job.task_identifier, payload, {
runAt: helpers.job.run_at,
queueName: helpers.job.queue_name ?? undefined,
- priority: helpers.job.priority - 1,
+ priority: helpers.job.priority,
jobKey: helpers.job.key ?? undefined,
flags: Object.keys(helpers.job.flags ?? []),
+ maxAttempts: helpers.job.max_attempts,
});
}
@@ -676,239 +673,3 @@ function removeUndefinedKeys(obj: T): T {
}
return obj;
}
-
-export interface GraphileRateLimiter {
- forbiddenFlags(): Promise;
- wrapTask(t: Task, rescheduler: Task): Task;
- setMaxSizeForFlag(flag: string, maxSize: number): Promise;
- delMaxSizeForFlag(flag: string): Promise;
-}
-
-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;
-
- afterTask(
- setKey: string,
- maxSizeKey: string,
- forbiddenFlagsKey: string,
- jobId: string,
- timestamp: string,
- windowSize: string,
- forbiddenFlag: string,
- maxSize: string,
- callback?: Callback
- ): Result;
- }
-}
-
-export type RedisGraphileRateLimiterOptions = {
- redis?: RedisOptions;
- cluster?: {
- startupNodes: ClusterNode[];
- options?: ClusterOptions;
- };
- defaultConcurrency?: number;
- windowSize?: number;
- prefix?: string;
-};
-
-const FORBIDDEN_FLAG_KEY = "rl:forbiddenFlags";
-
-// TODO: we need to somehow seed and update the rate limit for each flag in Redis
-export class RedisGraphileRateLimiter implements GraphileRateLimiter {
- private redis: Redis | Cluster;
- private defaultMaxSize: number;
- private windowSize: number;
-
- constructor(options?: RedisGraphileRateLimiterOptions) {
- 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)
- redis.call('SREM', forbiddenFlagsKey, forbiddenFlag)
-
- return true
-else
- redis.call('SADD', forbiddenFlagsKey, forbiddenFlag)
-
- return false
-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 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);
- }
-
- const beforeResults = await Promise.allSettled(
- flags.map(async (flag) => this.#callBeforeTask(flag, String(helpers.job.id)))
- );
-
- // If any of the beforeTask calls returned false, then we need to re-schedule the task and return
- if (beforeResults.some((result) => result.status === "rejected")) {
- return await rescheduler(payload, helpers);
- }
-
- if (
- beforeResults.some(
- (result) => result.status === "fulfilled" && result.value?.results === null
- )
- ) {
- return await rescheduler(payload, helpers);
- }
-
- 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) {
- try {
- const now = performance.now();
- const results = 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 {
- results,
- durationInMs,
- };
- } catch (error) {
- logger.error("Failed to call beforeTask", { error, flag, 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 });
- }
- }
-}
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 f1907b3de..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;
@@ -60,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,
},
@@ -83,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
@@ -94,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;
@@ -103,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) => ({
@@ -133,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/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 3801568a5..d2c5867fb 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
@@ -15,31 +15,39 @@ export function ListPagination({ list, className }: { list: RunList; className?:
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.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/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 bee4dccec..f7a4cdba7 100644
--- a/apps/webapp/app/services/jobs/registerJob.server.ts
+++ b/apps/webapp/app/services/jobs/registerJob.server.ts
@@ -13,6 +13,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;
@@ -104,32 +105,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({
@@ -141,57 +138,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: {
@@ -199,9 +168,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..59c9f8f09
--- /dev/null
+++ b/apps/webapp/app/services/runExecutionRateLimiter.server.ts
@@ -0,0 +1,400 @@
+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: {
+ username: env.REDIS_USERNAME,
+ password: env.REDIS_PASSWORD,
+ tls: {},
+ enableAutoPipelining: true,
+ },
+ dnsLookup: (address, callback) => callback(null, address),
+ },
+ },
+ 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,
+ },
+ defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
+ });
+ }
+ }
+}
diff --git a/apps/webapp/app/services/runs/createRun.server.ts b/apps/webapp/app/services/runs/createRun.server.ts
index 728460ae3..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,
@@ -54,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 e7e2f97d9..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 { JobRun } from "@trigger.dev/database";
+import {
+ ConcurrencyLimitGroup,
+ JobRun,
+ JobVersion,
+ RuntimeEnvironment,
+} from "@trigger.dev/database";
import { generateErrorMessage } from "zod-error";
import { eventRecordToApiJson } from "~/api.server";
import {
@@ -39,6 +44,8 @@ import { ResumeTaskService } from "../tasks/resumeTask.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];
@@ -86,7 +93,12 @@ export class PerformRunExecutionV3Service {
}
static async enqueue(
- run: JobRun,
+ run: JobRun & {
+ version: JobVersion & {
+ environment: RuntimeEnvironment;
+ concurrencyLimitGroup?: ConcurrencyLimitGroup | null;
+ };
+ },
priority: RunExecutionPriority,
tx: PrismaClientOrTransaction,
options: {
@@ -104,8 +116,8 @@ export class PerformRunExecutionV3Service {
tx,
runAt: options.runAt,
jobKey: `job_run:EXECUTE_JOB:${run.id}`,
- maxAttempts: options.skipRetrying ? 1 : undefined,
- flags: [`rl:executions:${run.organizationId}`],
+ maxAttempts: options.skipRetrying ? env.DEFAULT_DEV_ENV_EXECUTION_ATTEMPTS : undefined,
+ flags: executionRateLimiter?.flagsForRun(run, run.version) ?? [],
priority: priority === "initial" ? 0 : -1,
}
);
@@ -166,6 +178,7 @@ export class PerformRunExecutionV3Service {
projectId: run.projectId,
jobId: run.jobId,
runId: run.id,
+ concurrencyLimitGroupId: run.version.concurrencyLimitGroupId,
});
forceYieldCoordinator.registerRun(run.id);
@@ -183,6 +196,7 @@ export class PerformRunExecutionV3Service {
projectId: run.projectId,
jobId: run.jobId,
runId: run.id,
+ concurrencyLimitGroupId: run.version.concurrencyLimitGroupId,
});
forceYieldCoordinator.deregisterRun(run.id);
diff --git a/apps/webapp/app/services/runs/resumeRun.server.ts b/apps/webapp/app/services/runs/resumeRun.server.ts
index afc08751d..6b99035ff 100644
--- a/apps/webapp/app/services/runs/resumeRun.server.ts
+++ b/apps/webapp/app/services/runs/resumeRun.server.ts
@@ -118,7 +118,7 @@ export class ResumeRunService {
async #executeRun(run: FoundRun, priority: RunExecutionPriority) {
await PerformRunExecutionV3Service.enqueue(run, priority, this.#prismaClient, {
- skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
+ skipRetrying: run.version.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
}
@@ -147,7 +147,12 @@ async function findRun(prisma: PrismaClientOrTransaction, id: string) {
return await prisma.jobRun.findUnique({
where: { id },
include: {
- environment: true,
+ version: {
+ include: {
+ environment: true,
+ concurrencyLimitGroup: true,
+ },
+ },
},
});
}
diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts
index fd3f58423..fea528d41 100644
--- a/apps/webapp/app/services/worker.server.ts
+++ b/apps/webapp/app/services/worker.server.ts
@@ -3,11 +3,7 @@ import { ScheduledPayloadSchema, addMissingVersionField } from "@trigger.dev/cor
import { z } from "zod";
import { prisma } from "~/db.server";
import { env } from "~/env.server";
-import {
- GraphileRateLimiter,
- RedisGraphileRateLimiter,
- ZodWorker,
-} from "~/platform/zodWorker.server";
+import { ZodWorker } from "~/platform/zodWorker.server";
import { sendEmail } from "./email.server";
import { IndexEndpointService } from "./endpoints/indexEndpoint.server";
import { PerformEndpointIndexService } from "./endpoints/performEndpointIndexService";
@@ -31,6 +27,7 @@ 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";
const workerCatalog = {
indexEndpoint: z.object({
@@ -437,7 +434,7 @@ function getExecutionWorkerQueue() {
},
shutdownTimeoutInMs: env.GRACEFUL_SHUTDOWN_TIMEOUT,
schema: executionWorkerCatalog,
- rateLimiter: getRateLimiter(),
+ rateLimiter: executionRateLimiter,
tasks: {
performRunExecutionV2: {
priority: 0, // smaller number = higher priority
@@ -514,43 +511,4 @@ function getTaskOperationWorkerQueue() {
});
}
-function getRateLimiter(): GraphileRateLimiter | undefined {
- if (env.REDIS_HOST && env.REDIS_PORT) {
- if (env.REDIS_READER_HOST) {
- return new RedisGraphileRateLimiter({
- 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: "tr:gw:",
- scaleReads: "slave",
- redisOptions: {
- username: env.REDIS_USERNAME,
- password: env.REDIS_PASSWORD,
- tls: {},
- enableAutoPipelining: true,
- },
- dnsLookup: (address, callback) => callback(null, address),
- },
- },
- defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
- });
- } else {
- return new RedisGraphileRateLimiter({
- redis: {
- keyPrefix: "tr:gw:",
- port: env.REDIS_PORT,
- host: env.REDIS_HOST,
- username: env.REDIS_USERNAME,
- password: env.REDIS_PASSWORD,
- enableAutoPipelining: true,
- },
- defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
- });
- }
- }
-}
-
export { executionWorker, workerQueue, taskOperationWorker };
diff --git a/apps/webapp/app/utils/pathBuilder.ts b/apps/webapp/app/utils/pathBuilder.ts
index 241eea6c5..8b8b743fb 100644
--- a/apps/webapp/app/utils/pathBuilder.ts
+++ b/apps/webapp/app/utils/pathBuilder.ts
@@ -120,6 +120,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/packages/core/src/schemas/api.ts b/packages/core/src/schemas/api.ts
index 2159146e4..2b6affdcb 100644
--- a/packages/core/src/schemas/api.ts
+++ b/packages/core/src/schemas/api.ts
@@ -219,6 +219,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(),
@@ -230,6 +235,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;
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 3e058ad7b..0ca05a079 100644
--- a/packages/database/prisma/schema.prisma
+++ b/packages/database/prisma/schema.prisma
@@ -323,6 +323,7 @@ model RuntimeEnvironment {
scheduleSources ScheduleSource[]
ExternalAccount ExternalAccount[]
httpEndpointEnvironments TriggerHttpEndpointEnvironment[]
+ concurrencyLimitGroups ConcurrencyLimitGroup[]
@@unique([projectId, slug, orgMemberId])
@@unique([projectId, shortcode])
@@ -477,12 +478,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
@@ -521,6 +526,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
@@ -731,8 +753,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?
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 aeb14016d..82414705e 100644
--- a/packages/trigger-sdk/src/triggerClient.ts
+++ b/packages/trigger-sdk/src/triggerClient.ts
@@ -79,6 +79,7 @@ const registerSourceEvent: EventSpecification = {
import EventEmitter from "node:events";
import * as packageJson from "../package.json";
+import { ConcurrencyLimit, ConcurrencyLimitOptions } from "./concurrencyLimit";
export type TriggerClientOptions = {
/** The `id` property is used to uniquely identify the client.
@@ -622,6 +623,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);
@@ -1440,6 +1445,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 c0eaba19b..4083f9d94 100644
--- a/perf/src/index.ts
+++ b/perf/src/index.ts
@@ -172,7 +172,7 @@ async function mainSerial() {
async function mainConcurrency() {
const batches = 1;
const concurrency = 10;
- const eventsPer = 10;
+ const eventsPer = 5;
console.log("Preparing perf tests...");
diff --git a/perf/src/trigger.ts b/perf/src/trigger.ts
index 102fbd45d..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,6 +18,101 @@ 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, 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-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",
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"],