feat: Сервисы и модели RabbitMq разнесены согласно архитектуры

This commit is contained in:
Mikhail Trubnikov
2026-04-14 12:01:49 +10:00
parent 734c2e3c99
commit e892aa12d3
196 changed files with 589 additions and 1232 deletions

View File

@@ -1,11 +0,0 @@
namespace PARR.BLL.Contracts
{
/// <summary>
/// Действия при генерации шаблонов
/// </summary>
public enum GenerateTemplateActionsEnum
{
create,
deactivate
}
}

View File

@@ -1,26 +0,0 @@
namespace PARR.BLL.Contracts.Interfaces
{
/// <summary>
/// Поля настроек MQ, используется для appsettings.json
/// </summary>
public interface IMqSettings
{
string HostName { get; set; }
string QueueName { get; set; }
string User { get; set; }
string Password { get; set; }
/// <summary>
/// Кол-во сообщений которые может принимать consumer за один раз. 0 или null - без ограничений
/// </summary>
ushort? PrefetchCount { get; set; }
}
public class MqSettingsBase : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
}

View File

@@ -1,22 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, данные из АИХИТ
/// </summary>
public class AihitDataMq
{
public required string EK { get; set; }
public string? IP { get; set; }
public string? DB { get; set; }
public string? APP { get; set; }
public string? CKBS { get; set; }
public string? IB { get; set; }
public string? SI { get; set; }
public string? SM { get; set; }
public string? OS { get; set; }
public string? WorkGroup { get; set; }
public required string Status { get; set; }
public required string ResponseArea { get; set; }
public required string WorkGroupResponseArea { get; set; }
}
}

View File

@@ -1,18 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, данные из АИХИТ для любой ХП
/// </summary>
public class AihitMainDataMq
{
/// <summary>
/// ЭК
/// </summary>
public required string Name { get; set; }
/// <summary>
/// Список полей ЭК, в формате - имя поля / значение
/// </summary>
public Dictionary<string, string?> Properties { get; set; } = new Dictionary<string, string?>();
}
}

View File

@@ -1,50 +0,0 @@
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, запрос на генерацию шаблона
/// </summary>
public class GeneratorTemplateMq
{
/// <summary>
/// Id работы
/// </summary>
public Guid WorkId { get; set; }
/// <summary>
/// Id приложения
/// </summary>
public Guid ApplicationId { get; set; }
/// <summary>
/// Действие: create, deactivate
/// </summary>
public required string Action { get; set; }
/// <summary>
/// ЭК, принимает регулярное выражение. Например: *, ВРТ-*-ДВС
/// </summary>
public required string Ek { get; set; }
/// <summary>
/// Статус ЭК. Допустимые значения: 0 - все статусы, 1,2,3,4,5,6,7,9
/// </summary>
public int StatusEk { get; set; }
/// <summary>
/// Активировать шаблон
/// </summary>
public bool? IsActiveTemplate { get; set; }
/// <summary>
/// Активировать расписание
/// </summary>
public bool? IsActiveSchedule { get; set; }
/// <summary>
/// Инициатор запроса к генератору
/// </summary>
public HistoryInitiator? HistoryInitiator { get; set; }
}
}

View File

@@ -1,14 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Ответ от API RabbitMQ
/// </summary>
public class MqApiResult
{
public bool IsSuccess { get; set; }
public object? Response { get; set; }
public Exception? Exception { get; set; }
}
}

View File

@@ -1,18 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Результат отправки очереди в MQ
/// </summary>
public class MqSendResult
{
/// <summary>
/// Результат отправки
/// </summary>
public bool IsSuccess { get; set; }
/// <summary>
/// Неотправленные сообщения
/// </summary>
public string[]? NotSendMessages { get; set; }
}
}

View File

@@ -1,14 +0,0 @@
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Запрос на обновление всех nextRun для JobGroup
/// </summary>
public class NextRunUpdateMq
{
public Guid JobGroupId { get; set; }
public required HistoryInitiator Initiator { get; set; }
}
}

View File

@@ -1,10 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, запрос на изменение Наряда в ЕСПП
/// </summary>
public class OrderManageMq
{
public Guid OrderId { get; set; }
}
}

View File

