diff --git a/.env.example b/.env.example index 48170e8f365..e6e12aef310 100644 --- a/.env.example +++ b/.env.example @@ -24,6 +24,9 @@ NODE_ENV=development # FROM_EMAIL= # REPLY_TO_EMAIL= +# Required for upgrading graphile-worker to v0.14 when graceful shutdown is impossible or impractical. +# FAIL_LOCKED_JOBS_FOR_MIGRATION=true + # CLOUD VARIABLES POSTHOG_PROJECT_KEY= PLAIN_API_KEY= diff --git a/apps/webapp/app/components/code/JSONEditor.tsx b/apps/webapp/app/components/code/JSONEditor.tsx index 7d5ef0c2378..5cebb9f8101 100644 --- a/apps/webapp/app/components/code/JSONEditor.tsx +++ b/apps/webapp/app/components/code/JSONEditor.tsx @@ -1,5 +1,6 @@ import { json as jsonLang } from "@codemirror/lang-json"; import type { ViewUpdate } from "@codemirror/view"; +import type { Text } from "@codemirror/state"; import { CheckIcon, ClipboardIcon } from "@heroicons/react/20/solid"; import type { ReactCodeMirrorProps, UseCodeMirror } from "@uiw/react-codemirror"; import { useCodeMirror } from "@uiw/react-codemirror"; @@ -68,6 +69,7 @@ export function JSONEditor(opts: JSONEditorProps) { theme: darkTheme(), indentWithTab: false, basicSetup, + selection: getDefaultSelection(defaultValue), onChange, onUpdate, }; @@ -83,10 +85,16 @@ export function JSONEditor(opts: JSONEditorProps) { //if the defaultValue changes update the editor useEffect(() => { if (view !== undefined) { - if (view.state.doc.toString() === defaultValue) return; - view.dispatch({ - changes: { from: 0, to: view.state.doc.length, insert: defaultValue }, - }); + if (view.state.doc.toString() !== defaultValue) { + view.dispatch({ + changes: { from: 0, to: view.state.doc.length, insert: defaultValue }, + selection: getDefaultSelection(defaultValue), + }); + } else { + view.dispatch({ + selection: getDefaultSelection(defaultValue), + }); + } } }, [defaultValue, view]); @@ -150,3 +158,20 @@ export function JSONEditor(opts: JSONEditorProps) { ); } + +function isMultiline(content: string | Text) { + if (typeof content === "string") { + return content.includes("\n"); + } else { + return content.lines > 1; + } +} + +/** For multiline content, gets end of penultimate line. Otherwise, position `0`. */ +function getDefaultSelection(content: string | Text) { + if (!isMultiline(content)) { + return { anchor: 0 }; + } else { + return { anchor: content.length > 2 ? content.length - 2 : 0 }; + } +} diff --git a/apps/webapp/app/components/run/TriggerDetail.tsx b/apps/webapp/app/components/run/TriggerDetail.tsx index aee57b22eac..4516eb9e3e1 100644 --- a/apps/webapp/app/components/run/TriggerDetail.tsx +++ b/apps/webapp/app/components/run/TriggerDetail.tsx @@ -15,17 +15,19 @@ import { DisplayProperty } from "@trigger.dev/core"; export function TriggerDetail({ trigger, + payload, event, properties, }: { trigger: DetailedEvent; + payload: string; event: { title: string; icon: string; }; properties: DisplayProperty[]; }) { - const { id, name, payload, context, timestamp, deliveredAt } = trigger; + const { id, name, context, timestamp, deliveredAt } = trigger; return ( diff --git a/apps/webapp/app/entry.server.tsx b/apps/webapp/app/entry.server.tsx index c270926697f..6ed9e2a3fad 100644 --- a/apps/webapp/app/entry.server.tsx +++ b/apps/webapp/app/entry.server.tsx @@ -164,5 +164,12 @@ function logError(error: unknown, request?: Request) { ); } } + console.error(error); + + if (error instanceof Error && error.message === "division by zero") { + console.log("⚠️ possible graphile-worker migration issue detected"); + console.log("⚠️ set FAIL_LOCKED_JOBS_FOR_MIGRATION=true if this persists"); + console.log("⚠️ see: https://trigger.dev/docs/documentation/guides/self-hosting/graphile-migration"); + } } diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 4e229834c42..b27248c8a51 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -41,6 +41,7 @@ const EnvironmentSchema = z.object({ WORKER_ENABLED: z.string().default("true"), EXECUTION_WORKER_ENABLED: z.string().default("true"), GRACEFUL_SHUTDOWN_TIMEOUT: z.coerce.number().int().default(60000), + FAIL_LOCKED_JOBS_FOR_MIGRATION: z.string().default("false"), }); export type Environment = z.infer; diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts index 2a9e823ff96..daef241033f 100644 --- a/apps/webapp/app/platform/zodWorker.server.ts +++ b/apps/webapp/app/platform/zodWorker.server.ts @@ -1,7 +1,7 @@ import type { CronItem, CronItemOptions, - Job as GraphileJob, + DbJob as GraphileJob, Runner as GraphileRunner, JobHelpers, RunnerOptions, @@ -29,8 +29,8 @@ const RawCronPayloadSchema = z.object({ const GraphileJobSchema = z.object({ id: z.coerce.string(), - queue_name: z.string().nullable(), - task_identifier: z.string(), + job_queue_id: z.number().nullable(), + task_id: z.number(), payload: z.unknown(), priority: z.number(), run_at: z.coerce.date(), @@ -67,7 +67,7 @@ type RecurringTaskPayload = { export type ZodRecurringTasks = { [key: string]: { - pattern: string; + match: string; options?: CronItemOptions; handler: (payload: RecurringTaskPayload, job: GraphileJob) => Promise; }; @@ -411,7 +411,7 @@ export class ZodWorker { if (this.#cleanup) { cronItems.push({ - pattern: this.#cleanup.frequencyExpression, + match: this.#cleanup.frequencyExpression, identifier: CLEANUP_TASK_NAME, task: CLEANUP_TASK_NAME, options: this.#cleanup.taskOptions, @@ -420,7 +420,7 @@ export class ZodWorker { if (this.#reporter) { cronItems.push({ - pattern: "50 * * * *", // Every hour at 50 minutes past the hour + match: "50 * * * *", // Every hour at 50 minutes past the hour identifier: REPORTER_TASK_NAME, task: REPORTER_TASK_NAME, }); @@ -432,7 +432,7 @@ export class ZodWorker { for (const [key, task] of Object.entries(this.#recurringTasks)) { const cronItem: CronItem = { - pattern: task.pattern, + match: task.match, identifier: key, task: key, options: task.options, diff --git a/apps/webapp/app/presenters/RunPresenter.server.ts b/apps/webapp/app/presenters/RunPresenter.server.ts index 34934647d0d..239b4fde8a5 100644 --- a/apps/webapp/app/presenters/RunPresenter.server.ts +++ b/apps/webapp/app/presenters/RunPresenter.server.ts @@ -78,6 +78,7 @@ export class RunPresenter { slug: run.environment.slug, }, event: this.#prepareEventData(run.event), + payload: run.payload, tasks, runConnections: run.runConnections, missingConnections: run.missingConnections, @@ -104,6 +105,7 @@ export class RunPresenter { query({ id, userId }: RunOptions) { return this.#prismaClient.jobRun.findFirst({ select: { + payload: true, id: true, number: true, status: true, diff --git a/apps/webapp/app/presenters/TestJobPresenter.server.ts b/apps/webapp/app/presenters/TestJobPresenter.server.ts index c87b035c8db..af016a89014 100644 --- a/apps/webapp/app/presenters/TestJobPresenter.server.ts +++ b/apps/webapp/app/presenters/TestJobPresenter.server.ts @@ -1,4 +1,4 @@ -import { User } from "@trigger.dev/database"; +import { Prisma, User } from "@trigger.dev/database"; import { replacements } from "@trigger.dev/core"; import { PrismaClient, prisma } from "~/db.server"; import { Job } from "~/models/job.server"; @@ -74,6 +74,7 @@ export class TestJobPresenter { createdAt: true, number: true, status: true, + payload: true, event: { select: { payload: true, @@ -111,7 +112,7 @@ export class TestJobPresenter { alias.version.examples.map((example) => ({ ...example, icon: example.icon ?? undefined, - payload: example.payload ? JSON.stringify(example.payload, exampleReplacer, 2) : undefined, + payload: prettyJsonValue(example.payload, exampleReplacer), })) ); @@ -127,13 +128,16 @@ export class TestJobPresenter { ), })), examples, - runs: job.runs.map((r) => ({ - id: r.id, - number: r.number, - status: r.status, - created: r.createdAt, - payload: r.event.payload ? JSON.stringify(r.event.payload, null, 2) : undefined, - })), + runs: job.runs.map((r) => { + const payload = r.payload ?? r.event.payload; + return { + id: r.id, + number: r.number, + status: r.status, + created: r.createdAt, + payload: prettyJsonValue(payload), + }; + }), }; } } @@ -153,3 +157,16 @@ function exampleReplacer(key: string, value: any) { return value; } + +function prettyJsonValue( + value: Prisma.JsonValue, + replacer?: (key: string, value: any) => any, + space = 2 +) { + if (value === null) { + return; + } + + const pretty = JSON.stringify(value, replacer, space); + return pretty === "{}" ? "{\n \n}" : pretty; +} diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx index dcdd2b7e652..7a4cc43c8ae 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.runs.$runParam.trigger/route.tsx @@ -28,5 +28,14 @@ export default function Page() { const job = useJob(); const run = useRun(); - return ; + const payload = run.payload !== null ? JSON.stringify(run.payload, null, 2) : trigger.payload; + + return ( + + ); } diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx index d109230a3dd..9051c207370 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.jobs.$jobParam.test/route.tsx @@ -117,7 +117,7 @@ export const handle: Handle = { breadcrumb: (match) => , }; -const startingJson = "{\n\n}"; +const startingJson = "{\n \n}"; export default function Page() { const { environments, runs, examples } = useTypedLoaderData(); @@ -202,7 +202,6 @@ export default function Page() { //deselect the example if it's been edited if (selectedCodeSampleId) { if (v !== selectedCodeSample) { - setDefaultJson(v); setSelectedCodeSampleId(undefined); } } diff --git a/apps/webapp/app/services/events/deliverBatchedEvent.server.ts b/apps/webapp/app/services/events/deliverBatchedEvent.server.ts new file mode 100644 index 00000000000..2f530bfcdba --- /dev/null +++ b/apps/webapp/app/services/events/deliverBatchedEvent.server.ts @@ -0,0 +1,206 @@ +import type { EventDispatcher, EventRecord } from "@trigger.dev/database"; +import type { EventFilter } from "@trigger.dev/core"; +import { EventFilterSchema, eventFilterMatches } from "@trigger.dev/core"; +import { $transaction, PrismaClientOrTransaction, prisma } from "~/db.server"; +import { logger } from "~/services/logger.server"; +import { workerQueue } from "../worker.server"; + +export class DeliverBatchedEventService { + #prismaClient: PrismaClientOrTransaction; + + constructor(prismaClient: PrismaClientOrTransaction = prisma) { + this.#prismaClient = prismaClient; + } + + public async call(ids: string[]) { + await $transaction( + this.#prismaClient, + async (tx) => { + const eventRecords = await tx.eventRecord.findMany({ + where: { + id: { in: ids }, + }, + include: { + environment: { + include: { + organization: true, + project: true, + }, + }, + }, + }); + + if (!eventRecords.length) { + throw new Error("No event records found."); + } + + const environmentId = eventRecords[0].environmentId; + + if (!eventRecords.every((event) => event.environmentId === environmentId)) { + throw new Error("Cross-environment batched events should not exist."); + } + + const unique = (val: unknown, i: number, array: unknown[]) => array.indexOf(val) === i; + + const uniqueEventNames = eventRecords.map((event) => event.name).filter(unique); + const uniqueEventSources = eventRecords.map((event) => event.source).filter(unique); + + const nameSourceCombinations: { name: string; source: string }[] = []; + + for (let nameIndex = 0; nameIndex < uniqueEventNames.length; nameIndex++) { + for (let sourceIndex = 0; sourceIndex < uniqueEventSources.length; sourceIndex++) { + nameSourceCombinations.push({ + name: uniqueEventNames[nameIndex], + source: uniqueEventSources[sourceIndex], + }); + } + } + + type InvocableDispatcher = { dispatcherId: string; eventRecordIds: string[] }; + + const batchDispatchersToInvoke: InvocableDispatcher[] = []; + const nonBatchDispatchersToInvoke: InvocableDispatcher[] = []; + + for (const combination of nameSourceCombinations) { + const matchingEvents = eventRecords.filter( + (event) => event.name === combination.name && event.source === combination.source + ); + + if (!matchingEvents.length) { + continue; + } + + const possibleEventDispatchers = await tx.eventDispatcher.findMany({ + where: { + environmentId, + event: { + has: combination.name, + }, + source: combination.source, + enabled: true, + manual: false, + }, + }); + + logger.debug("Found possible event dispatchers", { + possibleEventDispatchers, + eventRecords: matchingEvents.map((event) => event.id), + }); + + // filter events that match dispatcher filters + for (const eventDispatcher of possibleEventDispatchers) { + const filteredEvents = matchingEvents.filter((event) => + this.#evaluateEventRule(eventDispatcher, event) + ); + + // don't invoke dispatchers without any events + if (!filteredEvents.length) { + continue; + } + + const invocableDispatcher = { + dispatcherId: eventDispatcher.id, + eventRecordIds: filteredEvents.map((event) => event.id), + }; + + if (eventDispatcher.batch) { + batchDispatchersToInvoke.push(invocableDispatcher); + } else { + nonBatchDispatchersToInvoke.push(invocableDispatcher); + } + } + } + + const eventRecordIds = eventRecords.map((event) => event.id); + + if (!batchDispatchersToInvoke.length && !nonBatchDispatchersToInvoke.length) { + logger.debug("No matching event dispatchers", { + eventRecords: eventRecordIds, + }); + + return; + } + + logger.debug("Found matching batch event dispatchers", { batchDispatchersToInvoke }); + + await Promise.all( + batchDispatchersToInvoke.map((dispatcher) => { + return workerQueue.enqueue( + "events.invokeDispatcher", + { + id: dispatcher.dispatcherId, + eventRecordIds: dispatcher.eventRecordIds, + }, + { tx } + ); + }) + ); + + logger.debug("Found matching non-batch event dispatchers", { nonBatchDispatchersToInvoke }); + + await Promise.all( + nonBatchDispatchersToInvoke.map(async (dispatcher) => { + // sequentially enqueue single events to preserve order + for (const eventRecordId of dispatcher.eventRecordIds) { + await workerQueue.enqueue( + "events.invokeDispatcher", + { + id: dispatcher.dispatcherId, + eventRecordIds: [eventRecordId], + }, + { tx } + ); + } + }) + ); + + await tx.eventRecord.updateMany({ + where: { + id: { in: eventRecordIds }, + }, + data: { + deliveredAt: new Date(), + }, + }); + }, + { timeout: 10000 } + ); + } + + #evaluateEventRule(dispatcher: EventDispatcher, eventRecord: EventRecord): boolean { + if (!dispatcher.payloadFilter && !dispatcher.contextFilter) { + return true; + } + + const payloadFilter = EventFilterSchema.safeParse(dispatcher.payloadFilter ?? {}); + + const contextFilter = EventFilterSchema.safeParse(dispatcher.contextFilter ?? {}); + + if (!payloadFilter.success || !contextFilter.success) { + logger.error("Invalid event filter", { + payloadFilter, + contextFilter, + }); + return false; + } + + const eventMatcher = new EventMatcher(eventRecord); + + return eventMatcher.matches({ + payload: payloadFilter.data, + context: contextFilter.data, + }); + } +} + +export class EventMatcher { + event: EventRecord; + + constructor(event: EventRecord) { + this.event = event; + } + + public matches(filter: EventFilter) { + return eventFilterMatches(this.event, filter); + } +} diff --git a/apps/webapp/app/services/events/deliverEvent.server.ts b/apps/webapp/app/services/events/deliverEvent.server.ts index 32e13d7212c..6e9113fef5a 100644 --- a/apps/webapp/app/services/events/deliverEvent.server.ts +++ b/apps/webapp/app/services/events/deliverEvent.server.ts @@ -70,7 +70,7 @@ export class DeliverEventService { "events.invokeDispatcher", { id: eventDispatcher.id, - eventRecordId: eventRecord.id, + eventRecordIds: [eventRecord.id], }, { tx } ) diff --git a/apps/webapp/app/services/events/ingestSendEvent.server.ts b/apps/webapp/app/services/events/ingestSendEvent.server.ts index 9e4a254bbad..579461d944a 100644 --- a/apps/webapp/app/services/events/ingestSendEvent.server.ts +++ b/apps/webapp/app/services/events/ingestSendEvent.server.ts @@ -10,6 +10,7 @@ type UpdateEventInput = { existingEventLog: EventRecord; reqEvent: RawEvent; deliverAt?: Date; + batchKey?: string; }; type CreateEventInput = { @@ -19,6 +20,7 @@ type CreateEventInput = { deliverAt?: Date; sourceContext?: { id: string; metadata?: any }; externalAccount?: ExternalAccount; + batchKey?: string; }; const EVENT_UPDATE_THRESHOLD_WINDOW_IN_MSECS = 5 * 1000; // 5 seconds @@ -81,7 +83,13 @@ export class IngestSendEvent { }); const eventLog = await (existingEventLog - ? this.updateEvent({ tx, existingEventLog, reqEvent: event, deliverAt }) + ? this.updateEvent({ + tx, + existingEventLog, + reqEvent: event, + deliverAt, + batchKey: options?.batchKey, + }) : this.createEvent({ tx, event, @@ -89,6 +97,7 @@ export class IngestSendEvent { deliverAt, sourceContext, externalAccount, + batchKey: options?.batchKey, })); return eventLog; @@ -116,6 +125,7 @@ export class IngestSendEvent { deliverAt, sourceContext, externalAccount, + batchKey, }: CreateEventInput) { const eventLog = await tx.eventRecord.create({ data: { @@ -134,12 +144,18 @@ export class IngestSendEvent { }, }); - await this.enqueueWorkerEvent(tx, eventLog); + await this.enqueueWorkerEvent(tx, eventLog, batchKey); return eventLog; } - private async updateEvent({ tx, existingEventLog, reqEvent, deliverAt }: UpdateEventInput) { + private async updateEvent({ + tx, + existingEventLog, + reqEvent, + deliverAt, + batchKey, + }: UpdateEventInput) { if (!this.shouldUpdateEvent(existingEventLog)) { logger.debug(`not updating event for event id: ${existingEventLog.eventId}`); return existingEventLog; @@ -159,7 +175,7 @@ export class IngestSendEvent { }, }); - await this.enqueueWorkerEvent(tx, updatedEventLog); + await this.enqueueWorkerEvent(tx, updatedEventLog, batchKey); return updatedEventLog; } @@ -170,16 +186,38 @@ export class IngestSendEvent { return eventLog.deliverAt >= thresholdTime; } - private async enqueueWorkerEvent(tx: PrismaClientOrTransaction, eventLog: EventRecord) { + private async enqueueWorkerEvent( + tx: PrismaClientOrTransaction, + eventLog: EventRecord, + batchKey?: string + ) { if (this.deliverEvents) { // Produce a message to the event bus - await workerQueue.enqueue( - "deliverEvent", - { - id: eventLog.id, - }, - { runAt: eventLog.deliverAt, tx, jobKey: `event:${eventLog.id}` } - ); + if (batchKey) { + await workerQueue.enqueue( + "deliverBatchedEvent", + // Batch jobs require array payloads. + [ + { + id: eventLog.id, + }, + ], + { + runAt: eventLog.deliverAt, + tx, + jobKey: `event:${eventLog.environmentId}:${batchKey}`, + jobKeyMode: "preserve_run_at", + } + ); + } else { + await workerQueue.enqueue( + "deliverEvent", + { + id: eventLog.id, + }, + { runAt: eventLog.deliverAt, tx, jobKey: `event:${eventLog.id}` } + ); + } } } } diff --git a/apps/webapp/app/services/events/invokeDispatcher.server.ts b/apps/webapp/app/services/events/invokeDispatcher.server.ts index f43412d19a2..0440e27abf4 100644 --- a/apps/webapp/app/services/events/invokeDispatcher.server.ts +++ b/apps/webapp/app/services/events/invokeDispatcher.server.ts @@ -26,7 +26,7 @@ export class InvokeDispatcherService { this.#prismaClient = prismaClient; } - public async call(id: string, eventRecordId: string) { + public async call(id: string, eventRecordIds: string[]) { const eventDispatcher = await this.#prismaClient.eventDispatcher.findUniqueOrThrow({ where: { id, @@ -49,15 +49,28 @@ export class InvokeDispatcherService { return; } - const eventRecord = await this.#prismaClient.eventRecord.findUniqueOrThrow({ + const foundEventRecords = await this.#prismaClient.eventRecord.findMany({ where: { - id: eventRecordId, + id: { in: eventRecordIds }, }, }); + if (!foundEventRecords.length) { + throw new Error("No event records found."); + } + + if (eventRecordIds.length !== foundEventRecords.length) { + logger.warn("Event record counts do not match", { + expected: eventRecordIds.length, + found: foundEventRecords.length, + }); + } + + const foundEventRecordIds = foundEventRecords.map((event) => event.id); + logger.debug("Invoking event dispatcher", { eventDispatcher, - eventRecord: eventRecord.id, + eventRecords: foundEventRecordIds, }); const dispatchable = DispatchableSchema.safeParse(eventDispatcher.dispatchable); @@ -85,10 +98,11 @@ export class InvokeDispatcherService { const createRunService = new CreateRunService(this.#prismaClient); await createRunService.call({ - eventId: eventRecord.id, + eventIds: foundEventRecordIds, job: jobVersion.job, version: jobVersion, environment: eventDispatcher.environment, + isBatched: eventDispatcher.batch, }); break; @@ -135,7 +149,7 @@ export class InvokeDispatcherService { const createRunService = new CreateRunService(this.#prismaClient); await createRunService.call({ - eventId: eventRecord.id, + eventIds: foundEventRecordIds, job: job, version: latestJobVersion, environment: eventDispatcher.environment, diff --git a/apps/webapp/app/services/jobs/registerJob.server.ts b/apps/webapp/app/services/jobs/registerJob.server.ts index 8e75e5376c8..e2e3135aa31 100644 --- a/apps/webapp/app/services/jobs/registerJob.server.ts +++ b/apps/webapp/app/services/jobs/registerJob.server.ts @@ -331,6 +331,7 @@ export class RegisterJobService { id: jobVersion.id, }, dispatchableId: job.id, + batch: !!trigger.batch, }, update: { event: @@ -343,6 +344,7 @@ export class RegisterJobService { id: jobVersion.id, }, enabled: true, + batch: !!trigger.batch, }, }); diff --git a/apps/webapp/app/services/jobs/testJob.server.ts b/apps/webapp/app/services/jobs/testJob.server.ts index ab95d10856c..d96cbd5cf57 100644 --- a/apps/webapp/app/services/jobs/testJob.server.ts +++ b/apps/webapp/app/services/jobs/testJob.server.ts @@ -101,7 +101,7 @@ export class TestJobService { return await createRunService.call({ environment, - eventId: eventLog.id, + eventIds: [eventLog.id], job: version.job, version, }); diff --git a/apps/webapp/app/services/runs/createRun.server.ts b/apps/webapp/app/services/runs/createRun.server.ts index 7c10897726f..ec6bbc8e771 100644 --- a/apps/webapp/app/services/runs/createRun.server.ts +++ b/apps/webapp/app/services/runs/createRun.server.ts @@ -3,6 +3,7 @@ import { $transaction, PrismaClientOrTransaction } from "~/db.server"; import { prisma } from "~/db.server"; import { workerQueue } from "~/services/worker.server"; import type { AuthenticatedEnvironment } from "../apiAuth.server"; +import { logger } from "../logger.server"; export class CreateRunService { #prismaClient: PrismaClientOrTransaction; @@ -13,15 +14,21 @@ export class CreateRunService { public async call({ environment, - eventId, job, version, + eventIds, + isBatched = false, }: { environment: AuthenticatedEnvironment; - eventId: string; job: Job; version: JobVersion; + eventIds: string[]; + isBatched?: boolean; }) { + if (!eventIds.length) { + return; + } + const endpoint = await this.#prismaClient.endpoint.findUniqueOrThrow({ where: { id: version.endpointId, @@ -34,12 +41,23 @@ export class CreateRunService { }, }); - const eventRecord = await this.#prismaClient.eventRecord.findUniqueOrThrow({ + const eventRecords = await this.#prismaClient.eventRecord.findMany({ where: { - id: eventId, + id: { in: eventIds }, }, }); + if (!eventRecords.length) { + throw new Error("No event records found."); + } + + if (eventRecords.length !== eventIds.length) { + logger.warn("Event record counts do not match", { + expected: eventIds.length, + found: eventRecords.length, + }); + } + return await $transaction(this.#prismaClient, async (tx) => { // Get the current max number for the given jobId const latestJob = await tx.jobRun.findFirst({ @@ -53,6 +71,8 @@ export class CreateRunService { // Increment the number for the new execution const newNumber = (latestJob?.number ?? 0) + 1; + const firstEvent = eventRecords[0]; + // Create the new execution with the incremented number const run = await tx.jobRun.create({ data: { @@ -60,16 +80,20 @@ export class CreateRunService { preprocess: version.preprocessRuns, jobId: job.id, versionId: version.id, - eventId: eventId, + eventId: firstEvent.id, + eventIds: eventRecords.map((event) => event.id), environmentId: environment.id, organizationId: environment.organizationId, + payload: isBatched + ? eventRecords.map((event) => event.payload) ?? [{}] + : firstEvent.payload ?? {}, projectId: environment.projectId, endpointId: endpoint.id, queueId: jobQueue.id, - externalAccountId: eventRecord.externalAccountId - ? eventRecord.externalAccountId + externalAccountId: firstEvent.externalAccountId + ? firstEvent.externalAccountId : undefined, - isTest: eventRecord.isTest, + isTest: firstEvent.isTest, internal: job.internal, }, }); diff --git a/apps/webapp/app/services/runs/performRunExecutionV2.server.ts b/apps/webapp/app/services/runs/performRunExecutionV2.server.ts index f22d13edb1a..220f0ab39a2 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV2.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV2.server.ts @@ -77,6 +77,7 @@ export class PerformRunExecutionV2Service { const { response, parser } = await client.preprocessRunRequest({ event, + payload: run.payload as any, job: { id: run.version.job.slug, version: run.version.version, @@ -426,6 +427,7 @@ export class PerformRunExecutionV2Service { return { event, + payload: run.payload as any, job: { id: run.version.job.slug, version: run.version.version, @@ -465,6 +467,7 @@ export class PerformRunExecutionV2Service { return { event, + payload: run.payload as any, job: { id: run.version.job.slug, version: run.version.version, diff --git a/apps/webapp/app/services/runs/reRun.server.ts b/apps/webapp/app/services/runs/reRun.server.ts index 9a604c16b17..8c36fe54b08 100644 --- a/apps/webapp/app/services/runs/reRun.server.ts +++ b/apps/webapp/app/services/runs/reRun.server.ts @@ -69,7 +69,7 @@ export class ReRunService { organization: existingRun.organization, project: existingRun.project, }, - eventId: eventLog.id, + eventIds: [eventLog.id], job: existingRun.job, version: existingRun.version, }); diff --git a/apps/webapp/app/services/sources/deliverHttpSourceRequest.server.ts b/apps/webapp/app/services/sources/deliverHttpSourceRequest.server.ts index 2f002e1017a..bc4107d78fc 100644 --- a/apps/webapp/app/services/sources/deliverHttpSourceRequest.server.ts +++ b/apps/webapp/app/services/sources/deliverHttpSourceRequest.server.ts @@ -65,7 +65,7 @@ export class DeliverHttpSourceRequestService { httpSourceRequest.endpoint.url ); - const { response, events, metadata } = await clientApi.deliverHttpSourceRequest({ + const { response, events, metadata, options } = await clientApi.deliverHttpSourceRequest({ key: httpSourceRequest.source.key, dynamicId: httpSourceRequest.source.dynamicTrigger?.slug, secret: secret.secret, @@ -108,6 +108,9 @@ export class DeliverHttpSourceRequestService { httpSourceRequest.environment, event, { + batchKey: options?.batchKey, + deliverAt: options?.deliverAt, + deliverAfter: options?.deliverAfter, accountId: httpSourceRequest.source.externalAccount?.identifier, }, httpSourceRequest.source.dynamicSourceId diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts index cb9b4d6b0e9..287c5b211eb 100644 --- a/apps/webapp/app/services/worker.server.ts +++ b/apps/webapp/app/services/worker.server.ts @@ -9,6 +9,7 @@ import { IndexEndpointService } from "./endpoints/indexEndpoint.server"; import { PerformEndpointIndexService } from "./endpoints/performEndpointIndexService"; import { RecurringEndpointIndexService } from "./endpoints/recurringEndpointIndex.server"; import { DeliverEventService } from "./events/deliverEvent.server"; +import { DeliverBatchedEventService } from "./events/deliverBatchedEvent.server"; import { InvokeDispatcherService } from "./events/invokeDispatcher.server"; import { integrationAuthRepository } from "./externalApis/integrationAuthRepository.server"; import { IntegrationConnectionCreatedService } from "./externalApis/integrationConnectionCreated.server"; @@ -61,9 +62,10 @@ const workerCatalog = { ]) ), deliverEvent: z.object({ id: z.string() }), + deliverBatchedEvent: z.array(z.object({ id: z.string() })), "events.invokeDispatcher": z.object({ id: z.string(), - eventRecordId: z.string(), + eventRecordIds: z.array(z.string()), }), "events.deliverScheduled": z.object({ id: z.string(), @@ -121,6 +123,10 @@ if (env.NODE_ENV === "production") { } export async function init() { + if (env.FAIL_LOCKED_JOBS_FOR_MIGRATION === "true") { + await failLockedJobsForMigration(); + } + if (env.WORKER_ENABLED === "true") { await workerQueue.initialize(); } @@ -130,6 +136,64 @@ export async function init() { } } +/** Helper for graphile-worker v0.14.0 migration */ +async function failLockedJobsForMigration() { + console.log("⚠️ failing locked jobs for migration"); + + const graphileWorkerSchema = env.WORKER_SCHEMA; + + const migrationQueryResult = await prisma.$queryRawUnsafe(` + SELECT id FROM graphile_worker.migrations + ORDER BY id DESC LIMIT 1`); + + const MigrationQueryResultSchema = z.array(z.object({ id: z.number() })); + + const migrationResults = MigrationQueryResultSchema.parse(migrationQueryResult); + + if (!migrationResults.length) { + console.log("⚠️ nothing to do, no migrations applied yet"); + console.log("⚠️ unset FAIL_LOCKED_JOBS_FOR_MIGRATION to hide these messages"); + return; + } + + const latestMigration = migrationResults[0].id; + + // the first v0.14.0 migration has ID 11 + if (latestMigration > 10) { + console.log("⚠️ nothing to do, graphile-worker already upgraded"); + console.log("⚠️ unset FAIL_LOCKED_JOBS_FOR_MIGRATION to hide these messages"); + return; + } + + const errorMessage = "Failing locked jobs for migration"; + + const failJobsResult = await prisma.$queryRawUnsafe>( + `WITH j AS ( + UPDATE ${graphileWorkerSchema}.jobs + SET + last_error = $1::text, + run_at = greatest(now(), run_at) + (exp(least(attempts, 10)) * interval '1 second'), + locked_by = null, + locked_at = null + WHERE locked_at is not null + AND locked_at > now() - interval '4 hours' + RETURNING * + ), queues AS ( + UPDATE ${graphileWorkerSchema}.job_queues + SET + locked_by = null, + locked_at = null + FROM j + WHERE job_queues.queue_name = j.queue_name + ) + SELECT * FROM j;`, + errorMessage + ); + + console.log(`⚠️ failed ${failJobsResult.length} locked jobs`); + console.log("⚠️ you can now unset FAIL_LOCKED_JOBS_FOR_MIGRATION"); +} + function getWorkerQueue() { return new ZodWorker({ name: "workerQueue", @@ -152,7 +216,7 @@ function getWorkerQueue() { recurringTasks: { // Run this every 5 minutes autoIndexProductionEndpoints: { - pattern: "*/5 * * * *", + match: "*/5 * * * *", handler: async (payload, job) => { const service = new RecurringEndpointIndexService(); @@ -161,7 +225,7 @@ function getWorkerQueue() { }, // Run this every hour purgeOldIndexings: { - pattern: "0 * * * *", + match: "0 * * * *", handler: async (payload, job) => { // Delete indexings that are older than 7 days await prisma.endpointIndex.deleteMany({ @@ -182,7 +246,7 @@ function getWorkerQueue() { handler: async (payload, job) => { const service = new InvokeDispatcherService(); - await service.call(payload.id, payload.eventRecordId); + await service.call(payload.id, payload.eventRecordIds); }, }, "events.deliverScheduled": { @@ -307,6 +371,15 @@ function getWorkerQueue() { await service.call(payload.id); }, }, + deliverBatchedEvent: { + priority: 0, // smaller number = higher priority + maxAttempts: 5, + handler: async (payload, job) => { + const service = new DeliverBatchedEventService(); + + await service.call(payload.map((p) => p.id)); + }, + }, refreshOAuthToken: { priority: 8, // smaller number = higher priority maxAttempts: 7, diff --git a/apps/webapp/package.json b/apps/webapp/package.json index 5cf5eacefa4..b35b0212482 100644 --- a/apps/webapp/package.json +++ b/apps/webapp/package.json @@ -75,7 +75,7 @@ "emails": "workspace:*", "express": "^4.18.1", "framer-motion": "^10.12.11", - "graphile-worker": "^0.13.0", + "graphile-worker": "0.14.0-rc.0", "highlight.run": "^7.3.4", "humanize-duration": "^3.27.3", "intl-parse-accept-language": "^1.0.0", diff --git a/docs/documentation/guides/self-hosting/graphile-migration.mdx b/docs/documentation/guides/self-hosting/graphile-migration.mdx new file mode 100644 index 00000000000..87ddffba334 --- /dev/null +++ b/docs/documentation/guides/self-hosting/graphile-migration.mdx @@ -0,0 +1,35 @@ +--- +title: "Graphile migration" +description: "The `graphile-worker` migration may require some manual input" +--- + +## See any warnings? + + + Graphile Migration Error + + +A warning like this may have brought you here. If it hasn't, you probably **don't** need this! + +## Semi-automatic migration steps + +This SHOULD be a safe process, but it's recommended to always have backups of your database. + +1. Set `FAIL_LOCKED_JOBS_FOR_MIGRATION=true` as you do with other env vars +2. Restart your server + +The environment variable won't have any effect if the migration was already successful. To remove additional log output, you may want to unset it again afterwards. + +## Alternative: Graceful shutdown + +The above issue only arises if there are any locked jobs when (automatic) migration is attempted. + +To prevent this, please shut your server down _gracefully_. How you do this will depend on your platform. This may involve: + +- Sending SIGTERM via the `kill` command +- Using `docker stop` or `docker-compose down` +- Equivalent methods on your managed platform + +## Problems + +If you run into any other issues, you may have to temporarily scale down to using only _a single instance_. All other instructions above still apply. \ No newline at end of file diff --git a/docs/images/graphile-migration-error.png b/docs/images/graphile-migration-error.png new file mode 100644 index 00000000000..6e9c6c1c2d3 Binary files /dev/null and b/docs/images/graphile-migration-error.png differ diff --git a/docs/mint.json b/docs/mint.json index 5abf8aa9e92..37110caf557 100644 --- a/docs/mint.json +++ b/docs/mint.json @@ -199,6 +199,7 @@ "documentation/guides/self-hosting/flyio", "documentation/guides/self-hosting/render", "documentation/guides/self-hosting/supabase", + "documentation/guides/self-hosting/graphile-migration", "documentation/guides/tunneling-platform" ] }, diff --git a/integrations/airtable/src/index.ts b/integrations/airtable/src/index.ts index bc9ff240261..c4b5b8e6f3f 100644 --- a/integrations/airtable/src/index.ts +++ b/integrations/airtable/src/index.ts @@ -10,9 +10,16 @@ import { type RunTaskOptions, type TriggerIntegration, } from "@trigger.dev/sdk"; -import AirtableSDK from "airtable"; +import AirtableSDK, { Error as AirtableApiError } from "airtable"; import { Base } from "./base"; -import { Webhooks, createWebhookEventSource } from "./webhooks"; +import * as events from "./events"; +import { + WebhookChangeType, + WebhookDataType, + Webhooks, + createTrigger, + createWebhookEventSource, +} from "./webhooks"; export * from "./types"; export * from "./base"; @@ -105,7 +112,7 @@ export class Airtable implements TriggerIntegration { ...(options ?? {}), connectionKey: this._connectionKey, }, - errorCallback + errorCallback ?? onError ); } @@ -113,20 +120,53 @@ export class Airtable implements TriggerIntegration { return new Base(this.runTask.bind(this), baseId); } - //todo these require batch support because they send too many events - // onTableChanges(params: { - // baseId: string; - // tableId?: string; - // changeTypes?: WebhookChangeType[]; - // dataTypes?: WebhookDataType[]; - // }) { - // return createTrigger(this.source, events.onTableChanged, params, { - // changeTypes: params.changeTypes, - // dataTypes: ["tableData", "tableFields", "tableMetadata"], - // }); - // } + onTableChanges(params: { + baseId: string; + tableId?: string; + changeTypes?: WebhookChangeType[]; + dataTypes?: WebhookDataType[]; + }) { + return createTrigger(this.source, events.onTableChanged, params, { + changeTypes: params.changeTypes, + dataTypes: ["tableData", "tableFields", "tableMetadata"], + }); + } webhooks() { return new Webhooks(this.runTask.bind(this)); } } + +function isAirtableApiError(error: unknown): error is AirtableApiError { + if (typeof error !== "object" || error === null) { + return false; + } + + const airtableError = error as AirtableApiError; + + return ( + typeof airtableError.error === "string" && + typeof airtableError.message === "string" && + typeof airtableError.statusCode === "number" + ); +} + +export function onError(error: unknown): ReturnType { + if (!isAirtableApiError(error)) { + return; + } + + if (error.statusCode === 429) { + // see: https://airtable.com/developers/web/api/rate-limits + return { + retryAt: new Date(Date.now() + 30 * 1000), + }; + } + + if (error.statusCode >= 400 && error.statusCode < 500) { + // see: https://airtable.com/developers/web/api/errors#user-error-codes + return { + skipRetrying: true, + }; + } +} diff --git a/integrations/airtable/src/webhooks.ts b/integrations/airtable/src/webhooks.ts index 786d9ef7e4f..3e03b17928a 100644 --- a/integrations/airtable/src/webhooks.ts +++ b/integrations/airtable/src/webhooks.ts @@ -6,7 +6,7 @@ import { IntegrationTaskKey, Logger, } from "@trigger.dev/sdk"; -import AirtableSDK from "airtable"; +import AirtableSDK, { Error as AirtableApiError } from "airtable"; import { z } from "zod"; import * as events from "./events"; import { Airtable, AirtableRunTask } from "./index"; @@ -46,6 +46,31 @@ type WebhookSpecification = { }; }; +const AirtableErrorBodySchema = z + .union([ + z.object({ + error: z.string(), + }), + z.object({ + error: z.object({ + type: z.string(), + message: z.string().optional(), + }), + }), + ]) + .transform((body) => { + if (typeof body.error === "string") { + return { + type: body.error, + }; + } else { + return { + type: body.error.type, + message: body.error.message, + }; + } + }); + const apiUrl = "https://api.airtable.com/v0/bases"; export class Webhooks { @@ -85,14 +110,7 @@ export class Webhooks { }); if (!response.ok) { - const errorText = await response - .text() - .then((t) => t) - .catch((e) => "No body"); - - throw new Error( - `Failed to create webhook: ${response.status} ${response.statusText}\n${errorText}` - ); + await handleWebhookError(response, "WEBHOOK_CREATE") } const webhook = await response.json(); @@ -123,7 +141,7 @@ export class Webhooks { }); if (!response.ok) { - throw new Error(`Failed to list webhooks: ${response.statusText}`); + await handleWebhookError(response, "WEBHOOK_LIST") } const webhook = await response.json(); @@ -153,7 +171,7 @@ export class Webhooks { }); if (!response.ok) { - throw new Error(`Failed to delete webhook: ${response.statusText}`); + await handleWebhookError(response, "WEBHOOK_DELETE") } }, { @@ -200,13 +218,15 @@ export function createTrigger( dataTypes: WebhookDataType[]; changeTypes?: WebhookChangeType[]; fromSources?: WebhookFromSource[]; - } + }, + batch = true ): CreateTriggersResult { return new ExternalSourceTrigger({ event, params, source, options, + batch, }); } @@ -244,6 +264,7 @@ export function createWebhookEventSource( id: "airtable.webhook", schema: z.object({ baseId: z.string(), tableId: z.string().optional() }), optionSchema: z.object({ + changeTypes: z.array(WebhookChangeTypeSchema).optional(), dataTypes: z.array(WebhookDataTypeSchema), fromSources: z.array(WebhookFromSourceSchema).optional(), }), @@ -271,7 +292,9 @@ export function createWebhookEventSource( const specification: WebhookSpecification = { filters: { dataTypes: options.dataTypes.desired as WebhookDataType[], - changeTypes: options.event.desired as WebhookChangeType[], + changeTypes: options.changeTypes + ? (options.changeTypes.desired as WebhookChangeType[]) + : ["add", "remove", "update"], fromSources: (options.fromSources?.desired ?? [ "client", "anonymousUser", @@ -406,6 +429,10 @@ async function webhookHandler(event: HandlerEvent<"HTTP">, logger: Logger, integ })) : [], metadata: response?.cursor ? { cursor: response.cursor } : undefined, + options: { + batchKey: "airtable-webhooks", + deliverAfter: 10, + }, }; } @@ -458,3 +485,21 @@ async function getPayload( const webhook = await response.json(); return ListWebhooksResponseSchema.parse(webhook); } + +async function handleWebhookError(response: Response, errorType: string) { + const rawErrorBody = await response.json(); + + const parsedErrorBody = AirtableErrorBodySchema.safeParse(rawErrorBody); + + if (!parsedErrorBody.success) { + throw new AirtableApiError( + `${errorType}_PARSE_ERROR`, + `${response.statusText}:\n${rawErrorBody}`, + response.status + ); + } + + const { type, message } = parsedErrorBody.data; + + throw new AirtableApiError(type, message ?? response.statusText, response.status); +} \ No newline at end of file diff --git a/packages/core/src/schemas/api.ts b/packages/core/src/schemas/api.ts index 4e16cbfc55d..3a44224339b 100644 --- a/packages/core/src/schemas/api.ts +++ b/packages/core/src/schemas/api.ts @@ -123,6 +123,13 @@ export const TriggerSourceSchema = z.object({ const HttpSourceResponseMetadataSchema = DeserializedJsonSchema; export type HttpSourceResponseMetadata = z.infer; +const HttpSourceResponseOptionsSchema = z.object({ + batchKey: z.string().optional(), + deliverAt: z.coerce.date().optional(), + deliverAfter: z.number().int().optional(), +}) +export type HttpSourceResponseOptions = z.infer; + export const HandleTriggerSourceSchema = z.object({ key: z.string(), secret: z.string(), @@ -408,6 +415,8 @@ export type ApiEventLog = z.infer; /** Options to control the delivery of the event */ export const SendEventOptionsSchema = z.object({ + /** An optional string to enable event batching. Events with a shared `batchKey` will be delivered together as an array of payloads. */ + batchKey: z.string().optional(), /** An optional Date when you want the event to trigger Jobs. The event will be sent to the platform immediately but won't be acted upon until the specified time. */ @@ -453,6 +462,7 @@ export type RunSourceContext = z.infer; export const RunJobBodySchema = z.object({ event: ApiEventLogSchema, + payload: DeserializedJsonSchema, job: z.object({ id: z.string(), version: z.string(), @@ -563,6 +573,7 @@ export type RunJobResponse = z.infer; export const PreprocessRunBodySchema = z.object({ event: ApiEventLogSchema, + payload: DeserializedJsonSchema, job: z.object({ id: z.string(), version: z.string(), @@ -778,6 +789,7 @@ export const HttpSourceResponseSchema = z.object({ response: NormalizedResponseSchema, events: z.array(RawEventSchema), metadata: HttpSourceResponseMetadataSchema.optional(), + options: HttpSourceResponseOptionsSchema.optional(), }); export const RegisterTriggerBodySchemaV1 = z.object({ diff --git a/packages/core/src/schemas/triggers.ts b/packages/core/src/schemas/triggers.ts index edc5256bc6e..04984d6af10 100644 --- a/packages/core/src/schemas/triggers.ts +++ b/packages/core/src/schemas/triggers.ts @@ -30,6 +30,7 @@ export const DynamicTriggerMetadataSchema = z.object({ export const StaticTriggerMetadataSchema = z.object({ type: z.literal("static"), + batch: z.boolean().optional(), title: z.union([z.string(), z.array(z.string())]), properties: z.array(DisplayPropertySchema).optional(), rule: EventRuleSchema, diff --git a/packages/database/package.json b/packages/database/package.json index 6c82f458825..4755ee2d29c 100644 --- a/packages/database/package.json +++ b/packages/database/package.json @@ -16,6 +16,7 @@ "db:migrate:dev": "prisma migrate dev", "db:migrate:dev:create": "prisma migrate dev --create-only", "db:migrate:deploy": "prisma migrate deploy", + "db:push": "prisma db push", "db:studio": "prisma studio", "typecheck": "tsc --noEmit" } diff --git a/packages/database/prisma/migrations/20231019123537_add_batched_events/migration.sql b/packages/database/prisma/migrations/20231019123537_add_batched_events/migration.sql new file mode 100644 index 00000000000..d5ffe12e0e7 --- /dev/null +++ b/packages/database/prisma/migrations/20231019123537_add_batched_events/migration.sql @@ -0,0 +1,6 @@ +-- AlterTable +ALTER TABLE "EventDispatcher" ADD COLUMN "batch" BOOLEAN NOT NULL DEFAULT false; + +-- AlterTable +ALTER TABLE "JobRun" ADD COLUMN "eventIds" TEXT[], +ADD COLUMN "payload" JSONB; diff --git a/packages/database/prisma/schema.prisma b/packages/database/prisma/schema.prisma index 7f1d219181c..919cfff6e8a 100644 --- a/packages/database/prisma/schema.prisma +++ b/packages/database/prisma/schema.prisma @@ -613,6 +613,7 @@ model EventDispatcher { payloadFilter Json? contextFilter Json? manual Boolean @default(false) + batch Boolean @default(false) dispatchableId String dispatchable Json @@ -686,6 +687,10 @@ model JobRun { event EventRecord @relation(fields: [eventId], references: [id], onDelete: Cascade, onUpdate: Cascade) eventId String + eventIds String[] + + payload Json? + organization Organization @relation(fields: [organizationId], references: [id], onDelete: Cascade, onUpdate: Cascade) organizationId String diff --git a/packages/trigger-sdk/src/io.ts b/packages/trigger-sdk/src/io.ts index b4dacfc5e0f..1a4c0af5d35 100644 --- a/packages/trigger-sdk/src/io.ts +++ b/packages/trigger-sdk/src/io.ts @@ -330,6 +330,9 @@ export class IO { text: event.name, }, ...(event?.id ? [{ label: "ID", text: event.id }] : []), + ...(options?.batchKey ? [{ label: "Batch Key", text: options.batchKey }] : []), + ...(options?.deliverAfter ? [{ label: "Deliver After", text: String(options.deliverAfter) }] : []), + ...(options?.deliverAt ? [{ label: "Deliver At", text: options.deliverAt.toISOString() }] : []), ], } ); diff --git a/packages/trigger-sdk/src/triggerClient.ts b/packages/trigger-sdk/src/triggerClient.ts index b1766f16429..dec83be8c2f 100644 --- a/packages/trigger-sdk/src/triggerClient.ts +++ b/packages/trigger-sdk/src/triggerClient.ts @@ -8,6 +8,7 @@ import { HandleTriggerSource, HttpSourceRequestHeadersSchema, HttpSourceResponseMetadata, + HttpSourceResponseOptions, IndexEndpointResponse, InitializeTriggerBodySchema, IntegrationConfig, @@ -108,6 +109,7 @@ export class TriggerClient { events: Array; response?: NormalizedResponse; metadata?: HttpSourceResponseMetadata; + options?: HttpSourceResponseOptions; } | void> > = {}; #registeredDynamicTriggers: Record< @@ -406,7 +408,7 @@ export class TriggerClient { metadata: inputMetadata, }; - const { response, events, metadata } = await this.#handleHttpSourceRequest( + const { response, events, metadata, options } = await this.#handleHttpSourceRequest( source, sourceRequest ); @@ -417,6 +419,7 @@ export class TriggerClient { events, response, metadata, + options, }, headers: this.#standardResponseHeaders, }; @@ -677,7 +680,7 @@ export class TriggerClient { async #preprocessRun(body: PreprocessRunBody, job: Job>, any>) { const context = this.#createPreprocessRunContext(body); - const parsedPayload = job.trigger.event.parsePayload(body.event.payload ?? {}); + const parsedPayload = job.trigger.event.parsePayload(body.payload ?? body.event.payload ?? {}); const properties = job.trigger.event.runProperties?.(parsedPayload) ?? []; @@ -740,7 +743,7 @@ export class TriggerClient { try { const output = await runLocalStorage.runWith({ io, ctx: context }, () => { return job.options.run( - job.trigger.event.parsePayload(body.event.payload ?? {}), + job.trigger.event.parsePayload(body.payload ?? body.event.payload ?? {}), ioWithConnections, context ); @@ -867,6 +870,7 @@ export class TriggerClient { response: NormalizedResponse; events: SendEvent[]; metadata?: HttpSourceResponseMetadata; + options?: HttpSourceResponseOptions; }> { this.#internalLogger.debug("Handling HTTP source request", { source, @@ -962,6 +966,7 @@ export class TriggerClient { }, }, metadata: results.metadata, + options: results.options, }; } diff --git a/packages/trigger-sdk/src/triggers/eventTrigger.ts b/packages/trigger-sdk/src/triggers/eventTrigger.ts index 6c2fc68b5cf..16eca32dfab 100644 --- a/packages/trigger-sdk/src/triggers/eventTrigger.ts +++ b/packages/trigger-sdk/src/triggers/eventTrigger.ts @@ -10,6 +10,7 @@ type EventTriggerOptions> = name?: string | string[]; source?: string; filter?: EventFilter; + batch?: boolean; }; export class EventTrigger> @@ -25,6 +26,7 @@ export class EventTrigger> return { type: "static", title: this.#options.name ?? this.#options.event.title, + batch: !!this.#options.batch, rule: { event: this.#options.name ?? this.#options.event.name, source: this.#options.source ?? "trigger.dev", @@ -74,6 +76,8 @@ type TriggerOptions = { * ``` */ filter?: EventFilter; + /** Receive events as an array of payloads. (Will only receive more than one event at a time if events are also sent with batching enabled.) */ + batch?: boolean; examples?: EventSpecificationExample[]; }; @@ -87,6 +91,7 @@ export function eventTrigger( return new EventTrigger({ name: options.name, filter: options.filter, + batch: options.batch, event: { name: options.name, title: "Event", diff --git a/packages/trigger-sdk/src/triggers/externalSource.ts b/packages/trigger-sdk/src/triggers/externalSource.ts index 01006564089..aa18e54410f 100644 --- a/packages/trigger-sdk/src/triggers/externalSource.ts +++ b/packages/trigger-sdk/src/triggers/externalSource.ts @@ -3,6 +3,7 @@ import { EventFilter, HandleTriggerSource, HttpSourceResponseMetadata, + HttpSourceResponseOptions, Logger, NormalizedResponse, RegisterTriggerSource, @@ -136,6 +137,7 @@ type HandlerFunction< events: SendEvent[]; response?: NormalizedResponse; metadata?: HttpSourceResponseMetadata; + options?: HttpSourceResponseOptions; } | void>; type KeyFunction = (params: TParams) => string; @@ -269,6 +271,7 @@ export type ExternalSourceTriggerOptions< source: TEventSource; params: ExternalSourceParams; options: TriggerOptionRecord; + batch?: boolean; }; export class ExternalSourceTrigger< @@ -286,6 +289,7 @@ export class ExternalSourceTrigger< return { type: "static", title: "External Source", + batch: !!this.options.batch, rule: { event: this.event.name, payload: deepMergeFilters( diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 3738f06f98a..a729dcff116 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -145,7 +145,7 @@ importers: eslint-config-prettier: ^8.5.0 express: ^4.18.1 framer-motion: ^10.12.11 - graphile-worker: ^0.13.0 + graphile-worker: 0.14.0-rc.0 highlight.run: ^7.3.4 humanize-duration: ^3.27.3 intl-parse-accept-language: ^1.0.0 @@ -246,7 +246,7 @@ importers: emails: link:../../packages/emails express: 4.18.2 framer-motion: 10.12.11_biqbaboplfbrettd7655fr4n2y - graphile-worker: 0.13.0 + graphile-worker: 0.14.0-rc.0 highlight.run: 7.3.4 humanize-duration: 3.27.3 intl-parse-accept-language: 1.0.0 @@ -2259,11 +2259,6 @@ packages: resolution: {integrity: sha512-mM4COjgZox8U+JcXQwPijIZLElkgEpO5rsERVDJTc2qfCDfERyob6k5WegS14SX18IIjv+XD+GrqNumY5JRCDw==} engines: {node: '>=6.9.0'} - /@babel/helper-validator-identifier/7.19.1: - resolution: {integrity: sha512-awrNfaMtnHUr653GgGEs++LlAvW6w+DcPrOliSMXWCKo597CwL5Acf/wWdNkf/tfEQE3mjkeD1YOVZOUV/od1w==} - engines: {node: '>=6.9.0'} - dev: true - /@babel/helper-validator-identifier/7.22.15: resolution: {integrity: sha512-4E/F9IIEi8WR94324mbDUMo074YTheJmd7eZF5vITTeYchqAi6sYXRLHUVsmkdmY4QjfKTcB2jB7dVP3NaBElQ==} engines: {node: '>=6.9.0'} @@ -5179,7 +5174,7 @@ packages: engines: {node: '>=6.9.0'} dependencies: '@babel/helper-string-parser': 7.21.5 - '@babel/helper-validator-identifier': 7.19.1 + '@babel/helper-validator-identifier': 7.22.15 to-fast-properties: 2.0.0 dev: true @@ -18795,7 +18790,7 @@ packages: eslint-import-resolver-webpack: optional: true dependencies: - '@typescript-eslint/parser': 5.59.6_eslint@8.42.0 + '@typescript-eslint/parser': 5.59.6_binxsscxvozjxebftqdoazsxm4 debug: 3.2.7 eslint: 8.42.0 eslint-import-resolver-node: 0.3.7 @@ -18880,7 +18875,7 @@ packages: '@typescript-eslint/parser': optional: true dependencies: - '@typescript-eslint/parser': 5.59.6_eslint@8.42.0 + '@typescript-eslint/parser': 5.59.6_binxsscxvozjxebftqdoazsxm4 array-includes: 3.1.6 array.prototype.flat: 1.3.1 array.prototype.flatmap: 1.3.1 @@ -21019,9 +21014,9 @@ packages: /graphemer/1.4.0: resolution: {integrity: sha512-EtKwoO6kxCL9WO5xipiHTZlSzBm7WLT627TqC/uVRd0HKmq8NXyebnNYxDoBi7wt8eTWrUrKXCOVaFq9x1kgag==} - /graphile-worker/0.13.0: - resolution: {integrity: sha512-8Hl5XV6hkabZRhYzvbUfvjJfPFR5EPxYRVWlzQC2rqYHrjULTLBgBYZna5R9ukbnsbWSvn4vVrzOBIOgIC1jjw==} - engines: {node: '>=10.0.0'} + /graphile-worker/0.14.0-rc.0: + resolution: {integrity: sha512-KV6nP1ljKlo/qROkB4RJS19OdnrM1yzWjDruf0BgLWx6gDTT+FECeHAFC+gxsvShm/rTyDB4OIPwqMMu/oiOqg==} + engines: {node: '>=14.0.0'} hasBin: true dependencies: '@graphile/logger': 0.2.0 @@ -21031,8 +21026,8 @@ packages: cosmiconfig: 7.1.0 json5: 2.2.3 pg: 8.10.0 - tslib: 2.4.1 - yargs: 16.2.0 + tslib: 2.6.2 + yargs: 17.7.2 transitivePeerDependencies: - pg-native dev: false @@ -32325,6 +32320,7 @@ packages: string-width: 4.2.3 y18n: 5.0.8 yargs-parser: 20.2.9 + dev: true /yargs/17.1.1: resolution: {integrity: sha512-c2k48R0PwKIqKhPMWjeiF6y2xY/gPMUlro0sgxqXpbOIohWiLNXWslsootttv7E1e73QPAMQSg5FeySbVcpsPQ==} diff --git a/references/job-catalog/src/airtable.ts b/references/job-catalog/src/airtable.ts index 703046bf3ee..78cc0b4007d 100644 --- a/references/job-catalog/src/airtable.ts +++ b/references/job-catalog/src/airtable.ts @@ -77,47 +77,46 @@ client.defineJob({ }, }); -//todo webhooks require batch support -// client.defineJob({ -// id: "airtable-delete-webhooks", -// name: "Airtable Example 2: webhook admin", -// version: "0.1.0", -// trigger: eventTrigger({ -// name: "airtable.example", -// schema: z.object({ -// baseId: z.string(), -// deleteWebhooks: z.boolean().optional(), -// }), -// }), -// integrations: { -// airtable, -// }, -// run: async (payload, io, ctx) => { -// const webhooks = await io.airtable.webhooks().list("list webhooks", { baseId: payload.baseId }); - -// if (payload.deleteWebhooks === true) { -// for (const webhook of webhooks.webhooks) { -// await io.airtable.webhooks().delete(`delete webhook: ${webhook.id}`, { -// baseId: payload.baseId, -// webhookId: webhook.id, -// }); -// } -// } -// }, -// }); +client.defineJob({ + id: "airtable-delete-webhooks", + name: "Airtable Example 2: webhook admin", + version: "0.1.0", + trigger: eventTrigger({ + name: "airtable.example", + schema: z.object({ + baseId: z.string(), + deleteWebhooks: z.boolean().optional(), + }), + }), + integrations: { + airtable, + }, + run: async (payload, io, ctx) => { + const webhooks = await io.airtable.webhooks().list("list webhooks", { baseId: payload.baseId }); + + if (payload.deleteWebhooks === true) { + for (const webhook of webhooks.webhooks) { + await io.airtable.webhooks().delete(`delete webhook: ${webhook.id}`, { + baseId: payload.baseId, + webhookId: webhook.id, + }); + } + } + }, +}); //todo changes the structure of the trigger so it's airtable.base("val").onChanges({ -// client.defineJob({ -// id: "airtable-on-table", -// name: "Airtable Example: onTable", -// version: "0.1.0", -// trigger: airtable.onTableChanges({ -// baseId: "appSX6ly4nZGfdUSy", -// tableId: "tblr5BReu2yeOMk7n", -// }), -// run: async (payload, io, ctx) => { -// await io.logger.log(`transaction number ${payload.baseTransactionNumber}`); -// }, -// }); +client.defineJob({ + id: "airtable-on-table", + name: "Airtable Example: onTable", + version: "0.1.0", + trigger: airtable.onTableChanges({ + baseId: "appSX6ly4nZGfdUSy", + tableId: "tblr5BReu2yeOMk7n", + }), + run: async (payload, io, ctx) => { + // await io.logger.log(`transaction number ${payload.baseTransactionNumber}`); + }, +}); createExpressServer(client); diff --git a/references/job-catalog/src/events.ts b/references/job-catalog/src/events.ts index 7bd9d9a0882..662da1db024 100644 --- a/references/job-catalog/src/events.ts +++ b/references/job-catalog/src/events.ts @@ -73,4 +73,51 @@ client.defineJob({ }, }); +client.defineJob({ + id: "send-batched-event", + name: "Send Batched Event", + version: "1.0.0", + trigger: eventTrigger({ + name: "batch.send", + }), + run: async (payload, io, ctx) => { + // unbatched + await io.sendEvent( + "unbatched", + { name: "new.message", payload: { message: "Not batching this!" } }, + { deliverAfter: 5 } + ); + + // batched + await io.sendEvent( + "batched-1", + { name: "new.message", payload: { message: "Message 1" } }, + { batchKey: "amazing-batchkey", deliverAfter: 10 } + ); + await io.sendEvent( + "batched-2", + { name: "new.message", payload: { message: "Message 2" } }, + { batchKey: "amazing-batchkey", deliverAfter: 200 } // deliverAfter: 10 - should be preserved + ); + await io.sendEvent( + "batched-3", + { name: "new.message", payload: { message: "Message 3" } }, + { batchKey: "amazing-batchkey", deliverAfter: 300 } // deliverAfter: 10 - should be preserved + ); + }, +}); + +client.defineJob({ + id: "receive-batched-event", + name: "Receive Batched Event", + version: "1.0.0", + trigger: eventTrigger({ + name: "new.message", + schema: z.object({ message: z.string() }), + }), + run: async (payload, io, ctx) => { + return payload; + }, +}); + createExpressServer(client);