using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using PARR.DAL.Context; using PARR.EsppTemplateSync.Services; using PARR.EsppTemplateSync.Settings; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Runtime.Serialization.DataContracts; using System.Text; namespace PARR.EsppTemplateSync { internal class TemplateMQSyncer : ITemplateSyncer { private readonly ILogger logger; private readonly IManager manager; private readonly GlobalSettings globalSettings; private readonly IServiceProvider serviceProvider; private IConnection connection; private RabbitMQ.Client.IModel channel; public TemplateMQSyncer( ILogger logger, IManager manager, GlobalSettings globalSettings, IServiceProvider serviceProvider ) { this.logger = logger; this.manager = manager; this.globalSettings = globalSettings; this.serviceProvider = serviceProvider; 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); GetMessages(); } 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); // ------------ //using var scope = serviceProvider.CreateScope(); //var services = scope.ServiceProvider; //var manager = services.GetRequiredService(); //// ------------ //if (manager != null) 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} закрыто === === ==="); } } }