Files
parr_api/PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs
2026-04-24 09:49:13 +10:00

272 lines
11 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.Task.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;
return true;
}
logger.LogDebug("Задача {TaskId} уже захвачена другим воркером", task.Id);
return false;
}
/// <summary>
/// Обработка успешного выполнения
/// </summary>
/// <param name="task"></param>
/// <param name="result"></param>
/// <returns></returns>
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
}
/// <summary>
/// Обработка ошибки: запись в TaskErrors, проверка MaxRetries.
/// </summary>
/// <param name="task"></param>
/// <param name="ex"></param>
/// <returns></returns>
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);
}
/// <summary>
/// Планирование повторной попытки (отправка в очередь с задержкой)
/// </summary>
/// <param name="task"></param>
/// <returns></returns>
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);
}
/// <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 System.Threading.Tasks.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 System.Threading.Tasks.Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask;
}
}