Compare commits
33 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| dfb00e84e9 | |||
| b5a64545bb | |||
| 3cf3eaff48 | |||
| 4ca758a3de | |||
| da81537009 | |||
| 7c1b13ba3b | |||
| 4b654e5b3b | |||
| 726a50ace7 | |||
| 9deffb67c4 | |||
| ff04bf44ee | |||
| e740297829 | |||
| a31705e198 | |||
| f0bdf53364 | |||
| 9638499163 | |||
| 9104c252b7 | |||
| 72cf345602 | |||
| d5b8f8299d | |||
| dd22b88d25 | |||
| 62b9c5879c | |||
| 1b0973fbc1 | |||
| fc78854ed4 | |||
| ed15bb0618 | |||
| ced9034428 | |||
| 7ec329d5d6 | |||
| 81180999de | |||
| feaa79ecf3 | |||
| d34dc0847c | |||
| 1b7c7520b4 | |||
| 34d77e830c | |||
| 87c302bf26 | |||
| d462c901f5 | |||
| 55c3d79cc8 | |||
| ce901c92e4 |
@@ -51,11 +51,6 @@ export function NoIntegrationSheet({
|
||||
)}
|
||||
</SheetHeader>
|
||||
<SheetBody>
|
||||
<Callout variant="info">
|
||||
We don’t 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>
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+121
@@ -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>
|
||||
);
|
||||
}
|
||||
+8
-2
@@ -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);
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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) {}
|
||||
}
|
||||
|
||||
@@ -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,5 +1,5 @@
|
||||
```typescript Wait example
|
||||
new Job(client, {
|
||||
client.defineJob({
|
||||
id: "delay-job",
|
||||
name: "Delay Job",
|
||||
version: "0.0.1",
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
```typescript
|
||||
new Job(client, {
|
||||
client.defineJob({
|
||||
//... other options
|
||||
integrations: {
|
||||
slack,
|
||||
|
||||
@@ -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>
|
||||
|
||||
@@ -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>
|
||||
|
||||
@@ -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">
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 |
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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[]
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
@@ -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:
|
||||
|
||||

|
||||
|
||||
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
|
||||
},
|
||||
});
|
||||
```
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
@@ -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"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
@@ -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>
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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;
|
||||
};
|
||||
};
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
},
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
},
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
},
|
||||
],
|
||||
};
|
||||
},
|
||||
};
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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,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",
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,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();
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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
Reference in New Issue
Block a user