Merge branch 'main' into worker-upgrade

This commit is contained in:
nicktrn
2023-11-29 09:37:13 +00:00
133 changed files with 2539 additions and 796 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/sdk": patch
---
Fix: `Key-Value Store` keys will now be URI encoded
-12
View File
@@ -1,12 +0,0 @@
---
"@trigger.dev/integration-kit": patch
"@trigger.dev/airtable": patch
"@trigger.dev/shopify": patch
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
"@trigger.dev/cli": patch
---
- Simplify `Webhook Triggers` and use the new HTTP Endpoints
- Add a `Key-Value Store` for use in and outside of Jobs
- Add a `@trigger.dev/shopify` package
-6
View File
@@ -1,6 +0,0 @@
---
"@trigger.dev/shopify": patch
"@trigger.dev/sdk": patch
---
Fix `@trigger.dev/shopify` imports, enhance docs, and suppress HTTP Endpoint warnings
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/openai": patch
---
Adding additional assistant tasks
+6
View File
@@ -0,0 +1,6 @@
---
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
---
Feature: Run execution concurrency limits
+7
View File
@@ -1,5 +1,12 @@
# proxy
## 0.0.2
### Patch Changes
- Updated dependencies [067e19fe]
- @trigger.dev/core@2.2.8
## 0.0.1
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "proxy",
"version": "0.0.1",
"version": "0.0.2",
"private": true,
"scripts": {
"deploy": "wrangler deploy",
+22 -1
View File
@@ -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 (
<Table fullWidth>
<TableHeader>
<TableRow>
<TableHeaderCell>Env</TableHeaderCell>
<TableHeaderCell>Last Run</TableHeaderCell>
{displayStyle === "long" && <TableHeaderCell>Concurrency</TableHeaderCell>}
<TableHeaderCell alignment="right">Version</TableHeaderCell>
<TableHeaderCell alignment="right">Status</TableHeaderCell>
</TableRow>
@@ -42,6 +46,23 @@ export function JobStatusTable({ environments }: JobStatusTableProps) {
<TableCell>
{environment.lastRun ? <DateTime date={environment.lastRun} /> : "Never Run"}
</TableCell>
{displayStyle === "long" && (
<TableCell>
{environment.concurrencyLimitGroup ? (
<span className="flex items-center gap-1">
<span>{environment.concurrencyLimitGroup.name}</span>
<span className="text-gray-400">
({environment.concurrencyLimitGroup.concurrencyLimit})
</span>
</span>
) : typeof environment.concurrencyLimit === "number" ? (
<span className="text-gray-400">{environment.concurrencyLimit}</span>
) : (
<span className="text-gray-400">Not specified</span>
)}
</TableCell>
)}
<TableCell alignment="right">{environment.version}</TableCell>
<TableCell alignment="right">
<ActiveBadge active={environment.enabled} />
@@ -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"
/>
<SideMenuItem
name="Runs"
icon="runs"
iconColor="text-teal-500"
to={projectRunsPath(organization, project)}
/>
<SideMenuItem
name="Triggers"
icon="trigger"
@@ -304,8 +304,9 @@ export const Button = forwardRef<HTMLButtonElement, ButtonPropsType>(
}
);
type LinkPropsType = Pick<LinkProps, "to" | "target"> & React.ComponentProps<typeof ButtonContent>;
export const LinkButton = ({ to, ...props }: LinkPropsType) => {
type LinkPropsType = Pick<LinkProps, "to" | "target" | "onClick"> &
React.ComponentProps<typeof ButtonContent>;
export const LinkButton = ({ to, onClick, ...props }: LinkPropsType) => {
const innerRef = useRef<HTMLAnchorElement>(null);
if (props.shortcut) {
useShortcutKeys({
@@ -324,6 +325,7 @@ export const LinkButton = ({ to, ...props }: LinkPropsType) => {
href={to.toString()}
ref={innerRef}
className={cn("group outline-none", props.fullWidth ? "w-full" : "")}
onClick={onClick}
>
<ButtonContent {...props} />
</ExtLink>
@@ -334,6 +336,7 @@ export const LinkButton = ({ to, ...props }: LinkPropsType) => {
to={to}
ref={innerRef}
className={cn("group outline-none", props.fullWidth ? "w-full" : "")}
onClick={onClick}
>
<ButtonContent {...props} />
</Link>
+25 -27
View File
@@ -10,13 +10,14 @@ import {
useNavigate,
useNavigation,
} from "@remix-run/react";
import { JobRunStatus, RuntimeEnvironmentType } from "@trigger.dev/database";
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { useMemo } from "react";
import { usePathName } from "~/hooks/usePathName";
import type { RunBasicStatus } from "~/models/jobRun.server";
import { ViewRun } from "~/presenters/RunPresenter.server";
import { cancelSchema } from "~/routes/resources.runs.$runId.cancel";
import { schema } from "~/routes/resources.runs.$runId.rerun";
import { formatDuration } from "~/utils";
import { formatDuration, formatDurationMilliseconds } from "~/utils";
import { cn } from "~/utils/cn";
import { runCompletedPath, runTaskPath, runTriggerPath } from "~/utils/pathBuilder";
import { CodeBlock } from "../code/CodeBlock";
@@ -38,14 +39,7 @@ import {
} from "../primitives/PageHeader";
import { Paragraph } from "../primitives/Paragraph";
import { Popover, PopoverContent, PopoverTrigger } from "../primitives/Popover";
import {
RunBasicStatus,
RunStatusIcon,
RunStatusLabel,
hasFinished,
runBasicStatus,
runStatusTitle,
} from "../runs/RunStatuses";
import { RunStatusIcon, RunStatusLabel, runStatusTitle } from "../runs/RunStatuses";
import {
RunPanel,
RunPanelBody,
@@ -95,8 +89,6 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
}
}, [pathName]);
const basicStatus = runBasicStatus(run.status);
return (
<PageContainer>
<PageHeader>
@@ -106,7 +98,9 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
to: paths.back,
text: "Runs",
}}
title={`Run #${run.number}`}
title={
typeof run.number === "number" ? `Run #${run.number}` : `Run ${run.id.slice(0, 8)}`
}
/>
<PageButtons>
{run.isTest && (
@@ -115,15 +109,15 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
Test run
</span>
)}
{showRerun && hasFinished(run.status) && (
{showRerun && run.isFinished && (
<RerunPopover
runId={run.id}
runsPath={paths.runsPath}
environmentType={run.environment.type}
status={basicStatus}
status={run.basicStatus}
/>
)}
{!hasFinished(run.status) && <CancelRun runId={run.id} />}
{!run.isFinished && <CancelRun runId={run.id} />}
</PageButtons>
</PageTitleRow>
<PageInfoRow>
@@ -146,7 +140,17 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
<PageInfoProperty
icon={"clock"}
label={"Duration"}
value={formatDuration(run.startedAt, run.completedAt)}
value={formatDuration(run.startedAt, run.completedAt, { style: "short" })}
/>
<PageInfoProperty
icon={"hourglass"}
label={"Execution Time"}
value={formatDurationMilliseconds(run.executionDuration, { style: "short" })}
/>
<PageInfoProperty
icon={"list-numbers"}
label={"Execution Count"}
value={run.executionCount}
/>
</PageInfoGroup>
<PageInfoGroup alignment="right">
@@ -211,10 +215,10 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
);
})
) : (
<BlankTasks status={run.status} basicStatus={basicStatus} />
<BlankTasks status={run.basicStatus} />
)}
</div>
{(basicStatus === "COMPLETED" || basicStatus === "FAILED") && (
{(run.basicStatus === "COMPLETED" || run.basicStatus === "FAILED") && (
<div>
<Header2 className={cn("mb-2")}>Run Summary</Header2>
<RunPanel
@@ -285,14 +289,8 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
);
}
function BlankTasks({
status,
basicStatus,
}: {
status: JobRunStatus;
basicStatus: RunBasicStatus;
}) {
switch (basicStatus) {
function BlankTasks({ status }: { status: RunBasicStatus }) {
switch (status) {
default:
case "COMPLETED":
return <Paragraph variant="small">There were no tasks for this run.</Paragraph>;
+21 -43
View File
@@ -3,6 +3,7 @@ import {
CheckCircleIcon,
ClockIcon,
ExclamationTriangleIcon,
PauseCircleIcon,
WrenchIcon,
XCircleIcon,
} from "@heroicons/react/24/solid";
@@ -10,18 +11,6 @@ import type { JobRunStatus } from "@trigger.dev/database";
import { cn } from "~/utils/cn";
import { Spinner } from "../primitives/Spinner";
export function hasFinished(status: JobRunStatus): boolean {
return (
status === "SUCCESS" ||
status === "FAILURE" ||
status === "ABORTED" ||
status === "TIMED_OUT" ||
status === "CANCELED" ||
status === "UNRESOLVED_AUTH" ||
status === "INVALID_PAYLOAD"
);
}
export function RunStatus({ status }: { status: JobRunStatus }) {
return (
<span className="flex items-center gap-1">
@@ -40,49 +29,26 @@ export function RunStatusIcon({ status, className }: { status: JobRunStatus; cla
case "SUCCESS":
return <CheckCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "PENDING":
case "WAITING_TO_CONTINUE":
return <ClockIcon className={cn(runStatusClassNameColor(status), className)} />;
case "QUEUED":
return <ClockIcon className={cn(runStatusClassNameColor(status), className)} />;
case "WAITING_TO_EXECUTE":
return <PauseCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "PREPROCESSING":
case "STARTED":
case "EXECUTING":
return <Spinner className={cn(runStatusClassNameColor(status), className)} />;
case "FAILURE":
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "TIMED_OUT":
return <ExclamationTriangleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "UNRESOLVED_AUTH":
case "FAILURE":
case "ABORTED":
case "INVALID_PAYLOAD":
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "WAITING_ON_CONNECTIONS":
return <WrenchIcon className={cn(runStatusClassNameColor(status), className)} />;
case "ABORTED":
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "PREPROCESSING":
return <Spinner className={cn(runStatusClassNameColor(status), className)} />;
case "CANCELED":
return <NoSymbolIcon className={cn(runStatusClassNameColor(status), className)} />;
}
}
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":
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}`);
@@ -99,7 +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 "FAILURE":
return "Failed";
case "TIMED_OUT":
@@ -130,9 +101,12 @@ export function runStatusClassNameColor(status: JobRunStatus): string {
case "PENDING":
return "text-slate-500";
case "STARTED":
case "EXECUTING":
case "WAITING_TO_CONTINUE":
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":
@@ -147,5 +121,9 @@ export function runStatusClassNameColor(status: JobRunStatus): string {
return "text-blue-500";
case "CANCELED":
return "text-slate-500";
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
}
}
}
+24 -7
View File
@@ -1,7 +1,7 @@
import { StopIcon } from "@heroicons/react/24/outline";
import { CheckIcon } from "@heroicons/react/24/solid";
import { JobRunStatus, RuntimeEnvironmentType } from "@trigger.dev/database";
import { formatDuration } from "~/utils";
import { formatDuration, formatDurationMilliseconds } from "~/utils";
import { EnvironmentLabel } from "../environments/EnvironmentLabel";
import { DateTime } from "../primitives/DateTime";
import { Paragraph } from "../primitives/Paragraph";
@@ -20,14 +20,16 @@ import { RunStatus } from "./RunStatuses";
type RunTableItem = {
id: string;
number: number;
number: number | null;
environment: {
type: RuntimeEnvironmentType;
};
job: { title: string; slug: string };
status: JobRunStatus;
startedAt: Date | null;
completedAt: Date | null;
createdAt: Date | null;
executionDuration: number;
version: string;
isTest: boolean;
};
@@ -35,6 +37,7 @@ type RunTableItem = {
type RunsTableProps = {
total: number;
hasFilters: boolean;
showJob?: boolean;
runs: RunTableItem[];
isLoading?: boolean;
runsParentPath: string;
@@ -45,6 +48,7 @@ export function RunsTable({
hasFilters,
runs,
isLoading = false,
showJob = false,
runsParentPath,
}: RunsTableProps) {
return (
@@ -52,10 +56,12 @@ export function RunsTable({
<TableHeader>
<TableRow>
<TableHeaderCell>Run</TableHeaderCell>
{showJob && <TableHeaderCell>Job</TableHeaderCell>}
<TableHeaderCell>Env</TableHeaderCell>
<TableHeaderCell>Status</TableHeaderCell>
<TableHeaderCell>Started</TableHeaderCell>
<TableHeaderCell>Duration</TableHeaderCell>
<TableHeaderCell>Exec Time</TableHeaderCell>
<TableHeaderCell>Test</TableHeaderCell>
<TableHeaderCell>Version</TableHeaderCell>
<TableHeaderCell>Created at</TableHeaderCell>
@@ -66,19 +72,24 @@ export function RunsTable({
</TableHeader>
<TableBody>
{total === 0 && !hasFilters ? (
<TableBlankRow colSpan={8}>
<NoRuns title="No Runs found for this Job" />
<TableBlankRow colSpan={showJob ? 10 : 9}>
<NoRuns title="No Runs found" />
</TableBlankRow>
) : runs.length === 0 ? (
<TableBlankRow colSpan={8}>
<TableBlankRow colSpan={showJob ? 10 : 9}>
<NoRuns title="No Runs match your filters" />
</TableBlankRow>
) : (
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 (
<TableRow key={run.id}>
<TableCell to={path}>#{run.number}</TableCell>
<TableCell to={path}>
{typeof run.number === "number" ? `#${run.number}` : "-"}
</TableCell>
{showJob && <TableCell to={path}>{run.job.slug}</TableCell>}
<TableCell to={path}>
<EnvironmentLabel environment={run.environment} />
</TableCell>
@@ -93,6 +104,11 @@ export function RunsTable({
style: "short",
})}
</TableCell>
<TableCell to={path}>
{formatDurationMilliseconds(run.executionDuration, {
style: "short",
})}
</TableCell>
<TableCell to={path}>
{run.isTest ? (
<CheckIcon className="h-4 w-4 text-slate-400" />
@@ -121,6 +137,7 @@ export function RunsTable({
</Table>
);
}
function NoRuns({ title }: { title: string }) {
return (
<div className="flex items-center justify-center">
+13 -8
View File
@@ -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,18 @@ 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(),
REDIS_TLS_DISABLED: 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<typeof EnvironmentSchema>;
+45
View File
@@ -0,0 +1,45 @@
import type { JobRun, JobRunStatus } from "@trigger.dev/database";
const COMPLETED_STATUSES: Array<JobRun["status"]> = [
"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}`);
}
}
}
@@ -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,
});
}
+27 -1
View File
@@ -115,6 +115,11 @@ export type ZodWorkerCleanupOptions = {
type ZodWorkerReporter = (event: string, properties: Record<string, any>) => Promise<void>;
export interface ZodWorkerRateLimiter {
forbiddenFlags(): Promise<string[]>;
wrapTask(t: Task, rescheduler: Task): Task;
}
export type ZodWorkerOptions<TMessageCatalog extends MessageCatalogSchema> = {
name: string;
runnerOptions: RunnerOptions;
@@ -125,6 +130,7 @@ export type ZodWorkerOptions<TMessageCatalog extends MessageCatalogSchema> = {
cleanup?: ZodWorkerCleanupOptions;
reporter?: ZodWorkerReporter;
shutdownTimeoutInMs?: number;
rateLimiter?: ZodWorkerRateLimiter;
};
export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
@@ -137,6 +143,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
#runner?: GraphileRunner;
#cleanup: ZodWorkerCleanupOptions | undefined;
#reporter?: ZodWorkerReporter;
#rateLimiter?: ZodWorkerRateLimiter;
#shutdownTimeoutInMs?: number;
#shuttingDown = false;
#workerUtils?: WorkerUtils;
@@ -150,6 +157,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
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
}
@@ -175,6 +183,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
noHandleSignals: true,
taskList: this.#createTaskListFromTasks(),
parsedCronItems,
forbiddenFlags: this.#rateLimiter?.forbiddenFlags.bind(this.#rateLimiter),
});
if (!this.#runner) {
@@ -525,7 +534,11 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
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 ?? {})) {
@@ -555,6 +568,19 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
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[] = [];
@@ -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 });
@@ -6,14 +6,15 @@ export type Direction = z.infer<typeof DirectionSchema>;
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<ReturnType<RunListPresenter["call"]>>;
@@ -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,
@@ -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,
@@ -83,6 +86,8 @@ export class RunPresenter {
runConnections: run.runConnections,
missingConnections: run.missingConnections,
error: runError,
executionDuration: run.executionDuration,
executionCount: run.executionCount,
};
}
@@ -114,6 +119,8 @@ export class RunPresenter {
properties: true,
output: true,
payload: true,
executionCount: true,
executionDuration: true,
version: {
select: {
version: true,
@@ -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() {
</LinkButton>
</PageButtons>
</PageTitleRow>
<PageInfoRow>
<PageInfoGroup alignment="right">
<Paragraph variant="extra-small" className="text-slate-600">
UID: {organization.id}
</Paragraph>
</PageInfoGroup>
</PageInfoRow>
</PageHeader>
<PageBody>
<ul className="grid grid-cols-1 gap-4 md:grid-cols-2 lg:grid-cols-3 xl:grid-cols-4">
@@ -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();
@@ -22,31 +22,39 @@ export function ListPagination({
function NextButton({ cursor }: { cursor?: string }) {
const path = useCursorPath(cursor, "forward");
return path ? (
return (
<LinkButton
to={path}
to={path ?? "#"}
variant={"tertiary/small"}
TrailingIcon="chevron-right"
className="flex items-center"
className={cn(
"flex items-center",
!path && "cursor-default opacity-50 group-hover:bg-transparent group-hover:text-slate-800"
)}
onClick={(e) => !path && e.preventDefault()}
>
Next
</LinkButton>
) : null;
);
}
function PreviousButton({ cursor }: { cursor?: string }) {
const path = useCursorPath(cursor, "backward");
return path ? (
return (
<LinkButton
to={path}
to={path ?? "#"}
variant={"tertiary/small"}
LeadingIcon="chevron-left"
className="flex items-center"
className={cn(
"flex items-center",
!path && "cursor-default opacity-50 group-hover:bg-transparent group-hover:text-slate-800"
)}
onClick={(e) => !path && e.preventDefault()}
>
Prev
</LinkButton>
) : null;
);
}
function useCursorPath(cursor: string | undefined, direction: Direction) {
@@ -72,8 +72,8 @@ export default function Page() {
<div className={cn("grid h-fit gap-4", open ? "grid-cols-2" : "grid-cols-1")}>
<div>
<div className="mb-2 flex items-center justify-end gap-x-2">
<ListPagination list={list} />
<HelpTrigger title="How do I run my Job?" />
<ListPagination list={list} />
</div>
<RunsTable
total={list.runs.length}
@@ -24,7 +24,7 @@ export default function Page() {
const project = useProject();
return (
<Help defaultOpen>
<Help>
{(open) => (
<div className={cn("grid h-fit gap-4", open ? "grid-cols-2" : "grid-cols-1")}>
<div className="w-full">
@@ -32,7 +32,7 @@ export default function Page() {
<Header2 className="mb-2 flex items-center gap-1">Environments</Header2>
<HelpTrigger title="How do disable a Job?" />
</div>
<JobStatusTable environments={job.environments} />
<JobStatusTable environments={job.environments} displayStyle="long" />
<div className="mt-4 flex w-full items-center justify-end gap-x-3">
{job.status === "ACTIVE" && (
<Paragraph variant="small">
@@ -297,7 +297,9 @@ export default function Page() {
label={<DateTime date={run.created} />}
description={
<>
Run #{run.number}{" "}
{typeof run.number === "number"
? `Run #${run.number}`
: `Run ${run.id.slice(0, 8)}`}
<span className={runStatusClassNameColor(run.status)}>
{runStatusTitle(run.status).toLocaleLowerCase()}
</span>
@@ -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<typeof loader>();
const navigation = useNavigation();
const isLoading = navigation.state !== "idle";
const organization = useOrganization();
const project = useProject();
return (
<PageContainer>
<PageHeader hideBorder>
<PageTitleRow>
<PageTitle title={`${project.name} Runs`} />
<PageButtons>
<LinkButton
LeadingIcon={"docs"}
to={docsPath("documentation/concepts/runs")}
variant="secondary/small"
>
Run documentation
</LinkButton>
</PageButtons>
</PageTitleRow>
<PageDescription>All job runs in this project</PageDescription>
</PageHeader>
<PageBody scrollable={false}>
<div className="h-full overflow-y-auto p-4 scrollbar-thin scrollbar-track-transparent scrollbar-thumb-slate-700">
<div className="mb-2 flex items-center justify-end gap-x-2">
<ListPagination list={list} />
</div>
<RunsTable
total={list.runs.length}
hasFilters={false}
showJob={true}
runs={list.runs}
isLoading={isLoading}
runsParentPath={projectPath(organization, project)}
/>
<ListPagination list={list} className="mt-2 justify-end" />
</div>
</PageBody>
</PageContainer>
);
}
@@ -97,6 +97,7 @@ export default function NewOrganizationPage() {
{...conform.input(orgName, { type: "text" })}
placeholder="Your Organization name"
icon="organization"
autoFocus
/>
<Hint>E.g. your company name or your workspace name.</Hint>
<FormError id={orgName.errorId}>{orgName.error}</FormError>
+5 -5
View File
@@ -42,14 +42,14 @@ export async function action({ request, params }: ActionFunctionArgs) {
const store = new KeyValueStore(authenticatedEnv);
const { key } = parsedParams.data;
const decodedKey = decodeURIComponent(parsedParams.data.key);
try {
switch (parsedMethod.data) {
case "DELETE": {
const deleted = await store.delete(key);
const deleted = await store.delete(decodedKey);
return json({ action: "DELETE", key, deleted });
return json({ action: "DELETE", key: decodedKey, deleted });
}
case "PUT": {
const value = await request.text();
@@ -65,9 +65,9 @@ export async function action({ request, params }: ActionFunctionArgs) {
);
}
const setValue = await store.set(key, value);
const setValue = await store.set(decodedKey, value);
return json({ action: "SET", key, value: setValue });
return json({ action: "SET", key: decodedKey, value: setValue });
}
default: {
assertExhaustive(parsedMethod.data);
@@ -324,7 +324,7 @@ export class PerformEndpointIndexService {
for (const webhook of webhooks) {
try {
await this.#registerWebhookService.call(endpoint, webhook);
indexStats.webhooks++;
indexStats.webhooks = indexStats.webhooks ?? 0 + 1;
} catch (error) {
logger.error("Failed to register webhook", {
endpointId: endpoint.id,
@@ -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}
)
`;
}
@@ -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<string>();
if (examples) {
@@ -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<void>;
putJobVersionConcurrencyLimit(jobVersion: JobVersion, env: RuntimeEnvironment): Promise<void>;
setMaxSizeForFlag(flag: string, maxSize: number): Promise<void>;
delMaxSizeForFlag(flag: string): Promise<void>;
flagsForRun(
run: JobRun,
version: JobVersion & {
environment: RuntimeEnvironment;
concurrencyLimitGroup?: ConcurrencyLimitGroup;
}
): string[];
}
declare module "ioredis" {
interface RedisCommander<Context> {
beforeTask(
setKey: string,
maxSizeKey: string,
forbiddenFlagsKey: string,
jobId: string,
timestamp: string,
windowSize: string,
forbiddenFlag: string,
maxSize: string,
callback?: Callback<string>
): Result<number | null, Context>;
rollbackBeforeTask(keys: number, ...args: string[]): Result<string, Context>;
afterTask(
setKey: string,
maxSizeKey: string,
forbiddenFlagsKey: string,
jobId: string,
timestamp: string,
windowSize: string,
forbiddenFlag: string,
maxSize: string,
callback?: Callback<string>
): Result<number | null, Context>;
}
}
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<string[]> {
return this.redis.smembers(FORBIDDEN_FLAG_KEY);
}
async putConcurrencyLimitGroup(
concurrencyLimitGroup: ConcurrencyLimitGroup,
env: RuntimeEnvironment
): Promise<void> {
await this.setMaxSizeForFlag(
this.flagForConcurrencyLimitGroup(concurrencyLimitGroup, env),
concurrencyLimitGroup.concurrencyLimit
);
}
async putJobVersionConcurrencyLimit(
jobVersion: JobVersion,
env: RuntimeEnvironment
): Promise<void> {
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<void> {
await this.redis.set(`${flag}:maxSize`, String(maxSize));
}
async delMaxSizeForFlag(flag: string): Promise<void> {
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,
...(env.REDIS_TLS_DISABLED === "true" ? {} : { tls: {} })
},
defaultConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
});
}
}
}
@@ -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;
@@ -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 }
);
@@ -35,12 +35,6 @@ export class CreateRunService {
},
});
const jobQueue = await this.#prismaClient.jobQueue.findUniqueOrThrow({
where: {
id: version.queueId,
},
});
const eventRecords = await this.#prismaClient.eventRecord.findMany({
where: {
id: { in: eventIds },
@@ -54,22 +48,8 @@ export class CreateRunService {
const firstEvent = eventRecords[0];
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,
@@ -78,7 +58,6 @@ export class CreateRunService {
organizationId: environment.organizationId,
projectId: environment.projectId,
endpointId: endpoint.id,
queueId: jobQueue.id,
payload: JSON.stringify(
eventRecords.length > 1
? eventRecords.map((event) => event.payload) ?? [{}]
@@ -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<Awaited<ReturnType<typeof findRun>>>;
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": {
@@ -646,6 +535,9 @@ export class PerformRunExecutionV3Service {
executionDuration: {
increment: durationInMs,
},
executionCount: {
increment: 1,
},
},
});
@@ -663,17 +555,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,
},
},
});
@@ -746,7 +639,6 @@ export class PerformRunExecutionV3Service {
case "ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.error ?? undefined,
"FAILURE",
@@ -756,7 +648,6 @@ export class PerformRunExecutionV3Service {
case "INVALID_PAYLOAD": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.errors,
"INVALID_PAYLOAD",
@@ -776,7 +667,6 @@ export class PerformRunExecutionV3Service {
case "UNRESOLVED_AUTH_ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.issues,
"UNRESOLVED_AUTH",
@@ -807,14 +697,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);
});
}
@@ -824,14 +707,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);
});
}
@@ -841,14 +717,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);
});
}
@@ -862,7 +731,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}`,
@@ -877,6 +745,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -894,9 +763,7 @@ export class PerformRunExecutionV3Service {
},
});
await enqueueRunExecutionV3(run, tx, {
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
await ResumeRunService.enqueue(run, tx);
});
}
@@ -912,6 +779,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -935,9 +803,7 @@ export class PerformRunExecutionV3Service {
},
});
await enqueueRunExecutionV3(run, tx, {
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
await ResumeRunService.enqueue(run, tx);
});
}
@@ -983,9 +849,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);
});
}
@@ -1037,6 +901,7 @@ export class PerformRunExecutionV3Service {
status: "WAITING",
run: {
update: {
status: "WAITING_TO_CONTINUE",
executionDuration: {
increment: durationInMs,
},
@@ -1056,8 +921,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;
@@ -1066,7 +930,6 @@ export class PerformRunExecutionV3Service {
if (executionDuration >= run.organization.maximumExecutionTimePerRunInMs) {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Execution timed out after ${
@@ -1114,7 +977,6 @@ export class PerformRunExecutionV3Service {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Function timeout detected in ${
@@ -1149,102 +1011,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<string, any>): Promise<void> {
async #failRunExecutionWithRetry(
run: FoundRun,
lastAttempt: boolean,
output: Record<string, any>
): Promise<void> {
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<string, any>,
status: "FAILURE" | "ABORTED" | "TIMED_OUT" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE",
durationInMs: number = 0
): Promise<void> {
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 {
@@ -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<Awaited<ReturnType<typeof findRun>>>;
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,
},
},
},
});
}
@@ -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<Awaited<ReturnType<typeof findRun>>>;
type RunConnectionsByKey = Awaited<ReturnType<typeof createRunConnections>>;
@@ -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);
}
@@ -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<ReturnType<typeof findTask>>;
@@ -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) {
@@ -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"
+17 -1
View File
@@ -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";
import { GraphileMigrationHelperService } from "./db/graphileMigrationHelper.server";
import { DispatchChunkerService } from "./events/dispatchChunker.server";
@@ -109,6 +111,9 @@ const workerCatalog = {
expireDispatcher: z.object({
id: z.string(),
}),
resumeRun: z.object({
id: z.string(),
}),
};
const executionWorkerCatalog = {
@@ -263,7 +268,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();
@@ -456,6 +460,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);
},
},
@@ -477,6 +490,7 @@ function getExecutionWorkerQueue() {
},
shutdownTimeoutInMs: env.GRACEFUL_SHUTDOWN_TIMEOUT,
schema: executionWorkerCatalog,
rateLimiter: executionRateLimiter,
tasks: {
performRunExecutionV2: {
priority: 0, // smaller number = higher priority
@@ -489,6 +503,7 @@ function getExecutionWorkerQueue() {
reason: payload.reason,
resumeTaskId: payload.resumeTaskId,
isRetry: payload.isRetry,
lastAttempt: job.max_attempts === job.attempts,
});
},
},
@@ -505,6 +520,7 @@ function getExecutionWorkerQueue() {
id: payload.id,
reason: payload.reason,
isRetry: false,
lastAttempt: job.max_attempts === job.attempts,
},
driftInMs
);
+3 -12
View File
@@ -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})`
}
+4
View File
@@ -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`;
}
+1
View File
@@ -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",
+19
View File
@@ -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
+19
View File
@@ -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
+3
View File
@@ -36,6 +36,9 @@
<ParamField body="enabled" type="boolean">
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`.
</ParamField>
<ParamField body="concurrencyLimit" type="number | ConcurrencyLimit">
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.
</ParamField>
<ParamField body="onSuccess" type="function">
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.
</ParamField>
+46 -1
View File
@@ -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](#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.
+47 -1
View File
@@ -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
<CardGroup cols={2}>
+63 -1
View File
@@ -33,6 +33,54 @@ const assistant = await io.openai.beta.assistants.create("create-assistant", {
});
```
### `update()`
Update an assistant. [Official OpenAI Docs](https://platform.openai.com/docs/api-reference/assistants/modifyAssistant)
```ts example.ts
const file = await io.openai.files.createAndWaitForProcessing("upload-file", {
purpose: "assistants",
file: fs.createReadStream("./fixtures/mydata.csv"),
});
const assistantId = "asst_abc123";
const assistant = await io.openai.beta.assistants.update("update-assistant", assistantId, {
file_ids: [file.id], // add a file to the assistant
});
```
### `list()`
List assistants. [Official OpenAI Docs](https://platform.openai.com/docs/api-reference/assistants/listAssistants)
```ts example.ts
const assistants = await io.openai.beta.assistants.list("list");
// with pagination
const assistants = await io.openai.beta.assistants.list("list", {
limit: 10,
order: "desc",
after: "asst_abc123",
});
```
### `retrieve()`
Retrieve an assistant. [Official OpenAI Docs](https://platform.openai.com/docs/api-reference/assistants/getAssistant)
```ts example.ts
const assistant = await io.openai.beta.assistants.retrieve("get-assistant", "asst_abc123");
```
### `del()`
Delete an assistant. [Official OpenAI Docs](https://platform.openai.com/docs/api-reference/assistants/deleteAssistant)
```ts example.ts
const deletedAssistant = await io.openai.beta.assistants.del("delete-assistant", "asst_abc123");
```
## Threads
Create threads that assistants can interact with. [Official OpenAI docs](https://platform.openai.com/docs/api-reference/threads/createThread)
@@ -132,12 +180,26 @@ Create messages within threads. [Official OpenAI docs](https://platform.openai.c
### `list()`
List all messages in a thread.
List messages in a thread.
```ts example.ts
const messages = await io.openai.beta.threads.messages.list("list-messages", "thread_abc123");
// with pagination
const messages = await io.openai.beta.threads.messages.list("list-messages", "thread_abc123", {
limit: 10,
order: "desc",
after: "message_abc123",
});
```
If you want to list all messages in a thread, you can use the `listAll()` helper:
```ts example.ts
const messages = await io.openai.beta.threads.messages.listAll("list-messages", "thread_abc123");
```
This will automatically paginate through all messages in the thread and return them as a single array.
### `create()`
Create a message. [Official OpenAI Docs](https://platform.openai.com/docs/api-reference/messages/createMessage)
+39 -11
View File
@@ -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",
@@ -254,7 +256,10 @@
"pages": [
{
"group": "Airtable",
"pages": ["integrations/apis/airtable", "integrations/apis/airtable-tasks"]
"pages": [
"integrations/apis/airtable",
"integrations/apis/airtable-tasks"
]
},
{
"group": "GitHub",
@@ -280,16 +285,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",
@@ -301,7 +315,10 @@
},
{
"group": "Slack",
"pages": ["integrations/apis/slack", "integrations/apis/slack-tasks"]
"pages": [
"integrations/apis/slack",
"integrations/apis/slack-tasks"
]
},
"integrations/apis/stripe",
{
@@ -345,6 +362,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"
]
}
@@ -389,7 +407,10 @@
"sdk/dynamictrigger/constructor",
{
"group": "Instance methods",
"pages": ["sdk/dynamictrigger/register", "sdk/dynamictrigger/unregister"]
"pages": [
"sdk/dynamictrigger/register",
"sdk/dynamictrigger/unregister"
]
}
]
},
@@ -400,7 +421,10 @@
"sdk/dynamicschedule/constructor",
{
"group": "Instance methods",
"pages": ["sdk/dynamicschedule/register", "sdk/dynamicschedule/unregister"]
"pages": [
"sdk/dynamicschedule/register",
"sdk/dynamicschedule/unregister"
]
}
]
},
@@ -412,7 +436,9 @@
},
{
"group": "HTTP Reference",
"pages": ["sdk/api-reference/events/create-an-event"]
"pages": [
"sdk/api-reference/events/create-an-event"
]
},
{
"group": "React SDK",
@@ -426,7 +452,9 @@
},
{
"group": "Overview",
"pages": ["examples/introduction"]
"pages": [
"examples/introduction"
]
}
],
"footerSocials": {
@@ -439,4 +467,4 @@
"apiKey": "phc_hwYmedO564b3Ik8nhA4Csrb5SueY0EwFJWCbseGwWW"
}
}
}
}
@@ -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.
<RequestExample>
```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,
});
```
</RequestExample>
## Parameters
<ParamField body="id" type="string" required>
The ID of the concurrency limit group.
</ParamField>
<ParamField body="limit" type="number" required>
The maximum number of concurrent executions allowed for this group.
</ParamField>
+12
View File
@@ -1,5 +1,17 @@
# @trigger.dev/airtable
## 2.2.8
### Patch Changes
- 067e19fe: - Simplify `Webhook Triggers` and use the new HTTP Endpoints
- Add a `Key-Value Store` for use in and outside of Jobs
- Add a `@trigger.dev/shopify` package
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/airtable",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for airtable",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"airtable": "^0.12.1",
"zod": "3.22.3"
},
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/github
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/github",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official GitHub integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -29,8 +29,8 @@
"@octokit/request": "^6.2.5",
"@octokit/request-error": "^4.0.1",
"@octokit/webhooks": "^10.4.0",
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"octokit": "^2.0.14",
"zod": "3.22.3"
},
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/linear
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/linear",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for @linear/sdk",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -27,8 +27,8 @@
},
"dependencies": {
"@linear/sdk": "^8.0.0",
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"zod": "3.22.3"
},
"engines": {
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/slack
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/openai",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official OpenAI integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -31,8 +31,8 @@
},
"dependencies": {
"openai": "^4.16.1",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.7"
"@trigger.dev/sdk": "workspace:^2.2.8",
"@trigger.dev/integration-kit": "workspace:^2.2.8"
},
"engines": {
"node": ">=18.0.0"
+171 -2
View File
@@ -2,13 +2,13 @@ import { IntegrationTaskKey, Prettify } from "@trigger.dev/sdk";
import { OpenAIRunTask } from "./index";
import { OpenAIIntegrationOptions, OpenAIRequestOptions } from "./types";
import OpenAI from "openai";
import { createTaskOutputProperties, handleOpenAIError } from "./taskUtils";
import { createTaskOutputProperties, handleOpenAIError, isRequestOptions } from "./taskUtils";
export class Assistants {
constructor(
private runTask: OpenAIRunTask,
private options: OpenAIIntegrationOptions
) {}
) { }
async create(
key: IntegrationTaskKey,
@@ -54,4 +54,173 @@ export class Assistants {
handleOpenAIError
);
}
async update(
key: IntegrationTaskKey,
id: string,
params: Prettify<OpenAI.Beta.AssistantUpdateParams>,
options: OpenAIRequestOptions = {}
): Promise<OpenAI.Beta.Assistant> {
return this.runTask(
key,
async (client, task) => {
const { data, response } = await client.beta.assistants
.update(id, params, {
idempotencyKey: task.idempotencyKey,
...options,
})
.withResponse();
const outputProperties = createTaskOutputProperties(undefined, response.headers);
task.outputProperties = [
...(outputProperties ?? []),
{
label: "assistantId",
text: data.id,
},
];
return data;
},
{
name: "Update Assistant",
params,
properties: [
...(params.model ? [{ label: "model", text: params.model }] : []),
...(params.name ? [{ label: "name", text: params.name }] : []),
...(params.file_ids && params.file_ids.length > 0
? [{ label: "files", text: params.file_ids.join(", ") }]
: []),
],
},
handleOpenAIError
);
}
list(
key: IntegrationTaskKey,
params?: Prettify<OpenAI.Beta.AssistantListParams>,
options?: OpenAIRequestOptions,
): Promise<OpenAI.Beta.Assistant[]>;
list(
key: IntegrationTaskKey,
options?: OpenAIRequestOptions,
): Promise<OpenAI.Beta.Assistant[]>;
async list(
key: IntegrationTaskKey,
params: Prettify<OpenAI.Beta.AssistantListParams> | OpenAIRequestOptions = {},
options: OpenAIRequestOptions | undefined = undefined
): Promise<OpenAI.Beta.Assistant[]> {
return this.runTask(
key,
async (client, task) => {
if (isRequestOptions(params)) {
const { data, response } = await client.beta.assistants
.list({
idempotencyKey: task.idempotencyKey,
...params,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return data.data;
}
const { data, response } = await client.beta.assistants
.list(params, {
idempotencyKey: task.idempotencyKey,
...options,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return data.data;
},
{
name: "List Assistants",
params,
properties: !isRequestOptions(params) ? [
...(params.before ? [{ label: "before", text: params.before }] : []),
...(params.order ? [{ label: "order", text: params.order }] : []),
...(params.after ? [{ label: "after", text: params.after }] : []),
...(params.limit ? [{ label: "limit", text: String(params.limit) }] : []),
] : [],
},
handleOpenAIError
);
}
async del(
key: IntegrationTaskKey,
id: string,
options: OpenAIRequestOptions = {}
): Promise<OpenAI.Beta.AssistantDeleted> {
return this.runTask(
key,
async (client, task) => {
const { data, response } = await client.beta.assistants
.del(id, {
idempotencyKey: task.idempotencyKey,
...options,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return data;
},
{
name: "Delete Assistant",
params: {
id,
},
properties: [
{
label: "assistantId",
text: id,
},
],
},
handleOpenAIError
);
}
async retrieve(
key: IntegrationTaskKey,
id: string,
options: OpenAIRequestOptions = {}
): Promise<OpenAI.Beta.Assistant> {
return this.runTask(
key,
async (client, task) => {
const { data, response } = await client.beta.assistants
.retrieve(id, {
idempotencyKey: task.idempotencyKey,
...options,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return data;
},
{
name: "Retrieve Assistant",
params: {
id,
},
properties: [
{
label: "assistantId",
text: id,
},
],
},
handleOpenAIError
);
}
}
+54 -25
View File
@@ -60,11 +60,11 @@ function createTaskUsageProperties(
},
...("completion_tokens" in usage
? [
{
label: "Completion Usage",
text: String(usage.completion_tokens),
},
]
{
label: "Completion Usage",
text: String(usage.completion_tokens),
},
]
: []),
];
}
@@ -83,35 +83,35 @@ function createTaskRateLimitProperties(headers: Headers | undefined) {
return [
...(remainingRequests
? [
{
label: "Remaining Requests",
text: remainingRequests ?? "Unknown",
},
]
{
label: "Remaining Requests",
text: remainingRequests ?? "Unknown",
},
]
: []),
...(resetRequests
? [
{
label: "Reset Requests",
text: resetRequests ?? "Unknown",
},
]
{
label: "Reset Requests",
text: resetRequests ?? "Unknown",
},
]
: []),
...(remainingTokens
? [
{
label: "Remaining Tokens",
text: remainingTokens ?? "Unknown",
},
]
{
label: "Remaining Tokens",
text: remainingTokens ?? "Unknown",
},
]
: []),
...(resetTokens
? [
{
label: "Reset Tokens",
text: resetTokens ?? "Unknown",
},
]
{
label: "Reset Tokens",
text: resetTokens ?? "Unknown",
},
]
: []),
];
}
@@ -282,3 +282,32 @@ export const backgroundTaskRetries: FetchRetryOptions = {
randomize: true,
},
};
type KeysEnum<T> = { [P in keyof Required<T>]: true };
const requestOptionsKeys: KeysEnum<OpenAIRequestOptions> = {
method: true,
path: true,
query: true,
headers: true,
idempotencyKey: true,
};
export const isRequestOptions = (obj: unknown): obj is OpenAIRequestOptions => {
return (
typeof obj === 'object' &&
obj !== null &&
!isEmptyObj(obj) &&
Object.keys(obj).every((k) => hasOwn(requestOptionsKeys, k))
);
};
function isEmptyObj(obj: Object | null | undefined): boolean {
if (!obj) return true;
for (const _k in obj) return false;
return true;
}
function hasOwn(obj: Object, key: string): boolean {
return Object.prototype.hasOwnProperty.call(obj, key);
}
+64 -8
View File
@@ -7,6 +7,7 @@ import {
createBackgroundFetchUrl,
createTaskOutputProperties,
handleOpenAIError,
isRequestOptions,
} from "./taskUtils";
import { RunSubmitToolOutputsParams } from "openai/resources/beta/threads/runs/runs";
import { ThreadUpdateParams } from "openai/resources/beta/threads/threads";
@@ -15,7 +16,7 @@ export class Threads {
constructor(
private runTask: OpenAIRunTask,
private options: OpenAIIntegrationOptions
) {}
) { }
/**
* Create a thread and run it in one task.
@@ -261,7 +262,7 @@ class Runs {
constructor(
private runTask: OpenAIRunTask,
private options: OpenAIIntegrationOptions
) {}
) { }
/**
* Creates a run and waits for it to complete by polling in the background.
@@ -551,15 +552,70 @@ class Messages {
constructor(
private runTask: OpenAIRunTask,
private options: OpenAIIntegrationOptions
) {}
) { }
/**
* Returns messages for a given thread.
*/
list(
key: IntegrationTaskKey,
threadId: string,
params?: Prettify<OpenAI.Beta.Threads.MessageListParams>,
options?: OpenAIRequestOptions
): Promise<OpenAI.Beta.Threads.ThreadMessage[]>
list(
key: IntegrationTaskKey,
threadId: string,
options?: OpenAIRequestOptions
): Promise<OpenAI.Beta.Threads.ThreadMessage[]>
async list(
key: IntegrationTaskKey,
threadId: string,
params: Prettify<OpenAI.Beta.AssistantListParams> | OpenAIRequestOptions = {},
options: OpenAIRequestOptions | undefined = undefined
): Promise<OpenAI.Beta.Threads.ThreadMessage[]> {
return this.runTask(
key,
async (client, task, io) => {
if (isRequestOptions(params)) {
const { data: page, response } = await client.beta.threads.messages
.list(threadId, {
idempotencyKey: task.idempotencyKey,
...params,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return page.data;
}
const { data: page, response } = await client.beta.threads.messages
.list(threadId, params, {
idempotencyKey: task.idempotencyKey,
...options,
})
.withResponse();
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
return page.data;
},
{
name: "List Messages",
properties: [{ label: "threadId", text: threadId }],
},
handleOpenAIError
);
}
/**
* Returns all messages for a given thread.
*/
async list(
async listAll(
key: IntegrationTaskKey,
threadId: string,
options: OpenAIRequestOptions = {}
options: OpenAIRequestOptions = {},
): Promise<OpenAI.Beta.Threads.ThreadMessage[]> {
return this.runTask(
key,
@@ -573,8 +629,8 @@ class Messages {
const allMessages = [];
for await (const fineTuningJob of page) {
allMessages.push(fineTuningJob);
for await (const message of page) {
allMessages.push(message);
}
task.outputProperties = createTaskOutputProperties(undefined, response.headers);
@@ -582,7 +638,7 @@ class Messages {
return allMessages;
},
{
name: "List Messages",
name: "List All Messages",
properties: [{ label: "threadId", text: threadId }],
},
handleOpenAIError
+1 -1
View File
@@ -34,4 +34,4 @@ export type OpenAIRequestOptions = {
path?: string;
headers?: OpenAIHeaders;
idempotencyKey?: string;
};
};
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/plain
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/plain",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official Plain.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"@team-plain/typescript-sdk": "^2.7.0"
},
"engines": {
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/replicate
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/replicate",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for replicate",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"replicate": "^0.18.1",
"zod": "3.22.3"
},
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/resend
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/resend",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official Resend.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"resend": "^2.0.0"
},
"engines": {
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/sendgrid
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/sendgrid",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for @sendgrid/mail",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -27,8 +27,8 @@
},
"dependencies": {
"@sendgrid/mail": "^7.7.0",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.7"
"@trigger.dev/sdk": "workspace:^2.2.8",
"@trigger.dev/integration-kit": "workspace:^2.2.8"
},
"engines": {
"node": ">=16.8.0"
+14
View File
@@ -0,0 +1,14 @@
# @trigger.dev/shopify
## 2.2.8
### Patch Changes
- 067e19fe: - Simplify `Webhook Triggers` and use the new HTTP Endpoints
- Add a `Key-Value Store` for use in and outside of Jobs
- Add a `@trigger.dev/shopify` package
- 096151c0: Fix `@trigger.dev/shopify` imports, enhance docs, and suppress HTTP Endpoint warnings
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/shopify",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for @shopify/shopify-api",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -27,8 +27,8 @@
},
"dependencies": {
"@shopify/shopify-api": "^8.0.2",
"@trigger.dev/sdk": "workspace:^2.2.6",
"@trigger.dev/integration-kit": "workspace:^2.2.6",
"@trigger.dev/sdk": "workspace:^2.2.8",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"zod": "3.22.3"
},
"engines": {
+8
View File
@@ -1,5 +1,13 @@
# @trigger.dev/slack
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+2 -2
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/slack",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official Slack integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,7 +25,7 @@
},
"dependencies": {
"@slack/web-api": "^6.8.1",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.8",
"zod": "3.22.3"
},
"engines": {
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/stripe
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/stripe",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for stripe",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"stripe": "^12.14.0",
"zod": "3.22.3"
},
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/supabase
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/supabase",
"version": "2.2.7",
"version": "2.2.8",
"description": "Trigger.dev integration for @supabase/supabase-js",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -27,8 +27,8 @@
},
"dependencies": {
"@supabase/supabase-js": "^2.26.0",
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"supabase-management-js": "^0.1.4",
"zod": "3.22.3"
},
+9
View File
@@ -1,5 +1,14 @@
# @trigger.dev/typeform
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/integration-kit@2.2.8
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/typeform",
"version": "2.2.7",
"version": "2.2.8",
"description": "The official Typeform integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,8 +25,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.7",
"@trigger.dev/sdk": "workspace:^2.2.7",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"@typeform/api-client": "^1.8.0",
"zod": "3.22.3"
},
+8
View File
@@ -1,5 +1,13 @@
# @trigger.dev/astro
## 2.2.8
### Patch Changes
- Updated dependencies [067e19fe]
- Updated dependencies [096151c0]
- @trigger.dev/sdk@2.2.8
## 2.2.7
### Patch Changes
+2 -2
View File
@@ -1,7 +1,7 @@
{
"name": "@trigger.dev/astro",
"description": "An Astro-native integration for Trigger.dev background jobs platform",
"version": "2.2.7",
"version": "2.2.8",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
"files": [
@@ -20,7 +20,7 @@
"build:tsup": "tsup"
},
"peerDependencies": {
"@trigger.dev/sdk": "workspace:^2.2.7"
"@trigger.dev/sdk": "workspace:^2.2.8"
},
"devDependencies": {
"astro": "^3.0.12",
+10
View File
@@ -1,5 +1,15 @@
# create-trigger
## 2.2.8
### Patch Changes
- 067e19fe: - Simplify `Webhook Triggers` and use the new HTTP Endpoints
- Add a `Key-Value Store` for use in and outside of Jobs
- Add a `@trigger.dev/shopify` package
- Updated dependencies [067e19fe]
- @trigger.dev/core@2.2.8
## 2.2.7
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/cli",
"version": "2.2.7",
"version": "2.2.8",
"description": "The Trigger.dev CLI",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
+2
View File
@@ -1,5 +1,7 @@
# @trigger.dev/core-backend
## 2.2.8
## 2.2.7
## 2.2.6
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/core-backend",
"version": "2.2.7",
"version": "2.2.8",
"description": "Core code used across `@trigger.dev/sdk` and Trigger.dev server",
"license": "MIT",
"main": "./dist/index.js",
+8
View File
@@ -1,5 +1,13 @@
# internal-platform
## 2.2.8
### Patch Changes
- 067e19fe: - Simplify `Webhook Triggers` and use the new HTTP Endpoints
- Add a `Key-Value Store` for use in and outside of Jobs
- Add a `@trigger.dev/shopify` package
## 2.2.7
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/core",
"version": "2.2.7",
"version": "2.2.8",
"description": "Core code used across the Trigger.dev SDK and platform",
"license": "MIT",
"main": "./dist/index.js",
+2 -1
View File
@@ -10,7 +10,8 @@ export async function requestFilterMatches(
return false;
}
if (filter.headers && !eventFilterMatches(clonedRequest.headers, filter.headers)) {
const headersObj = Object.fromEntries(clonedRequest.headers.entries());
if (filter.headers && !eventFilterMatches(headersObj, filter.headers)) {
return false;
}
+7 -2
View File
@@ -273,6 +273,11 @@ export const QueueOptionsSchema = z.object({
export type QueueOptions = z.infer<typeof QueueOptionsSchema>;
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<typeof JobMetadataSchema>;
@@ -398,7 +404,7 @@ export type EndpointIndexError = z.infer<typeof EndpointIndexErrorSchema>;
const IndexEndpointStatsSchema = z.object({
jobs: z.number(),
sources: z.number(),
webhooks: z.number(),
webhooks: z.number().optional(),
dynamicTriggers: z.number(),
dynamicSchedules: z.number(),
disabledJobs: z.number().default(0),
@@ -880,7 +886,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(),
});
+3
View File
@@ -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({
@@ -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';
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "JobRun" ALTER COLUMN "number" DROP NOT NULL;
@@ -0,0 +1,7 @@
-- CreateTable
CREATE TABLE "JobCounter" (
"jobId" TEXT NOT NULL,
"lastNumber" INTEGER NOT NULL DEFAULT 0,
CONSTRAINT "JobCounter_pkey" PRIMARY KEY ("jobId")
);
@@ -0,0 +1,10 @@
-- This is an empty migration.
INSERT INTO
"JobCounter" ("jobId", "lastNumber")
SELECT
"jobId",
MAX(number)
FROM
"JobRun"
GROUP BY
"jobId";
@@ -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;
@@ -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;

Some files were not shown because too many files have changed in this diff Show More