2026-04-21 16:54:39 +10:00
|
|
|
|
using Microsoft.EntityFrameworkCore;
|
|
|
|
|
|
using Microsoft.Extensions.Logging;
|
|
|
|
|
|
using PARR.Core.Repositories.Interfaces.TaskRepositories;
|
|
|
|
|
|
using PARR.Domain.Entities.TaskEntities;
|
|
|
|
|
|
using PARR.Domain.Enums;
|
|
|
|
|
|
|
2026-04-30 15:10:35 +10:00
|
|
|
|
namespace PARR.Core.Services.TaskServices.Handlers
|
2026-04-21 16:54:39 +10:00
|
|
|
|
{
|
|
|
|
|
|
/// <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)
|
2026-04-22 16:13:23 +10:00
|
|
|
|
{
|
|
|
|
|
|
// Логируем, но не считаем ошибкой
|
|
|
|
|
|
logger.LogDebug("Задача {TaskId} не в Pending (статус {Status}), пропускаем", task.Id, task.StatusCode);
|
2026-04-21 16:54:39 +10:00
|
|
|
|
return false;
|
2026-04-22 16:13:23 +10:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
|
|
|
|
|
|
// Атомарный захват задачи
|
|
|
|
|
|
var affectedRows = await taskRepository.TaskCaptureAsync(task.Id);
|
|
|
|
|
|
|
|
|
|
|
|
if (affectedRows > 0)
|
|
|
|
|
|
{
|
|
|
|
|
|
// Смог захватить задачу
|
|
|
|
|
|
// Обновим локальное состояние
|
|
|
|
|
|
task.StatusCode = TaskItemStatusEnum.Processing;
|
|
|
|
|
|
task.DateModified = DateTimeOffset.UtcNow;
|
|
|
|
|
|
|
2026-05-20 14:26:25 +10:00
|
|
|
|
// Ставим процент 0
|
|
|
|
|
|
task.ProgressPercent = 0;
|
|
|
|
|
|
await taskRepository.UpdateProgressAsync(task.Id, 0);
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
return true;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-04-22 16:13:23 +10:00
|
|
|
|
logger.LogDebug("Задача {TaskId} уже захвачена другим воркером", task.Id);
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
return false;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
|
/// Обработка успешного выполнения
|
|
|
|
|
|
/// </summary>
|
|
|
|
|
|
/// <param name="task"></param>
|
|
|
|
|
|
/// <param name="result"></param>
|
|
|
|
|
|
/// <returns></returns>
|
2026-05-20 14:26:25 +10:00
|
|
|
|
private async Task HandleSuccessAsync(TaskItem task, HandlerResult result)
|
2026-04-21 16:54:39 +10:00
|
|
|
|
{
|
|
|
|
|
|
logger.LogInformation("Задача {TaskId} успешно выполнена", task.Id);
|
|
|
|
|
|
|
|
|
|
|
|
task.StatusCode = TaskItemStatusEnum.Success;
|
|
|
|
|
|
task.ProcessedAt = DateTimeOffset.UtcNow;
|
|
|
|
|
|
task.DateModified = DateTimeOffset.UtcNow;
|
|
|
|
|
|
|
2026-05-20 14:26:25 +10:00
|
|
|
|
// Ставим 100%
|
|
|
|
|
|
task.ProgressPercent = 100;
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
//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>
|
2026-05-20 14:26:25 +10:00
|
|
|
|
private async Task HandleErrorAsync(TaskItem task, Exception? ex)
|
2026-04-21 16:54:39 +10:00
|
|
|
|
{
|
2026-05-20 14:26:25 +10:00
|
|
|
|
// задача в БД, у которой можно получить актуальный процент прогресса
|
|
|
|
|
|
var taskInDb = await taskRepository.Get()
|
|
|
|
|
|
.AsNoTracking()
|
|
|
|
|
|
.FirstOrDefaultAsync(t=>t.Id==task.Id);
|
|
|
|
|
|
|
|
|
|
|
|
if (taskInDb != null)
|
|
|
|
|
|
{
|
|
|
|
|
|
//так как в БД мы обновлям атомарно прогресс, то пришедший нам task не знает фактическое значение прогресса
|
|
|
|
|
|
task.ProgressPercent = taskInDb.ProgressPercent;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
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>
|
2026-05-20 14:26:25 +10:00
|
|
|
|
protected virtual async Task ScheduleRetryAsync(TaskItem task)
|
2026-04-21 16:54:39 +10:00
|
|
|
|
{
|
|
|
|
|
|
// Здесь нужна интеграция с ITaskQueueService
|
|
|
|
|
|
// Для этого можно использовать событие или callback
|
|
|
|
|
|
// В простой реализации — выбрасываем событие, которое ловит воркер
|
|
|
|
|
|
|
|
|
|
|
|
// Вариант: сохранить флаг, а воркер после ProcessAsync проверит и отправит
|
|
|
|
|
|
|
2026-04-22 16:13:23 +10:00
|
|
|
|
//todo: !!! Вот это делаем дальше!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!
|
|
|
|
|
|
|
2026-04-21 16:54:39 +10:00
|
|
|
|
//todo:!!!!!!!!!
|
2026-04-24 09:49:13 +10:00
|
|
|
|
|
|
|
|
|
|
// Вообще это можно не использовать, так как если рэббит упал а мы еще раз тут же отправляем с задержкой, зачем туда отправлять если рэббит лежит?
|
|
|
|
|
|
|
|
|
|
|
|
await System.Threading.Tasks.Task.CompletedTask;
|
|
|
|
|
|
|
|
|
|
|
|
// await System.Threading.Tasks.Task.Delay(50);
|
2026-04-21 16:54:39 +10:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// <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>
|
2026-05-20 14:26:25 +10:00
|
|
|
|
protected virtual Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask;
|
2026-04-21 16:54:39 +10:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
|
/// Хук при ошибке (можно переопределить в наследнике).
|
|
|
|
|
|
/// </summary>
|
|
|
|
|
|
/// <param name="task"></param>
|
|
|
|
|
|
/// <param name="ex"></param>
|
|
|
|
|
|
/// <param name="attemptNumber"></param>
|
|
|
|
|
|
/// <returns></returns>
|
2026-05-20 14:26:25 +10:00
|
|
|
|
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);
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-04-21 16:54:39 +10:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|