From 2cce68b5f76e51909c2f9f4912e6671b6a8357d7 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Fri, 24 Nov 2023 09:30:48 +0000 Subject: [PATCH] =?UTF-8?q?Revert=20"Remove=20the=20pg=20listen=20code=20t?= =?UTF-8?q?o=20see=20if=20it=E2=80=99s=20causing=20DB=20issues"=20(#749)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This reverts commit e8e7c116d192095bd2e0d61a196538fd5a5d9039. --- apps/webapp/app/platform/zodWorker.server.ts | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/apps/webapp/app/platform/zodWorker.server.ts b/apps/webapp/app/platform/zodWorker.server.ts index f917645d3..347dd5981 100644 --- a/apps/webapp/app/platform/zodWorker.server.ts +++ b/apps/webapp/app/platform/zodWorker.server.ts @@ -14,6 +14,7 @@ import { run as graphileRun, parseCronItems } from "graphile-worker"; import omit from "lodash.omit"; import { z } from "zod"; import { PrismaClient, PrismaClientOrTransaction } from "~/db.server"; +import { PgListenService } from "~/services/db/pgListen.server"; import { workerLogger as logger, trace } from "~/services/logger.server"; export interface MessageCatalogSchema { @@ -166,6 +167,21 @@ export class ZodWorker { 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.on("trigger:graphile:migrate", async ({ latestMigration }) => { + 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 }) => {