ЦИТК-901 (Добавление поддержки протоколов MQTT и KAFKA) #2

Merged
Mim merged 2 commits from Dollerok/P8-ExchangeService:master into master 2024-09-27 11:42:05 +03:00
2 changed files with 6 additions and 6 deletions
Showing only changes of commit ffb12c8550 - Show all commits

View File

@ -705,15 +705,15 @@ class InQueue extends EventEmitter {
} }
//Закрытие подключений //Закрытие подключений
async stopConnections() { stopConnections() {
//Если у нас есть соединения с MQTT //Если у нас есть соединения с MQTT
if (this.mqttConnections.length !== 0) { if (this.mqttConnections.length !== 0) {
//Закрываем их //Закрываем их
for (let connection of this.mqttConnections) { for (let connection of this.mqttConnections) {
try { try {
await connection.end(); connection.end();
} catch (e) { } catch (e) {
await this.logger.error(`Ошибка завершения MQTT подключения: ${makeErrorText(e)}`); this.logger.error(`Ошибка завершения MQTT подключения: ${makeErrorText(e)}`);
} }
} }
} }
@ -722,9 +722,9 @@ class InQueue extends EventEmitter {
//Закрываем их //Закрываем их
for (let connection of this.kafkaConnections) { for (let connection of this.kafkaConnections) {
try { try {
await connection.disconnect(); connection.disconnect();
} catch (e) { } catch (e) {
await this.logger.error(`Ошибка завершения Kafka подключения: ${makeErrorText(e)}`); this.logger.error(`Ошибка завершения Kafka подключения: ${makeErrorText(e)}`);
} }
} }
} }

View File

@ -68,7 +68,7 @@ const subscribeMQTT = async ({ settings, service, processMessage, logger }) => {
logger.error(`Соединение с MQTT потеряно (${sBroker})`); logger.error(`Соединение с MQTT потеряно (${sBroker})`);
}); });
//Прослушиваем восстановление соединения //Прослушиваем восстановление соединения
client.on("reconnect", () => { client.on("connect", () => {
//Сообщим о восстановлении соединения //Сообщим о восстановлении соединения
logger.info(`Соединение с MQTT восстановлено (${sBroker})`); logger.info(`Соединение с MQTT восстановлено (${sBroker})`);
}); });