using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using PARR.Common.Domain; using PARR.Constants; using PARR.DAL.NextRunServices; using PARR.DAL.Services.Interfaces; using PARR.DAL.Services.Interfaces.Job; using PARR.NextRun.Settings; namespace PARR.NextRun { /// /// Расчет nextRun по заданию из очереди /// internal class NextRunRabbitService : INextRunRabbitService { private readonly ILogger logger; private readonly IMqService mqService; private readonly WorkerSettings workerSettings; private readonly ITransformService transformService; private readonly IServiceProvider serviceProvider; public NextRunRabbitService( ILogger logger, IMqService 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 query = transformService.GetModelFromJson(msg); if (query == null) return; if (!await IsValidJobGroupAsync(query.JobGroupId)) { logger.LogError("Не найдена группа работ с id: {jobGroupId}", query.JobGroupId); return; } 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); } } /// /// Обновить все NextRun в JobGroup /// /// /// private async Task UpdateNextRunAsync(NextRunUpdateMq queryMq) { using var scope = serviceProvider.CreateScope(); var templateService = scope.ServiceProvider.GetRequiredService(); var nextRunService = scope.ServiceProvider.GetRequiredService(); // Берем шаблоны только в статусе Used var templates = await templateService.Get() .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; 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); } } } }