oex2003/evolution-api
0
1import { PrismaRepository } from '@api/repository/repository.service';2import { WAMonitoringService } from '@api/services/monitor.service';3import { configService, Log, Rabbitmq } from '@config/env.config';4import { Logger } from '@config/logger.config';5import * as amqp from 'amqplib/callback_api';6 7import { EmitData, EventController, EventControllerInterface } from '../event.controller';8 9export class RabbitmqController extends EventController implements EventControllerInterface {10 public amqpChannel: amqp.Channel | null = null;11 private amqpConnection: amqp.Connection | null = null;12 private readonly logger = new Logger('RabbitmqController');13 private reconnectAttempts = 0;14 private maxReconnectAttempts = 10;15 private reconnectDelay = 5000; // 5 seconds16 private isReconnecting = false;17 18 constructor(prismaRepository: PrismaRepository, waMonitor: WAMonitoringService) {19 super(prismaRepository, waMonitor, configService.get<Rabbitmq>('RABBITMQ')?.ENABLED, 'rabbitmq');20 }21 22 public async init(): Promise<void> {23 if (!this.status) {24 return;25 }26 27 await this.connect();28 }29 30 private async connect(): Promise<void> {31 return new Promise<void>((resolve, reject) => {32 const uri = configService.get<Rabbitmq>('RABBITMQ').URI;33 const frameMax = configService.get<Rabbitmq>('RABBITMQ').FRAME_MAX;34 const rabbitmqExchangeName = configService.get<Rabbitmq>('RABBITMQ').EXCHANGE_NAME;35 36 const url = new URL(uri);37 const connectionOptions = {38 protocol: url.protocol.slice(0, -1),39 hostname: url.hostname,40 port: url.port || 5672,41 username: url.username || 'guest',42 password: url.password || 'guest',43 vhost: url.pathname.slice(1) || '/',44 frameMax: frameMax,45 heartbeat: 30, // Add heartbeat of 30 seconds46 };47 48 amqp.connect(connectionOptions, (error, connection) => {49 if (error) {50 this.logger.error({51 local: 'RabbitmqController.connect',52 message: 'Failed to connect to RabbitMQ',53 error: error.message || error,54 });55 reject(error);56 return;57 }58 59 // Connection event handlers60 connection.on('error', (err) => {61 this.logger.error({62 local: 'RabbitmqController.connectionError',63 message: 'RabbitMQ connection error',64 error: err.message || err,65 });66 this.handleConnectionLoss();67 });68 69 connection.on('close', () => {70 this.logger.warn('RabbitMQ connection closed');71 this.handleConnectionLoss();72 });73 74 connection.createChannel((channelError, channel) => {75 if (channelError) {76 this.logger.error({77 local: 'RabbitmqController.createChannel',78 message: 'Failed to create RabbitMQ channel',79 error: channelError.message || channelError,80 });81 reject(channelError);82 return;83 }84 85 // Channel event handlers86 channel.on('error', (err) => {87 this.logger.error({88 local: 'RabbitmqController.channelError',89 message: 'RabbitMQ channel error',90 error: err.message || err,91 });92 this.handleConnectionLoss();93 });94 95 channel.on('close', () => {96 this.logger.warn('RabbitMQ channel closed');97 this.handleConnectionLoss();98 });99 100 const exchangeName = rabbitmqExchangeName;101 102 channel.assertExchange(exchangeName, 'topic', {103 durable: true,104 autoDelete: false,105 });106 107 this.amqpConnection = connection;108 this.amqpChannel = channel;109 this.reconnectAttempts = 0; // Reset reconnect attempts on successful connection110 this.isReconnecting = false;111 112 this.logger.info('AMQP initialized successfully');113 114 resolve();115 });116 });117 })118 .then(() => {119 if (configService.get<Rabbitmq>('RABBITMQ')?.GLOBAL_ENABLED) {120 this.initGlobalQueues();121 }122 })123 .catch((error) => {124 this.logger.error({125 local: 'RabbitmqController.init',126 message: 'Failed to initialize AMQP',127 error: error.message || error,128 });129 this.scheduleReconnect();130 throw error;131 });132 }133 134 private handleConnectionLoss(): void {135 if (this.isReconnecting) {136 return; // Already attempting to reconnect137 }138 139 this.amqpChannel = null;140 this.amqpConnection = null;141 this.scheduleReconnect();142 }143 144 private scheduleReconnect(): void {145 if (this.reconnectAttempts >= this.maxReconnectAttempts) {146 this.logger.error(147 `Maximum reconnect attempts (${this.maxReconnectAttempts}) reached. Stopping reconnection attempts.`,148 );149 return;150 }151 152 if (this.isReconnecting) {153 return; // Already scheduled154 }155 156 this.isReconnecting = true;157 this.reconnectAttempts++;158 159 const delay = this.reconnectDelay * Math.pow(2, Math.min(this.reconnectAttempts - 1, 5)); // Exponential backoff with max delay160 161 this.logger.info(162 `Scheduling RabbitMQ reconnection attempt ${this.reconnectAttempts}/${this.maxReconnectAttempts} in ${delay}ms`,163 );164 165 setTimeout(async () => {166 try {167 this.logger.info(168 `Attempting to reconnect to RabbitMQ (attempt ${this.reconnectAttempts}/${this.maxReconnectAttempts})`,169 );170 await this.connect();171 this.logger.info('Successfully reconnected to RabbitMQ');172 } catch (error) {173 this.logger.error({174 local: 'RabbitmqController.scheduleReconnect',175 message: `Reconnection attempt ${this.reconnectAttempts} failed`,176 error: error.message || error,177 });178 this.isReconnecting = false;179 this.scheduleReconnect();180 }181 }, delay);182 }183 184 private set channel(channel: amqp.Channel) {185 this.amqpChannel = channel;186 }187 188 public get channel(): amqp.Channel {189 return this.amqpChannel;190 }191 192 private async ensureConnection(): Promise<boolean> {193 if (!this.amqpChannel) {194 this.logger.warn('AMQP channel is not available, attempting to reconnect...');195 if (!this.isReconnecting) {196 this.scheduleReconnect();197 }198 return false;199 }200 return true;201 }202 203 public async emit({204 instanceName,205 origin,206 event,207 data,208 serverUrl,209 dateTime,210 sender,211 apiKey,212 integration,213 }: EmitData): Promise<void> {214 if (integration && !integration.includes('rabbitmq')) {215 return;216 }217 218 if (!this.status) {219 return;220 }221 222 if (!(await this.ensureConnection())) {223 this.logger.warn(`Failed to emit event ${event} for instance ${instanceName}: No AMQP connection`);224 return;225 }226 227 const instanceRabbitmq = await this.get(instanceName);228 const rabbitmqLocal = instanceRabbitmq?.events;229 const rabbitmqGlobal = configService.get<Rabbitmq>('RABBITMQ').GLOBAL_ENABLED;230 const rabbitmqEvents = configService.get<Rabbitmq>('RABBITMQ').EVENTS;231 const prefixKey = configService.get<Rabbitmq>('RABBITMQ').PREFIX_KEY;232 const rabbitmqExchangeName = configService.get<Rabbitmq>('RABBITMQ').EXCHANGE_NAME;233 const we = event.replace(/[.-]/gm, '_').toUpperCase();234 const logEnabled = configService.get<Log>('LOG').LEVEL.includes('WEBHOOKS');235 236 const message = {237 event,238 instance: instanceName,239 data,240 server_url: serverUrl,241 date_time: dateTime,242 sender,243 apikey: apiKey,244 };245 246 if (instanceRabbitmq?.enabled && this.amqpChannel) {247 if (Array.isArray(rabbitmqLocal) && rabbitmqLocal.includes(we)) {248 const exchangeName = instanceName ?? rabbitmqExchangeName;249 250 let retry = 0;251 252 while (retry < 3) {253 try {254 await this.amqpChannel.assertExchange(exchangeName, 'topic', {255 durable: true,256 autoDelete: false,257 });258 259 const eventName = event.replace(/_/g, '.').toLowerCase();260 261 const queueName = `${instanceName}.${eventName}`;262 263 await this.amqpChannel.assertQueue(queueName, {264 durable: true,265 autoDelete: false,266 arguments: {267 'x-queue-type': 'quorum',268 },269 });270 271 await this.amqpChannel.bindQueue(queueName, exchangeName, eventName);272 273 await this.amqpChannel.publish(exchangeName, event, Buffer.from(JSON.stringify(message)));274 275 if (logEnabled) {276 const logData = {277 local: `${origin}.sendData-RabbitMQ`,278 ...message,279 };280 281 this.logger.log(logData);282 }283 284 break;285 } catch (error) {286 this.logger.error({287 local: 'RabbitmqController.emit',288 message: `Error publishing local RabbitMQ message (attempt ${retry + 1}/3)`,289 error: error.message || error,290 });291 retry++;292 if (retry >= 3) {293 this.handleConnectionLoss();294 }295 }296 }297 }298 }299 300 if (rabbitmqGlobal && rabbitmqEvents[we] && this.amqpChannel) {301 const exchangeName = rabbitmqExchangeName;302 303 let retry = 0;304 305 while (retry < 3) {306 try {307 await this.amqpChannel.assertExchange(exchangeName, 'topic', {308 durable: true,309 autoDelete: false,310 });311 312 const queueName = prefixKey313 ? `${prefixKey}.${event.replace(/_/g, '.').toLowerCase()}`314 : event.replace(/_/g, '.').toLowerCase();315 316 await this.amqpChannel.assertQueue(queueName, {317 durable: true,318 autoDelete: false,319 arguments: {320 'x-queue-type': 'quorum',321 },322 });323 324 await this.amqpChannel.bindQueue(queueName, exchangeName, event);325 326 await this.amqpChannel.publish(exchangeName, event, Buffer.from(JSON.stringify(message)));327 328 if (logEnabled) {329 const logData = {330 local: `${origin}.sendData-RabbitMQ-Global`,331 ...message,332 };333 334 this.logger.log(logData);335 }336 337 break;338 } catch (error) {339 this.logger.error({340 local: 'RabbitmqController.emit',341 message: `Error publishing global RabbitMQ message (attempt ${retry + 1}/3)`,342 error: error.message || error,343 });344 retry++;345 if (retry >= 3) {346 this.handleConnectionLoss();347 }348 }349 }350 }351 }352 353 private async initGlobalQueues(): Promise<void> {354 this.logger.info('Initializing global queues');355 356 if (!(await this.ensureConnection())) {357 this.logger.error('Cannot initialize global queues: No AMQP connection');358 return;359 }360 361 const rabbitmqExchangeName = configService.get<Rabbitmq>('RABBITMQ').EXCHANGE_NAME;362 const events = configService.get<Rabbitmq>('RABBITMQ').EVENTS;363 const prefixKey = configService.get<Rabbitmq>('RABBITMQ').PREFIX_KEY;364 365 if (!events) {366 this.logger.warn('No events to initialize on AMQP');367 return;368 }369 370 const eventKeys = Object.keys(events);371 372 for (const event of eventKeys) {373 if (events[event] === false) continue;374 375 try {376 const queueName =377 prefixKey !== ''378 ? `${prefixKey}.${event.replace(/_/g, '.').toLowerCase()}`379 : `${event.replace(/_/g, '.').toLowerCase()}`;380 const exchangeName = rabbitmqExchangeName;381 382 await this.amqpChannel.assertExchange(exchangeName, 'topic', {383 durable: true,384 autoDelete: false,385 });386 387 await this.amqpChannel.assertQueue(queueName, {388 durable: true,389 autoDelete: false,390 arguments: {391 'x-queue-type': 'quorum',392 },393 });394 395 await this.amqpChannel.bindQueue(queueName, exchangeName, event);396 397 this.logger.info(`Global queue initialized: ${queueName}`);398 } catch (error) {399 this.logger.error({400 local: 'RabbitmqController.initGlobalQueues',401 message: `Failed to initialize global queue for event ${event}`,402 error: error.message || error,403 });404 this.handleConnectionLoss();405 break;406 }407 }408 }409}410 