Files
parr_api/PARR.TemplateDistributor/MqTemplateDistributor.cs

75 lines
2.7 KiB
C#
Raw Permalink Normal View History

using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using PARR.Core.Common.Interfaces;
using PARR.Core.Common.Interfaces.RabbitServices;
using PARR.Domain.Common.Rabbit.Messages;
using PARR.TemplateDistributor.Services;
using PARR.TemplateDistributor.Settings;
namespace PARR.TemplateDistributor
{
internal class MqTemplateDistributor : IMqTemplateDistributor
{
private readonly MqSettings mqSettings;
private readonly IRabbitService mqService;
private readonly ILogger<MqTemplateDistributor> logger;
private readonly ITransformService transformService;
private readonly IValidatorService validatorService;
private readonly IServiceProvider serviceProvider;
public MqTemplateDistributor(MqSettings mqSettings,
IRabbitService mqService,
ILogger<MqTemplateDistributor> logger,
ITransformService transformService,
IValidatorService validatorService,
IServiceProvider serviceProvider
)
{
this.mqSettings = mqSettings;
this.mqService = mqService;
this.logger = logger;
this.transformService = transformService;
this.validatorService = validatorService;
this.serviceProvider = serviceProvider;
}
public async Task StartAsync()
{
var isConnected = await mqService.InitConsumerAsync(mqSettings, UpdateScheduleAsync);
if (!isConnected)
throw new Exception("Ошибка при подключении к RabbitMq");
}
public async Task StopAsync()
{
await mqService.DisposeAsync();
}
private async Task UpdateScheduleAsync(string msg)
{
logger.LogInformation($"Получили запрос: {msg}");
var query = transformService.GetModelFromJson<TemplateDistributorMq>(msg);
if (query == null)
return;
if (!await validatorService.IsValidJobGroupAsync(query.JobGroupId))
{
logger.LogError("Не корректные параметры группы работ {jobGroupId}. Не буду ничего делать.", query.JobGroupId);
return;
}
using (var scope = serviceProvider.CreateScope())
{
var service = scope.ServiceProvider.GetService<ITemplateDistributor>();
if (service == null)
throw new Exception($"Не найден сервис: {nameof(ITemplateDistributor)}");
await service.DistributeAsync(query);
}
}
}
}