From 7da2265d1cbb9023772691461b12a0ad46d6864a Mon Sep 17 00:00:00 2001 From: Mikhail Trubnikov Date: Wed, 22 Apr 2026 16:13:23 +1000 Subject: [PATCH] =?UTF-8?q?feat(api,=20Core,=20Domain,=20Dal):=20=D0=94?= =?UTF-8?q?=D0=BB=D1=8F=20entity=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB=D0=B5?= =?UTF-8?q?=D0=B9,=20=D0=B4=D0=BB=D1=8F=20DateModified=20=D1=81=D0=BE?= =?UTF-8?q?=D0=B7=D0=B4=D0=B0=D0=BD=20=D0=B0=D1=82=D1=80=D0=B8=D0=B1=D1=83?= =?UTF-8?q?=D1=82=20ManualControl=20-=20=D0=BE=D1=82=D0=BA=D0=BB=D1=8E?= =?UTF-8?q?=D1=87=D0=B0=D0=B5=D1=82=20=D0=B0=D0=B2=D1=82=D0=BE=D0=BC=D0=B0?= =?UTF-8?q?=D1=82=D0=B8=D1=87=D0=B5=D1=81=D0=BA=D0=BE=D0=B5=20=D1=83=D0=BF?= =?UTF-8?q?=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=B4=D0=B0?= =?UTF-8?q?=D1=82=D0=BE=D0=B9.=20TaskReconciliationService=20-=20=D1=81?= =?UTF-8?q?=D0=B5=D1=80=D0=B2=D0=B8=D1=81=20=D0=BF=D0=BE=20=D1=83=D0=BF?= =?UTF-8?q?=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D1=8E=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=B4=D0=B0=D1=87=D0=B0=D0=BC=D0=B8=20=D0=B2=20=D1=81=D1=82?= =?UTF-8?q?=D0=B0=D1=82=D1=83=D1=81=D0=B0=D1=85=20=D0=BE=D1=88=D0=B8=D0=B1?= =?UTF-8?q?=D0=BA=D0=B8,=20=D0=B7=D0=B0=D0=B2=D0=B8=D1=81=D0=BB=D0=B0.=20I?= =?UTF-8?q?TaskMqSettingsProvider=20-=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B0?= =?UTF-8?q?=D0=B9=D0=B4=D0=B5=D1=80=20=D0=B4=D0=BB=D1=8F=20=D1=83=D0=BF?= =?UTF-8?q?=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D1=8F=20=D0=BD=D0=B0?= =?UTF-8?q?=D1=81=D1=82=D1=80=D0=BE=D0=B9=D0=BA=D0=B0=D0=BC=D0=B8=20MQ=20?= =?UTF-8?q?=D0=B2=20=D0=B7=D0=B0=D0=B2=D0=B8=D1=81=D0=B8=D0=BC=D0=BE=D1=81?= =?UTF-8?q?=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; } + } +}