From 85798c0956344c3a3c99740f188eb2cb0e3742c1 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Wed, 20 Aug 2025 14:10:22 +1000 Subject: [PATCH] =?UTF-8?q?feat(bll,=20all):=20=D0=BE=D0=B1=D0=BD=D0=BE?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=B0=20=D0=B1=D0=B8=D0=B1=D0=BB=D0=B8?= =?UTF-8?q?=D0=BE=D1=82=D0=B5=D0=BA=D0=B0=20RabbitMq.Client=20=D0=B4=D0=BE?= =?UTF-8?q?=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D0=B8=207.1.2.=20=D0=9E=D0=B1?= =?UTF-8?q?=D0=BD=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD=20MqServiceV2.=20=D0=92=20?= =?UTF-8?q?IMqSettings=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD?= =?UTF-8?q?=20=D0=BF=D0=B0=D1=80=D0=B0=D0=BC=D0=B5=D1=82=D1=80=20PrefetchC?= =?UTF-8?q?ount=20-=20=D0=BC=D0=B0=D0=BA=D1=81=20=D0=BA=D0=BE=D0=BB-=D0=B2?= =?UTF-8?q?=D0=BE=20=D0=BE=D0=B4=D0=BD=D0=BE=D0=B2=D1=80=D0=B5=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D0=BD=D1=8B=D1=85=20=D0=BE=D0=B1=D1=80=D0=B0=D0=B1=D0=BE?= =?UTF-8?q?=D1=82=D0=BE=D0=BA=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D0=B5?= =?UTF-8?q?=D0=B9.=20=D0=9E=D0=B1=D0=BD=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD?= =?UTF-8?q?=D1=8B=20=D1=81=D0=B2=D1=8F=D0=B7=D0=B0=D0=BD=D0=BD=D1=8B=D0=B5?= =?UTF-8?q?=20=D0=BF=D1=80=D0=BE=D0=B5=D0=BA=D1=82=D1=8B=20=D0=BF=D0=BE?= =?UTF-8?q?=D0=B4=20=D0=BD=D0=BE=D0=B2=D1=8B=D0=B9=20MqService.=20=D0=92?= =?UTF-8?q?=20=D0=BF=D1=80=D0=BE=D0=B5=D0=BA=D1=82=D0=B0=D1=85=20Aihit,=20?= =?UTF-8?q?SyncTemplate,=20SyncSchedule=20=D1=83=D1=81=D1=82=D0=B0=D0=BD?= =?UTF-8?q?=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD=20=D0=BF=D0=B0=D1=80=D0=B0=D0=BC?= =?UTF-8?q?=D0=B5=D1=82=D1=80=20=D0=B2=20=D0=BA=D0=BE=D0=BD=D1=84=D0=B8?= =?UTF-8?q?=D0=B3=20PrefetchCount=20=3D100?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PARR.AIHITLoader/AihitLoader.cs | 4 +- PARR.AIHITLoader/Settings/MqSettings.cs | 1 + PARR.AIHITMainLoader/AihitMainLoader.cs | 4 +- PARR.AIHITMainLoader/Settings/MqSettings.cs | 2 + PARR.AIHITMainSyncer/AihitMainSyncer.cs | 8 +- PARR.AIHITMainSyncer/IAihitMainSyncer.cs | 5 +- PARR.AIHITMainSyncer/Settings/MqSettings.cs | 2 + PARR.AIHITSyncer/AihitSyncer.cs | 9 +- PARR.AIHITSyncer/IAihitSyncer.cs | 4 +- PARR.AIHITSyncer/Settings/MqSettings.cs | 1 + PARR.AIHITSyncerWorker/Worker.cs | 4 +- PARR.AIHITSyncerWorker/appsettings.json | 3 +- .../Controllers/V1/DistributorController.cs | 2 +- .../V1/GeneratorTemplateController.cs | 5 +- .../V1/TemplateActivatorController.cs | 2 +- PARR.API/PARR.API.csproj | 1 - PARR.API/Settings/MqSettings.cs | 4 + PARR.BLL/Contracts/Interfaces/IMqSettings.cs | 4 + PARR.BLL/PARR.BLL.csproj | 2 +- PARR.BLL/ParrBllInstaller.cs | 5 +- .../Services/Implementations/MqService.cs | 217 +++++++++--------- .../Services/Implementations/MqServiceV2.cs | 203 ++++++++++++++++ PARR.BLL/Services/Interfaces/IMqService.cs | 6 +- PARR.EsppOrderManager/EsppOrderManager.cs | 9 +- PARR.EsppOrderManager/IEsppOrderManager.cs | 4 +- PARR.EsppOrderManager/Settings/MqSettings.cs | 1 + PARR.EsppOrderManagerWorker/Worker.cs | 4 +- PARR.EsppScheduleSync/IScheduleSyncher.cs | 4 +- PARR.EsppScheduleSync/ScheduleSyncher.cs | 37 ++- .../Settings/GlobalSettings.cs | 1 + PARR.EsppScheduleSyncWorker/Worker.cs | 5 +- PARR.EsppScheduleSyncWorker/appsettings.json | 3 +- PARR.EsppSync/SyncService.cs | 184 ++++++++------- PARR.EsppTemplateSync/ITemplateSyncer.cs | 4 +- .../Settings/GlobalSettings.cs | 1 + PARR.EsppTemplateSync/TemplateFileSyncer.cs | 13 +- PARR.EsppTemplateSync/TemplateMQSyncer.cs | 34 +-- PARR.EsppTemplateSyncWorker/Worker.cs | 4 +- PARR.EsppTemplateSyncWorker/appsettings.json | 3 +- PARR.GeneratorTemplates/GeneratorTemplate.cs | 9 +- PARR.GeneratorTemplates/IGeneratorTemplate.cs | 4 +- .../Settings/MqSettings.cs | 1 + PARR.GeneratorTemplatesWorker/Worker.cs | 4 +- PARR.JobAutoControl/JobAutoControlManager.cs | 15 +- PARR.JobAutoControl/Settings/MqSettings.cs | 2 + PARR.Master/Services/OrderMasterService.cs | 2 +- PARR.Master/Settings/MqSettings.cs | 1 + PARR.MockData/EsppDataImporter.cs | 71 +++--- PARR.MockData/PARR.MockData.csproj | 1 - .../IMqTemplateActivator.cs | 4 +- PARR.TemplateActivator/MqTemplateActivator.cs | 9 +- PARR.TemplateActivator/Settings/MqSettings.cs | 1 + PARR.TemplateActivatorWorker/Worker.cs | 4 +- .../IMqTemplateDistributor.cs | 4 +- .../MqTemplateDistributor.cs | 9 +- .../Settings/MqSettings.cs | 1 + PARR.TemplateDistributorWorker/Worker.cs | 4 +- PARR.TemplateGeneratorWorker/Worker.cs | 7 +- 58 files changed, 621 insertions(+), 341 deletions(-) create mode 100644 PARR.BLL/Services/Implementations/MqServiceV2.cs diff --git a/PARR.AIHITLoader/AihitLoader.cs b/PARR.AIHITLoader/AihitLoader.cs index 99a01ddf..1c7822b8 100644 --- a/PARR.AIHITLoader/AihitLoader.cs +++ b/PARR.AIHITLoader/AihitLoader.cs @@ -64,7 +64,7 @@ namespace PARR.AIHITLoader for (var i = 0; i < parts; i++) { var batch = preparedData.Skip(i * loaderSettings.PackageSize).Take(loaderSettings.PackageSize).ToList(); - var sendResult = mqService.Send(mqSettings, batch.ToArray()); + var sendResult = await mqService.SendAsync(mqSettings, batch.ToArray()); if (sendResult.IsSuccess) { logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count()}"); @@ -75,7 +75,7 @@ namespace PARR.AIHITLoader } } - await Task.CompletedTask; + // await Task.CompletedTask; } /// diff --git a/PARR.AIHITLoader/Settings/MqSettings.cs b/PARR.AIHITLoader/Settings/MqSettings.cs index 6f4682b4..d4b4fac7 100644 --- a/PARR.AIHITLoader/Settings/MqSettings.cs +++ b/PARR.AIHITLoader/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.AIHITLoader.Settings 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.AIHITMainLoader/AihitMainLoader.cs b/PARR.AIHITMainLoader/AihitMainLoader.cs index 2d007f5e..3fa57b1e 100644 --- a/PARR.AIHITMainLoader/AihitMainLoader.cs +++ b/PARR.AIHITMainLoader/AihitMainLoader.cs @@ -77,7 +77,7 @@ namespace PARR.AIHITMainLoader for (var i = 0; i < parts; i++) { var batch = preparedData.Skip(i * loaderSettings.PackageSize).Take(loaderSettings.PackageSize).ToList(); - var sendResult = mqService.Send(mqSettings, batch.ToArray()); + var sendResult = await mqService.SendAsync(mqSettings, batch.ToArray()); if (sendResult.IsSuccess) { logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count()}"); @@ -87,7 +87,7 @@ namespace PARR.AIHITMainLoader logger.LogError($"Ошибка при передаче данных в Rabbit, не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}"); } - await Task.CompletedTask; + // await Task.CompletedTask; } diff --git a/PARR.AIHITMainLoader/Settings/MqSettings.cs b/PARR.AIHITMainLoader/Settings/MqSettings.cs index 984822c9..e2ac1260 100644 --- a/PARR.AIHITMainLoader/Settings/MqSettings.cs +++ b/PARR.AIHITMainLoader/Settings/MqSettings.cs @@ -11,5 +11,7 @@ namespace PARR.AIHITMainLoader.Settings public string User { get; set; } = string.Empty; public string Password { get; set; } = string.Empty; + + public ushort? PrefetchCount { get; set; } = 0; } } diff --git a/PARR.AIHITMainSyncer/AihitMainSyncer.cs b/PARR.AIHITMainSyncer/AihitMainSyncer.cs index 71c163fa..dc0531f1 100644 --- a/PARR.AIHITMainSyncer/AihitMainSyncer.cs +++ b/PARR.AIHITMainSyncer/AihitMainSyncer.cs @@ -25,11 +25,11 @@ namespace PARR.AIHITMainSyncer this.mqSettings = mqSettings; this.serviceProvider = serviceProvider; } - public void Start() + public async Task StartAsync() { logger.LogInformation("Запуск сервиса синхронизации данных из АИХИТ и ПАРР."); - var isConnected = mqService.InitConsumer(mqSettings, async (string msg) => + var isConnected = await mqService.InitConsumerAsync(mqSettings, async (string msg) => { using (var scope = serviceProvider.CreateScope()) { @@ -47,9 +47,9 @@ namespace PARR.AIHITMainSyncer } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } } } diff --git a/PARR.AIHITMainSyncer/IAihitMainSyncer.cs b/PARR.AIHITMainSyncer/IAihitMainSyncer.cs index 5b975fd9..a0ce02c9 100644 --- a/PARR.AIHITMainSyncer/IAihitMainSyncer.cs +++ b/PARR.AIHITMainSyncer/IAihitMainSyncer.cs @@ -2,9 +2,8 @@ { public interface IAihitMainSyncer { - void Start(); + Task StartAsync(); - - void Stop(); + Task StopAsync(); } } diff --git a/PARR.AIHITMainSyncer/Settings/MqSettings.cs b/PARR.AIHITMainSyncer/Settings/MqSettings.cs index 3e0c3ad9..01a22e0b 100644 --- a/PARR.AIHITMainSyncer/Settings/MqSettings.cs +++ b/PARR.AIHITMainSyncer/Settings/MqSettings.cs @@ -11,5 +11,7 @@ namespace PARR.AIHITMainSyncer.Settings public string User { get; set; } = string.Empty; public string Password { get; set; } = string.Empty; + + public ushort? PrefetchCount { get; set; } = 0; } } diff --git a/PARR.AIHITSyncer/AihitSyncer.cs b/PARR.AIHITSyncer/AihitSyncer.cs index 8bfe42d2..9ce7adaa 100644 --- a/PARR.AIHITSyncer/AihitSyncer.cs +++ b/PARR.AIHITSyncer/AihitSyncer.cs @@ -3,6 +3,7 @@ using Microsoft.Extensions.Logging; using PARR.AIHITSyncer.Services; using PARR.AIHITSyncer.Settings; using PARR.BLL.Services.Interfaces; +using System.Threading.Tasks; namespace PARR.AIHITSyncer { @@ -27,12 +28,12 @@ namespace PARR.AIHITSyncer this.serviceProvider = serviceProvider; } - public void Start() + public async Task StartAsync() { logger.LogInformation("Запуск сервиса синхронизации данных из АИХИТ и ПАРР."); //var isConnected = mqService.InitConsumer(mqSettings, SyncDataAsync); - var isConnected = mqService.InitConsumer(mqSettings, async (string msg) => + var isConnected = await mqService.InitConsumerAsync(mqSettings, async (string msg) => { using (var scope = serviceProvider.CreateScope()) { @@ -48,9 +49,9 @@ namespace PARR.AIHITSyncer if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } } diff --git a/PARR.AIHITSyncer/IAihitSyncer.cs b/PARR.AIHITSyncer/IAihitSyncer.cs index 340a8178..dd42d05c 100644 --- a/PARR.AIHITSyncer/IAihitSyncer.cs +++ b/PARR.AIHITSyncer/IAihitSyncer.cs @@ -2,7 +2,7 @@ { public interface IAihitSyncer { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.AIHITSyncer/Settings/MqSettings.cs b/PARR.AIHITSyncer/Settings/MqSettings.cs index f0c697cb..66150547 100644 --- a/PARR.AIHITSyncer/Settings/MqSettings.cs +++ b/PARR.AIHITSyncer/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.AIHITSyncer.Settings 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.AIHITSyncerWorker/Worker.cs b/PARR.AIHITSyncerWorker/Worker.cs index ebdb869e..d119b0ea 100644 --- a/PARR.AIHITSyncerWorker/Worker.cs +++ b/PARR.AIHITSyncerWorker/Worker.cs @@ -16,13 +16,13 @@ namespace PARR.AIHITSyncerWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - aihitSyncer.Start(); + await aihitSyncer.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - aihitSyncer.Stop(); + aihitSyncer.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.AIHITSyncerWorker/appsettings.json b/PARR.AIHITSyncerWorker/appsettings.json index 70444c16..5ec258f6 100644 --- a/PARR.AIHITSyncerWorker/appsettings.json +++ b/PARR.AIHITSyncerWorker/appsettings.json @@ -33,6 +33,7 @@ "HostName": "10.99.253.216", "QueueName": "parr-aihit-data", "User": "aihit_syncer", - "Password": "afsdgDGjhfiU76496as42!@" + "Password": "afsdgDGjhfiU76496as42!@", + "PrefetchCount": 100 } } diff --git a/PARR.API/Controllers/V1/DistributorController.cs b/PARR.API/Controllers/V1/DistributorController.cs index 0a529956..f102f938 100644 --- a/PARR.API/Controllers/V1/DistributorController.cs +++ b/PARR.API/Controllers/V1/DistributorController.cs @@ -53,7 +53,7 @@ namespace PARR.API.Controllers.V1 var msg = JsonSerializer.Serialize(requestToMq); - var sendResult = mqService.Send(mqSettings.TemplateDistributor, new[] { msg }); + var sendResult = await mqService.SendAsync(mqSettings.TemplateDistributor, new[] { msg }); if (sendResult.IsSuccess) return Created("", new Response(null, true, new List(), "Отправлен запрос на перераспределение регламентных работ.")); diff --git a/PARR.API/Controllers/V1/GeneratorTemplateController.cs b/PARR.API/Controllers/V1/GeneratorTemplateController.cs index a2be4f98..4a6540a8 100644 --- a/PARR.API/Controllers/V1/GeneratorTemplateController.cs +++ b/PARR.API/Controllers/V1/GeneratorTemplateController.cs @@ -54,7 +54,8 @@ namespace PARR.API.Controllers.V1 request.Ek.ForEach(ek => { - request.Status.ForEach(status => + //todo: проверить в ForEach(async status - отрабатывает ли нормально async, возможно надо переделать в обычный foreach + request.Status.ForEach(async status => { var obj = new GeneratorTemplateMq { @@ -70,7 +71,7 @@ namespace PARR.API.Controllers.V1 var msg = JsonSerializer.Serialize(obj); - var _result = mqService.Send(mqSettings.GenerateTemplates, new[] { msg }); + var _result = await mqService.SendAsync(mqSettings.GenerateTemplates, new[] { msg }); if (_result.IsSuccess == false) result = false; diff --git a/PARR.API/Controllers/V1/TemplateActivatorController.cs b/PARR.API/Controllers/V1/TemplateActivatorController.cs index 6e3f6407..79ebacf6 100644 --- a/PARR.API/Controllers/V1/TemplateActivatorController.cs +++ b/PARR.API/Controllers/V1/TemplateActivatorController.cs @@ -74,7 +74,7 @@ namespace PARR.API.Controllers.V1 var msg = JsonSerializer.Serialize(mqRequest); - var sendResult = mqService.Send(mqSettings.TemplateActivator, new[] { msg }); + var sendResult = await mqService.SendAsync(mqSettings.TemplateActivator, new[] { msg }); if (!sendResult.IsSuccess) return BadRequest(new Response(false, new List { new ErrorModel { Message = $"Ошибка при отправке данных." } })); diff --git a/PARR.API/PARR.API.csproj b/PARR.API/PARR.API.csproj index 56cb3651..51b89c58 100644 --- a/PARR.API/PARR.API.csproj +++ b/PARR.API/PARR.API.csproj @@ -26,7 +26,6 @@ all runtime; build; native; contentfiles; analyzers; buildtransitive - diff --git a/PARR.API/Settings/MqSettings.cs b/PARR.API/Settings/MqSettings.cs index 228178cb..971596b0 100644 --- a/PARR.API/Settings/MqSettings.cs +++ b/PARR.API/Settings/MqSettings.cs @@ -16,6 +16,7 @@ namespace PARR.API.Settings 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; } public class MqTemplateDistributor : IMqSettings @@ -24,6 +25,7 @@ namespace PARR.API.Settings 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; } public class MqTemplateActivator : IMqSettings @@ -32,6 +34,7 @@ namespace PARR.API.Settings 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; } public class MqStatistics @@ -45,5 +48,6 @@ namespace PARR.API.Settings 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.BLL/Contracts/Interfaces/IMqSettings.cs b/PARR.BLL/Contracts/Interfaces/IMqSettings.cs index ef0536c9..aa6f9aa1 100644 --- a/PARR.BLL/Contracts/Interfaces/IMqSettings.cs +++ b/PARR.BLL/Contracts/Interfaces/IMqSettings.cs @@ -9,5 +9,9 @@ string QueueName { get; set; } string User { get; set; } string Password { get; set; } + /// + /// Кол-во сообщений которые может принимать consumer за один раз. 0 или null - без ограничений + /// + ushort? PrefetchCount { get; set; } } } diff --git a/PARR.BLL/PARR.BLL.csproj b/PARR.BLL/PARR.BLL.csproj index f5b59d30..719f7fcf 100644 --- a/PARR.BLL/PARR.BLL.csproj +++ b/PARR.BLL/PARR.BLL.csproj @@ -12,7 +12,7 @@ - + diff --git a/PARR.BLL/ParrBllInstaller.cs b/PARR.BLL/ParrBllInstaller.cs index 358827d9..79f992eb 100644 --- a/PARR.BLL/ParrBllInstaller.cs +++ b/PARR.BLL/ParrBllInstaller.cs @@ -17,7 +17,10 @@ namespace PARR.BLL services.AddTransient(); - services.AddTransient(); + + //services.AddTransient(); + services.AddTransient(); + services.AddHttpClient(); services.AddTransient(); services.AddTransient(); diff --git a/PARR.BLL/Services/Implementations/MqService.cs b/PARR.BLL/Services/Implementations/MqService.cs index 3f00dfe7..2b9a4be7 100644 --- a/PARR.BLL/Services/Implementations/MqService.cs +++ b/PARR.BLL/Services/Implementations/MqService.cs @@ -4,18 +4,19 @@ using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using RabbitMQ.Client; using RabbitMQ.Client.Events; -using System.Runtime; using System.Text; namespace PARR.BLL.Services.Implementations { - internal class MqService : IMqService + // RabbitMqClient 6.5.0 + + internal class MqService// : IMqService { private readonly ILogger logger; private IConnection? connection; - private IModel? channel; + //private IModel? channel; public MqService(ILogger logger) { @@ -24,144 +25,148 @@ namespace PARR.BLL.Services.Implementations public MqSendResult Send(IMqSettings mqSettings, string[] msgList) { - var factory = new ConnectionFactory - { - HostName = mqSettings.HostName, - UserName = mqSettings.User, - Password = mqSettings.Password - }; + return new MqSendResult { IsSuccess = false, NotSendMessages = new string[0] }; - // индекс текущей отправки в очередь из массива msgList - // нужен для формирования списка неотправленных сообщений - var sendIdx = 0; + //var factory = new ConnectionFactory + //{ + // HostName = mqSettings.HostName, + // UserName = mqSettings.User, + // Password = mqSettings.Password + //}; - try - { - using (var connection = factory.CreateConnection()) - using (var channel = connection.CreateModel()) - { - //https://www.rabbitmq.com/lazy-queues.html - //Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память - //Формируем соответствующий аргумент - var args = new Dictionary { { "x-queue-mode", "lazy" } }; + //// индекс текущей отправки в очередь из массива msgList + //// нужен для формирования списка неотправленных сообщений + //var sendIdx = 0; - //Объявляем очередь с которой будем работать. - //Если такой очереди ещё нет, то создатся. - //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. - //В целом эти параметры можно посмотреть в админке RabbitMQ + //try + //{ + // using (var connection = factory.CreateConnection()) + // using (var channel = connection.CreateModel()) + // { + // //https://www.rabbitmq.com/lazy-queues.html + // //Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память + // //Формируем соответствующий аргумент + // var args = new Dictionary { { "x-queue-mode", "lazy" } }; - channel.QueueDeclare( - queue: mqSettings.QueueName, - durable: true, - exclusive: false, - autoDelete: false, - arguments: args - ); + // //Объявляем очередь с которой будем работать. + // //Если такой очереди ещё нет, то создатся. + // //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. + // //В целом эти параметры можно посмотреть в админке RabbitMQ - var props = channel.CreateBasicProperties(); - //храним на диске - props.DeliveryMode = 2; - //время жизни, мс - //props.Expiration = "60000"; + // channel.QueueDeclare( + // queue: mqSettings.QueueName, + // durable: true, + // exclusive: false, + // autoDelete: false, + // arguments: args + // ); - foreach (var msg in msgList) - { - var body = Encoding.UTF8.GetBytes(msg); - channel.BasicPublish(exchange: "", routingKey: mqSettings.QueueName, basicProperties: props, body: body); + // var props = channel.CreateBasicProperties(); + // //храним на диске + // props.DeliveryMode = 2; + // //время жизни, мс + // //props.Expiration = "60000"; - sendIdx++; - logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}"); - } - } + // foreach (var msg in msgList) + // { + // var body = Encoding.UTF8.GetBytes(msg); + // channel.BasicPublish(exchange: "", routingKey: mqSettings.QueueName, basicProperties: props, body: body); - return new MqSendResult { IsSuccess = true }; - } - catch (Exception ex) - { - var notSendMsg = new List(); + // sendIdx++; + // logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}"); + // } + // } - // формируем список неотправленных элементов - for (int i = sendIdx; i < msgList.Length; i++) - notSendMsg.Add(msgList[i]); + // return new MqSendResult { IsSuccess = true }; + //} + //catch (Exception ex) + //{ + // var notSendMsg = new List(); - logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}"); + // // формируем список неотправленных элементов + // for (int i = sendIdx; i < msgList.Length; i++) + // notSendMsg.Add(msgList[i]); - return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() }; - } + // logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}"); + + // return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() }; + //} } public bool InitConsumer(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler) { logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}"); - var facory = new ConnectionFactory - { - HostName = mqSettings.HostName, - UserName = mqSettings.User, - Password = mqSettings.Password, - AutomaticRecoveryEnabled = true, - DispatchConsumersAsync = true - }; + return false; - try - { - //using var connection = facory.CreateConnection(); - //using var channel = connection.CreateModel(); + //var facory = new ConnectionFactory + //{ + // HostName = mqSettings.HostName, + // UserName = mqSettings.User, + // Password = mqSettings.Password, + // AutomaticRecoveryEnabled = true, + // DispatchConsumersAsync = true + //}; - connection = facory.CreateConnection(); - channel = connection.CreateModel(); + //try + //{ + // //using var connection = facory.CreateConnection(); + // //using var channel = connection.CreateModel(); - var args = new Dictionary { { "x-queue-mode", "lazy" } }; + // connection = facory.CreateConnection(); + // channel = connection.CreateModel(); - channel.QueueDeclare( - queue: mqSettings.QueueName, - durable: true, - exclusive: false, - autoDelete: false, - arguments: args - ); + // var args = new Dictionary { { "x-queue-mode", "lazy" } }; - logger.LogInformation($"Соединение с RabbitMQ установлено: {mqSettings.HostName}, {mqSettings.QueueName}"); + // channel.QueueDeclare( + // queue: mqSettings.QueueName, + // durable: true, + // exclusive: false, + // autoDelete: false, + // arguments: args + // ); - channel.CallbackException += (ch, ea) => - { - var _channel = (IModel?)ch; - // может тут если канал упал переоткрывать его - logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт: {_channel?.IsOpen}"); - }; + // logger.LogInformation($"Соединение с RabbitMQ установлено: {mqSettings.HostName}, {mqSettings.QueueName}"); - var consumer = new AsyncEventingBasicConsumer(channel); - consumer.Received += async (ch, ea) => - { - var content = Encoding.UTF8.GetString(ea.Body.ToArray()); - logger.LogDebug($"Получено сообщение: {content}"); + // channel.CallbackException += (ch, ea) => + // { + // var _channel = (IModel?)ch; + // // может тут если канал упал переоткрывать его + // logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт: {_channel?.IsOpen}"); + // }; - await messageHandler.Invoke(content); + // var consumer = new AsyncEventingBasicConsumer(channel); + // consumer.Received += async (ch, ea) => + // { + // var content = Encoding.UTF8.GetString(ea.Body.ToArray()); + // logger.LogDebug($"Получено сообщение: {content}"); - channel.BasicAck(ea.DeliveryTag, false); - await Task.Yield(); - }; + // await messageHandler.Invoke(content); - string consumerTag = channel.BasicConsume(mqSettings.QueueName, false, consumer); + // channel.BasicAck(ea.DeliveryTag, false); + // await Task.Yield(); + // }; - return true; + // string consumerTag = channel.BasicConsume(mqSettings.QueueName, false, consumer); - } - catch (Exception ex) - { - logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); + // return true; - return false; - } + //} + //catch (Exception ex) + //{ + // logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); + + // return false; + //} } public void Dispose() { - channel?.Close(); - connection?.Close(); + //channel?.Close(); + //connection?.Close(); - channel?.Dispose(); - connection?.Dispose(); + //channel?.Dispose(); + //connection?.Dispose(); } } } diff --git a/PARR.BLL/Services/Implementations/MqServiceV2.cs b/PARR.BLL/Services/Implementations/MqServiceV2.cs new file mode 100644 index 00000000..e7d863c4 --- /dev/null +++ b/PARR.BLL/Services/Implementations/MqServiceV2.cs @@ -0,0 +1,203 @@ +using Microsoft.Extensions.Logging; +using PARR.BLL.Contracts.Interfaces; +using PARR.BLL.Domain.Mq; +using PARR.BLL.Services.Interfaces; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using System.Text; + +namespace PARR.BLL.Services.Implementations +{ + internal class MqServiceV2 : IMqService + { + private readonly ILogger logger; + + private IConnection? consumerConnection; + private IChannel? consumerChannel; + + // реализация для RabbitMQ.Client 7.1.2 + + public MqServiceV2(ILogger logger) + { + this.logger = logger; + } + + + public async Task InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler) + { + logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}"); + + var factory = new ConnectionFactory + { + HostName = mqSettings.HostName, + UserName = mqSettings.User, + Password = mqSettings.Password, + AutomaticRecoveryEnabled = true + }; + + try + { + consumerConnection = await factory.CreateConnectionAsync(); + consumerChannel = await consumerConnection.CreateChannelAsync(); + + // кол-во сообщений которые можно обрабатывать за раз + if (mqSettings.PrefetchCount.HasValue) + { + await consumerChannel.BasicQosAsync(0, mqSettings.PrefetchCount.Value, false); + } + + + // создаем очередь, вдруг ее еще нет + await QueueDeclareAsync(consumerChannel, mqSettings); + + //await consumerChannel.QueueBindAsync(mqSettings.QueueName, "default", "routingKey"); + + + consumerChannel.CallbackExceptionAsync += async (ch, ea) => + { + var _channel = (IChannel?)ch; + // может тут если канал упал переоткрывать его + + logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт?: {_channel?.IsOpen}"); + + //if (consumerChannel.IsClosed) + //{ + // //todo: этот момент проверить, успешно ли он переоткроет канал и очередь будет дальше обрабатываться + // await consumerChannel.CloseAsync(); + // consumerChannel = await consumerConnection.CreateChannelAsync(); + + // logger.LogInformation(ea.Exception, $"Переоткрыл канал RabbitMq."); + //} + }; + + var consumer = new AsyncEventingBasicConsumer(consumerChannel); + //consumer.HandleChannelShutdownAsync = async (c, ea) =>{ }; + + consumer.ReceivedAsync += async (ch, ea) => + { + var content = Encoding.UTF8.GetString(ea.Body.ToArray()); + logger.LogDebug($"Получено сообщение: {content}"); + + await messageHandler.Invoke(content); + + await consumerChannel.BasicAckAsync(ea.DeliveryTag, false); + }; + + var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer); + + return true; + } + catch (Exception ex) + { + logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); + return false; + } + } + + + public async Task SendAsync(IMqSettings mqSettings, string[] msgList) + { + var factory = new ConnectionFactory + { + HostName = mqSettings.HostName, + UserName = mqSettings.User, + Password = mqSettings.Password, + AutomaticRecoveryEnabled = true + }; + + // индекс текущей отправки в очередь из массива msgList + // нужен для формирования списка неотправленных сообщений + var sendIdx = 0; + + try + { + using (var connection = await factory.CreateConnectionAsync()) + using (var channel = await connection.CreateChannelAsync()) + { + // создаем очередь, вдруг ее еще нет + await QueueDeclareAsync(channel, mqSettings); + + var props = new BasicProperties(); + //храним на диске + props.DeliveryMode = DeliveryModes.Persistent; + //время жизни, мс + props.Expiration = "60000"; + props.ContentType = "text/plain";//"application/json"; + + foreach (var msg in msgList) + { + var body = Encoding.UTF8.GetBytes(msg); + await channel.BasicPublishAsync( + exchange: "", + routingKey: mqSettings.QueueName, + //TODO: проверить mandatory: true, может указать его в false + mandatory: true, + basicProperties: props, + body: body + ); + + sendIdx++; + logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}"); + } + } + + return new MqSendResult { IsSuccess = true }; + } + catch (Exception ex) + { + var notSendMsg = new List(); + + // формируем список неотправленных элементов + for (int i = sendIdx; i < msgList.Length; i++) + notSendMsg.Add(msgList[i]); + + logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}"); + + return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() }; + } + } + + private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings) + { + //https://www.rabbitmq.com/lazy-queues.html + // lazy - ленивой очереди больше нет + var args = new Dictionary { { "x-queue-mode", "lazy" } }; + + //Объявляем очередь с которой будем работать. + //Если такой очереди ещё нет, то создатся. + //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. + //В целом эти параметры можно посмотреть в админке RabbitMQ + await channel.QueueDeclareAsync( + queue: mqSettings.QueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: args + ); + } + + + public async ValueTask DisposeAsync() + { + if (consumerChannel != null) + { + if (consumerChannel.IsOpen) + { + await consumerChannel.CloseAsync(); + } + + await consumerChannel.DisposeAsync(); + } + + if (consumerConnection != null) + { + if (consumerConnection.IsOpen) + { + await consumerConnection.CloseAsync(); + } + + await consumerConnection.DisposeAsync(); + } + } + } +} diff --git a/PARR.BLL/Services/Interfaces/IMqService.cs b/PARR.BLL/Services/Interfaces/IMqService.cs index d96bfb32..52344182 100644 --- a/PARR.BLL/Services/Interfaces/IMqService.cs +++ b/PARR.BLL/Services/Interfaces/IMqService.cs @@ -5,9 +5,9 @@ namespace PARR.BLL.Services.Interfaces { public delegate Task MqMessageHandlerDelegate(string msg); - public interface IMqService : IDisposable + public interface IMqService : IAsyncDisposable // IDisposable { - bool InitConsumer(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler); - MqSendResult Send(IMqSettings mqSettings, string[] msgList); + Task InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler); + Task SendAsync(IMqSettings mqSettings, string[] msgList); } } diff --git a/PARR.EsppOrderManager/EsppOrderManager.cs b/PARR.EsppOrderManager/EsppOrderManager.cs index 46ecc3a3..acf74eb6 100644 --- a/PARR.EsppOrderManager/EsppOrderManager.cs +++ b/PARR.EsppOrderManager/EsppOrderManager.cs @@ -8,6 +8,7 @@ using PARR.DAL.Models; using PARR.DAL.Services.Interfaces; using PARR.EsppOrderManager.Services; using PARR.EsppOrderManager.Settings; +using System.Threading.Tasks; namespace PARR.EsppOrderManager { @@ -40,11 +41,11 @@ namespace PARR.EsppOrderManager this.resultTemplateSettings = resultTemplateSettings; } - public void Start() + public async Task StartAsync() { // логика работы описана в документации - var isConnected = mqService.InitConsumer(mqSettings, ManageOrderAsync); + var isConnected = await mqService.InitConsumerAsync(mqSettings, ManageOrderAsync); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); @@ -53,9 +54,9 @@ namespace PARR.EsppOrderManager //await ManageOrderAsync("{\"OrderId\":\"49191af5-c313-4388-9720-a61fd13f15a8\"}"); } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } diff --git a/PARR.EsppOrderManager/IEsppOrderManager.cs b/PARR.EsppOrderManager/IEsppOrderManager.cs index bb4a1f5a..90a0a572 100644 --- a/PARR.EsppOrderManager/IEsppOrderManager.cs +++ b/PARR.EsppOrderManager/IEsppOrderManager.cs @@ -2,7 +2,7 @@ { public interface IEsppOrderManager { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.EsppOrderManager/Settings/MqSettings.cs b/PARR.EsppOrderManager/Settings/MqSettings.cs index abb27739..1a8c55bb 100644 --- a/PARR.EsppOrderManager/Settings/MqSettings.cs +++ b/PARR.EsppOrderManager/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.EsppOrderManager.Settings 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.EsppOrderManagerWorker/Worker.cs b/PARR.EsppOrderManagerWorker/Worker.cs index aa43fd05..3f17ce76 100644 --- a/PARR.EsppOrderManagerWorker/Worker.cs +++ b/PARR.EsppOrderManagerWorker/Worker.cs @@ -15,12 +15,12 @@ namespace PARR.EsppOrderManagerWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - esppOrderManager.Start(); + await esppOrderManager.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - esppOrderManager.Stop(); + esppOrderManager.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.EsppScheduleSync/IScheduleSyncher.cs b/PARR.EsppScheduleSync/IScheduleSyncher.cs index 23ef722d..76d8afda 100644 --- a/PARR.EsppScheduleSync/IScheduleSyncher.cs +++ b/PARR.EsppScheduleSync/IScheduleSyncher.cs @@ -2,7 +2,7 @@ { public interface IScheduleSyncher { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.EsppScheduleSync/ScheduleSyncher.cs b/PARR.EsppScheduleSync/ScheduleSyncher.cs index bfa1c7ac..6a499c98 100644 --- a/PARR.EsppScheduleSync/ScheduleSyncher.cs +++ b/PARR.EsppScheduleSync/ScheduleSyncher.cs @@ -20,7 +20,7 @@ namespace PARR.EsppScheduleSync private readonly ISyncService syncService; private readonly SettingsFromDb settingsFromDb; private readonly IServiceProvider serviceProvider; - private readonly INextRunModifierService nextRunModifierService; + //private readonly INextRunModifierService nextRunModifierService; public ScheduleSyncher( ILogger logger, @@ -28,8 +28,8 @@ namespace PARR.EsppScheduleSync IMqService mqService, ISyncService syncService, SettingsFromDb settingsFromDb, - IServiceProvider serviceProvider, - INextRunModifierService nextRunModifierService + IServiceProvider serviceProvider + //INextRunModifierService nextRunModifierService ) { this.logger = logger; @@ -38,7 +38,7 @@ namespace PARR.EsppScheduleSync this.syncService = syncService; this.settingsFromDb = settingsFromDb; this.serviceProvider = serviceProvider; - this.nextRunModifierService = nextRunModifierService; + //this.nextRunModifierService = nextRunModifierService; if (globalSettings.MqSettings == null) { logger.LogError("Нет секции настроек хранилища. MqSettings, EsppTemplates"); @@ -47,9 +47,9 @@ namespace PARR.EsppScheduleSync } - public void Start() + public async Task StartAsync() { - var isConnected = mqService.InitConsumer(globalSettings!.MqSettings!, SyncScheduleAsync); + var isConnected = await mqService.InitConsumerAsync(globalSettings!.MqSettings!, SyncScheduleAsync); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); @@ -58,9 +58,9 @@ namespace PARR.EsppScheduleSync } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); logger.LogInformation($"=== === === Соединение с очередью {globalSettings.MqSettings!.QueueName} закрыто === === ==="); } @@ -80,6 +80,19 @@ namespace PARR.EsppScheduleSync /// private EsppObjectSchedule ConvertDbObjToEsppObj(Template template) { + //дефолтное значение, изменится в сервисе nextRunModifierService + var nextRunByAccountRobotTimeZone = DateTimeOffset.MinValue; + + using (var scope = serviceProvider.CreateScope()) + { + var nextRunModifierService = scope.ServiceProvider.GetService(); + if (nextRunModifierService == null) + throw new Exception($"Не найден сервис: {nameof(INextRunModifierService)}"); + + nextRunByAccountRobotTimeZone = nextRunModifierService.GetNextRunByAccountRobotTimeZone(template.NextRun); + } + + var esppObjectFromDb = new EsppObjectSchedule { TemplateName = template.Name, @@ -92,8 +105,10 @@ namespace PARR.EsppScheduleSync //Мы решили, что для всех расписаний "Нет исключений", если что-то поменяется, тут нужно переделать TypeV60calendar = settingsFromDb.ScheduleExclude == "Нет исключений" ? "NONE" : "", //Scheduled = EsppScheduleHelpers.GetNextRun(template.NextRun), - Scheduled = EsppScheduleHelpers.GetNextRun(nextRunModifierService.GetNextRunByAccountRobotTimeZone(template.NextRun)), - BasisTime = EsppScheduleHelpers.GetGenerationTime(nextRunModifierService.GetNextRunByAccountRobotTimeZone(template.NextRun)), + //Scheduled = EsppScheduleHelpers.GetNextRun(nextRunModifierService.GetNextRunByAccountRobotTimeZone(template.NextRun)), + //BasisTime = EsppScheduleHelpers.GetGenerationTime(nextRunModifierService.GetNextRunByAccountRobotTimeZone(template.NextRun)), + Scheduled = EsppScheduleHelpers.GetNextRun(nextRunByAccountRobotTimeZone), + BasisTime = EsppScheduleHelpers.GetGenerationTime(nextRunByAccountRobotTimeZone), Timezone = settingsFromDb.ScheduleTimezone, //Мы решили, что для всех расписаний "Отсутствует дата завершения", если что-то поменяется, тут нужно переделать TerminationType = settingsFromDb.ScheduleRepeatRange == "Отсутствует дата завершения" ? "forever" : "", @@ -151,7 +166,7 @@ namespace PARR.EsppScheduleSync //} //else //{ - esppObject.Dayofmonth = esppSchedule.Values.First(t => t.Order == 0).Value.EsppExportValue; + esppObject.Dayofmonth = esppSchedule.Values.First(t => t.Order == 0).Value.EsppExportValue; //} break; case EsppSchTypeScheduleEnum.Monthly2: diff --git a/PARR.EsppScheduleSync/Settings/GlobalSettings.cs b/PARR.EsppScheduleSync/Settings/GlobalSettings.cs index a22e87a5..ab1dda14 100644 --- a/PARR.EsppScheduleSync/Settings/GlobalSettings.cs +++ b/PARR.EsppScheduleSync/Settings/GlobalSettings.cs @@ -14,5 +14,6 @@ namespace PARR.EsppScheduleSync.Settings 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.EsppScheduleSyncWorker/Worker.cs b/PARR.EsppScheduleSyncWorker/Worker.cs index 28e88dfb..e6fd9794 100644 --- a/PARR.EsppScheduleSyncWorker/Worker.cs +++ b/PARR.EsppScheduleSyncWorker/Worker.cs @@ -1,4 +1,3 @@ -using PARR.DAL.TransformServices; using PARR.EsppScheduleSync; namespace PARR.EsppScheduleSyncWorker @@ -16,12 +15,12 @@ namespace PARR.EsppScheduleSyncWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - scheduleSyncher.Start(); + await scheduleSyncher.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - scheduleSyncher.Stop(); + scheduleSyncher.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.EsppScheduleSyncWorker/appsettings.json b/PARR.EsppScheduleSyncWorker/appsettings.json index b87ba7d4..bd835bb8 100644 --- a/PARR.EsppScheduleSyncWorker/appsettings.json +++ b/PARR.EsppScheduleSyncWorker/appsettings.json @@ -32,7 +32,8 @@ "HostName": "parr-rabbitmq", "QueueName": "parr-espp-schedulers", "User": "espp_schedulers_reader", - "Password": "KjdhGLJDshgd&^%84S2" + "Password": "KjdhGLJDshgd&^%84S2", + "PrefetchCount": 100 }, "ParsingSeparator": "<|>" } diff --git a/PARR.EsppSync/SyncService.cs b/PARR.EsppSync/SyncService.cs index 8fddb68f..3d38ac66 100644 --- a/PARR.EsppSync/SyncService.cs +++ b/PARR.EsppSync/SyncService.cs @@ -12,21 +12,25 @@ namespace PARR.EsppSync internal class SyncService : ISyncService where EsppObject : class, IEsppObject { private readonly ILogger> logger; - private readonly ITemplateService templateService; - private readonly IRobotConfigurationService robotConfigurationService; - private readonly IShortcodesService shortcodesService; + private readonly IServiceProvider serviceProvider; + + //private readonly ITemplateService templateService; + //private readonly IRobotConfigurationService robotConfigurationService; + //private readonly IShortcodesService shortcodesService; public SyncService( ILogger> logger, - ITemplateService templateService, - IRobotConfigurationService robotConfigurationService, - IShortcodesService shortcodesService + IServiceProvider serviceProvider + //ITemplateService templateService, + //IRobotConfigurationService robotConfigurationService, + //IShortcodesService shortcodesService ) { this.logger = logger; - this.templateService = templateService; - this.robotConfigurationService = robotConfigurationService; - this.shortcodesService = shortcodesService; + this.serviceProvider = serviceProvider; + //this.templateService = templateService; + //this.robotConfigurationService = robotConfigurationService; + //this.shortcodesService = shortcodesService; } @@ -52,106 +56,114 @@ namespace PARR.EsppSync return; } - try + using (var scope = serviceProvider.CreateScope()) { - var template = await templateService.GetTemplateByNameAsync(esppObject.TemplateName); + var templateService = GetServiceInScope(scope); + var robotConfigurationService = GetServiceInScope(scope); + var shortcodesService = GetServiceInScope(scope); - //todo: существует в ЕСПП но отсутствует в ПАРР. Может его деактивировать или еще что-то сделать. Пока просто пропустим - if (template == null) + + try { - logger.LogWarning($"Найден объект в ЕСПП с именем шаблона {esppObject.TemplateName} незарегистрированный в ПАРР."); + var template = await templateService.GetTemplateByNameAsync(esppObject.TemplateName); - return; - } - else - { - var dbObjectInEsppObject = converterToEsppObject.Invoke(template); - - //Проверяем наличие Shortcode в полях объекта из БД - var properties = dbObjectInEsppObject.GetType().GetProperties(); - - foreach (PropertyInfo property in properties) + //todo: существует в ЕСПП но отсутствует в ПАРР. Может его деактивировать или еще что-то сделать. Пока просто пропустим + if (template == null) { - var value = property.GetValue(dbObjectInEsppObject)?.ToString(); + logger.LogWarning($"Найден объект в ЕСПП с именем шаблона {esppObject.TemplateName} незарегистрированный в ПАРР."); - if (value != null && shortcodesService.isAnyShortcodes(value)) - property.SetValue(dbObjectInEsppObject, await shortcodesService.ApplyShortcodesAsync(value, template.UnitId, template.JobId)); - } - - bool isChanged = false; - //TODO: FIX ME Please, BRO - //bool isChanged; - - // если в БД isActive == false, то синхронизировать только по полям из IEsppObject - - if (dbObjectInEsppObject.IsActive == false) - { - logger.LogDebug($"Объект деактивирован в ПАРР. Сравниваем только обязательные поля. {esppObject.TemplateName}"); - var lightDbObj = new EsppLightObject(dbObjectInEsppObject); - var lightEsppObject = new EsppLightObject(esppObject); - - isChanged = IsChanged(lightEsppObject, lightDbObj); + return; } else { - logger.LogDebug($"Объект активирован в ПАРР. Сравниваем все поля. {esppObject.TemplateName}"); - isChanged = IsChanged(esppObject, dbObjectInEsppObject); - } + var dbObjectInEsppObject = converterToEsppObject.Invoke(template); - if (isChanged) - { - logger.LogInformation($"Есть изменения, требуется обновление. {esppObject.TemplateName}"); + //Проверяем наличие Shortcode в полях объекта из БД + var properties = dbObjectInEsppObject.GetType().GetProperties(); - var config = robotConfigurationService.GetFromTemplateByRobotCode(esppObject.Robot, ref template); - - // Если предыдущий статус был Create или Ок, то ставим ему Update - // пусть даже Create завершился с ошибками, но раз он уже есть в ЕСПП, то изменим на Update, сбросим все счетчики и пусть попробует обновить и исправить все - if (config.TaskStatusCode != (int)TaskStatusEnum.Updating) + foreach (PropertyInfo property in properties) { - SetUpdateStatus(ref template, robotConfigurationService, esppObject.Robot); + var value = property.GetValue(dbObjectInEsppObject)?.ToString(); - if (!await templateService.CommitAsync()) - logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); - else - logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Updating}, Template {template.Name}, Robot: {esppObject.Robot}"); + if (value != null && shortcodesService.isAnyShortcodes(value)) + property.SetValue(dbObjectInEsppObject, await shortcodesService.ApplyShortcodesAsync(value, template.UnitId, template.JobId)); + } + + bool isChanged = false; + //TODO: FIX ME Please, BRO + //bool isChanged; + + // если в БД isActive == false, то синхронизировать только по полям из IEsppObject + + if (dbObjectInEsppObject.IsActive == false) + { + logger.LogDebug($"Объект деактивирован в ПАРР. Сравниваем только обязательные поля. {esppObject.TemplateName}"); + var lightDbObj = new EsppLightObject(dbObjectInEsppObject); + var lightEsppObject = new EsppLightObject(esppObject); + + isChanged = IsChanged(lightEsppObject, lightDbObj); } else { - // Если пред статус был Update, то ничего не делаем, так его и оставляем, не сбрасывам кол-во попыток и ошибок - logger.LogInformation($"Есть изменения в Template {template.Name}, но предыдущий статус TaskStatusCode: {(TaskStatusEnum)config.TaskStatusCode}. Не меняем статус, будем разбираться вручную."); + logger.LogDebug($"Объект активирован в ПАРР. Сравниваем все поля. {esppObject.TemplateName}"); + isChanged = IsChanged(esppObject, dbObjectInEsppObject); } - #region Old logic - - //SetUpdateStatus(ref template, robotConfigurationService, esppObject.Robot); - - //if (!await templateService.CommitAsync()) - // logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); - //else - // logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Updating}, Template {template.Name}, Robot: {esppObject.Robot}"); - #endregion - }//надо ли проверять если не изменился, но был статус Updating не понятно. Доверяем роботу пока, что после окончания работ он точно сообщит - else - { - logger.LogInformation($"Нет изменений, обновление не требуется. {esppObject.TemplateName}"); - - //если все поля совпали - //проверяем, какой был статус предыдущий статус в БД, если он был не Ок, то ставим ему ОК - var robotConfig = robotConfigurationService.GetFromTemplateByRobotCode(esppObject.Robot, ref template); - if (robotConfig.TaskStatusCode != (int)TaskStatusEnum.Ok) + if (isChanged) { - robotConfigurationService.ChangeTaskStatus(TaskStatusEnum.Ok, ref robotConfig); - if (!await templateService.CommitAsync()) - logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); + logger.LogInformation($"Есть изменения, требуется обновление. {esppObject.TemplateName}"); + + var config = robotConfigurationService.GetFromTemplateByRobotCode(esppObject.Robot, ref template); + + // Если предыдущий статус был Create или Ок, то ставим ему Update + // пусть даже Create завершился с ошибками, но раз он уже есть в ЕСПП, то изменим на Update, сбросим все счетчики и пусть попробует обновить и исправить все + if (config.TaskStatusCode != (int)TaskStatusEnum.Updating) + { + SetUpdateStatus(ref template, robotConfigurationService, esppObject.Robot); + + if (!await templateService.CommitAsync()) + logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); + else + logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Updating}, Template {template.Name}, Robot: {esppObject.Robot}"); + } else - logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Ok}, Template {template.Name}, Robot: {esppObject.Robot}"); + { + // Если пред статус был Update, то ничего не делаем, так его и оставляем, не сбрасывам кол-во попыток и ошибок + logger.LogInformation($"Есть изменения в Template {template.Name}, но предыдущий статус TaskStatusCode: {(TaskStatusEnum)config.TaskStatusCode}. Не меняем статус, будем разбираться вручную."); + } + + #region Old logic + + //SetUpdateStatus(ref template, robotConfigurationService, esppObject.Robot); + + //if (!await templateService.CommitAsync()) + // logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); + //else + // logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Updating}, Template {template.Name}, Robot: {esppObject.Robot}"); + #endregion + }//надо ли проверять если не изменился, но был статус Updating не понятно. Доверяем роботу пока, что после окончания работ он точно сообщит + else + { + logger.LogInformation($"Нет изменений, обновление не требуется. {esppObject.TemplateName}"); + + //если все поля совпали + //проверяем, какой был статус предыдущий статус в БД, если он был не Ок, то ставим ему ОК + var robotConfig = robotConfigurationService.GetFromTemplateByRobotCode(esppObject.Robot, ref template); + if (robotConfig.TaskStatusCode != (int)TaskStatusEnum.Ok) + { + robotConfigurationService.ChangeTaskStatus(TaskStatusEnum.Ok, ref robotConfig); + if (!await templateService.CommitAsync()) + logger.LogError($"Не удалось изменить запись Template {template.Name}, Robot: {esppObject.Robot}"); + else + logger.LogInformation($"Установлен принудительный статус {TaskStatusEnum.Ok}, Template {template.Name}, Robot: {esppObject.Robot}"); + } } } } - } - catch (Exception ex) - { - logger.LogError(ex, $"Ошибка синхронизации объекта АСУ ЕСПП {esppObject.TemplateName}"); + catch (Exception ex) + { + logger.LogError(ex, $"Ошибка синхронизации объекта АСУ ЕСПП {esppObject.TemplateName}"); + } } } diff --git a/PARR.EsppTemplateSync/ITemplateSyncer.cs b/PARR.EsppTemplateSync/ITemplateSyncer.cs index 197e9173..d4bf58ac 100644 --- a/PARR.EsppTemplateSync/ITemplateSyncer.cs +++ b/PARR.EsppTemplateSync/ITemplateSyncer.cs @@ -2,7 +2,7 @@ { public interface ITemplateSyncer { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.EsppTemplateSync/Settings/GlobalSettings.cs b/PARR.EsppTemplateSync/Settings/GlobalSettings.cs index 1fbb53f1..5b0037ed 100644 --- a/PARR.EsppTemplateSync/Settings/GlobalSettings.cs +++ b/PARR.EsppTemplateSync/Settings/GlobalSettings.cs @@ -17,6 +17,7 @@ namespace PARR.EsppTemplateSync.Settings 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; } internal class StorageSettings diff --git a/PARR.EsppTemplateSync/TemplateFileSyncer.cs b/PARR.EsppTemplateSync/TemplateFileSyncer.cs index 1b44ac00..01c8d7b9 100644 --- a/PARR.EsppTemplateSync/TemplateFileSyncer.cs +++ b/PARR.EsppTemplateSync/TemplateFileSyncer.cs @@ -39,7 +39,7 @@ namespace PARR.EsppTemplateSync filterExtensions = string.Join(",", storageSettings.EsppTemplates!.AllowedExtensions).Replace(".", "*."); } - public void Start() + public void StartAsync() { //смотрим, доступна ли папка if (!Directory.Exists(dirPath)) @@ -80,7 +80,7 @@ namespace PARR.EsppTemplateSync } - public void Stop() + public void StopAsync() { //watcher?.Dispose(); logger.LogInformation($"=== === === FileWatcher остановлен === === ==="); @@ -132,5 +132,14 @@ namespace PARR.EsppTemplateSync //} } + Task ITemplateSyncer.StartAsync() + { + throw new NotImplementedException(); + } + + Task ITemplateSyncer.StopAsync() + { + throw new NotImplementedException(); + } } } diff --git a/PARR.EsppTemplateSync/TemplateMQSyncer.cs b/PARR.EsppTemplateSync/TemplateMQSyncer.cs index 78152e67..7588870f 100644 --- a/PARR.EsppTemplateSync/TemplateMQSyncer.cs +++ b/PARR.EsppTemplateSync/TemplateMQSyncer.cs @@ -15,21 +15,25 @@ namespace PARR.EsppTemplateSync private readonly GlobalSettings globalSettings; private readonly IMqService mqService; private readonly SettingsFromDb settingsFromDb; - private readonly IServiceProvider serviceProvider; + private readonly ISyncService syncService; + + //private readonly IServiceProvider serviceProvider; public TemplateMQSyncer( ILogger logger, GlobalSettings globalSettings, IMqService mqService, SettingsFromDb settingsFromDb, - IServiceProvider serviceProvider + ISyncService syncService + //IServiceProvider serviceProvider ) { this.logger = logger; this.globalSettings = globalSettings; this.mqService = mqService; this.settingsFromDb = settingsFromDb; - this.serviceProvider = serviceProvider; + this.syncService = syncService; + //this.serviceProvider = serviceProvider; if (globalSettings.MqSettings == null) { logger.LogError("Нет секции настроек хранилища. MqSettings, EsppTemplates"); @@ -37,10 +41,10 @@ namespace PARR.EsppTemplateSync } } - public void Start() + public async Task StartAsync() { - var isConnected = mqService.InitConsumer(globalSettings!.MqSettings!, SyncTemplateAsync); + var isConnected = await mqService.InitConsumerAsync(globalSettings!.MqSettings!, SyncTemplateAsync); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); @@ -48,9 +52,9 @@ namespace PARR.EsppTemplateSync logger.LogInformation($"Запущена проверка очереди {globalSettings.MqSettings!.QueueName}."); } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); logger.LogInformation($"=== === === Соединение с очередью {globalSettings.MqSettings!.QueueName} закрыто === === ==="); } @@ -58,16 +62,16 @@ namespace PARR.EsppTemplateSync private async Task SyncTemplateAsync(string str) { - using (var scope = serviceProvider.CreateScope()) - { - var syncService = scope.ServiceProvider.GetService>(); + //using (var scope = serviceProvider.CreateScope()) + //{ + // var syncService = scope.ServiceProvider.GetService>(); - if (syncService == null) - throw new Exception($"Не смог получить серивс {nameof(ITemplateSyncer)}"); + // if (syncService == null) + // throw new Exception($"Не смог получить серивс {nameof(ITemplateSyncer)}"); + + await syncService.SyncEsppObjectAsync(str, ParseStrToEsppObject, ConvertDbObjToEsppObj); + //} - await syncService.SyncEsppObjectAsync(str, ParseStrToEsppObject, ConvertDbObjToEsppObj); - } - } diff --git a/PARR.EsppTemplateSyncWorker/Worker.cs b/PARR.EsppTemplateSyncWorker/Worker.cs index ccada57b..327df192 100644 --- a/PARR.EsppTemplateSyncWorker/Worker.cs +++ b/PARR.EsppTemplateSyncWorker/Worker.cs @@ -20,12 +20,12 @@ namespace PARR.EsppTemplateSyncWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - templateSyncer.Start(); + await templateSyncer.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - templateSyncer.Stop(); + templateSyncer.StopAsync().Wait(); return base.StopAsync(cancellationToken); } } diff --git a/PARR.EsppTemplateSyncWorker/appsettings.json b/PARR.EsppTemplateSyncWorker/appsettings.json index d402b5ff..6d56933d 100644 --- a/PARR.EsppTemplateSyncWorker/appsettings.json +++ b/PARR.EsppTemplateSyncWorker/appsettings.json @@ -39,7 +39,8 @@ "HostName": "parr-rabbitmq", "QueueName": "parr-espp-templates", "User": "espp_templates_reader", - "Password": "P@ssReaderPtk202!" + "Password": "P@ssReaderPtk202!", + "PrefetchCount": 100 }, "StorageSettings": { "CheckIntervalSeconds": 30 diff --git a/PARR.GeneratorTemplates/GeneratorTemplate.cs b/PARR.GeneratorTemplates/GeneratorTemplate.cs index 09f67b42..ea09ab68 100644 --- a/PARR.GeneratorTemplates/GeneratorTemplate.cs +++ b/PARR.GeneratorTemplates/GeneratorTemplate.cs @@ -6,6 +6,7 @@ using PARR.BLL.Services.Interfaces; using PARR.DAL.Models; using PARR.GeneratorTemplates.Services; using PARR.GeneratorTemplates.Settings; +using System.Threading.Tasks; namespace PARR.GeneratorTemplates { @@ -35,7 +36,7 @@ namespace PARR.GeneratorTemplates this.serviceProvider = serviceProvider; } - public void Start() + public async Task StartAsync() { #region test // string msg = "{\"WorkId\":\"91219bec-e5f7-4a1a-b4c9-c135035d2498\",\"ApplicationId\":\"13044a1f-fa0c-476f-b501-2b8d06bac9f8\",\"Action\":\"create\",\"Ek\":\"ВРТ*ДВС\",\"StatusEk\":0}"; @@ -44,15 +45,15 @@ namespace PARR.GeneratorTemplates // await GenerateTemplates(msg); #endregion - var isConnected = mqService.InitConsumer(mqSettings, GenerateTemplates); + var isConnected = await mqService.InitConsumerAsync(mqSettings, GenerateTemplates); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } diff --git a/PARR.GeneratorTemplates/IGeneratorTemplate.cs b/PARR.GeneratorTemplates/IGeneratorTemplate.cs index 89c2f408..d0a33e71 100644 --- a/PARR.GeneratorTemplates/IGeneratorTemplate.cs +++ b/PARR.GeneratorTemplates/IGeneratorTemplate.cs @@ -2,7 +2,7 @@ { public interface IGeneratorTemplate { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.GeneratorTemplates/Settings/MqSettings.cs b/PARR.GeneratorTemplates/Settings/MqSettings.cs index 97d78349..3dbbb1ac 100644 --- a/PARR.GeneratorTemplates/Settings/MqSettings.cs +++ b/PARR.GeneratorTemplates/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.GeneratorTemplates.Settings 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.GeneratorTemplatesWorker/Worker.cs b/PARR.GeneratorTemplatesWorker/Worker.cs index e5e37dc9..18d0c73d 100644 --- a/PARR.GeneratorTemplatesWorker/Worker.cs +++ b/PARR.GeneratorTemplatesWorker/Worker.cs @@ -15,7 +15,7 @@ namespace PARR.GeneratorTemplatesWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - generatorTemplate.Start(); + await generatorTemplate.StartAsync(); //while (!stoppingToken.IsCancellationRequested) //{ // _logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now); @@ -25,7 +25,7 @@ namespace PARR.GeneratorTemplatesWorker public override Task StopAsync(CancellationToken cancellationToken) { - generatorTemplate.Stop(); + generatorTemplate.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.JobAutoControl/JobAutoControlManager.cs b/PARR.JobAutoControl/JobAutoControlManager.cs index ebc12ada..4454f96d 100644 --- a/PARR.JobAutoControl/JobAutoControlManager.cs +++ b/PARR.JobAutoControl/JobAutoControlManager.cs @@ -87,7 +87,8 @@ namespace PARR.JobAutoControl { job.JobEkMasks.ToList().ForEach(ek => { - job.JobAutoControlInEkStatuses.ToList().ForEach(ekStatus => + //todo: проверить, нормально ли отрабатывает async в Foreach + job.JobAutoControlInEkStatuses.ToList().ForEach(async ekStatus => { var request = new GeneratorTemplateMq { @@ -103,7 +104,7 @@ namespace PARR.JobAutoControl var msg = JsonSerializer.Serialize(request); - var result = mqService.Send(mqSettings.Generator, new[] { msg }); + var result =await mqService.SendAsync(mqSettings.Generator, new[] { msg }); if (result.IsSuccess) logger.LogInformation($"Отправлен запрос на генерацию: {msg}"); @@ -135,7 +136,8 @@ namespace PARR.JobAutoControl if (!jobsToDeactivate.Any()) return; - jobsToDeactivate.ForEach(applicationInWorkId => + //todo: проверить, нормально ли отрабатывает async в foreach + jobsToDeactivate.ForEach(async applicationInWorkId => { var request = new TemplateActivatorMq { @@ -146,7 +148,7 @@ namespace PARR.JobAutoControl var msg = JsonSerializer.Serialize(request); - var result = mqService.Send(mqSettings.Activator, new[] { msg }); + var result = await mqService.SendAsync(mqSettings.Activator, new[] { msg }); if (result.IsSuccess) logger.LogInformation($"Отправлен запрос на деактивацию: {msg}"); @@ -174,7 +176,8 @@ namespace PARR.JobAutoControl if (!jobsToActivate.Any()) return; - jobsToActivate.ForEach(applicationInWorkId => + //todo: проверить, нормально ли отрабатывает async в foreach + jobsToActivate.ForEach(async applicationInWorkId => { var request = new TemplateActivatorMq { @@ -185,7 +188,7 @@ namespace PARR.JobAutoControl var msg = JsonSerializer.Serialize(request); - var result = mqService.Send(mqSettings.Activator, new[] { msg }); + var result = await mqService.SendAsync(mqSettings.Activator, new[] { msg }); if (result.IsSuccess) logger.LogInformation($"Отправлен запрос на активацию: {msg}"); diff --git a/PARR.JobAutoControl/Settings/MqSettings.cs b/PARR.JobAutoControl/Settings/MqSettings.cs index 6d201ccc..097fe9c0 100644 --- a/PARR.JobAutoControl/Settings/MqSettings.cs +++ b/PARR.JobAutoControl/Settings/MqSettings.cs @@ -14,6 +14,7 @@ namespace PARR.JobAutoControl.Settings 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; } internal class Activator : IMqSettings @@ -22,6 +23,7 @@ namespace PARR.JobAutoControl.Settings 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.Master/Services/OrderMasterService.cs b/PARR.Master/Services/OrderMasterService.cs index 55c2b106..a64e44aa 100644 --- a/PARR.Master/Services/OrderMasterService.cs +++ b/PARR.Master/Services/OrderMasterService.cs @@ -66,7 +66,7 @@ namespace PARR.Master.Services if (objs.Any()) { - var sendResult = mqService.Send(mqSettings, objs.ToArray()); + var sendResult = await mqService.SendAsync(mqSettings, objs.ToArray()); if (sendResult.IsSuccess) { logger.LogInformation($"Наряды переданы в RabbitMQ: {objs.Count()}"); diff --git a/PARR.Master/Settings/MqSettings.cs b/PARR.Master/Settings/MqSettings.cs index 9df5a304..826593e4 100644 --- a/PARR.Master/Settings/MqSettings.cs +++ b/PARR.Master/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.Master.Settings 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.MockData/EsppDataImporter.cs b/PARR.MockData/EsppDataImporter.cs index 06463c72..c33c843c 100644 --- a/PARR.MockData/EsppDataImporter.cs +++ b/PARR.MockData/EsppDataImporter.cs @@ -6,44 +6,45 @@ namespace PARR.MockData { internal class EsppDataImporter { - public void Import() { - var factory = new ConnectionFactory() { HostName = "10.99.253.216" }; - factory.UserName = "espp_templates_robot"; - factory.Password = "P@ssRwEsppTempl202"; + public void Import() + { + //var factory = new ConnectionFactory() { HostName = "10.99.253.216" }; + //factory.UserName = "espp_templates_robot"; + //factory.Password = "P@ssRwEsppTempl202"; - string[] rows = EsppExport.data.Replace("\r", "").Split('\n'); - using (var connection = factory.CreateConnection()) - using (var channel = connection.CreateModel()) - { - //https://www.rabbitmq.com/lazy-queues.html - //Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память - //Формируем соответствующий аргумент - Dictionary argums = new Dictionary(); - argums.Add("x-queue-mode", "lazy"); + //string[] rows = EsppExport.data.Replace("\r", "").Split('\n'); + //using (var connection = factory.CreateConnection()) + //using (var channel = connection.CreateModel()) + //{ + // //https://www.rabbitmq.com/lazy-queues.html + // //Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память + // //Формируем соответствующий аргумент + // Dictionary argums = new Dictionary(); + // argums.Add("x-queue-mode", "lazy"); - //Объявляем очередь с которой будем работать. Если такой очереди ещё нет, то создатся. Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. В целом эти параметры можно посмотреть в админке RabbitMQ - channel.QueueDeclare(queue: "parr-espp-templates", - durable: true,//If you want to make sure your messages are persisted, you need to create a durable queue. To do this, simply set the durable flag to true when creating the queue. - exclusive: false, - autoDelete: false, - arguments: argums); - //перебираем строки и отправляем строки сообщениями - for (int i = 0; i < rows.Count(); i++) - { - var body = Encoding.UTF8.GetBytes(rows[i]); - IBasicProperties props = channel.CreateBasicProperties(); - //If you don’t mark your messages as persistent, then RabbitMQ will store them in memory. This is fine if you’re just playing around or building a prototype. But if you’re building a production system, then you need to make sure that your messages are persisted to disk so that they’re not lost if the server crashes. - //Говорим менеджеру RabbitMq при получении хранить сообщения на диске - props.DeliveryMode = 2; - //время жизни в мс. на этапе отладки сделаем коротким чтобы не устраивать помойку - props.Expiration = "60000"; + // //Объявляем очередь с которой будем работать. Если такой очереди ещё нет, то создатся. Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. В целом эти параметры можно посмотреть в админке RabbitMQ + // channel.QueueDeclare(queue: "parr-espp-templates", + // durable: true,//If you want to make sure your messages are persisted, you need to create a durable queue. To do this, simply set the durable flag to true when creating the queue. + // exclusive: false, + // autoDelete: false, + // arguments: argums); + // //перебираем строки и отправляем строки сообщениями + // for (int i = 0; i < rows.Count(); i++) + // { + // var body = Encoding.UTF8.GetBytes(rows[i]); + // IBasicProperties props = channel.CreateBasicProperties(); + // //If you don’t mark your messages as persistent, then RabbitMQ will store them in memory. This is fine if you’re just playing around or building a prototype. But if you’re building a production system, then you need to make sure that your messages are persisted to disk so that they’re not lost if the server crashes. + // //Говорим менеджеру RabbitMq при получении хранить сообщения на диске + // props.DeliveryMode = 2; + // //время жизни в мс. на этапе отладки сделаем коротким чтобы не устраивать помойку + // props.Expiration = "60000"; - channel.BasicPublish(exchange: "", - routingKey: "parr-espp-templates",//совпадает с именем очереди иначе не принимает - basicProperties: props, - body: body); - } - } + // channel.BasicPublish(exchange: "", + // routingKey: "parr-espp-templates",//совпадает с именем очереди иначе не принимает + // basicProperties: props, + // body: body); + // } + //} } } diff --git a/PARR.MockData/PARR.MockData.csproj b/PARR.MockData/PARR.MockData.csproj index f050d586..ec3023c3 100644 --- a/PARR.MockData/PARR.MockData.csproj +++ b/PARR.MockData/PARR.MockData.csproj @@ -10,7 +10,6 @@ - diff --git a/PARR.TemplateActivator/IMqTemplateActivator.cs b/PARR.TemplateActivator/IMqTemplateActivator.cs index 02eae893..9c681f9b 100644 --- a/PARR.TemplateActivator/IMqTemplateActivator.cs +++ b/PARR.TemplateActivator/IMqTemplateActivator.cs @@ -2,6 +2,6 @@ public interface IMqTemplateActivator { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } diff --git a/PARR.TemplateActivator/MqTemplateActivator.cs b/PARR.TemplateActivator/MqTemplateActivator.cs index c5ea3d4f..b08e303a 100644 --- a/PARR.TemplateActivator/MqTemplateActivator.cs +++ b/PARR.TemplateActivator/MqTemplateActivator.cs @@ -3,6 +3,7 @@ using Microsoft.Extensions.Logging; using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using PARR.Constants; +using System.Threading.Tasks; namespace PARR.TemplateActivator; @@ -31,9 +32,9 @@ internal class MqTemplateActivator : IMqTemplateActivator this.validatorService = validatorService; this.serviceProvider = serviceProvider; } - public void Start() + public async Task StartAsync() { - var isConnected = mqService.InitConsumer(mqSettings, UpdateTemplateStatusAsync); + var isConnected = await mqService.InitConsumerAsync(mqSettings, UpdateTemplateStatusAsync); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); @@ -64,8 +65,8 @@ internal class MqTemplateActivator : IMqTemplateActivator } } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } } diff --git a/PARR.TemplateActivator/Settings/MqSettings.cs b/PARR.TemplateActivator/Settings/MqSettings.cs index 426a9831..0183195e 100644 --- a/PARR.TemplateActivator/Settings/MqSettings.cs +++ b/PARR.TemplateActivator/Settings/MqSettings.cs @@ -8,5 +8,6 @@ internal class MqSettings : IMqSettings 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.TemplateActivatorWorker/Worker.cs b/PARR.TemplateActivatorWorker/Worker.cs index 4888141e..7aee0811 100644 --- a/PARR.TemplateActivatorWorker/Worker.cs +++ b/PARR.TemplateActivatorWorker/Worker.cs @@ -15,12 +15,12 @@ namespace PARR.TemplateActivatorWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - mqTemplateActivator.Start(); + await mqTemplateActivator.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - mqTemplateActivator.Stop(); + mqTemplateActivator.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.TemplateDistributor/IMqTemplateDistributor.cs b/PARR.TemplateDistributor/IMqTemplateDistributor.cs index 53146f75..76c50066 100644 --- a/PARR.TemplateDistributor/IMqTemplateDistributor.cs +++ b/PARR.TemplateDistributor/IMqTemplateDistributor.cs @@ -2,7 +2,7 @@ { public interface IMqTemplateDistributor { - void Start(); - void Stop(); + Task StartAsync(); + Task StopAsync(); } } diff --git a/PARR.TemplateDistributor/MqTemplateDistributor.cs b/PARR.TemplateDistributor/MqTemplateDistributor.cs index 5d55d238..7e88735d 100644 --- a/PARR.TemplateDistributor/MqTemplateDistributor.cs +++ b/PARR.TemplateDistributor/MqTemplateDistributor.cs @@ -4,6 +4,7 @@ using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using PARR.TemplateDistributor.Services; using PARR.TemplateDistributor.Settings; +using System.Threading.Tasks; namespace PARR.TemplateDistributor { @@ -31,18 +32,18 @@ namespace PARR.TemplateDistributor this.validatorService = validatorService; this.serviceProvider = serviceProvider; } - public void Start() + public async Task StartAsync() { - var isConnected = mqService.InitConsumer(mqSettings, UpdateScheduleAsync); + var isConnected = await mqService.InitConsumerAsync(mqSettings, UpdateScheduleAsync); if (!isConnected) throw new Exception("Ошибка при подключении к RabbitMq"); } - public void Stop() + public async Task StopAsync() { - mqService.Dispose(); + await mqService.DisposeAsync(); } private async Task UpdateScheduleAsync(string msg) diff --git a/PARR.TemplateDistributor/Settings/MqSettings.cs b/PARR.TemplateDistributor/Settings/MqSettings.cs index 80c9a3b4..6ae427c7 100644 --- a/PARR.TemplateDistributor/Settings/MqSettings.cs +++ b/PARR.TemplateDistributor/Settings/MqSettings.cs @@ -8,5 +8,6 @@ namespace PARR.TemplateDistributor.Settings 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.TemplateDistributorWorker/Worker.cs b/PARR.TemplateDistributorWorker/Worker.cs index 5870802b..364b5f6e 100644 --- a/PARR.TemplateDistributorWorker/Worker.cs +++ b/PARR.TemplateDistributorWorker/Worker.cs @@ -15,12 +15,12 @@ namespace PARR.TemplateDistributorWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - mqTemplateDistributor.Start(); + await mqTemplateDistributor.StartAsync(); } public override Task StopAsync(CancellationToken cancellationToken) { - mqTemplateDistributor.Stop(); + mqTemplateDistributor.StopAsync().Wait(); return base.StopAsync(cancellationToken); } diff --git a/PARR.TemplateGeneratorWorker/Worker.cs b/PARR.TemplateGeneratorWorker/Worker.cs index 9372ddb4..3f26262a 100644 --- a/PARR.TemplateGeneratorWorker/Worker.cs +++ b/PARR.TemplateGeneratorWorker/Worker.cs @@ -35,19 +35,16 @@ namespace PARR.TemplateGeneratorWorker protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - while (!stoppingToken.IsCancellationRequested) - { - var isConnected = mqService.InitConsumer(mqSettings, GenerateTemplateAsync); + var isConnected = await mqService.InitConsumerAsync(mqSettings, GenerateTemplateAsync); if (!isConnected) throw new Exception(" RabbitMq"); - } } public override Task StopAsync(CancellationToken cancellationToken) { - mqService.Dispose(); + mqService.DisposeAsync().AsTask().Wait(); return base.StopAsync(cancellationToken); }