From 6602a4e4988837b18433b89a114e97a3b2b76de4 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 26 Oct 2023 14:37:43 +0000 Subject: [PATCH 1/4] Add notification system --- .../webapp/app/services/db/pgListen.server.ts | 48 +++++++++++++++++++ .../webapp/app/services/db/pgNotify.server.ts | 23 +++++++++ apps/webapp/package.json | 1 + pnpm-lock.yaml | 2 + 4 files changed, 74 insertions(+) create mode 100644 apps/webapp/app/services/db/pgListen.server.ts create mode 100644 apps/webapp/app/services/db/pgNotify.server.ts diff --git a/apps/webapp/app/services/db/pgListen.server.ts b/apps/webapp/app/services/db/pgListen.server.ts new file mode 100644 index 00000000000..74ecc26e775 --- /dev/null +++ b/apps/webapp/app/services/db/pgListen.server.ts @@ -0,0 +1,48 @@ +import { logger } from "~/services/logger.server"; +import { Logger } from "@trigger.dev/core"; +import type { PoolClient } from "pg"; + +export class PgListenService { + #poolClient: PoolClient; + #logger: Logger; + #loggerNamespace: string; + + constructor(poolClient: PoolClient, loggerNamespace?: string, loggerInstance?: Logger) { + this.#poolClient = poolClient; + this.#logger = loggerInstance ?? logger; + this.#loggerNamespace = loggerNamespace ?? ""; + } + + public async call(channelName: string, callback: (payload: string) => Promise) { + this.#logDebug("Registering notification handler", { channelName }); + + const isValidChannel = channelName.match(/^[a-zA-Z0-9:-_]+$/); + + if (!isValidChannel) { + throw new Error(`Invalid channel name: ${channelName}`); + } + + this.#poolClient.query(`LISTEN "${channelName}"`).then(null, (error) => { + this.#logDebug("LISTEN error", error); + }); + + this.#poolClient.on("notification", async (notification) => { + if (notification.channel !== channelName) { + return; + } + + this.#logDebug("Notification received", { notification }); + + if (!notification.payload) { + return; + } + + await callback(notification.payload); + }); + } + + #logDebug(message: string, args?: any) { + const namespace = this.#loggerNamespace ? `[${this.#loggerNamespace}]` : ""; + this.#logger.debug(`[pgListen]${namespace} ${message}`, args); + } +} diff --git a/apps/webapp/app/services/db/pgNotify.server.ts b/apps/webapp/app/services/db/pgNotify.server.ts new file mode 100644 index 00000000000..2b41e09adf1 --- /dev/null +++ b/apps/webapp/app/services/db/pgNotify.server.ts @@ -0,0 +1,23 @@ +import { SerializableJson } from "@trigger.dev/core"; +import { PrismaClient, prisma } from "~/db.server"; +import { logger } from "~/services/logger.server"; + +export class PgNotifyService { + #prismaClient: PrismaClient; + + constructor(prismaClient: PrismaClient = prisma) { + this.#prismaClient = prismaClient; + } + + public async call(channelName: string, payload: SerializableJson) { + this.#logDebug("Sending notification", { channelName, notifyPayload: payload }); + + await this.#prismaClient.$executeRaw` + SELECT pg_notify(${channelName}, ${JSON.stringify(payload)}) + `; + } + + #logDebug(message: string, args?: any) { + logger.debug(`[pgNotify] ${message}`, args); + } +} diff --git a/apps/webapp/package.json b/apps/webapp/package.json index 54a59d4cdd8..5cbb0ec2d50 100644 --- a/apps/webapp/package.json +++ b/apps/webapp/package.json @@ -65,6 +65,7 @@ "@trigger.dev/core": "workspace:*", "@trigger.dev/database": "workspace:*", "@trigger.dev/sdk": "workspace:*", + "@types/pg": "8.6.6", "@uiw/react-codemirror": "^4.19.5", "class-variance-authority": "^0.5.2", "clsx": "^1.2.1", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 3f9c5b43ba1..db16395d5ad 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -121,6 +121,7 @@ importers: '@types/morgan': ^1.9.3 '@types/node': ^18.11.15 '@types/node-fetch': ^2.6.2 + '@types/pg': 8.6.6 '@types/prismjs': ^1.26.0 '@types/qs': ^6.9.7 '@types/react': 18.2.17 @@ -236,6 +237,7 @@ importers: '@trigger.dev/core': link:../../packages/core '@trigger.dev/database': link:../../packages/database '@trigger.dev/sdk': link:../../packages/trigger-sdk + '@types/pg': 8.6.6 '@uiw/react-codemirror': 4.19.5_th22fcplkuhrqjnlojwclcaim4 class-variance-authority: 0.5.2_typescript@5.2.2 clsx: 1.2.1 From 24c5bf1c47dfd94ba2f10eeadc0ec426f5315411 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 26 Oct 2023 15:00:36 +0000 Subject: [PATCH 2/4] Graceful worker shutdown on migration notification --- apps/webapp/app/platform/zodWorker.server.ts | 27 +++++++++++++++++++- apps/webapp/app/services/worker.server.ts | 5 ++++ 2 files changed, 31 insertions(+), 1 deletion(-) diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts index d4480286df0..236b386786f 100644 --- a/apps/webapp/app/platform/zodWorker.server.ts +++ b/apps/webapp/app/platform/zodWorker.server.ts @@ -15,6 +15,8 @@ import omit from "lodash.omit"; import { z } from "zod"; import { PrismaClient, PrismaClientOrTransaction } from "~/db.server"; import { workerLogger as logger } from "~/services/logger.server"; +import { PgListenService } from "~/services/db/pgListen.server"; +import { safeJsonParse } from "~/utils/json"; export interface MessageCatalogSchema { [key: string]: z.ZodFirstPartySchemaTypes | z.ZodDiscriminatedUnion; @@ -164,8 +166,31 @@ export class ZodWorker { this.#logDebug("pool:create", { attempts }); }); - this.#runner?.events.on("pool:listen:success", ({ workerPool, client }) => { + this.#runner?.events.on("pool:listen:success", async ({ workerPool, client }) => { this.#logDebug("pool:listen:success"); + + // hijack client instance to listen and react to incoming NOTIFY events + const pgListen = new PgListenService(client, this.#name, logger); + + await pgListen.call("trigger:graphile:migrate", async (payload) => { + const parsedPayload = safeJsonParse(payload); + + const MigrationNotificationPayloadSchema = z.object({ + latestMigration: z.number(), + }); + + const { latestMigration } = MigrationNotificationPayloadSchema.parse(parsedPayload); + + this.#logDebug("Detected incoming migration", { latestMigration }); + + if (latestMigration > 10) { + // already migrated past v0.14 - nothing to do + return; + } + + // simulate SIGTERM to trigger graceful shutdown + this._handleSignal("SIGTERM"); + }); }); this.#runner?.events.on("pool:listen:error", ({ error }) => { diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts index 15f2534a27f..dc77abc7654 100644 --- a/apps/webapp/app/services/worker.server.ts +++ b/apps/webapp/app/services/worker.server.ts @@ -21,6 +21,7 @@ import { DeliverHttpSourceRequestService } from "./sources/deliverHttpSourceRequ import { PerformTaskOperationService } from "./tasks/performTaskOperation.server"; import { ProcessCallbackTimeoutService } from "./tasks/processCallbackTimeout"; import { ProbeEndpointService } from "./endpoints/probeEndpoint.server"; +import { PgNotifyService } from "./db/pgNotify.server"; const workerCatalog = { indexEndpoint: z.object({ @@ -121,6 +122,10 @@ if (env.NODE_ENV === "production") { } export async function init() { + // const pgNotify = new PgNotifyService(); + // await pgNotify.call("trigger:graphile:migrate", { latestMigration: 10 }); + // await new Promise((resolve) => setTimeout(resolve, 10000)) + if (env.WORKER_ENABLED === "true") { await workerQueue.initialize(); } From b159a5890c55c63881c5babddd87ac5b2ee55f92 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 26 Oct 2023 20:18:49 +0000 Subject: [PATCH 3/4] Add notification catalog --- apps/webapp/app/platform/zodWorker.server.ts | 10 +------ .../webapp/app/services/db/pgListen.server.ts | 26 ++++++++++++++++--- .../webapp/app/services/db/pgNotify.server.ts | 8 ++++-- apps/webapp/app/services/db/types.ts | 11 ++++++++ 4 files changed, 40 insertions(+), 15 deletions(-) create mode 100644 apps/webapp/app/services/db/types.ts diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts index 236b386786f..496dae542d4 100644 --- a/apps/webapp/app/platform/zodWorker.server.ts +++ b/apps/webapp/app/platform/zodWorker.server.ts @@ -172,15 +172,7 @@ export class ZodWorker { // hijack client instance to listen and react to incoming NOTIFY events const pgListen = new PgListenService(client, this.#name, logger); - await pgListen.call("trigger:graphile:migrate", async (payload) => { - const parsedPayload = safeJsonParse(payload); - - const MigrationNotificationPayloadSchema = z.object({ - latestMigration: z.number(), - }); - - const { latestMigration } = MigrationNotificationPayloadSchema.parse(parsedPayload); - + await pgListen.call("trigger:graphile:migrate", async ({ latestMigration }) => { this.#logDebug("Detected incoming migration", { latestMigration }); if (latestMigration > 10) { diff --git a/apps/webapp/app/services/db/pgListen.server.ts b/apps/webapp/app/services/db/pgListen.server.ts index 74ecc26e775..104209d18d7 100644 --- a/apps/webapp/app/services/db/pgListen.server.ts +++ b/apps/webapp/app/services/db/pgListen.server.ts @@ -1,6 +1,9 @@ -import { logger } from "~/services/logger.server"; -import { Logger } from "@trigger.dev/core"; import type { PoolClient } from "pg"; +import { z } from "zod"; +import { Logger } from "@trigger.dev/core"; +import { logger } from "~/services/logger.server"; +import { NotificationCatalog, NotificationChannel, notificationCatalog } from "./types"; +import { safeJsonParse } from "~/utils/json"; export class PgListenService { #poolClient: PoolClient; @@ -13,7 +16,10 @@ export class PgListenService { this.#loggerNamespace = loggerNamespace ?? ""; } - public async call(channelName: string, callback: (payload: string) => Promise) { + public async call( + channelName: TChannel, + callback: (payload: z.infer) => Promise + ) { this.#logDebug("Registering notification handler", { channelName }); const isValidChannel = channelName.match(/^[a-zA-Z0-9:-_]+$/); @@ -37,7 +43,19 @@ export class PgListenService { return; } - await callback(notification.payload); + const payload = safeJsonParse(notification.payload); + + const parsedPayload = notificationCatalog[channelName].safeParse(payload); + + if (!parsedPayload.success) { + throw new Error( + `Failed to parse notification payload: ${channelName} - ${JSON.stringify( + parsedPayload.error + )}` + ); + } + + await callback(parsedPayload.data); }); } diff --git a/apps/webapp/app/services/db/pgNotify.server.ts b/apps/webapp/app/services/db/pgNotify.server.ts index 2b41e09adf1..2a80294ebd9 100644 --- a/apps/webapp/app/services/db/pgNotify.server.ts +++ b/apps/webapp/app/services/db/pgNotify.server.ts @@ -1,6 +1,7 @@ -import { SerializableJson } from "@trigger.dev/core"; +import { z } from "zod"; import { PrismaClient, prisma } from "~/db.server"; import { logger } from "~/services/logger.server"; +import { NotificationCatalog, NotificationChannel } from "./types"; export class PgNotifyService { #prismaClient: PrismaClient; @@ -9,7 +10,10 @@ export class PgNotifyService { this.#prismaClient = prismaClient; } - public async call(channelName: string, payload: SerializableJson) { + public async call( + channelName: TChannel, + payload: z.infer + ) { this.#logDebug("Sending notification", { channelName, notifyPayload: payload }); await this.#prismaClient.$executeRaw` diff --git a/apps/webapp/app/services/db/types.ts b/apps/webapp/app/services/db/types.ts new file mode 100644 index 00000000000..2e5b6d222fd --- /dev/null +++ b/apps/webapp/app/services/db/types.ts @@ -0,0 +1,11 @@ +import { z } from "zod"; + +export const notificationCatalog = { + "trigger:graphile:migrate": z.object({ + latestMigration: z.number(), + }), +}; + +export type NotificationCatalog = typeof notificationCatalog; + +export type NotificationChannel = keyof NotificationCatalog; From 590f9c5a7f93461e2e9b52120f8d281435c48fef Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Mon, 30 Oct 2023 17:05:44 +0000 Subject: [PATCH 4/4] Rename pgListen call to on --- apps/webapp/app/services/db/pgListen.server.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/webapp/app/services/db/pgListen.server.ts b/apps/webapp/app/services/db/pgListen.server.ts index 104209d18d7..94474f69f04 100644 --- a/apps/webapp/app/services/db/pgListen.server.ts +++ b/apps/webapp/app/services/db/pgListen.server.ts @@ -16,7 +16,7 @@ export class PgListenService { this.#loggerNamespace = loggerNamespace ?? ""; } - public async call( + public async on( channelName: TChannel, callback: (payload: z.infer) => Promise ) {