180 lines
7.7 KiB
C#
180 lines
7.7 KiB
C#
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);
|
||
// }
|
||
//}
|
||
}
|
||
}
|