Compare commits

...

38 Commits

Author SHA1 Message Date
Eric Allam c04cdfde8c Fix pnpm lock file 2023-12-07 13:00:32 +00:00
github-actions[bot] 750c1ff1e3 chore: Update version for release (#783)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-12-07 12:59:53 +00:00
Eric Allam 0f9f010206 @trigger.dev/airtable: Removed named imports (#782) 2023-12-07 12:53:58 +00:00
Eric Allam 11f8ff1bb4 Fixing endless retrying of non-successful run executions
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 0s
🚀 Publish Trigger.dev Docker / units (push) Failing after 0s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 0s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-12-06 21:01:41 +00:00
Eric Allam 85b3352764 Fix pnpm lock file 2023-12-06 20:05:49 +00:00
github-actions[bot] 9fb8dceef4 chore: Update version for release (#781)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-12-06 20:05:29 +00:00
Eric Allam 0c14e4cdfe Stop importing subpaths in @trigger.dev/airtable (#780) 2023-12-06 20:02:37 +00:00
D-K-P 2099c6308e Latest lock file 2023-12-06 15:45:43 +00:00
github-actions[bot] 08e6cad28a chore: Update version for release (#779)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-12-06 15:42:36 +00:00
Eric Allam f3efcc0c28 Moved Logger to core-backend, no longer importing node:buffer in core/react (#778) 2023-12-06 15:37:50 +00:00
Eric Allam 37ef335b66 Resolve pnpm lock file 2023-12-06 11:29:58 +00:00
github-actions[bot] 33d555d00e chore: Update version for release (#776)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-12-06 11:28:06 +00:00
Eric Allam 17f6f29d05 Feature: Support multiple runtimes other than Node.js (#774)
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 0s
🚀 Publish Trigger.dev Docker / units (push) Failing after 5s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 5s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-12-06 10:23:37 +00:00
Eric Allam 08f7c639ef Maximum background fetch timeout is now 5 minutes
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 1s
🚀 Publish Trigger.dev Docker / units (push) Failing after 0s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 0s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-12-04 15:02:15 +00:00
github-actions[bot] 1567239718 chore: Update version for release (#771)
* chore: Update version for release

* Release lockfile

---------

Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Co-authored-by: nicktrn <55853254+nicktrn@users.noreply.github.com>
2023-12-04 14:14:41 +00:00
nicktrn de652c1dfb Fix: Shopify tasks and KV get() return types (#770)
* Improve KV get types

* Fix shopify find and all task return types

* Add changeset
2023-12-04 12:37:55 +00:00
Eric Allam 1f3733b70f Fix resuming timed out or auto-yielded executions
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 5s
🚀 Publish Trigger.dev Docker / units (push) Failing after 21s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 21s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-12-04 09:20:40 +00:00
Dan b5aea6c534 Replaced showcase links -> api (#764)
* Updated the examples section with new links / changed showcase link to project showcase

* Added links to the api pages for all of our integrations

* Deleted card from the airtable tasks page
2023-12-01 17:35:54 +00:00
Eric Allam 0769dc4315 Fixed the event names for stripe docs 2023-12-01 08:42:05 +00:00
Eric Allam 5d00fc7cdb Update pnpm lock file for release 2023-11-30 21:54:43 +00:00
github-actions[bot] 00b0c3e02e chore: Update version for release (#765)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-11-30 21:54:13 +00:00
Eric Allam 7e3a82ef47 Issue #747: Add Stripe invoice triggers 2023-11-30 21:48:45 +00:00
James Ritchie 5dda6cd16c Added the border to the page header on the Runs page
🚀 Publish Trigger.dev Docker / units (push) Failing after 4s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 3s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 21s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-11-30 10:32:38 +00:00
Matt Aitken 76b7fb2337 The status API endpoints return only the original set of statuses 2023-11-30 10:31:27 +00:00
Matt Aitken 68cbfd8d23 Log out the run status 2023-11-30 10:29:48 +00:00
Eric Allam 41a49f6bb2 Make the run list query more efficient
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 0s
🚀 Publish Trigger.dev Docker / units (push) Failing after 4s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 4s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-11-29 19:28:15 +00:00
Eric Allam 2f755158b4 Updated pnpm lock file 2023-11-29 11:25:27 +00:00
github-actions[bot] 419b93809d chore: Update version for release (#757)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-11-29 11:24:58 +00:00
Eric Allam bd4bc51daa Update plain (#762)
* updated the plain SDK

* Updated the action to support the new Threads SDK

* latest lock file

---------

Co-authored-by: James Ritchie <james@jamesritchie.co.uk>
2023-11-29 10:58:06 +00:00
Eric Allam ff540c9e4a Revert "Update plain (#761)"
🚀 Publish Trigger.dev Docker / units (push) Failing after 0s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 3s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 5s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
This reverts commit 3f217ff3e0.
2023-11-29 10:51:13 +00:00
Eric Allam 3f217ff3e0 Update plain (#761)
🚀 Publish Trigger.dev Docker / units (push) Failing after 0s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 0s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 5s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
* updated the plain SDK

* Updated the action to support the new Threads SDK

---------

Co-authored-by: James Ritchie <james@jamesritchie.co.uk>
2023-11-29 10:49:37 +00:00
Eric Allam 9cb39bf7d7 Added Redis example env vars 2023-11-29 10:34:54 +00:00
Eric Allam ca05d5f603 @trigger.dev/openai: Adding additional assistant tasks 2023-11-28 22:48:31 +00:00
Hemachandar 4dc46cbbe4 fix: transform Headers object to Record type (#760) 2023-11-28 18:42:50 +00:00
nicktrn 1dcd87a2aa URI encode keys (#758) 2023-11-28 18:05:15 +00:00
Eric Allam c4cb98af5c Clarified the run execution limits 2023-11-28 16:23:35 +00:00
Eric Allam 6ebd435e81 Feature: Run execution concurrency limits (#750)
* WIP execution concurrency controls implemented via Redis

- Split up resuming a run and executing a run
- Added some new statuses to better show what is going on in a run
- Removed preprocessing runs

* WIP

* Convert to using ZSETs and adding env vars

* Removed unused import

* Improve run number generation using advistory locks, and only on start

* More execution concurrency stuff

* Add support for job concurrency limits and concurrency limit groups

* Create wild-swans-battle.md

* Increase slots refresh timeout to 10s

* Try to fix Redis connection issues

* Don’t be so strict about the APP_ENV

* Add the blank tls option to the normal redis client as well

* Add docs
2023-11-28 16:21:06 +00:00
Eric Allam caf203c084 Made webhook index stats optional
🚀 Publish Trigger.dev Docker / units (push) Failing after 4s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 4s
🚀 Publish Trigger.dev Docker / e2e (push) Failing after 19s
🚀 Publish Trigger.dev Docker / publish (push) Has been skipped
2023-11-28 15:04:24 +00:00
273 changed files with 7388 additions and 3474 deletions
+5
View File
@@ -12,6 +12,11 @@ APP_ENV=development
APP_ORIGIN=http://localhost:3030
NODE_ENV=development
# Redis is used for concurrency control
# REDIS_HOST="localhost"
# REDIS_PORT="6379"
# REDIS_TLS_DISABLED="true"
# OPTIONAL VARIABLES
# This is used for validating emails that are allowed to log in. Every email that do not match this regex will be rejected.
# WHITELISTED_EMAILS="authorized@yahoo\.com|authorized@gmail\.com"
+10
View File
@@ -22,6 +22,16 @@ jobs:
node-version: 18
cache: "pnpm"
- name: ⎔ Setup Deno
uses: denoland/setup-deno@v1
with:
deno-version: v1.x
- name: ⎔ Setup bun
uses: oven-sh/setup-bun@v1
with:
bun-version: "1.0.15"
- name: 📥 Download deps
run: pnpm install --frozen-lockfile
+1 -1
View File
@@ -1,4 +1,4 @@
{
"deno.enablePaths": ["references/deno-reference"],
"deno.enablePaths": ["references/deno-reference", "runtime_tests/tests/deno"],
"debug.toolBarLocation": "commandCenter"
}
+45
View File
@@ -1,5 +1,50 @@
# proxy
## 0.0.9
### Patch Changes
- @trigger.dev/core@2.3.3
## 0.0.8
### Patch Changes
- @trigger.dev/core@2.3.2
## 0.0.7
### Patch Changes
- Updated dependencies [f3efcc0c]
- @trigger.dev/core@2.3.1
## 0.0.6
### Patch Changes
- Updated dependencies [17f6f29d]
- @trigger.dev/core@2.3.0
## 0.0.5
### Patch Changes
- @trigger.dev/core@2.2.11
## 0.0.4
### Patch Changes
- @trigger.dev/core@2.2.10
## 0.0.3
### Patch Changes
- Updated dependencies [6ebd435e]
- @trigger.dev/core@2.2.9
## 0.0.2
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "proxy",
"version": "0.0.2",
"version": "0.0.9",
"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">
+1 -1
View File
@@ -8,4 +8,4 @@ export const EXECUTE_JOB_RETRY_LIMIT = 10;
export const MAX_RUN_YIELDED_EXECUTIONS = 100;
export const RUN_CHUNK_EXECUTION_BUFFER = 350;
export const MAX_RUN_CHUNK_EXECUTION_LIMIT = 120000; // 2 minutes
export const RESPONSE_TIMEOUT_STATUS_CODES = [408, 504];
export const VERCEL_RESPONSE_TIMEOUT_STATUS_CODES = [408, 504];
+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>;
+23 -3
View File
@@ -1,4 +1,4 @@
import { RESPONSE_TIMEOUT_STATUS_CODES } from "~/consts";
import { VERCEL_RESPONSE_TIMEOUT_STATUS_CODES } from "~/consts";
import { prisma } from "~/db.server";
import { Prettify } from "~/lib.es5";
@@ -20,13 +20,33 @@ export async function findEndpoint(id: string) {
});
}
export function detectResponseIsTimeout(response?: Response) {
export function detectResponseIsTimeout(rawBody: string, response?: Response) {
if (!response) {
return false;
}
return (
RESPONSE_TIMEOUT_STATUS_CODES.includes(response.status) ||
isResponseVercelTimeout(response) ||
isResponseDenoDeployTimeout(rawBody, response) ||
isResponseCloudflareTimeout(rawBody, response)
);
}
function isResponseCloudflareTimeout(rawBody: string, response: Response) {
return (
response.status === 503 &&
rawBody.includes("Worker exceeded resource limits") &&
typeof response.headers.get("cf-ray") === "string"
);
}
function isResponseVercelTimeout(response: Response) {
return (
VERCEL_RESPONSE_TIMEOUT_STATUS_CODES.includes(response.status) ||
response.headers.get("x-vercel-error") === "FUNCTION_INVOCATION_TIMEOUT"
);
}
function isResponseDenoDeployTimeout(rawBody: string, response: Response) {
return response.status === 502 && rawBody.includes("TIME_LIMIT");
}
+56
View File
@@ -0,0 +1,56 @@
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}`);
}
}
}
export function runOriginalStatus(status: JobRunStatus) {
switch (status) {
case "EXECUTING":
case "WAITING_TO_CONTINUE":
case "WAITING_TO_EXECUTE":
return "STARTED";
default:
return 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
@@ -94,6 +94,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;
@@ -104,6 +109,7 @@ export type ZodWorkerOptions<TMessageCatalog extends MessageCatalogSchema> = {
cleanup?: ZodWorkerCleanupOptions;
reporter?: ZodWorkerReporter;
shutdownTimeoutInMs?: number;
rateLimiter?: ZodWorkerRateLimiter;
};
export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
@@ -116,6 +122,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
#runner?: GraphileRunner;
#cleanup: ZodWorkerCleanupOptions | undefined;
#reporter?: ZodWorkerReporter;
#rateLimiter?: ZodWorkerRateLimiter;
#shutdownTimeoutInMs?: number;
#shuttingDown = false;
@@ -128,6 +135,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
}
@@ -151,6 +159,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
noHandleSignals: true,
taskList: this.#createTaskListFromTasks(),
parsedCronItems,
forbiddenFlags: this.#rateLimiter?.forbiddenFlags.bind(this.#rateLimiter),
});
if (!this.#runner) {
@@ -395,7 +404,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 ?? {})) {
@@ -425,6 +438,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 });
@@ -1,19 +1,20 @@
import { PrismaClient, prisma } from "~/db.server";
import { z } from "zod";
import { PrismaClient, prisma } from "~/db.server";
import { DirectionSchema } from "~/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam._index/route";
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,9 +32,45 @@ export class RunListPresenter {
projectSlug,
direction = "forward",
cursor,
pageSize = DEFAULT_PAGE_SIZE,
}: RunListOptions) {
const directionMultiplier = direction === "forward" ? 1 : -1;
// Find the organization that the user is a member of
const organization = await this.#prismaClient.organization.findFirstOrThrow({
where: {
slug: organizationSlug,
members: { some: { userId } },
},
});
// Find the project scoped to the organization
const project = await this.#prismaClient.project.findFirstOrThrow({
where: {
slug: projectSlug,
organizationId: organization.id,
},
});
// Find all runtimeEnvironments that the user has access to
const environments = await this.#prismaClient.runtimeEnvironment.findMany({
where: {
projectId: project.id,
OR: [
{ orgMember: { userId } },
{ orgMemberId: null },
]
}
});
const job = jobSlug ? await this.#prismaClient.job.findFirstOrThrow({
where: {
slug: jobSlug,
projectId: project.id,
},
}) : undefined;
const runs = await this.#prismaClient.jobRun.findMany({
select: {
id: true,
@@ -41,6 +78,7 @@ export class RunListPresenter {
startedAt: true,
completedAt: true,
createdAt: true,
executionDuration: true,
isTest: true,
status: true,
environment: {
@@ -59,41 +97,34 @@ export class RunListPresenter {
version: true,
},
},
job: {
select: {
slug: true,
title: true,
},
},
},
where: {
job: {
slug: jobSlug,
},
project: {
slug: projectSlug,
},
organization: { slug: organizationSlug, members: { some: { userId } } },
environment: {
OR: [
{
orgMember: null,
},
{
orgMember: {
userId,
},
},
],
jobId: job?.id,
projectId: project.id,
organizationId: organization.id,
environmentId: {
in: environments.map((environment) => environment.id),
},
},
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
? {
id: cursor,
}
id: cursor,
}
: 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 +133,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 +156,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 +165,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,
@@ -82,6 +85,8 @@ export class RunPresenter {
runConnections: run.runConnections,
missingConnections: run.missingConnections,
error: runError,
executionDuration: run.executionDuration,
executionCount: run.executionCount,
};
}
@@ -112,6 +117,8 @@ export class RunPresenter {
isTest: true,
properties: true,
output: true,
executionCount: true,
executionDuration: true,
version: {
select: {
version: true,
@@ -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>
<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>
@@ -3,6 +3,7 @@ import { json } from "@remix-run/server-runtime";
import { GetEvent } from "@trigger.dev/core";
import { z } from "zod";
import { prisma } from "~/db.server";
import { runOriginalStatus } from "~/models/jobRun.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { apiCors } from "~/utils/apiCors";
@@ -49,7 +50,7 @@ function toJSON(eventRecord: FoundEventRecord): GetEvent {
updatedAt: eventRecord.updatedAt,
runs: eventRecord.runs.map((run) => ({
id: run.id,
status: run.status,
status: runOriginalStatus(run.status),
startedAt: run.startedAt,
completedAt: run.completedAt,
})),
@@ -3,6 +3,7 @@ import { json } from "@remix-run/server-runtime";
import { JobRunStatusRecordSchema } from "@trigger.dev/core";
import { z } from "zod";
import { prisma } from "~/db.server";
import { runOriginalStatus } from "~/models/jobRun.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { apiCors } from "~/utils/apiCors";
@@ -66,7 +67,7 @@ export async function loader({ request, params }: LoaderFunctionArgs) {
json({
run: {
id: run.id,
status: run.status,
status: runOriginalStatus(run.status),
output: run.output,
},
statuses: parsedStatuses,
+2 -1
View File
@@ -1,6 +1,7 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { z } from "zod";
import { runOriginalStatus } from "~/models/jobRun.server";
import { ApiRunPresenter } from "~/presenters/ApiRunPresenter.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { apiCors } from "~/utils/apiCors";
@@ -79,7 +80,7 @@ export async function loader({ request, params }: LoaderFunctionArgs) {
request,
json({
id: jobRun.id,
status: jobRun.status,
status: runOriginalStatus(jobRun.status),
startedAt: jobRun.startedAt,
updatedAt: jobRun.updatedAt,
completedAt: jobRun.completedAt,
+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);
+38 -53
View File
@@ -1,12 +1,6 @@
import { parse } from "@conform-to/zod";
import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import {
ComponentDividerSpacingSize,
ComponentSpacerSize,
ComponentTextColor,
ComponentTextSize,
PlainClient,
} from "@team-plain/typescript-sdk";
import { PlainClient, uiComponent } from "@team-plain/typescript-sdk";
import { inspect } from "util";
import { z } from "zod";
import { env } from "~/env.server";
@@ -68,6 +62,8 @@ export async function action({ request }: ActionFunctionArgs) {
onCreate: {
externalId: user.id,
fullName: user.name ?? "",
// TODO - Optional: set 'first name' on user
// shortName: ''
email: {
email: user.email,
isVerified: true,
@@ -76,6 +72,8 @@ export async function action({ request }: ActionFunctionArgs) {
onUpdate: {
externalId: { value: user.id },
fullName: { value: user.name ?? "" },
// TODO - see above
// shortName: { value: "" },
email: {
email: user.email,
isVerified: true,
@@ -96,63 +94,50 @@ export async function action({ request }: ActionFunctionArgs) {
}
const title = feedbackTypeLabel[submission.value.feedbackType as FeedbackType];
const upsertTimelineEntryRes = await client.upsertCustomTimelineEntry({
customerId: upsertCustomerRes.data.customer.id,
const createThreadRes = await client.createThread({
customerIdentifier: {
customerId: upsertCustomerRes.data.customer.id,
},
title,
components: [
{
componentText: {
text: `New ${title} reported by ${user.name} (${user.email})`,
},
},
{
componentDivider: {
dividerSpacingSize: ComponentDividerSpacingSize.M,
},
},
{
componentText: {
textSize: ComponentTextSize.S,
textColor: ComponentTextColor.Muted,
text: "Page",
},
},
{
componentText: {
text: submission.value.path,
},
},
{
componentSpacer: {
spacerSize: ComponentSpacerSize.M,
},
},
{
componentText: {
textSize: ComponentTextSize.S,
textColor: ComponentTextColor.Muted,
text: "Message",
},
},
{
componentText: {
text: submission.value.message,
},
},
uiComponent.text({
text: `New ${title} reported by ${user.name} (${user.email})`,
}),
uiComponent.divider({ spacingSize: "M" }),
uiComponent.text({
size: "S",
color: "MUTED",
text: "Page",
}),
uiComponent.text({
text: submission.value.path,
}),
uiComponent.spacer({ size: "M" }),
uiComponent.text({
size: "S",
color: "MUTED",
text: "Message",
}),
uiComponent.text({
text: submission.value.message,
}),
],
changeCustomerStatusToActive: true,
sendCustomTimelineEntryCreatedNotification: true,
// TODO: Optional: set labels on threads here on creation
// labelTypeIds: [],
// TODO: Optional: set the priority (0 is urgent, 3 is low)
// priority: 0,
});
if (upsertTimelineEntryRes.error) {
if (createThreadRes.error) {
console.error(
inspect(upsertTimelineEntryRes.error, {
inspect(createThreadRes.error, {
showHidden: false,
depth: null,
colors: true,
})
);
submission.error.message = upsertTimelineEntryRes.error.message;
submission.error.message = createThreadRes.error.message;
return json(submission);
}
@@ -1,9 +1,9 @@
import { Logger } from "@trigger.dev/core-backend";
import type { PoolClient } from "pg";
import { z } from "zod";
import { Logger } from "@trigger.dev/core";
import { logger } from "~/services/logger.server";
import { NotificationCatalog, NotificationChannel, notificationCatalog } from "./types";
import { safeJsonParse } from "~/utils/json";
import { NotificationCatalog, NotificationChannel, notificationCatalog } from "./types";
export class PgListenService {
#poolClient: PoolClient;
@@ -444,6 +444,7 @@ function addStandardRequestOptions(options: RequestInit) {
...options.headers,
"user-agent": "triggerdotdev-server/2.0.0",
"x-trigger-version": API_VERSIONS.LAZY_LOADED_CACHED_TASKS,
accept: "application/json",
},
};
}
@@ -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,
@@ -1,4 +1,4 @@
import { MAX_RUN_CHUNK_EXECUTION_LIMIT, RESPONSE_TIMEOUT_STATUS_CODES } from "~/consts";
import { MAX_RUN_CHUNK_EXECUTION_LIMIT } from "~/consts";
import { prisma, PrismaClient } from "~/db.server";
import { EndpointApi } from "../endpointApi.server";
import { logger } from "../logger.server";
@@ -46,8 +46,10 @@ export class ProbeEndpointService {
},
});
const rawBody = await response.text();
// If the response is a 200, or it was a timeout, we can assume the endpoint is up and update the runChunkExecutionLimit
if (response.status === 200 || detectResponseIsTimeout(response)) {
if (response.status === 200 || detectResponseIsTimeout(rawBody, response)) {
await this.#prismaClient.endpoint.update({
where: {
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}
)
`;
}
@@ -9,7 +9,7 @@ const supabase = new SupabaseManagement({
id: "__SLUG__",
});
new Job(client, {
client.defineJob({
id: "on-new-todos",
name: "On New Todos",
version: "0.1.1",
@@ -32,7 +32,7 @@ const supabase = new SupabaseManagement({
apiKey: process.env.SUPABASE_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "on-new-todos",
name: "On New Todos",
version: "0.1.1",
@@ -136,7 +136,7 @@ const supabase = new Supabase<Database>({
supabaseKey: process.env.SUPABASE_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "on-new-users",
name: "On New Users",
version: "0.1.1",
@@ -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) {
+2 -2
View File
@@ -1,5 +1,5 @@
import type { LogLevel } from "@trigger.dev/core";
import { Logger } from "@trigger.dev/core";
import type { LogLevel } from "@trigger.dev/core-backend";
import { Logger } from "@trigger.dev/core-backend";
import { sensitiveDataReplacer } from "./sensitiveDataReplacer";
import { AsyncLocalStorage } from "async_hooks";
@@ -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 }
);
@@ -31,12 +31,6 @@ export class CreateRunService {
},
});
const jobQueue = await this.#prismaClient.jobQueue.findUniqueOrThrow({
where: {
id: version.queueId,
},
});
const eventRecord = await this.#prismaClient.eventRecord.findUniqueOrThrow({
where: {
id: eventId,
@@ -44,22 +38,8 @@ export class CreateRunService {
});
return await $transaction(this.#prismaClient, async (tx) => {
// Get the current max number for the given jobId
const latestJob = await tx.jobRun.findFirst({
where: { jobId: job.id },
orderBy: { id: "desc" },
select: {
number: true,
},
});
// Increment the number for the new execution
const newNumber = (latestJob?.number ?? 0) + 1;
// Create the new execution with the incremented number
const run = await tx.jobRun.create({
data: {
number: newNumber,
preprocess: version.preprocessRuns,
jobId: job.id,
versionId: version.id,
@@ -68,7 +48,6 @@ export class CreateRunService {
organizationId: environment.organizationId,
projectId: environment.projectId,
endpointId: endpoint.id,
queueId: jobQueue.id,
externalAccountId: eventRecord.externalAccountId
? eventRecord.externalAccountId
: undefined,
@@ -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,14 +196,20 @@ 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({
message: `Connection could not be established to the endpoint (${run.endpoint.url})`,
});
return await this.#failRunExecutionWithRetry(
run,
input.lastAttempt,
{
message: `Connection could not be established to the endpoint (${run.endpoint.url})`,
},
durationInMs
);
}
// Update the endpoint version if it has changed
@@ -386,6 +290,8 @@ export class PerformRunExecutionV3Service {
status: response.status,
runId: run.id,
endpoint: run.endpoint.url,
headers: rawHeaders,
rawBody,
});
const errorBody = safeJsonZodParse(errorParser, rawBody);
@@ -393,14 +299,14 @@ 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,
durationInMs
);
}
}
@@ -408,7 +314,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`,
@@ -418,18 +323,22 @@ export class PerformRunExecutionV3Service {
);
} else {
// If the error is a timeout, we should mark this execution as succeeded (by not throwing an error) and enqueue a new execution
if (detectResponseIsTimeout(response)) {
if (detectResponseIsTimeout(rawBody, response)) {
return await this.#resumeRunExecutionAfterTimeout(
this.#prismaClient,
run,
input,
durationInMs,
executionCount
durationInMs
);
} else {
return await this.#failRunExecutionWithRetry({
message: `Endpoint responded with ${response.status} status code`,
});
return await this.#failRunExecutionWithRetry(
run,
input.lastAttempt,
{
message: `Endpoint responded with ${response.status} status code`,
},
durationInMs
);
}
}
}
@@ -439,7 +348,6 @@ export class PerformRunExecutionV3Service {
if (!safeBody) {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
{
message: "Endpoint responded with invalid JSON",
@@ -452,7 +360,6 @@ export class PerformRunExecutionV3Service {
if (!safeBody.success) {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
{
message: generateErrorMessage(safeBody.error.issues),
@@ -491,7 +398,6 @@ export class PerformRunExecutionV3Service {
break;
}
case "CANCELED": {
await this.#cancelExecution(run);
break;
}
case "UNRESOLVED_AUTH_ERROR": {
@@ -644,6 +550,9 @@ export class PerformRunExecutionV3Service {
executionDuration: {
increment: durationInMs,
},
executionCount: {
increment: 1,
},
},
});
@@ -661,17 +570,18 @@ export class PerformRunExecutionV3Service {
run: FoundRun,
data: RunJobResumeWithTask,
durationInMs: number,
executionCount: number = 1
executionCountIncrement: number = 1
) {
return await $transaction(this.#prismaClient, async (tx) => {
await tx.jobRun.update({
where: { id: run.id },
data: {
status: "WAITING_TO_CONTINUE",
executionDuration: {
increment: durationInMs,
},
executionCount: {
increment: executionCount,
increment: executionCountIncrement,
},
},
});
@@ -744,7 +654,6 @@ export class PerformRunExecutionV3Service {
case "ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.error ?? undefined,
"FAILURE",
@@ -754,7 +663,6 @@ export class PerformRunExecutionV3Service {
case "INVALID_PAYLOAD": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.errors,
"INVALID_PAYLOAD",
@@ -774,7 +682,6 @@ export class PerformRunExecutionV3Service {
case "UNRESOLVED_AUTH_ERROR": {
return await this.#failRunExecution(
this.#prismaClient,
"EXECUTE_JOB",
run,
childError.issues,
"UNRESOLVED_AUTH",
@@ -805,14 +712,7 @@ export class PerformRunExecutionV3Service {
});
}
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.error ?? undefined,
"FAILURE",
durationInMs
);
await this.#failRunExecution(tx, execution, data.error ?? undefined, "FAILURE", durationInMs);
});
}
@@ -822,14 +722,7 @@ export class PerformRunExecutionV3Service {
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.issues,
"UNRESOLVED_AUTH",
durationInMs
);
await this.#failRunExecution(tx, execution, data.issues, "UNRESOLVED_AUTH", durationInMs);
});
}
@@ -839,14 +732,7 @@ export class PerformRunExecutionV3Service {
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.errors,
"INVALID_PAYLOAD",
durationInMs
);
await this.#failRunExecution(tx, execution, data.errors, "INVALID_PAYLOAD", durationInMs);
});
}
@@ -860,7 +746,6 @@ export class PerformRunExecutionV3Service {
if (run.yieldedExecutions.length + 1 > MAX_RUN_YIELDED_EXECUTIONS) {
return await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Run has yielded too many times, the maximum is ${MAX_RUN_YIELDED_EXECUTIONS}`,
@@ -875,6 +760,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -892,9 +778,7 @@ export class PerformRunExecutionV3Service {
},
});
await enqueueRunExecutionV3(run, tx, {
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
await ResumeRunService.enqueue(run, tx);
});
}
@@ -910,6 +794,7 @@ export class PerformRunExecutionV3Service {
id: run.id,
},
data: {
status: "WAITING_TO_EXECUTE",
executionDuration: {
increment: durationInMs,
},
@@ -933,9 +818,7 @@ export class PerformRunExecutionV3Service {
},
});
await enqueueRunExecutionV3(run, tx, {
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
await ResumeRunService.enqueue(run, tx);
});
}
@@ -968,6 +851,7 @@ export class PerformRunExecutionV3Service {
],
},
forceYieldImmediately: false,
status: "WAITING_TO_EXECUTE",
},
select: {
executionCount: true,
@@ -981,9 +865,7 @@ export class PerformRunExecutionV3Service {
output: data.output ? (JSON.parse(data.output) as any) : undefined,
});
await enqueueRunExecutionV3(run, tx, {
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
await ResumeRunService.enqueue(run, tx);
});
}
@@ -1035,6 +917,7 @@ export class PerformRunExecutionV3Service {
status: "WAITING",
run: {
update: {
status: "WAITING_TO_CONTINUE",
executionDuration: {
increment: durationInMs,
},
@@ -1054,8 +937,7 @@ export class PerformRunExecutionV3Service {
prisma: PrismaClientOrTransaction,
run: FoundRun,
input: PerformRunExecutionV3Input,
durationInMs: number,
executionCount: number
durationInMs: number
) {
await $transaction(prisma, async (tx) => {
const executionDuration = run.executionDuration + durationInMs;
@@ -1064,7 +946,6 @@ export class PerformRunExecutionV3Service {
if (executionDuration >= run.organization.maximumExecutionTimePerRunInMs) {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Execution timed out after ${
@@ -1112,7 +993,6 @@ export class PerformRunExecutionV3Service {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Function timeout detected in ${
@@ -1133,6 +1013,9 @@ export class PerformRunExecutionV3Service {
executionDuration: {
increment: durationInMs,
},
executionCount: {
increment: 1,
},
endpoint: {
update: {
// Never allow the execution limit to be less than 10 seconds or more than MAX_RUN_CHUNK_EXECUTION_LIMIT
@@ -1143,106 +1026,92 @@ export class PerformRunExecutionV3Service {
},
},
forceYieldImmediately: false,
status: "WAITING_TO_EXECUTE",
},
});
// 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> {
throw new Error(JSON.stringify(output));
async #failRunExecutionWithRetry(
run: FoundRun,
lastAttempt: boolean,
output: Record<string, any>,
durationInMs: number = 0
): Promise<void> {
if (lastAttempt) {
return await this.#failRunExecution(this.#prismaClient, run, output);
}
const updatedJob = await this.#prismaClient.jobRun.update({
where: { id: run.id },
data: {
status: "WAITING_TO_EXECUTE",
executionFailureCount: {
increment: 1,
},
},
});
if (updatedJob.executionFailureCount >= 10) {
return await this.#failRunExecution(this.#prismaClient, run, output);
}
// Use the job.executionFailureCount to determine how long to wait before retrying, using an exponential backoff
const runAt = new Date(Date.now() + Math.pow(1.5, updatedJob.executionFailureCount) * 500); // 500ms, 750ms, 1125ms, 1687ms, 2531ms, 3796ms, 5694ms, 8541ms, 12812ms, 19218ms
await ResumeRunService.enqueue(run, this.#prismaClient, runAt);
}
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,
},
executionCount: {
increment: 1,
},
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);
}
@@ -214,8 +214,8 @@ export class PerformTaskOperationService {
const abortController = new AbortController();
// calculate the actual timeout. If timeoutInMs is undefined, we use the default of 120s
// Also make sure the timeout is at least 1s, but not bigger than 120s
const actualTimeoutInMs = Math.min(Math.max(timeout?.durationInMs ?? 120000, 1000), 120000);
// Also make sure the timeout is at least 1s, but not bigger than 300s
const actualTimeoutInMs = Math.min(Math.max(timeout?.durationInMs ?? 120000, 1000), 300000);
const timeoutId = setTimeout(() => {
abortController.abort();
@@ -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";
const workerCatalog = {
@@ -97,6 +99,9 @@ const workerCatalog = {
expireDispatcher: z.object({
id: z.string(),
}),
resumeRun: z.object({
id: z.string(),
}),
};
const executionWorkerCatalog = {
@@ -225,7 +230,6 @@ function getWorkerQueue() {
"events.invokeDispatcher": {
priority: 0, // smaller number = higher priority
maxAttempts: 6,
queueName: (payload) => `dispatcher:${payload.id}`, // use a queue for a dispatcher so runs are created sequentially
handler: async (payload, job) => {
const service = new InvokeDispatcherService();
@@ -412,6 +416,15 @@ function getWorkerQueue() {
handler: async (payload) => {
const service = new ExpireDispatcherService();
return await service.call(payload.id);
},
},
resumeRun: {
priority: 0,
maxAttempts: 10,
handler: async (payload, job) => {
const service = new ResumeRunService();
return await service.call(payload.id);
},
},
@@ -433,6 +446,7 @@ function getExecutionWorkerQueue() {
},
shutdownTimeoutInMs: env.GRACEFUL_SHUTDOWN_TIMEOUT,
schema: executionWorkerCatalog,
rateLimiter: executionRateLimiter,
tasks: {
performRunExecutionV2: {
priority: 0, // smaller number = higher priority
@@ -445,6 +459,7 @@ function getExecutionWorkerQueue() {
reason: payload.reason,
resumeTaskId: payload.resumeTaskId,
isRetry: payload.isRetry,
lastAttempt: job.max_attempts === job.attempts,
});
},
},
@@ -461,6 +476,7 @@ function getExecutionWorkerQueue() {
id: payload.id,
reason: payload.reason,
isRetry: false,
lastAttempt: job.max_attempts === job.attempts,
},
driftInMs
);
+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`;
}
+3 -2
View File
@@ -62,7 +62,7 @@
"@remix-run/server-runtime": "2.1.0",
"@remix-run/v1-meta": "^0.1.3",
"@tabler/icons-react": "^2.39.0",
"@team-plain/typescript-sdk": "^2.2.0",
"@team-plain/typescript-sdk": "^3.5.0",
"@trigger.dev/companyicons": "^1.5.32",
"@trigger.dev/core": "workspace:*",
"@trigger.dev/core-backend": "workspace:*",
@@ -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",
@@ -187,4 +188,4 @@
"engines": {
"node": ">=16.0.0"
}
}
}
+27 -9
View File
@@ -1,19 +1,37 @@
{
"extends": "./node18.json",
"compilerOptions": {
"lib": ["DOM", "DOM.Iterable", "ES2019"],
"lib": [
"DOM",
"DOM.Iterable",
"ES2019"
],
"paths": {
"@trigger.dev/tsup/*": ["../../config-packages/tsup/src/*"],
"@trigger.dev/tsup": ["../../config-packages/tsup/src/index"],
"@trigger.dev/sdk/*": ["../../packages/trigger-sdk/src/*"],
"@trigger.dev/sdk": ["../../packages/trigger-sdk/src/index"],
"@trigger.dev/integration-kit/*": ["../../packages/integration-kit/src/*"],
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"]
"@trigger.dev/tsup/*": [
"../../config-packages/tsup/src/*"
],
"@trigger.dev/tsup": [
"../../config-packages/tsup/src/index"
],
"@trigger.dev/sdk/*": [
"../../packages/trigger-sdk/src/*"
],
"@trigger.dev/sdk": [
"../../packages/trigger-sdk/src/index"
],
"@trigger.dev/integration-kit/*": [
"../../packages/integration-kit/src/*"
],
"@trigger.dev/integration-kit": [
"../../packages/integration-kit/src/index"
]
},
"declaration": false,
"declarationMap": false,
"baseUrl": ".",
"stripInternal": true
},
"exclude": ["node_modules"]
}
"exclude": [
"node_modules"
]
}
+12 -2
View File
@@ -3,7 +3,17 @@
"version": "0.0.0",
"private": true,
"license": "MIT",
"main": "./src/index.ts",
"types": "./src/index.ts",
"dependencies": {
"esbuild": "^0.19.2",
"tsup": "^8.0.1"
},
"devDependencies": {
"tsup": "7.1.x"
"@types/node": "18",
"typescript": "^5.3.0"
},
"engines": {
"node": ">=18.0.0"
}
}
}
+1
View File
@@ -1,3 +1,4 @@
export { defineConfig } from "tsup";
export { deepMergeOptions } from "./utils";
export { options as integrationOptions } from "./integration";
export { options as packageOptions, defineConfig as defineConfigPackage } from "./package";
+39
View File
@@ -0,0 +1,39 @@
import { Plugin } from "esbuild";
import { Options, defineConfig as defineConfigTSUP } from "tsup";
const restoreNodeProtocolPlugin = (): Plugin => {
return {
name: "node-protocol-plugin-restorer",
setup(build) {
build.onResolve(
{
filter: /node:/,
},
async (args) => {
return { path: args.path, external: true };
}
);
},
};
};
export const options: Options = {
name: "main",
config: "tsconfig.json",
entry: ["./src/index.ts"],
outDir: "./dist",
platform: "node",
format: ["cjs", "esm"],
legacyOutput: false,
sourcemap: true,
clean: true,
bundle: true,
splitting: false,
dts: true,
treeshake: {
preset: "recommended",
},
esbuildPlugins: [restoreNodeProtocolPlugin()],
};
export const defineConfig = defineConfigTSUP(options);
+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>
+3 -3
View File
@@ -5,15 +5,15 @@ To begin, install the necessary packages in your Remix project directory. You ca
<CodeGroup>
```bash npm
npm i @trigger.dev/sdk @trigger.dev/remix
npm i @trigger.dev/sdk@latest @trigger.dev/remix@latest
```
```bash pnpm
pnpm install @trigger.dev/sdk @trigger.dev/remix
pnpm install @trigger.dev/sdk@latest @trigger.dev/remix@latest
```
```bash yarn
yarn add @trigger.dev/sdk @trigger.dev/remix
yarn add @trigger.dev/sdk@latest @trigger.dev/remix@latest
```
</CodeGroup>
+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}>
+433
View File
@@ -0,0 +1,433 @@
---
title: "Hono Quick Start"
sidebarTitle: "Hono.dev"
description: "Start creating Jobs in 5 minutes in your Hono project."
icon: "fire"
---
Hono is a fast & lightweight web framework built on top of Web Standards, and we support using Hono with Trigger.dev on Cloudflare Workers, Bun, Deno, and Node.js.
## Installing Required Packages
To begin, install the necessary packages in your Hono project:
<CodeGroup>
```bash npm
npm add @trigger.dev/sdk@latest @trigger.dev/hono@latest
```
```bash pnpm
pnpm add @trigger.dev/sdk@latest @trigger.dev/hono@latest
```
```bash yarn
yarn add @trigger.dev/sdk@latest @trigger.dev/hono@latest
```
</CodeGroup>
## Obtain the Development Server API Key
To locate your development Server API key, login to the [Trigger.dev
dashboard](https://cloud.trigger.dev) and select the Project you want to
connect to. Then click on the **Environments & API Keys** tab in the left menu.
You can copy your development Server API Key from the field at the top of this page.
(Your development key will start with `tr_dev_`).
## Configure Environment Variables
Add the following environment variables to your `.env` file (or `.dev.vars` file if you are using Cloudflare Workers):
```bash
TRIGGER_API_KEY=<development server api key>
TRIGGER_API_URL=https://api.trigger.dev # change this if you are self-hosting
```
Replace `<development server api key>` with the actual API key obtained from the previous step.
## Enable Node.js compatibility
If you are using Cloudflare Workers, you'll need to enable Node.js compatibility mode in your `wrangler.toml` file:
```toml
compatibility_flags = ["nodejs_compat"]
```
## Add Middelware to Your Hono app
Our `@trigger.dev/hono` package provides two different ways of configuring the necessary middleware needed to connect your Hono app to Trigger.dev. `addMiddleware` which should be used for Cloudflare Workers, and `createMiddleware` which can be used with Bun, Deno, and Node.js.
### Cloudflare Workers
Becase environment variables in Cloudflare Workers aren't available in the global scope, but are instead available only inside the fetch handler, we need to use the `addMiddleware` function to add the necessary middleware to your Hono app.
```ts
import { Hono } from "hono";
import { addMiddleware } from "@trigger.dev/hono";
import { TriggerClient } from "@trigger.dev/sdk";
const app = new Hono<{
Bindings: {
TRIGGER_API_KEY: string;
TRIGGER_API_URL: string;
};
}>();
addMiddleware(app, (env) => {
const client = new TriggerClient({
id: "hono-client",
apiKey: env.TRIGGER_API_KEY,
apiUrl: env.TRIGGER_API_URL,
});
return client;
});
// Your other routes here
export default app;
```
The second argument to `addMiddleware` is a function that receives the environment variables (either from `.dev.vars` in development or from Cloudflare when deployed), and returns a `TriggerClient` instance. This function will be called once per request.
If you want, you can extract our the function that creates the `TriggerClient` instance into a separate file, and import it into your `index.ts` file:
<CodeGroup>
```ts trigger-client.ts
import { TriggerClient } from "@trigger.dev/sdk";
export function triggerClient(apiKey: string, apiUrl: string) {
const client = new TriggerClient({
id: "hono-client",
apiKey,
apiUrl,
});
return client;
}
```
```ts index.ts
import { Hono } from "hono";
import { addMiddleware } from "@trigger.dev/hono";
import { triggerClient } from "./trigger-client";
const app = new Hono<{
Bindings: {
TRIGGER_API_KEY: string;
TRIGGER_API_URL: string;
};
}>();
addMiddleware(app, (env) => triggerClient(env.TRIGGER_API_KEY, env.TRIGGER_API_URL));
// Your other routes here
export default app;
```
</CodeGroup>
Now that you've created the `TriggerClient` and setup the middleware, you can add your first job:
<CodeGroup>
```ts trigger-client.ts
import { TriggerClient } from "@trigger.dev/sdk";
import { exampleJob } from "./jobs";
export function triggerClient(apiKey: string, apiUrl: string) {
const client = new TriggerClient({
id: "hono-client",
apiKey,
apiUrl,
});
exampleJob.attachToClient(client);
return client;
}
```
```ts jobs.ts
import { Job, invokeTrigger } from "@trigger.dev/sdk";
export const exampleJob = new Job({
id: "example-job",
name: "Example Job",
version: "0.0.1",
trigger: invokeTrigger(),
run: async (payload, io, ctx) => {
await io.logger.info("Hello world!", { payload });
return {
message: "Hello world!",
};
},
});
```
</CodeGroup>
As you can see above in `jobs.ts`, we define our first job using the `new Job` constructor from `@trigger.dev/sdk`, and then in `trigger-client.ts` we attach the job to the `TriggerClient` instance.
### Bun
If you are using Bun, you can use the `createMiddleware` function to create the necessary middleware to connect your Hono app to Trigger.dev and define your `TriggerClient` and jobs in the global scope:
```ts
import { createMiddleware } from "@trigger.dev/hono";
import { TriggerClient, invokeTrigger } from "@trigger.dev/sdk";
import { Hono } from "hono";
const app = new Hono();
const client = new TriggerClient({
id: "hono-client",
apiKey: Bun.env.TRIGGER_API_KEY, // Bun.env is available in the global scope
apiUrl: Bun.env.TRIGGER_API_URL, // Bun.env is available in the global scope
});
client.defineJob({
id: "example-job",
name: "Example Job",
version: "0.0.1",
trigger: invokeTrigger(),
run: async (payload, io, ctx) => {
await io.logger.info("Hello world!", { payload });
return {
message: "Hello world!",
};
},
});
app.use("/api/trigger", createMiddleware(client));
// The rest of your routes here
export default app;
```
### Deno
Deno works similarly to Bun, but the imports are slightly different. First, import the `@trigger.dev/sdk` and `@trigger.dev/hono` packages using [npm: specifiers](https://docs.deno.com/runtime/manual/node/npm_specifiers)
```ts index.ts
import { createMiddleware } from "npm:@trigger.dev/hono@latest";
import { TriggerClient, invokeTrigger } from "npm:@trigger.dev/sdk@latest";
import { Hono } from "npm:hono"; // Make sure to use the npm specifier for hono as well
```
Deno doesn't automatically load environment variables from a `.env` file, so you'll need to load them manually using the `dotenv` package:
```ts index.ts
import { load } from "https://deno.land/std@0.208.0/dotenv/mod.ts";
const env = await load();
```
Now we can create the `TriggerClient`, define our jobs, and create the middleware:
```ts index.ts
const app = new Hono();
const client = new TriggerClient({
id: "hono-client",
apiKey: env["TRIGGER_API_KEY"],
apiUrl: env["TRIGGER_API_URL"],
});
client.defineJob({
id: "example-job",
name: "Example Job",
version: "0.0.1",
trigger: invokeTrigger(),
run: async (payload, io, ctx) => {
await io.logger.info("Hello world!", { payload });
return {
message: "Hello world!",
};
},
});
app.use("/api/trigger", createMiddleware(client));
// The rest of your routes here
Deno.serve(app.fetch);
```
### Node.js
Node.js works very similarly to Deno and Bun, in that you can define the `TriggerClient` and jobs in the global scope, and then create the middleware and add it to your Hono app:
```ts index.ts
import "dotenv/config";
import { serve } from "@hono/node-server";
import { Hono } from "hono";
import { createMiddleware } from "@trigger.dev/hono";
import { TriggerClient, invokeTrigger } from "@trigger.dev/sdk";
const app = new Hono();
const client = new TriggerClient({
id: "hono-client",
apiKey: process.env.TRIGGER_API_KEY!,
apiUrl: process.env.TRIGGER_API_URL!,
});
client.defineJob({
id: "example-job",
name: "Example Job",
version: "0.0.1",
trigger: invokeTrigger(),
run: async (payload, io, ctx) => {
await io.logger.info("Hello world!", { payload });
return {
message: "Hello world!",
};
},
});
app.use("/api/trigger", createMiddleware(client));
// Your other routes here
serve(app, (info) => {
console.log(`Listening on port ${info.port}`);
});
```
## Running
### Cloudflare Workers
Run your Hono app locally, like you normally would. For example:
<CodeGroup>
```bash npm
npm run dev
```
```bash pnpm
pnpm run dev
```
```bash yarn
yarn run dev
```
</CodeGroup>
In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
</CodeGroup>
### Bun
Run your Hono app locally, like you normally would. For example:
```bash
bun run index.ts
```
In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787 -H localhost
```
</CodeGroup>
### Deno
Run your Hono app locally, like you normally would. For example:
```bash
deno run --allow-net --allow-read --watch index.ts
```
In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev --client-id hono-client -p 8000 -H 127.0.0.1
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8000 -H 127.0.0.1
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8000 -H 127.0.0.1
```
</CodeGroup>
### Node.js
Run your Hono app locally, like you normally would. For example:
<CodeGroup>
```bash npm
npm run dev
```
```bash pnpm
pnpm run dev
```
```bash yarn
yarn run dev
```
</CodeGroup>
In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev --client-id hono-client -p 8787
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev --client-id hono-client -p 8787
```
</CodeGroup>
+11 -11
View File
@@ -4,18 +4,17 @@ description: "Jobs and code examples you can use to get started."
---
<CardGroup>
<Card title="API catalog with code samples" icon="code" href="https://trigger.dev/apis">
Find code examples for the API you need that you can copy and paste into your projects.
</Card>
<Card
title="Browse our Jobs Showcase"
title="Browse our Project Showcase"
icon="rocket-launch"
href="https://trigger.dev/showcase"
color="#EC4899"
>
The showcase is our library of Jobs. Use them as they are, right out of the box, or customize
them to suit your needs.
</Card>
<Card title="API catalog with code samples" icon="code" href="https://trigger.dev/apis">
Find code examples for the API you need that you can copy and paste into your
projects.
Our library of full-stack projects. A great place to learn more and find inspiration for your
next project.
</Card>
</CardGroup>
@@ -27,9 +26,10 @@ To run them, simply follow the instructions in the README files linked below.
| Project Name | Description | Integrations | Author | Status |
| ---------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------- | --------------------------------------- | ------ |
| [OpenAI text summarizer](https://github.com/triggerdotdev/examples/tree/main/openai-text-summarizer) | An app which uses OpenAI to summarize an article and then post the result to Slack. | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) [Slack](https://trigger.dev/docs/integrations/apis/slack) | Trigger.dev | ✅ |
| [Supabase onboarding emails](https://github.com/triggerdotdev/examples/tree/main/supabase-onboarding-emails) | When a user signs up and confirms their email address, they will receive 3 "onboarding" emails over 2 days using Resend.com and Trigger.dev | [Supabase](https://trigger.dev/docs/integrations/apis/supabase) [Resend](https://trigger.dev/docs/integrations/apis/resend) | Trigger.dev | ✅ |
| [OpenAI text summarizer](https://trigger.dev/showcase/projects/openai-text-summarizer) | An app which uses OpenAI to summarize an article and then post the result to Slack. | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) [Slack](https://trigger.dev/docs/integrations/apis/slack) | Trigger.dev | ✅ |
| [Supabase onboarding emails](https://trigger.dev/showcase/projects/supabase-onboarding-emails) | When a user signs up and confirms their email address, they will receive 3 "onboarding" emails over 2 days using Resend.com and Trigger.dev | [Supabase](https://trigger.dev/docs/integrations/apis/supabase) [Resend](https://trigger.dev/docs/integrations/apis/resend) | Trigger.dev | ✅ |
| [Generate presentation titles using OpenAI](https://github.com/triggerdotdev/examples/tree/main/express-vanilla) | Generate presentation titles using OpenAI background jobs with Node.js, Express and Trigger.dev | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) | [lirantal](https://github.com/lirantal) | ✅ |
| [Send a basic email with Resend](https://github.com/triggerdotdev/examples/tree/main/resend-email-form) | Send a basic email from a form with Resend | [Resend](https://trigger.dev/docs/integrations/apis/resend) | Trigger.dev | ✅ |
| [Send a basic email with Resend](https://trigger.dev/showcase/projects/resend-email-form) | Send a basic email from a form with Resend | [Resend](https://trigger.dev/docs/integrations/apis/resend) | Trigger.dev | ✅ |
| [AI changelog generator](https://autochangelog.dev/) | Generates a changelog from your GitHub commits using OpenAI | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) [GitHub](https://trigger.dev/docs/integrations/apis/github) | Trigger.dev | ✅ |
| [AI avatar generator](https://trigger.dev/showcase/projects/avatar-generator) | Turn yourself into a superhero using AI | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) | Trigger.dev | ✅ |
| AI landing page copy generator | Copies your site and generates new copy using OpenAI | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) | Trigger.dev | 🛠️ |
| AI changelog generator | Generates a changelog from your GitHub commits using OpenAI | [OpenAI](https://trigger.dev/docs/integrations/apis/openai) [GitHub](https://trigger.dev/docs/integrations/apis/github) | Trigger.dev | 🛠️ |
@@ -161,3 +161,4 @@ client.defineJob({
},
});
```
+8 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Our Airtable integration allows you to easily connect to the Airtable API and perform tasks such as creating / updating / deleting single or multiple records in your tables.
<Card
title="Jobs Showcase - Airtable"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=airtable"
>
Check out pre-built Airtable jobs in our showcase.
</Card>
## Installation
To get started with our Airtable integration, you need to install the `@trigger.dev/airtable` packages. You can do this using `npm`, `pnpm`, or `yarn`:
@@ -70,6 +62,14 @@ const airtable = new Airtable({
Once you have set up a Airtable client, you can use it to create tasks.
{" "}
<Card title="Tasks" icon="sparkles" href="/integrations/apis/airtable-tasks">
Perform tasks such as creating / updating / deleting single or multiple records in table.
</Card>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/airtable">
Check out pre-built jobs using Airtable in our API section.
</Card>
+6 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Our GitHub integration allows you to create triggers and tasks that interact with GitHub. Trigger jobs when events happen, like when a new issue is added to a repo, or when a pull request is opened, etc. You can also perform tasks like creating issues, getting information about a repo, adding comments, and much more.
<Card
title="Jobs Showcase - GitHub"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=github"
>
Check out pre-built GitHub jobs in our showcase.
</Card>
## Installing the GitHub packages
To get started with our GitHub integration, you need to install the `@trigger.dev/github` packages. You can do this using `npm`, `pnpm`, or `yarn`:
@@ -80,3 +72,9 @@ Once you have set up a GitHub client, you can use it to create triggers and task
Perform tasks such as creating a new issue or a new comment.
</Card>
</CardGroup>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/github">
Check out pre-built jobs using GitHub in our API section.
</Card>
+6
View File
@@ -463,3 +463,9 @@ client.defineJob({
| `createProjectMilestone` | Creates a project milestone. |
| `issuePriorityValues` | Gets issue priority values and labels. |
| `viewer` | Gets the currently authenticated user. |
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/linear">
Check out pre-built jobs using Linear in our API section.
</Card>
+6 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Our OpenAI integration allows you to easily perform AI-powered tasks, such as summarizing text, answering questions, generating images, fine tuning and much more.
<Card
title="Jobs Showcase - OpenAI"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=openai"
>
Check out pre-built OpenAI jobs in our showcase.
</Card>
## Installing the OpenAI packages
<CodeGroup>
@@ -173,3 +165,9 @@ client.defineJob({
And you'll get the same experience in the Run Dashboard when viewing the logs:
![Perplexity.ai logs](/images/perplexity.png)
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/openai">
Check out pre-built jobs using OpenAI in our API section.
</Card>
+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)
+6 -8
View File
@@ -8,14 +8,6 @@ sidebarTitle: Overview & authentication
Plain is the customer support tool for technical teams and products.
It aims to bring engineering and customer service teams together by creating a modern opinionated platform that's fantastic to build with.
<Card
title="Jobs Showcase - Plain"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=plain"
>
Check out pre-built Plain jobs in our showcase.
</Card>
## Installing the Plain packages
<CodeGroup>
@@ -56,3 +48,9 @@ Once you have set up a Plain client, you can use it to create tasks.
Perform tasks such as creating/updating customers and adding timeline entries.
</Card>
</CardGroup>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/plain">
Check out pre-built jobs using Plain in our API section.
</Card>
+6
View File
@@ -168,3 +168,9 @@ client.defineJob({
| `paginate` | Pagination helper that returns an async generator. |
| `request` | Sends authenticated requests to the Replicate API. |
| `run` | Creates and waits for a prediction. |
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/replicate">
Check out pre-built jobs using Replicate in our API section.
</Card>
+6 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Resend is the email API for developers. With our Resend integration you can send email campaigns, transactional emails, and automated emails (drip campaigns) from your app.
<Card
title="Jobs Showcase - Resend"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=resend"
>
Check out pre-built Resend jobs in our showcase.
</Card>
## Installing the Resend packages
To get started with our Resend integration, you need to install the `@trigger.dev/resend` packages. You can do this using `npm`, `pnpm`, or `yarn`:
@@ -55,3 +47,9 @@ Once you have set up a Resend client, you can use it to create tasks.
<Card title="Tasks" icon="sparkles" href="/integrations/apis/resend-tasks">
Send emails with Resend.
</Card>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/resend">
Check out pre-built jobs using Resend in our API section.
</Card>
+6 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
SendGrid is a cloud-based SMTP provider that allows you to send email without having to maintain email servers. With our SendGrid integration you can send email campaigns, transactional emails, and automated emails (drip campaigns) from your app.
<Card
title="Jobs Showcase - SendGrid"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=sendgrid"
>
Check out pre-built SendGrid jobs in our showcase.
</Card>
## Installing the SendGrid packages
To get started with our SendGrid integration, you need to install the `@trigger.dev/sendgrid` packages. You can do this using `npm`, `pnpm`, or `yarn`:
@@ -57,3 +49,9 @@ Once you have set up a SendGrid client, you can use it to create tasks.
<Card title="Tasks" icon="sparkles" href="/integrations/apis/sendgrid-tasks">
Send emails with SendGrid.
</Card>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/sendgrid">
Check out pre-built jobs using SendGrid in our API section.
</Card>
+9 -11
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Our Shopify integration allows you to create triggers and tasks that interact with Shopify. Trigger jobs when events happen, like when a new product is added to a shop, or when an order is paid for, etc. You can also perform tasks like creating products, editing variants, getting information about an order, and a lot more.
{/* <Card
title="Jobs Showcase - Shopify"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=shopify"
>
Check out pre-built Shopify jobs in our showcase.
</Card> */}
## Installing the Shopify packages
To get started with our Shopify integration, you need to install the `@trigger.dev/shopify` packages. You can do this using `npm`, `pnpm`, or `yarn`:
@@ -43,13 +35,13 @@ It's **required** to import the correct Runtime Adapter for your platform. All e
```ts
// Import the Node.js adapter
import '@shopify/shopify-api/adapters/node';
import "@shopify/shopify-api/adapters/node";
// Import the CloudFlare Worker adapter
import '@shopify/shopify-api/adapters/cf-worker';
import "@shopify/shopify-api/adapters/cf-worker";
// Import the generic Web API adapter
import '@shopify/shopify-api/adapters/web-api';
import "@shopify/shopify-api/adapters/web-api";
```
You can then import and use `@trigger.dev/shopify` like any other integration:
@@ -102,3 +94,9 @@ Once you have set up a Shopify client, you can use it to create triggers and tas
Perform Tasks such as creating new variants, or editing orders, and more.
</Card>
</CardGroup>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/shopify">
Check out pre-built jobs using Shopify in our API section.
</Card>
-2
View File
@@ -67,8 +67,6 @@ client.defineJob({
Use their [Block kit builder](https://api.slack.com/block-kit), and then use the `blocks` property to send the message.
To see this in action, check out our 'Daily Slack alert for Linear issues' [example job](https://trigger.dev/showcase/jobs/linearIssuesDailySlackAlert).
```ts linearIssuesDailySlackAlert.ts
...
await io.slack.postMessage("post message", {
+6 -8
View File
@@ -7,14 +7,6 @@ sidebarTitle: Overview & authentication
Our Slack integration allows you to connect to the Slack API and post messages to Slack.
<Card
title="Jobs Showcase - Slack"
icon="rocket"
href="https://trigger.dev/showcase?tags=&apis=slack"
>
Check out pre-built Slack jobs in our showcase.
</Card>
## Installation
<CodeGroup>
@@ -58,3 +50,9 @@ Once you have set up a Slack client, you can use it to create tasks.
<Card title="Tasks" icon="sparkles" href="/integrations/apis/slack-tasks">
Perform tasks such as posting messages to a channel.
</Card>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/slack">
Check out pre-built jobs using Slack in our API section.
</Card>
+76 -55
View File
@@ -85,61 +85,76 @@ client.defineJob({
Available triggers are listed below:
| Function Name | Payload Object | Events | Aggregate Version |
| ---------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------ |
| `onCharge` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.succeeded`, `charge.failed`, `charge.captured`, `charge.refunded`, `charge.updated` | ✔️ |
| `onChargeSucceeded` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.succeeded` | `onCharge` |
| `onChargeFailed` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.failed` | `onCharge` |
| `onChargeCaptured` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.captured` | `onCharge` |
| `onChargeRefunded` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.refunded` | `onCharge` |
| `onChargeUpdated` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.updated` | `onCharge` |
| `onProduct` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.created`, `product.updated`, `product.deleted` | ✔️ |
| `onProductCreated` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.created` | `onProduct` |
| `onProductUpdated` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.updated` | `onProduct` |
| `onProductDeleted` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.deleted` | `onProduct` |
| `onPrice` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.created`, `price.updated`, `price.deleted` | ✔️ |
| `onPriceCreated` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.created` | `onPrice` |
| `onPriceUpdated` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.updated` | `onPrice` |
| `onPriceDeleted` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.deleted` | `onPrice` |
| `onCheckoutSession` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.completed`, `checkout.session.async_payment_succeeded`, `checkout.session.async_payment_failed`, `checkout.session.expired` | ✔️ |
| `onCheckoutSessionCompleted` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.completed` | `onCheckoutSession` |
| `onCheckoutSessionExpired` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.expired` | `onCheckoutSession` |
| `onCustomerSubscription` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.created`, `customer.subscription.updated`, `customer.subscription.deleted`, `customer.subscription.paused`, `customer.subscription.pending_updated_applied`, `customer.subscription.pending_update_expired`, `customer.subscription.resumed` | ✔️ |
| `onCustomerSubscriptionCreated` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.created` | `onCustomerSubscription` |
| `onCustomerSubscriptionUpdated` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.updated` | `onCustomerSubscription` |
| `onCustomerSubscriptionDeleted` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.deleted` | `onCustomerSubscription` |
| `onCustomerSubscriptionPaused` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.paused` | `onCustomerSubscription` |
| `onCustomerSubscriptionPending` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.pending_updated_applied`, `customer.subscription.pending_update_expired` | `onCustomerSubscription` |
| `onCustomerSubscriptionResumed` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.resumed` | `onCustomerSubscription` |
| `onCustomer` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.created`, `customer.updated`, `customer.deleted` | ✔️ |
| `onCustomerCreated` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.created` | `onCustomer` |
| `onCustomerUpdated` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.updated` | `onCustomer` |
| `onCustomerDeleted` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.deleted` | `onCustomer` |
| `onExternalAccount` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.created`, `account.external_account.updated`, `account.external_account.deleted` | ✔️ |
| `onExternalAccountCreated` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.created` | `onExternalAccount` |
| `onExternalAccountUpdated` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.updated` | `onExternalAccount` |
| `onExternalAccountDeleted` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.deleted` | `onExternalAccount` |
| `onPerson` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.created`, `account.person.updated`, `account.person.deleted` | ✔️ |
| `onPersonCreated` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.created` | `onPerson` |
| `onPersonUpdated` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.updated` | `onPerson` |
| `onPersonDeleted` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.deleted` | `onPerson` |
| `onAccountUpdated` | [Account](https://stripe.com/docs/api/events/types#account_object) | `account.updated` | N/A |
| `onPaymentIntent` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.created`, `payment_intent.succeeded`, `payment_intent.payment_failed`, `payment_intent.canceled`, `payment_intent.processing`, `payment_intent.amount_capturable_updated`, `payment_intent.requires_action`, `payment_intent.partially_funded` | ✔️ |
| `onPaymentIntentCreated` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.created` | `onPaymentIntent` |
| `onPaymentIntentSucceeded` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.succeeded` | `onPaymentIntent` |
| `onPaymentIntentPaymentFailed` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.payment_failed` | `onPaymentIntent` |
| `onPaymentIntentCanceled` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.canceled` | `onPaymentIntent` |
| `onPaymentIntentProcessing` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.processing` | `onPaymentIntent` |
| `onPaymentIntentRequiresAction` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.requires_action` | `onPaymentIntent` |
| `onPaymentIntentAmountCapturableUpdated` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.amount_capturable_updated` | `onPaymentIntent` |
| `onPaymentIntentPartiallyFunded` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.partially_funded` | `onPaymentIntent` |
| `onPayout` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.created`, `payout.updated`, `payout.canceled`, `payout.paid`, `payout.failed`, `payout.reconciliation_completed` | ✔️ |
| `onPayoutCreated` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.created` | `onPayout` |
| `onPayoutUpdated` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.updated` | `onPayout` |
| `onPayoutCanceled` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.canceled` | `onPayout` |
| `onPayoutPaid` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.paid` | `onPayout` |
| `onPayoutFailed` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.failed` | `onPayout` |
| `onPayoutReconciliationCompleted` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.reconciliation_completed` | `onPayout` |
| Function Name | Payload Object | Events | Aggregate Version |
| ---------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------ |
| `onCharge` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.succeeded`, `charge.failed`, `charge.captured`, `charge.refunded`, `charge.updated` | ✔️ |
| `onChargeSucceeded` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.succeeded` | `onCharge` |
| `onChargeFailed` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.failed` | `onCharge` |
| `onChargeCaptured` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.captured` | `onCharge` |
| `onChargeRefunded` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.refunded` | `onCharge` |
| `onChargeUpdated` | [Charge](https://stripe.com/docs/api/events/types#charge_object) | `charge.updated` | `onCharge` |
| `onProduct` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.created`, `product.updated`, `product.deleted` | ✔️ |
| `onProductCreated` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.created` | `onProduct` |
| `onProductUpdated` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.updated` | `onProduct` |
| `onProductDeleted` | [Product](https://stripe.com/docs/api/events/types#product_object) | `product.deleted` | `onProduct` |
| `onPrice` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.created`, `price.updated`, `price.deleted` | ✔️ |
| `onPriceCreated` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.created` | `onPrice` |
| `onPriceUpdated` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.updated` | `onPrice` |
| `onPriceDeleted` | [Price](https://stripe.com/docs/api/events/types#price_object) | `price.deleted` | `onPrice` |
| `onCheckoutSession` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.completed`, `checkout.session.async_payment_succeeded`, `checkout.session.async_payment_failed`, `checkout.session.expired` | ✔️ |
| `onCheckoutSessionCompleted` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.completed` | `onCheckoutSession` |
| `onCheckoutSessionExpired` | [CheckoutSession](https://stripe.com/docs/api/events/types#checkout_session_object) | `checkout.session.expired` | `onCheckoutSession` |
| `onCustomerSubscription` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.created`, `customer.subscription.updated`, `customer.subscription.deleted`, `customer.subscription.paused`, `customer.subscription.pending_updated_applied`, `customer.subscription.pending_update_expired`, `customer.subscription.resumed` | ✔️ |
| `onCustomerSubscriptionCreated` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.created` | `onCustomerSubscription` |
| `onCustomerSubscriptionUpdated` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.updated` | `onCustomerSubscription` |
| `onCustomerSubscriptionDeleted` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.deleted` | `onCustomerSubscription` |
| `onCustomerSubscriptionPaused` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.paused` | `onCustomerSubscription` |
| `onCustomerSubscriptionPending` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.pending_updated_applied`, `customer.subscription.pending_update_expired` | `onCustomerSubscription` |
| `onCustomerSubscriptionResumed` | [Subscription](https://stripe.com/docs/api/events/types#subscription_object) | `customer.subscription.resumed` | `onCustomerSubscription` |
| `onCustomer` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.created`, `customer.updated`, `customer.deleted` | ✔️ |
| `onCustomerCreated` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.created` | `onCustomer` |
| `onCustomerUpdated` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.updated` | `onCustomer` |
| `onCustomerDeleted` | [Customer](https://stripe.com/docs/api/events/types#customer_object) | `customer.deleted` | `onCustomer` |
| `onExternalAccount` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.created`, `account.external_account.updated`, `account.external_account.deleted` | ✔️ |
| `onExternalAccountCreated` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.created` | `onExternalAccount` |
| `onExternalAccountUpdated` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.updated` | `onExternalAccount` |
| `onExternalAccountDeleted` | [Card](https://stripe.com/docs/api/events/types#account_card_object) or [Bank Account](https://stripe.com/docs/api/events/types#account_bank_account_object) | `account.external_account.deleted` | `onExternalAccount` |
| `onPerson` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.created`, `account.person.updated`, `account.person.deleted` | ✔️ |
| `onPersonCreated` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.created` | `onPerson` |
| `onPersonUpdated` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.updated` | `onPerson` |
| `onPersonDeleted` | [Person](https://stripe.com/docs/api/events/types#person_object) | `account.person.deleted` | `onPerson` |
| `onAccountUpdated` | [Account](https://stripe.com/docs/api/events/types#account_object) | `account.updated` | N/A |
| `onPaymentIntent` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.created`, `payment_intent.succeeded`, `payment_intent.payment_failed`, `payment_intent.canceled`, `payment_intent.processing`, `payment_intent.amount_capturable_updated`, `payment_intent.requires_action`, `payment_intent.partially_funded` | ✔️ |
| `onPaymentIntentCreated` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.created` | `onPaymentIntent` |
| `onPaymentIntentSucceeded` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.succeeded` | `onPaymentIntent` |
| `onPaymentIntentPaymentFailed` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.payment_failed` | `onPaymentIntent` |
| `onPaymentIntentCanceled` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.canceled` | `onPaymentIntent` |
| `onPaymentIntentProcessing` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.processing` | `onPaymentIntent` |
| `onPaymentIntentRequiresAction` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.requires_action` | `onPaymentIntent` |
| `onPaymentIntentAmountCapturableUpdated` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.amount_capturable_updated` | `onPaymentIntent` |
| `onPaymentIntentPartiallyFunded` | [PaymentIntent](https://stripe.com/docs/api/events/types#payment_intent_object) | `payment_intent.partially_funded` | `onPaymentIntent` |
| `onPayout` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.created`, `payout.updated`, `payout.canceled`, `payout.paid`, `payout.failed`, `payout.reconciliation_completed` | ✔️ |
| `onPayoutCreated` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.created` | `onPayout` |
| `onPayoutUpdated` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.updated` | `onPayout` |
| `onPayoutCanceled` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.canceled` | `onPayout` |
| `onPayoutPaid` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.paid` | `onPayout` |
| `onPayoutFailed` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.failed` | `onPayout` |
| `onPayoutReconciliationCompleted` | [Payout](https://stripe.com/docs/api/events/types#payout_object) | `payout.reconciliation_completed` | `onPayout` |
| `onInvoice` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.created`, `invoice.finalized`, `invoice.finalization_failed`, `invoice.deleted`, `invoice.marked_uncollectible`, `invoice.paid`, `invoice.payment_action_required`, `invoice.payment_failed`, `invoice.payment_succeeded`, `invoice.sent`, `invoice.upcoming`, `invoice.voided`, `invoiceitem.created`, `invoiceitem.deleted` | ✔️ |
| `onInvoiceCreated` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.created` | `onInvoice` |
| `onInvoiceFinalized` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.finalized` | `onInvoice` |
| `onInvoiceFinalizationFailed` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.finalization_failed` | `onInvoice` |
| `onInvoiceDeleted` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.deleted` | `onInvoice` |
| `onInvoiceMarkedUncollectible` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.marked_uncollectible` | `onInvoice` |
| `onInvoicePaid` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.paid` | `onInvoice` |
| `onInvoicePaymentActionRequired` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.payment_action_required` | `onInvoice` |
| `onInvoicePaymentFailed` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.payment_failed` | `onInvoice` |
| `onInvoicePaymentSucceeded` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.payment_succeeded` | `onInvoice` |
| `onInvoiceSent` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.sent` | `onInvoice` |
| `onInvoiceUpcoming` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.upcoming` | `onInvoice` |
| `onInvoiceVoided` | [Invoice](https://stripe.com/docs/api/invoices/object) | `invoice.voided` | `onInvoice` |
| `onInvoiceItemCreated` | [InvoiceItem](https://stripe.com/docs/api/invoiceitems/object) | `invoiceitem.created` | N/A |
| `onInvoiceItemDeleted` | [InvoiceItem](https://stripe.com/docs/api/invoiceitems/object) | `invoiceitem.deleted` | N/A |
If there are any triggers missing that you'd like to see added, please [open a new GitHub Issue](https://github.com/triggerdotdev/trigger.dev/issues/new)
@@ -260,3 +275,9 @@ client.defineJob({
```
Make sure to pass the `idempotencyKey` to the underlying client to ensure that the API call is only executed once. This is only needed for mutating API calls.
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/stripe">
Check out pre-built jobs using Stripe in our API section.
</Card>
@@ -73,3 +73,9 @@ Our Supabase package supports two different integrations: One for the Supabase M
[service_role](https://supabase.com/docs/guides/api/api-keys#the-servicerole-key) key.
</Card>
</CardGroup>
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/supabase">
Check out pre-built jobs using Supabase in our API section.
</Card>
+6
View File
@@ -165,3 +165,9 @@ client.defineJob({
},
});
```
## Example jobs
<Card title="Code examples" icon="code" href="https://trigger.dev/apis/typeform">
Check out pre-built jobs using Typeform in our API section.
</Card>
+41 -12
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",
@@ -89,6 +91,7 @@
"documentation/quickstarts/astro",
"documentation/quickstarts/nuxt",
"documentation/quickstarts/sveltekit",
"documentation/quickstarts/hono",
"documentation/quickstarts/fastify"
]
},
@@ -139,7 +142,7 @@
"group": "Guides",
"pages": [
{
"group": "Platforms",
"group": "Frameworks",
"pages": [
"documentation/guides/platforms/nextjs",
"documentation/guides/platforms/express",
@@ -253,7 +256,10 @@
"pages": [
{
"group": "Airtable",
"pages": ["integrations/apis/airtable", "integrations/apis/airtable-tasks"]
"pages": [
"integrations/apis/airtable",
"integrations/apis/airtable-tasks"
]
},
{
"group": "GitHub",
@@ -279,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",
@@ -300,7 +315,10 @@
},
{
"group": "Slack",
"pages": ["integrations/apis/slack", "integrations/apis/slack-tasks"]
"pages": [
"integrations/apis/slack",
"integrations/apis/slack-tasks"
]
},
"integrations/apis/stripe",
{
@@ -344,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"
]
}
@@ -388,7 +407,10 @@
"sdk/dynamictrigger/constructor",
{
"group": "Instance methods",
"pages": ["sdk/dynamictrigger/register", "sdk/dynamictrigger/unregister"]
"pages": [
"sdk/dynamictrigger/register",
"sdk/dynamictrigger/unregister"
]
}
]
},
@@ -399,7 +421,10 @@
"sdk/dynamicschedule/constructor",
{
"group": "Instance methods",
"pages": ["sdk/dynamicschedule/register", "sdk/dynamicschedule/unregister"]
"pages": [
"sdk/dynamicschedule/register",
"sdk/dynamicschedule/unregister"
]
}
]
},
@@ -411,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",
@@ -425,7 +452,9 @@
},
{
"group": "Overview",
"pages": ["examples/introduction"]
"pages": [
"examples/introduction"
]
}
],
"footerSocials": {
@@ -438,4 +467,4 @@
"apiKey": "phc_hwYmedO564b3Ik8nhA4Csrb5SueY0EwFJWCbseGwWW"
}
}
}
}
+5 -9
View File
@@ -9,8 +9,8 @@ You can define a job by using the `TriggerClient.defineJob` instance method:
<RequestExample>
```ts client.defineJob
client.defineJob({
```ts Example
new Job({
id: "github-integration-on-issue",
name: "GitHub Integration - On Issue",
version: "0.1.0",
@@ -23,11 +23,11 @@ client.defineJob({
await io.logger.info("This is a simple log info message");
return { payload, ctx };
},
});
}).attachToClient(client);
```
```ts notifications
client.defineJob({
new Job({
id: "github-integration-on-issue",
name: "GitHub Integration - On Issue",
version: "0.1.0",
@@ -46,7 +46,7 @@ client.defineJob({
await io.logger.info("This is a simple log info message");
return { payload, ctx };
},
});
}).attachToClient(client);
```
</RequestExample>
@@ -55,10 +55,6 @@ client.defineJob({
## Parameters
<ParamField body="client" type="object" required>
An instance of [TriggerClient](/sdk/triggerclient) that is used to send events to the Trigger API.
</ParamField>
<Snippet file="jobs/options.mdx" />
## Returns
@@ -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>
+63
View File
@@ -1,5 +1,68 @@
# @trigger.dev/airtable
## 2.3.3
### Patch Changes
- 0f9f0102: Named exports don't work because Airtable is a CommonJS module
- @trigger.dev/integration-kit@2.3.3
- @trigger.dev/sdk@2.3.3
## 2.3.2
### Patch Changes
- 0c14e4cd: Fixed importing package subpath
- @trigger.dev/integration-kit@2.3.2
- @trigger.dev/sdk@2.3.2
## 2.3.1
### Patch Changes
- Updated dependencies [f3efcc0c]
- @trigger.dev/sdk@2.3.1
- @trigger.dev/integration-kit@2.3.1
## 2.3.0
### Minor Changes
- 17f6f29d: Support for Deno, Bun and Cloudflare workers, as well as conditionally exporting ESM versions of the package instead of just commonjs.
Cloudflare worker support requires the node compat flag turned on (https://developers.cloudflare.com/workers/runtime-apis/nodejs/)
### Patch Changes
- Updated dependencies [17f6f29d]
- @trigger.dev/integration-kit@2.3.0
- @trigger.dev/sdk@2.3.0
## 2.2.11
### Patch Changes
- de652c1d: Fix Shopify task types and KV `get()` return types
- Updated dependencies [de652c1d]
- @trigger.dev/sdk@2.2.11
- @trigger.dev/integration-kit@2.2.11
## 2.2.10
### Patch Changes
- @trigger.dev/integration-kit@2.2.10
- @trigger.dev/sdk@2.2.10
## 2.2.9
### Patch Changes
- Updated dependencies [1dcd87a2]
- Updated dependencies [6ebd435e]
- @trigger.dev/sdk@2.2.9
- @trigger.dev/integration-kit@2.2.9
## 2.2.8
### Patch Changes
+20 -9
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/airtable",
"version": "2.2.8",
"version": "2.3.3",
"description": "Trigger.dev integration for airtable",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -8,16 +8,15 @@
"access": "public"
},
"files": [
"dist/index.js",
"dist/index.d.ts",
"dist/index.js.map"
"dist"
],
"devDependencies": {
"@trigger.dev/tsconfig": "workspace:*",
"@types/node": "16.x",
"rimraf": "^3.0.2",
"tsup": "7.1.x",
"typescript": "4.9.4"
"tsup": "8.0.1",
"typescript": "^5.3.0",
"@trigger.dev/tsup": "workspace:*"
},
"scripts": {
"clean": "rimraf dist",
@@ -26,12 +25,24 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"@trigger.dev/integration-kit": "workspace:^2.3.3",
"@trigger.dev/sdk": "workspace:^2.3.3",
"airtable": "^0.12.1",
"zod": "3.22.3"
},
"engines": {
"node": ">=16.8.0"
}
},
"exports": {
".": {
"import": {
"types": "./dist/index.d.mts",
"default": "./dist/index.mjs"
},
"require": "./dist/index.js",
"types": "./dist/index.d.ts"
},
"./package.json": "./package.json"
},
"module": "./dist/index.mjs"
}
+5 -6
View File
@@ -1,14 +1,13 @@
import { DisplayProperty, IOWithIntegrations, IntegrationTaskKey } from "@trigger.dev/sdk";
import { FieldSet, Records, SelectOptions } from "airtable";
import { AirtableFieldSet, AirtableRecord, AirtableRunTask, CreateAirtableRecord } from ".";
import { QueryParams } from "airtable/lib/query_params";
import { DisplayProperty, IntegrationTaskKey } from "@trigger.dev/sdk";
import AirtableSDK from "airtable";
import { AirtableFieldSet, AirtableRecord, AirtableRunTask } from ".";
type TableParams<Params extends Record<string, unknown>> = {
tableName: string;
} & Params;
export type AirtableRecordsParams = TableParams<{}>;
export type AirtableRecords = Records<FieldSet>;
export type AirtableRecords = AirtableSDK.Records<AirtableSDK.FieldSet>;
export class Base {
constructor(
@@ -32,7 +31,7 @@ export class Table<TFields extends AirtableFieldSet> {
this.tableName = tableName;
}
getRecords(key: IntegrationTaskKey, params?: SelectOptions<TFields>) {
getRecords(key: IntegrationTaskKey, params?: AirtableSDK.SelectOptions<TFields>) {
return this.runTask(
key,
async (client) => {
+5 -12
View File
@@ -10,19 +10,12 @@ import {
type RunTaskOptions,
type TriggerIntegration,
} from "@trigger.dev/sdk";
import AirtableSDK, { Error as AirtableApiError } from "airtable";
import AirtableSDK from "airtable";
import { Base } from "./base";
import * as events from "./events";
import {
WebhookChangeType,
WebhookDataType,
Webhooks,
createWebhookSource,
createWebhookTrigger,
} from "./webhooks";
import { Webhooks, createWebhookSource } from "./webhooks";
export * from "./types";
export * from "./base";
export * from "./types";
export type AirtableIntegrationOptions = {
/** An ID for this client */
@@ -138,12 +131,12 @@ export class Airtable implements TriggerIntegration {
}
}
function isAirtableApiError(error: unknown): error is AirtableApiError {
function isAirtableApiError(error: unknown): error is AirtableSDK.Error {
if (typeof error !== "object" || error === null) {
return false;
}
const airtableError = error as AirtableApiError;
const airtableError = error as AirtableSDK.Error;
return (
typeof airtableError.error === "string" &&
-2
View File
@@ -1,5 +1,3 @@
import { Prettify } from "@trigger.dev/integration-kit";
export type AirtableFieldSet = {
[key: string]:
| undefined
+14 -5
View File
@@ -1,11 +1,12 @@
import { EventFilter, IntegrationTaskKey, verifyRequestSignature } from "@trigger.dev/sdk";
import AirtableSDK, { Error as AirtableApiError } from "airtable";
import AirtableSDK from "airtable";
import { z } from "zod";
import * as events from "./events";
import { Airtable, AirtableRunTask } from "./index";
import { ListWebhooksResponse, ListWebhooksResponseSchema } from "./schemas";
import { WebhookSource, WebhookTrigger } from "@trigger.dev/sdk/triggers/webhook";
import { registerJobNamespace } from "@trigger.dev/integration-kit/webhooks";
import { WebhookSource, WebhookTrigger } from "@trigger.dev/sdk";
import { registerJobNamespace } from "@trigger.dev/integration-kit";
import { Buffer } from "node:buffer";
const WebhookFromSourceSchema = z.union([
z.literal("formSubmission"),
@@ -315,6 +316,10 @@ export function createWebhookSource(
delete: async ({ io, ctx }) => {
const webhookId = await io.store.job.get<string>("get-webhook-id", "webhook-id");
if (!webhookId) {
throw new Error("Missing webhook ID for delete operation.");
}
await io.integration.webhooks().delete("delete-webhook", {
baseId: ctx.params?.baseId,
webhookId,
@@ -327,6 +332,10 @@ export function createWebhookSource(
`${registerJobNamespace(ctx.key)}:webhook-secret-base64`
);
if (!secretBase64) {
throw new Error("Missing secret for verification.");
}
return await verifyRequestSignature({
request,
headerName: "x-airtable-content-mac",
@@ -442,7 +451,7 @@ async function handleWebhookError(response: Response, errorType: string) {
const parsedErrorBody = AirtableErrorBodySchema.safeParse(rawErrorBody);
if (!parsedErrorBody.success) {
throw new AirtableApiError(
throw new AirtableSDK.Error(
`${errorType}_PARSE_ERROR`,
`${response.statusText}:\n${rawErrorBody}`,
response.status
@@ -451,5 +460,5 @@ async function handleWebhookError(response: Response, errorType: string) {
const { type, message } = parsedErrorBody.data;
throw new AirtableApiError(type, message ?? response.statusText, response.status);
throw new AirtableSDK.Error(type, message ?? response.statusText, response.status);
}
+3 -1
View File
@@ -7,7 +7,9 @@
"@trigger.dev/sdk/*": ["../../packages/trigger-sdk/src/*"],
"@trigger.dev/sdk": ["../../packages/trigger-sdk/src/index"],
"@trigger.dev/integration-kit/*": ["../../packages/integration-kit/src/*"],
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"]
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"],
"@trigger.dev/tsup/*": ["../../config-packages/tsup/src/*"],
"@trigger.dev/tsup": ["../../config-packages/tsup/src/index"]
},
"declaration": false,
"declarationMap": false,
+2 -23
View File
@@ -1,24 +1,3 @@
import { defineConfigPackage } from "@trigger.dev/tsup";
import { defineConfig } from "tsup";
export default defineConfig([
{
name: "main",
entry: ["./src/index.ts"],
outDir: "./dist",
platform: "node",
format: ["cjs"],
legacyOutput: true,
sourcemap: true,
clean: true,
bundle: true,
splitting: false,
dts: true,
treeshake: {
preset: "smallest",
},
esbuildPlugins: [],
external: ["http", "https", "util", "events", "tty", "os", "timers"],
},
]);
export default defineConfigPackage;
+60
View File
@@ -1,5 +1,65 @@
# @trigger.dev/github
## 2.3.3
### Patch Changes
- @trigger.dev/integration-kit@2.3.3
- @trigger.dev/sdk@2.3.3
## 2.3.2
### Patch Changes
- @trigger.dev/integration-kit@2.3.2
- @trigger.dev/sdk@2.3.2
## 2.3.1
### Patch Changes
- Updated dependencies [f3efcc0c]
- @trigger.dev/sdk@2.3.1
- @trigger.dev/integration-kit@2.3.1
## 2.3.0
### Minor Changes
- 17f6f29d: Support for Deno, Bun and Cloudflare workers, as well as conditionally exporting ESM versions of the package instead of just commonjs.
Cloudflare worker support requires the node compat flag turned on (https://developers.cloudflare.com/workers/runtime-apis/nodejs/)
### Patch Changes
- Updated dependencies [17f6f29d]
- @trigger.dev/integration-kit@2.3.0
- @trigger.dev/sdk@2.3.0
## 2.2.11
### Patch Changes
- Updated dependencies [de652c1d]
- @trigger.dev/sdk@2.2.11
- @trigger.dev/integration-kit@2.2.11
## 2.2.10
### Patch Changes
- @trigger.dev/integration-kit@2.2.10
- @trigger.dev/sdk@2.2.10
## 2.2.9
### Patch Changes
- Updated dependencies [1dcd87a2]
- Updated dependencies [6ebd435e]
- @trigger.dev/sdk@2.2.9
- @trigger.dev/integration-kit@2.2.9
## 2.2.8
### Patch Changes
+20 -8
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/github",
"version": "2.2.8",
"version": "2.3.3",
"description": "The official GitHub integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -8,9 +8,7 @@
"access": "public"
},
"files": [
"dist/index.js",
"dist/index.d.ts",
"dist/index.js.map"
"dist"
],
"devDependencies": {
"@octokit/types": "^9.2.3",
@@ -18,7 +16,9 @@
"@trigger.dev/tsconfig": "workspace:*",
"@types/node": "18",
"rimraf": "^3.0.2",
"tsup": "^6.5.0"
"tsup": "8.0.1",
"@trigger.dev/tsup": "workspace:*",
"typescript": "^5.3.0"
},
"scripts": {
"clean": "rimraf dist",
@@ -29,12 +29,24 @@
"@octokit/request": "^6.2.5",
"@octokit/request-error": "^4.0.1",
"@octokit/webhooks": "^10.4.0",
"@trigger.dev/integration-kit": "workspace:^2.2.8",
"@trigger.dev/sdk": "workspace:^2.2.8",
"@trigger.dev/integration-kit": "workspace:^2.3.3",
"@trigger.dev/sdk": "workspace:^2.3.3",
"octokit": "^2.0.14",
"zod": "3.22.3"
},
"engines": {
"node": ">=16.8.0"
}
},
"exports": {
".": {
"import": {
"types": "./dist/index.d.mts",
"default": "./dist/index.mjs"
},
"require": "./dist/index.js",
"types": "./dist/index.d.ts"
},
"./package.json": "./package.json"
},
"module": "./dist/index.mjs"
}
+1
View File
@@ -4,6 +4,7 @@ import type { Logger } from "@trigger.dev/sdk";
import { ExternalSource, HandlerEvent } from "@trigger.dev/sdk";
import { z } from "zod";
import { Github } from "./index";
import { Buffer } from "node:buffer";
type WebhookData = {
id: number;
+3 -1
View File
@@ -7,7 +7,9 @@
"@trigger.dev/sdk/*": ["../../packages/trigger-sdk/src/*"],
"@trigger.dev/sdk": ["../../packages/trigger-sdk/src/index"],
"@trigger.dev/integration-kit/*": ["../../packages/integration-kit/src/*"],
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"]
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"],
"@trigger.dev/tsup/*": ["../../config-packages/tsup/src/*"],
"@trigger.dev/tsup": ["../../config-packages/tsup/src/index"]
},
"declaration": false,
"declarationMap": false,
+2 -21
View File
@@ -1,22 +1,3 @@
import { defineConfig } from "tsup";
import { defineConfigPackage } from "@trigger.dev/tsup";
export default defineConfig([
{
name: "main",
entry: ["./src/index.ts"],
outDir: "./dist",
platform: "node",
format: ["cjs"],
legacyOutput: true,
sourcemap: true,
clean: true,
bundle: true,
splitting: false,
dts: true,
treeshake: {
preset: "smallest",
},
esbuildPlugins: [],
external: ["http", "https", "util", "events", "tty", "os", "timers"],
},
]);
export default defineConfigPackage;
+60
View File
@@ -1,5 +1,65 @@
# @trigger.dev/linear
## 2.3.3
### Patch Changes
- @trigger.dev/integration-kit@2.3.3
- @trigger.dev/sdk@2.3.3
## 2.3.2
### Patch Changes
- @trigger.dev/integration-kit@2.3.2
- @trigger.dev/sdk@2.3.2
## 2.3.1
### Patch Changes
- Updated dependencies [f3efcc0c]
- @trigger.dev/sdk@2.3.1
- @trigger.dev/integration-kit@2.3.1
## 2.3.0
### Minor Changes
- 17f6f29d: Support for Deno, Bun and Cloudflare workers, as well as conditionally exporting ESM versions of the package instead of just commonjs.
Cloudflare worker support requires the node compat flag turned on (https://developers.cloudflare.com/workers/runtime-apis/nodejs/)
### Patch Changes
- Updated dependencies [17f6f29d]
- @trigger.dev/integration-kit@2.3.0
- @trigger.dev/sdk@2.3.0
## 2.2.11
### Patch Changes
- Updated dependencies [de652c1d]
- @trigger.dev/sdk@2.2.11
- @trigger.dev/integration-kit@2.2.11
## 2.2.10
### Patch Changes
- @trigger.dev/integration-kit@2.2.10
- @trigger.dev/sdk@2.2.10
## 2.2.9
### Patch Changes
- Updated dependencies [1dcd87a2]
- Updated dependencies [6ebd435e]
- @trigger.dev/sdk@2.2.9
- @trigger.dev/integration-kit@2.2.9
## 2.2.8
### Patch Changes

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