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: