using Microsoft.Extensions.Logging; using PARR.EsppTemplateSync.Services; using PARR.EsppTemplateSync.Settings; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Text; namespace PARR.EsppTemplateSync { internal class TemplateMQSyncer : ITemplateSyncer { private readonly ILogger logger; private readonly IManager manager; private readonly GlobalSettings globalSettings; private IConnection? connection; private IModel? channel; public TemplateMQSyncer( ILogger logger, IManager manager, GlobalSettings globalSettings ) { this.logger = logger; this.manager = manager; this.globalSettings = globalSettings; if (globalSettings.MqSettings == null) { logger.LogError("Нет секции настроек хранилища. MqSettings, EsppTemplates"); throw new Exception("Нет секции настроек хранилища. MqSettings, EsppTemplates"); } } public void Start() { InitConnection(); logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}."); GetMessages();//TODO переименовать в какой-нибудь AttachListener } public void Stop() { 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.UserName; 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); } } }