ecd050bece
* 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>
177 lines
4.9 KiB
TypeScript
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)],
|
|
};
|
|
}
|