Files
parr_api/PARR.Core/Services/TaskServices/Handlers/BaseTaskHandler.cs

326 lines
14 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
{
/// <summary>
/// Базовый класс для всех обработчиков задач.
/// Реализует шаблонный метод ProcessAsync с общей логикой:
/// - Атомарный захват задачи
/// - Обработка исключений
/// - Обновление статусов
/// - Запись ошибок
/// - Логика retry
/// </summary>
public abstract class BaseTaskHandler
{
protected readonly ITaskRepository taskRepository;
protected readonly ITaskErrorRepository taskErrorRepository;
protected readonly ILogger<BaseTaskHandler> logger;
protected BaseTaskHandler(ITaskRepository taskRepository, ITaskErrorRepository taskErrorRepository, ILogger<BaseTaskHandler> logger)
{
this.taskRepository = taskRepository;
this.taskErrorRepository = taskErrorRepository;
this.logger = logger;
}
/// <summary>
/// Тип задачи, который обрабатывает этот хендлер.
/// Должен быть переопределен в конкретном классе.
/// </summary>
public abstract TaskTypeEnum SupportedType { get; }
public async Task<HandlerResult> 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<TaskItem?> LoadTaskAsync(Guid taskId)
{
return await taskRepository.Get()
.Include(t => t.TaskType)
.Include(t => t.TaskStatus)
.FirstOrDefaultAsync(t => t.Id == taskId);
}
/// <summary>
/// Атомарный захват задачи: Pending -> Processing.
/// Возвращает true, если удалось захватить.
/// </summary>
/// <param name="task"></param>
/// <returns></returns>
private async Task<bool> 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;
}
/// <summary>
/// Обработка успешного выполнения
/// </summary>
/// <param name="task"></param>
/// <param name="result"></param>
/// <returns></returns>
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
}
/// <summary>
/// Обработка ошибки: запись в TaskErrors, проверка MaxRetries.
/// </summary>
/// <param name="task"></param>
/// <param name="ex"></param>
/// <returns></returns>
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);
}
/// <summary>
/// Планирование повторной попытки (отправка в очередь с задержкой)
/// </summary>
/// <param name="task"></param>
/// <returns></returns>
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);
}
/// <summary>
/// Абстрактный метод - логика конкретного типа задачи.
/// Релизация в наследнике.
/// </summary>
/// <param name="task"></param>
/// <returns></returns>
protected abstract Task<HandlerResult> ExecuteInternalAsync(TaskItem task);
/// <summary>
/// Хук после успешного выполнения (можно переопределить в наследнике).
/// </summary>
/// <param name="task"></param>
/// <param name="result"></param>
/// <returns></returns>
protected virtual Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask;
/// <summary>
/// Хук при ошибке (можно переопределить в наследнике).
/// </summary>
/// <param name="task"></param>
/// <param name="ex"></param>
/// <param name="attemptNumber"></param>
/// <returns></returns>
protected virtual Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask;
/// <summary>
/// Обновляет процент выполнения задачи.
/// Обновлять процент не чаще 1 раза в 5 сек.
/// Вызывается из конкретного хендлера при желании показать прогресс.
/// </summary>
/// <param name="taskId">Id задачи</param>
/// <param name="percent">Процент (0-100). Значения >100 будут ограничены до 100.</param>
/// <returns></returns>
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);
}
}
}
}