From c9a43795842720d5d7fc0b5bed66f5a83217bab1 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Tue, 26 Sep 2023 11:31:01 +1000 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B5=D0=B4=D0=B5=D0=BB?= =?UTF-8?q?=D0=B0=D0=BB=20TemplateMqSyncer=20=D0=BD=D0=B0=20IMqService?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PARR.EsppTemplateSync/TemplateMQSyncer.cs | 94 +++-------------------- 1 file changed, 9 insertions(+), 85 deletions(-) diff --git a/PARR.EsppTemplateSync/TemplateMQSyncer.cs b/PARR.EsppTemplateSync/TemplateMQSyncer.cs index 3f6a9377..fdce03be 100644 --- a/PARR.EsppTemplateSync/TemplateMQSyncer.cs +++ b/PARR.EsppTemplateSync/TemplateMQSyncer.cs @@ -2,9 +2,6 @@ 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 { @@ -13,19 +10,19 @@ namespace PARR.EsppTemplateSync private readonly ILogger logger; private readonly IManager manager; private readonly GlobalSettings globalSettings; - - private IConnection? connection; - private IModel? channel; + private readonly IMqService mqService; public TemplateMQSyncer( ILogger logger, IManager manager, - GlobalSettings globalSettings + GlobalSettings globalSettings, + IMqService mqService ) { this.logger = logger; this.manager = manager; this.globalSettings = globalSettings; + this.mqService = mqService; if (globalSettings.MqSettings == null) { @@ -34,96 +31,23 @@ 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() { - InitConnection(); + var isConnected = mqService.InitConsumer(globalSettings!.MqSettings!, manager.ManageStringAsync); + + if (!isConnected) + throw new Exception("Ошибка при подключении к RabbitMq"); logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}."); - - GetMessages();//TODO переименовать в какой-нибудь AttachListener } public void Stop() { - channel?.Close(); - connection?.Close(); + mqService.Dispose(); 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); - } } }