using Microsoft.Extensions.Logging; using PARR.BLL.Contracts.Interfaces; using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Text; using System.Text.Encodings.Web; using System.Text.Json; namespace PARR.BLL.Services.Implementations { internal class MqServiceV2 : IMqService { private readonly ILogger logger; private IConnection? consumerConnection; private IChannel? consumerChannel; // реализация для RabbitMQ.Client 7.1.2 private static readonly JsonSerializerOptions jsonOptions = new JsonSerializerOptions { Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping, }; public MqServiceV2(ILogger logger) { this.logger = logger; } public async Task InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler) { logger.LogInformation("Устанавливаю соединение с RabbitMQ. HostName: {HostName}, QueueName: {QueueName}", mqSettings.HostName, mqSettings.QueueName); var factory = new ConnectionFactory { HostName = mqSettings.HostName, UserName = mqSettings.User, Password = mqSettings.Password, AutomaticRecoveryEnabled = true }; try { consumerConnection = await factory.CreateConnectionAsync(); consumerChannel = await consumerConnection.CreateChannelAsync(); // кол-во сообщений которые можно обрабатывать за раз if (mqSettings.PrefetchCount.HasValue) { await consumerChannel.BasicQosAsync(0, mqSettings.PrefetchCount.Value, false); } // создаем очередь, вдруг ее еще нет await QueueDeclareAsync(consumerChannel, mqSettings); //await consumerChannel.QueueBindAsync(mqSettings.QueueName, "default", "routingKey"); consumerChannel.CallbackExceptionAsync += async (ch, ea) => { var _channel = (IChannel?)ch; // может тут если канал упал переоткрывать его logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт?: {_channel?.IsOpen}"); //if (consumerChannel.IsClosed) //{ // //todo: этот момент проверить, успешно ли он переоткроет канал и очередь будет дальше обрабатываться // await consumerChannel.CloseAsync(); // consumerChannel = await consumerConnection.CreateChannelAsync(); // logger.LogInformation(ea.Exception, $"Переоткрыл канал RabbitMq."); //} }; var consumer = new AsyncEventingBasicConsumer(consumerChannel); //consumer.HandleChannelShutdownAsync = async (c, ea) =>{ }; consumer.ReceivedAsync += async (ch, ea) => { //var content = Encoding.UTF8.GetString(ea.Body.ToArray()); var content = Encoding.UTF8.GetString(ea.Body.Span); logger.LogDebug("Получено сообщение: {content}", content); try { await messageHandler.Invoke(content); await consumerChannel.BasicAckAsync(ea.DeliveryTag, false); } catch (Exception ex) { logger.LogError(ex, $"Ошибка при обработке сообщения: {content}"); //requeue: true - сообщение возвращаем в очередь в случае ошибки. Это немного опасно, // если ошибка в формате сообщения, то она никогда не устранится и вечно будет добавляться/удаляться в очередь //todo: тут можно реализовать DLQ с кол-вом попыток, после не успешных попыток, перемещать эти сообщения в другую очередь // сделал false, выбрасывать на свежий воздух куда подальше эти кривые сообщения которые обработал с ошибками await consumerChannel.BasicNackAsync(ea.DeliveryTag, false, requeue: false); //await consumerChannel.BasicNackAsync(ea.DeliveryTag, false, requeue: true); } }; var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer); return true; } catch (Exception ex) { logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); return false; } } public async Task SendAsync(IMqSettings mqSettings, List msgObjectList) { var msgStringList = msgObjectList.Select(t => JsonSerializer.Serialize(t, jsonOptions)).ToArray(); return await SendAsync(mqSettings, msgStringList); } public async Task SendAsync(IMqSettings mqSettings, string[] msgList) { var factory = new ConnectionFactory { HostName = mqSettings.HostName, UserName = mqSettings.User, Password = mqSettings.Password, AutomaticRecoveryEnabled = true }; // индекс текущей отправки в очередь из массива msgList // нужен для формирования списка неотправленных сообщений var sendIdx = 0; try { using (var connection = await factory.CreateConnectionAsync()) using (var channel = await connection.CreateChannelAsync()) { // создаем очередь, вдруг ее еще нет await QueueDeclareAsync(channel, mqSettings); var props = new BasicProperties(); //храним на диске props.DeliveryMode = DeliveryModes.Persistent; //время жизни, мс //props.Expiration = "60000"; //props.ContentType = "text/plain";//"application/json"; props.ContentType = "text/plain; charset=utf-8";//"application/json"; channel.BasicReturnAsync += async (sender, ea) => { var body = Encoding.UTF8.GetString(ea.Body.Span); logger.LogWarning($"Сообщение не было доставлено в очередь: {ea.Exchange} -> {ea.RoutingKey}. Тело: {body}"); }; foreach (var msg in msgList) { var body = Encoding.UTF8.GetBytes(msg); await channel.BasicPublishAsync( exchange: "", routingKey: mqSettings.QueueName, //TODO: проверить mandatory: true, может указать его в false mandatory: true, basicProperties: props, body: body ); sendIdx++; logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}"); } } return new MqSendResult { IsSuccess = true }; } catch (Exception ex) { var notSendMsg = new List(); // формируем список неотправленных элементов for (int i = sendIdx; i < msgList.Length; i++) notSendMsg.Add(msgList[i]); logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}"); return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() }; } } private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings) { //https://www.rabbitmq.com/lazy-queues.html // lazy - ленивой очереди больше нет //var args = new Dictionary { { "x-queue-mode", "lazy" } }; var args = new Dictionary { }; //Объявляем очередь с которой будем работать. //Если такой очереди ещё нет, то создатся. //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. //В целом эти параметры можно посмотреть в админке RabbitMQ await channel.QueueDeclareAsync( queue: mqSettings.QueueName, durable: true, exclusive: false, autoDelete: false, arguments: args ); } public async ValueTask DisposeAsync() { if (consumerChannel != null) { if (consumerChannel.IsOpen) { await consumerChannel.CloseAsync(); } await consumerChannel.DisposeAsync(); } if (consumerConnection != null) { if (consumerConnection.IsOpen) { await consumerConnection.CloseAsync(); } await consumerConnection.DisposeAsync(); } } } }