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>
125 lines
3.4 KiB
TypeScript
125 lines
3.4 KiB
TypeScript
import { z } from "zod";
|
|
import type { PrismaClient } from "~/db.server";
|
|
import { prisma } from "~/db.server";
|
|
import { EndpointApi } from "../endpointApi.server";
|
|
import { IngestSendEvent } from "../events/ingestSendEvent.server";
|
|
import { getSecretStore } from "../secrets/secretStore.server";
|
|
import { resolveApiConnection, resolveRunConnection } from "~/models/runConnection.server";
|
|
import { ConnectionAuth } from "@trigger.dev/sdk";
|
|
import { resolveSourceConnection } from "~/models/sourceConnection.server";
|
|
|
|
export class DeliverHttpSourceRequestService {
|
|
#prismaClient: PrismaClient;
|
|
|
|
constructor(prismaClient: PrismaClient = prisma) {
|
|
this.#prismaClient = prismaClient;
|
|
}
|
|
|
|
public async call(id: string) {
|
|
const httpSourceRequest = await this.#prismaClient.httpSourceRequestDelivery.findUniqueOrThrow({
|
|
where: { id },
|
|
include: {
|
|
endpoint: true,
|
|
environment: {
|
|
include: {
|
|
organization: true,
|
|
project: true,
|
|
},
|
|
},
|
|
source: {
|
|
include: {
|
|
secretReference: true,
|
|
dynamicTrigger: true,
|
|
externalAccount: true,
|
|
integration: {
|
|
include: {
|
|
connections: true,
|
|
},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
});
|
|
|
|
if (!httpSourceRequest.source.active) {
|
|
return;
|
|
}
|
|
|
|
const secretStore = getSecretStore(httpSourceRequest.source.secretReference.provider);
|
|
|
|
const secret = await secretStore.getSecret(
|
|
z.object({
|
|
secret: z.string(),
|
|
}),
|
|
httpSourceRequest.source.secretReference.key
|
|
);
|
|
|
|
if (!secret) {
|
|
throw new Error(`Secret not found for ${httpSourceRequest.source.key}`);
|
|
}
|
|
|
|
const auth = await resolveSourceConnection(this.#prismaClient, httpSourceRequest.source);
|
|
|
|
const clientApi = new EndpointApi(
|
|
httpSourceRequest.environment.apiKey,
|
|
httpSourceRequest.endpoint.url
|
|
);
|
|
|
|
const { response, events, metadata } = await clientApi.deliverHttpSourceRequest({
|
|
key: httpSourceRequest.source.key,
|
|
dynamicId: httpSourceRequest.source.dynamicTrigger?.slug,
|
|
secret: secret.secret,
|
|
params: httpSourceRequest.source.params,
|
|
data: httpSourceRequest.source.channelData,
|
|
request: {
|
|
url: httpSourceRequest.url,
|
|
method: httpSourceRequest.method,
|
|
headers: httpSourceRequest.headers as Record<string, string>,
|
|
rawBody: httpSourceRequest.body,
|
|
},
|
|
auth,
|
|
metadata: httpSourceRequest.source.metadata,
|
|
});
|
|
|
|
await this.#prismaClient.httpSourceRequestDelivery.update({
|
|
where: {
|
|
id,
|
|
},
|
|
data: {
|
|
deliveredAt: new Date(),
|
|
},
|
|
});
|
|
|
|
if (metadata) {
|
|
await this.#prismaClient.triggerSource.update({
|
|
where: {
|
|
id: httpSourceRequest.source.id,
|
|
},
|
|
data: {
|
|
metadata: metadata,
|
|
},
|
|
});
|
|
}
|
|
|
|
const ingestService = new IngestSendEvent();
|
|
|
|
for (const event of events) {
|
|
await ingestService.call(
|
|
httpSourceRequest.environment,
|
|
event,
|
|
{
|
|
accountId: httpSourceRequest.source.externalAccount?.identifier,
|
|
},
|
|
httpSourceRequest.source.dynamicSourceId
|
|
? {
|
|
id: httpSourceRequest.source.dynamicSourceId,
|
|
metadata: httpSourceRequest.source.dynamicSourceMetadata,
|
|
}
|
|
: undefined
|
|
);
|
|
}
|
|
|
|
return response;
|
|
}
|
|
}
|