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.API.sln b/PARR.API.sln index 0855d326..e8856eae 100644 --- a/PARR.API.sln +++ b/PARR.API.sln @@ -96,6 +96,10 @@ 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 +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 @@ -268,6 +272,14 @@ 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 + {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/Controllers/V1/Statistics/StatWorkloadController.cs b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs index c688fd7f..d9fedf61 100644 --- a/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs +++ b/PARR.API/Controllers/V1/Statistics/StatWorkloadController.cs @@ -1,8 +1,15 @@ 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.Core.Services.Task.Providers; using PARR.Domain.Common.Roles; +using PARR.Domain.Entities.Base.History; +using PARR.Domain.Enums; namespace PARR.API.Controllers.V1.Statistics { @@ -12,9 +19,25 @@ 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; + private readonly ITaskMqSettingsProvider taskMqSettingsProvider; + public StatWorkloadController( + ITaskManagementService taskManagementService, + IClientService clientService, + //MqSettings mqSettings, + ILogger logger, + ITaskMqSettingsProvider taskMqSettingsProvider + ) + { + this.taskManagementService = taskManagementService; + this.clientService = clientService; + //this.mqSettings = mqSettings; + this.logger = logger; + this.taskMqSettingsProvider = taskMqSettingsProvider; } @@ -25,8 +48,34 @@ 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(); + //todo: избавиться от try/catch + + try + { + //mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings); + + //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); + + 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.API/Program.cs b/PARR.API/Program.cs index 3c7e520e..28a7333d 100644 --- a/PARR.API/Program.cs +++ b/PARR.API/Program.cs @@ -3,6 +3,7 @@ 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.Infrastructure; @@ -21,14 +22,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/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 c003bf91..7bbd7f2b 100644 --- a/PARR.Core/DependencyInjection.cs +++ b/PARR.Core/DependencyInjection.cs @@ -3,6 +3,15 @@ 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.Task.Providers; +using PARR.Core.Services.Task.ReconciliationHosted; +using PARR.Core.Services.Workload.Implementations; +using PARR.Core.Services.Workload.Interfaces; +using PARR.Domain.Enums; using PARR.Domain.Settings; namespace PARR.Core @@ -35,12 +44,91 @@ namespace PARR.Core #endregion + #region Task + + // ---------> Регистрируем в AddTaskManagement <--------- + + //// Регистрация хендлеров (нужно регистировать каждый отдельно) + //services.AddScoped(); + //services.AddScoped(sp => sp.GetRequiredService());// регим для фабрики + //// Еще хэндлеры... + + //// Фабрика хэндлеров + //services.AddScoped(); + + //services.AddScoped(); + + //// Провайдер настроек MQ + //services.AddSingleton(); + + #endregion + #region Services + //services.AddScoped(); #endregion + #region IMessageConsumer services + + services.AddSingleton(); + + #endregion + + 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; + } + + + /// + /// Подключение сервисов для 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; } 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.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..e79e9ac3 --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs @@ -0,0 +1,271 @@ +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) + { + // Логируем, но не считаем ошибкой + logger.LogDebug("Задача {TaskId} не в Pending (статус {Status}), пропускаем", task.Id, task.StatusCode); + return false; + } + + + // Атомарный захват задачи + var affectedRows = await taskRepository.TaskCaptureAsync(task.Id); + + if (affectedRows > 0) + { + // Смог захватить задачу + // Обновим локальное состояние + task.StatusCode = TaskItemStatusEnum.Processing; + task.DateModified = DateTimeOffset.UtcNow; + + return true; + } + + logger.LogDebug("Задача {TaskId} уже захвачена другим воркером", task.Id); + + 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 проверит и отправит + + //todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!! + + //todo:!!!!!!!!! + + // Вообще это можно не использовать, так как если рэббит упал а мы еще раз тут же отправляем с задержкой, зачем туда отправлять если рэббит лежит? + + await System.Threading.Tasks.Task.CompletedTask; + + // 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..708ca3e5 --- /dev/null +++ b/PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs @@ -0,0 +1,39 @@ +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.DAL/TaskServices/TaskManagementService.cs b/PARR.Core/Services/Task/Implementations/TaskManagementService.cs similarity index 53% rename from PARR.DAL/TaskServices/TaskManagementService.cs rename to PARR.Core/Services/Task/Implementations/TaskManagementService.cs index dca781c5..c6b33a06 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,34 @@ 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); + + // так как задачу не смогли отправить в очередь повторно, обновим ей DateModified, относительно нее считается время обработки, и задача встанет на повтор позже сама + task.DateModified = DateTimeOffset.UtcNow; + await taskRepository.CommitAsync(); + } return task.Id; } @@ -95,7 +106,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 +120,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.Core/Services/Task/Implementations/TaskReconciliationService.cs b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs new file mode 100644 index 00000000..2a00d5ae --- /dev/null +++ b/PARR.Core/Services/Task/Implementations/TaskReconciliationService.cs @@ -0,0 +1,190 @@ +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 && + t.TypeCode == taskType.Code + // если это вновь созданная задача и она упала, то у нее может не быть DateModified, тогда сравним с DateCreated + && (t.DateModified ?? t.DateCreated) < 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.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.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/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 new file mode 100644 index 00000000..57cfbdef --- /dev/null +++ b/PARR.Core/Services/Workload/Implementations/WorkloadCacheBuilderService.cs @@ -0,0 +1,77 @@ +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.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); + + if (result.IsSuccess) + logger.LogInformation("Задача {TaskId} успешно обработана", queueMessage.TaskId); + else + logger.LogWarning("Задача {TaskId} завершена с ошибкой. Статус: {Status}", queueMessage.TaskId, result.ErrorMessage); + } + }); + + if (!isConnected) + throw new Exception("Ошибка при подключении к RabbitMq"); + } + + public async System.Threading.Tasks.Task StopProcessingAsync() + { + await rabbitService.DisposeAsync(); + } + + } +} 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 new file mode 100644 index 00000000..5b0e2e02 --- /dev/null +++ b/PARR.Core/Services/Workload/Interfaces/IWorkloadCacheBuilderService.cs @@ -0,0 +1,14 @@ +using PARR.Core.Common.Interfaces.RabbitServices.Workers; + +namespace PARR.Core.Services.Workload.Interfaces +{ + /// + /// Вызывается в воркере, формирование КЭШ для workload + /// + public interface IWorkloadCacheBuilderService : IMessageConsumer + { + + } +} + + 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..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; @@ -99,6 +101,9 @@ namespace PARR.DAL.Repositories.Base /// private void SetInitiator(IHistoryInitiator? initiator) { + if (initiator == null) + return; + logger.LogDebug("Устанавливаю инициатора для изменений"); // Задаем инициатора только для новых и измененных записей @@ -130,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); + } } } @@ -319,7 +342,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.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; } + } +} 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; } + } +} 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 new file mode 100644 index 00000000..01235ea1 --- /dev/null +++ b/PARR.TaskReconciliationWorker/PARR.TaskReconciliationWorker.csproj @@ -0,0 +1,27 @@ + + + + net7.0 + enable + enable + dotnet-PARR.TaskReconciliationWorker-90386488-d496-4f4d-8dcd-7dcc50882876 + Linux + + + + + + + + + + + + + + + + + + + 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..12954245 --- /dev/null +++ b/PARR.TaskReconciliationWorker/Properties/launchSettings.json @@ -0,0 +1,14 @@ +{ + "profiles": { + "PARR.TaskReconciliationWorker": { + "commandName": "Project", + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + }, + "dotnetRunMessages": true + }, + "Container (Dockerfile)": { + "commandName": "Docker" + } + } +} \ No newline at end of file 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" + } + } + } + } +} 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 new file mode 100644 index 00000000..6201a87b --- /dev/null +++ b/PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj @@ -0,0 +1,27 @@ + + + + net7.0 + enable + enable + dotnet-PARR.WorkloadBuilderWorker-5cefb704-873b-40e7-b8a0-0537ca29e5ed + Linux + + + + + + + + + + + + + + + + + + + diff --git a/PARR.WorkloadBuilderWorker/Program.cs b/PARR.WorkloadBuilderWorker/Program.cs new file mode 100644 index 00000000..46209388 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Program.cs @@ -0,0 +1,53 @@ +using Elastic.CommonSchema.Serilog; +using PARR.Core; +using PARR.DAL; +using PARR.Domain.Enums; +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()); +}); + +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 + + +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..a3943d83 --- /dev/null +++ b/PARR.WorkloadBuilderWorker/Properties/launchSettings.json @@ -0,0 +1,14 @@ +{ + "profiles": { + "PARR.WorkloadBuilderWorker": { + "commandName": "Project", + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + }, + "dotnetRunMessages": true + }, + "Container (Dockerfile)": { + "commandName": "Docker" + } + } +} \ No newline at end of file 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" + } + } +} 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