using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using PARR.Core.Common.Interfaces;
using PARR.Core.Common.Interfaces.RabbitServices;
using PARR.Core.Repositories.Interfaces.JobGroupRepositories;
using PARR.Domain.Common.Rabbit.Messages;
using PARR.NextRun.Services;
using PARR.NextRun.Settings;
namespace PARR.NextRun
{
///
/// Расчет nextRun по заданию из очереди
///
internal class NextRunRabbitService : INextRunRabbitService
{
private readonly ILogger logger;
private readonly IRabbitService mqService;
private readonly WorkerSettings workerSettings;
private readonly ITransformService transformService;
private readonly IServiceProvider serviceProvider;
public NextRunRabbitService(
ILogger logger,
IRabbitService mqService,
WorkerSettings workerSettings,
ITransformService transformService,
IServiceProvider serviceProvider
)
{
this.logger = logger;
this.mqService = mqService;
this.workerSettings = workerSettings;
this.transformService = transformService;
this.serviceProvider = serviceProvider;
}
public async Task StartAsync()
{
var isConnected = await mqService.InitConsumerAsync(workerSettings.MqSettings!, HandlerAsync);
if (!isConnected)
throw new Exception("Ошибка при подключении к RabbitMq");
logger.LogInformation("Запущена проверка очереди {QueueName}.", workerSettings.MqSettings!.QueueName);
}
public async Task StopAsync()
{
await mqService.DisposeAsync();
logger.LogInformation("Соединение с очередью {QueueName} закрыто", workerSettings.MqSettings!.QueueName);
}
private async Task HandlerAsync(string msg)
{
logger.LogInformation("Получили запрос: {message}", msg);
var queryMq = transformService.GetModelFromJson(msg);
if (queryMq == null)
return;
if (!await IsValidJobGroupAsync(queryMq.JobGroupId))
{
logger.LogError("Не найдена группа работ с id: {jobGroupId}", queryMq.JobGroupId);
return;
}
using (var scope = serviceProvider.CreateScope())
{
var nextRunUpdateService = scope.ServiceProvider.GetRequiredService();
await nextRunUpdateService.UpdateNextRunForTemplatesAsync(
// фильтр по JobGroupId
query => query.Where(t => t.Job!.GroupId == queryMq.JobGroupId),
() => queryMq.Initiator,
$"RabbitMq_JobGroupId_{queryMq.JobGroupId}"
);
}
//await UpdateNextRunAsync(query);
}
///
/// Существует ли JobGroup с таким Id?
///
///
///
private async Task IsValidJobGroupAsync(Guid jobGroupId)
{
using (var scope = serviceProvider.CreateScope())
{
var jobGroupService = scope.ServiceProvider.GetRequiredService();
return await jobGroupService.Get().AnyAsync(t => t.Id == jobGroupId);
}
}
}
}