From 332dc95f9edb5da44359d83a682c959c4da250f0 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Fri, 22 Sep 2023 09:43:56 +1000 Subject: [PATCH] =?UTF-8?q?TemplateMQSyncer=20=D0=B2=D0=B5=D1=80=D0=BD?= =?UTF-8?q?=D1=83=D0=BB=20=D0=BA=D0=B0=D0=BA=20=D0=B1=D1=8B=D0=BB,=20?= =?UTF-8?q?=D0=B1=D0=B5=D0=B7=20=D0=B8=D1=81=D0=BF=D0=BE=D0=BB=D1=8C=D0=B7?= =?UTF-8?q?=D0=BE=D0=B2=D0=B0=D0=BD=D0=B8=D1=8F=20IMqService?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Services/Implementations/MqService.cs | 6 +- PARR.BLL/Services/Interfaces/IMqService.cs | 2 +- PARR.EsppTemplateSync/TemplateMQSyncer.cs | 94 +++++++++++++++++-- PARR.GeneratorTemplates/GeneratorTemplate.cs | 7 +- docker-compose.espp-template-sync.yml | 2 +- 5 files changed, 99 insertions(+), 12 deletions(-) diff --git a/PARR.BLL/Services/Implementations/MqService.cs b/PARR.BLL/Services/Implementations/MqService.cs index 12720100..efaeec7b 100644 --- a/PARR.BLL/Services/Implementations/MqService.cs +++ b/PARR.BLL/Services/Implementations/MqService.cs @@ -91,7 +91,7 @@ namespace PARR.BLL.Services.Implementations } - public void InitConsumer(IMqSettings mqSettings) + public bool InitConsumer(IMqSettings mqSettings) { logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}"); @@ -140,10 +140,14 @@ namespace PARR.BLL.Services.Implementations channel.BasicConsume(mqSettings.QueueName, false, consumer); + return true; + } catch (Exception ex) { logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); + + return false; } } diff --git a/PARR.BLL/Services/Interfaces/IMqService.cs b/PARR.BLL/Services/Interfaces/IMqService.cs index 3db517c6..9e1359bb 100644 --- a/PARR.BLL/Services/Interfaces/IMqService.cs +++ b/PARR.BLL/Services/Interfaces/IMqService.cs @@ -9,7 +9,7 @@ namespace PARR.BLL.Services.Interfaces { event ReceivedHandler? Received; - void InitConsumer(IMqSettings mqSettings); + bool InitConsumer(IMqSettings mqSettings); MqSendResult Send(IMqSettings mqSettings, string[] msgList); } } diff --git a/PARR.EsppTemplateSync/TemplateMQSyncer.cs b/PARR.EsppTemplateSync/TemplateMQSyncer.cs index a7192008..3f6a9377 100644 --- a/PARR.EsppTemplateSync/TemplateMQSyncer.cs +++ b/PARR.EsppTemplateSync/TemplateMQSyncer.cs @@ -2,6 +2,9 @@ using PARR.BLL.Services.Interfaces; using PARR.EsppTemplateSync.Services; using PARR.EsppTemplateSync.Settings; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using System.Text; namespace PARR.EsppTemplateSync { @@ -10,19 +13,19 @@ namespace PARR.EsppTemplateSync private readonly ILogger logger; private readonly IManager manager; private readonly GlobalSettings globalSettings; - private readonly IMqService mqService; + + private IConnection? connection; + private IModel? channel; public TemplateMQSyncer( ILogger logger, IManager manager, - GlobalSettings globalSettings, - IMqService mqService + GlobalSettings globalSettings ) { this.logger = logger; this.manager = manager; this.globalSettings = globalSettings; - this.mqService = mqService; if (globalSettings.MqSettings == null) { @@ -31,19 +34,96 @@ namespace PARR.EsppTemplateSync } } + #region должно быть так: + //public void Start() + //{ + // mqService.Received += async (msg) => + // { + // await manager.ManageStringAsync(msg); + // }; + + // var isConnected = mqService.InitConsumer(globalSettings!.MqSettings!); + // if (!isConnected) + // throw new Exception("Ошибка при подключении к RabbitMq"); + + // logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}."); + //} + + //public void Stop() + //{ + // mqService.Dispose(); + + // logger.LogInformation($"=== === === Соединение с очередью {globalSettings.MqSettings!.QueueName} закрыто === === ==="); + //} + #endregion + public void Start() { - mqService.Received += async (msg) => await manager.ManageStringAsync(msg); - mqService.InitConsumer(globalSettings!.MqSettings!); + InitConnection(); logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}."); + + GetMessages();//TODO переименовать в какой-нибудь AttachListener } public void Stop() { - mqService.Dispose(); + channel?.Close(); + connection?.Close(); logger.LogInformation($"=== === === Соединение с очередью {globalSettings.MqSettings!.QueueName} закрыто === === ==="); } + + private void InitConnection() + { + logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {globalSettings.MqSettings!.HostName}"); + + var factory = new ConnectionFactory { HostName = globalSettings.MqSettings!.HostName }; + factory.UserName = globalSettings.MqSettings.User; + factory.Password = globalSettings.MqSettings.Password; + factory.AutomaticRecoveryEnabled = true; + factory.DispatchConsumersAsync = true; + + connection = factory.CreateConnection(); + channel = connection.CreateModel(); + + var args = new Dictionary + { + { "x-queue-mode", "lazy" } + }; + + channel.QueueDeclare(queue: globalSettings.MqSettings.QueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: args); + + logger.LogInformation("Соединение с RabbitMQ установлено."); + } + + private void GetMessages() + { + if (channel == null || connection == null) + { + logger.LogError("Отсутствует соедниение с RabbitMQ."); + return; + } + + var consumer = new AsyncEventingBasicConsumer(channel); + + consumer.Received += async (ch, ea) => + { + var content = Encoding.UTF8.GetString(ea.Body.ToArray()); + logger.LogDebug($"Получено сообщение: {content}"); + + await manager.ManageStringAsync(content); + + channel.BasicAck(ea.DeliveryTag, false); + await Task.Yield(); + }; + + channel.BasicConsume(globalSettings.MqSettings!.QueueName, false, consumer); + } } } + diff --git a/PARR.GeneratorTemplates/GeneratorTemplate.cs b/PARR.GeneratorTemplates/GeneratorTemplate.cs index 00b33020..a55c33d7 100644 --- a/PARR.GeneratorTemplates/GeneratorTemplate.cs +++ b/PARR.GeneratorTemplates/GeneratorTemplate.cs @@ -34,13 +34,16 @@ namespace PARR.GeneratorTemplates // проверить как ту дела с асинком // логгер работает в разных потоках все таки или нет? - + // поиск в существующих шаблонах, есть ли такие и есть ли у них эти работы // поиск в хостах эк которых еще нет в шаблонах и у них такие работы // генерим список и вставляем в шаблоны }; - mqService.InitConsumer(mqSettings); + var isConnected = mqService.InitConsumer(mqSettings); + + if (!isConnected) + throw new Exception("Ошибка при подключении к RabbitMq"); } diff --git a/docker-compose.espp-template-sync.yml b/docker-compose.espp-template-sync.yml index 640d7b7d..10e4685b 100644 --- a/docker-compose.espp-template-sync.yml +++ b/docker-compose.espp-template-sync.yml @@ -9,7 +9,7 @@ services: environment: - ASPNETCORE_ENVIRONMENT=Production - TZ=Europe/Moscow - - GlobalSettings__MqSettings__HostName=parr-rabbitmq + #- GlobalSettings__MqSettings__HostName=parr-rabbitmq #volumes: #- /var/log/parr-espp-template-sync:/app/log #- parr_storage_dev:/app/Data