From 10bf68709b74ba36a232f0de9baf1ea3979efcd5 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Fri, 27 Feb 2026 14:05:25 +1000 Subject: [PATCH] =?UTF-8?q?feat(nextRunWorker):=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=20=D0=B5=D1=89=D0=B5=20=D0=BE=D0=B4?= =?UTF-8?q?=D0=B8=D0=BD=20backgroundWorker,=20=D0=BA=D0=BE=D1=82=D0=BE?= =?UTF-8?q?=D1=80=D1=8B=D0=B9=20=D1=81=D0=BC=D0=BE=D1=82=D1=80=D0=B8=D1=82?= =?UTF-8?q?=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D1=8C,=20=D0=B8=20?= =?UTF-8?q?=D0=BF=D0=B5=D1=80=D0=B5=D1=81=D1=87=D0=B8=D1=82=D1=8B=D0=B2?= =?UTF-8?q?=D0=B0=D0=B5=D1=82=20nextRun=20=D0=B4=D0=BB=D1=8F=20jobGroupId?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PARR.BLL/Domain/Mq/NextRunUpdateMq.cs | 10 ++ PARR.EsppTemplateSync/TemplateMQSyncer.cs | 1 - ...nManager.cs => INextRunIntervalService.cs} | 2 +- PARR.NextRun/INextRunRabbitService.cs | 8 + PARR.NextRun/NextRunInstaller.cs | 3 +- ...unManager.cs => NextRunIntervalService.cs} | 13 +- PARR.NextRun/NextRunRabbitService.cs | 163 ++++++++++++++++++ PARR.NextRun/Settings/WorkerSettings.cs | 16 +- PARR.NextRunWorker/IntervalWorker.cs | 22 +++ PARR.NextRunWorker/Program.cs | 20 ++- PARR.NextRunWorker/RabbitWorker.cs | 28 +++ PARR.NextRunWorker/Worker.cs | 19 -- .../appsettings.Development.json | 5 + PARR.NextRunWorker/appsettings.json | 11 +- PARR.Test/NextRun/NextRunTest.cs | 4 + docker-compose.next-run.yml | 27 ++- 16 files changed, 319 insertions(+), 33 deletions(-) create mode 100644 PARR.BLL/Domain/Mq/NextRunUpdateMq.cs rename PARR.NextRun/{INextRunManager.cs => INextRunIntervalService.cs} (60%) create mode 100644 PARR.NextRun/INextRunRabbitService.cs rename PARR.NextRun/{NextRunManager.cs => NextRunIntervalService.cs} (94%) create mode 100644 PARR.NextRun/NextRunRabbitService.cs create mode 100644 PARR.NextRunWorker/IntervalWorker.cs create mode 100644 PARR.NextRunWorker/RabbitWorker.cs delete mode 100644 PARR.NextRunWorker/Worker.cs diff --git a/PARR.BLL/Domain/Mq/NextRunUpdateMq.cs b/PARR.BLL/Domain/Mq/NextRunUpdateMq.cs new file mode 100644 index 00000000..c8fd050f --- /dev/null +++ b/PARR.BLL/Domain/Mq/NextRunUpdateMq.cs @@ -0,0 +1,10 @@ +namespace PARR.BLL.Domain.Mq +{ + /// + /// Запрос на обновление всех nextRun для JobGroup + /// + public class NextRunUpdateMq + { + public Guid JobGroupId { get; set; } + } +} diff --git a/PARR.EsppTemplateSync/TemplateMQSyncer.cs b/PARR.EsppTemplateSync/TemplateMQSyncer.cs index 2f0cd7e0..90a593d4 100644 --- a/PARR.EsppTemplateSync/TemplateMQSyncer.cs +++ b/PARR.EsppTemplateSync/TemplateMQSyncer.cs @@ -40,7 +40,6 @@ namespace PARR.EsppTemplateSync public async Task StartAsync() { - var isConnected = await mqService.InitConsumerAsync(globalSettings!.MqSettings!, SyncTemplateAsync); if (!isConnected) diff --git a/PARR.NextRun/INextRunManager.cs b/PARR.NextRun/INextRunIntervalService.cs similarity index 60% rename from PARR.NextRun/INextRunManager.cs rename to PARR.NextRun/INextRunIntervalService.cs index 987a3141..e62d3ca1 100644 --- a/PARR.NextRun/INextRunManager.cs +++ b/PARR.NextRun/INextRunIntervalService.cs @@ -1,6 +1,6 @@ namespace PARR.NextRun { - public interface INextRunManager + public interface INextRunIntervalService { Task StartAsync(); } diff --git a/PARR.NextRun/INextRunRabbitService.cs b/PARR.NextRun/INextRunRabbitService.cs new file mode 100644 index 00000000..1a8c7bc9 --- /dev/null +++ b/PARR.NextRun/INextRunRabbitService.cs @@ -0,0 +1,8 @@ +namespace PARR.NextRun +{ + public interface INextRunRabbitService + { + Task StartAsync(); + Task StopAsync(); + } +} diff --git a/PARR.NextRun/NextRunInstaller.cs b/PARR.NextRun/NextRunInstaller.cs index 062b4d38..5a5bb1e2 100644 --- a/PARR.NextRun/NextRunInstaller.cs +++ b/PARR.NextRun/NextRunInstaller.cs @@ -17,7 +17,8 @@ namespace PARR.NextRun configuration.GetSection(nameof(WorkerSettings)).Bind(settings); services.AddSingleton(settings); - services.AddTransient(); + services.AddTransient(); + services.AddTransient(); } diff --git a/PARR.NextRun/NextRunManager.cs b/PARR.NextRun/NextRunIntervalService.cs similarity index 94% rename from PARR.NextRun/NextRunManager.cs rename to PARR.NextRun/NextRunIntervalService.cs index 555bf983..4cccea50 100644 --- a/PARR.NextRun/NextRunManager.cs +++ b/PARR.NextRun/NextRunIntervalService.cs @@ -11,16 +11,19 @@ using PARR.NextRun.Settings; namespace PARR.NextRun { - internal class NextRunManager : INextRunManager + /// + /// Расчет nextRun по интервалу + /// + internal class NextRunIntervalService : INextRunIntervalService { private readonly WorkerSettings workerSettings; - private readonly ILogger logger; + private readonly ILogger logger; private readonly IIntervalService intervalService; private readonly IServiceProvider serviceProvider; - public NextRunManager( + public NextRunIntervalService( WorkerSettings workerSettings, - ILogger logger, + ILogger logger, IIntervalService intervalService, IServiceProvider serviceProvider ) @@ -117,6 +120,6 @@ namespace PARR.NextRun } - + } } diff --git a/PARR.NextRun/NextRunRabbitService.cs b/PARR.NextRun/NextRunRabbitService.cs new file mode 100644 index 00000000..2f3b6adf --- /dev/null +++ b/PARR.NextRun/NextRunRabbitService.cs @@ -0,0 +1,163 @@ +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.JobGroupId); + } + + /// + /// Существует ли 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(Guid jobGroupId) + { + 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 == jobGroupId + && t.StatusTypeId == TemplateStatusTypeEnum.Used + ).ToListAsync(); + + logger.LogInformation("Найдено шаблонов {count} шт. в статусе Used в группе работ {jobGroupId}", templates.Count, 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(new HistoryInitiator { InitiatorComment = "Запрос из очереди на перерасчет всех nextRun для jobGroupId", InitiatorParrComponentId = ParrComponentsEnum.NextRun })) + { + logger.LogInformation("Успешно обновлены nextRun у {count} шаблонов, группа работ: {jobGroupId}", updatedTemplates, jobGroupId); + } + else + { + logger.LogError("При сохранении nextRun для шаблонов {count} шт, произошла ошибка при сохранении в БД. JobGroupId: {jobGroupId}", updatedTemplates, jobGroupId); + } + } + else + { + logger.LogInformation("Для группы работ {jobGroupId}, все nextRun актуальны. Нечего обновлять.", jobGroupId); + } + } + } +} diff --git a/PARR.NextRun/Settings/WorkerSettings.cs b/PARR.NextRun/Settings/WorkerSettings.cs index f91e9744..b39478c0 100644 --- a/PARR.NextRun/Settings/WorkerSettings.cs +++ b/PARR.NextRun/Settings/WorkerSettings.cs @@ -1,7 +1,21 @@ -namespace PARR.NextRun.Settings +using PARR.BLL.Contracts.Interfaces; + +namespace PARR.NextRun.Settings { internal class WorkerSettings { public TimeSpan RepeatEvery { get; set; } + + public MqSettings? MqSettings { get; set; } } + + internal class MqSettings : IMqSettings + { + public string HostName { get; set; } = string.Empty; + public string QueueName { get; set; } = string.Empty; + public string User { get; set; } = string.Empty; + public string Password { get; set; } = string.Empty; + public ushort? PrefetchCount { get; set; } = 0; + } + } diff --git a/PARR.NextRunWorker/IntervalWorker.cs b/PARR.NextRunWorker/IntervalWorker.cs new file mode 100644 index 00000000..89ade3a3 --- /dev/null +++ b/PARR.NextRunWorker/IntervalWorker.cs @@ -0,0 +1,22 @@ +using PARR.NextRun; + +namespace PARR.NextRunWorker +{ + public class IntervalWorker : BackgroundService + { + private readonly INextRunIntervalService nextRunIntervalService; + private readonly ILogger logger; + + public IntervalWorker(INextRunIntervalService nextRunIntervalService, ILogger logger) + { + this.nextRunIntervalService = nextRunIntervalService; + this.logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + logger.LogInformation("Запуск {workerName}", nameof(IntervalWorker)); + await nextRunIntervalService.StartAsync(); + } + } +} diff --git a/PARR.NextRunWorker/Program.cs b/PARR.NextRunWorker/Program.cs index 46df3c1b..004ee584 100644 --- a/PARR.NextRunWorker/Program.cs +++ b/PARR.NextRunWorker/Program.cs @@ -25,7 +25,25 @@ builder.Services.InstallNextRunServices(builder.Configuration); builder.Configuration.AddNextRunConfigurations(builder.Services); builder.Services.AddNextRunSettings(builder.Configuration); -builder.Services.AddHostedService(); +// Задается через переменные окружеия +var runRabbit = builder.Configuration.GetValue("ENABLE_RABBIT", true); +var runInterval = builder.Configuration.GetValue("ENABLE_INTERVAL", true); + +if (runRabbit) +{ + builder.Services.AddHostedService(); +} + +if (runInterval) +{ + builder.Services.AddHostedService(); +} + +// Настройка краша при ошибке. Если хоть один из воркеров упадет, то падает целиком проект +builder.Services.Configure(options => +{ + options.BackgroundServiceExceptionBehavior = BackgroundServiceExceptionBehavior.StopHost; +}); var host = builder.Build(); host.Run(); diff --git a/PARR.NextRunWorker/RabbitWorker.cs b/PARR.NextRunWorker/RabbitWorker.cs new file mode 100644 index 00000000..3423eee3 --- /dev/null +++ b/PARR.NextRunWorker/RabbitWorker.cs @@ -0,0 +1,28 @@ +using PARR.NextRun; + +namespace PARR.NextRunWorker +{ + public class RabbitWorker : BackgroundService + { + private readonly INextRunRabbitService nextRunRabbitService; + private readonly ILogger logger; + + public RabbitWorker(INextRunRabbitService nextRunRabbitService, ILogger logger) + { + this.nextRunRabbitService = nextRunRabbitService; + this.logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + logger.LogInformation("Запуск {workerName}", nameof(RabbitWorker)); + await nextRunRabbitService.StopAsync(); + } + + public override Task StopAsync(CancellationToken cancellationToken) + { + nextRunRabbitService.StopAsync().Wait(); + return base.StopAsync(cancellationToken); + } + } +} diff --git a/PARR.NextRunWorker/Worker.cs b/PARR.NextRunWorker/Worker.cs deleted file mode 100644 index 68f879fc..00000000 --- a/PARR.NextRunWorker/Worker.cs +++ /dev/null @@ -1,19 +0,0 @@ -using PARR.NextRun; - -namespace PARR.NextRunWorker -{ - public class Worker : BackgroundService - { - private readonly INextRunManager nextRunManager; - - public Worker(INextRunManager nextRunManager) - { - this.nextRunManager = nextRunManager; - } - - protected override async Task ExecuteAsync(CancellationToken stoppingToken) - { - await nextRunManager.StartAsync(); - } - } -} diff --git a/PARR.NextRunWorker/appsettings.Development.json b/PARR.NextRunWorker/appsettings.Development.json index 71c186b7..d349e11e 100644 --- a/PARR.NextRunWorker/appsettings.Development.json +++ b/PARR.NextRunWorker/appsettings.Development.json @@ -25,5 +25,10 @@ } } ] + }, + "WorkerSettings": { + "MqSettings": { + "HostName": "10.99.253.216" + } } } diff --git a/PARR.NextRunWorker/appsettings.json b/PARR.NextRunWorker/appsettings.json index d3b75e8c..d4a6b900 100644 --- a/PARR.NextRunWorker/appsettings.json +++ b/PARR.NextRunWorker/appsettings.json @@ -22,7 +22,14 @@ } } }, - "WorkerSettings": { - "RepeatEvery": "0:50:00" + "WorkerSettings": { + "RepeatEvery": "0:50:00", + "MqSettings": { + "HostName": "parr-rabbitmq", + "QueueName": "parr-next-run-updater", + "User": "next_run_updater_reader", + "Password": "LJFgsifgFidfsdvi123#!23", + "PrefetchCount": 100 } + } } diff --git a/PARR.Test/NextRun/NextRunTest.cs b/PARR.Test/NextRun/NextRunTest.cs index 48137d41..e01a2850 100644 --- a/PARR.Test/NextRun/NextRunTest.cs +++ b/PARR.Test/NextRun/NextRunTest.cs @@ -37,6 +37,10 @@ namespace PARR.Test.NextRun public async Task Test() { + //var nextRun = await nextRunServiceV2.GetNextRunForTemplateAsync(Guid.Parse("0bef9892-1672-40b7-a2bd-528f9bd1ef22"), false); + + + #region тестирование INextRunServiceV2 //рапсределить diff --git a/docker-compose.next-run.yml b/docker-compose.next-run.yml index d0de9560..e88401ce 100644 --- a/docker-compose.next-run.yml +++ b/docker-compose.next-run.yml @@ -2,11 +2,13 @@ version: '3.4' # NEXT RUN services: - parr-next-run: + parr-next-run-interval: image: harbor.dvgd.rzd/parr/parr-next-run:${tag:-latest} environment: - ASPNETCORE_ENVIRONMENT=Production - TZ=Europe/Moscow + - ENABLE_RABBIT=false + - ENABLE_INTERVAL=true logging: driver: fluentd options: @@ -15,7 +17,28 @@ services: fluentd-max-retries: '30' fluentd-async: 'true' fluentd-buffer-limit: '52428800' - tag: parr.next-run.serilog + tag: parr.next-run-interval.serilog + deploy: + replicas: 1 + networks: + - parr-network + + parr-next-run-rabbit: + image: harbor.dvgd.rzd/parr/parr-next-run:${tag:-latest} + environment: + - ASPNETCORE_ENVIRONMENT=Production + - TZ=Europe/Moscow + - ENABLE_RABBIT=true + - ENABLE_INTERVAL=false + logging: + driver: fluentd + options: + fluentd-address: dvgd-efk-01.dvgd.oao.rzd:24224 + fluentd-retry-wait: '10s' + fluentd-max-retries: '30' + fluentd-async: 'true' + fluentd-buffer-limit: '52428800' + tag: parr.next-run-rabbit.serilog deploy: replicas: 1 networks: