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.TaskServices.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; // Ставим процент 0 task.ProgressPercent = 0; await taskRepository.UpdateProgressAsync(task.Id, 0); return true; } logger.LogDebug("Задача {TaskId} уже захвачена другим воркером", task.Id); return false; } /// /// Обработка успешного выполнения /// /// /// /// private async Task HandleSuccessAsync(TaskItem task, HandlerResult result) { logger.LogInformation("Задача {TaskId} успешно выполнена", task.Id); task.StatusCode = TaskItemStatusEnum.Success; task.ProcessedAt = DateTimeOffset.UtcNow; task.DateModified = DateTimeOffset.UtcNow; // Ставим 100% task.ProgressPercent = 100; //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 Task HandleErrorAsync(TaskItem task, Exception? ex) { // задача в БД, у которой можно получить актуальный процент прогресса var taskInDb = await taskRepository.Get() .AsNoTracking() .FirstOrDefaultAsync(t=>t.Id==task.Id); if (taskInDb != null) { //так как в БД мы обновлям атомарно прогресс, то пришедший нам task не знает фактическое значение прогресса task.ProgressPercent = taskInDb.ProgressPercent; } 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 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 Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask; /// /// Хук при ошибке (можно переопределить в наследнике). /// /// /// /// /// protected virtual Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask; /// /// Обновляет процент выполнения задачи. /// Обновлять процент не чаще 1 раза в 5 сек. /// Вызывается из конкретного хендлера при желании показать прогресс. /// /// Id задачи /// Процент (0-100). Значения >100 будут ограничены до 100. /// protected async Task UpdateProgressAsync(Guid taskId, int percent) { if (percent < 0) { logger.LogWarning("Попытка установить отрицательный прогресс {Percent} для задачи {TaskId}", percent, taskId); percent = 0; } else if (percent > 100) { logger.LogWarning("Прогресс {Percent} для задачи {TaskId} превысил 100%. Записываем 100.", percent, taskId); percent = 100; } // Обновляем в БД var affectedRows = await taskRepository.UpdateProgressAsync(taskId, percent); if (affectedRows == 0) { // Задача не найдена или уже завершена — не считаем ошибкой logger.LogDebug("Не удалось обновить прогресс задачи {TaskId} (возможно, уже завершена)", taskId); } else { logger.LogDebug("Прогресс задачи {TaskId} обновлен: {Percent}%", taskId, percent); } } } }