@@ -1,43 +0,0 @@
using PARR.Constants;
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, активация/деактивация шаблонов/расписаний
/// </summary>
public class TemplateActivatorMq
{
/// <summary>
/// ИД объекта, JobGroupId, JobId, TemplateId
/// </summary>
public Guid Id { get; set; }
/// <summary>
/// Тип объекта
/// </summary>
public SyncTaskEntityTypeEnum EntityType { get; set; }
/// <summary>
/// Учитывать неиспользованные шаблоны
/// </summary>
public bool AllowUnused { get; set; } = false;
/// <summary>
/// Объект воздейстивя (шаблон, расписание)
/// </summary>
public TemplateActivatorEsppObjectTypesEnum EsppObject { get; set; }
/// <summary>
/// Новое состояние объекта
/// </summary>
public bool NewState { get; set; }
/// <summary>
/// Инициатор изменений
/// </summary>
public required HistoryInitiator HistoryInitiator { get; set; }
}
}

View File

@@ -1,14 +0,0 @@
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, распределить шаблоны для JobGroupId (для TemplateDistributor)
/// </summary>
public class TemplateDistributorMq
{
public Guid JobGroupId { get; set; }
public required HistoryInitiator Initiator { get; set; }
}
}

View File

@@ -1,46 +0,0 @@
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, простого создания Template
/// </summary>
public class TemplateGeneratorMq
{
/// <summary>
/// Id регламентной работы
/// </summary>
public Guid JobId { get; set; }
/// <summary>
/// Id Юнита(единицы обслуживания)/ЭК
/// </summary>
public Guid UnitId { get; set; }
/// <summary>
/// Актировать шаблон при инициализации
/// </summary>
public bool? IsActiveTemplate { get; set; }
/// <summary>
/// Актировать расписание при инициализации
/// </summary>
public bool? IsActiveSchedule { get; set; }
/// <summary>
/// Инициатор запроса к генератору
/// </summary>
public HistoryInitiator? HistoryInitiator { get; set; }
/// <summary>
/// Связанные ЭК для сгруппированного типа JobGroup
/// </summary>
public required List<Guid> UnitsInTemplate { get; set; }
/// <summary>
/// Индекс, используется в групповых шаблонах
/// </summary>
public int? Index { get; set; }
}
}

View File

@@ -1,13 +0,0 @@
using PARR.Constants;
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
public class TemplateMatcherMq
{
public Guid Id { get; set; }
public SyncTaskEntityTypeEnum EntityType { get; set; }
public TemplateMatcherActionEnum Action { get; set; }
public required HistoryInitiator Initiator { get; set; }
}
}

View File

@@ -1,13 +0,0 @@
namespace PARR.BLL.Domain.Mq
{
/// <summary>
/// Модель в MQ, генератора заданий на генерацию шаблонов
/// </summary>
public class TemplateTaskGeneratorMq
{
/// <summary>
/// Id регламентной работы
/// </summary>
public Guid JobId { get; set; }
}
}

View File

@@ -1,58 +0,0 @@
using PARR.Constants;
using PARR.Domain.Entities.Base.History;
namespace PARR.BLL.Domain.Mq
{
public class TemplateUpdaterMq
{
public Guid TemplateId { get; set; }
/// <summary>
/// Id регламентной работы
/// </summary>
public Guid JobId { get; set; }
public required string Name { get; set; }
/// <summary>
/// Актировать шаблон при инициализации
/// </summary>
public bool IsActiveTemplate { get; set; }
/// <summary>
/// Актировать расписание при инициализации
/// </summary>
public bool IsActiveSchedule { get; set; }
//public DateTimeOffset? LastRun { get; set; }
//public DateTimeOffset NextRun { get; set; }
/// <summary>
/// Это новый шаблон? (true - перемещаем существующий в новый Job) (false - шаблон остается в своем Job)
/// </summary>
public bool IsNew { get; set; }
/// <summary>
/// Id Юнита(единицы обслуживания)/ЭК
/// </summary>
public Guid UnitId { get; set; }
/// <summary>
/// Индекс, используется в групповых шаблонах
/// </summary>
public int? Index { get; set; }
public TemplateStatusTypeEnum StatusTypeId { get; set; }
/// <summary>
/// Инициатор запроса к генератору
/// </summary>
public required HistoryInitiator Initiator { get; set; }
/// <summary>
/// Связанные ЭК для сгруппированного типа JobGroup
/// </summary>
public required List<Guid> UnitsInTemplate { get; set; }
}
}

View File

@@ -19,9 +19,9 @@ namespace PARR.BLL
services.AddTransient<IFileService, FileService>();
//services.AddTransient<IMqService, MqService>();
services.AddTransient<IMqService, MqServiceV2>();
//services.AddTransient<IMqService, MqServiceV2>();
services.AddHttpClient<IMqAdminService, MqAdminService>();
//services.AddHttpClient<IMqAdminService, MqAdminService>();
services.AddTransient<ITransformService, TransformService>();
services.AddTransient<IIntervalService, IntervalService>();
services.AddTransient<ICalendarService, CalendarService>();

