feat(core): WorkloadCacheService - логика
This commit is contained in:
271
PARR.Core/Services/TaskServices/Handlers/BaseTaskHandler.cs
Normal file
271
PARR.Core/Services/TaskServices/Handlers/BaseTaskHandler.cs
Normal file
@@ -0,0 +1,271 @@
|
||||
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;
|
||||
|
||||
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;
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
using PARR.Core.Services.TaskServices.Handlers;
|
||||
using PARR.Domain.Enums;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Handlers.Factory
|
||||
{
|
||||
/// <summary>
|
||||
/// Фабрика обработчиков
|
||||
/// </summary>
|
||||
public interface ITaskHandlerFactory
|
||||
{
|
||||
/// <summary>
|
||||
/// Получить обработчик для типа задачи
|
||||
/// </summary>
|
||||
/// <param name="taskType"></param>
|
||||
/// <returns></returns>
|
||||
BaseTaskHandler GetHandler(TaskTypeEnum taskType);
|
||||
|
||||
/// <summary>
|
||||
/// Все зарегистрированные типы.
|
||||
/// </summary>
|
||||
/// <returns></returns>
|
||||
IEnumerable<TaskTypeEnum> GetSupportedTypes();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using PARR.Domain.Enums;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Handlers.Factory
|
||||
{
|
||||
public class TaskHandlerFactory : ITaskHandlerFactory
|
||||
{
|
||||
private readonly IServiceProvider serviceProvider;
|
||||
private readonly Dictionary<TaskTypeEnum, Type> handlers;
|
||||
|
||||
public TaskHandlerFactory(IServiceProvider serviceProvider, IEnumerable<BaseTaskHandler> handlers)
|
||||
{
|
||||
this.serviceProvider = serviceProvider;
|
||||
this.handlers = handlers.ToDictionary(t => t.SupportedType, t => t.GetType());
|
||||
}
|
||||
|
||||
public BaseTaskHandler GetHandler(TaskTypeEnum taskType)
|
||||
{
|
||||
if (!handlers.TryGetValue(taskType, out var handlerType))
|
||||
{
|
||||
throw new InvalidOperationException($"Не зарегистрирован обработчик для типа {taskType}");
|
||||
}
|
||||
|
||||
return (BaseTaskHandler)serviceProvider.GetRequiredService(handlerType);
|
||||
}
|
||||
|
||||
public IEnumerable<TaskTypeEnum> GetSupportedTypes()
|
||||
{
|
||||
return handlers.Keys;
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
49
PARR.Core/Services/TaskServices/Handlers/HandlerResult.cs
Normal file
49
PARR.Core/Services/TaskServices/Handlers/HandlerResult.cs
Normal file
@@ -0,0 +1,49 @@
|
||||
namespace PARR.Core.Services.TaskServices.Handlers
|
||||
{
|
||||
/// <summary>
|
||||
/// Результат выполнения задачи обработчиком.
|
||||
/// Используется для преедачи данных между BaseTaskHandler и конкретными хендлерами.
|
||||
/// </summary>
|
||||
public class HandlerResult
|
||||
{
|
||||
// TODO: может его перенести в Domain
|
||||
|
||||
/// <summary>
|
||||
/// Успешно ли выполнено
|
||||
/// </summary>
|
||||
public bool IsSuccess { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Данные результата (JSON, если нужно сохранить в БД)
|
||||
/// </summary>
|
||||
public string? ResultData { get; set; }
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Сообщение об ошибке (если IsSuccess = false)
|
||||
/// </summary>
|
||||
public string? ErrorMessage { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Исключение
|
||||
/// </summary>
|
||||
public Exception? Exception { get; set; }
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Создать успешный результат
|
||||
/// </summary>
|
||||
/// <param name="resultData"></param>
|
||||
/// <returns></returns>
|
||||
public static HandlerResult Success(string? resultData = null) => new() { IsSuccess = true, ResultData = resultData };
|
||||
|
||||
/// <summary>
|
||||
/// Создать результат с ошибкой
|
||||
/// </summary>
|
||||
/// <param name="errorMessage"></param>
|
||||
/// <param name="ex"></param>
|
||||
/// <returns></returns>
|
||||
public static HandlerResult Failure(string errorMessage, Exception? ex = null) => new() { IsSuccess = false, ErrorMessage = errorMessage, Exception = ex };
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
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>
|
||||
/// Обработчик задачаи формирования отчета по загруженности (формирование КЭШ).
|
||||
/// </summary>
|
||||
public class WorkloadReportHandler : BaseTaskHandler
|
||||
{
|
||||
|
||||
public WorkloadReportHandler(ITaskRepository taskRepository, ITaskErrorRepository taskErrorRepository, ILogger<WorkloadReportHandler> logger)
|
||||
: base(taskRepository, taskErrorRepository, logger) { }
|
||||
|
||||
public override TaskTypeEnum SupportedType => TaskTypeEnum.Workload;
|
||||
|
||||
protected override async Task<HandlerResult> ExecuteInternalAsync(TaskItem task)
|
||||
{
|
||||
logger.LogInformation("Генерация КЭШ для workload, taskId: {TaskId}", task.Id);
|
||||
|
||||
// 1. Десериализуем Payload
|
||||
//var payload = string.IsNullOrEmpty(task.Payload)
|
||||
// ? new ReportPayload()
|
||||
// : JsonSerializer.Deserialize<ReportPayload>(task.Payload);
|
||||
|
||||
//todo: тут логика построения отчета
|
||||
|
||||
|
||||
|
||||
await System.Threading.Tasks.Task.CompletedTask;
|
||||
|
||||
// return HandlerResult.Success();
|
||||
|
||||
return HandlerResult.Failure("Тут капец какая ошибка! Просто жуть", new Exception("А это я создал эксепшен!!!"));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,138 @@
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Common.Interfaces.RabbitServices;
|
||||
using PARR.Core.Repositories.Interfaces.TaskRepositories;
|
||||
using PARR.Core.Services.TaskServices.Interfaces;
|
||||
using PARR.Domain.Common.Rabbit.Messages;
|
||||
using PARR.Domain.Entities.Base.History;
|
||||
using PARR.Domain.Entities.TaskEntities;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
using System.Text.Encodings.Web;
|
||||
using System.Text.Json;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Implementations
|
||||
{
|
||||
internal class TaskManagementService : ITaskManagementService
|
||||
{
|
||||
private readonly ILogger<TaskManagementService> logger;
|
||||
private readonly ITaskTypeRepository taskTypeRepository;
|
||||
private readonly ITaskRepository taskRepository;
|
||||
private readonly IRabbitService rabbitService;
|
||||
|
||||
public TaskManagementService(
|
||||
ILogger<TaskManagementService> logger,
|
||||
ITaskTypeRepository taskTypeRepository,
|
||||
ITaskRepository taskRepository,
|
||||
IRabbitService rabbitService
|
||||
)
|
||||
{
|
||||
this.logger = logger;
|
||||
this.taskTypeRepository = taskTypeRepository;
|
||||
this.taskRepository = taskRepository;
|
||||
this.rabbitService = rabbitService;
|
||||
}
|
||||
|
||||
|
||||
public async Task<Guid> CreateTaskAsync<T>(TaskTypeEnum typeCode, T payload, IHistoryInitiator initiator, IMqSettings mqSettings)
|
||||
{
|
||||
var taskType = await taskTypeRepository.Get().AsNoTracking().FirstOrDefaultAsync(t => t.Code == typeCode);
|
||||
|
||||
if (taskType == null)
|
||||
{
|
||||
logger.LogError("Тип задачи {TypeCode} не найден в БД", typeCode);
|
||||
throw new InvalidOperationException($"Тип задачи typeCode не найден в БД");
|
||||
}
|
||||
|
||||
// Проверка IsSingleton: если задача уже активна — возвращаем её
|
||||
if (taskType.IsSingleton)
|
||||
{
|
||||
var existingTask = await GetActiveSingletonTaskAsync(typeCode);
|
||||
|
||||
if (existingTask != null)
|
||||
{
|
||||
logger.LogInformation("Задача типа {TypeCode} уже активна (id: {ExistingId}). Возвращаем существующую.", typeCode, existingTask.Id);
|
||||
return existingTask.Id;
|
||||
}
|
||||
}
|
||||
|
||||
var payloadStr = PayloadToString(payload);
|
||||
|
||||
// Создаем новую запись задачи
|
||||
var task = new TaskItem
|
||||
{
|
||||
Id = Guid.NewGuid(),
|
||||
TypeCode = typeCode,
|
||||
StatusCode = TaskItemStatusEnum.Pending,
|
||||
Payload = payloadStr,
|
||||
RetryCount = 0,
|
||||
ProcessedAt = null
|
||||
};
|
||||
|
||||
if (!await taskRepository.CreateAsync(task) || !await taskRepository.CommitAsync(initiator))
|
||||
{
|
||||
logger.LogError("Ошибка при сохранении задачи в БД");
|
||||
throw new DbUpdateException("Ошибка при сохранении задачи в БД");
|
||||
}
|
||||
|
||||
logger.LogInformation("Задача {TaskId} типа {TypeCode} сохранена в БД со статусом Pending", task.Id, typeCode);
|
||||
|
||||
var message = new TaskMessage { TaskId = task.Id, TypeCode = task.TypeCode };
|
||||
|
||||
var sendResult = await rabbitService.SendAsync(mqSettings, new List<object> { message });
|
||||
if (sendResult.IsSuccess)
|
||||
{
|
||||
logger.LogDebug("Задача {TaskId} типа {TypeCode} отправлена в очередь {QueueName}", task.Id, typeCode, mqSettings.QueueName);
|
||||
}
|
||||
else
|
||||
{
|
||||
// Не пробрасываем исключение дальше — задача создана, просто не в очереди.
|
||||
// задача останется в БД со статусом pending, и Reconciliation Task позже ее возмет в работу
|
||||
logger.LogError("Не удалось опубликовать задачу {TaskId} в очередь. Задача осталась в БД со статусом Pending.", task.Id);
|
||||
|
||||
// так как задачу не смогли отправить в очередь повторно, обновим ей DateModified, относительно нее считается время обработки, и задача встанет на повтор позже сама
|
||||
task.DateModified = DateTimeOffset.UtcNow;
|
||||
await taskRepository.CommitAsync();
|
||||
}
|
||||
|
||||
return task.Id;
|
||||
}
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Получить активную singleton задачу указанного типа
|
||||
/// </summary>
|
||||
/// <param name="typeCode"></param>
|
||||
/// <returns></returns>
|
||||
private async Task<TaskItem?> GetActiveSingletonTaskAsync(TaskTypeEnum typeCode)
|
||||
{
|
||||
return await taskRepository.Get().AsNoTracking()
|
||||
.FirstOrDefaultAsync(t => t.TypeCode == typeCode && (
|
||||
t.StatusCode == TaskItemStatusEnum.Pending
|
||||
|| t.StatusCode == TaskItemStatusEnum.Processing
|
||||
));
|
||||
}
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Payload конвертировать в string
|
||||
/// </summary>
|
||||
/// <typeparam name="T"></typeparam>
|
||||
/// <param name="payload"></param>
|
||||
/// <returns></returns>
|
||||
private string? PayloadToString<T>(T payload)
|
||||
{
|
||||
var jsonOptions = new JsonSerializerOptions
|
||||
{
|
||||
Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping
|
||||
};
|
||||
|
||||
if (payload == null)
|
||||
return null;
|
||||
|
||||
return JsonSerializer.Serialize(payload, jsonOptions);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,190 @@
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Common.Interfaces.RabbitServices;
|
||||
using PARR.Core.Repositories.Interfaces.TaskRepositories;
|
||||
using PARR.Core.Services.TaskServices.Interfaces;
|
||||
using PARR.Core.Services.TaskServices.Models;
|
||||
using PARR.Core.Services.TaskServices.Providers;
|
||||
using PARR.Domain.Common.Rabbit.Messages;
|
||||
using PARR.Domain.Entities.TaskEntities;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Implementations
|
||||
{
|
||||
internal class TaskReconciliationService : ITaskReconciliationService
|
||||
{
|
||||
private readonly ILogger<TaskReconciliationService> 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<TaskReconciliationService> 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<ReconciliationReport> 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;
|
||||
}
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Поиск зависших задач для конкретного типа.
|
||||
/// </summary>
|
||||
/// <param name="taskType"></param>
|
||||
/// <returns></returns>
|
||||
private async Task<List<TaskItem>> 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 &&
|
||||
t.TypeCode == taskType.Code
|
||||
// если это вновь созданная задача и она упала, то у нее может не быть DateModified, тогда сравним с DateCreated
|
||||
&& (t.DateModified ?? t.DateCreated) < timeoutThreshold
|
||||
&& (
|
||||
// Зависла в Processing
|
||||
t.StatusCode == TaskItemStatusEnum.Processing
|
||||
//Зависли в Pending
|
||||
|| (t.StatusCode == TaskItemStatusEnum.Pending /*&& t.RetryCount > 0*/)
|
||||
)
|
||||
)
|
||||
.OrderBy(t => t.DateModified) // Сначала самые старые
|
||||
.Take(reconciliationSettings.MaxTaskPerRun) // Защита от перегрузки
|
||||
.ToListAsync();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Обработка одной зависшей задачи.
|
||||
/// </summary>
|
||||
/// <param name="task"></param>
|
||||
/// <param name="taskType"></param>
|
||||
/// <param name="report"></param>
|
||||
/// <returns></returns>
|
||||
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<object> { 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++;
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
using PARR.Domain.Entities.Base.History;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Interfaces
|
||||
{
|
||||
/// <summary>
|
||||
/// Управления задачами (очередями)
|
||||
/// </summary>
|
||||
public interface ITaskManagementService
|
||||
{
|
||||
/// <summary>
|
||||
/// Создать новую задачу в БД и опубликовать в очередь
|
||||
/// </summary>
|
||||
/// <typeparam name="T"></typeparam>
|
||||
/// <param name="typeCode">Тип задачи</param>
|
||||
/// <param name="payload">Payload</param>
|
||||
/// <param name="initiator">Инициатор</param>
|
||||
/// <param name="mqSettings">Настройки очереди</param>
|
||||
/// <returns></returns>
|
||||
Task<Guid> CreateTaskAsync<T>(TaskTypeEnum typeCode, T payload, IHistoryInitiator initiator, IMqSettings mqSettings);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
using PARR.Core.Services.TaskServices.Models;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Interfaces
|
||||
{
|
||||
/// <summary>
|
||||
/// Сервис для проверки и исправления зависших задач.
|
||||
/// Запускается фоновой задачей по расписанию.
|
||||
/// </summary>
|
||||
internal interface ITaskReconciliationService
|
||||
{
|
||||
/// <summary>
|
||||
/// Выполнить один проход проверки
|
||||
/// </summary>
|
||||
/// <returns></returns>
|
||||
Task<ReconciliationReport> RunAsync();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
namespace PARR.Core.Services.TaskServices.Models
|
||||
{
|
||||
/// <summary>
|
||||
/// Отчет о выполнении Reconciliation Job.
|
||||
/// Используется для мониторинга и логирования.
|
||||
/// </summary>
|
||||
internal class ReconciliationReport
|
||||
{
|
||||
/// <summary>
|
||||
/// Время начала выполнения
|
||||
/// </summary>
|
||||
public DateTimeOffset StartedAt { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Время завершения выполнения
|
||||
/// </summary>
|
||||
public DateTimeOffset CompletedAt { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Сколько задач проверено
|
||||
/// </summary>
|
||||
public int CheckedCount { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Сколько задач исправлено (переведено в retry или failed)
|
||||
/// </summary>
|
||||
public int FixedCount { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Сколько задач отправлено на retry
|
||||
/// </summary>
|
||||
public int RetriedCount { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Сколько задач переведено в финальный failed
|
||||
/// </summary>
|
||||
public int FailedCount { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Сообщения/ошибки (для логирования)
|
||||
/// </summary>
|
||||
public List<string> Messages { get; set; } = new List<string>();
|
||||
|
||||
/// <summary>
|
||||
/// Длительность выполнения
|
||||
/// </summary>
|
||||
public TimeSpan Duration => CompletedAt - StartedAt;
|
||||
|
||||
/// <summary>
|
||||
/// Были ли какие-то действия (повторные отправки, перевод в ошибку)
|
||||
/// </summary>
|
||||
public bool HasChanged => FixedCount > 0;
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Providers
|
||||
{
|
||||
/// <summary>
|
||||
/// Провайдер настроек очереди для задач Tasks
|
||||
/// </summary>
|
||||
public interface ITaskMqSettingsProvider
|
||||
{
|
||||
/// <summary>
|
||||
/// Получить настройки очереди для типа задач
|
||||
/// </summary>
|
||||
/// <param name="typeCode"></param>
|
||||
/// <returns></returns>
|
||||
IMqSettings GetSettings(TaskTypeEnum typeCode);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.Providers
|
||||
{
|
||||
internal class TaskMqSettingsProvider : ITaskMqSettingsProvider
|
||||
{
|
||||
private readonly Func<TaskTypeEnum, IMqSettings> settingsFactory;
|
||||
|
||||
public TaskMqSettingsProvider(
|
||||
Func<TaskTypeEnum, IMqSettings> settingsFactory
|
||||
)
|
||||
{
|
||||
this.settingsFactory = settingsFactory;
|
||||
}
|
||||
|
||||
public IMqSettings GetSettings(TaskTypeEnum typeCode)
|
||||
{
|
||||
return settingsFactory(typeCode);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Common.Interfaces;
|
||||
using PARR.Core.Services.TaskServices.Interfaces;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.TaskServices.ReconciliationHosted
|
||||
{
|
||||
/// <summary>
|
||||
/// Фоновая задача для периодического запуска Reconciliation Job.
|
||||
/// Обработка зависших задач.
|
||||
/// </summary>
|
||||
internal class ReconciliationHostedService : BackgroundService
|
||||
{
|
||||
// Это реализация воркера, подключается напрямую в воркер
|
||||
|
||||
private readonly ILogger<ReconciliationHostedService> logger;
|
||||
private readonly IReconciliationSettings reconciliationSettings;
|
||||
private readonly IServiceProvider serviceProvider;
|
||||
private readonly IIntervalService intervalService;
|
||||
|
||||
public ReconciliationHostedService(
|
||||
ILogger<ReconciliationHostedService> logger,
|
||||
IReconciliationSettings reconciliationSettings,
|
||||
IServiceProvider serviceProvider,
|
||||
IIntervalService intervalService
|
||||
)
|
||||
{
|
||||
this.logger = logger;
|
||||
this.reconciliationSettings = reconciliationSettings;
|
||||
this.serviceProvider = serviceProvider;
|
||||
this.intervalService = intervalService;
|
||||
}
|
||||
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
await intervalService.IntervalInitAsync(async () =>
|
||||
{
|
||||
await using (var scope = serviceProvider.CreateAsyncScope())
|
||||
{
|
||||
var reconciliationService = scope.ServiceProvider.GetRequiredService<ITaskReconciliationService>();
|
||||
|
||||
// Запускаем проверку
|
||||
var report = await reconciliationService.RunAsync();
|
||||
|
||||
if (report.HasChanged)
|
||||
logger.LogInformation("Reconciliation Job выполнен. Проверено: {Checked}, Исправлено: {Fixed}", report.CheckedCount, report.FixedCount);
|
||||
}
|
||||
|
||||
}, reconciliationSettings.CheckInterval);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user