From 2a6fd7177d83890d34f3e271fff83dcb6b08bc45 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Wed, 15 Apr 2026 16:48:23 +1000 Subject: [PATCH 1/6] =?UTF-8?q?feat(core):=20TaskManagementService=20-=20?= =?UTF-8?q?=D0=BE=D1=82=D0=BF=D1=80=D0=B0=D0=B2=D0=BA=D0=B0=20=D1=81=D0=BE?= =?UTF-8?q?=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../V1/Statistics/StatWorkloadController.cs | 48 ++++++++++++++- PARR.Core/DependencyInjection.cs | 3 + PARR.Core/PARR.Core.csproj | 1 - .../Implementations}/TaskManagementService.cs | 60 +++++++++++-------- .../Interfaces}/ITaskManagementService.cs | 2 +- PARR.DAL/ParrDalInstaller.cs | 3 - PARR.DAL/Repositories/Base/BaseRepository.cs | 3 + .../Common/Rabbit/Messages/TaskMessage.cs | 16 +++++ 8 files changed, 103 insertions(+), 33 deletions(-) rename {PARR.DAL/TaskServices => PARR.Core/Services/Task/Implementations}/TaskManagementService.cs (57%) rename {PARR.DAL/TaskServices => PARR.Core/Services/Task/Interfaces}/ITaskManagementService.cs (95%) create mode 100644 PARR.Domain/Common/Rabbit/Messages/TaskMessage.cs diff --git a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs index c688fd7f..f4c2a764 100644 --- a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs +++ b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs @@ -1,8 +1,14 @@ using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.Mvc; using PARR.API.Contracts.V1; +using PARR.API.Contracts.V1.Responses.Base; using PARR.API.Controllers.V1.Base; +using PARR.API.Services.Interfaces; +using PARR.API.Settings; +using PARR.Core.Services.Task.Interfaces; using PARR.Domain.Common.Roles; +using PARR.Domain.Entities.Base.History; +using PARR.Domain.Enums; namespace PARR.API.Controllers.V1.Statistics { @@ -12,9 +18,22 @@ namespace PARR.API.Controllers.V1.Statistics [Authorize(Roles = ParrRoles.EsppRobot.RoleOrAdmin)] public class StatWorkloadController : BaseApiController { - public StatWorkloadController() - { + private readonly ITaskManagementService taskManagementService; + private readonly IClientService clientService; + private readonly MqSettings mqSettings; + private readonly ILogger logger; + public StatWorkloadController( + ITaskManagementService taskManagementService, + IClientService clientService, + MqSettings mqSettings, + ILogger logger + ) + { + this.taskManagementService = taskManagementService; + this.clientService = clientService; + this.mqSettings = mqSettings; + this.logger = logger; } @@ -25,8 +44,31 @@ namespace PARR.API.Controllers.V1.Statistics [HttpPost(ApiRoutes.Workload.Build)] public async Task Build() { + var initiator = new HistoryInitiator + { + InitiatorIp = clientService.GetClientIp()?.ToString(), + InitiatorParrComponentId = ParrComponentsEnum.Api, + InitiatorComment = "Отправлен запрос из API на формирование отчета о загруженности" + }; - return Ok(); + try + { + mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings); + + if (queueSettings == null) + { + logger.LogError("Не удалось получить настройки очереди для TaskType: {TaskType}", TaskTypeEnum.Workload); + throw new InvalidOperationException($"Не удалось получить настройки очереди для TaskTypeEnum.Workload"); + } + + var taskId = await taskManagementService.CreateTaskAsync(TaskTypeEnum.Workload, default, initiator, queueSettings); + + return Ok(new Response(null, true, null!, $"Создана задача {taskId} на формирование отчета")); + } + catch (Exception ex) + { + return BadRequest(new Response(false, new List { new ErrorModel { Message = ex.Message } })); + } } diff --git a/PARR.Core/DependencyInjection.cs b/PARR.Core/DependencyInjection.cs index c003bf91..be87155e 100644 --- a/PARR.Core/DependencyInjection.cs +++ b/PARR.Core/DependencyInjection.cs @@ -3,6 +3,8 @@ using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using PARR.Core.Common.Implementations; using PARR.Core.Common.Interfaces; +using PARR.Core.Services.Task.Implementations; +using PARR.Core.Services.Task.Interfaces; using PARR.Domain.Settings; namespace PARR.Core @@ -37,6 +39,7 @@ namespace PARR.Core #region Services + services.AddScoped(); //services.AddScoped(); #endregion diff --git a/PARR.Core/PARR.Core.csproj b/PARR.Core/PARR.Core.csproj index aef64738..a7524906 100644 --- a/PARR.Core/PARR.Core.csproj +++ b/PARR.Core/PARR.Core.csproj @@ -8,7 +8,6 @@ - diff --git a/PARR.DAL/TaskServices/TaskManagementService.cs b/PARR.Core/Services/Task/Implementations/TaskManagementService.cs similarity index 57% rename from PARR.DAL/TaskServices/TaskManagementService.cs rename to PARR.Core/Services/Task/Implementations/TaskManagementService.cs index dca781c5..636f2d52 100644 --- a/PARR.DAL/TaskServices/TaskManagementService.cs +++ b/PARR.Core/Services/Task/Implementations/TaskManagementService.cs @@ -1,6 +1,9 @@ using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Logging; +using PARR.Core.Common.Interfaces.RabbitServices; using PARR.Core.Repositories.Interfaces.TaskRepositories; +using PARR.Core.Services.Task.Interfaces; +using PARR.Domain.Common.Rabbit.Messages; using PARR.Domain.Entities.Base.History; using PARR.Domain.Entities.TaskEntities; using PARR.Domain.Enums; @@ -8,29 +11,32 @@ using PARR.Domain.Settings; using System.Text.Encodings.Web; using System.Text.Json; -namespace PARR.DAL.TaskServices +namespace PARR.Core.Services.Task.Implementations { internal class TaskManagementService : ITaskManagementService { private readonly ILogger logger; - private readonly ITaskTypeRepository taskTypeService; - private readonly ITaskRepository taskService; + private readonly ITaskTypeRepository taskTypeRepository; + private readonly ITaskRepository taskRepository; + private readonly IRabbitService rabbitService; public TaskManagementService( ILogger logger, - ITaskTypeRepository taskTypeService, - ITaskRepository taskService + ITaskTypeRepository taskTypeRepository, + ITaskRepository taskRepository, + IRabbitService rabbitService ) { this.logger = logger; - this.taskTypeService = taskTypeService; - this.taskService = taskService; + this.taskTypeRepository = taskTypeRepository; + this.taskRepository = taskRepository; + this.rabbitService = rabbitService; } public async Task CreateTaskAsync(TaskTypeEnum typeCode, T payload, IHistoryInitiator initiator, IMqSettings mqSettings) { - var taskType = await taskTypeService.Get().AsNoTracking().FirstOrDefaultAsync(t => t.Code == typeCode); + var taskType = await taskTypeRepository.Get().AsNoTracking().FirstOrDefaultAsync(t => t.Code == typeCode); if (taskType == null) { @@ -45,10 +51,7 @@ namespace PARR.DAL.TaskServices if (existingTask != null) { - logger.LogInformation( - "Задача типа {TypeCode} уже активна (id: {ExistingId}). Возвращаем существующую.", - typeCode, existingTask.Id); - + logger.LogInformation("Задача типа {TypeCode} уже активна (id: {ExistingId}). Возвращаем существующую.", typeCode, existingTask.Id); return existingTask.Id; } } @@ -63,26 +66,30 @@ namespace PARR.DAL.TaskServices StatusCode = TaskItemStatusEnum.Pending, Payload = payloadStr, RetryCount = 0, - ProcessedAt = null, - InitiatorIp = initiator.InitiatorIp, - InitiatorParrComponentId = initiator.InitiatorParrComponentId, - InitiatorComment = initiator.InitiatorComment + ProcessedAt = null }; - if (!await taskService.CreateAsync(task) || !await taskService.CommitAsync()) + if (!await taskRepository.CreateAsync(task) || !await taskRepository.CommitAsync(initiator)) { logger.LogError("Ошибка при сохранении задачи в БД"); - return default; + throw new DbUpdateException("Ошибка при сохранении задачи в БД"); } logger.LogInformation("Задача {TaskId} типа {TypeCode} сохранена в БД со статусом Pending", task.Id, typeCode); - // Публикуме задачу в очередь. Если вдруг даже не получится, - // задача останется в БД со статусом pending, и Reconciliation Task позже ее возмет в работу + var message = new TaskMessage { TaskId = task.Id, TypeCode = task.TypeCode }; - // todo: опубликовать в очередь - // вопросы, зачем scoped? может можно AddTransient? - // нужно придумать шаблонную модель для очереди + var sendResult = await rabbitService.SendAsync(mqSettings, new List { message }); + if (sendResult.IsSuccess) + { + logger.LogDebug("Задача {TaskId} типа {TypeCode} отправлена в очередь {QueueName}", task.Id, typeCode, mqSettings.QueueName); + } + else + { + // Не пробрасываем исключение дальше — задача создана, просто не в очереди. + // задача останется в БД со статусом pending, и Reconciliation Task позже ее возмет в работу + logger.LogError("Не удалось опубликовать задачу {TaskId} в очередь. Задача осталась в БД со статусом Pending.", task.Id); + } return task.Id; } @@ -95,7 +102,7 @@ namespace PARR.DAL.TaskServices /// private async Task GetActiveSingletonTaskAsync(TaskTypeEnum typeCode) { - return await taskService.Get().AsNoTracking() + return await taskRepository.Get().AsNoTracking() .FirstOrDefaultAsync(t => t.TypeCode == typeCode && ( t.StatusCode == TaskItemStatusEnum.Pending || t.StatusCode == TaskItemStatusEnum.Processing @@ -109,13 +116,16 @@ namespace PARR.DAL.TaskServices /// /// /// - private string PayloadToString(T payload) + private string? PayloadToString(T payload) { var jsonOptions = new JsonSerializerOptions { Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; + if (payload == null) + return null; + return JsonSerializer.Serialize(payload, jsonOptions); } diff --git a/PARR.DAL/TaskServices/ITaskManagementService.cs b/PARR.Core/Services/Task/Interfaces/ITaskManagementService.cs similarity index 95% rename from PARR.DAL/TaskServices/ITaskManagementService.cs rename to PARR.Core/Services/Task/Interfaces/ITaskManagementService.cs index 161c264e..6f8a06a5 100644 --- a/PARR.DAL/TaskServices/ITaskManagementService.cs +++ b/PARR.Core/Services/Task/Interfaces/ITaskManagementService.cs @@ -2,7 +2,7 @@ using PARR.Domain.Enums; using PARR.Domain.Settings; -namespace PARR.DAL.TaskServices +namespace PARR.Core.Services.Task.Interfaces { /// /// Управления задачами (очередями) diff --git a/PARR.DAL/ParrDalInstaller.cs b/PARR.DAL/ParrDalInstaller.cs index 938002f7..c582903c 100644 --- a/PARR.DAL/ParrDalInstaller.cs +++ b/PARR.DAL/ParrDalInstaller.cs @@ -22,7 +22,6 @@ using PARR.DAL.Services.Interfaces; using PARR.DAL.Services.Interfaces.Job; using PARR.DAL.Services.Interfaces.Schedule; using PARR.DAL.Services.Interfaces.Unit; -using PARR.DAL.TaskServices; using PARR.Domain.Settings; namespace PARR.DAL @@ -154,8 +153,6 @@ namespace PARR.DAL services.AddScoped(); services.AddScoped(); - services.AddTransient(); - #endregion //services.AddTransient(); diff --git a/PARR.DAL/Repositories/Base/BaseRepository.cs b/PARR.DAL/Repositories/Base/BaseRepository.cs index 32e213a8..c5f41def 100644 --- a/PARR.DAL/Repositories/Base/BaseRepository.cs +++ b/PARR.DAL/Repositories/Base/BaseRepository.cs @@ -99,6 +99,9 @@ namespace PARR.DAL.Repositories.Base /// private void SetInitiator(IHistoryInitiator? initiator) { + if (initiator == null) + return; + logger.LogDebug("Устанавливаю инициатора для изменений"); // Задаем инициатора только для новых и измененных записей diff --git a/PARR.Domain/Common/Rabbit/Messages/TaskMessage.cs b/PARR.Domain/Common/Rabbit/Messages/TaskMessage.cs new file mode 100644 index 00000000..1a96585c --- /dev/null +++ b/PARR.Domain/Common/Rabbit/Messages/TaskMessage.cs @@ -0,0 +1,16 @@ +using PARR.Domain.Enums; +using System.Text.Json.Serialization; + +namespace PARR.Domain.Common.Rabbit.Messages +{ + /// + /// Модель для отправки сообщений через ITaskManagement + /// + public record TaskMessage + { + public Guid TaskId { get; init; } + + [JsonConverter(typeof(JsonStringEnumConverter))] + public TaskTypeEnum TypeCode { get; init; } + } +} From 001c09edaad36cb3a82ee46c48382cbbb06345a8 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Tue, 21 Apr 2026 16:54:39 +1000 Subject: [PATCH 2/6] =?UTF-8?q?feat(dal,=20core,=20workloadBuilderWorker):?= =?UTF-8?q?=20=D0=A1=D0=B5=D1=80=D0=B2=D0=B8=D1=81=D1=8B=20=D0=BF=D0=BE=20?= =?UTF-8?q?=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=B5=20=D1=81=20=D0=BE=D1=87?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D0=B4=D1=8F=D0=BC=D0=B8=20(=D0=B7=D0=B0?= =?UTF-8?q?=D0=B4=D0=B0=D1=87=D0=B8=D0=B0=D0=BC=D0=B8).=20=D0=97=D0=B0?= =?UTF-8?q?=D0=B3=D0=BE=D1=82=D0=BE=D0=B2=D0=BA=D0=B0=20WorkloadBuilderWor?= =?UTF-8?q?ker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PARR.API.sln | 6 + .../V1/Statistics/StatWorkloadController.cs | 2 + .../Workers/IMessageConsumer.cs | 13 + PARR.Core/DependencyInjection.cs | 26 +- .../TaskRepositories/ITaskRepository.cs | 6 + .../Services/Task/Handlers/BaseTaskHandler.cs | 258 ++++++++++++++++++ .../Handlers/Factory/ITaskHandlerFactory.cs | 23 ++ .../Handlers/Factory/TaskHandlerFactory.cs | 33 +++ .../Services/Task/Handlers/HandlerResult.cs | 49 ++++ .../Task/Handlers/WorkloadReportHandler.cs | 37 +++ .../WorkloadCacheBuilderService.cs | 84 ++++++ .../IWorkloadCacheBuilderService.cs | 11 + PARR.DAL/PARR.DAL.csproj | 1 + PARR.DAL/Repositories/Base/BaseRepository.cs | 3 +- .../TaskRepositories/TaskRepository.cs | 27 +- .../PARR.WorkloadBuilderWorker.csproj | 25 ++ PARR.WorkloadBuilderWorker/Program.cs | 45 +++ .../Properties/launchSettings.json | 11 + .../Settings/WorkerSettings.cs | 9 + PARR.WorkloadBuilderWorker/Worker.cs | 29 ++ .../appsettings.Development.json | 35 +++ PARR.WorkloadBuilderWorker/appsettings.json | 33 +++ 22 files changed, 762 insertions(+), 4 deletions(-) create mode 100644 PARR.Core/Common/Interfaces/RabbitServices/Workers/IMessageConsumer.cs create mode 100644 PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs create mode 100644 PARR.Core/Services/Task/Handlers/Factory/ITaskHandlerFactory.cs create mode 100644 PARR.Core/Services/Task/Handlers/Factory/TaskHandlerFactory.cs create mode 100644 PARR.Core/Services/Task/Handlers/HandlerResult.cs create mode 100644 PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs create mode 100644 PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs create mode 100644 PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs create mode 100644 PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj create mode 100644 PARR.WorkloadBuilderWorker/Program.cs create mode 100644 PARR.WorkloadBuilderWorker/Properties/launchSettings.json create mode 100644 PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs create mode 100644 PARR.WorkloadBuilderWorker/Worker.cs create mode 100644 PARR.WorkloadBuilderWorker/appsettings.Development.json create mode 100644 PARR.WorkloadBuilderWorker/appsettings.json diff --git a/PARR.API.sln b/PARR.API.sln index 0855d326..9c8bcd5a 100644 --- a/PARR.API.sln +++ b/PARR.API.sln @@ -96,6 +96,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.Domain", "PARR.Domain\ EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.Infrastructure", "PARR.Infrastructure\PARR.Infrastructure.csproj", "{8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.WorkloadBuilderWorker", "PARR.WorkloadBuilderWorker\PARR.WorkloadBuilderWorker.csproj", "{835BD6CF-024E-47F5-BB02-C904AC9733D9}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -268,6 +270,10 @@ Global {8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Debug|Any CPU.Build.0 = Debug|Any CPU {8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Release|Any CPU.ActiveCfg = Release|Any CPU {8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Release|Any CPU.Build.0 = Release|Any CPU + {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Debug|Any CPU.Build.0 = Debug|Any CPU + {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.ActiveCfg = Release|Any CPU + {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE diff --git a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs index f4c2a764..626ed11f 100644 --- a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs +++ b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs @@ -51,6 +51,8 @@ namespace PARR.API.Controllers.V1.Statistics InitiatorComment = "Отправлен запрос из API на формирование отчета о загруженности" }; + //todo: избавиться от try/catch + try { mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings); diff --git a/PARR.Core/Common/Interfaces/RabbitServices/Workers/IMessageConsumer.cs b/PARR.Core/Common/Interfaces/RabbitServices/Workers/IMessageConsumer.cs new file mode 100644 index 00000000..8dff4cb8 --- /dev/null +++ b/PARR.Core/Common/Interfaces/RabbitServices/Workers/IMessageConsumer.cs @@ -0,0 +1,13 @@ +using PARR.Domain.Settings; + +namespace PARR.Core.Common.Interfaces.RabbitServices.Workers +{ + /// + /// Интерфейс для сервисов которые используются воркерами по подписке на очереди MQ + /// + public interface IMessageConsumer + { + Task ProcessMessagesAsync(IMqSettings mqSettings); + Task StopProcessingAsync(); + } +} diff --git a/PARR.Core/DependencyInjection.cs b/PARR.Core/DependencyInjection.cs index be87155e..84bf1e37 100644 --- a/PARR.Core/DependencyInjection.cs +++ b/PARR.Core/DependencyInjection.cs @@ -3,8 +3,12 @@ using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using PARR.Core.Common.Implementations; using PARR.Core.Common.Interfaces; +using PARR.Core.Services.Task.Handlers; +using PARR.Core.Services.Task.Handlers.Factory; using PARR.Core.Services.Task.Implementations; using PARR.Core.Services.Task.Interfaces; +using PARR.Core.Services.Workload.Implementations; +using PARR.Core.Services.Workload.Interfaces; using PARR.Domain.Settings; namespace PARR.Core @@ -37,13 +41,33 @@ namespace PARR.Core #endregion - #region Services + #region Task + + // Регистрация хендлеров (нужно регистировать каждый отдельно) + services.AddScoped(); + services.AddScoped(sp => sp.GetRequiredService());// регим для фабрики + // Еще хэндлеры... + + // Фабрика хэндлеров + services.AddScoped(); services.AddScoped(); + + #endregion + + #region Services + + //services.AddScoped(); #endregion + #region IMessageConsumer services + + services.AddSingleton(); + + #endregion + return services; } diff --git a/PARR.Core/Repositories/Interfaces/TaskRepositories/ITaskRepository.cs b/PARR.Core/Repositories/Interfaces/TaskRepositories/ITaskRepository.cs index af5c2bda..5bf52170 100644 --- a/PARR.Core/Repositories/Interfaces/TaskRepositories/ITaskRepository.cs +++ b/PARR.Core/Repositories/Interfaces/TaskRepositories/ITaskRepository.cs @@ -5,5 +5,11 @@ namespace PARR.Core.Repositories.Interfaces.TaskRepositories { public interface ITaskRepository : IBaseRepository { + /// + /// Атомарный захват задачи (взять в работу) + /// + /// + /// Возвращает кол-во обработанных строк + Task TaskCaptureAsync(Guid taskId); } } diff --git a/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs new file mode 100644 index 00000000..f8287d5f --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs @@ -0,0 +1,258 @@ +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; +using PARR.Core.Repositories.Interfaces.TaskRepositories; +using PARR.Domain.Entities.TaskEntities; +using PARR.Domain.Enums; + +namespace PARR.Core.Services.Task.Handlers +{ + /// + /// Базовый класс для всех обработчиков задач. + /// Реализует шаблонный метод ProcessAsync с общей логикой: + /// - Атомарный захват задачи + /// - Обработка исключений + /// - Обновление статусов + /// - Запись ошибок + /// - Логика retry + /// + public abstract class BaseTaskHandler + { + protected readonly ITaskRepository taskRepository; + protected readonly ITaskErrorRepository taskErrorRepository; + protected readonly ILogger logger; + + protected BaseTaskHandler(ITaskRepository taskRepository, ITaskErrorRepository taskErrorRepository, ILogger logger) + { + this.taskRepository = taskRepository; + this.taskErrorRepository = taskErrorRepository; + this.logger = logger; + } + + /// + /// Тип задачи, который обрабатывает этот хендлер. + /// Должен быть переопределен в конкретном классе. + /// + public abstract TaskTypeEnum SupportedType { get; } + + + public async Task ProcessAsync(Guid taskId) + { + //todo: может принимать не taskId, а модель TaskMessage + + logger.LogInformation("Начало обработки задачи {TaskId} (тип {Type})", taskId, SupportedType); + + // Загружаем задачу из БД (с конфигурацией типа) + var task = await LoadTaskAsync(taskId); + + if (task == null) + { + logger.LogWarning("Задача {TaskId} не найдена в БД", taskId); + return HandlerResult.Failure("Не найдена задача"); + } + + // Атомарный захват задачи (Pending -> Processing) + var captured = await TryCaptureTaskAsync(task); + if (!captured) + { + logger.LogWarning("Не удалось захватить задачу {TaskId} (уже обрабатывается)", taskId); + return HandlerResult.Failure("Задача уже обрабатывается"); + } + + // Выполняем основную логику (реализация в наследнике) + try + { + var result = await ExecuteInternalAsync(task); + + if (result.IsSuccess) + await HandleSuccessAsync(task, result); + else + await HandleErrorAsync(task, result.Exception); + + return result; + } + catch (Exception ex) + { + logger.LogError(ex, "Непредвиденная ошибка при обработке задачи {TaskId}", taskId); + await HandleErrorAsync(task, ex); + return HandlerResult.Failure(ex.Message, ex); + } + } + + + private async Task LoadTaskAsync(Guid taskId) + { + return await taskRepository.Get() + .Include(t => t.TaskType) + .Include(t => t.TaskStatus) + .FirstOrDefaultAsync(t => t.Id == taskId); + } + + + /// + /// Атомарный захват задачи: Pending -> Processing. + /// Возвращает true, если удалось захватить. + /// + /// + /// + private async Task TryCaptureTaskAsync(TaskItem task) + { + if (task.StatusCode != TaskItemStatusEnum.Pending) + return false; + + // Атомарный захват задачи + var affectedRows = await taskRepository.TaskCaptureAsync(task.Id); + + if (affectedRows > 0) + { + // Смог захватить задачу + // Обновим локальное состояние + task.StatusCode = TaskItemStatusEnum.Processing; + task.DateModified = DateTimeOffset.UtcNow; + + return true; + } + + return false; + } + + + /// + /// Обработка успешного выполнения + /// + /// + /// + /// + private async System.Threading.Tasks.Task HandleSuccessAsync(TaskItem task, HandlerResult result) + { + logger.LogInformation("Задача {TaskId} успешно выполнена", task.Id); + + task.StatusCode = TaskItemStatusEnum.Success; + task.ProcessedAt = DateTimeOffset.UtcNow; + task.DateModified = DateTimeOffset.UtcNow; + + //todo: если есть результат, его можно сохранить, если надо + if (!string.IsNullOrEmpty(result.ResultData)) + { + //task.ResultPayload = result.ResultData; + } + + // Хук для наследников (например отправка уведомлений) + await OnSuccessAsync(task, result); + + var commitResult = await taskRepository.CommitAsync(); + + //todo: тут можно обработать commitResult + } + + + /// + /// Обработка ошибки: запись в TaskErrors, проверка MaxRetries. + /// + /// + /// + /// + private async System.Threading.Tasks.Task HandleErrorAsync(TaskItem task, Exception? ex) + { + var attemptNumber = task.RetryCount + 1; + var errorMessage = ex?.Message ?? "Неизвестная ошибка"; + var stackTrace = ex?.StackTrace; + + logger.LogError(ex, "Ошибка при обработке задачи {TaskId} (попытка {Attempt})", task.Id, attemptNumber); + + // Записываем ошибку в историю + var taskError = new TaskError + { + Id = Guid.NewGuid(), + TaskId = task.Id, + AttemptNumber = attemptNumber, + ErrorMessage = errorMessage, + StackTrace = stackTrace + }; + + await taskErrorRepository.CreateAsync(taskError); + + // Хук для наследников, например отправка уведомления + await OnErrorAsync(task, ex, attemptNumber); + + // Проверяем лимит попыток + var maxRetries = task.TaskType?.MaxRetries ?? 3;// дефолт 3 + + if (attemptNumber >= maxRetries) + { + // Лимит исчерпан, финальный Filed + logger.LogError("Превышен лимит попыток ({Max} для задачи {TaskId})", maxRetries, task.Id); + + task.StatusCode = TaskItemStatusEnum.Failed; + task.ProcessedAt = DateTimeOffset.UtcNow; + task.DateModified = DateTimeOffset.UtcNow; + } + else + { + // Остались попытки, временный Failed + logger.LogInformation("Задача {TaskId} будет повторена (попытка {Attempt} из {Max})", task.Id, attemptNumber, maxRetries); + + // Временно + task.StatusCode = TaskItemStatusEnum.Failed; + task.DateModified = DateTimeOffset.UtcNow; + } + + task.RetryCount = attemptNumber; + + //todo: может обработать commitResult? + var commitResult = await taskRepository.CommitAsync(); + + // Если есть попытки - отправляем в очередь на retry + if (attemptNumber < maxRetries) + await ScheduleRetryAsync(task); + } + + + /// + /// Планирование повторной попытки (отправка в очередь с задержкой) + /// + /// + /// + protected virtual async System.Threading.Tasks.Task ScheduleRetryAsync(TaskItem task) + { + // Здесь нужна интеграция с ITaskQueueService + // Для этого можно использовать событие или callback + // В простой реализации — выбрасываем событие, которое ловит воркер + + // Вариант: сохранить флаг, а воркер после ProcessAsync проверит и отправит + // Это будет реализовано в BackgroundService + + //todo:!!!!!!!!! + await System.Threading.Tasks.Task.Delay(50); + } + + + /// + /// Абстрактный метод - логика конкретного типа задачи. + /// Релизация в наследнике. + /// + /// + /// + protected abstract Task ExecuteInternalAsync(TaskItem task); + + + /// + /// Хук после успешного выполнения (можно переопределить в наследнике). + /// + /// + /// + /// + protected virtual System.Threading.Tasks.Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask; + + + /// + /// Хук при ошибке (можно переопределить в наследнике). + /// + /// + /// + /// + /// + protected virtual System.Threading.Tasks.Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask; + + + } +} diff --git a/PARR.Core/Services/Task/Handlers/Factory/ITaskHandlerFactory.cs b/PARR.Core/Services/Task/Handlers/Factory/ITaskHandlerFactory.cs new file mode 100644 index 00000000..013b95ef --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/Factory/ITaskHandlerFactory.cs @@ -0,0 +1,23 @@ +using PARR.Domain.Enums; + +namespace PARR.Core.Services.Task.Handlers.Factory +{ + /// + /// Фабрика обработчиков + /// + public interface ITaskHandlerFactory + { + /// + /// Получить обработчик для типа задачи + /// + /// + /// + BaseTaskHandler GetHandler(TaskTypeEnum taskType); + + /// + /// Все зарегистрированные типы. + /// + /// + IEnumerable GetSupportedTypes(); + } +} diff --git a/PARR.Core/Services/Task/Handlers/Factory/TaskHandlerFactory.cs b/PARR.Core/Services/Task/Handlers/Factory/TaskHandlerFactory.cs new file mode 100644 index 00000000..b6294859 --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/Factory/TaskHandlerFactory.cs @@ -0,0 +1,33 @@ +using Microsoft.Extensions.DependencyInjection; +using PARR.Domain.Enums; + +namespace PARR.Core.Services.Task.Handlers.Factory +{ + public class TaskHandlerFactory : ITaskHandlerFactory + { + private readonly IServiceProvider serviceProvider; + private readonly Dictionary handlers; + + public TaskHandlerFactory(IServiceProvider serviceProvider, IEnumerable handlers) + { + this.serviceProvider = serviceProvider; + this.handlers = handlers.ToDictionary(t => t.SupportedType, t => t.GetType()); + } + + public BaseTaskHandler GetHandler(TaskTypeEnum taskType) + { + if (!handlers.TryGetValue(taskType, out var handlerType)) + { + throw new InvalidOperationException($"Не зарегистрирован обработчик для типа {taskType}"); + } + + return (BaseTaskHandler)serviceProvider.GetRequiredService(handlerType); + } + + public IEnumerable GetSupportedTypes() + { + return handlers.Keys; + } + + } +} diff --git a/PARR.Core/Services/Task/Handlers/HandlerResult.cs b/PARR.Core/Services/Task/Handlers/HandlerResult.cs new file mode 100644 index 00000000..932e77b2 --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/HandlerResult.cs @@ -0,0 +1,49 @@ +namespace PARR.Core.Services.Task.Handlers +{ + /// + /// Результат выполнения задачи обработчиком. + /// Используется для преедачи данных между BaseTaskHandler и конкретными хендлерами. + /// + public class HandlerResult + { + // TODO: может его перенести в Domain + + /// + /// Успешно ли выполнено + /// + public bool IsSuccess { get; set; } + + /// + /// Данные результата (JSON, если нужно сохранить в БД) + /// + public string? ResultData { get; set; } + + + /// + /// Сообщение об ошибке (если IsSuccess = false) + /// + public string? ErrorMessage { get; set; } + + /// + /// Исключение + /// + public Exception? Exception { get; set; } + + + /// + /// Создать успешный результат + /// + /// + /// + public static HandlerResult Success(string? resultData = null) => new() { IsSuccess = true, ResultData = resultData }; + + /// + /// Создать результат с ошибкой + /// + /// + /// + /// + public static HandlerResult Failure(string errorMessage, Exception? ex = null) => new() { IsSuccess = false, ErrorMessage = errorMessage, Exception = ex }; + + } +} diff --git a/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs b/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs new file mode 100644 index 00000000..ac0662cc --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs @@ -0,0 +1,37 @@ +using Microsoft.Extensions.Logging; +using PARR.Core.Repositories.Interfaces.TaskRepositories; +using PARR.Domain.Entities.TaskEntities; +using PARR.Domain.Enums; + +namespace PARR.Core.Services.Task.Handlers +{ + /// + /// Обработчик задачаи формирования отчета по загруженности (формирование КЭШ). + /// + public class WorkloadReportHandler : BaseTaskHandler + { + + public WorkloadReportHandler(ITaskRepository taskRepository, ITaskErrorRepository taskErrorRepository, ILogger logger) + : base(taskRepository, taskErrorRepository, logger) { } + + public override TaskTypeEnum SupportedType => TaskTypeEnum.Workload; + + protected override async Task ExecuteInternalAsync(TaskItem task) + { + logger.LogInformation("Генерация КЭШ для workload, taskId: {TaskId}", task.Id); + + // 1. Десериализуем Payload + //var payload = string.IsNullOrEmpty(task.Payload) + // ? new ReportPayload() + // : JsonSerializer.Deserialize(task.Payload); + + //todo: тут логика построения отчета + + await System.Threading.Tasks.Task.CompletedTask; + + // return HandlerResult.Success(); + + return HandlerResult.Failure("Тут капец какая ошибка! Просто жуть", new Exception("А это я создал эксепшен!!!")); + } + } +} diff --git a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs new file mode 100644 index 00000000..00fd1877 --- /dev/null +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs @@ -0,0 +1,84 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using PARR.Core.Common.Interfaces; +using PARR.Core.Common.Interfaces.RabbitServices; +using PARR.Core.Services.Task.Handlers; +using PARR.Core.Services.Task.Handlers.Factory; +using PARR.Core.Services.Workload.Interfaces; +using PARR.Domain.Common.Rabbit.Messages; +using PARR.Domain.Enums; +using PARR.Domain.Settings; + +namespace PARR.Core.Services.Workload.Implementations +{ + internal class WorkloadCacheBuilderService : IWorkloadCacheBuilderService + { + private readonly ILogger logger; + private readonly IRabbitService rabbitService; + private readonly ITransformService transformService; + private readonly IServiceProvider serviceProvider; + + public WorkloadCacheBuilderService( + ILogger logger, + IRabbitService rabbitService, + ITransformService transformService, + IServiceProvider serviceProvider + ) + { + this.logger = logger; + this.rabbitService = rabbitService; + this.transformService = transformService; + this.serviceProvider = serviceProvider; + } + + + public async System.Threading.Tasks.Task ProcessMessagesAsync(IMqSettings mqSettings) + { + // этот метод вызывается в воркере + + var isConnected = await rabbitService.InitConsumerAsync(mqSettings, async (msg) => + { + logger.LogInformation($"Получили запрос: {msg}"); + + var queueMessage = transformService.GetModelFromJson(msg); + if (queueMessage == null) + return; + + // Проверим тип сообщения, что он соответствует текущему сервису + if (queueMessage.TypeCode != TaskTypeEnum.Workload) + { + logger.LogError("Несоответствие типа: в очереди {Queue} получена задача типа {Actual}", mqSettings.QueueName, queueMessage.TypeCode); + return; + } + + using (var scope = serviceProvider.CreateScope()) + { + var handlerFactory = scope.ServiceProvider.GetRequiredService(); + var handler = handlerFactory.GetHandler(queueMessage.TypeCode); + + var result = await handler.ProcessAsync(queueMessage.TaskId); + + // Если ошибка и есть retry — отправляем в очередь с задержкой + if (!result.IsSuccess && handler is BaseTaskHandler baseHandler) + { + // TODO: + // Если ошибка и есть retry — отправляем в очередь с задержкой + // (это упрощенная реализация, в идеале — через ITaskQueueService) + + // Здесь нужна логика retry с задержкой + // Можно реализовать через отдельный сервис + } + } + }); + + if (!isConnected) + throw new Exception("Ошибка при подключении к RabbitMq"); + } + + public async System.Threading.Tasks.Task StopProcessingAsync() + { + await rabbitService.DisposeAsync(); + } + + } +} diff --git a/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs new file mode 100644 index 00000000..0d658045 --- /dev/null +++ b/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs @@ -0,0 +1,11 @@ +using PARR.Core.Common.Interfaces.RabbitServices.Workers; + +namespace PARR.Core.Services.Workload.Interfaces +{ + public interface IWorkloadCacheBuilderService : IMessageConsumer + { + + } +} + + diff --git a/PARR.DAL/PARR.DAL.csproj b/PARR.DAL/PARR.DAL.csproj index 234b2959..f7b981c5 100644 --- a/PARR.DAL/PARR.DAL.csproj +++ b/PARR.DAL/PARR.DAL.csproj @@ -7,6 +7,7 @@ + all runtime; build; native; contentfiles; analyzers; buildtransitive diff --git a/PARR.DAL/Repositories/Base/BaseRepository.cs b/PARR.DAL/Repositories/Base/BaseRepository.cs index c5f41def..a620c9ec 100644 --- a/PARR.DAL/Repositories/Base/BaseRepository.cs +++ b/PARR.DAL/Repositories/Base/BaseRepository.cs @@ -322,7 +322,8 @@ namespace PARR.DAL.Repositories.Base { logger.LogDebug("Начинаю создание объекта типа {EntityType}", typeof(T).Name); - obj.DateCreated = DateTimeOffset.UtcNow; + if (obj.DateCreated == DateTimeOffset.MinValue) + obj.DateCreated = DateTimeOffset.UtcNow; try { diff --git a/PARR.DAL/Repositories/TaskRepositories/TaskRepository.cs b/PARR.DAL/Repositories/TaskRepositories/TaskRepository.cs index 142983dc..4b638af4 100644 --- a/PARR.DAL/Repositories/TaskRepositories/TaskRepository.cs +++ b/PARR.DAL/Repositories/TaskRepositories/TaskRepository.cs @@ -1,8 +1,11 @@ -using Microsoft.Extensions.Logging; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; using PARR.Core.Repositories.Interfaces.TaskRepositories; using PARR.DAL.Context; using PARR.DAL.Repositories.Base; +using PARR.Domain.Constants; using PARR.Domain.Entities.TaskEntities; +using PARR.Domain.Enums; namespace PARR.DAL.Repositories.TaskRepositories { @@ -10,7 +13,27 @@ namespace PARR.DAL.Repositories.TaskRepositories { public TaskRepository(DataContext dataContext, ILogger logger) : base(logger, dataContext) { } + public async Task TaskCaptureAsync(Guid taskId) + { + //var affectedRows = await EntityContext.Database.ExecuteSqlInterpolatedAsync($@" + // UPDATE ""{DatabaseSchemas.Task}"".""Tasks"" + // SET ""{nameof(TaskItem.StatusCode)}""={(int)TaskItemStatusEnum.Processing}, + // ""{nameof(TaskItem.DateModified)}""={DateTimeOffset.UtcNow} + // WHERE ""{nameof(TaskItem.Id)}""={taskId} + // AND ""{nameof(TaskItem.StatusCode)}""={(int)TaskItemStatusEnum.Pending} + //"); + + var sql = $@" + UPDATE ""{DatabaseSchemas.Task}"".""Tasks"" + SET ""{nameof(TaskItem.StatusCode)}"" = {(int)TaskItemStatusEnum.Processing}, + ""{nameof(TaskItem.DateModified)}"" = {{0}} + WHERE ""{nameof(TaskItem.Id)}"" = {{1}} + AND ""{nameof(TaskItem.StatusCode)}"" = {(int)TaskItemStatusEnum.Pending}"; + + var affectedRows = await EntityContext.Database.ExecuteSqlRawAsync(sql, DateTimeOffset.UtcNow, taskId); + + return affectedRows; + } - // метод атомарного взятия в работу } } diff --git a/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj b/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj new file mode 100644 index 00000000..fb4489ea --- /dev/null +++ b/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj @@ -0,0 +1,25 @@ + + + + net7.0 + enable + enable + dotnet-PARR.WorkloadBuilderWorker-5cefb704-873b-40e7-b8a0-0537ca29e5ed + + + + + + + + + + + + + + + + + + diff --git a/PARR.WorkloadBuilderWorker/Program.cs b/PARR.WorkloadBuilderWorker/Program.cs new file mode 100644 index 00000000..94d5ff05 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Program.cs @@ -0,0 +1,45 @@ +using Elastic.CommonSchema.Serilog; +using PARR.Core; +using PARR.DAL; +using PARR.DAL.Services.Interfaces.Job; +using PARR.Infrastructure; +using PARR.WorkloadBuilderWorker; +using PARR.WorkloadBuilderWorker.Settings; +using Serilog; + +var builder = Host.CreateApplicationBuilder(); + +builder.Services.AddLogging(config => +{ + config.ClearProviders(); + + var logger = new LoggerConfiguration(); + + if (builder.Environment.IsProduction()) + logger.WriteTo.Console(new EcsTextFormatter()); + else + logger.WriteTo.Console(); + + logger.ReadFrom.Configuration(builder.Configuration); + + config.AddSerilog(logger.CreateLogger()); +}); + +#region PARR modules +builder.Services.InstallDalServices(builder.Configuration); +builder.Services.AddCoreServices(builder.Configuration); +builder.Services.AddInfrastructureServices(builder.Configuration); + +builder.Configuration.AddDalConfigurations(builder.Services); +builder.Services.AddDallSettings(builder.Configuration); +#endregion + +var workerSettings = new WorkerSettings(); +builder.Configuration.GetSection(nameof(WorkerSettings)).Bind(workerSettings); +builder.Services.AddSingleton(workerSettings); + + +builder.Services.AddHostedService(); + +var host = builder.Build(); +host.Run(); diff --git a/PARR.WorkloadBuilderWorker/Properties/launchSettings.json b/PARR.WorkloadBuilderWorker/Properties/launchSettings.json new file mode 100644 index 00000000..1187fd6b --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Properties/launchSettings.json @@ -0,0 +1,11 @@ +{ + "profiles": { + "PARR.WorkloadBuilderWorker": { + "commandName": "Project", + "dotnetRunMessages": true, + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs b/PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs new file mode 100644 index 00000000..5ab5e0b0 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs @@ -0,0 +1,9 @@ +using PARR.Domain.Settings; + +namespace PARR.WorkloadBuilderWorker.Settings +{ + public record WorkerSettings + { + public MqSettingsBase MqSettings { get; init; } = new(); + } +} diff --git a/PARR.WorkloadBuilderWorker/Worker.cs b/PARR.WorkloadBuilderWorker/Worker.cs new file mode 100644 index 00000000..730de234 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Worker.cs @@ -0,0 +1,29 @@ +using PARR.Core.Services.Workload.Interfaces; +using PARR.WorkloadBuilderWorker.Settings; + +namespace PARR.WorkloadBuilderWorker +{ + public class Worker : BackgroundService + { + private readonly IWorkloadCacheBuilderService workloadCacheBuilderService; + private readonly WorkerSettings workerSettings; + + public Worker(ILogger logger, IWorkloadCacheBuilderService workloadCacheBuilderService, WorkerSettings WorkerSettings) + { + this.workloadCacheBuilderService = workloadCacheBuilderService; + this.workerSettings = WorkerSettings; + } + + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + await workloadCacheBuilderService.ProcessMessagesAsync(workerSettings.MqSettings); + } + + public override async Task StopAsync(CancellationToken cancellationToken) + { + await workloadCacheBuilderService.StopProcessingAsync(); + await base.StopAsync(cancellationToken); + } + } +} diff --git a/PARR.WorkloadBuilderWorker/appsettings.Development.json b/PARR.WorkloadBuilderWorker/appsettings.Development.json new file mode 100644 index 00000000..7d01cc26 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/appsettings.Development.json @@ -0,0 +1,35 @@ +{ + "ConnectionStrings": { + "RedisConnection": "10.99.253.216:6379,password=ParrP@ssPtk202MMdevDvs" + }, + "Logging": { + "LogLevel": { + "Default": "Debug", + "Microsoft.Hosting.Lifetime": "Information" + } + }, + "Serilog": { + "MinimumLevel": { + "Default": "Debug", + "Override": { + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + }, + "WriteTo": [ + { + "Name": "File", + "Args": { + "path": "log/log-.txt", + "rollingInterval": "Day" + } + } + ] + }, + "WorkerSettings": { + "MqSettings": { + "HostName": "10.99.253.216" + } + } + +} diff --git a/PARR.WorkloadBuilderWorker/appsettings.json b/PARR.WorkloadBuilderWorker/appsettings.json new file mode 100644 index 00000000..496c9590 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/appsettings.json @@ -0,0 +1,33 @@ +{ + "ConnectionStrings": { + "DefaultConnection": "Server=10.99.253.184;Database=parr;User Id=app_parr; Password=PosdfkhT&)%sdfligL&%5546;", + "RedisConnection": "parr-redis:6379,password=ParrP@ssPtk202MMdevDvs" + }, + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft.EntityFrameworkCore": "Error", + "Microsoft.EntityFrameworkCore.Database.Command": "Warning", + "Microsoft.AspNetCore": "Warning" + } + }, + "Serilog": { + "MinimumLevel": { + "Default": "Information", + "Override": { + "Microsoft": "Warning", + "Microsoft.EntityFrameworkCore": "Error", + "Microsoft.EntityFrameworkCore.Database.Command": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } + }, + "WorkerSettings": { + "MqSettings": { + "HostName": "parr-rabbitmq", + "QueueName": "parr-task-workload", + "User": "task_workload_reader", + "Password": "sdjfgIYFdifyKJFDhvsad@!314d" + } + } +} From 7da2265d1cbb9023772691461b12a0ad46d6864a Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Wed, 22 Apr 2026 16:13:23 +1000 Subject: [PATCH 3/6] =?UTF-8?q?feat(api,=20Core,=20Domain,=20Dal):=20?= =?UTF-8?q?=D0=94=D0=BB=D1=8F=20entity=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB?= =?UTF-8?q?=D0=B5=D0=B9,=20=D0=B4=D0=BB=D1=8F=20DateModified=20=D1=81?= =?UTF-8?q?=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=20=D0=B0=D1=82=D1=80=D0=B8=D0=B1?= =?UTF-8?q?=D1=83=D1=82=20ManualControl=20-=20=D0=BE=D1=82=D0=BA=D0=BB?= =?UTF-8?q?=D1=8E=D1=87=D0=B0=D0=B5=D1=82=20=D0=B0=D0=B2=D1=82=D0=BE=D0=BC?= =?UTF-8?q?=D0=B0=D1=82=D0=B8=D1=87=D0=B5=D1=81=D0=BA=D0=BE=D0=B5=20=D1=83?= =?UTF-8?q?=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=B4?= =?UTF-8?q?=D0=B0=D1=82=D0=BE=D0=B9.=20TaskReconciliationService=20-=20?= =?UTF-8?q?=D1=81=D0=B5=D1=80=D0=B2=D0=B8=D1=81=20=D0=BF=D0=BE=20=D1=83?= =?UTF-8?q?=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D1=8E=20=D0=B7?= =?UTF-8?q?=D0=B0=D0=B4=D0=B0=D1=87=D0=B0=D0=BC=D0=B8=20=D0=B2=20=D1=81?= =?UTF-8?q?=D1=82=D0=B0=D1=82=D1=83=D1=81=D0=B0=D1=85=20=D0=BE=D1=88=D0=B8?= =?UTF-8?q?=D0=B1=D0=BA=D0=B8,=20=D0=B7=D0=B0=D0=B2=D0=B8=D1=81=D0=BB?= =?UTF-8?q?=D0=B0.=20ITaskMqSettingsProvider=20-=20=D0=BF=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=B0=D0=B9=D0=B4=D0=B5=D1=80=20=D0=B4=D0=BB=D1=8F=20?= =?UTF-8?q?=D1=83=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D1=8F=20?= =?UTF-8?q?=D0=BD=D0=B0=D1=81=D1=82=D1=80=D0=BE=D0=B9=D0=BA=D0=B0=D0=BC?= =?UTF-8?q?=D0=B8=20MQ=20=D0=B2=20=D0=B7=D0=B0=D0=B2=D0=B8=D1=81=D0=B8?= =?UTF-8?q?=D0=BC=D0=BE=D1=81=D1=82=D0=B8=20=D0=BE=D1=82=20TaskType?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../V1/Statistics/StatWorkloadController.cs | 25 ++- PARR.API/Program.cs | 21 +- PARR.Core/DependencyInjection.cs | 50 ++++- .../Services/Task/Handlers/BaseTaskHandler.cs | 11 ++ .../Task/Handlers/WorkloadReportHandler.cs | 2 + .../TaskReconciliationService.cs | 187 ++++++++++++++++++ .../Interfaces/ITaskReconciliationService.cs | 17 ++ .../Task/Models/ReconciliationReport.cs | 55 ++++++ .../Task/Providers/ITaskMqSettingsProvider.cs | 18 ++ .../Task/Providers/TaskMqSettingsProvider.cs | 22 +++ .../WorkloadCacheBuilderService.cs | 22 ++- PARR.DAL/Repositories/Base/BaseRepository.cs | 26 ++- .../Attributes/ManualControlAttribute.cs | 11 ++ .../Entities/Base/IBaseEntityDateModified.cs | 5 + PARR.Domain/Entities/TaskEntities/TaskItem.cs | 2 + .../Settings/IReconciliationSettings.cs | 21 ++ 16 files changed, 465 insertions(+), 30 deletions(-) create mode 100644 PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs create mode 100644 PARR.Core/Services/Task/Interfaces/ITaskReconciliationService.cs create mode 100644 PARR.Core/Services/Task/Models/ReconciliationReport.cs create mode 100644 PARR.Core/Services/Task/Providers/ITaskMqSettingsProvider.cs create mode 100644 PARR.Core/Services/Task/Providers/TaskMqSettingsProvider.cs create mode 100644 PARR.Domain/Entities/Attributes/ManualControlAttribute.cs create mode 100644 PARR.Domain/Settings/IReconciliationSettings.cs diff --git a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs index 626ed11f..d9fedf61 100644 --- a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs +++ b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs @@ -6,6 +6,7 @@ using PARR.API.Controllers.V1.Base; using PARR.API.Services.Interfaces; using PARR.API.Settings; using PARR.Core.Services.Task.Interfaces; +using PARR.Core.Services.Task.Providers; using PARR.Domain.Common.Roles; using PARR.Domain.Entities.Base.History; using PARR.Domain.Enums; @@ -20,20 +21,23 @@ namespace PARR.API.Controllers.V1.Statistics { private readonly ITaskManagementService taskManagementService; private readonly IClientService clientService; - private readonly MqSettings mqSettings; + //private readonly MqSettings mqSettings; private readonly ILogger logger; + private readonly ITaskMqSettingsProvider taskMqSettingsProvider; public StatWorkloadController( ITaskManagementService taskManagementService, IClientService clientService, - MqSettings mqSettings, - ILogger logger + //MqSettings mqSettings, + ILogger logger, + ITaskMqSettingsProvider taskMqSettingsProvider ) { this.taskManagementService = taskManagementService; this.clientService = clientService; - this.mqSettings = mqSettings; + //this.mqSettings = mqSettings; this.logger = logger; + this.taskMqSettingsProvider = taskMqSettingsProvider; } @@ -55,13 +59,14 @@ namespace PARR.API.Controllers.V1.Statistics try { - mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings); + //mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings); - if (queueSettings == null) - { - logger.LogError("Не удалось получить настройки очереди для TaskType: {TaskType}", TaskTypeEnum.Workload); - throw new InvalidOperationException($"Не удалось получить настройки очереди для TaskTypeEnum.Workload"); - } + //if (queueSettings == null) + //{ + // logger.LogError("Не удалось получить настройки очереди для TaskType: {TaskType}", TaskTypeEnum.Workload); + // throw new InvalidOperationException($"Не удалось получить настройки очереди для TaskTypeEnum.Workload"); + //} + var queueSettings = taskMqSettingsProvider.GetSettings(TaskTypeEnum.Workload); var taskId = await taskManagementService.CreateTaskAsync(TaskTypeEnum.Workload, default, initiator, queueSettings); diff --git a/PARR.API/Program.cs b/PARR.API/Program.cs index 3c7e520e..a4a8cca8 100644 --- a/PARR.API/Program.cs +++ b/PARR.API/Program.cs @@ -3,8 +3,11 @@ using FluentValidation; using Microsoft.AspNetCore.HttpOverrides; using PARR.API.Authentication; using PARR.API.Installers; +using PARR.API.Settings; using PARR.Core; using PARR.DAL; +using PARR.Domain.Enums; +using PARR.Domain.Settings; using PARR.Infrastructure; using Serilog; using System.Reflection; @@ -21,14 +24,30 @@ builder.Host.UseSerilog((context, config) => config.ReadFrom.Configuration(builder.Configuration); }); +builder.Services.InstallSettings(builder.Configuration); // Add services to the container. builder.Services.InstallDalServices(builder.Configuration); //---- new ---- builder.Services.AddCoreServices(builder.Configuration); +builder.Services.AddTaskManagement(mqSettingsFactory: (provider, taskType) => +{ + var taskSettingsMq = provider.GetRequiredService(); + + // swith для примера + //return taskType switch + //{ + // TaskTypeEnum.Workload => MqSettingsBase{ }, + // _ => throw new ArgumentException($"Неизвестный тип задания {taskType}. Не смог для него получить настройки") + //}; + + if (taskSettingsMq.Tasks.TryGetValue(taskType, out var settings)) + return settings; + + throw new ArgumentException($"Настройки для типа задания {taskType} не найдены в конфигурации."); +}); builder.Services.AddInfrastructureServices(builder.Configuration); //---- --- ---- builder.Services.InstallApiServices(builder.Configuration); -builder.Services.InstallSettings(builder.Configuration); builder.Services.AddAutoMapper(AppDomain.CurrentDomain.GetAssemblies()); // Dal configuration diff --git a/PARR.Core/DependencyInjection.cs b/PARR.Core/DependencyInjection.cs index 84bf1e37..80a5b702 100644 --- a/PARR.Core/DependencyInjection.cs +++ b/PARR.Core/DependencyInjection.cs @@ -7,8 +7,10 @@ using PARR.Core.Services.Task.Handlers; using PARR.Core.Services.Task.Handlers.Factory; using PARR.Core.Services.Task.Implementations; using PARR.Core.Services.Task.Interfaces; +using PARR.Core.Services.Task.Providers; using PARR.Core.Services.Workload.Implementations; using PARR.Core.Services.Workload.Interfaces; +using PARR.Domain.Enums; using PARR.Domain.Settings; namespace PARR.Core @@ -43,15 +45,20 @@ namespace PARR.Core #region Task - // Регистрация хендлеров (нужно регистировать каждый отдельно) - services.AddScoped(); - services.AddScoped(sp => sp.GetRequiredService());// регим для фабрики - // Еще хэндлеры... + // ---------> Регистрируем в AddTaskManagement <--------- - // Фабрика хэндлеров - services.AddScoped(); + //// Регистрация хендлеров (нужно регистировать каждый отдельно) + //services.AddScoped(); + //services.AddScoped(sp => sp.GetRequiredService());// регим для фабрики + //// Еще хэндлеры... - services.AddScoped(); + //// Фабрика хэндлеров + //services.AddScoped(); + + //services.AddScoped(); + + //// Провайдер настроек MQ + //services.AddSingleton(); #endregion @@ -71,6 +78,35 @@ namespace PARR.Core return services; } + /// + /// Регистрация сервисов по управлению задчами (Task) + /// + /// + /// Настройки MQ для всех TaskTypeEnum + /// + /// + public static IServiceCollection AddTaskManagement(this IServiceCollection services, Func mqSettingsFactory) + { + if (mqSettingsFactory == null) + throw new ArgumentNullException(nameof(mqSettingsFactory), "Ты забыл зарегистрировать настройки MQ для Task! Бестолочь! :)"); + + // Регистрация хендлеров (нужно регистировать каждый отдельно) + services.AddScoped(); + services.AddScoped(sp => sp.GetRequiredService());// регим для фабрики хэндлеров + // Еще хэндлеры... + + // Фабрика хэндлеров + services.AddScoped(); + + services.AddScoped(); + + // Провайдер настроек MQ + services.AddTransient>(sp => (taskType) => mqSettingsFactory(sp, taskType)); // регим фабрику настроек для ITaskMqSettingsProvider + services.AddSingleton(); + + return services; + } + private static IServiceCollection AddBllMapping(this IServiceCollection services) { diff --git a/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs index f8287d5f..7274863c 100644 --- a/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs +++ b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs @@ -97,7 +97,12 @@ namespace PARR.Core.Services.Task.Handlers private async Task TryCaptureTaskAsync(TaskItem task) { if (task.StatusCode != TaskItemStatusEnum.Pending) + { + // Логируем, но не считаем ошибкой + logger.LogDebug("Задача {TaskId} не в Pending (статус {Status}), пропускаем", task.Id, task.StatusCode); return false; + } + // Атомарный захват задачи var affectedRows = await taskRepository.TaskCaptureAsync(task.Id); @@ -112,6 +117,8 @@ namespace PARR.Core.Services.Task.Handlers return true; } + logger.LogDebug("Задача {TaskId} уже захвачена другим воркером", task.Id); + return false; } @@ -221,6 +228,10 @@ namespace PARR.Core.Services.Task.Handlers // Вариант: сохранить флаг, а воркер после ProcessAsync проверит и отправит // Это будет реализовано в BackgroundService + //todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! + + // тоже самое в WorkloadCacheBuilderService + //todo:!!!!!!!!! await System.Threading.Tasks.Task.Delay(50); } diff --git a/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs b/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs index ac0662cc..708ca3e5 100644 --- a/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs +++ b/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs @@ -27,6 +27,8 @@ namespace PARR.Core.Services.Task.Handlers //todo: тут логика построения отчета + + await System.Threading.Tasks.Task.CompletedTask; // return HandlerResult.Success(); diff --git a/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs new file mode 100644 index 00000000..68991449 --- /dev/null +++ b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs @@ -0,0 +1,187 @@ +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; +using PARR.Core.Common.Interfaces.RabbitServices; +using PARR.Core.Repositories.Interfaces.TaskRepositories; +using PARR.Core.Services.Task.Interfaces; +using PARR.Core.Services.Task.Models; +using PARR.Core.Services.Task.Providers; +using PARR.Domain.Common.Rabbit.Messages; +using PARR.Domain.Entities.TaskEntities; +using PARR.Domain.Enums; +using PARR.Domain.Settings; + +namespace PARR.Core.Services.Task.Implementations +{ + internal class TaskReconciliationService : ITaskReconciliationService + { + private readonly ILogger logger; + private readonly ITaskTypeRepository taskTypeRepository; + private readonly ITaskRepository taskRepository; + private readonly ITaskErrorRepository taskErrorRepository; + private readonly IReconciliationSettings reconciliationSettings; + private readonly IRabbitService rabbitService; + private readonly ITaskMqSettingsProvider taskMqSettingsProvider; + + public TaskReconciliationService( + ILogger logger, + ITaskTypeRepository taskTypeRepository, + ITaskRepository taskRepository, + ITaskErrorRepository taskErrorRepository, + IReconciliationSettings reconciliationSettings, + IRabbitService rabbitService, + ITaskMqSettingsProvider taskMqSettingsProvider + ) + { + this.logger = logger; + this.taskTypeRepository = taskTypeRepository; + this.taskRepository = taskRepository; + this.taskErrorRepository = taskErrorRepository; + this.reconciliationSettings = reconciliationSettings; + this.rabbitService = rabbitService; + this.taskMqSettingsProvider = taskMqSettingsProvider; + } + + public async Task RunAsync() + { + var report = new ReconciliationReport + { + StartedAt = DateTimeOffset.UtcNow + }; + + logger.LogInformation("Начало выполнения Reconciliation Job."); + + // Загружаем все типы задач (для получения максимальной длительности) + var taskTypes = await taskTypeRepository.Get().AsNoTracking().ToDictionaryAsync(t => t.Code, t => t); + + if (taskTypes.Count == 0) + { + logger.LogWarning("Типы задач не настроены в БД (таблица {TaskTypes})", nameof(TaskType)); + report.Messages.Add("TaskTypes пустые. Не настроены в БД."); + return report; + } + + // Для каждого типа ищем зависшие задачи + foreach (var taskType in taskTypes.Values) + { + var stuckTasks = await FindStuckTasksAsync(taskType); + report.CheckedCount += stuckTasks.Count; + + foreach (var task in stuckTasks) + { + await ProcessStuckTaskAsync(task, taskType, report); + } + } + + report.CompletedAt = DateTimeOffset.UtcNow; + + if (report.HasChanged) + logger.LogInformation("Reconciliation Job завершен. Проверено: {Checked}, Исправлено: {Fixed}, Retry: {Retried}, Failed: {Failed}. Длительность: {Duration}", + report.CheckedCount, report.FixedCount, report.RetriedCount, report.FailedCount, report.Duration); + else + logger.LogDebug("Reconciliation Job завершен. Зависших задач не найдено"); + + return report; + } + + + /// + /// Поиск зависших задач для конкретного типа. + /// + /// + /// + private async Task> FindStuckTasksAsync(TaskType taskType) + { + // Вычисляем порог: задача зависла, если выполняется дольше чем MaxExecutionTimeMinutes + var timeoutThreshold = DateTimeOffset.UtcNow.AddMinutes(-taskType.MaxExecutionTimeMinutes); + + logger.LogDebug("Поиск зависших задач типа {Type}. Порог: {Threshold}", taskType.Code, timeoutThreshold); + + return await taskRepository.Get() + .Where(t => + t.TypeCode == taskType.Code && t.DateModified < timeoutThreshold && + ( + // Зависла в Processing + t.StatusCode == TaskItemStatusEnum.Processing + //Зависли в Pending + || (t.StatusCode == TaskItemStatusEnum.Pending && t.RetryCount > 0) + ) + ) + .OrderBy(t => t.DateModified) // Сначала самые старые + .Take(reconciliationSettings.MaxTaskPerRun) // Защита от перегрузки + .ToListAsync(); + } + + /// + /// Обработка одной зависшей задачи. + /// + /// + /// + /// + /// + private async System.Threading.Tasks.Task ProcessStuckTaskAsync(TaskItem task, TaskType taskType, ReconciliationReport report) + { + logger.LogWarning("Обнаружена одна зависшая задача {TaskId} (тип {Type}). Последнее обновление: {LastModified}", task.Id, task.TypeCode, task.DateModified); + + // Записываем ошибку в историю + var error = new TaskError + { + Id = Guid.NewGuid(), + TaskId = task.Id, + AttemptNumber = task.RetryCount + 1, // показываем планируемую попытку + ErrorMessage = $"Превышен порог выполнения задачи. LastDateModified: {task.DateModified}" + }; + + await taskErrorRepository.CreateAsync(error); + await taskErrorRepository.CommitAsync(); + + // Если лимит исчерпан - ставим Failed + if (task.RetryCount >= taskType.MaxRetries) + { + task.StatusCode = TaskItemStatusEnum.Failed; + task.ProcessedAt = DateTimeOffset.UtcNow; + task.DateModified = DateTimeOffset.UtcNow; + + await taskRepository.CommitAsync(); + + report.FailedCount++; + report.Messages.Add($"Задача {task.Id}: закончились попытки."); + + return; + } + + // Сначала пытаемся отправить в очередь + var mqSettings = taskMqSettingsProvider.GetSettings(taskType.Code); + + var message = new TaskMessage { TaskId = task.Id, TypeCode = task.TypeCode }; + var sendResult = await rabbitService.SendAsync(mqSettings, new List { message }); + + if (!sendResult.IsSuccess) + { + // Очередь недоступна — НЕ меняем БД, НЕ тратим RetryCount + // Задача останется в текущем статусе и будет найдена в следующий раз + + logger.LogError("Не удалось отправить задачу {TaskId} в очередь. Повторим в следующей итерации.", task.Id); + report.Messages.Add($"Задача {task.Id}: не смогла отправиться в очередь, обработается в следующей итерации."); + return; // Выходим без сохранения изменений + } + + // Очередь успешна - обновляем БД + // Ставим задачу в ожидание + task.StatusCode = TaskItemStatusEnum.Pending; + + // Может получиться так, что воркер возьмет задачу быстрее, чем в БД у нее установится статус Pending. + // Если это произойдет, ничего страшного, воркер это не отработает. Задача опять попадет в ProcessStuckTaskAsync и повторно отправится в очередь. + + task.RetryCount++; + task.DateModified = DateTimeOffset.UtcNow; + + await taskRepository.CommitAsync(); + + report.RetriedCount++; + report.Messages.Add($"Задача {task.Id}: опубликована в очередь, попытка #{task.RetryCount}"); + report.FixedCount++; + } + + + } +} diff --git a/PARR.Core/Services/Task/Interfaces/ITaskReconciliationService.cs b/PARR.Core/Services/Task/Interfaces/ITaskReconciliationService.cs new file mode 100644 index 00000000..69d6e02f --- /dev/null +++ b/PARR.Core/Services/Task/Interfaces/ITaskReconciliationService.cs @@ -0,0 +1,17 @@ +using PARR.Core.Services.Task.Models; + +namespace PARR.Core.Services.Task.Interfaces +{ + /// + /// Сервис для проверки и исправления зависших задач. + /// Запускается фоновой задачей по расписанию. + /// + internal interface ITaskReconciliationService + { + /// + /// Выполнить один проход проверки + /// + /// + Task RunAsync(); + } +} diff --git a/PARR.Core/Services/Task/Models/ReconciliationReport.cs b/PARR.Core/Services/Task/Models/ReconciliationReport.cs new file mode 100644 index 00000000..c910409d --- /dev/null +++ b/PARR.Core/Services/Task/Models/ReconciliationReport.cs @@ -0,0 +1,55 @@ +namespace PARR.Core.Services.Task.Models +{ + /// + /// Отчет о выполнении Reconciliation Job. + /// Используется для мониторинга и логирования. + /// + internal class ReconciliationReport + { + /// + /// Время начала выполнения + /// + public DateTimeOffset StartedAt { get; set; } + + /// + /// Время завершения выполнения + /// + public DateTimeOffset CompletedAt { get; set; } + + /// + /// Сколько задач проверено + /// + public int CheckedCount { get; set; } + + /// + /// Сколько задач исправлено (переведено в retry или failed) + /// + public int FixedCount { get; set; } + + /// + /// Сколько задач отправлено на retry + /// + public int RetriedCount { get; set; } + + /// + /// Сколько задач переведено в финальный failed + /// + public int FailedCount { get; set; } + + /// + /// Сообщения/ошибки (для логирования) + /// + public List Messages { get; set; } = new List(); + + /// + /// Длительность выполнения + /// + public TimeSpan Duration => CompletedAt - StartedAt; + + /// + /// Были ли какие-то действия (повторные отправки, перевод в ошибку) + /// + public bool HasChanged => FixedCount > 0; + + } +} diff --git a/PARR.Core/Services/Task/Providers/ITaskMqSettingsProvider.cs b/PARR.Core/Services/Task/Providers/ITaskMqSettingsProvider.cs new file mode 100644 index 00000000..8b67a6b1 --- /dev/null +++ b/PARR.Core/Services/Task/Providers/ITaskMqSettingsProvider.cs @@ -0,0 +1,18 @@ +using PARR.Domain.Enums; +using PARR.Domain.Settings; + +namespace PARR.Core.Services.Task.Providers +{ + /// + /// Провайдер настроек очереди для задач Tasks + /// + public interface ITaskMqSettingsProvider + { + /// + /// Получить настройки очереди для типа задач + /// + /// + /// + IMqSettings GetSettings(TaskTypeEnum typeCode); + } +} diff --git a/PARR.Core/Services/Task/Providers/TaskMqSettingsProvider.cs b/PARR.Core/Services/Task/Providers/TaskMqSettingsProvider.cs new file mode 100644 index 00000000..5a3f48b6 --- /dev/null +++ b/PARR.Core/Services/Task/Providers/TaskMqSettingsProvider.cs @@ -0,0 +1,22 @@ +using PARR.Domain.Enums; +using PARR.Domain.Settings; + +namespace PARR.Core.Services.Task.Providers +{ + internal class TaskMqSettingsProvider : ITaskMqSettingsProvider + { + private readonly Func settingsFactory; + + public TaskMqSettingsProvider( + Func settingsFactory + ) + { + this.settingsFactory = settingsFactory; + } + + public IMqSettings GetSettings(TaskTypeEnum typeCode) + { + return settingsFactory(typeCode); + } + } +} diff --git a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs index 00fd1877..57715328 100644 --- a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs @@ -58,16 +58,20 @@ namespace PARR.Core.Services.Workload.Implementations var result = await handler.ProcessAsync(queueMessage.TaskId); - // Если ошибка и есть retry — отправляем в очередь с задержкой - if (!result.IsSuccess && handler is BaseTaskHandler baseHandler) - { - // TODO: - // Если ошибка и есть retry — отправляем в очередь с задержкой - // (это упрощенная реализация, в идеале — через ITaskQueueService) + //// Если ошибка и есть retry — отправляем в очередь с задержкой + //if (!result.IsSuccess && handler is BaseTaskHandler baseHandler) + //{ + // // TODO: + // // Если ошибка и есть retry — отправляем в очередь с задержкой + // // (это упрощенная реализация, в идеале — через ITaskQueueService) - // Здесь нужна логика retry с задержкой - // Можно реализовать через отдельный сервис - } + // // Здесь нужна логика retry с задержкой + // // Можно реализовать через отдельный сервис + + // //todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! -> + // // тоже самое в BaseTaskHandler + + //} } }); diff --git a/PARR.DAL/Repositories/Base/BaseRepository.cs b/PARR.DAL/Repositories/Base/BaseRepository.cs index a620c9ec..7eedce14 100644 --- a/PARR.DAL/Repositories/Base/BaseRepository.cs +++ b/PARR.DAL/Repositories/Base/BaseRepository.cs @@ -1,9 +1,11 @@ using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.ChangeTracking; +using Microsoft.EntityFrameworkCore.Metadata.Internal; using Microsoft.Extensions.Logging; using PARR.Core.Repositories.Base; using PARR.DAL.Context; using PARR.Domain.Common.Pagination; +using PARR.Domain.Entities.Attributes; using PARR.Domain.Entities.Base; using PARR.Domain.Entities.Base.History; using PARR.Domain.Entities.Base.History.Base; @@ -133,10 +135,28 @@ namespace PARR.DAL.Repositories.Base /// private void DateModifiedResolver(EntityEntry obj) { - if (obj.Entity is IBaseEntityDateModified) + if (obj.Entity is IBaseEntityDateModified entity) { - logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name); - (obj.Entity as IBaseEntityDateModified)!.DateModified = DateTimeOffset.UtcNow; + //logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name); + //(obj.Entity as IBaseEntityDateModified)!.DateModified = DateTimeOffset.UtcNow; + + var entityType = entity.GetType(); + + //Ищем DateModified + var properyInfo = entityType.GetProperty(nameof(IBaseEntityDateModified.DateModified)); + + // Проверяем наличие атрибута ManualControlAttribute на этом свойстве + bool isManual = properyInfo?.GetCustomAttributes(typeof(ManualControlAttribute), false).Any() ?? false; + + if (!isManual) + { + logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name); + entity.DateModified = DateTimeOffset.UtcNow; + } + else + { + logger.LogDebug("Пропуск обновления DateModified (ManualControl) для {EntityType}", entityType.Name); + } } } diff --git a/PARR.Domain/Entities/Attributes/ManualControlAttribute.cs b/PARR.Domain/Entities/Attributes/ManualControlAttribute.cs new file mode 100644 index 00000000..8e9f1d1e --- /dev/null +++ b/PARR.Domain/Entities/Attributes/ManualControlAttribute.cs @@ -0,0 +1,11 @@ +namespace PARR.Domain.Entities.Attributes +{ + /// + /// Отключает автоматическое управление DateModified. + /// + [AttributeUsage(AttributeTargets.Property)] + public class ManualControlAttribute : Attribute + { + // позже, можно применять к другим свойствам, если придумаем зачем и реализуем логику + } +} diff --git a/PARR.Domain/Entities/Base/IBaseEntityDateModified.cs b/PARR.Domain/Entities/Base/IBaseEntityDateModified.cs index 0c3447c8..b3a4ade3 100644 --- a/PARR.Domain/Entities/Base/IBaseEntityDateModified.cs +++ b/PARR.Domain/Entities/Base/IBaseEntityDateModified.cs @@ -2,6 +2,11 @@ { public interface IBaseEntityDateModified { + /// + /// Дата обновления объекта в БД. + /// Обновляется автоматически при комите. + /// Чтобы отключить автоматическое управление, применить атрибут [ManualControl] + /// public DateTimeOffset? DateModified { get; set; } } } diff --git a/PARR.Domain/Entities/TaskEntities/TaskItem.cs b/PARR.Domain/Entities/TaskEntities/TaskItem.cs index 732ea4af..8571e309 100644 --- a/PARR.Domain/Entities/TaskEntities/TaskItem.cs +++ b/PARR.Domain/Entities/TaskEntities/TaskItem.cs @@ -1,5 +1,6 @@ using Microsoft.EntityFrameworkCore; using PARR.Domain.Constants; +using PARR.Domain.Entities.Attributes; using PARR.Domain.Entities.Base; using PARR.Domain.Entities.Base.History; using PARR.Domain.Enums; @@ -17,6 +18,7 @@ namespace PARR.Domain.Entities.TaskEntities public DateTimeOffset DateCreated { get; set; } + [ManualControl] public DateTimeOffset? DateModified { get; set; } public TaskTypeEnum TypeCode { get; set; } diff --git a/PARR.Domain/Settings/IReconciliationSettings.cs b/PARR.Domain/Settings/IReconciliationSettings.cs new file mode 100644 index 00000000..e37e38e3 --- /dev/null +++ b/PARR.Domain/Settings/IReconciliationSettings.cs @@ -0,0 +1,21 @@ +namespace PARR.Domain.Settings +{ + /// + /// Настройки поиска зависших заданий (Task, Reconciliation Job) + /// + public interface IReconciliationSettings + { + /// + /// Интервал запуска проверки. + /// Рекомендуемое значение: 5-10 минут. + /// + TimeSpan CheckInterval { get; set; } + + /// + /// Максимальное кол-во задач для обработки за один проход. + /// Защита от перегрузки при большом кол-ве зависших задач. + /// Рекомендуемое значение 100 + /// + public int MaxTaskPerRun { get; set; } + } +} From e8086b3a7dbf165b16bd0c42e8177e23fac9f91a Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Thu, 23 Apr 2026 16:08:16 +1000 Subject: [PATCH 4/6] =?UTF-8?q?feat(api,=20Core,=20TaskReconciliationWorke?= =?UTF-8?q?r):=20=D0=92=D0=BE=D1=80=D0=BA=D0=B5=D1=80=20=D0=BE=D0=B1=D1=80?= =?UTF-8?q?=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B8=20=D0=B7=D0=B0=D0=B2=D0=B8?= =?UTF-8?q?=D1=81=D1=88=D0=B8=D1=85=20=D0=B7=D0=B0=D0=B4=D0=B0=D1=87.=20?= =?UTF-8?q?=D0=94=D0=BE=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=BB?= =?UTF-8?q?=D0=BE=D0=B3=D0=B8=D0=BA=D0=B8.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PARR.API.sln | 6 +++ PARR.API/Program.cs | 2 - PARR.Core/DependencyInjection.cs | 25 +++++++++ .../Implementations/TaskManagementService.cs | 4 ++ .../TaskReconciliationService.cs | 9 ++-- .../ReconciliationHostedService.cs | 54 +++++++++++++++++++ .../WorkloadCacheBuilderService.cs | 1 - .../PARR.TaskReconciliationWorker.csproj | 25 +++++++++ PARR.TaskReconciliationWorker/Program.cs | 49 +++++++++++++++++ .../Properties/launchSettings.json | 11 ++++ .../Settings/WorkerSettings.cs | 18 +++++++ .../appsettings.Development.json | 38 +++++++++++++ .../appsettings.json | 39 ++++++++++++++ 13 files changed, 275 insertions(+), 6 deletions(-) create mode 100644 PARR.Core/Services/Task/ReconciliationHosted/ReconciliationHostedService.cs create mode 100644 PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj create mode 100644 PARR.TaskReconciliationWorker/Program.cs create mode 100644 PARR.TaskReconciliationWorker/Properties/launchSettings.json create mode 100644 PARR.TaskReconciliationWorker/Settings/WorkerSettings.cs create mode 100644 PARR.TaskReconciliationWorker/appsettings.Development.json create mode 100644 PARR.TaskReconciliationWorker/appsettings.json diff --git a/PARR.API.sln b/PARR.API.sln index 9c8bcd5a..e8856eae 100644 --- a/PARR.API.sln +++ b/PARR.API.sln @@ -98,6 +98,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.Infrastructure", "PARR EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.WorkloadBuilderWorker", "PARR.WorkloadBuilderWorker\PARR.WorkloadBuilderWorker.csproj", "{835BD6CF-024E-47F5-BB02-C904AC9733D9}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.TaskReconciliationWorker", "PARR.TaskReconciliationWorker\PARR.TaskReconciliationWorker.csproj", "{DB701295-1696-4E1A-9088-D3AECF823BA6}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -274,6 +276,10 @@ Global {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Debug|Any CPU.Build.0 = Debug|Any CPU {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.ActiveCfg = Release|Any CPU {835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.Build.0 = Release|Any CPU + {DB701295-1696-4E1A-9088-D3AECF823BA6}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {DB701295-1696-4E1A-9088-D3AECF823BA6}.Debug|Any CPU.Build.0 = Debug|Any CPU + {DB701295-1696-4E1A-9088-D3AECF823BA6}.Release|Any CPU.ActiveCfg = Release|Any CPU + {DB701295-1696-4E1A-9088-D3AECF823BA6}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE diff --git a/PARR.API/Program.cs b/PARR.API/Program.cs index a4a8cca8..28a7333d 100644 --- a/PARR.API/Program.cs +++ b/PARR.API/Program.cs @@ -6,8 +6,6 @@ using PARR.API.Installers; using PARR.API.Settings; using PARR.Core; using PARR.DAL; -using PARR.Domain.Enums; -using PARR.Domain.Settings; using PARR.Infrastructure; using Serilog; using System.Reflection; diff --git a/PARR.Core/DependencyInjection.cs b/PARR.Core/DependencyInjection.cs index 80a5b702..7bbd7f2b 100644 --- a/PARR.Core/DependencyInjection.cs +++ b/PARR.Core/DependencyInjection.cs @@ -8,6 +8,7 @@ using PARR.Core.Services.Task.Handlers.Factory; using PARR.Core.Services.Task.Implementations; using PARR.Core.Services.Task.Interfaces; using PARR.Core.Services.Task.Providers; +using PARR.Core.Services.Task.ReconciliationHosted; using PARR.Core.Services.Workload.Implementations; using PARR.Core.Services.Workload.Interfaces; using PARR.Domain.Enums; @@ -108,6 +109,30 @@ namespace PARR.Core } + /// + /// Подключение сервисов для TaskReconciliation + /// + /// + /// + /// + public static IServiceCollection AddTaskReconciliation(this IServiceCollection services, IReconciliationSettings reconciliationSettings) + { + if (reconciliationSettings == null) + throw new ArgumentNullException(nameof(reconciliationSettings)); + + // Настройки + services.AddSingleton(reconciliationSettings); + + // Сервис + services.AddScoped(); + + // Воркер + services.AddHostedService(); + + return services; + } + + private static IServiceCollection AddBllMapping(this IServiceCollection services) { // AutoMapper diff --git a/PARR.Core/Services/Task/Implementations/TaskManagementService.cs b/PARR.Core/Services/Task/Implementations/TaskManagementService.cs index 636f2d52..c6b33a06 100644 --- a/PARR.Core/Services/Task/Implementations/TaskManagementService.cs +++ b/PARR.Core/Services/Task/Implementations/TaskManagementService.cs @@ -89,6 +89,10 @@ namespace PARR.Core.Services.Task.Implementations // Не пробрасываем исключение дальше — задача создана, просто не в очереди. // задача останется в БД со статусом pending, и Reconciliation Task позже ее возмет в работу logger.LogError("Не удалось опубликовать задачу {TaskId} в очередь. Задача осталась в БД со статусом Pending.", task.Id); + + // так как задачу не смогли отправить в очередь повторно, обновим ей DateModified, относительно нее считается время обработки, и задача встанет на повтор позже сама + task.DateModified = DateTimeOffset.UtcNow; + await taskRepository.CommitAsync(); } return task.Id; diff --git a/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs index 68991449..2a00d5ae 100644 --- a/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs +++ b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs @@ -98,12 +98,15 @@ namespace PARR.Core.Services.Task.Implementations return await taskRepository.Get() .Where(t => - t.TypeCode == taskType.Code && t.DateModified < timeoutThreshold && - ( + //t.TypeCode == taskType.Code && t.DateModified < timeoutThreshold && + t.TypeCode == taskType.Code + // если это вновь созданная задача и она упала, то у нее может не быть DateModified, тогда сравним с DateCreated + && (t.DateModified ?? t.DateCreated) < timeoutThreshold + && ( // Зависла в Processing t.StatusCode == TaskItemStatusEnum.Processing //Зависли в Pending - || (t.StatusCode == TaskItemStatusEnum.Pending && t.RetryCount > 0) + || (t.StatusCode == TaskItemStatusEnum.Pending /*&& t.RetryCount > 0*/) ) ) .OrderBy(t => t.DateModified) // Сначала самые старые diff --git a/PARR.Core/Services/Task/ReconciliationHosted/ReconciliationHostedService.cs b/PARR.Core/Services/Task/ReconciliationHosted/ReconciliationHostedService.cs new file mode 100644 index 00000000..f378f1ec --- /dev/null +++ b/PARR.Core/Services/Task/ReconciliationHosted/ReconciliationHostedService.cs @@ -0,0 +1,54 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using PARR.Core.Common.Interfaces; +using PARR.Core.Services.Task.Interfaces; +using PARR.Domain.Settings; + +namespace PARR.Core.Services.Task.ReconciliationHosted +{ + /// + /// Фоновая задача для периодического запуска Reconciliation Job. + /// Обработка зависших задач. + /// + internal class ReconciliationHostedService : BackgroundService + { + // Это реализация воркера, подключается напрямую в воркер + + private readonly ILogger logger; + private readonly IReconciliationSettings reconciliationSettings; + private readonly IServiceProvider serviceProvider; + private readonly IIntervalService intervalService; + + public ReconciliationHostedService( + ILogger logger, + IReconciliationSettings reconciliationSettings, + IServiceProvider serviceProvider, + IIntervalService intervalService + ) + { + this.logger = logger; + this.reconciliationSettings = reconciliationSettings; + this.serviceProvider = serviceProvider; + this.intervalService = intervalService; + } + + protected override async System.Threading.Tasks.Task ExecuteAsync(CancellationToken stoppingToken) + { + await intervalService.IntervalInitAsync(async () => + { + await using (var scope = serviceProvider.CreateAsyncScope()) + { + var reconciliationService = scope.ServiceProvider.GetRequiredService(); + + // Запускаем проверку + var report = await reconciliationService.RunAsync(); + + if (report.HasChanged) + logger.LogInformation("Reconciliation Job выполнен. Проверено: {Checked}, Исправлено: {Fixed}", report.CheckedCount, report.FixedCount); + } + + }, reconciliationSettings.CheckInterval); + } + } +} diff --git a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs index 57715328..14c25621 100644 --- a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs @@ -2,7 +2,6 @@ using Microsoft.Extensions.Logging; using PARR.Core.Common.Interfaces; using PARR.Core.Common.Interfaces.RabbitServices; -using PARR.Core.Services.Task.Handlers; using PARR.Core.Services.Task.Handlers.Factory; using PARR.Core.Services.Workload.Interfaces; using PARR.Domain.Common.Rabbit.Messages; diff --git a/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj b/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj new file mode 100644 index 00000000..b4d2b4c6 --- /dev/null +++ b/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj @@ -0,0 +1,25 @@ + + + + net7.0 + enable + enable + dotnet-PARR.TaskReconciliationWorker-90386488-d496-4f4d-8dcd-7dcc50882876 + + + + + + + + + + + + + + + + + + diff --git a/PARR.TaskReconciliationWorker/Program.cs b/PARR.TaskReconciliationWorker/Program.cs new file mode 100644 index 00000000..734837da --- /dev/null +++ b/PARR.TaskReconciliationWorker/Program.cs @@ -0,0 +1,49 @@ +using Elastic.CommonSchema.Serilog; +using PARR.Core; +using PARR.DAL; +using PARR.Infrastructure; +using PARR.TaskReconciliationWorker.Settings; +using Serilog; + +var builder = Host.CreateApplicationBuilder(); + +builder.Services.AddLogging(config => +{ + config.ClearProviders(); + + var logger = new LoggerConfiguration(); + + if (builder.Environment.IsProduction()) + logger.WriteTo.Console(new EcsTextFormatter()); + else + logger.WriteTo.Console(); + + logger.ReadFrom.Configuration(builder.Configuration); + + config.AddSerilog(logger.CreateLogger()); +}); + +// --- Настройки --- +var workerSettings = new WorkerSettings(); +builder.Configuration.GetSection(nameof(WorkerSettings)).Bind(workerSettings); + +builder.Services.InstallDalServices(builder.Configuration); +builder.Services.AddCoreServices(builder.Configuration); +builder.Services.AddInfrastructureServices(builder.Configuration); +builder.Configuration.AddDalConfigurations(builder.Services); +builder.Services.AddDallSettings(builder.Configuration); + +builder.Services.AddTaskManagement((provider, taskType) => +{ + // настройки очередей + //var workerSettings = provider.GetRequiredService(); + if (workerSettings.MqSettings.Tasks.TryGetValue(taskType, out var settings)) + return settings; + + throw new ArgumentException($"Настройки для типа задания {taskType} не найдены в конфигурации."); +}); + +builder.Services.AddTaskReconciliation(workerSettings); + +var host = builder.Build(); +host.Run(); \ No newline at end of file diff --git a/PARR.TaskReconciliationWorker/Properties/launchSettings.json b/PARR.TaskReconciliationWorker/Properties/launchSettings.json new file mode 100644 index 00000000..db49dcf0 --- /dev/null +++ b/PARR.TaskReconciliationWorker/Properties/launchSettings.json @@ -0,0 +1,11 @@ +{ + "profiles": { + "PARR.TaskReconciliationWorker": { + "commandName": "Project", + "dotnetRunMessages": true, + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/PARR.TaskReconciliationWorker/Settings/WorkerSettings.cs b/PARR.TaskReconciliationWorker/Settings/WorkerSettings.cs new file mode 100644 index 00000000..16976a84 --- /dev/null +++ b/PARR.TaskReconciliationWorker/Settings/WorkerSettings.cs @@ -0,0 +1,18 @@ +using PARR.Domain.Enums; +using PARR.Domain.Settings; + +namespace PARR.TaskReconciliationWorker.Settings +{ + internal class WorkerSettings : IReconciliationSettings + { + public TimeSpan CheckInterval { get; set; } = TimeSpan.FromMinutes(10); + public int MaxTaskPerRun { get; set; } = 100; + + public MqSettings MqSettings { get; set; } = new(); + } + + internal class MqSettings + { + public Dictionary Tasks { get; set; } = new(); + } +} diff --git a/PARR.TaskReconciliationWorker/appsettings.Development.json b/PARR.TaskReconciliationWorker/appsettings.Development.json new file mode 100644 index 00000000..7e71cd41 --- /dev/null +++ b/PARR.TaskReconciliationWorker/appsettings.Development.json @@ -0,0 +1,38 @@ +{ + "ConnectionStrings": { + "RedisConnection": "10.99.253.216:6379,password=ParrP@ssPtk202MMdevDvs" + }, + "Logging": { + "LogLevel": { + "Default": "Debug", + "Microsoft.Hosting.Lifetime": "Information" + } + }, + "Serilog": { + "MinimumLevel": { + "Default": "Debug", + "Override": { + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + }, + "WriteTo": [ + { + "Name": "File", + "Args": { + "path": "log/log-.txt", + "rollingInterval": "Day" + } + } + ] + }, + "WorkerSettings": { + "MqSettings": { + "Tasks": { + "Workload": { + "HostName": "10.99.253.216" + } + } + } + } +} diff --git a/PARR.TaskReconciliationWorker/appsettings.json b/PARR.TaskReconciliationWorker/appsettings.json new file mode 100644 index 00000000..05d4b3fe --- /dev/null +++ b/PARR.TaskReconciliationWorker/appsettings.json @@ -0,0 +1,39 @@ +{ + "ConnectionStrings": { + "DefaultConnection": "Server=10.99.253.184;Database=parr;User Id=app_parr; Password=PosdfkhT&)%sdfligL&%5546;", + "RedisConnection": "parr-redis:6379,password=ParrP@ssPtk202MMdevDvs" + }, + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft.EntityFrameworkCore": "Error", + "Microsoft.EntityFrameworkCore.Database.Command": "Warning", + "Microsoft.AspNetCore": "Warning" + } + }, + "Serilog": { + "MinimumLevel": { + "Default": "Information", + "Override": { + "Microsoft": "Warning", + "Microsoft.EntityFrameworkCore": "Error", + "Microsoft.EntityFrameworkCore.Database.Command": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } + }, + "WorkerSettings": { + "CheckInterval": "00:10:00", + "MaxTaskPerRun": 100, + "MqSettings": { + "Tasks": { + "Workload": { + "HostName": "parr-rabbitmq", + "QueueName": "parr-task-workload", + "User": "task_workload_writer", + "Password": "UFtsduifytIUR$^85382Wt1few" + } + } + } + } +} From ee119caf916d9a927d76f877fdc9d8ee82e5eec9 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Fri, 24 Apr 2026 09:49:13 +1000 Subject: [PATCH 5/6] feat(taskReconciliation, workloadBuilder): Dockerfile, ci-cd --- .gitlab-ci.yml | 100 ++++++++++++++++++ .../Services/Task/Handlers/BaseTaskHandler.cs | 10 +- .../WorkloadCacheBuilderService.cs | 18 +--- PARR.TaskReconciliationWorker/Dockerfile | 33 ++++++ .../PARR.TaskReconciliationWorker.csproj | 2 + .../Properties/launchSettings.json | 11 +- PARR.WorkloadBuilderWorker/Dockerfile | 33 ++++++ .../PARR.WorkloadBuilderWorker.csproj | 2 + PARR.WorkloadBuilderWorker/Program.cs | 18 +++- .../Properties/launchSettings.json | 11 +- docker-compose.task-reconciliation.yml | 27 +++++ docker-compose.workload-builder.yml | 27 +++++ stack.yml | 75 ------------- 13 files changed, 261 insertions(+), 106 deletions(-) create mode 100644 PARR.TaskReconciliationWorker/Dockerfile create mode 100644 PARR.WorkloadBuilderWorker/Dockerfile create mode 100644 docker-compose.task-reconciliation.yml create mode 100644 docker-compose.workload-builder.yml delete mode 100644 stack.yml diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index bf8f2f97..f7179161 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -18,6 +18,9 @@ variables: PROD_JOB_AUTO_CONTROL: "parr/parr-job-auto-control" PROD_TEMPLATE_MATCHER: "parr/parr-template-matcher" PROD_TEMPLATE_UPDATER: "parr/parr-template-updater" + PROD_WORKLOAD_BUILDER: "parr/parr-workload-builder" + PROD_TASK_RECONCILIATION: "parr/parr-task-reconciliation" + stages: - build @@ -899,5 +902,102 @@ prod_template_updater_deploy: parallel: matrix: - RUNNER: shell-api-swarm-01 + tags: + - ${RUNNER} + + + +### WORKLOAD BUILDER PROD ### +prod_workload_builder_build: + stage: build + only: + - /^wb[0-9]+\.[0-9]+\.[0-9]+$/ + except: + - branches + services: + - name: docker:20.10.21-dind + command: [ + "--insecure-registry=10.99.253.167:8090", + "--registry-mirror=http://10.99.253.167:8090", + "--insecure-registry=harbor.dvgd.rzd", + "--tls=false" + ] + variables: + DOCKER_HOST: tcp://docker:2375 + DOCKER_DRIVER: overlay2 + DOCKER_TLS_CERTDIR: "" + script: + - APP_VERSION=$(echo $CI_COMMIT_TAG | tr -d wb) + - IMAGE_VERSION=$(echo $CI_COMMIT_TAG | sed 's/wb/v/g') + - docker build -t $REPO/$PROD_WORKLOAD_BUILDER:$IMAGE_VERSION -t $REPO/$PROD_WORKLOAD_BUILDER:latest -t $PROD_WORKLOAD_BUILDER:$IMAGE_VERSION -t $PROD_WORKLOAD_BUILDER:latest --build-arg app_version=$APP_VERSION -f PARR.WorkloadBuilderWorker/Dockerfile . + - docker login -u $HARBOR_PUSH_USER -p $HARBOR_PUSH_PASS $REPO + - docker push --all-tags $REPO/$PROD_WORKLOAD_BUILDER + tags: + - docker + + +prod_workload_builder_deploy: + stage: deploy + environment: + name: parr-workload-builder + only: + - /^wb[0-9]+\.[0-9]+\.[0-9]+$/ + except: + - branches + script: + - IMAGE_VERSION=$(echo $CI_COMMIT_TAG | sed 's/wb/v/g') + - docker login -u $HARBOR_PULL_USER -p $HARBOR_PULL_PASS $REPO + - tag=$IMAGE_VERSION docker stack deploy -c docker-compose.workload-builder.yml parr-workload-builder --with-registry-auth + parallel: + matrix: + - RUNNER: shell-api-swarm-01 + tags: + - ${RUNNER} + + +### TASK RECONCILIATION PROD ### +prod_task_reconciliation_build: + stage: build + only: + - /^tr[0-9]+\.[0-9]+\.[0-9]+$/ + except: + - branches + services: + - name: docker:20.10.21-dind + command: [ + "--insecure-registry=10.99.253.167:8090", + "--registry-mirror=http://10.99.253.167:8090", + "--insecure-registry=harbor.dvgd.rzd", + "--tls=false" + ] + variables: + DOCKER_HOST: tcp://docker:2375 + DOCKER_DRIVER: overlay2 + DOCKER_TLS_CERTDIR: "" + script: + - APP_VERSION=$(echo $CI_COMMIT_TAG | tr -d tr) + - IMAGE_VERSION=$(echo $CI_COMMIT_TAG | sed 's/tr/v/g') + - docker build -t $REPO/$PROD_TASK_RECONCILIATION:$IMAGE_VERSION -t $REPO/$PROD_TASK_RECONCILIATION:latest -t $PROD_TASK_RECONCILIATION:$IMAGE_VERSION -t $PROD_TASK_RECONCILIATION:latest --build-arg app_version=$APP_VERSION -f PARR.TaskReconciliationWorker/Dockerfile . + - docker login -u $HARBOR_PUSH_USER -p $HARBOR_PUSH_PASS $REPO + - docker push --all-tags $REPO/$PROD_TASK_RECONCILIATION + tags: + - docker + + +prod_task_reconciliation_deploy: + stage: deploy + environment: + name: parr-task-reconciliation + only: + - /^tr[0-9]+\.[0-9]+\.[0-9]+$/ + except: + - branches + script: + - IMAGE_VERSION=$(echo $CI_COMMIT_TAG | sed 's/tr/v/g') + - docker login -u $HARBOR_PULL_USER -p $HARBOR_PULL_PASS $REPO + - tag=$IMAGE_VERSION docker stack deploy -c docker-compose.task-reconciliation.yml parr-task-reconciliation --with-registry-auth + parallel: + matrix: + - RUNNER: shell-api-swarm-01 tags: - ${RUNNER} \ No newline at end of file diff --git a/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs index 7274863c..e79e9ac3 100644 --- a/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs +++ b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs @@ -226,14 +226,16 @@ namespace PARR.Core.Services.Task.Handlers // В простой реализации — выбрасываем событие, которое ловит воркер // Вариант: сохранить флаг, а воркер после ProcessAsync проверит и отправит - // Это будет реализовано в BackgroundService //todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! - // тоже самое в WorkloadCacheBuilderService - //todo:!!!!!!!!! - await System.Threading.Tasks.Task.Delay(50); + + // Вообще это можно не использовать, так как если рэббит упал а мы еще раз тут же отправляем с задержкой, зачем туда отправлять если рэббит лежит? + + await System.Threading.Tasks.Task.CompletedTask; + + // await System.Threading.Tasks.Task.Delay(50); } diff --git a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs index 14c25621..57cfbdef 100644 --- a/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs @@ -57,20 +57,10 @@ namespace PARR.Core.Services.Workload.Implementations var result = await handler.ProcessAsync(queueMessage.TaskId); - //// Если ошибка и есть retry — отправляем в очередь с задержкой - //if (!result.IsSuccess && handler is BaseTaskHandler baseHandler) - //{ - // // TODO: - // // Если ошибка и есть retry — отправляем в очередь с задержкой - // // (это упрощенная реализация, в идеале — через ITaskQueueService) - - // // Здесь нужна логика retry с задержкой - // // Можно реализовать через отдельный сервис - - // //todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! -> - // // тоже самое в BaseTaskHandler - - //} + if (result.IsSuccess) + logger.LogInformation("Задача {TaskId} успешно обработана", queueMessage.TaskId); + else + logger.LogWarning("Задача {TaskId} завершена с ошибкой. Статус: {Status}", queueMessage.TaskId, result.ErrorMessage); } }); diff --git a/PARR.TaskReconciliationWorker/Dockerfile b/PARR.TaskReconciliationWorker/Dockerfile new file mode 100644 index 00000000..76032c28 --- /dev/null +++ b/PARR.TaskReconciliationWorker/Dockerfile @@ -0,0 +1,33 @@ +# See https://aka.ms/customizecontainer to learn how to customize your debug container and how Visual Studio uses this Dockerfile to build your images for faster debugging. + +# This stage is used when running from VS in fast mode (Default for Debug configuration) +FROM 10.99.253.167:8090/dotnet/runtime:7.0 AS base +WORKDIR /app + + +# This stage is used to build the service project +FROM 10.99.253.167:8090/dotnet/sdk:7.0 AS build +ARG BUILD_CONFIGURATION=Release +WORKDIR /src +COPY ["NuGet.config", "."] +COPY ["PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj", "PARR.TaskReconciliationWorker/"] +COPY ["PARR.Core/PARR.Core.csproj", "PARR.Core/"] +COPY ["PARR.Domain/PARR.Domain.csproj", "PARR.Domain/"] +COPY ["PARR.DAL/PARR.DAL.csproj", "PARR.DAL/"] +COPY ["PARR.Infrastructure/PARR.Infrastructure.csproj", "PARR.Infrastructure/"] +RUN dotnet restore "./PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj" +COPY . . +WORKDIR "/src/PARR.TaskReconciliationWorker" +RUN dotnet build "./PARR.TaskReconciliationWorker.csproj" -c $BUILD_CONFIGURATION -o /app/build + +# This stage is used to publish the service project to be copied to the final stage +FROM build AS publish +ARG app_version=0.0.0-default +ARG BUILD_CONFIGURATION=Release +RUN dotnet publish "./PARR.TaskReconciliationWorker.csproj" -c $BUILD_CONFIGURATION -o /app/publish /p:UseAppHost=false /p:Version=$app_version + +# This stage is used in production or when running from VS in regular mode (Default when not using the Debug configuration) +FROM base AS final +WORKDIR /app +COPY --from=publish /app/publish . +ENTRYPOINT ["dotnet", "PARR.TaskReconciliationWorker.dll"] \ No newline at end of file diff --git a/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj b/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj index b4d2b4c6..01235ea1 100644 --- a/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj +++ b/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj @@ -5,11 +5,13 @@ enable enable dotnet-PARR.TaskReconciliationWorker-90386488-d496-4f4d-8dcd-7dcc50882876 + Linux + diff --git a/PARR.TaskReconciliationWorker/Properties/launchSettings.json b/PARR.TaskReconciliationWorker/Properties/launchSettings.json index db49dcf0..12954245 100644 --- a/PARR.TaskReconciliationWorker/Properties/launchSettings.json +++ b/PARR.TaskReconciliationWorker/Properties/launchSettings.json @@ -1,11 +1,14 @@ -{ +{ "profiles": { "PARR.TaskReconciliationWorker": { "commandName": "Project", - "dotnetRunMessages": true, "environmentVariables": { "DOTNET_ENVIRONMENT": "Development" - } + }, + "dotnetRunMessages": true + }, + "Container (Dockerfile)": { + "commandName": "Docker" } } -} +} \ No newline at end of file diff --git a/PARR.WorkloadBuilderWorker/Dockerfile b/PARR.WorkloadBuilderWorker/Dockerfile new file mode 100644 index 00000000..de02d494 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Dockerfile @@ -0,0 +1,33 @@ +# See https://aka.ms/customizecontainer to learn how to customize your debug container and how Visual Studio uses this Dockerfile to build your images for faster debugging. + +# This stage is used when running from VS in fast mode (Default for Debug configuration) +FROM 10.99.253.167:8090/dotnet/runtime:7.0 AS base +WORKDIR /app + + +# This stage is used to build the service project +FROM 10.99.253.167:8090/dotnet/sdk:7.0 AS build +ARG BUILD_CONFIGURATION=Release +WORKDIR /src +COPY ["NuGet.config", "."] +COPY ["PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj", "PARR.WorkloadBuilderWorker/"] +COPY ["PARR.Core/PARR.Core.csproj", "PARR.Core/"] +COPY ["PARR.Domain/PARR.Domain.csproj", "PARR.Domain/"] +COPY ["PARR.DAL/PARR.DAL.csproj", "PARR.DAL/"] +COPY ["PARR.Infrastructure/PARR.Infrastructure.csproj", "PARR.Infrastructure/"] +RUN dotnet restore "./PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj" +COPY . . +WORKDIR "/src/PARR.WorkloadBuilderWorker" +RUN dotnet build "./PARR.WorkloadBuilderWorker.csproj" -c $BUILD_CONFIGURATION -o /app/build + +# This stage is used to publish the service project to be copied to the final stage +FROM build AS publish +ARG app_version=0.0.0-default +ARG BUILD_CONFIGURATION=Release +RUN dotnet publish "./PARR.WorkloadBuilderWorker.csproj" -c $BUILD_CONFIGURATION -o /app/publish /p:UseAppHost=false /p:Version=$app_version + +# This stage is used in production or when running from VS in regular mode (Default when not using the Debug configuration) +FROM base AS final +WORKDIR /app +COPY --from=publish /app/publish . +ENTRYPOINT ["dotnet", "PARR.WorkloadBuilderWorker.dll"] \ No newline at end of file diff --git a/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj b/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj index fb4489ea..6201a87b 100644 --- a/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj +++ b/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj @@ -5,11 +5,13 @@ enable enable dotnet-PARR.WorkloadBuilderWorker-5cefb704-873b-40e7-b8a0-0537ca29e5ed + Linux + diff --git a/PARR.WorkloadBuilderWorker/Program.cs b/PARR.WorkloadBuilderWorker/Program.cs index 94d5ff05..46209388 100644 --- a/PARR.WorkloadBuilderWorker/Program.cs +++ b/PARR.WorkloadBuilderWorker/Program.cs @@ -1,7 +1,7 @@ using Elastic.CommonSchema.Serilog; using PARR.Core; using PARR.DAL; -using PARR.DAL.Services.Interfaces.Job; +using PARR.Domain.Enums; using PARR.Infrastructure; using PARR.WorkloadBuilderWorker; using PARR.WorkloadBuilderWorker.Settings; @@ -25,19 +25,27 @@ builder.Services.AddLogging(config => config.AddSerilog(logger.CreateLogger()); }); +var workerSettings = new WorkerSettings(); +builder.Configuration.GetSection(nameof(WorkerSettings)).Bind(workerSettings); +builder.Services.AddSingleton(workerSettings); + #region PARR modules builder.Services.InstallDalServices(builder.Configuration); builder.Services.AddCoreServices(builder.Configuration); builder.Services.AddInfrastructureServices(builder.Configuration); +builder.Services.AddTaskManagement((provider, taskType) => +{ + // Тут обрабатываем только задачи типа Workload + if (taskType == TaskTypeEnum.Workload) + return workerSettings.MqSettings; + + throw new ArgumentException($"Настройки для типа задания {taskType} не найдены в конфигурации."); +}); builder.Configuration.AddDalConfigurations(builder.Services); builder.Services.AddDallSettings(builder.Configuration); #endregion -var workerSettings = new WorkerSettings(); -builder.Configuration.GetSection(nameof(WorkerSettings)).Bind(workerSettings); -builder.Services.AddSingleton(workerSettings); - builder.Services.AddHostedService(); diff --git a/PARR.WorkloadBuilderWorker/Properties/launchSettings.json b/PARR.WorkloadBuilderWorker/Properties/launchSettings.json index 1187fd6b..a3943d83 100644 --- a/PARR.WorkloadBuilderWorker/Properties/launchSettings.json +++ b/PARR.WorkloadBuilderWorker/Properties/launchSettings.json @@ -1,11 +1,14 @@ -{ +{ "profiles": { "PARR.WorkloadBuilderWorker": { "commandName": "Project", - "dotnetRunMessages": true, "environmentVariables": { "DOTNET_ENVIRONMENT": "Development" - } + }, + "dotnetRunMessages": true + }, + "Container (Dockerfile)": { + "commandName": "Docker" } } -} +} \ No newline at end of file diff --git a/docker-compose.task-reconciliation.yml b/docker-compose.task-reconciliation.yml new file mode 100644 index 00000000..f9134e87 --- /dev/null +++ b/docker-compose.task-reconciliation.yml @@ -0,0 +1,27 @@ +version: '3.4' + +# task reconciliation +services: + parr-task-reconciliation: + image: harbor.dvgd.rzd/parr/parr-task-reconciliation:${tag:-latest} + environment: + - ASPNETCORE_ENVIRONMENT=Production + - TZ=Europe/Moscow + 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.task-reconciliation.serilog + deploy: + replicas: 1 + networks: + - parr-network + +networks: + parr-network: + driver: overlay + external: true diff --git a/docker-compose.workload-builder.yml b/docker-compose.workload-builder.yml new file mode 100644 index 00000000..6994cb4b --- /dev/null +++ b/docker-compose.workload-builder.yml @@ -0,0 +1,27 @@ +version: '3.4' + +# workload builder +services: + parr-workload-builder: + image: harbor.dvgd.rzd/parr/parr-workload-builder:${tag:-latest} + environment: + - ASPNETCORE_ENVIRONMENT=Production + - TZ=Europe/Moscow + 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.workload-builder.serilog + deploy: + replicas: 1 + networks: + - parr-network + +networks: + parr-network: + driver: overlay + external: true diff --git a/stack.yml b/stack.yml deleted file mode 100644 index 257fdac0..00000000 --- a/stack.yml +++ /dev/null @@ -1,75 +0,0 @@ -version: "3.4" - -# !!!not relevant!!! - -# services: -# parr-espp-template-sync: -# image: harbor.dvgd.rzd/parr/parr-espp-template-sync:${tag:-latest} -# container_name: parr-espp-template-sync -# restart: always -# environment: -# - ASPNETCORE_ENVIRONMENT=Production -# - TZ=Europe/Moscow -# #volumes: -# #- /var/log/parr-espp-template-sync:/app/log -# #- parr_storage_dev:/app/Data -# deploy: -# replicas: 2 - -# rabbitmq: -# image: rabbitmq:3.12.4-management -# container_name: rabbitmq -# hostname: parr-rabbitmq -# environment: -# - RABBITMQ_DEFAULT_USER=rmuser -# - RABBITMQ_DEFAULT_PASS=rmpassword -# - RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS=-rabbit log_levels [{connection,error},{default,error}] disk_free_limit 2147483648 -# volumes: -# - parr_rabbit:/var/lib/rabbitmq -# ports: -# - 15672:15672 -# - 5672:5672 -# deploy: -# replicas: 1 - - -# volumes: -# parr_storage_dev: -# name: parr_storage_dev -# driver: local -# driver_opts: -# type: cifs -# o: "addr=dvgd-appfs-01,username=parruserdev,password=P@ssDevParr202MM" -# device: //dvgd-appfs-01/parr_storage_dev -# parr_rabbit: -# name: parr_rabbit -# driver_opts: -# type: "nfs" -# o: "addr=dvgd-appfs-01.dvgd.oao.rzd,soft,rw,nfsvers=4" -# device: ":/nfs/dvgd-swarm-01/parr_rabbit/data" - - - # # --- NGINX --- - # nginx-2: - # image: nginx:latest - # ports: - # - '8089:80' - # deploy: - # replicas: 5 - # update_config: - # parallelism: 2 - # order: start-first - # failure_action: rollback - # delay: 5s - # rollback_config: - # parallelism: 0 - # order: stop-first - # restart_policy: - # condition: any - # delay: 5s - # max_attempts: 3 - # window: 120s - # healthcheck: - # test: ["CMD", "service", "nginx", "status"] - # networks: - # - nginx-2 \ No newline at end of file From 7ddc8b1bed1c8452efac72e1f0ea1b2b85f24780 Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Wed, 29 Apr 2026 14:27:02 +1000 Subject: [PATCH 6/6] =?UTF-8?q?feat(core):=20=D0=97=D0=B0=D0=B3=D0=BE?= =?UTF-8?q?=D1=82=D0=BE=D0=B2=D0=BA=D0=B0=20WorkloadCacheService?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Implementations/WorkloadCacheService.cs | 53 +++++++++++++++++++ .../IWorkloadCacheBuilderService.cs | 3 ++ 2 files changed, 56 insertions(+) create mode 100644 PARR.Core/Services/Workload/Implementations/WorkloadCacheService.cs diff --git a/PARR.Core/Services/Workload/Implementations/WorkloadCacheService.cs b/PARR.Core/Services/Workload/Implementations/WorkloadCacheService.cs new file mode 100644 index 00000000..61d7dc95 --- /dev/null +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheService.cs @@ -0,0 +1,53 @@ +using Microsoft.Extensions.Logging; +using PARR.Core.Common.Interfaces; + +namespace PARR.Core.Services.Workload.Implementations +{ + /// + /// Класс управления КЭШем для workload + /// + internal class WorkloadCacheService + { + private readonly ILogger logger; + private readonly IRedisCacheService redisCacheService; + + public WorkloadCacheService( + ILogger logger, + IRedisCacheService redisCacheService + //IShortcodeService + ) + { + this.logger = logger; + this.redisCacheService = redisCacheService; + } + + + /// + /// Создать КЭШ шаблонов. + /// КЭШ рабочих групп и зон ответственности + /// + /// + public async System.Threading.Tasks.Task CreateTemplateCacheAsync() + { + await DeleteTemplateCacheAsync(); + + //1. Получить из БД все шаблоны - НЕТ РЕПЫ + //2. Получить по каждому шаблону ЗО и РГ + //3. Сохранить в Redis + + //redisCacheService + } + + /// + /// Удалить КЭШ шаблонов. + /// + /// + private async System.Threading.Tasks.Task DeleteTemplateCacheAsync() + { + + } + + + + } +} diff --git a/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs b/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs index 0d658045..5b0e2e02 100644 --- a/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs +++ b/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs @@ -2,6 +2,9 @@ namespace PARR.Core.Services.Workload.Interfaces { + /// + /// Вызывается в воркере, формирование КЭШ для workload + /// public interface IWorkloadCacheBuilderService : IMessageConsumer {