/
githubmirror
/
novu
Обзор
Документация
Войти
/
githubmirror
/
novu
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
next
libs/application-generic/src/modules/queues.module.ts
126 строк
4 KB
Himanshu Garg
refactor(api, worker): remove Cloudflare Scheduler integration and related configurations fixes NV-8518 (#12216)
04 авг 2026, 15:17
Не верифицирован
04 авг 2026, 15:17
7d0c0bd
Код
Авторство
О чём код?
import { DynamicModule, Module, OnApplicationShutdown, Provider } from '@nestjs/common'; import { CommunityOrganizationRepository, MessageRepository } from '@novu/dal'; import { JobTopicNameEnum } from '@novu/shared'; import { featureFlagsService } from '../custom-providers'; import { ActiveJobsMetricQueueServiceHealthIndicator, InboundParseQueueServiceHealthIndicator, StandardQueueServiceHealthIndicator, SubscriberProcessQueueHealthIndicator, WebSocketsQueueServiceHealthIndicator, WorkflowQueueServiceHealthIndicator, } from '../health'; import { ReadinessService, SocketWorkerService, SqsService, WorkflowInMemoryProviderService } from '../services'; import { ActiveJobsMetricQueueService, InboundParseQueueService, StandardQueueService, SubscriberProcessQueueService, WebSocketsQueueService, WorkflowQueueService, } from '../services/queues'; import { ActiveJobsMetricWorkerService } from '../services/workers'; const memoryQueueService = { provide: WorkflowInMemoryProviderService, useFactory: async () => { const memoryService = new WorkflowInMemoryProviderService(); await memoryService.initialize(); return memoryService; }, }; const INTERNAL_MODULE_PROVIDERS = [memoryQueueService, featureFlagsService]; const BASE_PROVIDERS: Provider[] = [ReadinessService, CommunityOrganizationRepository, SqsService]; @Module({ providers: [], exports: [], }) export class QueuesModule implements OnApplicationShutdown { static forRoot(entities: JobTopicNameEnum[] = []): DynamicModule { if (!entities.length) { entities = Object.values(JobTopicNameEnum); } const healthIndicators = []; const tokenList = []; const DYNAMIC_PROVIDERS = [...BASE_PROVIDERS]; for (const entity of entities) { switch (entity) { case JobTopicNameEnum.INBOUND_PARSE_MAIL: healthIndicators.push(InboundParseQueueServiceHealthIndicator); tokenList.push(InboundParseQueueService); DYNAMIC_PROVIDERS.push(InboundParseQueueService, InboundParseQueueServiceHealthIndicator); break; case JobTopicNameEnum.WORKFLOW: healthIndicators.push(WorkflowQueueServiceHealthIndicator); tokenList.push(WorkflowQueueService); DYNAMIC_PROVIDERS.push(WorkflowQueueService, WorkflowQueueServiceHealthIndicator); break; case JobTopicNameEnum.WEB_SOCKETS: healthIndicators.push(WebSocketsQueueServiceHealthIndicator); tokenList.push(WebSocketsQueueService); DYNAMIC_PROVIDERS.push( MessageRepository, SocketWorkerService, WebSocketsQueueService, WebSocketsQueueServiceHealthIndicator ); break; case JobTopicNameEnum.STANDARD: healthIndicators.push(StandardQueueServiceHealthIndicator); tokenList.push(StandardQueueService); DYNAMIC_PROVIDERS.push(StandardQueueService, StandardQueueServiceHealthIndicator); break; case JobTopicNameEnum.PROCESS_SUBSCRIBER: healthIndicators.push(SubscriberProcessQueueHealthIndicator); tokenList.push(SubscriberProcessQueueService); DYNAMIC_PROVIDERS.push(SubscriberProcessQueueService, SubscriberProcessQueueHealthIndicator); break; case JobTopicNameEnum.ACTIVE_JOBS_METRIC: healthIndicators.push(ActiveJobsMetricQueueServiceHealthIndicator); tokenList.push(ActiveJobsMetricQueueService); DYNAMIC_PROVIDERS.push( ActiveJobsMetricQueueService, ActiveJobsMetricQueueServiceHealthIndicator, ActiveJobsMetricWorkerService ); break; default: break; } } DYNAMIC_PROVIDERS.push({ provide: 'BULLMQ_LIST', useFactory: (...args: any[]) => { return args; }, inject: tokenList, }); DYNAMIC_PROVIDERS.push({ provide: 'QUEUE_HEALTH_INDICATORS', useFactory: (...args: any[]) => { return args; }, inject: healthIndicators, }); return { module: QueuesModule, providers: [...DYNAMIC_PROVIDERS, ...INTERNAL_MODULE_PROVIDERS], exports: [...DYNAMIC_PROVIDERS], }; } constructor(private workflowInMemoryProviderService: WorkflowInMemoryProviderService) {} async onApplicationShutdown() { await this.workflowInMemoryProviderService.shutdown(); } }