diff --git a/apps/server/src/integrations/queue/queue.module.ts b/apps/server/src/integrations/queue/queue.module.ts index eeb74cbb3..509dbdf47 100644 --- a/apps/server/src/integrations/queue/queue.module.ts +++ b/apps/server/src/integrations/queue/queue.module.ts @@ -1,105 +1,20 @@ import { Global, Module } from '@nestjs/common'; import { BullModule } from '@nestjs/bullmq'; import { EnvironmentService } from '../environment/environment.service'; -import { createRetryStrategy, parseRedisUrl } from '../../common/helpers'; -import { QueueName } from './constants'; import { GeneralQueueProcessor } from './processors/general-queue.processor'; +import { + bullConfigFactory, + createQueueRegistrations, +} from './queue.registrations'; @Global() @Module({ imports: [ BullModule.forRootAsync({ - useFactory: (environmentService: EnvironmentService) => { - const redisConfig = parseRedisUrl(environmentService.getRedisUrl()); - return { - connection: { - host: redisConfig.host, - port: redisConfig.port, - password: redisConfig.password, - db: redisConfig.db, - family: redisConfig.family, - retryStrategy: createRetryStrategy(), - }, - defaultJobOptions: { - attempts: 3, - backoff: { - type: 'exponential', - delay: 20 * 1000, - }, - removeOnComplete: { - count: 200, - }, - removeOnFail: { - count: 100, - }, - }, - }; - }, + useFactory: bullConfigFactory, inject: [EnvironmentService], }), - BullModule.registerQueue({ - name: QueueName.EMAIL_QUEUE, - }), - BullModule.registerQueue({ - name: QueueName.ATTACHMENT_QUEUE, - }), - BullModule.registerQueue({ - name: QueueName.GENERAL_QUEUE, - }), - BullModule.registerQueue({ - name: QueueName.BILLING_QUEUE, - }), - BullModule.registerQueue({ - name: QueueName.FILE_TASK_QUEUE, - defaultJobOptions: { - removeOnComplete: true, - removeOnFail: true, - attempts: 1, - }, - }), - BullModule.registerQueue({ - name: QueueName.SEARCH_QUEUE, - defaultJobOptions: { - removeOnComplete: true, - removeOnFail: true, - attempts: 2, - }, - }), - BullModule.registerQueue({ - name: QueueName.AI_QUEUE, - defaultJobOptions: { - removeOnComplete: true, - removeOnFail: true, - attempts: 1, - }, - }), - BullModule.registerQueue({ - name: QueueName.HISTORY_QUEUE, - defaultJobOptions: { - removeOnComplete: true, - removeOnFail: true, - attempts: 2, - }, - }), - BullModule.registerQueue({ - name: QueueName.NOTIFICATION_QUEUE, - }), - BullModule.registerQueue({ - name: QueueName.AUDIT_QUEUE, - defaultJobOptions: { - removeOnComplete: true, - removeOnFail: true, - attempts: 3, - }, - }), - BullModule.registerQueue({ - name: QueueName.BASE_QUEUE, - defaultJobOptions: { - attempts: 2, - removeOnComplete: { count: 200 }, - removeOnFail: { count: 100 }, - }, - }), + ...createQueueRegistrations(), ], exports: [BullModule], providers: [GeneralQueueProcessor], diff --git a/apps/server/src/integrations/queue/queue.registrations.ts b/apps/server/src/integrations/queue/queue.registrations.ts new file mode 100644 index 000000000..bb0655177 --- /dev/null +++ b/apps/server/src/integrations/queue/queue.registrations.ts @@ -0,0 +1,97 @@ +import { BullModule } from '@nestjs/bullmq'; +import { EnvironmentService } from '../environment/environment.service'; +import { createRetryStrategy, parseRedisUrl } from '../../common/helpers'; +import { QueueName } from './constants'; + +export const bullConfigFactory = (environmentService: EnvironmentService) => { + const redisConfig = parseRedisUrl(environmentService.getRedisUrl()); + return { + connection: { + host: redisConfig.host, + port: redisConfig.port, + password: redisConfig.password, + db: redisConfig.db, + family: redisConfig.family, + retryStrategy: createRetryStrategy(), + }, + defaultJobOptions: { + attempts: 3, + backoff: { + type: 'exponential', + delay: 20 * 1000, + }, + removeOnComplete: { + count: 200, + }, + removeOnFail: { + count: 100, + }, + }, + }; +}; + +export const createQueueRegistrations = () => [ + BullModule.registerQueue({ + name: QueueName.EMAIL_QUEUE, + }), + BullModule.registerQueue({ + name: QueueName.ATTACHMENT_QUEUE, + }), + BullModule.registerQueue({ + name: QueueName.GENERAL_QUEUE, + }), + BullModule.registerQueue({ + name: QueueName.BILLING_QUEUE, + }), + BullModule.registerQueue({ + name: QueueName.FILE_TASK_QUEUE, + defaultJobOptions: { + removeOnComplete: true, + removeOnFail: true, + attempts: 1, + }, + }), + BullModule.registerQueue({ + name: QueueName.SEARCH_QUEUE, + defaultJobOptions: { + removeOnComplete: true, + removeOnFail: true, + attempts: 2, + }, + }), + BullModule.registerQueue({ + name: QueueName.AI_QUEUE, + defaultJobOptions: { + removeOnComplete: true, + removeOnFail: true, + attempts: 1, + }, + }), + BullModule.registerQueue({ + name: QueueName.HISTORY_QUEUE, + defaultJobOptions: { + removeOnComplete: true, + removeOnFail: true, + attempts: 2, + }, + }), + BullModule.registerQueue({ + name: QueueName.NOTIFICATION_QUEUE, + }), + BullModule.registerQueue({ + name: QueueName.AUDIT_QUEUE, + defaultJobOptions: { + removeOnComplete: true, + removeOnFail: true, + attempts: 3, + }, + }), + BullModule.registerQueue({ + name: QueueName.BASE_QUEUE, + defaultJobOptions: { + attempts: 2, + removeOnComplete: { count: 200 }, + removeOnFail: { count: 100 }, + }, + }), +];