From e7f08139a9e03c212bbd464c2b3a7d433e09d137 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 22 Dec 2022 11:24:36 +0000 Subject: [PATCH] Add the CONNECTED_SLOT_CONNECTED internal event, to register a webhook after a connection slot gets a connection --- .../app/services/messageBroker.server.ts | 30 +++++++++++++++--- .../webhooks/registerWebhook.server.ts | 31 +++++++++++++++++-- examples/github-webhook/src/index.ts | 2 +- 3 files changed, 56 insertions(+), 7 deletions(-) diff --git a/apps/webapp/app/services/messageBroker.server.ts b/apps/webapp/app/services/messageBroker.server.ts index 5093300a2..a74920900 100644 --- a/apps/webapp/app/services/messageBroker.server.ts +++ b/apps/webapp/app/services/messageBroker.server.ts @@ -13,6 +13,7 @@ import { z } from "zod"; import { prisma } from "~/db.server"; import { env } from "~/env.server"; import { findRegisteredWebhookById } from "~/models/registeredWebhook.server"; +import { findWorkflowConnectionSlotById } from "~/models/workflowConnectionSlot.server"; import { completeWorkflowRun, failWorkflowRun, @@ -184,6 +185,10 @@ const InternalCatalog = { data: z.object({ id: z.string() }), properties: z.object({}), }, + CONNECTION_SLOT_CONNECTED: { + data: z.object({ id: z.string() }), + properties: z.object({}), + }, }; async function createInternalPubSub() { @@ -199,16 +204,33 @@ async function createInternalPubSub() { }, schema: InternalCatalog, handlers: { - REGISTERED_WEBHOOK_CREATED: async (id, data, properties) => { - const webhook = await findRegisteredWebhookById(data.id); + CONNECTION_SLOT_CONNECTED: async (id, data, properties) => { + const slot = await findWorkflowConnectionSlotById(data.id); - if (!webhook) { + if (!slot) { + return true; + } + + if (!slot.connection) { + return true; + } + + if (!slot.registeredWebhook) { return true; } const registerWebhookService = new RegisterWebhook(); - const isRegistered = await registerWebhookService.call(webhook); + const isRegistered = await registerWebhookService.call( + slot.registeredWebhook + ); + + return isRegistered; // Returning true will mean we don't retry + }, + REGISTERED_WEBHOOK_CREATED: async (id, data, properties) => { + const registerWebhookService = new RegisterWebhook(); + + const isRegistered = await registerWebhookService.call(data.id); return isRegistered; // Returning true will mean we don't retry }, diff --git a/apps/webapp/app/services/webhooks/registerWebhook.server.ts b/apps/webapp/app/services/webhooks/registerWebhook.server.ts index c4aecdada..98a5a31b9 100644 --- a/apps/webapp/app/services/webhooks/registerWebhook.server.ts +++ b/apps/webapp/app/services/webhooks/registerWebhook.server.ts @@ -1,7 +1,9 @@ +import type { RegisteredWebhook } from ".prisma/client"; +import { github } from "internal-integrations"; import type { PrismaClient } from "~/db.server"; import { prisma } from "~/db.server"; import type { RegisteredWebhookWithRelationships } from "~/models/registeredWebhook.server"; -import { github } from "internal-integrations"; +import { findRegisteredWebhookById } from "~/models/registeredWebhook.server"; import { pizzly } from "../pizzly.server"; import { originOrProxyUrl } from "../webhookProxy.server"; @@ -12,7 +14,13 @@ export class RegisterWebhook { this.#prismaClient = prismaClient; } - public async call(webhook: RegisteredWebhookWithRelationships) { + public async call(idOrWebhook: string | RegisteredWebhook) { + const webhook = await this.#findWebhook(idOrWebhook); + + if (!webhook) { + return true; + } + if (webhook.status === "CONNECTED") { return true; } @@ -56,6 +64,25 @@ export class RegisterWebhook { return true; } + async #findWebhook( + idOrWebhook: string | RegisteredWebhook + ): Promise { + const webhook = + typeof idOrWebhook === "string" + ? await findRegisteredWebhookById(idOrWebhook) + : await findRegisteredWebhookById(idOrWebhook.id); + + if (!webhook) { + return; + } + + if (!webhook.connectionSlot.connection) { + return; + } + + return webhook; + } + async #registerWebhookWithConnection( serviceIdentifier: string, accessToken: string, diff --git a/examples/github-webhook/src/index.ts b/examples/github-webhook/src/index.ts index 65776224a..2955b1566 100644 --- a/examples/github-webhook/src/index.ts +++ b/examples/github-webhook/src/index.ts @@ -2,7 +2,7 @@ import { Trigger } from "@trigger.dev/sdk"; import { github } from "@trigger.dev/integrations"; const trigger = new Trigger({ - id: "github-webhook-7", + id: "github-webhook-8", name: "GitHub Issue changes for jsonhero-web", apiKey: "trigger_dev_zC25mKNn6c0q", endpoint: "ws://localhost:8889/ws",