Files
parr_api/PARR.EsppTemplateSync/TemplateMQSyncer.cs

130 lines
4.5 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.BLL.Services.Interfaces;
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 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");
}
}
#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();
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.User;
factory.Password = globalSettings.MqSettings.Password;
factory.AutomaticRecoveryEnabled = true;
factory.DispatchConsumersAsync = true;
connection = factory.CreateConnection();
channel = connection.CreateModel();
var args = new Dictionary<string, object>
{
{ "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);
}
}
}