View File

@@ -1,7 +1,7 @@
using Microsoft.Extensions.Logging;
using PARR.BLL.Services.Interfaces;
using PARR.Common;
using PARR.Constants;
using PARR.Domain.Enums;
namespace PARR.BLL.Services.Implementations
{

View File

@@ -1,78 +0,0 @@
using Microsoft.Extensions.Logging;
using PARR.BLL.Contracts.Interfaces;
using PARR.BLL.Domain.Mq;
using PARR.BLL.Services.Interfaces;
using System.Net;
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json;
namespace PARR.BLL.Services.Implementations
{
internal class MqAdminService : IMqAdminService
{
private readonly HttpClient httpClient;
private readonly ILogger<MqAdminService> logger;
public MqAdminService(HttpClient httpClient, ILogger<MqAdminService> logger)
{
this.httpClient = httpClient;
this.logger = logger;
}
public async Task<MqApiResult> GetConnectionsAsync(IMqSettings mqSettings)
{
var url = $"http://{mqSettings.HostName}:15672/api/connections";
return await SendRequestAsync(mqSettings, url);
}
public async Task<MqApiResult> GetQueuesAsync(IMqSettings mqSettings)
{
var url = $"http://{mqSettings.HostName}:15672/api/queues";
return await SendRequestAsync(mqSettings, url);
}
private async Task<MqApiResult> SendRequestAsync(IMqSettings mqSettings, string url)
{
SetAuthorizationHeaders(mqSettings);
try
{
var result = await httpClient.GetAsync(url);
if (result.StatusCode == HttpStatusCode.OK)
{
var content = await result.Content.ReadAsStringAsync();
var json = JsonSerializer.Deserialize<object>(content);
return new MqApiResult { IsSuccess = true, Response = json };
}
else
{
logger.LogError($"Ошибка при запросе к RabbitMq, url: {url}, StatusCode: {result.StatusCode}");
return new MqApiResult { IsSuccess = false };
}
}
catch (Exception ex)
{
logger.LogError(ex, $"Ошибка при запросе к RabbitMQ, url: {url}");
return new MqApiResult { IsSuccess = false, Exception = ex };
}
}
private void SetAuthorizationHeaders(IMqSettings mqSettings)
{
httpClient.DefaultRequestHeaders.Clear();
httpClient.DefaultRequestHeaders.Authorization =
new AuthenticationHeaderValue("Basic", Convert.ToBase64String(Encoding.ASCII.GetBytes($"{mqSettings.User}:{mqSettings.Password}")));
}
}
}

View File

@@ -1,172 +0,0 @@
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;
namespace PARR.BLL.Services.Implementations
{
// RabbitMqClient 6.5.0
internal class MqService// : IMqService
{
private readonly ILogger<MqService> logger;
private IConnection? connection;
//private IModel? channel;
public MqService(ILogger<MqService> logger)
{
this.logger = logger;
}
public MqSendResult Send(IMqSettings mqSettings, string[] msgList)
{
return new MqSendResult { IsSuccess = false, NotSendMessages = new string[0] };
//var factory = new ConnectionFactory
//{
// HostName = mqSettings.HostName,
// UserName = mqSettings.User,
// Password = mqSettings.Password
//};
//// индекс текущей отправки в очередь из массива msgList
//// нужен для формирования списка неотправленных сообщений
//var sendIdx = 0;
//try
//{
// using (var connection = factory.CreateConnection())
// using (var channel = connection.CreateModel())
// {
// //https://www.rabbitmq.com/lazy-queues.html
// //Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память
// //Формируем соответствующий аргумент
// var args = new Dictionary<string, object> { { "x-queue-mode", "lazy" } };
// //Объявляем очередь с которой будем работать.
// //Если такой очереди ещё нет, то создатся.
// //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка.
// //В целом эти параметры можно посмотреть в админке RabbitMQ
// channel.QueueDeclare(
// queue: mqSettings.QueueName,
// durable: true,
// exclusive: false,
// autoDelete: false,
// arguments: args
// );
// var props = channel.CreateBasicProperties();
// //храним на диске
// props.DeliveryMode = 2;
// //время жизни, мс
// //props.Expiration = "60000";
// foreach (var msg in msgList)
// {
// var body = Encoding.UTF8.GetBytes(msg);
// channel.BasicPublish(exchange: "", routingKey: mqSettings.QueueName, basicProperties: props, body: body);
// sendIdx++;
// logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}");
// }
// }
// return new MqSendResult { IsSuccess = true };
//}
//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)}");
// return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() };
//}
}
public bool InitConsumer(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler)
{
logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}");
return false;
//var facory = new ConnectionFactory
//{
// HostName = mqSettings.HostName,
// UserName = mqSettings.User,
// Password = mqSettings.Password,
// AutomaticRecoveryEnabled = true,
// DispatchConsumersAsync = true
//};
//try
//{
// //using var connection = facory.CreateConnection();
// //using var channel = connection.CreateModel();
// connection = facory.CreateConnection();
// channel = connection.CreateModel();
// var args = new Dictionary<string, object> { { "x-queue-mode", "lazy" } };
// channel.QueueDeclare(
// queue: mqSettings.QueueName,
// durable: true,
// exclusive: false,
// autoDelete: false,
// arguments: args
// );
// logger.LogInformation($"Соединение с RabbitMQ установлено: {mqSettings.HostName}, {mqSettings.QueueName}");
// channel.CallbackException += (ch, ea) =>
// {
// var _channel = (IModel?)ch;
// // может тут если канал упал переоткрывать его
// logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт: {_channel?.IsOpen}");
// };
// var consumer = new AsyncEventingBasicConsumer(channel);
// consumer.Received += async (ch, ea) =>
// {
// var content = Encoding.UTF8.GetString(ea.Body.ToArray());
// logger.LogDebug($"Получено сообщение: {content}");
// await messageHandler.Invoke(content);
// channel.BasicAck(ea.DeliveryTag, false);
// await Task.Yield();
// };
// string consumerTag = channel.BasicConsume(mqSettings.QueueName, false, consumer);
// return true;
//}
//catch (Exception ex)
//{
// logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}");
// return false;
//}
}
public void Dispose()
{
//channel?.Close();
//connection?.Close();
//channel?.Dispose();
//connection?.Dispose();
}
}
}

