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); - } } }