Files
Matt Aitken ecd050bece Airtable integration with webhook and integration changes (#399)
* CLI create-integration command now accepts an Open AI api key

* Create integration docs separated into multiple pages

* Initial Airtable integration commit, with OpenAI generated code

* OAuth page coming soon

* Export DisplayProperty from the SDK

* TSConfig made to match GitHub’s with paths

* getRecords

* Removed duplicate Stripe job from the catalog

* Renamed Airtable apiKey option to token

* First Airtable job

* Export Collaborator and Attachment field types

* A typesafe example that uses runTask

* WIP on new integration tasks… not working yet

* Attempt with class

* Revert "Attempt with class"

This reverts commit 93a48330019f754c3216c5b49964fa4b0218bd3f.

* WIP changing how tasks work

* Mock of async local storage

* Moved client creation from constructor

* New approach with a clone method on TriggerIntegration

* Added runTask to Airtable which is used by integration tasks

* Added the Airtable icon and connection when using runTask

* base().table() is working

* runTask options moved to the 3rd param, made optional with optional name

* Added some generic arguments

* Added generic type to table

* Removed old comment

* We don’t need to repeat the icon

* getRecords and getRecord now returning the right data and types

* Creating records

* Update records

* Delete records

* The internal properties of integrations are now hidden by the TypeScript types

* Sprinkled a Prettify in there

* Improved the types

* Added Airtable to the integration catalog

* Early work on Airtable webhook registration

* More progress with webhooks

* connectionKey needs to be cloned for webhooks to work

* connectionKey needs to be cloned for webhooks to work

* It was unclear that the ActivateSourceService was using a graphileJob id

* ActivateSourceService optionally takes a jobId, if missing it generate a unique id

* When retrying trigger registration, don’t pass an id so it is generated

* Removed Airtable webhooks tasks from the job-catalog example

* Added TriggerSourceOption, removed TriggerSourceEvent

* WIP with new ExternalSource options

* ExternalSourceTrigger setup

* DynamicTrigger changed to options, will need some more work

* filter gets options passed to it

* SourceMetadata v2 renamed to SourceMetadataV2, kept original

* Started versioning the backend

* Moved param order on io.getEvent and io.cancelEvent

* The runTask stuff that allows unknown to work is back

* Indexing for v1 and v2, with version on “activateSource” schema

* Added todos, to deal with Airtable SDK calls inside the webhook handler

* “deliverHttpSourceRequest” queueName changed to the source id so they process in order

* ActivateSource changes to deal with old and new data formats

* Update existing TriggerSources to v2

* Fix for dynamic.ts typescript errors, need to revisit this later

* UpdateSourceService v1 and v2, with new v2 API endpoint

* Removed unused imports

* More progress on v1 and v2

* Airtable webhooks are now triggering a job

* Moved webhooks to a new file

* You can do API calls in the webhook handler now, Airtable webhook data is being processed

* Airtable events coming through

* Defined the Airtable table payload type

* TriggerSource metadata is being stored and used

* Removed some logs

* Added filtering and don’t allow any webhooks that use automated sources

* Resend switched to new integration

* Moved Resend test jobs to the catalog, and tested it worked

* WIP on Slack, there are compile errors

* Created a generic type that strips out indexes

* Slack updated to use new integration

* SendGrid migrated over

* Integration runTask is now allowing regular types

* Changed io.runTask types so it only allows Json-able types

* OpenAI models tasks working

* Added Airtable changes to runTask

* Removed the index signature crap from the Slack integration

* Don’t need to cast the callback result

* Updated Resend

* Re-ordered runTask params

* WIP on openai

* onAccountUpdated is Connect only

* Removed RunTaskResult

* Handle Resend errors, the official SDK doesn’t expose them properly at the moment

* Removed OmitIndexSignature

* OpenAI converted to new integration, with backwards compatible functions

* Put the openai catalog back to what it was originally

* Export a standard retry with backoff, to be used

* Use the standard exponential backoff in the integrations

* Retry options moved earlier so they can be overriden by a task

* GitHub tasks migrated

* Added sources, fixed one bundling issue

* Added GitHub jobs to catalog

* Remove duplicate options

* Deduplicate events

* Removed duplicate Job

* Switched Plain over

* Set the Plain icon

* Converted Stripe over

* Supabase adapted

* Typeform working

* Added dynamic-schedule to catalog

* Added background-fetch job catalog

* Created dynamic-triggers catalog file

* Fixed old general file with runTask param order

* Dynamic triggers working

* SendGrid updated to use the same tsconfig as other integrations

* Removed Airtable webhook, until we have batch support

* Added OAuth airtable auth example

* Created beta changeset tag

* Beta changesets for most packages

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2023-09-05 14:21:42 +01:00

177 lines
4.9 KiB
TypeScript

import {
REGISTER_SOURCE_EVENT_V2,
REGISTER_SOURCE_EVENT_V1,
RegisterTriggerSource,
RegisterSourceEventV1,
RegisterSourceEventV2,
RegisterSourceEventOptions,
RegisteredOptionsDiff,
} from "@trigger.dev/core";
import type { SecretReference, TriggerSource, TriggerSourceOption } from "@trigger.dev/database";
import { z } from "zod";
import type { PrismaClient } from "~/db.server";
import { prisma } from "~/db.server";
import { env } from "~/env.server";
import type { AuthenticatedEnvironment } from "../apiAuth.server";
import { IngestSendEvent } from "../events/ingestSendEvent.server";
import { getSecretStore } from "../secrets/secretStore.server";
import { nanoid } from "nanoid";
export class ActivateSourceService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(id: string, jobId?: string, orphanedOptions?: Record<string, string[]>) {
const triggerSource = await this.#prismaClient.triggerSource.findUniqueOrThrow({
where: {
id,
},
include: {
endpoint: true,
environment: {
include: {
organization: true,
project: true,
},
},
options: true,
secretReference: true,
},
});
const eventId = `${id}:${jobId ?? nanoid()}`;
// TODO: support more channels
switch (triggerSource.channel) {
case "HTTP": {
await this.#activateHttpSource(
triggerSource.environment,
triggerSource,
triggerSource.options,
triggerSource.secretReference,
eventId,
orphanedOptions
);
}
}
}
async #activateHttpSource(
environment: AuthenticatedEnvironment,
triggerSource: TriggerSource,
options: Array<TriggerSourceOption>,
secretReference: SecretReference,
eventId: string,
orphanedOptions?: Record<string, string[]>
) {
const secretStore = getSecretStore(secretReference.provider);
const httpSecret = await secretStore.getSecret(
z.object({
secret: z.string(),
}),
secretReference.key
);
if (!httpSecret) {
throw new Error("HTTP Secret not found");
}
const service = new IngestSendEvent();
const source: RegisterTriggerSource = {
key: triggerSource.key,
active: triggerSource.active,
secret: httpSecret.secret,
data: triggerSource.channelData as any,
channel: {
type: "HTTP",
url: `${env.APP_ORIGIN}/api/v1/sources/http/${triggerSource.id}`,
},
};
switch (triggerSource.version) {
case "1": {
const events = triggerSource.active
? options.filter((e) => e.registered).map((e) => e.value)
: options.map((e) => e.value);
const missingEvents = triggerSource.active
? options.filter((e) => !e.registered).map((e) => e.value)
: [];
const orphanedEvents = orphanedOptions
? Object.values(orphanedOptions).flatMap((vals) => vals)
: [];
const payload: RegisterSourceEventV1 = {
id: triggerSource.id,
source,
events,
missingEvents,
orphanedEvents,
};
await service.call(environment, {
id: eventId,
name: REGISTER_SOURCE_EVENT_V1,
source: "trigger.dev",
payload,
});
break;
}
case "2": {
//group the options by the name
const optionsRecord = options.reduce((acc, option) => {
if (!acc[option.name]) {
acc[option.name] = [];
}
acc[option.name].push(option);
return acc;
}, {} as Record<string, Array<TriggerSourceOption>>);
//for each of the optionsRecord, create the diff
const payloadOptions = Object.entries(optionsRecord).reduce(
(acc, [key, value]) => ({
...acc,
[key]: getOptionsDiff(triggerSource.active, orphanedOptions?.[key] ?? [], value),
}),
{} as Record<string, RegisteredOptionsDiff>
) as RegisterSourceEventOptions;
const payload: RegisterSourceEventV2 = {
id: triggerSource.id,
source,
options: payloadOptions,
};
await service.call(environment, {
id: eventId,
name: REGISTER_SOURCE_EVENT_V2,
source: "trigger.dev",
payload,
});
}
}
}
}
function getOptionsDiff(
sourceIsActive: boolean,
orphaned: string[],
options: Array<TriggerSourceOption>
): RegisteredOptionsDiff {
const desired = sourceIsActive
? options.filter((e) => e.registered).map((e) => e.value)
: options.map((e) => e.value);
const missing = sourceIsActive ? options.filter((e) => !e.registered).map((e) => e.value) : [];
return {
desired: [...new Set(desired)],
missing: [...new Set(missing)],
orphaned: [...new Set(orphaned)],
};
}