Files
parr_api/PARR.NextRun/NextRunRabbitService.cs

180 lines
7.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.DAL.Services.Interfaces.Job;
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<IJobGroupService>();
return await jobGroupService.Get().AnyAsync(t => t.Id == jobGroupId);
}
}
///// <summary>
///// Обновить все NextRun в JobGroup
///// </summary>
///// <param name="jobGroupId"></param>
///// <returns></returns>
//private async Task UpdateNextRunAsync(NextRunUpdateMq queryMq)
//{
// using var scope = serviceProvider.CreateScope();
// var templateService = scope.ServiceProvider.GetRequiredService<ITemplateService>();
// var nextRunService = scope.ServiceProvider.GetRequiredService<INextRunService>();
// var robotConfigurationService = scope.ServiceProvider.GetRequiredService<IRobotConfigurationService>();
// // Берем шаблоны только в статусе Used
// var templates = await templateService.Get()
// .Include(t => t.RobotConfigurations)
// .Where(t =>
// t.Job!.GroupId == queryMq.JobGroupId
// && t.StatusTypeId == TemplateStatusTypeEnum.Used
// ).ToListAsync();
// logger.LogInformation("Найдено шаблонов {count} шт. в статусе Used в группе работ {jobGroupId}", templates.Count, queryMq.JobGroupId);
// if (!templates.Any())
// return;
// var updatedTemplates = 0;
// foreach (var template in templates)
// {
// var newNextRun = await nextRunService.GetNextRunForTemplateAsync(template.Id, false);
// if (newNextRun == null)
// {
// logger.LogError("При расчете nextRun для шаблона {templateId}, {templateName} вернулся null", template.Id, template.Name);
// continue;
// }
// if (template.NextRun != newNextRun)
// {
// logger.LogInformation("Обновлен nextRun для шаблона {templateId}, {templateName}, newNextRun: {newNextRun}, oldNextRun: {oldNextRun}",
// template.Id, template.Name, newNextRun, template.NextRun);
// template.LastRun = template.NextRun;
// template.NextRun = newNextRun.Value;
// //так как nextRun обновился, пробуем поставить задание на обновление
// var config = robotConfigurationService.GetFromTemplateByRobotCode(RobotsEnum.ScheduleOrder, template);
// robotConfigurationService.SetUpdateTaskStatusIfAllow(config);
// updatedTemplates++;
// }
// else
// {
// logger.LogDebug("Не требуется обновлять nextRun для шаблона {templateId}, {templateName}. Рассчитанный и исходный равны. NextRun: {NextRun}",
// template.Id, template.Name, template.NextRun);
// }
// }
// if (updatedTemplates > 0)
// {
// if (await templateService.CommitAsync(queryMq.Initiator))
// {
// logger.LogInformation("Успешно обновлены nextRun у {count} шаблонов, группа работ: {jobGroupId}", updatedTemplates, queryMq.JobGroupId);
// }
// else
// {
// logger.LogError("При сохранении nextRun для шаблонов {count} шт, произошла ошибка при сохранении в БД. JobGroupId: {jobGroupId}", updatedTemplates, queryMq.JobGroupId);
// }
// }
// else
// {
// logger.LogInformation("Для группы работ {jobGroupId}, все nextRun актуальны. Нечего обновлять.", queryMq.JobGroupId);
// }
//}
}
}