Files
parr_api/PARR.NextRun/NextRunRabbitService.cs

104 lines
3.7 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
{
/// <summary>
/// Расчет nextRun по заданию из очереди
/// </summary>
internal class NextRunRabbitService : INextRunRabbitService
{
private readonly ILogger<NextRunRabbitService> logger;
private readonly IRabbitService mqService;
private readonly WorkerSettings workerSettings;
private readonly ITransformService transformService;
private readonly IServiceProvider serviceProvider;
public NextRunRabbitService(
ILogger<NextRunRabbitService> 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<NextRunUpdateMq>(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<INextRunUpdateService>();
await nextRunUpdateService.UpdateNextRunForTemplatesAsync(
// фильтр по JobGroupId
query => query.Where(t => t.Job!.GroupId == queryMq.JobGroupId),
() => queryMq.Initiator,
$"RabbitMq_JobGroupId_{queryMq.JobGroupId}"
);
}
//await UpdateNextRunAsync(query);
}
/// <summary>
/// Существует ли JobGroup с таким Id?
/// </summary>
/// <param name="jobGroupId"></param>
/// <returns></returns>
private async Task<bool> IsValidJobGroupAsync(Guid jobGroupId)
{
using (var scope = serviceProvider.CreateScope())
{
var jobGroupService = scope.ServiceProvider.GetRequiredService<IJobGroupRepository>();
return await jobGroupService.Get().AnyAsync(t => t.Id == jobGroupId);
}
}
}
}