Team Ai
Apppublic

oex2003/evolution-api

sourceHugging Faceupdated 1y agoView on Hugging Face
0likes
rabbitmq.controller.ts410 linesDownload Raw Back to rabbitmq
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