Files
parr_api/PARR.EsppTemplateSync/TemplateMQSyncer.cs

89 lines
3.4 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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<TemplateMQSyncer> logger;
private readonly IManager manager;
private readonly GlobalSettings globalSettings;
private IConnection connection;
private RabbitMQ.Client.IModel channel;
public TemplateMQSyncer(
ILogger<TemplateMQSyncer> 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<String, Object> args = new Dictionary<String, Object>();
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} закрыто === === ===");
}
}
}