2025-08-20 14:10:22 +10:00
|
|
|
|
using Microsoft.Extensions.Logging;
|
2026-04-14 16:27:27 +10:00
|
|
|
|
using PARR.Core.Common.Interfaces.RabbitServices;
|
2026-04-14 12:01:49 +10:00
|
|
|
|
using PARR.Domain.Common.Rabbit;
|
|
|
|
|
|
using PARR.Domain.Settings;
|
2025-08-20 14:10:22 +10:00
|
|
|
|
using RabbitMQ.Client;
|
|
|
|
|
|
using RabbitMQ.Client.Events;
|
|
|
|
|
|
using System.Text;
|
2026-01-22 10:35:27 +10:00
|
|
|
|
using System.Text.Encodings.Web;
|
|
|
|
|
|
using System.Text.Json;
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
2026-04-14 12:01:49 +10:00
|
|
|
|
namespace PARR.Infrastructure.Rabbit
|
2025-08-20 14:10:22 +10:00
|
|
|
|
{
|
2026-04-14 12:01:49 +10:00
|
|
|
|
internal class RabbitService : IRabbitService
|
2025-08-20 14:10:22 +10:00
|
|
|
|
{
|
2026-04-14 12:01:49 +10:00
|
|
|
|
private readonly ILogger<RabbitService> logger;
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
|
|
|
|
|
private IConnection? consumerConnection;
|
|
|
|
|
|
private IChannel? consumerChannel;
|
|
|
|
|
|
|
|
|
|
|
|
// реализация для RabbitMQ.Client 7.1.2
|
|
|
|
|
|
|
2026-01-22 10:35:27 +10:00
|
|
|
|
private static readonly JsonSerializerOptions jsonOptions = new JsonSerializerOptions
|
|
|
|
|
|
{
|
|
|
|
|
|
Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping,
|
|
|
|
|
|
};
|
|
|
|
|
|
|
2026-04-14 12:01:49 +10:00
|
|
|
|
public RabbitService(ILogger<RabbitService> logger)
|
2025-08-20 14:10:22 +10:00
|
|
|
|
{
|
|
|
|
|
|
this.logger = logger;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public async Task<bool> InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler)
|
|
|
|
|
|
{
|
2025-12-30 11:57:15 +10:00
|
|
|
|
logger.LogInformation("Устанавливаю соединение с RabbitMQ. HostName: {HostName}, QueueName: {QueueName}", mqSettings.HostName, mqSettings.QueueName);
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
|
|
|
|
|
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) =>
|
|
|
|
|
|
{
|
2025-11-17 15:36:24 +10:00
|
|
|
|
//var content = Encoding.UTF8.GetString(ea.Body.ToArray());
|
|
|
|
|
|
var content = Encoding.UTF8.GetString(ea.Body.Span);
|
2025-12-30 11:57:15 +10:00
|
|
|
|
logger.LogDebug("Получено сообщение: {content}", content);
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
2025-11-17 15:36:24 +10:00
|
|
|
|
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);
|
|
|
|
|
|
}
|
2025-08-20 14:10:22 +10:00
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer);
|
|
|
|
|
|
|
|
|
|
|
|
return true;
|
|
|
|
|
|
}
|
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
|
{
|
|
|
|
|
|
logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}");
|
|
|
|
|
|
return false;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 16:12:52 +10:00
|
|
|
|
public async Task<RabbitSendResult> SendAsync(IMqSettings mqSettings, IEnumerable<object> msgObjectList)
|
2026-01-22 10:35:27 +10:00
|
|
|
|
{
|
|
|
|
|
|
var msgStringList = msgObjectList.Select(t => JsonSerializer.Serialize(t, jsonOptions)).ToArray();
|
|
|
|
|
|
|
|
|
|
|
|
return await SendAsync(mqSettings, msgStringList);
|
|
|
|
|
|
}
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
2026-04-14 12:01:49 +10:00
|
|
|
|
public async Task<RabbitSendResult> SendAsync(IMqSettings mqSettings, string[] msgList)
|
2025-08-20 14:10:22 +10:00
|
|
|
|
{
|
|
|
|
|
|
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;
|
|
|
|
|
|
//время жизни, мс
|
2025-08-20 16:09:44 +10:00
|
|
|
|
//props.Expiration = "60000";
|
2026-01-22 10:35:27 +10:00
|
|
|
|
//props.ContentType = "text/plain";//"application/json";
|
|
|
|
|
|
props.ContentType = "text/plain; charset=utf-8";//"application/json";
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
2025-11-17 15:36:24 +10:00
|
|
|
|
channel.BasicReturnAsync += async (sender, ea) =>
|
|
|
|
|
|
{
|
|
|
|
|
|
var body = Encoding.UTF8.GetString(ea.Body.Span);
|
|
|
|
|
|
logger.LogWarning($"Сообщение не было доставлено в очередь: {ea.Exchange} -> {ea.RoutingKey}. Тело: {body}");
|
|
|
|
|
|
};
|
|
|
|
|
|
|
2025-08-20 14:10:22 +10:00
|
|
|
|
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}");
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-04-14 12:01:49 +10:00
|
|
|
|
return new RabbitSendResult { IsSuccess = true };
|
2025-08-20 14:10:22 +10:00
|
|
|
|
}
|
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
|
{
|
|
|
|
|
|
var notSendMsg = new List<string>();
|
|
|
|
|
|
|
|
|
|
|
|
// формируем список неотправленных элементов
|
|
|
|
|
|
for (int i = sendIdx; i < msgList.Length; i++)
|
|
|
|
|
|
notSendMsg.Add(msgList[i]);
|
|
|
|
|
|
|
|
|
|
|
|
logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}");
|
|
|
|
|
|
|
2026-04-14 12:01:49 +10:00
|
|
|
|
return new RabbitSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() };
|
2025-08-20 14:10:22 +10:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2025-11-17 15:36:24 +10:00
|
|
|
|
|
2025-08-20 14:10:22 +10:00
|
|
|
|
private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings)
|
|
|
|
|
|
{
|
|
|
|
|
|
//https://www.rabbitmq.com/lazy-queues.html
|
|
|
|
|
|
// lazy - ленивой очереди больше нет
|
2025-11-17 15:36:24 +10:00
|
|
|
|
//var args = new Dictionary<string, object?> { { "x-queue-mode", "lazy" } };
|
|
|
|
|
|
var args = new Dictionary<string, object?> { };
|
2025-08-20 14:10:22 +10:00
|
|
|
|
|
|
|
|
|
|
//Объявляем очередь с которой будем работать.
|
|
|
|
|
|
//Если такой очереди ещё нет, то создатся.
|
|
|
|
|
|
//Если нет, нужно параметры типа 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();
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|