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