From 74db2de1bc8567dfa2d96d5f87fe9dd09fbd2467 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 5 Jun 2024 19:08:00 +0100 Subject: [PATCH] Use graphile strategy 0 (no named queues) and remove all named queues --- .../sources/handleHttpSource.server.ts | 1 - apps/webapp/app/services/worker.server.ts | 32 ++++++++----------- .../alerts/performDeploymentAlerts.server.ts | 4 +-- .../alerts/performTaskAttemptAlerts.server.ts | 4 +-- .../services/bulk/performBulkAction.server.ts | 1 - .../app/v3/services/resumeBatchRun.server.ts | 2 +- patches/graphile-worker@0.16.6.patch | 11 ++++++- pnpm-lock.yaml | 6 ++-- 8 files changed, 30 insertions(+), 31 deletions(-) diff --git a/apps/webapp/app/services/sources/handleHttpSource.server.ts b/apps/webapp/app/services/sources/handleHttpSource.server.ts index ed96c90b2..2e39e4bd8 100644 --- a/apps/webapp/app/services/sources/handleHttpSource.server.ts +++ b/apps/webapp/app/services/sources/handleHttpSource.server.ts @@ -61,7 +61,6 @@ export class HandleHttpSourceService { id: delivery.id, }, { - queueName: `deliver:${triggerSource.id}`, tx, maxAttempts: triggerSource.environment.type === RuntimeEnvironmentType.DEVELOPMENT ? 1 : undefined, diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts index dc313451e..b50336b54 100644 --- a/apps/webapp/app/services/worker.server.ts +++ b/apps/webapp/app/services/worker.server.ts @@ -4,7 +4,21 @@ import { z } from "zod"; import { prisma } from "~/db.server"; import { env } from "~/env.server"; import { ZodWorker } from "~/platform/zodWorker.server"; +import { eventRepository } from "~/v3/eventRepository.server"; +import { RequeueTaskRunService } from "~/v3/requeueTaskRun.server"; +import { DeliverAlertService } from "~/v3/services/alerts/deliverAlert.server"; +import { PerformDeploymentAlertsService } from "~/v3/services/alerts/performDeploymentAlerts.server"; +import { PerformTaskAttemptAlertsService } from "~/v3/services/alerts/performTaskAttemptAlerts.server"; +import { PerformBulkActionService } from "~/v3/services/bulk/performBulkAction.server"; +import { ExecuteTasksWaitingForDeployService } from "~/v3/services/executeTasksWaitingForDeploy"; import { IndexDeploymentService } from "~/v3/services/indexDeployment.server"; +import { ResumeBatchRunService } from "~/v3/services/resumeBatchRun.server"; +import { ResumeTaskDependencyService } from "~/v3/services/resumeTaskDependency.server"; +import { ResumeTaskRunDependenciesService } from "~/v3/services/resumeTaskRunDependencies.server"; +import { RetryAttemptService } from "~/v3/services/retryAttempt.server"; +import { TimeoutDeploymentService } from "~/v3/services/timeoutDeployment.server"; +import { TriggerScheduledTaskService } from "~/v3/services/triggerScheduledTask.server"; +import { GraphileMigrationHelperService } from "./db/graphileMigrationHelper.server"; import { ExpireDispatcherService } from "./dispatchers/expireDispatcher.server"; import { InvokeEphemeralDispatcherService } from "./dispatchers/invokeEphemeralEventDispatcher.server"; import { sendEmail } from "./email.server"; @@ -30,22 +44,6 @@ import { DeliverWebhookRequestService } from "./sources/deliverWebhookRequest.se import { PerformTaskOperationService } from "./tasks/performTaskOperation.server"; import { ProcessCallbackTimeoutService } from "./tasks/processCallbackTimeout.server"; import { ResumeTaskService } from "./tasks/resumeTask.server"; -import { ResumeTaskRunDependenciesService } from "~/v3/services/resumeTaskRunDependencies.server"; -import { ResumeBatchRunService } from "~/v3/services/resumeBatchRun.server"; -import { ResumeTaskDependencyService } from "~/v3/services/resumeTaskDependency.server"; -import { TimeoutDeploymentService } from "~/v3/services/timeoutDeployment.server"; -import { eventRepository } from "~/v3/eventRepository.server"; -import { ExecuteTasksWaitingForDeployService } from "~/v3/services/executeTasksWaitingForDeploy"; -import { TriggerScheduledTaskService } from "~/v3/services/triggerScheduledTask.server"; -import { PerformTaskAttemptAlertsService } from "~/v3/services/alerts/performTaskAttemptAlerts.server"; -import { DeliverAlertService } from "~/v3/services/alerts/deliverAlert.server"; -import { PerformDeploymentAlertsService } from "~/v3/services/alerts/performDeploymentAlerts.server"; -import { GraphileMigrationHelperService } from "./db/graphileMigrationHelper.server"; -import { PerformBulkActionService } from "~/v3/services/bulk/performBulkAction.server"; -import { CancelTaskRunService } from "~/v3/services/cancelTaskRun.server"; -import { ReplayTaskRunService } from "~/v3/services/replayTaskRun.server"; -import { RequeueTaskRunService } from "~/v3/requeueTaskRun.server"; -import { RetryAttemptService } from "~/v3/services/retryAttempt.server"; const workerCatalog = { indexEndpoint: z.object({ @@ -363,7 +361,6 @@ function getWorkerQueue() { deliverHttpSourceRequest: { priority: 0, // smaller number = higher priority maxAttempts: 14, - queueName: (payload) => `sources:${payload.id}`, handler: async (payload, job) => { const service = new DeliverHttpSourceRequestService(); @@ -373,7 +370,6 @@ function getWorkerQueue() { deliverWebhookRequest: { priority: 0, // smaller number = higher priority maxAttempts: 14, - queueName: (payload) => `webhooks:${payload.id}`, handler: async (payload, job) => { const service = new DeliverWebhookRequestService(); diff --git a/apps/webapp/app/v3/services/alerts/performDeploymentAlerts.server.ts b/apps/webapp/app/v3/services/alerts/performDeploymentAlerts.server.ts index 8f8b945e6..5c41cf8ea 100644 --- a/apps/webapp/app/v3/services/alerts/performDeploymentAlerts.server.ts +++ b/apps/webapp/app/v3/services/alerts/performDeploymentAlerts.server.ts @@ -58,9 +58,7 @@ export class PerformDeploymentAlertsService extends BaseService { }, }); - await DeliverAlertService.enqueue(alert.id, tx, { - queueName: `alert-channel:${alertChannel.id}`, - }); + await DeliverAlertService.enqueue(alert.id, tx); }); } diff --git a/apps/webapp/app/v3/services/alerts/performTaskAttemptAlerts.server.ts b/apps/webapp/app/v3/services/alerts/performTaskAttemptAlerts.server.ts index b0599c613..38486ff4c 100644 --- a/apps/webapp/app/v3/services/alerts/performTaskAttemptAlerts.server.ts +++ b/apps/webapp/app/v3/services/alerts/performTaskAttemptAlerts.server.ts @@ -59,9 +59,7 @@ export class PerformTaskAttemptAlertsService extends BaseService { }, }); - await DeliverAlertService.enqueue(alert.id, tx, { - queueName: `alert-channel:${alertChannel.id}`, - }); + await DeliverAlertService.enqueue(alert.id, tx); }); } diff --git a/apps/webapp/app/v3/services/bulk/performBulkAction.server.ts b/apps/webapp/app/v3/services/bulk/performBulkAction.server.ts index 26d97a1db..07190cdeb 100644 --- a/apps/webapp/app/v3/services/bulk/performBulkAction.server.ts +++ b/apps/webapp/app/v3/services/bulk/performBulkAction.server.ts @@ -88,7 +88,6 @@ export class PerformBulkActionService extends BaseService { }, { jobKey: `performBulkActionItem:${bulkActionItemId}`, - queueName: `bulkActionItem:${groupId}`, } ); } diff --git a/apps/webapp/app/v3/services/resumeBatchRun.server.ts b/apps/webapp/app/v3/services/resumeBatchRun.server.ts index 1d8a8b174..481306b2f 100644 --- a/apps/webapp/app/v3/services/resumeBatchRun.server.ts +++ b/apps/webapp/app/v3/services/resumeBatchRun.server.ts @@ -107,7 +107,7 @@ export class ResumeBatchRunService extends BaseService { { tx, runAt, - queueName: `resumeBatchRun-${batchRunId}`, + jobKey: `resumeBatchRun-${batchRunId}`, } ); } diff --git a/patches/graphile-worker@0.16.6.patch b/patches/graphile-worker@0.16.6.patch index 4d7a4d7a3..68f01bab9 100644 --- a/patches/graphile-worker@0.16.6.patch +++ b/patches/graphile-worker@0.16.6.patch @@ -1,7 +1,16 @@ diff --git a/dist/sql/getJob.js b/dist/sql/getJob.js -index 70cc9b49d7d08c8dd32214f15c463b2a568abd15..bb3fe15317f139edd2b24c5046512813ae88f36e 100644 +index 70cc9b49d7d08c8dd32214f15c463b2a568abd15..5bcb50a4544046e56f6d5dd70e96e26e59998ad5 100644 --- a/dist/sql/getJob.js +++ b/dist/sql/getJob.js +@@ -61,7 +61,7 @@ async function getJob(compiledSharedOptions, withPgClient, tasks, workerId, flag + * + * I recommend you either use strat 0 if you can, or strat 2 otherwise. + */ +- const strategy = 2; ++ const strategy = 0; + const queueClause = strategy === 0 + ? `and jobs.job_queue_id is null` + : strategy === 1 @@ -153,6 +153,13 @@ with j as ( const name = !preparedStatements ? undefined diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e845445bc..fb39ecc2d 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -12,7 +12,7 @@ patchedDependencies: hash: e6nctogrhpxoivwiwy37ersfu4 path: patches/engine.io-parser@5.2.2.patch graphile-worker@0.16.6: - hash: jucduuwnhi5xpqrafqjynulvhe + hash: hdpetta7btqcc7xb5wfkcnanoa path: patches/graphile-worker@0.16.6.patch tsup@8.0.1: hash: a5ztaafw5l4qfghy2hjjuynb34 @@ -464,7 +464,7 @@ importers: version: 10.12.11(react-dom@18.2.0)(react@18.2.0) graphile-worker: specifier: 0.16.6 - version: 0.16.6(patch_hash=jucduuwnhi5xpqrafqjynulvhe)(typescript@5.2.2) + version: 0.16.6(patch_hash=hdpetta7btqcc7xb5wfkcnanoa)(typescript@5.2.2) highlight.run: specifier: ^7.3.4 version: 7.3.4 @@ -23580,7 +23580,7 @@ packages: - supports-color dev: false - /graphile-worker@0.16.6(patch_hash=jucduuwnhi5xpqrafqjynulvhe)(typescript@5.2.2): + /graphile-worker@0.16.6(patch_hash=hdpetta7btqcc7xb5wfkcnanoa)(typescript@5.2.2): resolution: {integrity: sha512-e7gGYDmGqzju2l83MpzX8vNG/lOtVJiSzI3eZpAFubSxh/cxs7sRrRGBGjzBP1kNG0H+c95etPpNRNlH65PYhw==} engines: {node: '>=14.0.0'} hasBin: true