View File

@@ -1,239 +0,0 @@
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<MqServiceV2> 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<MqServiceV2> logger)
{
this.logger = logger;
}
public async Task<bool> 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<MqSendResult> SendAsync(IMqSettings mqSettings, List<object> msgObjectList)
{
var msgStringList = msgObjectList.Select(t => JsonSerializer.Serialize(t, jsonOptions)).ToArray();
return await SendAsync(mqSettings, msgStringList);
}
public async Task<MqSendResult> 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<string>();
// формируем список неотправленных элементов
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<string, object?> { { "x-queue-mode", "lazy" } };
var args = new Dictionary<string, object?> { };
//Объявляем очередь с которой будем работать.
//Если такой очереди ещё нет, то создатся.
//Если нет, нужно параметры типа 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();
}
}
}
}

View File

@@ -1,4 +1,4 @@
using PARR.Constants;
using PARR.Domain.Enums;
namespace PARR.BLL.Services.Interfaces
{

View File

@@ -1,22 +0,0 @@
using PARR.BLL.Contracts.Interfaces;
using PARR.BLL.Domain.Mq;
namespace PARR.BLL.Services.Interfaces
{
public interface IMqAdminService
{
/// <summary>
/// Получить статистику соединений из RabbitMQ
/// </summary>
/// <param name="mqSettings"></param>
/// <returns></returns>
Task<MqApiResult> GetConnectionsAsync(IMqSettings mqSettings);
/// <summary>
/// Получить статистику по очередям из RabbitMQ
/// </summary>
/// <param name="mqSettings"></param>
/// <returns></returns>
Task<MqApiResult> GetQueuesAsync(IMqSettings mqSettings);
}
}

View File

@@ -1,29 +0,0 @@
using PARR.BLL.Contracts.Interfaces;
using PARR.BLL.Domain.Mq;
namespace PARR.BLL.Services.Interfaces
{
public delegate Task MqMessageHandlerDelegate(string msg);
public interface IMqService : IAsyncDisposable // IDisposable
{
Task<bool> InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler);
/// <summary>
/// Отправить сообщение в очередь используя список строк
/// !!! Избавиться от этого метода, вместо него использовать со списком объектов !!!
/// </summary>
/// <param name="mqSettings"></param>
/// <param name="msgList"></param>
/// <returns></returns>
Task<MqSendResult> SendAsync(IMqSettings mqSettings, string[] msgList);
/// <summary>
/// Отправить сообщение в очередь используя список объектов
/// </summary>
/// <param name="mqSettings"></param>
/// <param name="msgObjectList"></param>
/// <returns></returns>
Task<MqSendResult> SendAsync(IMqSettings mqSettings, List<object> msgObjectList);
}
}