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 RabbitMQ.Client.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"); } 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(); Dictionary args = new Dictionary(); args.Add("x-queue-mode", "lazy"); channel.QueueDeclare(queue: globalSettings.MqSettings.QueueName, durable: true, exclusive: false, autoDelete: false, arguments: args); } public async Task StartAsync() { logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}. Интервал: {globalSettings.MqSettings!.CheckIntervalSeconds} секунд."); CancellationTokenSource cancelTokenSource = new CancellationTokenSource(); CancellationToken token = cancelTokenSource.Token; Task task = new Task(() => { GetMessages(); }, token); task.Start(); await task.WaitAsync(token); } private void GetMessages() { var consumer = new AsyncEventingBasicConsumer(channel); consumer.Received += async (ch, ea) => { var content = Encoding.UTF8.GetString(ea.Body.ToArray()); // Обрабатываем полученное сообщение await manager.ManageStringAsync(content); channel.BasicAck(ea.DeliveryTag, false); await Task.Yield(); }; channel.BasicConsume(globalSettings.MqSettings!.QueueName, false, consumer); } public void Stop() { channel.Close(); connection.Close(); logger.LogInformation($"=== === === Соединение с очередью {globalSettings.MqSettings!.QueueName} закрыто === === ==="); } } }