Compare commits

...

33 Commits

Author SHA1 Message Date
Eric Allam dfb00e84e9 Fixing pnpm lock file 2023-08-08 07:35:25 +01:00
github-actions[bot] b5a64545bb chore: Update version for release (#273)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-08-08 07:32:24 +01:00
Eric Allam 3cf3eaff48 @trigger.dev/supabase: You can now trigger on multiple database events in the same job 2023-08-07 17:27:49 +01:00
Matt Aitken 4ca758a3de job-catalog: added initial setup and CLI build step
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-08-07 15:07:58 +01:00
Matt Aitken da81537009 Update job-catalog instructions
`pnpm run trigger:dev` should be `pnpm run dev:trigger`
2023-08-07 14:57:08 +01:00
Eric Allam 7c1b13ba3b Increase webhook delivery attempts to the max 25 attempts over 3 days 2023-08-07 14:25:23 +01:00
D-K-P 4b654e5b3b Fixed broken docs link 2023-08-07 13:29:11 +01:00
Eric Allam 726a50ace7 Point to the workspace version of @trigger.dev/sdk 2023-08-07 13:17:36 +01:00
github-actions[bot] 9deffb67c4 chore: Update version for release (#267)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-08-07 13:08:27 +01:00
Matt Aitken ff04bf44ee CLI: detect package manager using artifacts if they exist (#261)
* Detect package manager from artifacts (e.g. package-lock.json) and fallback to process.env.npm_config_user_agent

* Added a changeset
2023-08-07 13:01:22 +01:00
Matt Aitken e740297829 Cli: improve Next.js project detection (#262)
* Detect presence of a “next” dependency

* Use a strict undefined check instead

* Changeset: Detect Next.js project by looking at dependencies, not next.config.js

* Read a package json file

* First check for next.config file, otherwise use next dependency

* Update changeset description
2023-08-07 13:01:03 +01:00
D-K-P a31705e198 Added text summarizer / github issue reminder jobs to the examples table 2023-08-07 12:58:43 +01:00
Eric Allam f0bdf53364 Auto index production endpoints every 10 minutes 2023-08-07 12:49:37 +01:00
Matt Aitken 9638499163 Add your first endpoint in the dashboard (#269)
* Created the sheet

* Started work on the resource route

* Created ValidateCreateEndpointService and the EndpointValidateApi (name TBC)

* The TriggerClient responds with the id

* Tidied some stuff up

* Alternative: EndpointApi doesn’t take endpointSlug. Instead just Ping() does

* Create cool-snakes-deny.md

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2023-08-07 11:28:34 +01:00
Matt Aitken 9104c252b7 Telemetry (new user) (#268)
* Removed old telemetry

* Started adding TriggerClient telemetry

* Added a log if the telemetry client is created, event renamed to “user.created”
2023-08-07 11:28:13 +01:00
Matt Aitken 72cf345602 Added “thumbsRating” to docs pages 2023-08-07 10:50:02 +01:00
jemiluv8 d5b8f8299d update cli init to create index file in examples folder (#225)
* update cli init to create index file in examples folder

* add patch changeset

* Create silly-baboons-join.md

* remove comment in router file
add comment in jobs index to instruct them to

* rename examplesIndex to examplesIndexFileName

* import jobs module in routeContent for createTriggerPageRoute just as we did for createTriggerAppRoute

---------

Co-authored-by: Matt Aitken <matt@mattaitken.com>
2023-08-06 16:42:32 +01:00
Matt Aitken dd22b88d25 Latest lockfile… 🙄 2023-08-03 22:52:46 +01:00
github-actions[bot] 62b9c5879c chore: Update version for release (#248) 2023-08-03 22:47:23 +01:00
Eric Allam 1b0973fbc1 Fixed error when configuring a new endpoint failed and fixed ping to now throw an error when parsing JSON
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-08-03 16:02:16 +01:00
Eric Allam fc78854ed4 @trigger.dev/cli: Don't require the TRIGGER_API_URL to be set (it has a default) 2023-08-03 14:49:46 +01:00
Eric Allam ed15bb0618 Another small supabase tweak 2023-08-03 13:09:56 +01:00
Eric Allam ced9034428 Another small supabase doc fix 2023-08-03 13:08:42 +01:00
Eric Allam 7ec329d5d6 Fixed broken image in supabase docs 2023-08-03 13:06:23 +01:00
Eric Allam 81180999de The dev command should use a POST request when doing the PING to the local server 2023-08-03 13:01:08 +01:00
Eric Allam feaa79ecf3 @trigger.dev/supabase: Automatically enable database webhooks when using triggers 2023-08-03 13:00:29 +01:00
Eric Allam d34dc0847c More OpenAI examples 2023-08-03 09:55:07 +01:00
Eric Allam 1b7c7520b4 new Job -> client.defineJob 2023-08-03 09:53:15 +01:00
Eric Allam 34d77e830c Incorporate improved OpenAI docs from old repo 2023-08-03 09:51:25 +01:00
Eric Allam 87c302bf26 Document how to enable supabase database webhooks 2023-08-03 09:39:58 +01:00
Eric Allam d462c901f5 A couple of quickstart tweaks 2023-08-02 17:16:39 +01:00
D-K-P 55c3d79cc8 Deleted the info box on the integration pages 2023-08-02 17:11:41 +01:00
Eric Allam ce901c92e4 Update @trigger.dev/supabase to the latest supabase-management-js 2023-08-02 17:01:02 +01:00
116 changed files with 2176 additions and 852 deletions
@@ -51,11 +51,6 @@ export function NoIntegrationSheet({
)}
</SheetHeader>
<SheetBody>
<Callout variant="info">
We dont have an Integration for the {api.name} API yet but you can request one by
clicking the button above. In the meantime, connect to {api.name} using one of the
methods below.
</Callout>
<CustomHelp name={api.name} />
</SheetBody>
</SheetContent>
+2
View File
@@ -20,6 +20,8 @@ const EnvironmentSchema = z.object({
.default(process.env.NODE_ENV),
SECRET_STORE: SecretStoreOptionsSchema.default("DATABASE"),
POSTHOG_PROJECT_KEY: z.string().optional(),
TELEMETRY_TRIGGER_API_KEY: z.string().optional(),
TELEMETRY_TRIGGER_API_URL: z.string().optional(),
HIGHLIGHT_PROJECT_ID: z.string().optional(),
AUTH_GITHUB_CLIENT_ID: z.string().optional(),
AUTH_GITHUB_CLIENT_SECRET: z.string().optional(),
+99 -1
View File
@@ -1,4 +1,6 @@
import type {
CronItem,
CronItemOptions,
Job as GraphileJob,
Runner as GraphileRunner,
JobHelpers,
@@ -7,7 +9,7 @@ import type {
TaskList,
TaskSpec,
} from "graphile-worker";
import { run as graphileRun } from "graphile-worker";
import { run as graphileRun, parseCronItems } from "graphile-worker";
import omit from "lodash.omit";
import { z } from "zod";
@@ -18,6 +20,13 @@ export interface MessageCatalogSchema {
[key: string]: z.ZodFirstPartySchemaTypes | z.ZodDiscriminatedUnion<any, any>;
}
const RawCronPayloadSchema = z.object({
_cron: z.object({
ts: z.coerce.date(),
backfilled: z.boolean(),
}),
});
const GraphileJobSchema = z.object({
id: z.coerce.string(),
queue_name: z.string().nullable(),
@@ -50,6 +59,19 @@ export type ZodTasks<TConsumerSchema extends MessageCatalogSchema> = {
};
};
type RecurringTaskPayload = {
ts: Date;
backfilled: boolean;
};
export type ZodRecurringTasks = {
[key: string]: {
pattern: string;
options?: CronItemOptions;
handler: (payload: RecurringTaskPayload, job: GraphileJob) => Promise<void>;
};
};
export type ZodWorkerEnqueueOptions = TaskSpec & {
tx?: PrismaClientOrTransaction;
};
@@ -59,6 +81,7 @@ export type ZodWorkerOptions<TMessageCatalog extends MessageCatalogSchema> = {
prisma: PrismaClient;
schema: TMessageCatalog;
tasks: ZodTasks<TMessageCatalog>;
recurringTasks?: ZodRecurringTasks;
};
export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
@@ -66,6 +89,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
#prisma: PrismaClient;
#runnerOptions: RunnerOptions;
#tasks: ZodTasks<TMessageCatalog>;
#recurringTasks?: ZodRecurringTasks;
#runner?: GraphileRunner;
constructor(options: ZodWorkerOptions<TMessageCatalog>) {
@@ -73,6 +97,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
this.#prisma = options.prisma;
this.#runnerOptions = options.runnerOptions;
this.#tasks = options.tasks;
this.#recurringTasks = options.recurringTasks;
}
public async initialize(): Promise<boolean> {
@@ -84,9 +109,12 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
runnerOptions: this.#runnerOptions,
});
const parsedCronItems = parseCronItems(this.#createCronItemsFromRecurringTasks());
this.#runner = await graphileRun({
...this.#runnerOptions,
taskList: this.#createTaskListFromTasks(),
parsedCronItems,
});
if (!this.#runner) {
@@ -192,9 +220,38 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
taskList[key] = task;
}
for (const [key] of Object.entries(this.#recurringTasks ?? {})) {
const task: Task = (payload, helpers) => {
return this.#handleRecurringTask(key, payload, helpers);
};
taskList[key] = task;
}
return taskList;
}
#createCronItemsFromRecurringTasks() {
const cronItems: CronItem[] = [];
if (!this.#recurringTasks) {
return cronItems;
}
for (const [key, task] of Object.entries(this.#recurringTasks)) {
const cronItem: CronItem = {
pattern: task.pattern,
identifier: key,
task: key,
options: task.options,
};
cronItems.push(cronItem);
}
return cronItems;
}
async #handleMessage<K extends keyof TMessageCatalog>(
typeName: K,
rawPayload: unknown,
@@ -226,4 +283,45 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
await task.handler(payload, job);
}
async #handleRecurringTask(
typeName: string,
rawPayload: unknown,
helpers: JobHelpers
): Promise<void> {
const job = helpers.job;
logger.debug("Received recurring task, calling handler", {
type: String(typeName),
payload: rawPayload,
job,
});
const recurringTask = this.#recurringTasks?.[typeName];
if (!recurringTask) {
throw new Error(`No recurring task for message type: ${String(typeName)}`);
}
const parsedPayload = RawCronPayloadSchema.safeParse(rawPayload);
if (!parsedPayload.success) {
throw new Error(
`Failed to parse recurring task payload: ${JSON.stringify(parsedPayload.error)}`
);
}
const payload = parsedPayload.data;
try {
await recurringTask.handler(payload._cron, job);
} catch (error) {
logger.error("Failed to handle recurring task", {
error,
payload,
});
throw error;
}
}
}
@@ -0,0 +1,121 @@
import { conform, useForm, useInputEvent } from "@conform-to/react";
import { parse } from "@conform-to/zod";
import { useFetcher } from "@remix-run/react";
import { InlineCode } from "~/components/code/InlineCode";
import { Button, ButtonContent } from "~/components/primitives/Buttons";
import { FormError } from "~/components/primitives/FormError";
import { Header1, Header2 } from "~/components/primitives/Headers";
import { Hint } from "~/components/primitives/Hint";
import { Input } from "~/components/primitives/Input";
import { InputGroup } from "~/components/primitives/InputGroup";
import { Paragraph } from "~/components/primitives/Paragraph";
import {
Sheet,
SheetBody,
SheetContent,
SheetHeader,
SheetTrigger,
} from "~/components/primitives/Sheet";
import { TextLink } from "~/components/primitives/TextLink";
import { docsPath } from "~/utils/pathBuilder";
import { bodySchema } from "../resources.projects.$projectId.endpoint";
import { RuntimeEnvironment, RuntimeEnvironmentType } from "@trigger.dev/database";
import {
Select,
SelectContent,
SelectGroup,
SelectItem,
SelectTrigger,
SelectValue,
} from "~/components/primitives/Select";
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
import { useRef, useState } from "react";
type FirstEndpointSheetProps = {
projectId: string;
environments: { id: string; type: RuntimeEnvironmentType }[];
};
export function FirstEndpointSheet({ projectId, environments }: FirstEndpointSheetProps) {
const setEndpointUrlFetcher = useFetcher();
const [form, { url, environmentId }] = useForm({
id: "new-endpoint-url",
lastSubmission: setEndpointUrlFetcher.data,
onValidate({ formData }) {
return parse(formData, { schema: bodySchema });
},
});
const loadingEndpointUrl = setEndpointUrlFetcher.state !== "idle";
return (
<Sheet>
<SheetTrigger>
<ButtonContent variant={"primary/medium"}>Add your first endpoint</ButtonContent>
</SheetTrigger>
<SheetContent size="lg">
<SheetHeader>
<div>
<Header1>Add your first endpoint</Header1>
<Paragraph variant="small">
We recommend you use{" "}
<TextLink href={docsPath("documentation/guides/cli")}>the CLI</TextLink> when working
in development.
</Paragraph>
</div>
</SheetHeader>
<SheetBody>
<setEndpointUrlFetcher.Form
method="post"
action={`/resources/projects/${projectId}/endpoint`}
{...form.props}
>
<InputGroup className="mb-4 max-w-none">
<Header2>Environment type</Header2>
<SelectGroup>
<Select name={"environmentId"} defaultValue={environments[0].id}>
<SelectTrigger size="secondary/small">
<SelectValue placeholder="Select environment" className="m-0 p-0" /> Environment
</SelectTrigger>
<SelectContent>
{environments.map((environment) => (
<SelectItem key={environment.id} value={environment.id}>
<EnvironmentLabel environment={environment} />
</SelectItem>
))}
</SelectContent>
</Select>
</SelectGroup>
<FormError id={environmentId.errorId}>{environmentId.error}</FormError>
</InputGroup>
<InputGroup className="max-w-none">
<Header2>Endpoint URL</Header2>
<div className="flex items-center">
<Input
className="rounded-r-none"
{...conform.input(url, { type: "url" })}
placeholder="URL for your Trigger API route"
/>
<Button
type="submit"
variant="primary/medium"
className="rounded-l-none"
disabled={loadingEndpointUrl}
LeadingIcon={loadingEndpointUrl ? "spinner-white" : undefined}
>
{loadingEndpointUrl ? "Saving" : "Save"}
</Button>
</div>
<FormError id={url.errorId}>{url.error}</FormError>
<FormError id={form.errorId}>{form.error}</FormError>
<Hint>
This is the URL of your Trigger API route, Typically this would be:{" "}
<InlineCode variant="extra-small">https://yourdomain.com/api/trigger</InlineCode>.
</Hint>
</InputGroup>
</setEndpointUrlFetcher.Form>
</SheetBody>
</SheetContent>
</Sheet>
);
}
@@ -7,7 +7,7 @@ import { EnvironmentLabel, environmentTitle } from "~/components/environments/En
import { HowToUseApiKeysAndEndpoints } from "~/components/helpContent/HelpContentText";
import { PageBody, PageContainer } from "~/components/layout/AppLayout";
import { BreadcrumbLink } from "~/components/navigation/NavBar";
import { ButtonContent } from "~/components/primitives/Buttons";
import { Button, ButtonContent } from "~/components/primitives/Buttons";
import { ClipboardField } from "~/components/primitives/ClipboardField";
import { DateTime } from "~/components/primitives/DateTime";
import { Header2, Header3 } from "~/components/primitives/Headers";
@@ -39,6 +39,7 @@ import { requestUrl } from "~/utils/requestUrl.server";
import { RuntimeEnvironmentType } from "../../../../../packages/database/src";
import { ConfigureEndpointSheet } from "./ConfigureEndpointSheet";
import { Badge } from "~/components/primitives/Badge";
import { FirstEndpointSheet } from "./FirstEndpointSheet";
export const loader = async ({ request, params }: LoaderArgs) => {
const userId = await requireUserId(request);
@@ -202,7 +203,12 @@ export default function Page() {
</div>
))
) : (
<Paragraph>You have no clients yet</Paragraph>
<>
<Paragraph>Add your first endpoint</Paragraph>
<Paragraph>
<FirstEndpointSheet projectId={project.id} environments={environments} />
</Paragraph>
</>
)}
</div>
{selectedEndpoint && (
@@ -8,7 +8,7 @@ import { ProjectsMenu } from "~/components/navigation/ProjectsMenu";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { ProjectPresenter } from "~/presenters/ProjectPresenter.server";
import { analytics } from "~/services/analytics.server";
import { telemetry } from "~/services/telemetry.server";
import { requireUserId } from "~/services/session.server";
import { Handle } from "~/utils/handle";
import { projectPath } from "~/utils/pathBuilder";
@@ -33,7 +33,7 @@ export const loader = async ({ request, params }: LoaderArgs) => {
});
}
analytics.project.identify({ project });
telemetry.project.identify({ project });
return typedjson({
project,
@@ -6,7 +6,7 @@ import invariant from "tiny-invariant";
import { RouteErrorDisplay } from "~/components/ErrorDisplay";
import { useOrganization } from "~/hooks/useOrganizations";
import { getOrganizationFromSlug } from "~/models/organization.server";
import { analytics } from "~/services/analytics.server";
import { telemetry } from "~/services/telemetry.server";
import { commitCurrentOrgSession, setCurrentOrg } from "~/services/currentOrganization.server";
import { requireUserId } from "~/services/session.server";
import { organizationPath } from "~/utils/pathBuilder";
@@ -25,7 +25,7 @@ export const loader = async ({ request, params }: LoaderArgs) => {
throw new Response("Not Found", { status: 404 });
}
analytics.organization.identify({ organization });
telemetry.organization.identify({ organization });
const session = await setCurrentOrg(organization.slug, request);
@@ -57,6 +57,12 @@ export async function action({ request, params }: ActionArgs) {
return json(submission);
}
return json(e, { status: 400 });
if (e instanceof Error) {
submission.error.url = `${e.name}: ${e.message}`;
} else {
submission.error.url = "Unknown error";
}
return json(submission, { status: 400 });
}
}
@@ -0,0 +1,71 @@
import { parse } from "@conform-to/zod";
import { ActionArgs, json } from "@remix-run/server-runtime";
import { z } from "zod";
import { prisma } from "~/db.server";
import {
CreateEndpointError,
CreateEndpointService,
} from "~/services/endpoints/createEndpoint.server";
import { requireUserId } from "~/services/session.server";
import { RuntimeEnvironmentTypeSchema } from "@trigger.dev/core";
import { env } from "process";
import { ValidateCreateEndpointService } from "~/services/endpoints/validateCreateEndpoint.server";
const ParamsSchema = z.object({
projectId: z.string(),
});
export const bodySchema = z.object({
environmentId: z.string(),
url: z.string().url("Must be a valid URL"),
});
export async function action({ request, params }: ActionArgs) {
const userId = await requireUserId(request);
const { projectId } = ParamsSchema.parse(params);
const formData = await request.formData();
const submission = parse(formData, { schema: bodySchema });
if (!submission.value || submission.intent !== "submit") {
return json(submission);
}
try {
const environment = await prisma.runtimeEnvironment.findUnique({
include: {
organization: true,
project: true,
},
where: {
id: submission.value.environmentId,
},
});
if (!environment) {
submission.error.environmentId = "Environment not found";
return json(submission);
}
const service = new ValidateCreateEndpointService();
const result = await service.call({
url: submission.value.url,
environment,
});
return json(submission);
} catch (e) {
if (e instanceof CreateEndpointError) {
submission.error.url = e.message;
return json(submission);
}
if (e instanceof Error) {
submission.error.url = `${e.name}: ${e.message}`;
} else {
submission.error.url = "Unknown error";
}
return json(submission, { status: 400 });
}
}
@@ -1,364 +0,0 @@
import { PostHog } from "posthog-node";
import { env } from "~/env.server";
import type { Organization } from "~/models/organization.server";
import type { Project } from "~/models/project.server";
import type { RuntimeEnvironment } from "~/models/runtimeEnvironment.server";
import type { User } from "~/models/user.server";
class BehaviouralAnalytics {
client: PostHog | undefined = undefined;
constructor(apiKey?: string) {
if (!apiKey) {
console.log("No PostHog API key, so analytics won't track");
return;
}
this.client = new PostHog(apiKey, { host: "https://app.posthog.com" });
}
user = {
identify: ({ user, isNewUser }: { user: User; isNewUser: boolean }) => {
if (this.client === undefined) return;
this.client.identify({
distinctId: user.id,
properties: {
email: user.email,
name: user.name,
authenticationMethod: user.authenticationMethod,
admin: user.admin,
createdAt: user.createdAt,
isNewUser,
},
});
if (isNewUser) {
this.#capture({
userId: user.id,
event: "user created",
eventProperties: {
email: user.email,
name: user.name,
authenticationMethod: user.authenticationMethod,
admin: user.admin,
createdAt: user.createdAt,
},
});
}
},
};
organization = {
identify: ({ organization }: { organization: Organization }) => {
if (this.client === undefined) return;
this.client.groupIdentify({
groupType: "organization",
groupKey: organization.id,
properties: {
name: organization.title,
slug: organization.slug,
createdAt: organization.createdAt,
updatedAt: organization.updatedAt,
},
});
},
new: ({
userId,
organization,
organizationCount,
}: {
userId: string;
organization: Organization;
organizationCount: number;
}) => {
if (this.client === undefined) return;
this.#capture({
userId,
event: "organization created",
organizationId: organization.id,
eventProperties: {
id: organization.id,
slug: organization.slug,
title: organization.title,
createdAt: organization.createdAt,
updatedAt: organization.updatedAt,
},
userProperties: {
organizationCount: organizationCount,
},
});
},
};
project = {
identify: ({ project }: { project: Project }) => {
if (this.client === undefined) return;
this.client.groupIdentify({
groupType: "project",
groupKey: project.id,
properties: {
name: project.name,
createdAt: project.createdAt,
updatedAt: project.updatedAt,
},
});
},
new: ({
userId,
organizationId,
project,
}: {
userId: string;
organizationId: string;
project: Project;
}) => {
if (this.client === undefined) return;
this.#capture({
userId,
event: "project created",
organizationId,
eventProperties: {
id: project.id,
title: project.name,
createdAt: project.createdAt,
updatedAt: project.updatedAt,
},
});
},
};
//todo Job
// workflow = {
// identify: ({ workflow }: { workflow: Workflow }) => {
// if (this.client === undefined) return;
// this.client.groupIdentify({
// groupType: "workflow",
// groupKey: workflow.id,
// properties: {
// name: workflow.title,
// slug: workflow.slug,
// packageJson: workflow.packageJson,
// jsonSchema: workflow.jsonSchema,
// createdAt: workflow.createdAt,
// updatedAt: workflow.updatedAt,
// organizationId: workflow.organizationId,
// type: workflow.type,
// status: workflow.status,
// externalSourceId: workflow.externalSourceId,
// service: workflow.service,
// eventNames: workflow.eventNames,
// disabledAt: workflow.disabledAt,
// archivedAt: workflow.archivedAt,
// isArchived: workflow.isArchived,
// triggerTtlInSeconds: workflow.triggerTtlInSeconds,
// },
// });
// },
// new: ({
// userId,
// organizationId,
// workflow,
// workflowCount,
// }: {
// userId: string;
// organizationId: string;
// workflow: Workflow;
// workflowCount: number;
// }) => {
// if (this.client === undefined) return;
// this.#capture({
// userId,
// event: "workflow created",
// organizationId: organizationId,
// jobId: workflow.id,
// eventProperties: {
// id: workflow.id,
// slug: workflow.slug,
// title: workflow.title,
// packageJson: workflow.packageJson,
// jsonSchema: workflow.jsonSchema,
// createdAt: workflow.createdAt,
// updatedAt: workflow.updatedAt,
// organizationId: workflow.organizationId,
// type: workflow.type,
// status: workflow.status,
// externalSourceId: workflow.externalSourceId,
// service: workflow.service,
// eventNames: workflow.eventNames,
// disabledAt: workflow.disabledAt,
// archivedAt: workflow.archivedAt,
// isArchived: workflow.isArchived,
// triggerTtlInSeconds: workflow.triggerTtlInSeconds,
// },
// userProperties: {
// workflowCount: workflowCount,
// },
// });
// },
// };
// workflowRun = {
// new: ({
// userId,
// organizationId,
// workflowId,
// workflowRun,
// environmentType,
// runCount,
// }: {
// userId: string;
// organizationId: string;
// workflowId: string;
// workflowRun: WorkflowRun;
// environmentType: string;
// runCount: number;
// }) => {
// if (this.client === undefined) return;
// this.#capture({
// userId,
// event: "workflow run created",
// eventProperties: {
// id: workflowRun.id,
// workflowId: workflowRun.workflowId,
// environmentId: workflowRun.environmentId,
// environmentType,
// eventRuleId: workflowRun.eventRuleId,
// eventId: workflowRun.eventId,
// error: workflowRun.error,
// status: workflowRun.status,
// attemptCount: workflowRun.attemptCount,
// createdAt: workflowRun.createdAt,
// updatedAt: workflowRun.updatedAt,
// startedAt: workflowRun.startedAt,
// finishedAt: workflowRun.finishedAt,
// timedOutAt: workflowRun.timedOutAt,
// timedOutReason: workflowRun.timedOutReason,
// isTest: workflowRun.isTest,
// },
// userProperties: {
// runCount: runCount,
// },
// organizationId: organizationId,
// jobId: workflowId,
// environmentId: workflowRun.environmentId,
// });
// },
// };
environment = {
identify: ({ environment }: { environment: RuntimeEnvironment }) => {
if (this.client === undefined) return;
this.client.groupIdentify({
groupType: "environment",
groupKey: environment.id,
properties: {
name: environment.slug,
slug: environment.slug,
organizationId: environment.organizationId,
createdAt: environment.createdAt,
updatedAt: environment.updatedAt,
},
});
},
};
telemetry = {
capture: ({
userId,
event,
properties,
organizationId,
environmentId,
}: {
userId: string;
event: string;
properties: Record<string | number, any>;
organizationId?: string;
environmentId?: string;
}) => {
this.#capture({
userId,
event,
eventProperties: properties,
organizationId,
environmentId,
});
},
};
#capture(event: CaptureEvent) {
if (this.client === undefined) return;
let groups: Record<string, string> = {};
if (event.organizationId) {
groups = {
...groups,
organization: event.organizationId,
};
}
if (event.projectId) {
groups = {
...groups,
project: event.projectId,
};
}
if (event.jobId) {
groups = {
...groups,
workflow: event.jobId,
};
}
if (event.environmentId) {
groups = {
...groups,
environment: event.environmentId,
};
}
let properties: Record<string, any> = {};
if (event.eventProperties) {
properties = {
...properties,
...event.eventProperties,
};
}
if (event.userProperties) {
properties = {
...properties,
$set: event.userProperties,
};
}
if (event.userOnceProperties) {
properties = {
...properties,
$set_once: event.userOnceProperties,
};
}
const eventData = {
distinctId: event.userId,
event: event.event,
properties,
groups,
};
this.client.capture(eventData);
}
}
type CaptureEvent = {
userId: string;
event: string;
organizationId?: string;
projectId?: string;
jobId?: string;
environmentId?: string;
eventProperties?: Record<string, any>;
userProperties?: Record<string, any>;
userOnceProperties?: Record<string, any>;
};
export const analytics = new BehaviouralAnalytics(env.POSTHOG_PROJECT_KEY);
+79 -14
View File
@@ -13,8 +13,10 @@ import {
RegisterTriggerBodySchema,
RunJobBody,
RunJobResponseSchema,
ValidateResponse,
ValidateResponseSchema,
} from "@trigger.dev/core";
import { safeBodyFromResponse } from "~/utils/json";
import { safeBodyFromResponse, safeParseBodyFromResponse } from "~/utils/json";
import { logger } from "./logger.server";
export class EndpointApiError extends Error {
@@ -25,20 +27,15 @@ export class EndpointApiError extends Error {
}
}
// TODO: this should work with tunnelling
export class EndpointApi {
constructor(
private apiKey: string,
private url: string,
private id: string
) {}
constructor(private apiKey: string, private url: string) {}
async ping(): Promise<PongResponse> {
async ping(endpointId: string): Promise<PongResponse> {
const response = await safeFetch(this.url, {
method: "POST",
headers: {
"x-trigger-api-key": this.apiKey,
"x-trigger-endpoint-id": this.id,
"x-trigger-endpoint-id": endpointId,
"x-trigger-action": "PING",
},
});
@@ -73,13 +70,23 @@ export class EndpointApi {
};
}
const anyBody = await response.json();
const pongResponse = await safeParseBodyFromResponse(response, PongResponseSchema);
logger.debug("ping() response from endpoint", {
body: anyBody,
});
if (!pongResponse) {
return {
ok: false,
error: `Could not parse response from endpoint. Make sure it points to the correct URL (you might be missing /api/trigger)`,
};
}
return PongResponseSchema.parse(anyBody);
if (!pongResponse.success) {
return {
ok: false,
error: `Endpoint ${this.url} responded with error: ${pongResponse.error.message}`,
};
}
return pongResponse.data;
}
async indexEndpoint() {
@@ -265,6 +272,64 @@ export class EndpointApi {
return HttpSourceResponseSchema.parse(anyBody);
}
async validate(): Promise<ValidateResponse> {
const response = await safeFetch(this.url, {
method: "POST",
headers: {
"x-trigger-api-key": this.apiKey,
"x-trigger-action": "VALIDATE",
},
});
if (!response) {
return {
ok: false,
error: `Could not connect to endpoint ${this.url}`,
};
}
if (response.status === 401) {
const body = await safeBodyFromResponse(response, ErrorWithStackSchema);
if (body) {
return {
ok: false,
error: body.message,
} as const;
}
return {
ok: false,
error: `Trigger API key is invalid`,
} as const;
}
if (!response.ok) {
return {
ok: false,
error: `Could not connect to endpoint ${this.url}. Status code: ${response.status}`,
};
}
const validateResponse = await safeParseBodyFromResponse(response, ValidateResponseSchema);
if (!validateResponse) {
return {
ok: false,
error: `Could not parse response from endpoint. Make sure it points to the correct URL (you might be missing /api/trigger)`,
};
}
if (!validateResponse.success) {
return {
ok: false,
error: `Endpoint ${this.url} responded with error: ${validateResponse.error.message}`,
};
}
return validateResponse.data;
}
}
async function safeFetch(url: string, options: RequestInit) {
@@ -34,9 +34,9 @@ export class CreateEndpointService {
}) {
const endpointUrl = this.#normalizeEndpointUrl(url);
const client = new EndpointApi(environment.apiKey, endpointUrl, id);
const client = new EndpointApi(environment.apiKey, endpointUrl);
const pong = await client.ping();
const pong = await client.ping(id);
if (!pong.ok) {
throw new CreateEndpointError("FAILED_PING", pong.error);
@@ -28,7 +28,7 @@ export class IndexEndpointService {
const endpoint = await findEndpoint(id);
// Make a request to the endpoint to fetch a list of jobs
const client = new EndpointApi(endpoint.environment.apiKey, endpoint.url, endpoint.slug);
const client = new EndpointApi(endpoint.environment.apiKey, endpoint.url);
const indexResponse = await client.indexEndpoint();
@@ -0,0 +1,44 @@
import { PrismaClient, prisma } from "~/db.server";
import { logger } from "../logger.server";
import { workerQueue } from "../worker.server";
import { RuntimeEnvironmentType } from "@trigger.dev/database";
export class RecurringEndpointIndexService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(ts: Date) {
// Find all production endpoints that haven't been indexed in the last 10 minutes
const currentTimestamp = ts.getTime();
const endpoints = await this.#prismaClient.endpoint.findMany({
where: {
environment: {
type: RuntimeEnvironmentType.PRODUCTION,
},
indexings: {
none: {
createdAt: {
gt: new Date(currentTimestamp - 10 * 60 * 1000),
},
},
},
},
});
logger.debug("Found endpoints that haven't been indexed in the last 10 minutes", {
count: endpoints.length,
});
// Enqueue each endpoint for indexing
for (const endpoint of endpoints) {
await workerQueue.enqueue("indexEndpoint", {
id: endpoint.id,
source: "INTERNAL",
});
}
}
}
@@ -0,0 +1,100 @@
import { customAlphabet } from "nanoid";
import { $transaction, prisma, PrismaClient } from "~/db.server";
import { env } from "~/env.server";
import { AuthenticatedEnvironment } from "../apiAuth.server";
import { workerQueue } from "../worker.server";
import { CreateEndpointError } from "./createEndpoint.server";
import { EndpointApi } from "../endpointApi.server";
const indexingHookIdentifier = customAlphabet("0123456789abcdefghijklmnopqrstuvxyz", 10);
export class ValidateCreateEndpointService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call({ environment, url }: { environment: AuthenticatedEnvironment; url: string }) {
const endpointUrl = this.#normalizeEndpointUrl(url);
const client = new EndpointApi(environment.apiKey, endpointUrl);
const validationResult = await client.validate();
if (!validationResult.ok) {
throw new Error(validationResult.error);
}
try {
const result = await $transaction(this.#prismaClient, async (tx) => {
const endpoint = await tx.endpoint.upsert({
where: {
environmentId_slug: {
environmentId: environment.id,
slug: validationResult.endpointId,
},
},
create: {
environment: {
connect: {
id: environment.id,
},
},
organization: {
connect: {
id: environment.organizationId,
},
},
project: {
connect: {
id: environment.projectId,
},
},
slug: validationResult.endpointId,
url: endpointUrl,
indexingHookIdentifier: indexingHookIdentifier(),
},
update: {
url: endpointUrl,
},
});
// Kick off process to fetch the jobs for this endpoint
await workerQueue.enqueue(
"indexEndpoint",
{
id: endpoint.id,
source: "INTERNAL",
},
{ tx }
);
return endpoint;
});
return result;
} catch (error) {
if (error instanceof Error) {
throw new CreateEndpointError("FAILED_UPSERT", error.message);
} else {
throw new CreateEndpointError("FAILED_UPSERT", "Something went wrong");
}
}
}
// If the endpoint URL points to localhost, and the RUNTIME_PLATFORM is docker-compose, then we need to rewrite the host to host.docker.internal
// otherwise we shouldn't change anything
#normalizeEndpointUrl(url: string) {
if (env.RUNTIME_PLATFORM === "docker-compose") {
const urlObj = new URL(url);
if (urlObj.hostname === "localhost") {
urlObj.hostname = "host.docker.internal";
return urlObj.toString();
}
}
return url;
}
}
+2 -2
View File
@@ -1,5 +1,5 @@
import type { User } from "~/models/user.server";
import { analytics } from "./analytics.server";
import { telemetry } from "./telemetry.server";
export async function postAuthentication({
user,
@@ -10,5 +10,5 @@ export async function postAuthentication({
loginMethod: User["authenticationMethod"];
isNewUser: boolean;
}) {
analytics.user.identify({ user, isNewUser });
telemetry.user.identify({ user, isNewUser });
}
@@ -54,7 +54,7 @@ export class PerformRunExecutionService {
async #executePreprocessing(execution: FoundRunExecution) {
const { run } = execution;
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url, run.endpoint.slug);
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url);
const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId });
const startedAt = new Date();
@@ -189,7 +189,7 @@ export class PerformRunExecutionService {
return;
}
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url, run.endpoint.slug);
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url);
const event = ApiEventLogSchema.parse({ ...run.event, id: run.eventId });
const startedAt = new Date();
@@ -55,8 +55,7 @@ export class DeliverHttpSourceRequestService {
const clientApi = new EndpointApi(
httpSourceRequest.environment.apiKey,
httpSourceRequest.endpoint.url,
httpSourceRequest.endpoint.slug
httpSourceRequest.endpoint.url
);
const { response, events } = await clientApi.deliverHttpSourceRequest({
@@ -0,0 +1,239 @@
import { TriggerClient } from "@trigger.dev/sdk";
import { PostHog } from "posthog-node";
import { env } from "~/env.server";
import type { Organization } from "~/models/organization.server";
import type { Project } from "~/models/project.server";
import type { User } from "~/models/user.server";
type Options = {
postHogApiKey?: string;
trigger?: {
apiKey: string;
apiUrl: string;
};
};
class Telemetry {
#posthogClient: PostHog | undefined = undefined;
#triggerClient: TriggerClient | undefined = undefined;
constructor({ postHogApiKey, trigger }: Options) {
if (postHogApiKey) {
this.#posthogClient = new PostHog(postHogApiKey, { host: "https://app.posthog.com" });
} else {
console.log("No PostHog API key, so analytics won't track");
}
if (trigger) {
this.#triggerClient = new TriggerClient({
id: "triggerdotdev",
apiKey: trigger.apiKey,
apiUrl: trigger.apiUrl,
});
console.log("Created telemetry TriggerClient");
}
}
user = {
identify: ({ user, isNewUser }: { user: User; isNewUser: boolean }) => {
if (this.#posthogClient === undefined) return;
this.#posthogClient.identify({
distinctId: user.id,
properties: {
email: user.email,
name: user.name,
authenticationMethod: user.authenticationMethod,
admin: user.admin,
createdAt: user.createdAt,
isNewUser,
},
});
if (isNewUser) {
this.#capture({
userId: user.id,
event: "user created",
eventProperties: {
email: user.email,
name: user.name,
authenticationMethod: user.authenticationMethod,
admin: user.admin,
createdAt: user.createdAt,
},
});
this.#triggerClient?.sendEvent({
name: "user.created",
payload: {
userId: user.id,
},
});
}
},
};
organization = {
identify: ({ organization }: { organization: Organization }) => {
if (this.#posthogClient === undefined) return;
this.#posthogClient.groupIdentify({
groupType: "organization",
groupKey: organization.id,
properties: {
name: organization.title,
slug: organization.slug,
createdAt: organization.createdAt,
updatedAt: organization.updatedAt,
},
});
},
new: ({
userId,
organization,
organizationCount,
}: {
userId: string;
organization: Organization;
organizationCount: number;
}) => {
if (this.#posthogClient === undefined) return;
this.#capture({
userId,
event: "organization created",
organizationId: organization.id,
eventProperties: {
id: organization.id,
slug: organization.slug,
title: organization.title,
createdAt: organization.createdAt,
updatedAt: organization.updatedAt,
},
userProperties: {
organizationCount: organizationCount,
},
});
},
};
project = {
identify: ({ project }: { project: Project }) => {
if (this.#posthogClient === undefined) return;
this.#posthogClient.groupIdentify({
groupType: "project",
groupKey: project.id,
properties: {
name: project.name,
createdAt: project.createdAt,
updatedAt: project.updatedAt,
},
});
},
new: ({
userId,
organizationId,
project,
}: {
userId: string;
organizationId: string;
project: Project;
}) => {
if (this.#posthogClient === undefined) return;
this.#capture({
userId,
event: "project created",
organizationId,
eventProperties: {
id: project.id,
title: project.name,
createdAt: project.createdAt,
updatedAt: project.updatedAt,
},
});
},
};
#capture(event: CaptureEvent) {
if (this.#posthogClient === undefined) return;
let groups: Record<string, string> = {};
if (event.organizationId) {
groups = {
...groups,
organization: event.organizationId,
};
}
if (event.projectId) {
groups = {
...groups,
project: event.projectId,
};
}
if (event.jobId) {
groups = {
...groups,
workflow: event.jobId,
};
}
if (event.environmentId) {
groups = {
...groups,
environment: event.environmentId,
};
}
let properties: Record<string, any> = {};
if (event.eventProperties) {
properties = {
...properties,
...event.eventProperties,
};
}
if (event.userProperties) {
properties = {
...properties,
$set: event.userProperties,
};
}
if (event.userOnceProperties) {
properties = {
...properties,
$set_once: event.userOnceProperties,
};
}
const eventData = {
distinctId: event.userId,
event: event.event,
properties,
groups,
};
this.#posthogClient.capture(eventData);
}
}
type CaptureEvent = {
userId: string;
event: string;
organizationId?: string;
projectId?: string;
jobId?: string;
environmentId?: string;
eventProperties?: Record<string, any>;
userProperties?: Record<string, any>;
userOnceProperties?: Record<string, any>;
};
export const telemetry = new Telemetry({
postHogApiKey: env.POSTHOG_PROJECT_KEY,
trigger:
env.TELEMETRY_TRIGGER_API_KEY && env.TELEMETRY_TRIGGER_API_URL
? {
apiKey: env.TELEMETRY_TRIGGER_API_KEY,
apiUrl: env.TELEMETRY_TRIGGER_API_URL,
}
: undefined,
});
@@ -45,7 +45,7 @@ export class InitializeTriggerService {
},
});
const clientApi = new EndpointApi(environment.apiKey, endpoint.url, endpoint.slug);
const clientApi = new EndpointApi(environment.apiKey, endpoint.url);
const registerMetadata = await clientApi.initializeTrigger(dynamicTrigger.slug, payload.params);
+27 -1
View File
@@ -6,6 +6,7 @@ import { env } from "~/env.server";
import { ZodWorker } from "~/platform/zodWorker.server";
import { sendEmail } from "./email.server";
import { IndexEndpointService } from "./endpoints/indexEndpoint.server";
import { RecurringEndpointIndexService } from "./endpoints/recurringEndpointIndex.server";
import { DeliverEventService } from "./events/deliverEvent.server";
import { InvokeDispatcherService } from "./events/invokeDispatcher.server";
import { integrationAuthRepository } from "./externalApis/integrationAuthRepository.server";
@@ -95,6 +96,31 @@ function getWorkerQueue() {
pollInterval: 1000,
},
schema: workerCatalog,
recurringTasks: {
// Run this every 5 minutes
autoIndexProductionEndpoints: {
pattern: "*/5 * * * *",
handler: async (payload, job) => {
const service = new RecurringEndpointIndexService();
await service.call(payload.ts);
},
},
// Run this every hour
purgeOldIndexings: {
pattern: "0 * * * *",
handler: async (payload, job) => {
// Delete indexings that are older than 7 days
await prisma.endpointIndex.deleteMany({
where: {
createdAt: {
lt: new Date(Date.now() - 7 * 24 * 60 * 60 * 1000),
},
},
});
},
},
},
tasks: {
"events.invokeDispatcher": {
maxAttempts: 3,
@@ -154,7 +180,7 @@ function getWorkerQueue() {
},
},
deliverHttpSourceRequest: {
maxAttempts: 5,
maxAttempts: 25,
handler: async (payload, job) => {
const service = new DeliverHttpSourceRequestService();
+17
View File
@@ -43,3 +43,20 @@ export async function safeBodyFromResponse<T>(
return parsedJson.data;
}
}
export async function safeParseBodyFromResponse<T>(
response: Response,
schema: z.Schema<T>
): Promise<z.SafeParseReturnType<unknown, T> | undefined> {
try {
const unknownJson = await response.json();
if (!unknownJson) {
return;
}
const parsedJson = schema.safeParse(unknownJson);
return parsedJson;
} catch (error) {}
}
+1
View File
@@ -61,6 +61,7 @@
"@trigger.dev/companyicons": "^1.5.14",
"@trigger.dev/database": "workspace:*",
"@trigger.dev/core": "workspace:*",
"@trigger.dev/sdk": "workspace:2.0.5",
"@uiw/react-codemirror": "^4.19.5",
"class-variance-authority": "^0.5.2",
"clsx": "^1.2.1",
+1 -1
View File
@@ -1,5 +1,5 @@
```typescript Wait example
new Job(client, {
client.defineJob({
id: "delay-job",
name: "Delay Job",
version: "0.0.1",
+1 -1
View File
@@ -1,5 +1,5 @@
```typescript
new Job(client, {
client.defineJob({
//... other options
integrations: {
slack,
+7 -21
View File
@@ -4,9 +4,8 @@ description: "Integrations make it easy to use APIs in your Jobs"
---
<Note>
You can use any API in your Jobs by using existing Node.js SDKs or HTTP
requests. Integrations just make it much easier especially when you want to
use OAuth. And you get great logging.
You can use any API in your Jobs by using existing Node.js SDKs or HTTP requests. Integrations
just make it much easier especially when you want to use OAuth. And you get great logging.
</Note>
An Integration is a package you install that makes it easy to work with a specific API. They:
@@ -35,7 +34,7 @@ const slack = new Slack({
id: "slack",
});
new Job(client, {
client.defineJob({
id: "alert-on-new-github-issues",
name: "Alert on new GitHub issues",
version: "0.1.1",
@@ -84,29 +83,16 @@ You can use OAuth to authenticate your internal team with an Integration or to a
## References
<CardGroup>
<Card
title="Integrations Dashboard"
icon="sidebar"
href="documentation/guides/integrations"
>
The Integrations Dashboard allows you to manage your Integrations and setup
OAuth.
<Card title="Integrations Dashboard" icon="sidebar" href="documentation/guides/integrations">
The Integrations Dashboard allows you to manage your Integrations and setup OAuth.
</Card>
<Card
title="Trigger.dev Connect"
icon="user-plus"
href="/documentation/concepts/connect"
>
<Card title="Trigger.dev Connect" icon="user-plus" href="/documentation/concepts/connect">
Authenticate your users with an Integration using Trigger.dev Connect.
</Card>
<Card title="View Integrations" icon="grid-2" href="/integrations">
Trigger.dev integrates with a wide range of services.
</Card>
<Card
title="Create an Integration"
icon="square-plus"
href="/integrations/create"
>
<Card title="Create an Integration" icon="square-plus" href="/integrations/create">
Create an Integration for your own use or as a public package.
</Card>
</CardGroup>
+2 -6
View File
@@ -17,7 +17,7 @@ A Job is made up of a few things:
```ts
//Job definition uses the client
new Job(client, {
client.defineJob({
// 1. Metadata
id: "event-1",
name: "Run when the foo.bar event happens",
@@ -51,11 +51,7 @@ Events [trigger](/documentation/concepts/triggers) Jobs. Jobs generate a [Run](/
<Card title="Job SDK reference" icon="wrench" href="/sdk/job">
Detailed SDK reference for Jobs.
</Card>
<Card
title="Managing Jobs Dashboard"
icon="globe"
href="/documentation/guides/managing-jobs"
>
<Card title="Managing Jobs Dashboard" icon="globe" href="/documentation/guides/managing-jobs">
Viewing and managing your Jobs in the Dashboard.
</Card>
</CardGroup>
+2 -6
View File
@@ -10,7 +10,7 @@ description: "When a [Job](/documentation/concepts/jobs) is [Triggered](/documen
A Run is a record of the execution of a Job. It is created from `run()` function of a Job.
```ts
new Job(client, {
client.defineJob({
id: "event-1",
name: "Run when the foo.bar event happens",
version: "0.0.1",
@@ -64,11 +64,7 @@ The `context` object gives you access to information about the current Run, Job,
## References
<CardGroup cols={2}>
<Card
title="Viewing Runs Dashboard"
icon="globe"
href="/documentation/guides/viewing-runs"
>
<Card title="Viewing Runs Dashboard" icon="globe" href="/documentation/guides/viewing-runs">
View all Runs for a Job, all the way down to individual Tasks.
</Card>
<Card title="`io` SDK Reference" icon="wrench" href="/sdk/io">
+7 -23
View File
@@ -10,7 +10,7 @@ description: "Tasks are individual building blocks of a Run."
In the `run()` function you can use regular code and you can use Tasks.
```ts
new Job(client, {
client.defineJob({
id: "new-user",
name: "Run when a new user signs up",
version: "0.0.1",
@@ -44,13 +44,9 @@ new Job(client, {
await io.wait("wait", 60 * 60 * 3); // wait for 3 hours
// You can wrap your own code in a Task, for retrying, resumability and logging
const response = await io.runTask(
"my-task",
{ name: "My Task" },
async () => {
return await longRunningCode(payload.userId);
}
);
const response = await io.runTask("my-task", { name: "My Task" }, async () => {
return await longRunningCode(payload.userId);
});
return response;
},
@@ -76,28 +72,16 @@ The first param of all Tasks is a `key`. This is a unique identifier for the Tas
## References
<CardGroup cols={2}>
<Card
title="Resumability"
icon="clock"
href="/documentation/concepts/resumability"
>
<Card title="Resumability" icon="clock" href="/documentation/concepts/resumability">
Runs can be very long-running. Learn how we handle this.
</Card>
<Card
title="Integrations"
icon="grid-2"
href="/documentation/concepts/integrations"
>
<Card title="Integrations" icon="grid-2" href="/documentation/concepts/integrations">
Integrations utilize Tasks.
</Card>
<Card title="`io` SDK Reference" icon="wrench" href="/sdk/io">
The `io` object allows you to easily run a Task yourself.
</Card>
<Card
title="Viewing Runs Dashboard"
icon="globe"
href="/documentation/guides/viewing-runs"
>
<Card title="Viewing Runs Dashboard" icon="globe" href="/documentation/guides/viewing-runs">
View all Runs for a Job, all the way down to individual Tasks.
</Card>
</CardGroup>
@@ -17,7 +17,7 @@ const dynamicSchedule = new DynamicSchedule(client, {
});
//2. create a Job that is attached to the dynamic schedule
new Job(client, {
client.defineJob({
id: "user-dynamicinterval",
name: "User Dynamic Interval",
version: "0.1.1",
@@ -41,7 +41,7 @@ async function registerUserCronJob(userId: string, userSchedule: string) {
}
//5. Register inside other Jobs
new Job(client, {
client.defineJob({
id: "register-dynamicinterval",
name: "Register Dynamic Interval",
version: "0.1.1",
@@ -77,7 +77,7 @@ const dynamicOnIssueOpenedTrigger = new DynamicTrigger(client, {
});
//2. create a Job that is attached to the dynamic trigger
new Job(client, {
client.defineJob({
id: "listen-for-dynamic-trigger",
name: "Listen for dynamic trigger",
version: "0.1.1",
@@ -87,9 +87,7 @@ new Job(client, {
},
run: async (payload, io, ctx) => {
await io.slack.postMessage("Slack 📝", {
text: `New Issue opened on repo: ${
payload.issue.html_url
}. \n\n${JSON.stringify(ctx)}`,
text: `New Issue opened on repo: ${payload.issue.html_url}. \n\n${JSON.stringify(ctx)}`,
channel: "C04GWUTDC3W",
});
},
@@ -105,7 +103,7 @@ async function registerRepo(owner: string, repo: string) {
}
//4. Register inside other Jobs
new Job(client, {
client.defineJob({
id: "new-repo",
name: "New repo",
version: "0.1.1",
@@ -23,7 +23,7 @@ You can always start out by using `z.any()` as your schema, and then later on yo
## Example
```ts
new Job(client, {
client.defineJob({
id: "new-user-slack",
name: "New user slack message",
version: "0.1.0",
@@ -58,9 +58,8 @@ new Job(client, {
```
<Note>
You can subscribe to the same event from multiple different Jobs. This is
useful if you want to send an event to multiple different services or if you
want to keep each Job small and simple.
You can subscribe to the same event from multiple different Jobs. This is useful if you want to
send an event to multiple different services or if you want to keep each Job small and simple.
</Note>
## Sending events
@@ -84,7 +83,7 @@ await client.sendEvent({
You can use `io.sendEvent()` to send events from inside a Job run, to trigger another. [View the SDK reference](/sdk/io/sendevent).
```ts
new Job(client, {
client.defineJob({
id: "event-1",
name: "Run when the foo.bar event happens",
version: "0.0.1",
@@ -15,7 +15,7 @@ This job will run every 60 seconds, starting 60 seconds after this Job is first
```ts
import { Job, intervalTrigger } from "@trigger.dev/sdk";
new Job(client, {
client.defineJob({
id: "scheduled-job-1",
name: "Scheduled Job 1",
version: "0.1.1",
@@ -43,7 +43,7 @@ This job will run at 2:30pm every Monday. You can get help with [CRON syntax](ht
```ts
import { Job, cronTrigger } from "@trigger.dev/sdk";
new Job(client, {
client.defineJob({
id: "scheduled-job-2",
name: "Scheduled Job 2",
version: "0.1.1",
@@ -32,7 +32,7 @@ const github = new Github({
token: process.env.GITHUB_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "critical-issue-alert",
name: "Critical Issue Alert",
version: "0.1.0",
@@ -42,7 +42,7 @@ const slack = new Slack({
id: "slack",
});
new Job(client, {
client.defineJob({
id: "critical-issue-alert",
name: "Critical Issue Alert",
version: "0.1.0",
@@ -46,7 +46,7 @@ There are two way to use Integrations in a Job:
This example automatically assigns "matt-aitken" to any new issue in the `trigger.dev` repo (lucky him).
```ts
new Job(client, {
client.defineJob({
id: "assign-on-issue-opened",
name: "Assign on Issue Opened",
version: "0.1.0",
@@ -35,7 +35,7 @@ There are two way to use Integrations in a Job:
This example send a Slack message when someone stars the `trigger.dev` GitHub repo 🤩.
```ts
new Job(client, {
client.defineJob({
id: "star-slack-notification",
name: "New Star Slack Notification",
version: "0.1.0",
+3 -3
View File
@@ -13,7 +13,7 @@ We use it [extensively](https://github.com/search?q=repo%3Atriggerdotdev%2Ftrigg
But there are a few places where we ask you to provide us with a Zod schema, for example when defining your own [events](/documentation/concepts/triggers/events):
```ts
new Job(client, {
client.defineJob({
id: "new-user",
name: "New user",
version: "0.1.0",
@@ -36,8 +36,8 @@ new Job(client, {
So it will help to know a little about Zod and how to use it. We definitely recommend the well written [Zod README](https://github.com/colinhacks/zod#readme) but we've included a short primer below.
<Tip>
Wherever we require you to pass in a Zod schema, you can always start with
`z.any()` which accepts `any` type and then add more strict validations later.
Wherever we require you to pass in a Zod schema, you can always start with `z.any()` which accepts
`any` type and then add more strict validations later.
</Tip>
## Basic Usage
+4 -9
View File
@@ -62,7 +62,7 @@ yarn dlx @trigger.dev/cli@latest init
It will ask you a few questions
1. Are you using the [Trigger.dev Cloud](https://trigger.dev) or [self-hosting](/documentation/guides/self-hosting)? You're probably using the cloud.
1. Are you using the [Trigger.dev Cloud](https://cloud.trigger.dev) or [self-hosting](/documentation/guides/self-hosting)?
2. Enter your development API key. Enter the key you copied earlier.
3. Enter a unique ID for your endpoint (you can just use the default by hitting enter)
@@ -167,7 +167,7 @@ In there is this Job:
```typescript
//Job definition uses the client
new Job(client, {
client.defineJob({
// 1. Metadata
id: "example-job",
name: "Example Job",
@@ -212,11 +212,7 @@ Congratulations, you should get redirected so you can see your first Run!
## What's next?
<CardGroup cols={2}>
<Card
title="Write your first Job"
icon="hexagon-plus"
href="/documentation/guides/create-a-job"
>
<Card title="Write your first Job" icon="hexagon-plus" href="/documentation/guides/create-a-job">
A Guide for how to create your first real Job
</Card>
<Card
@@ -227,8 +223,7 @@ Congratulations, you should get redirected so you can see your first Run!
Learn more about how Trigger.dev works and how it can help you.
</Card>
<Card title="Examples" icon="slot-machine" href="/examples">
One of the quickest ways to learn how Trigger.dev works is to view some
example Jobs.
One of the quickest ways to learn how Trigger.dev works is to view some example Jobs.
</Card>
<Card title="Get help" icon="hire-a-helper" href="/documentation/get-help">
Struggling getting setup or have a question? We're here to help.
+7 -4
View File
@@ -4,21 +4,24 @@ description: "An ever-growing list of example Jobs which you can use to get star
---
<Info>
If you are using integrations, you'll need set up authentication either using
OAuth or API keys / access tokens. You can find out how to do that in the
[integrations section](/integrations).
If you are using integrations, you'll need set up authentication either using OAuth or API keys /
access tokens. You can find out how to do that in the [integrations section](/integrations).
</Info>
Click the links below to view the Job code. You can also easily test these Jobs by following the instructions in the README of each example.
| Job (code in link) | Description | Integrations used |
| ----------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------------------------- |
| [Basic delay](https://github.com/triggerdotdev/examples/blob/main/delays/src/jobs/delayJob.ts) | Logs a message to the console, waits for 5 minutes, and then logs another message. | N/A |
| [Basic interval](https://github.com/triggerdotdev/examples/blob/main/scheduled/src/jobs/interval.ts) | This Job will run every 60 seconds, starting 60 seconds after this Job is first indexed. | N/A |
| [Cron scheduled interval](https://github.com/triggerdotdev/examples/blob/main/scheduled/src/jobs/cronScheduled.ts) | A scheduled Job which runs at 2:30pm every Monday. | N/A |
| [OpenAI text summarizer](https://github.com/triggerdotdev/examples/blob/main/openai-text-summarizer/src/jobs/textSummarizer.ts) | Summarizes a block of text, pulling out the most unique and helpful points using OpenAI GPT-3.5 turbo. | [OpenAI](/integrations/apis/openai) |
| [Tell me a joke using OpenAI](https://github.com/triggerdotdev/examples/blob/main/openai/src/jobs/tellMeAJoke.ts) | Generates a random joke using OpenAI GPT 3.5. | [OpenAI](/integrations/apis/openai) |
| [Generate an image using OpenAI](https://github.com/triggerdotdev/examples/blob/main/openai/src/jobs/generateHedgehogImages.ts) | Generates a random image of a hedgehog using OpenAI DALL-E. | [OpenAI](/integrations/apis/openai) |
| [GitHub issue reminder](https://github.com/triggerdotdev/examples/blob/main/github-issue-reminder/jobs/githubIssue.ts) | Sends a Slack message to a channel if a GitHub issue is left open for 24 hours | [GitHub](/integrations/apis/github), [Slack](/integrations/apis/slack) |
| [Github new star alert in Slack](https://github.com/triggerdotdev/examples/blob/main/github/src/jobs/newStarToSlack.ts) | When a repo is starred, a message is sent to a Slack channel with the name and URL of the GitHub user who starred the repo, and the updated Stargazers count. | [GitHub](/integrations/apis/github), [Slack](/integrations/apis/slack) |
| [Add a custom label to a GitHub issue when it is created](https://github.com/triggerdotdev/examples/blob/main/github/src/jobs/onIssueOpened.ts) | When a new GitHub issue is opened it adds a "Bug" label to it. | [GitHub](/integrations/apis/github) |
| [GitHub new star alert](https://github.com/triggerdotdev/examples/blob/main/github/src/jobs/newStarAlert.ts) | When a repo is starred a message is logged with the new Stargazers count. | [GitHub](/integrations/apis/github) |
| [Github new star alert in Slack](https://github.com/triggerdotdev/examples/blob/main/github/src/jobs/newStarToSlack.ts) | When a repo is starred, a message is sent to a Slack channel with the name and URL of the GitHub user who starred the repo, and the updated Stargazers count. | [GitHub](/integrations/apis/github), [Slack](/integrations/apis/slack) |
| [Send a Slack message when an event is received](https://github.com/triggerdotdev/examples/blob/main/slack/src/jobs/sendSlackMessage.ts) | Sends a Slack message to a specific channel when an event is received. | [Slack](/integrations/apis/slack) |
| [Send an email using Resend](https://github.com/triggerdotdev/examples/blob/main/resend/src/jobs/resendBasicEmail.ts) | Send a basic email using Resend | [Resend](/integrations/apis/resend) |
Binary file not shown.

After

Width:  |  Height:  |  Size: 60 KiB

+1 -1
View File
@@ -23,7 +23,7 @@ title: Tasks
## Usage
```ts
new Job(client, {
client.defineJob({
id: "github-integration-on-issue-opened",
name: "GitHub Integration - On Issue Opened",
version: "0.1.0",
+1 -1
View File
@@ -29,7 +29,7 @@ const github = new Github({
token: process.env.GITHUB_TOKEN!,
});
new Job(client, {
client.defineJob({
id: "github-integration-on-issue",
name: "GitHub Integration - On Issue",
version: "0.1.0",
+11 -16
View File
@@ -45,8 +45,7 @@ const github2 = new Github({
<CardGroup cols={2}>
<Card title="Triggers" icon="stars" href="/integrations/apis/github-triggers">
Trigger Jobs when events happen in GitHub, such as a new commit or a new
issue.
Trigger Jobs when events happen in GitHub, such as a new commit or a new issue.
</Card>
<Card title="Tasks" icon="sparkles" href="/integrations/apis/github-tasks">
Perform tasks such as creating a new issue or a new comment.
@@ -58,8 +57,8 @@ const github2 = new Github({
You can use the underlying client to do anything Octokit supports. In this example we create a project card when a new issue is opened..
<Info>
View [the official GitHub docs](https://docs.github.com/en/rest) for
everything that is supported{" "}
View [the official GitHub docs](https://docs.github.com/en/rest) for everything that is
supported{" "}
</Info>
```ts
@@ -70,7 +69,7 @@ const github = new Github({
token: process.env.GITHUB_TOKEN!,
});
new Job(client, {
client.defineJob({
id: "alert-on-new-github-issues",
name: "Alert on new GitHub issues",
version: "0.1.1",
@@ -84,17 +83,13 @@ new Job(client, {
},
run: async (payload, io, ctx) => {
//wrap the SDK call in runTask
const { data } = await io.runTask(
"create-card",
{ name: "Create card" },
async () => {
//create a project card using the underlying client
return io.github.client.rest.projects.createCard({
column_id: 123,
note: "test",
});
}
);
const { data } = await io.runTask("create-card", { name: "Create card" }, async () => {
//create a project card using the underlying client
return io.github.client.rest.projects.createCard({
column_id: 123,
note: "test",
});
});
//log the url of the created card
await io.logger.info(data.url);
+239 -20
View File
@@ -2,10 +2,17 @@
title: Introduction
---
Trigger.dev provides seamless integration with OpenAI, enabling developers to harness the power of AI
language models in their serverless applications. With Trigger.dev's background tasks, long-running
OpenAI completions become possible, even within the constraints of serverless timeouts.
<Snippet file="integration-getting-started.mdx" />
## Installation
To get started with the OpenAI integration on Trigger.dev, you need to install the `@trigger.dev/openai` package.
You can do this using npm, pnpm, or yarn:
<CodeGroup>
```bash npm
@@ -13,7 +20,7 @@ npm install @trigger.dev/openai@latest
```
```bash pnpm
pnpm install @trigger.dev/openai@latest
pnpm add @trigger.dev/openai@latest
```
```bash yarn
@@ -24,7 +31,8 @@ yarn add @trigger.dev/openai@latest
## Authentication
OpenAI supports API Keys
To use the OpenAI API with Trigger.dev, you'll need an API Key from OpenAI.
If you don't have one yet, you can obtain it from the [OpenAI dashboard](https://platform.openai.com/account/api-keys).
```ts
import { OpenAI } from "@trigger.dev/openai";
@@ -35,10 +43,12 @@ const openai = new OpenAI({
});
```
## Example
## Usage
Include the OpenAI integration in your Trigger.dev job:
```ts
new Job(client, {
client.defineJob({
id: "openai-tasks",
name: "OpenAI Tasks",
version: "0.0.1",
@@ -50,20 +60,11 @@ new Job(client, {
openai,
},
run: async (payload, io, ctx) => {
const response = await io.openai.backgroundCreateChatCompletion(
"background-chat-completion",
{
model: "gpt-3.5-turbo",
messages: [
{
role: "user",
content: "Create a good programming joke about background jobs",
},
],
}
);
await io.logger.info("choices", response.choices);
// Now you can access the OpenAI tasks through the io object
await io.openai.createCompletion("completion", {
model: "davinci",
prompt: "Once upon a time",
});
},
});
```
@@ -75,9 +76,9 @@ Tasks that are marked as "long-running" can last longer than your serverless tim
| Function Name | Description | Long-running? |
| -------------------------------- | ------------------------------------------------------------------------- | ------------- |
| `createCompletion` | Generates text completions given a prompt. |
| `backgroundCreateCompletion` | Generates text completions in the background. | ✔ |
| `backgroundCreateCompletion` | Generates text completions in the background. | ✔ |
| `createChatCompletion` | Generates text completions in a conversational context. |
| `backgroundCreateChatCompletion` | Generates text completions in a conversational context in the background. | ✔ |
| `backgroundCreateChatCompletion` | Generates text completions in a conversational context in the background. | ✔ |
| `retrieveModel` | Retrieves a specific model by ID. |
| `listModels` | Lists the available models. |
| `createEdit` | Edits a given text prompt. |
@@ -92,3 +93,221 @@ Tasks that are marked as "long-running" can last longer than your serverless tim
| `cancelFineTune` | Cancels a specific fine-tune by ID. |
| `listFineTuneEvents` | Lists the events for a specific fine-tune by ID. |
| `deleteFineTune` | Deletes a specific fine-tune by ID. |
## Examples
### Generate a joke
Here's an example of how to use the OpenAI integration in a Trigger.dev job.
In this example, we'll create a background task to generate a programming joke.
```ts
client.defineJob({
id: "openai-tasks",
name: "OpenAI Tasks",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.tasks",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const response = await io.openai.backgroundCreateChatCompletion("background-chat-completion", {
model: "gpt-3.5-turbo",
messages: [
{
role: "user",
content: "Create a good programming joke about background jobs",
},
],
});
await io.logger.info("choices", response.choices);
},
});
```
### Generate Code Snippets
In this example, we'll leverage Trigger.dev's background task to generate code snippets for
a given programming task:
```ts
client.defineJob({
id: "openai-tasks",
name: "OpenAI Tasks",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.tasks",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const programmingTask = `Create a function that checks if a string is a palindrome.`;
const response = await io.openai.backgroundCreateCompletion("background-completion", {
model: "gpt-3.5-turbo",
prompt: `Coding task: ${programmingTask}\n\n`,
});
await io.logger.info("codeSnippet", response.choices[0]?.text);
},
});
```
### Summarize Text
We'll use Trigger.dev's background task to summarize a lengthy article:
```ts
client.defineJob({
id: "openai-tasks",
name: "OpenAI Tasks",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.tasks",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const articleToSummarize = `Lorem ipsum. olor sit amet, consectetur adipiscing elit.
Sed nec aliquet sapien. Pellentesque vitae nisi id purus luctus tincidunt.
Proin condimentum malesuada turpis, eget tincidunt mauris viverra in.`;
const response = await io.openai.backgroundCreateCompletion("background-completion", {
model: "gpt-3.5-turbo",
prompt: `Please summarize the following article:\n\n${articleToSummarize}`,
});
await io.logger.info("summary", response.choices[0]?.text);
},
});
```
### Draft Email Response
we'll use Trigger.dev's background task to draft an email response based on a given email content:
```ts
client.defineJob({
id: "openai-tasks",
name: "OpenAI Tasks",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.tasks",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const emailContent = `Dear John,
Thank you for your inquiry. We appreciate your interest in our products.
I have reviewed your request, and I'm pleased to inform you that we can
accommodate your requirements. Please find the attached proposal for your
reference. If you have any further questions, feel free to ask.
Best regards,
Jane Doe`;
const response = await io.openai.backgroundCreateChatCompletion("background-chat-completion", {
model: "gpt-3.5-turbo",
messages: [
{
role: "user",
content: emailContent,
},
{
role: "assistant",
content: "Draft a suitable response to the email above.",
},
],
});
await io.logger.info("draftedEmailResponse", response.choices[0]?.text);
},
});
```
### Chatbot Counseling Session
This job represents a simulated AI counseling session. Leveraging OpenAI's ability to understand context and generate human-like text, it forms empathetic responses to user inputs. Such a system could be part of a mental wellness app.
```ts
client.defineJob({
id: "openai-chatbot-counseling",
name: "Chatbot Counseling Session",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.startCounselingSession",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const response = await io.openai.backgroundCreateChatCompletion(
"background-counseling-chat-completion",
{
model: "gpt-3.5-turbo",
messages: [
{
role: "system",
content: "You are a helpful and empathetic AI counselor.",
},
{
role: "user",
content: "I've been feeling really stressed out lately.",
},
],
}
);
await io.logger.info("counseling session", response.choices);
},
});
```
### AI Roleplay Game Session
This job creates a fantasy AI role-playing game. It could be fun for interactive storytelling or game development contexts.
```ts
client.defineJob({
id: "openai-roleplay-game-session",
name: "AI Roleplay Game Session",
version: "0.0.1",
trigger: eventTrigger({
name: "openai.startRoleplayGameSession",
schema: z.object({}),
}),
integrations: {
openai,
},
run: async (payload, io, ctx) => {
const response = await io.openai.backgroundCreateChatCompletion(
"background-roleplay-game-session-chat-completion",
{
model: "gpt-3.5-turbo",
messages: [
{
role: "system",
content: "You are an intelligent guide in a fantasy role-playing game.",
},
{
role: "user",
content: "I embark on a quest for the enchanted crown. What's the first step?",
},
],
}
);
await io.logger.info("roleplay game session", response.choices);
},
});
```
+26 -29
View File
@@ -50,7 +50,7 @@ export const plain = new Plain({
apiKey: process.env.PLAIN_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "plain-playground",
name: "Plain Playground",
version: "0.1.1",
@@ -87,37 +87,34 @@ new Job(client, {
customerId: customer.id,
});
const timelineEntry = await io.plain.upsertCustomTimelineEntry(
"upsert-timeline-entry",
{
customerId: customer.id,
title: "My timeline entry",
components: [
{
componentText: {
text: `This is a nice title`,
},
const timelineEntry = await io.plain.upsertCustomTimelineEntry("upsert-timeline-entry", {
customerId: customer.id,
title: "My timeline entry",
components: [
{
componentText: {
text: `This is a nice title`,
},
{
componentDivider: {
dividerSpacingSize: ComponentDividerSpacingSize.M,
},
},
{
componentDivider: {
dividerSpacingSize: ComponentDividerSpacingSize.M,
},
{
componentText: {
textSize: ComponentTextSize.S,
textColor: ComponentTextColor.Muted,
text: "External id",
},
},
{
componentText: {
textSize: ComponentTextSize.S,
textColor: ComponentTextColor.Muted,
text: "External id",
},
{
componentText: {
text: foundCustomer?.externalId ?? "",
},
},
{
componentText: {
text: foundCustomer?.externalId ?? "",
},
],
}
);
},
],
});
},
});
```
@@ -145,7 +142,7 @@ export const plain = new Plain({
apiKey: process.env.PLAIN_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "plain-client",
name: "Plain Client",
version: "0.1.0",
+1 -1
View File
@@ -51,7 +51,7 @@ const resend = new Resend({
apiKey: process.env.RESEND_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "send-resend-email",
name: "Send Resend Email",
version: "0.1.0",
+1 -1
View File
@@ -37,7 +37,7 @@ const slack = new Slack({
## Example
```ts
new Job(client, {
client.defineJob({
id: "slack-test",
name: "Slack test",
version: "0.0.1",
+15 -23
View File
@@ -7,8 +7,8 @@ description: "Interact with your Supabase project using the Supabase JS Client."
Our `@trigger.dev/supabase` package provides an integration that wraps the [@supabase/supabase-js](https://github.com/supabase/supabase-js) package, allowing you to run tasks to interact with your Supabase project.
<Note>
If you want to trigger jobs based on changes in your Supabase database, you'll
need to use the [Supabase Management API](../management) integration
If you want to trigger jobs based on changes in your Supabase database, you'll need to use the
[Supabase Management API](/integrations/apis/supabase/management) integration
</Note>
## Usage
@@ -26,8 +26,7 @@ const supabase = new Supabase({
```
<Warning>
Never expose the `service_role` key in a browser or anywhere where a user can
see it.
Never expose the `service_role` key in a browser or anywhere where a user can see it.
</Warning>
You can then use the `supabase` integration to run tasks in your jobs:
@@ -39,21 +38,17 @@ client.defineJob({
supabase,
},
run: async (payload, io, ctx) => {
const { data: users, error } = await io.supabase.runTask(
"find-users",
async (db) => {
return db.from("users").select("*");
}
);
const { data: todos, error } = await io.supabase.runTask("find-todos", async (db) => {
return db.from("todos").select("*");
});
},
});
```
<Note>
By using `runTask` instead of the `@supabase/supabase-js` client directly
inside your job run, you'll be able to create tasks that can be run
idempotently and also retried. For more, see our guide on
[Resumability](http://localhost:3050/documentation/concepts/resumability)
By using `runTask` instead of the `@supabase/supabase-js` client directly inside your job run,
you'll be able to create tasks that can be run idempotently and also retried. For more, see our
guide on [Resumability](http://localhost:3050/documentation/concepts/resumability)
</Note>
You can also choose to throw an error if the query fails and abort the job run:
@@ -65,8 +60,8 @@ client.defineJob({
supabase,
},
run: async (payload, io, ctx) => {
const users = await io.supabase.runTask("find-users", async (db) => {
const { data, error } = await db.from("users").select("*");
const todos = await io.supabase.runTask("find-todos", async (db) => {
const { data, error } = await db.from("todos").select("*");
if (error) throw error;
@@ -83,10 +78,7 @@ The `db` object passed to the callback is an instance of the [@supabase/supabase
- [Invoking Functions](https://supabase.com/docs/reference/javascript/functions-invoke)
- [Storage](https://supabase.com/docs/reference/javascript/storage-createbucket)
<Warning>
Currently we do not support Supabase Realtime (such as subscribing to a
channel)
</Warning>
<Warning>Currently we do not support Supabase Realtime (such as subscribing to a channel)</Warning>
## Typescript Support
@@ -108,15 +100,15 @@ client.defineJob({
supabase,
},
run: async (payload, io, ctx) => {
const users = await io.supabase.runTask("find-users", async (db) => {
const { data, error } = await db.from("users").select("*");
const todos = await io.supabase.runTask("find-todos", async (db) => {
const { data, error } = await db.from("todos").select("*");
if (error) throw error;
return data;
});
// users is now typed as User[] instead of any[]
// todos is now typed as Todo[] instead of any[]
},
});
```
+52 -20
View File
@@ -85,6 +85,28 @@ For a full list of available tasks, see the [Supabase Management API](https://su
The `SupabaseManagement` integration also provides the ability to trigger jobs based on changes in your Supabase database through the use of [Supabase Database Webhooks](https://supabase.com/docs/guides/database/webhooks).
### Enable Database Webhooks
<Info>
Manually enabling database webhooks are only needed if you are using `@trigger.dev/supabase` at
version `2.0.2` or earlier. If you are using `2.0.3` or later, this is done automatically for you.
</Info>
Currently the Supabase Management API does not provide a way to enable database webhooks, so you'll need to do this manually.
You can do this by visiting your [Database Webhooks settings](https://supabase.com/dashboard/project/_/database/hooks) and clicking the "Enable webhooks" button:
![Enable Database Webhooks](/images/supabase-enable-webhooks.png)
You'll have to do this for each Supabase project you want to use webhooks with.
<Note>
You don't actually need to create any webhooks yourself, our integration will take care of that
part for you.
</Note>
### Usage
To use this feature, you'll first initialize a `db` instance, passing in your Supabase project [ID](https://supabase.com/dashboard/project/_/settings/api) (or URL):
```ts
@@ -104,7 +126,7 @@ client.defineJob({
id: "supabase-trigger",
name: "Supabase Trigger",
trigger: db.onInserted({
table: "users",
table: "todos",
}),
run: async (payload, io, ctx) => {
// payload is the database webhook body (see https://supabase.com/docs/guides/database/webhooks#payload)
@@ -115,20 +137,6 @@ client.defineJob({
You can add additional filters to the trigger by passing a `filter` object:
```ts
client.defineJob({
id: "supabase-trigger",
name: "Supabase Trigger",
trigger: db.onUpdated({
table: "users",
filter: {
country: ["USA", "Canada"], // This will only trigger the job if the user.country is USA or Canada
},
}),
run: async (payload, io, ctx) => {
// payload is the database webhook body (see https://supabase.com/docs/guides/database/webhooks#payload)
},
});
client.defineJob({
id: "supabase-trigger",
name: "Supabase Trigger",
@@ -150,10 +158,34 @@ client.defineJob({
});
```
You can also listen for multiple different events using the `on` trigger:
```ts
client.defineJob({
id: "supabase-trigger",
name: "Supabase Trigger",
trigger: db.on({
table: "todos",
events: ["INSERT", "UPDATE"] // Trigger on both insert and update events
filter: {
record: {
is_completed: [false],
},
},
}),
run: async (payload, io, ctx) => {
if (payload.type === "INSERT") {
// payload will be typed as the INSERT payload
} else {
// payload will be typed as the UPDATE payload
}
},
});
```
<Note>
We will only create at most 1 database webhook per table, to limit resource
usage when writing to your database. This means we cannot support scoping
updated triggers to specific columns.
We will only create at most 1 database webhook per table, to limit resource usage when writing to
your database. This means we cannot support scoping updated triggers to specific columns.
</Note>
### Typescript Support
@@ -175,10 +207,10 @@ client.defineJob({
id: "supabase-trigger",
name: "Supabase Trigger",
trigger: db.onUpdated({
table: "users",
table: "todos",
}),
run: async (payload, io, ctx) => {
// payload.record and payload.old_record are now correctly typed to match the users table
// payload.record and payload.old_record are now correctly typed to match the todos table
},
});
```
+6 -9
View File
@@ -50,7 +50,7 @@ export const typeform = new Typeform({
token: process.env.TYPEFORM_API_KEY!,
});
new Job(client, {
client.defineJob({
id: "do-something-on-new-responses",
name: "Send a message to slack on new responses",
version: "0.1.1",
@@ -90,7 +90,7 @@ const typeform = new Typeform({
token: process.env.TYPEFORM_PAT!,
});
new Job(client, {
client.defineJob({
id: "typeform-tasks",
name: "Typeform Tasks",
version: "0.1.0",
@@ -110,12 +110,9 @@ new Job(client, {
pageSize: 50,
});
const allResponses = await io.typeform.getAllResponses(
"get-all-responses",
{
uid: payload.formId,
}
);
const allResponses = await io.typeform.getAllResponses("get-all-responses", {
uid: payload.formId,
});
},
});
```
@@ -133,7 +130,7 @@ const typeform = new Typeform({
token: process.env.TYPEFORM_PAT!,
});
new Job(client, {
client.defineJob({
id: "typeform-client",
name: "Typeform Client",
version: "0.1.0",
+22 -31
View File
@@ -71,9 +71,9 @@ Once you've created your Integration package, you can start developing it. In th
This is the entry point of the Integration package. It exports a main "integration" class that implements the `TriggerIntegration` interface. For example, the `@trigger.dev/github` Integration exports a `Github` class that implements.
<Tip>
We're adopting the naming convention of naming the class after the service,
without a suffix or prefix. We prefer the exported name be `Slack` instead of
something like `SlackIntegration` or `SlackConnector`
We're adopting the naming convention of naming the class after the service, without a suffix or
prefix. We prefer the exported name be `Slack` instead of something like `SlackIntegration` or
`SlackConnector`
</Tip>
<Accordion title="Example: OpenAI">
@@ -84,9 +84,7 @@ import { Configuration, OpenAIApi } from "openai";
import * as tasks from "./tasks";
import { OpenAIIntegrationOptions } from "./types";
export class OpenAI
implements TriggerIntegration<IntegrationClient<OpenAIApi, typeof tasks>>
{
export class OpenAI implements TriggerIntegration<IntegrationClient<OpenAIApi, typeof tasks>> {
client: IntegrationClient<OpenAIApi, typeof tasks>;
constructor(private options: OpenAIIntegrationOptions) {
@@ -121,19 +119,18 @@ export class OpenAI
The `TriggerIntegration` interface requires three properties to be implemented:
<ParamField body="id" type="string" required>
The `id` that uniquely identifies the Integration. This should always be
passed through the constructor options.
The `id` that uniquely identifies the Integration. This should always be passed through the
constructor options.
</ParamField>
<ParamField body="metadata" type="object" required>
<Expandable title="properties">
<ParamField body="id" type="string" required>
A unique identifier for the Integration. For example, the OpenAI
Integration has an id of `"openai"`.
A unique identifier for the Integration. For example, the OpenAI Integration has an id of
`"openai"`.
</ParamField>
<ParamField body="name" type="string" required>
The name of the Integration. For example, the OpenAI Integration has a
name of `"OpenAI"`.
The name of the Integration. For example, the OpenAI Integration has a name of `"OpenAI"`.
</ParamField>
</Expandable>
</ParamField>
@@ -194,11 +191,7 @@ For example, here is the `getForm` authenticated task defined in the `@trigger.d
import type { AuthenticatedTask } from "@trigger.dev/sdk";
import type { GetFormParams, GetFormResponse, TypeformSDK } from "./types";
export const getForm: AuthenticatedTask<
TypeformSDK,
GetFormParams,
GetFormResponse
> = {
export const getForm: AuthenticatedTask<TypeformSDK, GetFormParams, GetFormResponse> = {
init: (params) => {
return {
name: "Get Form",
@@ -232,7 +225,7 @@ export type GetFormResponse = Prettify<Typeform.Form>;
```
```ts usage.ts
new Job(client, {
client.defineJob({
id: "typeform-playground",
name: "Typeform Playground",
version: "0.1.1",
@@ -258,9 +251,9 @@ The first thing to notice is the explicit typing of the `getForm` export as an `
If you take a look at the `usage.ts` file above, you can see how this task is used in a job. The `io.typeform.getForm` function is typed as returning `Promise<GetFormResponse>` and the `params` argument is typed as `GetFormParams`.
<Note>
Notice how the params are the _second_ argument to `getForm`, that's because
the first argument is always the task key. See our [Keys and Resumability
docs](/documentation/concepts/resumability) for more on why this is important
Notice how the params are the _second_ argument to `getForm`, that's because the first argument is
always the task key. See our [Keys and Resumability docs](/documentation/concepts/resumability)
for more on why this is important
</Note>
#### `run` function
@@ -268,8 +261,8 @@ If you take a look at the `usage.ts` file above, you can see how this task is us
The `run` function is the main function that will be called when the task is run. It's an async function that takes up to 5 arguments:
<ParamField body="params" type="type parameter" required>
The input params that were passed to the task. This is the second argument to
the `getForm` function in the example above.
The input params that were passed to the task. This is the second argument to the `getForm`
function in the example above.
</ParamField>
<ParamField body="client" type="type parameter" required>
@@ -286,9 +279,9 @@ The `run` function is the main function that will be called when the task is run
</ParamField>
<ParamField body="auth" type="ConnectionAuth">
If for some reason you need to access the auth object that was used to seed
the SDK client, you can access it here. The `AuthenticatedTask` generic type
takes an optional 4th type parameter that allows you to specify the auth type
If for some reason you need to access the auth object that was used to seed the SDK client, you
can access it here. The `AuthenticatedTask` generic type takes an optional 4th type parameter that
allows you to specify the auth type
</ParamField>
#### `init` function
@@ -296,8 +289,8 @@ The `run` function is the main function that will be called when the task is run
The `init` function is used to initialize the task. It's a synchronous function that takes a single argument:
<ParamField body="params" type="type parameter" required>
The input params that were passed to the task. This is the second argument to
the `getForm` function in the example above.
The input params that were passed to the task. This is the second argument to the `getForm`
function in the example above.
</ParamField>
#### `onError` function
@@ -466,9 +459,7 @@ export const backgroundCreateCompletion: AuthenticatedTask<
headers: {
"Content-Type": "application/json",
Authorization: redactString`Bearer ${auth.apiKey}`,
...(auth.organization
? { "OpenAI-Organization": auth.organization }
: {}),
...(auth.organization ? { "OpenAI-Organization": auth.organization } : {}),
},
body: JSON.stringify(params),
}
+10 -27
View File
@@ -23,7 +23,8 @@
},
"feedback": {
"suggestEdit": true,
"raiseIssue": true
"raiseIssue": true,
"thumbsRating": true
},
"topbarCtaButton": {
"type": "github",
@@ -154,10 +155,7 @@
},
{
"group": "Overview",
"pages": [
"integrations/introduction",
"integrations/create"
]
"pages": ["integrations/introduction", "integrations/create"]
},
{
"group": "Integrations",
@@ -180,22 +178,16 @@
},
{
"group": "OpenAI",
"pages": [
"integrations/apis/openai"
]
"pages": ["integrations/apis/openai"]
},
"integrations/apis/plain",
{
"group": "Resend",
"pages": [
"integrations/apis/resend"
]
"pages": ["integrations/apis/resend"]
},
{
"group": "Slack",
"pages": [
"integrations/apis/slack"
]
"pages": ["integrations/apis/slack"]
},
"integrations/apis/typeform"
]
@@ -249,10 +241,7 @@
"sdk/dynamictrigger/constructor",
{
"group": "Instance methods",
"pages": [
"sdk/dynamictrigger/register",
"sdk/dynamictrigger/unregister"
]
"pages": ["sdk/dynamictrigger/register", "sdk/dynamictrigger/unregister"]
}
]
},
@@ -263,10 +252,7 @@
"sdk/dynamicschedule/constructor",
{
"group": "Instance methods",
"pages": [
"sdk/dynamicschedule/register",
"sdk/dynamicschedule/unregister"
]
"pages": ["sdk/dynamicschedule/register", "sdk/dynamicschedule/unregister"]
}
]
},
@@ -287,10 +273,7 @@
},
{
"group": "Overview",
"pages": [
"examples/introduction",
"examples/examples-repository"
]
"pages": ["examples/introduction", "examples/examples-repository"]
}
],
"footerSocials": {
@@ -303,4 +286,4 @@
"apiKey": "phc_hwYmedO564b3Ik8nhA4Csrb5SueY0EwFJWCbseGwWW"
}
}
}
}
+3 -4
View File
@@ -23,8 +23,7 @@ A useful tool when writing CRON expressions is [crontab guru](https://crontab.gu
<ResponseField name="options" type="object" required>
<Expandable title="options" defaultOpen>
<ResponseField name="cron" type="string" required>
A CRON expression that defines the schedule. Note that the timezone used
is always UTC.
A CRON expression that defines the schedule. Note that the timezone used is always UTC.
</ResponseField>
</Expandable>
</ResponseField>
@@ -32,7 +31,7 @@ A useful tool when writing CRON expressions is [crontab guru](https://crontab.gu
<RequestExample>
```typescript 9am UTC everyday
new Job(client, {
client.defineJob({
id: "scheduled-job-1",
name: "Scheduled Job 1",
version: "0.1.1",
@@ -51,7 +50,7 @@ new Job(client, {
```
```typescript First day of month
new Job(client, {
client.defineJob({
id: "scheduled-job-2",
name: "Scheduled Job 2",
version: "0.1.1",
+2 -2
View File
@@ -39,7 +39,7 @@ const dynamicSchedule = new DynamicSchedule(client, {
});
//2. create a Job that is attached to the dynamic schedule
new Job(client, {
client.defineJob({
id: "user-dynamicinterval",
name: "User Dynamic Interval",
version: "0.1.1",
@@ -63,7 +63,7 @@ async function registerUserCronJob(userId: string, userSchedule: string) {
}
//5. Register inside other Jobs
new Job(client, {
client.defineJob({
id: "register-dynamicinterval",
name: "Register Dynamic Interval",
version: "0.1.1",
+3 -5
View File
@@ -40,7 +40,7 @@ const dynamicOnIssueOpenedTrigger = new DynamicTrigger(client, {
});
//2. create a Job that is attached to the dynamic trigger
new Job(client, {
client.defineJob({
id: "listen-for-dynamic-trigger",
name: "Listen for dynamic trigger",
version: "0.1.1",
@@ -50,9 +50,7 @@ new Job(client, {
},
run: async (payload, io, ctx) => {
await io.slack.postMessage("Slack 📝", {
text: `New Issue opened on repo: ${
payload.issue.html_url
}. \n\n${JSON.stringify(ctx)}`,
text: `New Issue opened on repo: ${payload.issue.html_url}. \n\n${JSON.stringify(ctx)}`,
channel: "C04GWUTDC3W",
});
},
@@ -68,7 +66,7 @@ async function registerRepo(owner: string, repo: string) {
}
//4. Register inside other Jobs
new Job(client, {
client.defineJob({
id: "new-repo",
name: "New repo",
version: "0.1.1",
+1 -1
View File
@@ -53,7 +53,7 @@ You can have multiple Jobs that subscribe to the same event, they will all trigg
```typescript eventTrigger()
//this Job subscribes to an event called new.user
new Job(client, {
client.defineJob({
id: "job-2",
name: "Second job",
version: "0.0.1",
+1 -1
View File
@@ -30,7 +30,7 @@ If you wish to Run a Job at an exact time or less frequently than once pr day yo
<RequestExample>
```typescript Every 5 minutes
new Job(client, {
client.defineJob({
id: "scheduled-job-1",
name: "Scheduled Job 1",
version: "0.1.1",
+6 -10
View File
@@ -20,9 +20,8 @@ This is used inside the OpenAI Integration for Tasks like `backgroundCreateChatC
The HTTP method to use for the request.
</ResponseField>
<ResponseField name="headers" type="object">
Any headers to send with the request. Note that you can use
[redactString](sdk/redactString) to prevent sensitive information from being
stored (e.g. in the logs), like API keys and tokens.
Any headers to send with the request. Note that you can use [redactString](sdk/redactString) to
prevent sensitive information from being stored (e.g. in the logs), like API keys and tokens.
</ResponseField>
<ResponseField name="body" type="string | ArrayBuffer">
The body of the request.
@@ -84,19 +83,16 @@ An individual retrying strategy can be one of two types:
<Expandable title="headers strategy">
<ResponseField name="type" type="headers" required>
The `headers` strategy retries the request using info from the response
headers.
The `headers` strategy retries the request using info from the response headers.
</ResponseField>
<ResponseField name="limitHeader" type="string">
The header to use to determine the maximum number of times to retry the
request.
The header to use to determine the maximum number of times to retry the request.
</ResponseField>
<ResponseField name="remainingHeader" type="string">
The header to use to determine the number of remaining retries.
</ResponseField>
<ResponseField name="resetHeader" type="string">
The header to use to determine the time when the number of remaining retries
will be reset.
The header to use to determine the time when the number of remaining retries will be reset.
</ResponseField>
</Expandable>
@@ -111,7 +107,7 @@ A `Promise` that resolves after the specified amount of time.
<RequestExample>
```typescript backgroundFetch example
new Job(client, {
client.defineJob({
id: "background-fetch-job",
name: "Background fetch Job",
version: "0.0.1",
+4 -6
View File
@@ -19,9 +19,8 @@ description: "`io.registerCron()` allows you to register a [DynamicSchedule](/sd
<Expandable title="options" defaultOpen>
<ResponseField name="cron" type="string" required>
A CRON expression that defines the schedule. A useful tool when writing CRON
expressions is [crontab guru](https://crontab.guru). Note that the timezone
used is UTC.
A CRON expression that defines the schedule. A useful tool when writing CRON expressions is
[crontab guru](https://crontab.guru). Note that the timezone used is UTC.
</ResponseField>
</Expandable>
@@ -32,8 +31,7 @@ description: "`io.registerCron()` allows you to register a [DynamicSchedule](/sd
A Promise that resolves to an object with the following fields:
<ResponseField name="id" type="string" required>
A unique id for the interval. This is used to identify and unregister the
interval later.
A unique id for the interval. This is used to identify and unregister the interval later.
</ResponseField>
<ResponseField name="metadata" type="any" required>
Any additional metadata about the interval.
@@ -63,7 +61,7 @@ A Promise that resolves to an object with the following fields:
<RequestExample>
```typescript
new Job(client, {
client.defineJob({
id: "my-job",
name: "My job",
version: "0.1.1",
+2 -3
View File
@@ -30,8 +30,7 @@ description: "`io.registerInterval()` allows you to register a [DynamicSchedule]
A Promise that resolves to an object with the following fields:
<ResponseField name="id" type="string" required>
A unique id for the interval. This is used to identify and unregister the
interval later.
A unique id for the interval. This is used to identify and unregister the interval later.
</ResponseField>
<ResponseField name="metadata" type="any" required>
Any additional metadata about the interval.
@@ -61,7 +60,7 @@ A Promise that resolves to an object with the following fields:
<RequestExample>
```typescript
new Job(client, {
client.defineJob({
id: "my-job",
name: "My job",
version: "0.1.1",
+7 -13
View File
@@ -8,16 +8,13 @@ description: "`io.registerTrigger()` allows you to register a [DynamicTrigger](/
<Snippet file="stable-key-param.mdx" />
<ResponseField name="dynamicTrigger" type="DynamicTrigger" required>
A [DynamicTrigger](/sdk/dynamictrigger) that will trigger any Jobs it's
attached.
A [DynamicTrigger](/sdk/dynamictrigger) that will trigger any Jobs it's attached.
</ResponseField>
<ResponseField name="id" type="string" required>
A unique id for the registration. This is used to identify and unregister
later.
A unique id for the registration. This is used to identify and unregister later.
</ResponseField>
<ResponseField name="params" type="object" required>
The params for the DynamicTrigger. These will vary depending on the type of
the DynamicTrigger.
The params for the DynamicTrigger. These will vary depending on the type of the DynamicTrigger.
</ResponseField>
## Returns
@@ -25,8 +22,7 @@ description: "`io.registerTrigger()` allows you to register a [DynamicTrigger](/
A Promise that resolves to an object with the following fields:
<ResponseField name="id" type="string" required>
A unique id for the registration. This is used to identify and unregister
later.
A unique id for the registration. This is used to identify and unregister later.
</ResponseField>
<ResponseField name="key" type="string" required>
The key of the registration.
@@ -43,7 +39,7 @@ const dynamicOnIssueOpenedTrigger = new DynamicTrigger(client, {
});
//2. create a Job that is attached to the dynamic trigger
new Job(client, {
client.defineJob({
id: "listen-for-dynamic-trigger",
name: "Listen for dynamic trigger",
version: "0.1.1",
@@ -53,15 +49,13 @@ new Job(client, {
},
run: async (payload, io, ctx) => {
await io.slack.postMessage("Slack 📝", {
text: `New Issue opened on repo: ${
payload.issue.html_url
}. \n\n${JSON.stringify(ctx)}`,
text: `New Issue opened on repo: ${payload.issue.html_url}. \n\n${JSON.stringify(ctx)}`,
channel: "C04GWUTDC3W",
});
},
});
new Job(client, {
client.defineJob({
id: "new-repo",
name: "New repo",
version: "0.1.1",
+2 -2
View File
@@ -122,7 +122,7 @@ A Promise that resolves with the returned value of the callback.
<RequestExample>
```typescript Run a task
new Job(client, {
client.defineJob({
id: "alert-on-new-github-issues",
name: "Alert on new GitHub issues",
version: "0.1.1",
@@ -155,7 +155,7 @@ new Job(client, {
```
```typescript onError callback
new Job(client, {
client.defineJob({
id: "custom-error-handling",
name: "Custom Error handling",
version: "0.1.1",
+3 -4
View File
@@ -13,8 +13,7 @@ Use [eventTrigger()](/sdk/eventtrigger) on a Job to listen for events.
<Snippet file="stable-key-param.mdx" />
<ResponseField name="seconds" type="number" required>
The number of seconds to wait. This can be very long, serverless timeouts are
not an issue.
The number of seconds to wait. This can be very long, serverless timeouts are not an issue.
</ResponseField>
<Snippet file="send-event-params.mdx" />
@@ -27,7 +26,7 @@ Use [eventTrigger()](/sdk/eventtrigger) on a Job to listen for events.
```typescript Send an event
//this Job sends an event that triggers the second job
new Job(client, {
client.defineJob({
id: "job-1",
name: "First job",
version: "0.0.1",
@@ -45,7 +44,7 @@ new Job(client, {
},
});
new Job(client, {
client.defineJob({
id: "job-2",
name: "Second job",
version: "0.0.1",
+1 -1
View File
@@ -32,7 +32,7 @@ You have two options:
<RequestExample>
```typescript Using io.try()
new Job(client, {
client.defineJob({
id: "get-repo-info",
name: "GitHub get repo info",
version: "0.1.0",
+4 -5
View File
@@ -8,12 +8,11 @@ description: "`io.unregisterCron()` allows you to unregister a [DynamicSchedule]
<Snippet file="stable-key-param.mdx" />
<ResponseField name="dynamicSchedule" type="DynamicSchedule" required>
A [DynamicSchedule](/sdk/dynamicschedule) that will trigger any Jobs it's
attached to on a regular interval.
A [DynamicSchedule](/sdk/dynamicschedule) that will trigger any Jobs it's attached to on a regular
interval.
</ResponseField>
<ResponseField name="id" type="string" required>
A unique id for the schedule. This is used to identify and unregister the
schedule later.
A unique id for the schedule. This is used to identify and unregister the schedule later.
</ResponseField>
## Returns
@@ -27,7 +26,7 @@ A Promise with the following shape:
<RequestExample>
```typescript
new Job(client, {
client.defineJob({
id: "unregister-job",
name: "Unregister dynamic schedule",
version: "0.1.1",
+4 -5
View File
@@ -8,12 +8,11 @@ description: "`io.unregisterInterval()` allows you to unregister a [DynamicSched
<Snippet file="stable-key-param.mdx" />
<ResponseField name="dynamicSchedule" type="DynamicSchedule" required>
A [DynamicSchedule](/sdk/dynamicschedule) that will trigger any Jobs it's
attached to on a regular interval.
A [DynamicSchedule](/sdk/dynamicschedule) that will trigger any Jobs it's attached to on a regular
interval.
</ResponseField>
<ResponseField name="id" type="string" required>
A unique id for the interval. This is used to identify and unregister the
interval later.
A unique id for the interval. This is used to identify and unregister the interval later.
</ResponseField>
## Returns
@@ -27,7 +26,7 @@ A Promise with the following shape:
<RequestExample>
```typescript
new Job(client, {
client.defineJob({
id: "unregister-job",
name: "Unregister dynamic schedule",
version: "0.1.1",
+2 -3
View File
@@ -8,8 +8,7 @@ description: "`io.unregisterTrigger()` allows you to unregister a [DynamicTrigge
<Snippet file="stable-key-param.mdx" />
<ResponseField name="dynamicTrigger" type="DynamicTrigger" required>
A [DynamicTrigger](/sdk/dynamictrigger) that will trigger any Jobs it's
attached to.
A [DynamicTrigger](/sdk/dynamictrigger) that will trigger any Jobs it's attached to.
</ResponseField>
<ResponseField name="id" type="string" required>
A unique id for the trigger. This is used to identify and unregister it later.
@@ -26,7 +25,7 @@ A Promise with the following shape:
<RequestExample>
```typescript
new Job(client, {
client.defineJob({
id: "unregister-job",
name: "Unregister dynamic trigger",
version: "0.1.1",
+1 -1
View File
@@ -32,7 +32,7 @@ You must rethrow the error if this function returns `true`.
<RequestExample>
```typescript Using io.try()
new Job(client, {
client.defineJob({
id: "get-repo-info",
name: "GitHub get repo info",
version: "0.1.0",
+13 -17
View File
@@ -12,7 +12,7 @@ By far the most important thing to understand is the constructor.
<RequestExample>
```ts cronTrigger
new Job(client, {
client.defineJob({
id: "slack-kpi-summary",
name: "Slack kpi summary",
version: "0.1.1",
@@ -35,7 +35,7 @@ new Job(client, {
```
```ts webhook
new Job(client, {
client.defineJob({
id: "github-integration-on-issue",
name: "GitHub Integration - On Issue",
version: "0.1.0",
@@ -52,7 +52,7 @@ new Job(client, {
```
```ts event
new Job(client, {
client.defineJob({
id: "openai-joke",
name: "OpenAI Joke",
version: "0.0.1",
@@ -66,18 +66,15 @@ new Job(client, {
openai,
},
run: async (payload, io, ctx) => {
const joke = await io.openai.backgroundCreateChatCompletion(
"generate-jokes",
{
model: "gpt-3.5-turbo",
messages: [
{
role: "user",
content: payload.jokePrompt,
},
],
}
);
const joke = await io.openai.backgroundCreateChatCompletion("generate-jokes", {
model: "gpt-3.5-turbo",
messages: [
{
role: "user",
content: payload.jokePrompt,
},
],
});
return joke.choices;
},
@@ -89,8 +86,7 @@ new Job(client, {
## Parameters
<ParamField body="client" type="object" required>
An instance of [TriggerClient](/sdk/triggerclient) that is used to send events
to the Trigger API.
An instance of [TriggerClient](/sdk/triggerclient) that is used to send events to the Trigger API.
</ParamField>
<ParamField body="options" type="object" required>
+11 -1
View File
@@ -2,8 +2,18 @@
This project is meant to be used to create a catalog of jobs, usually to test something in an integration or the SDK.
## Setup
You will need to create a `.env` file. You can duplicate the `.env.example` file and set your local `TRIGGER_API_KEY` value.
### Running
You need to build the CLI:
```sh
pnpm run build --filter @trigger.dev/cli
```
Each file in `src` is a separate set of jobs that can be run separately. For example, the `src/stripe.ts` file can be run with:
```sh
@@ -15,7 +25,7 @@ This will open up a local server using `express` on port 8080. Then in a new ter
```sh
cd examples/job-catalog
pnpm run trigger:dev
pnpm run dev:trigger
```
### Adding a new file
+2
View File
@@ -4,6 +4,7 @@
"private": true,
"scripts": {
"stripe": "nodemon --watch src/stripe.ts -r tsconfig-paths/register -r dotenv/config src/stripe.ts",
"supabase": "nodemon --watch src/supabase.ts -r tsconfig-paths/register -r dotenv/config src/supabase.ts",
"dev:trigger": "trigger-cli dev --port 8080"
},
"dependencies": {
@@ -16,6 +17,7 @@
"@trigger.dev/slack": "workspace:*",
"@trigger.dev/stripe": "workspace:*",
"@trigger.dev/typeform": "workspace:*",
"@trigger.dev/supabase": "workspace:*",
"@types/node": "20.4.2",
"typescript": "5.1.6",
"zod": "3.21.4"
+124
View File
@@ -0,0 +1,124 @@
import { TriggerClient } from "@trigger.dev/sdk";
import { createExpressServer } from "@trigger.dev/express";
import { Supabase, SupabaseManagement } from "@trigger.dev/supabase";
const supabaseManagement = new SupabaseManagement({
id: "supabase-management",
apiKey: process.env["SUPABASE_API_KEY"]!,
});
const triggers = supabaseManagement.db<Database>(process.env["SUPABASE_ID"]!);
const supabase = new Supabase({
id: "supabase",
supabaseKey: process.env["SUPABASE_SERVICE_ROLE_KEY"]!,
supabaseUrl: process.env["SUPABASE_URL"]!,
});
export const client = new TriggerClient({
id: "job-catalog",
apiKey: process.env["TRIGGER_API_KEY"],
apiUrl: process.env["TRIGGER_API_URL"],
verbose: false,
ioLogLocalEnabled: true,
});
createExpressServer(client);
client.defineJob({
id: "supabase-management-example-1",
name: "Supabase Management Example 1",
version: "0.1.0",
trigger: triggers.onInserted({
table: "todos",
}),
run: async (payload, io, ctx) => {},
});
client.defineJob({
id: "supabase-management-example-2",
name: "Supabase Management Example 2",
version: "0.1.0",
trigger: triggers.onUpdated({
table: "todos",
}),
run: async (payload, io, ctx) => {},
});
client.defineJob({
id: "supabase-management-example-on",
name: "Supabase Management Example On",
version: "0.1.0",
trigger: triggers.on({
table: "todos",
events: ["INSERT", "UPDATE"],
}),
integrations: {
supabase,
},
run: async (payload, io, ctx) => {
const user = await io.supabase.runTask("fetch-user", async (db) => {
const { data, error } = await db.auth.admin.getUserById(payload.record.user_id);
if (error) {
throw error;
}
return data.user;
});
return user;
},
});
export type Json = string | number | boolean | null | { [key: string]: Json | undefined } | Json[];
export interface Database {
public: {
Tables: {
todos: {
Row: {
id: number;
inserted_at: string;
is_complete: boolean | null;
task: string | null;
user_id: string;
};
Insert: {
id?: number;
inserted_at?: string;
is_complete?: boolean | null;
task?: string | null;
user_id: string;
};
Update: {
id?: number;
inserted_at?: string;
is_complete?: boolean | null;
task?: string | null;
user_id?: string;
};
Relationships: [
{
foreignKeyName: "todos_user_id_fkey";
columns: ["user_id"];
referencedRelation: "users";
referencedColumns: ["id"];
},
];
};
};
Views: {
[_ in never]: never;
};
Functions: {
[_ in never]: never;
};
Enums: {
[_ in never]: never;
};
CompositeTypes: {
[_ in never]: never;
};
};
}
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/github
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/github",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official GitHub integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -29,8 +29,8 @@
"@octokit/request": "^6.2.5",
"@octokit/request-error": "^4.0.1",
"@octokit/webhooks": "^10.4.0",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.5",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"octokit": "^2.0.14",
"zod": "3.21.4"
},
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/slack
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/openai",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official OpenAI integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,8 +25,8 @@
},
"dependencies": {
"openai": "^3.3.0",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.5",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"zod": "3.21.4"
},
"engines": {
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/plain
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/plain",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official Plain.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"@trigger.dev/sdk": "workspace:^2.0.5",
"@team-plain/typescript-sdk": "^2.7.0"
},
"engines": {
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/resend
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/resend",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official Resend.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"@trigger.dev/sdk": "workspace:^2.0.5",
"resend": "^0.9.1"
},
"engines": {
+19
View File
@@ -1,5 +1,24 @@
# @trigger.dev/slack
## 2.0.5
### Patch Changes
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+2 -2
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/slack",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official Slack integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,7 +25,7 @@
},
"dependencies": {
"@slack/web-api": "^6.8.1",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.5",
"zod": "3.21.4"
},
"engines": {
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/stripe
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/stripe",
"version": "2.0.2",
"version": "2.0.5",
"description": "Trigger.dev integration for stripe",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"@trigger.dev/sdk": "workspace:^2.0.5",
"stripe": "^12.14.0",
"zod": "3.21.4"
},
+25
View File
@@ -1,5 +1,30 @@
# @trigger.dev/supabase
## 2.0.5
### Patch Changes
- 3cf3eaff: You can now trigger on multiple database events in the same job
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- ce901c92: Update supabase-management-js to 0.1.3
- feaa79ec: Automatically enable database webhooks when using triggers
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+5 -5
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/supabase",
"version": "2.0.2",
"version": "2.0.5",
"description": "Trigger.dev integration for @supabase/supabase-js",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,12 +26,12 @@
},
"dependencies": {
"@supabase/supabase-js": "^2.26.0",
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.2",
"supabase-management-js": "^0.1.2",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"@trigger.dev/sdk": "workspace:^2.0.5",
"supabase-management-js": "^0.1.4",
"zod": "3.21.4"
},
"engines": {
"node": ">=18.0.0"
}
}
}
+197 -2
View File
@@ -7,6 +7,7 @@ import {
IntegrationClient,
Logger,
TriggerIntegration,
isTriggerError,
} from "@trigger.dev/sdk";
import { SupabaseManagementAPI } from "supabase-management-js";
import { z } from "zod";
@@ -33,6 +34,91 @@ class SupabaseDatabase<Database = any> {
private projectRef: string
) {}
/**
* The function `on` creates a trigger for when a record is inserted, updated, or deleted on a
* specific table in a database schema.
* @param params - The `params` parameter is an object that contains the following properties:
* @param params.table - The `table` property is a string that specifies the name of the table
* that the trigger will be created for.
* @param params.events - The `events` property is an array of events that specifies the events
* that the trigger will be called for. The events that can be specified are `INSERT`, `UPDATE`, or `DELETE`.
* By default, the trigger will be called for all events.
* @param params.schema - The `schema` property is a string that specifies the name of the schema
* that the trigger will be created for. If the schema is not specified, the default schema will
* be used. (public)
* @param params.filter - The `filter` property is an object that specifies the filter that will
* be used to determine if the trigger should be called. If the filter is not specified, the
* trigger will be called for all records.
*
* @example
*
* ```ts
* const supabase = new SupabaseManagement({ id: "supabase" });
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
*
* client.defineJob({
* trigger: database.on({
* table: "todos",
* events: ["INSERTED", "UPDATED"],
* schema: "public",
* filter: {
* record: { is_completed: [false] },
* },
* }),
* })
* ```
*/
on<
SchemaName extends string & keyof Database = "public" extends keyof Database
? "public"
: string & keyof Database,
Schema extends GenericSchema = Database[SchemaName] extends GenericSchema
? Database[SchemaName]
: any,
TTableName extends string & keyof Schema["Tables"] = string & keyof Schema["Tables"],
TTable extends Schema["Tables"][TTableName] = Schema["Tables"][TTableName],
TEvents extends WebhookEvents[] = ["INSERT", "UPDATE", "DELETE"],
>(params: { table: TTableName; events?: TEvents; schema?: SchemaName; filter?: EventFilter }) {
return createTrigger<Prettify<UnionPayloads<TEvents, TTableName, SchemaName, TTable["Row"]>>>(
this.integration.source,
{
event: params.events ?? ["INSERT", "UPDATE", "DELETE"],
projectRef: this.projectRef,
...params,
}
);
}
/**
* The function `onInserted` creates a trigger for when a new record is inserted into a specific
* table in a database schema.
* @param params - The `params` parameter is an object that contains the following properties:
* @param params.table - The `table` property is a string that specifies the name of the table
* that the trigger will be created for.
* @param params.schema - The `schema` property is a string that specifies the name of the schema
* that the trigger will be created for. If the schema is not specified, the default schema will
* be used. (public)
* @param params.filter - The `filter` property is an object that specifies the filter that will
* be used to determine if the trigger should be called. If the filter is not specified, the
* trigger will be called for all records.
*
* @example
*
* ```ts
* const supabase = new SupabaseManagement({ id: "supabase" });
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
*
* client.defineJob({
* trigger: database.onInserted({
* table: "todos",
* schema: "public",
* filter: {
* record: { is_completed: [false] },
* },
* }),
* })
* ```
*/
onInserted<
SchemaName extends string & keyof Database = "public" extends keyof Database
? "public"
@@ -56,6 +142,37 @@ class SupabaseDatabase<Database = any> {
});
}
/**
* The function `onUpdated` creates a trigger for when a new record is updated on a specific
* table in a database schema.
* @param params - The `params` parameter is an object that contains the following properties:
* @param params.table - The `table` property is a string that specifies the name of the table
* that the trigger will be created for.
* @param params.schema - The `schema` property is a string that specifies the name of the schema
* that the trigger will be created for. If the schema is not specified, the default schema will
* be used. (public)
* @param params.filter - The `filter` property is an object that specifies the filter that will
* be used to determine if the trigger should be called. If the filter is not specified, the
* trigger will be called for all records.
*
* @example
*
* ```ts
* const supabase = new SupabaseManagement({ id: "supabase" });
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
*
* client.defineJob({
* trigger: database.onUpdated({
* table: "todos",
* schema: "public",
* filter: {
* record: { completed: [true] },
* old_record: { completed: [false] },
* },
* }),
* })
* ```
*/
onUpdated<
SchemaName extends string & keyof Database = "public" extends keyof Database
? "public"
@@ -79,6 +196,36 @@ class SupabaseDatabase<Database = any> {
});
}
/**
* The function `onDeleted` creates a trigger for when a new record is deleted from a specific
* table in a database schema.
* @param params - The `params` parameter is an object that contains the following properties:
* @param params.table - The `table` property is a string that specifies the name of the table
* that the trigger will be created for.
* @param params.schema - The `schema` property is a string that specifies the name of the schema
* that the trigger will be created for. If the schema is not specified, the default schema will
* be used. (public)
* @param params.filter - The `filter` property is an object that specifies the filter that will
* be used to determine if the trigger should be called. If the filter is not specified, the
* trigger will be called for all records.
*
* @example
*
* ```ts
* const supabase = new SupabaseManagement({ id: "supabase" });
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
*
* client.defineJob({
* trigger: database.onDeleted({
* table: "todos",
* schema: "public",
* filter: {
* old_record: { is_completed: [true] },
* },
* }),
* })
* ```
*/
onDeleted<
SchemaName extends string & keyof Database = "public" extends keyof Database
? "public"
@@ -174,9 +321,44 @@ type WebhookEventSource = ReturnType<typeof createWebhookEventSource>;
type WebhookEvents = "INSERT" | "UPDATE" | "DELETE";
type WebhookEventPayloads<
TTableName extends string,
TSchemaName extends string = "public",
TRecord = any,
> = {
INSERT: {
table: TTableName;
record: Prettify<TRecord>;
type: "INSERT";
schema: TSchemaName;
old_record: null;
};
UPDATE: {
table: TTableName;
record: Prettify<TRecord>;
type: "UPDATE";
schema: TSchemaName;
old_record: Prettify<TRecord>;
};
DELETE: {
table: TTableName;
record: null;
type: "DELETE";
schema: TSchemaName;
old_record: Prettify<TRecord>;
};
};
type UnionPayloads<
T extends WebhookEvents[],
TTableName extends string,
TSchemaName extends string = "public",
TRecord = any,
> = WebhookEventPayloads<TTableName, TSchemaName, TRecord>[T[number]];
function createTrigger<TEvent extends any>(
source: WebhookEventSource,
params: { event: WebhookEvents; filter?: EventFilter } & {
params: { event: WebhookEvents | WebhookEvents[]; filter?: EventFilter } & {
projectRef: string;
table: string;
schema?: string;
@@ -189,7 +371,7 @@ function createTrigger<TEvent extends any>(
icon: "supabase",
filter: {
...params.filter,
type: [params.event],
type: typeof params.event === "string" ? [params.event] : params.event,
schema: [params.schema ?? "public"],
},
properties: [],
@@ -293,6 +475,19 @@ export function createWebhookEventSource(
const url = new URL(httpSource.url);
const id = url.pathname.split("/").pop() ?? randomUUID();
try {
await io.integration.enableDatabaseWebhooks("enable-webhooks", { ref: params.projectRef });
} catch (error) {
if (isTriggerError(error)) {
throw error;
}
await io.logger.info(
"Enabling database webhooks failed, probably because it is already enabled. Continuing...",
{ error }
);
}
// Create the trigger name using the last 12 characters of the id
const triggerName = `tr_${id.slice(-12)}`;
@@ -181,3 +181,27 @@ export const getPGConfig: AuthenticatedTask<
};
},
};
/** Enable Database Webhooks in project */
export const enableDatabaseWebhooks: AuthenticatedTask<
SupabaseManagementAPI,
{ ref: string },
void
> = {
run: async (params, client) => {
return client.enableWebhooks(params.ref);
},
init: (params) => {
return {
name: "Enable Database Webhooks",
params,
icon: "supabase",
properties: [
{
label: "Project",
text: params.ref,
},
],
};
},
};
+22
View File
@@ -1,5 +1,27 @@
# @trigger.dev/typeform
## 2.0.5
### Patch Changes
- @trigger.dev/integration-kit@2.0.5
- @trigger.dev/sdk@2.0.5
## 2.0.4
### Patch Changes
- Updated dependencies [96384991]
- @trigger.dev/sdk@2.0.4
- @trigger.dev/integration-kit@2.0.4
## 2.0.3
### Patch Changes
- @trigger.dev/integration-kit@2.0.3
- @trigger.dev/sdk@2.0.3
## 2.0.2
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/typeform",
"version": "2.0.2",
"version": "2.0.5",
"description": "The official Typeform integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
},
"dependencies": {
"@typeform/api-client": "^1.8.0",
"@trigger.dev/sdk": "workspace:^2.0.2",
"@trigger.dev/integration-kit": "workspace:^2.0.2",
"@trigger.dev/sdk": "workspace:^2.0.5",
"@trigger.dev/integration-kit": "workspace:^2.0.5",
"zod": "3.21.4"
},
"engines": {
+17
View File
@@ -1,5 +1,22 @@
# create-trigger
## 2.0.5
## 2.0.4
### Patch Changes
- ff04bf44: Detect package manager from artifacts if they exist
- d5b8f829: The cli init command creates a jobs/index file and that is used to import jobs
- e7402978: Detect Next.js project by looking at dependencies if can't find next.config.js
## 2.0.3
### Patch Changes
- fc78854e: Don't require the TRIGGER_API_URL to be set (it has a default)
- 81180999: The dev command should use a POST request when doing the PING to the local server
## 2.0.2
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/cli",
"version": "2.0.2",
"version": "2.0.5",
"description": "The Trigger.dev CLI",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
+8 -7
View File
@@ -11,6 +11,7 @@ import { pathExists, readFile } from "../utils/fileSystem.js";
import { logger } from "../utils/logger.js";
import { resolvePath } from "../utils/parseNameAndPath.js";
import { TriggerApi } from "../utils/triggerApi.js";
import { CLOUD_API_URL } from "../consts.js";
export const DevCommandOptionsSchema = z.object({
port: z.coerce.number(),
@@ -71,7 +72,7 @@ export async function devCommand(path: string, anyOptions: any) {
try {
await fetch(localEndpointHandlerUrl, {
method: "HEAD",
method: "POST",
headers: {
"x-trigger-api-key": apiKey,
"x-trigger-action": "PING",
@@ -250,26 +251,26 @@ async function getTriggerApiDetails(path: string, envFile: string) {
]);
if (!resolvedEnvFile) {
logger.error(`You must add TRIGGER_API_KEY and TRIGGER_API_URL to your ${envFile} file.`);
logger.error(`You must add TRIGGER_API_KEY to your ${envFile} file.`);
return;
}
const parsedEnvFile = dotenv.parse(resolvedEnvFile.content);
if (!parsedEnvFile.TRIGGER_API_KEY || !parsedEnvFile.TRIGGER_API_KEY) {
logger.error(`You must add TRIGGER_API_KEY and TRIGGER_API_URL to your ${envFile} file.`);
if (!parsedEnvFile) {
logger.error(`You must add TRIGGER_API_KEY to your ${envFile} file.`);
return;
}
const apiKey = parsedEnvFile.TRIGGER_API_KEY;
const apiUrl = parsedEnvFile.TRIGGER_API_URL;
if (!apiKey || !apiUrl) {
logger.error(`You must add TRIGGER_API_KEY and TRIGGER_API_URL to your ${envFile} file.`);
if (!apiKey) {
logger.error(`You must add TRIGGER_API_KEY to your ${envFile} file.`);
return;
}
return { apiKey, apiUrl, envFile: resolvedEnvFile.fileName };
return { apiKey, apiUrl: apiUrl ?? CLOUD_API_URL, envFile: resolvedEnvFile.fileName };
}
async function resolveEndpointUrl(apiUrl: string, port: number) {
+39 -12
View File
@@ -30,10 +30,12 @@ export type InitCommandOptions = {
type ResolvedOptions = Required<InitCommandOptions>;
export const initCommand = async (options: InitCommandOptions) => {
renderTitle();
telemetryClient.init.started(options);
const resolvedPath = resolvePath(options.projectPath);
await renderTitle(resolvedPath);
if (options.triggerUrl === CLOUD_TRIGGER_URL) {
logger.info(`✨ Initializing project in Trigger.dev Cloud`);
} else if (typeof options.triggerUrl === "string") {
@@ -42,7 +44,6 @@ export const initCommand = async (options: InitCommandOptions) => {
logger.info(`✨ Initializing Trigger.dev in project`);
}
const resolvedPath = resolvePath(options.projectPath);
// Detect if are are in a Next.js project
const isNextJsProject = await detectNextJsProject(resolvedPath);
@@ -436,10 +437,11 @@ async function createTriggerAppRoute(
const tsConfigPath = pathModule.join(projectPath, configFileName);
const { tsconfig } = await parse(tsConfigPath);
const extension = isTypescriptProject ? ".ts" : ".js";
const triggerFileName = `trigger${extension}`;
const examplesFileName = `examples${extension}`;
const routeFileName = `route${extension}`;
const extension = isTypescriptProject ? ".ts" : ".js"
const triggerFileName = `trigger${extension}`
const examplesFileName = `examples${extension}`
const examplesIndexFileName = `index${extension}`
const routeFileName = `route${extension}`
const pathAlias = getPathAlias(tsconfig, usesSrcDir);
const routePathPrefix = pathAlias ? pathAlias + "/" : "../../../";
@@ -448,8 +450,8 @@ async function createTriggerAppRoute(
import { createAppRoute } from "@trigger.dev/nextjs";
import { client } from "${routePathPrefix}trigger";
// Replace this with your own jobs
import "${routePathPrefix}jobs/examples";
import "${routePathPrefix}jobs";
//this route is used to send and receive data with Trigger.dev
export const { POST, dynamic } = createAppRoute(client);
@@ -489,6 +491,12 @@ client.defineJob({
});
`;
const examplesIndexContent = `
// import all your job files here
export * from "./examples"
`
const directories = pathModule.join(path, "app", "api", "trigger");
await fs.mkdir(directories, { recursive: true });
@@ -519,6 +527,11 @@ client.defineJob({
if (!exampleFileExists) {
await fs.writeFile(pathModule.join(exampleDirectories, examplesFileName), jobsContent);
await fs.writeFile(
pathModule.join(exampleDirectories, examplesIndexFileName),
examplesIndexContent
);
logger.success(
`✅ Created example job at ${usesSrcDir ? "src/" : ""}jobs/examples/examplesFileName`
);
@@ -539,14 +552,17 @@ async function createTriggerPageRoute(
const pathAlias = getPathAlias(tsconfig, usesSrcDir);
const routePathPrefix = pathAlias ? pathAlias + "/" : "../..";
const extension = isTypescriptProject ? ".ts" : ".js";
const triggerFileName = `trigger${extension}`;
const examplesFileName = `examples${extension}`;
const extension = isTypescriptProject ? ".ts" : ".js"
const triggerFileName = `trigger${extension}`
const examplesFileName = `examples${extension}`
const examplesIndexFileName = `index${extension}`
const routeContent = `
import { createPagesRoute } from "@trigger.dev/nextjs";
import { client } from "${routePathPrefix}trigger";
import "${routePathPrefix}jobs";
//this route is used to send and receive data with Trigger.dev
const { handler, config } = createPagesRoute(client);
export { config };
@@ -588,6 +604,12 @@ client.defineJob({
});
`;
const examplesIndexContent = `
// import all your job files here
export * from "./examples"
`
const directories = pathModule.join(path, "pages", "api");
await fs.mkdir(directories, { recursive: true });
@@ -620,6 +642,11 @@ client.defineJob({
if (!exampleFileExists) {
await fs.writeFile(pathModule.join(exampleDirectories, examplesFileName), jobsContent);
await fs.writeFile(
pathModule.join(exampleDirectories, examplesIndexFileName),
examplesIndexContent
);
logger.success(
`✅ Created example job at ${usesSrcDir ? "src/" : ""}jobs/examples/${examplesFileName}`
);
+2 -2
View File
@@ -2,7 +2,7 @@ import chalk from "chalk";
import { execa } from "execa";
import ora, { type Ora } from "ora";
import pathModule from "path";
import { getUserPkgManager, type PackageManager } from "./getUserPkgManager.js";
import { getUserPackageManager, type PackageManager } from "./getUserPkgManager.js";
import fs from "fs/promises";
import fetch from "node-fetch";
import { z } from "zod";
@@ -18,7 +18,7 @@ export type InstalledPackage = {
};
export async function addDependencies(projectDir: string, packages: Array<InstallPackage>) {
const pkgManager = getUserPkgManager();
const pkgManager = await getUserPackageManager(projectDir);
const spinner = ora("Adding @trigger.dev dependencies to package.json...").start();
+20 -6
View File
@@ -1,15 +1,29 @@
import fs from "fs/promises";
import pathModule from "path";
import { readPackageJson } from "./readPackageJson.js";
/** Detects if the project is a Next.js project at path */
export async function detectNextJsProject(path: string): Promise<boolean> {
// Checks for the presence of a next.config.js file
try {
// Check if next.config.js file exists in the given path
await fs.access(pathModule.join(path, "next.config.js"));
const hasNextConfigFile = await detectNextConfigFile(path);
if (hasNextConfigFile) {
return true;
} catch (error) {
// If next.config.js file doesn't exist, it's not a Next.js project
}
return await detectNextDependency(path);
}
async function detectNextConfigFile(path: string): Promise<boolean> {
return fs
.access(pathModule.join(path, "next.config.js"))
.then(() => true)
.catch(() => false);
}
async function detectNextDependency(path: string): Promise<boolean> {
const packageJsonContent = await readPackageJson(path);
if (!packageJsonContent) {
return false;
}
return packageJsonContent.dependencies?.next !== undefined;
}
+31 -2
View File
@@ -1,6 +1,17 @@
import pathModule from "path";
import { pathExists } from "./fileSystem.js";
export type PackageManager = "npm" | "pnpm" | "yarn";
export const getUserPkgManager: () => PackageManager = () => {
export async function getUserPackageManager(path: string): Promise<PackageManager> {
try {
return detectPackageManagerFromArtifacts(path);
} catch (error) {
return detectPackageManagerFromCurrentCommand();
}
}
function detectPackageManagerFromCurrentCommand(): PackageManager {
// This environment variable is set by npm and yarn but pnpm seems less consistent
const userAgent = process.env.npm_config_user_agent;
@@ -16,4 +27,22 @@ export const getUserPkgManager: () => PackageManager = () => {
// If no user agent is set, assume npm
return "npm";
}
};
}
async function detectPackageManagerFromArtifacts(path: string): Promise<PackageManager> {
const packageFiles = [
{ name: "yarn.lock", pm: "yarn" } as const,
{ name: "pnpm-lock.yaml", pm: "pnpm" } as const,
{ name: "package-lock.json", pm: "npm" } as const,
{ name: "npm-shrinkwrap.json", pm: "npm" } as const,
];
for (const { name, pm } of packageFiles) {
const exists = await pathExists(pathModule.join(path, name));
if (exists) {
return pm;
}
}
throw new Error("Could not detect package manager from artifacts");
}
@@ -1,4 +1,4 @@
import { getUserPkgManager, type PackageManager } from "./getUserPkgManager.js";
import { getUserPackageManager, type PackageManager } from "./getUserPkgManager.js";
import { logger } from "./logger.js";
import ora, { type Ora } from "ora";
import chalk from "chalk";
@@ -7,7 +7,7 @@ import { execa } from "execa";
export async function installDependencies(projectDir: string) {
logger.info("Installing dependencies...");
const pkgManager = getUserPkgManager();
const pkgManager = await getUserPackageManager(projectDir);
const installSpinner = await runInstallCommand(pkgManager, projectDir);
+10
View File
@@ -0,0 +1,10 @@
import pathModule from "path";
import { type PackageJson } from "type-fest";
import { readJSONFile } from "./fileSystem.js";
export async function readPackageJson(directory: string): Promise<PackageJson | undefined> {
const packageJsonPath = pathModule.join(directory, "package.json");
return readJSONFile(packageJsonPath)
.then((f) => f as PackageJson)
.catch(() => undefined);
}
+3 -3
View File
@@ -1,6 +1,6 @@
import gradient from "gradient-string";
import { TITLE_TEXT } from "../consts.js";
import { getUserPkgManager } from "./getUserPkgManager.js";
import { getUserPackageManager } from "./getUserPkgManager.js";
// colors brought in from vscode poimandres theme
const poimandresTheme = {
@@ -12,11 +12,11 @@ const poimandresTheme = {
yellow: "#fffac2",
};
export const renderTitle = () => {
export const renderTitle = async (projectDirectory: string) => {
const triggerGradient = gradient(Object.values(poimandresTheme));
// resolves weird behavior where the ascii is offset
const pkgManager = getUserPkgManager();
const pkgManager = await getUserPackageManager(projectDirectory);
if (pkgManager === "yarn" || pkgManager === "pnpm") {
console.log("");
}

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