feat(dal, core, workloadBuilderWorker): Сервисы по работе с очередями (задачиами). Заготовка WorkloadBuilderWorker
This commit is contained in:
@@ -96,6 +96,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.Domain", "PARR.Domain\
|
||||
EndProject
|
||||
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.Infrastructure", "PARR.Infrastructure\PARR.Infrastructure.csproj", "{8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}"
|
||||
EndProject
|
||||
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PARR.WorkloadBuilderWorker", "PARR.WorkloadBuilderWorker\PARR.WorkloadBuilderWorker.csproj", "{835BD6CF-024E-47F5-BB02-C904AC9733D9}"
|
||||
EndProject
|
||||
Global
|
||||
GlobalSection(SolutionConfigurationPlatforms) = preSolution
|
||||
Debug|Any CPU = Debug|Any CPU
|
||||
@@ -268,6 +270,10 @@ Global
|
||||
{8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Debug|Any CPU.Build.0 = Debug|Any CPU
|
||||
{8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Release|Any CPU.ActiveCfg = Release|Any CPU
|
||||
{8EB3BC72-2EFA-48AA-A305-C28ABEDB03AD}.Release|Any CPU.Build.0 = Release|Any CPU
|
||||
{835BD6CF-024E-47F5-BB02-C904AC9733D9}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
|
||||
{835BD6CF-024E-47F5-BB02-C904AC9733D9}.Debug|Any CPU.Build.0 = Debug|Any CPU
|
||||
{835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.ActiveCfg = Release|Any CPU
|
||||
{835BD6CF-024E-47F5-BB02-C904AC9733D9}.Release|Any CPU.Build.0 = Release|Any CPU
|
||||
EndGlobalSection
|
||||
GlobalSection(SolutionProperties) = preSolution
|
||||
HideSolutionNode = FALSE
|
||||
|
||||
@@ -51,6 +51,8 @@ namespace PARR.API.Controllers.V1.Statistics
|
||||
InitiatorComment = "Отправлен запрос из API на формирование отчета о загруженности"
|
||||
};
|
||||
|
||||
//todo: избавиться от try/catch
|
||||
|
||||
try
|
||||
{
|
||||
mqSettings.Tasks.TryGetValue(TaskTypeEnum.Workload, out var queueSettings);
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Common.Interfaces.RabbitServices.Workers
|
||||
{
|
||||
/// <summary>
|
||||
/// Интерфейс для сервисов которые используются воркерами по подписке на очереди MQ
|
||||
/// </summary>
|
||||
public interface IMessageConsumer
|
||||
{
|
||||
Task ProcessMessagesAsync(IMqSettings mqSettings);
|
||||
Task StopProcessingAsync();
|
||||
}
|
||||
}
|
||||
@@ -3,8 +3,12 @@ using Microsoft.Extensions.Configuration;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using PARR.Core.Common.Implementations;
|
||||
using PARR.Core.Common.Interfaces;
|
||||
using PARR.Core.Services.Task.Handlers;
|
||||
using PARR.Core.Services.Task.Handlers.Factory;
|
||||
using PARR.Core.Services.Task.Implementations;
|
||||
using PARR.Core.Services.Task.Interfaces;
|
||||
using PARR.Core.Services.Workload.Implementations;
|
||||
using PARR.Core.Services.Workload.Interfaces;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core
|
||||
@@ -37,13 +41,33 @@ namespace PARR.Core
|
||||
|
||||
#endregion
|
||||
|
||||
#region Services
|
||||
#region Task
|
||||
|
||||
// Регистрация хендлеров (нужно регистировать каждый отдельно)
|
||||
services.AddScoped<WorkloadReportHandler>();
|
||||
services.AddScoped<BaseTaskHandler>(sp => sp.GetRequiredService<WorkloadReportHandler>());// регим для фабрики
|
||||
// Еще хэндлеры...
|
||||
|
||||
// Фабрика хэндлеров
|
||||
services.AddScoped<ITaskHandlerFactory, TaskHandlerFactory>();
|
||||
|
||||
services.AddScoped<ITaskManagementService, TaskManagementService>();
|
||||
|
||||
#endregion
|
||||
|
||||
#region Services
|
||||
|
||||
|
||||
//services.AddScoped<IUserService, UserService>();
|
||||
|
||||
#endregion
|
||||
|
||||
#region IMessageConsumer services
|
||||
|
||||
services.AddSingleton<IWorkloadCacheBuilderService, WorkloadCacheBuilderService>();
|
||||
|
||||
#endregion
|
||||
|
||||
return services;
|
||||
}
|
||||
|
||||
|
||||
@@ -5,5 +5,11 @@ namespace PARR.Core.Repositories.Interfaces.TaskRepositories
|
||||
{
|
||||
public interface ITaskRepository : IBaseRepository<TaskItem>
|
||||
{
|
||||
/// <summary>
|
||||
/// Атомарный захват задачи (взять в работу)
|
||||
/// </summary>
|
||||
/// <param name="taskId"></param>
|
||||
/// <returns>Возвращает кол-во обработанных строк</returns>
|
||||
Task<int> TaskCaptureAsync(Guid taskId);
|
||||
}
|
||||
}
|
||||
|
||||
258
PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs
Normal file
258
PARR.Core/Services/Task/Handlers/BaseTaskHandler.cs
Normal file
@@ -0,0 +1,258 @@
|
||||
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)
|
||||
return false;
|
||||
|
||||
// Атомарный захват задачи
|
||||
var affectedRows = await taskRepository.TaskCaptureAsync(task.Id);
|
||||
|
||||
if (affectedRows > 0)
|
||||
{
|
||||
// Смог захватить задачу
|
||||
// Обновим локальное состояние
|
||||
task.StatusCode = TaskItemStatusEnum.Processing;
|
||||
task.DateModified = DateTimeOffset.UtcNow;
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
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 проверит и отправит
|
||||
// Это будет реализовано в BackgroundService
|
||||
|
||||
//todo:!!!!!!!!!
|
||||
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,23 @@
|
||||
using PARR.Domain.Enums;
|
||||
|
||||
namespace PARR.Core.Services.Task.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.Task.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/Task/Handlers/HandlerResult.cs
Normal file
49
PARR.Core/Services/Task/Handlers/HandlerResult.cs
Normal file
@@ -0,0 +1,49 @@
|
||||
namespace PARR.Core.Services.Task.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 };
|
||||
|
||||
}
|
||||
}
|
||||
37
PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs
Normal file
37
PARR.Core/Services/Task/Handlers/WorkloadReportHandler.cs
Normal file
@@ -0,0 +1,37 @@
|
||||
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>
|
||||
/// Обработчик задачаи формирования отчета по загруженности (формирование КЭШ).
|
||||
/// </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,84 @@
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Common.Interfaces;
|
||||
using PARR.Core.Common.Interfaces.RabbitServices;
|
||||
using PARR.Core.Services.Task.Handlers;
|
||||
using PARR.Core.Services.Task.Handlers.Factory;
|
||||
using PARR.Core.Services.Workload.Interfaces;
|
||||
using PARR.Domain.Common.Rabbit.Messages;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.Core.Services.Workload.Implementations
|
||||
{
|
||||
internal class WorkloadCacheBuilderService : IWorkloadCacheBuilderService
|
||||
{
|
||||
private readonly ILogger<WorkloadCacheBuilderService> logger;
|
||||
private readonly IRabbitService rabbitService;
|
||||
private readonly ITransformService transformService;
|
||||
private readonly IServiceProvider serviceProvider;
|
||||
|
||||
public WorkloadCacheBuilderService(
|
||||
ILogger<WorkloadCacheBuilderService> logger,
|
||||
IRabbitService rabbitService,
|
||||
ITransformService transformService,
|
||||
IServiceProvider serviceProvider
|
||||
)
|
||||
{
|
||||
this.logger = logger;
|
||||
this.rabbitService = rabbitService;
|
||||
this.transformService = transformService;
|
||||
this.serviceProvider = serviceProvider;
|
||||
}
|
||||
|
||||
|
||||
public async System.Threading.Tasks.Task ProcessMessagesAsync(IMqSettings mqSettings)
|
||||
{
|
||||
// этот метод вызывается в воркере
|
||||
|
||||
var isConnected = await rabbitService.InitConsumerAsync(mqSettings, async (msg) =>
|
||||
{
|
||||
logger.LogInformation($"Получили запрос: {msg}");
|
||||
|
||||
var queueMessage = transformService.GetModelFromJson<TaskMessage>(msg);
|
||||
if (queueMessage == null)
|
||||
return;
|
||||
|
||||
// Проверим тип сообщения, что он соответствует текущему сервису
|
||||
if (queueMessage.TypeCode != TaskTypeEnum.Workload)
|
||||
{
|
||||
logger.LogError("Несоответствие типа: в очереди {Queue} получена задача типа {Actual}", mqSettings.QueueName, queueMessage.TypeCode);
|
||||
return;
|
||||
}
|
||||
|
||||
using (var scope = serviceProvider.CreateScope())
|
||||
{
|
||||
var handlerFactory = scope.ServiceProvider.GetRequiredService<ITaskHandlerFactory>();
|
||||
var handler = handlerFactory.GetHandler(queueMessage.TypeCode);
|
||||
|
||||
var result = await handler.ProcessAsync(queueMessage.TaskId);
|
||||
|
||||
// Если ошибка и есть retry — отправляем в очередь с задержкой
|
||||
if (!result.IsSuccess && handler is BaseTaskHandler baseHandler)
|
||||
{
|
||||
// TODO:
|
||||
// Если ошибка и есть retry — отправляем в очередь с задержкой
|
||||
// (это упрощенная реализация, в идеале — через ITaskQueueService)
|
||||
|
||||
// Здесь нужна логика retry с задержкой
|
||||
// Можно реализовать через отдельный сервис
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
if (!isConnected)
|
||||
throw new Exception("Ошибка при подключении к RabbitMq");
|
||||
}
|
||||
|
||||
public async System.Threading.Tasks.Task StopProcessingAsync()
|
||||
{
|
||||
await rabbitService.DisposeAsync();
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
using PARR.Core.Common.Interfaces.RabbitServices.Workers;
|
||||
|
||||
namespace PARR.Core.Services.Workload.Interfaces
|
||||
{
|
||||
public interface IWorkloadCacheBuilderService : IMessageConsumer
|
||||
{
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Microsoft.EntityFrameworkCore.Relational" Version="7.0.20" />
|
||||
<PackageReference Include="Microsoft.EntityFrameworkCore.Tools" Version="7.0.20">
|
||||
<PrivateAssets>all</PrivateAssets>
|
||||
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
|
||||
|
||||
@@ -322,7 +322,8 @@ namespace PARR.DAL.Repositories.Base
|
||||
{
|
||||
logger.LogDebug("Начинаю создание объекта типа {EntityType}", typeof(T).Name);
|
||||
|
||||
obj.DateCreated = DateTimeOffset.UtcNow;
|
||||
if (obj.DateCreated == DateTimeOffset.MinValue)
|
||||
obj.DateCreated = DateTimeOffset.UtcNow;
|
||||
|
||||
try
|
||||
{
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Repositories.Interfaces.TaskRepositories;
|
||||
using PARR.DAL.Context;
|
||||
using PARR.DAL.Repositories.Base;
|
||||
using PARR.Domain.Constants;
|
||||
using PARR.Domain.Entities.TaskEntities;
|
||||
using PARR.Domain.Enums;
|
||||
|
||||
namespace PARR.DAL.Repositories.TaskRepositories
|
||||
{
|
||||
@@ -10,7 +13,27 @@ namespace PARR.DAL.Repositories.TaskRepositories
|
||||
{
|
||||
public TaskRepository(DataContext dataContext, ILogger<TaskRepository> logger) : base(logger, dataContext) { }
|
||||
|
||||
public async Task<int> TaskCaptureAsync(Guid taskId)
|
||||
{
|
||||
//var affectedRows = await EntityContext.Database.ExecuteSqlInterpolatedAsync($@"
|
||||
// UPDATE ""{DatabaseSchemas.Task}"".""Tasks""
|
||||
// SET ""{nameof(TaskItem.StatusCode)}""={(int)TaskItemStatusEnum.Processing},
|
||||
// ""{nameof(TaskItem.DateModified)}""={DateTimeOffset.UtcNow}
|
||||
// WHERE ""{nameof(TaskItem.Id)}""={taskId}
|
||||
// AND ""{nameof(TaskItem.StatusCode)}""={(int)TaskItemStatusEnum.Pending}
|
||||
//");
|
||||
|
||||
var sql = $@"
|
||||
UPDATE ""{DatabaseSchemas.Task}"".""Tasks""
|
||||
SET ""{nameof(TaskItem.StatusCode)}"" = {(int)TaskItemStatusEnum.Processing},
|
||||
""{nameof(TaskItem.DateModified)}"" = {{0}}
|
||||
WHERE ""{nameof(TaskItem.Id)}"" = {{1}}
|
||||
AND ""{nameof(TaskItem.StatusCode)}"" = {(int)TaskItemStatusEnum.Pending}";
|
||||
|
||||
var affectedRows = await EntityContext.Database.ExecuteSqlRawAsync(sql, DateTimeOffset.UtcNow, taskId);
|
||||
|
||||
return affectedRows;
|
||||
}
|
||||
|
||||
// метод атомарного взятия в работу
|
||||
}
|
||||
}
|
||||
|
||||
25
PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj
Normal file
25
PARR.WorkloadBuilderWorker/PARR.WorkloadBuilderWorker.csproj
Normal file
@@ -0,0 +1,25 @@
|
||||
<Project Sdk="Microsoft.NET.Sdk.Worker">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net7.0</TargetFramework>
|
||||
<Nullable>enable</Nullable>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<UserSecretsId>dotnet-PARR.WorkloadBuilderWorker-5cefb704-873b-40e7-b8a0-0537ca29e5ed</UserSecretsId>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Elastic.CommonSchema.Serilog" Version="8.6.1" />
|
||||
<PackageReference Include="Microsoft.Extensions.Hosting" Version="7.0.1" />
|
||||
<PackageReference Include="Serilog.Extensions.Hosting" Version="7.0.0" />
|
||||
<PackageReference Include="Serilog.Settings.Configuration" Version="7.0.1" />
|
||||
<PackageReference Include="Serilog.Sinks.Console" Version="4.1.0" />
|
||||
<PackageReference Include="Serilog.Sinks.File" Version="5.0.0" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\PARR.Core\PARR.Core.csproj" />
|
||||
<ProjectReference Include="..\PARR.DAL\PARR.DAL.csproj" />
|
||||
<ProjectReference Include="..\PARR.Domain\PARR.Domain.csproj" />
|
||||
<ProjectReference Include="..\PARR.Infrastructure\PARR.Infrastructure.csproj" />
|
||||
</ItemGroup>
|
||||
</Project>
|
||||
45
PARR.WorkloadBuilderWorker/Program.cs
Normal file
45
PARR.WorkloadBuilderWorker/Program.cs
Normal file
@@ -0,0 +1,45 @@
|
||||
using Elastic.CommonSchema.Serilog;
|
||||
using PARR.Core;
|
||||
using PARR.DAL;
|
||||
using PARR.DAL.Services.Interfaces.Job;
|
||||
using PARR.Infrastructure;
|
||||
using PARR.WorkloadBuilderWorker;
|
||||
using PARR.WorkloadBuilderWorker.Settings;
|
||||
using Serilog;
|
||||
|
||||
var builder = Host.CreateApplicationBuilder();
|
||||
|
||||
builder.Services.AddLogging(config =>
|
||||
{
|
||||
config.ClearProviders();
|
||||
|
||||
var logger = new LoggerConfiguration();
|
||||
|
||||
if (builder.Environment.IsProduction())
|
||||
logger.WriteTo.Console(new EcsTextFormatter());
|
||||
else
|
||||
logger.WriteTo.Console();
|
||||
|
||||
logger.ReadFrom.Configuration(builder.Configuration);
|
||||
|
||||
config.AddSerilog(logger.CreateLogger());
|
||||
});
|
||||
|
||||
#region PARR modules
|
||||
builder.Services.InstallDalServices(builder.Configuration);
|
||||
builder.Services.AddCoreServices(builder.Configuration);
|
||||
builder.Services.AddInfrastructureServices(builder.Configuration);
|
||||
|
||||
builder.Configuration.AddDalConfigurations(builder.Services);
|
||||
builder.Services.AddDallSettings(builder.Configuration);
|
||||
#endregion
|
||||
|
||||
var workerSettings = new WorkerSettings();
|
||||
builder.Configuration.GetSection(nameof(WorkerSettings)).Bind(workerSettings);
|
||||
builder.Services.AddSingleton(workerSettings);
|
||||
|
||||
|
||||
builder.Services.AddHostedService<Worker>();
|
||||
|
||||
var host = builder.Build();
|
||||
host.Run();
|
||||
11
PARR.WorkloadBuilderWorker/Properties/launchSettings.json
Normal file
11
PARR.WorkloadBuilderWorker/Properties/launchSettings.json
Normal file
@@ -0,0 +1,11 @@
|
||||
{
|
||||
"profiles": {
|
||||
"PARR.WorkloadBuilderWorker": {
|
||||
"commandName": "Project",
|
||||
"dotnetRunMessages": true,
|
||||
"environmentVariables": {
|
||||
"DOTNET_ENVIRONMENT": "Development"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
9
PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs
Normal file
9
PARR.WorkloadBuilderWorker/Settings/WorkerSettings.cs
Normal file
@@ -0,0 +1,9 @@
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.WorkloadBuilderWorker.Settings
|
||||
{
|
||||
public record WorkerSettings
|
||||
{
|
||||
public MqSettingsBase MqSettings { get; init; } = new();
|
||||
}
|
||||
}
|
||||
29
PARR.WorkloadBuilderWorker/Worker.cs
Normal file
29
PARR.WorkloadBuilderWorker/Worker.cs
Normal file
@@ -0,0 +1,29 @@
|
||||
using PARR.Core.Services.Workload.Interfaces;
|
||||
using PARR.WorkloadBuilderWorker.Settings;
|
||||
|
||||
namespace PARR.WorkloadBuilderWorker
|
||||
{
|
||||
public class Worker : BackgroundService
|
||||
{
|
||||
private readonly IWorkloadCacheBuilderService workloadCacheBuilderService;
|
||||
private readonly WorkerSettings workerSettings;
|
||||
|
||||
public Worker(ILogger<Worker> logger, IWorkloadCacheBuilderService workloadCacheBuilderService, WorkerSettings WorkerSettings)
|
||||
{
|
||||
this.workloadCacheBuilderService = workloadCacheBuilderService;
|
||||
this.workerSettings = WorkerSettings;
|
||||
}
|
||||
|
||||
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
await workloadCacheBuilderService.ProcessMessagesAsync(workerSettings.MqSettings);
|
||||
}
|
||||
|
||||
public override async Task StopAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
await workloadCacheBuilderService.StopProcessingAsync();
|
||||
await base.StopAsync(cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
35
PARR.WorkloadBuilderWorker/appsettings.Development.json
Normal file
35
PARR.WorkloadBuilderWorker/appsettings.Development.json
Normal file
@@ -0,0 +1,35 @@
|
||||
{
|
||||
"ConnectionStrings": {
|
||||
"RedisConnection": "10.99.253.216:6379,password=ParrP@ssPtk202MMdevDvs"
|
||||
},
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Debug",
|
||||
"Microsoft.Hosting.Lifetime": "Information"
|
||||
}
|
||||
},
|
||||
"Serilog": {
|
||||
"MinimumLevel": {
|
||||
"Default": "Debug",
|
||||
"Override": {
|
||||
"Microsoft": "Warning",
|
||||
"Microsoft.Hosting.Lifetime": "Information"
|
||||
}
|
||||
},
|
||||
"WriteTo": [
|
||||
{
|
||||
"Name": "File",
|
||||
"Args": {
|
||||
"path": "log/log-.txt",
|
||||
"rollingInterval": "Day"
|
||||
}
|
||||
}
|
||||
]
|
||||
},
|
||||
"WorkerSettings": {
|
||||
"MqSettings": {
|
||||
"HostName": "10.99.253.216"
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
33
PARR.WorkloadBuilderWorker/appsettings.json
Normal file
33
PARR.WorkloadBuilderWorker/appsettings.json
Normal file
@@ -0,0 +1,33 @@
|
||||
{
|
||||
"ConnectionStrings": {
|
||||
"DefaultConnection": "Server=10.99.253.184;Database=parr;User Id=app_parr; Password=PosdfkhT&)%sdfligL&%5546;",
|
||||
"RedisConnection": "parr-redis:6379,password=ParrP@ssPtk202MMdevDvs"
|
||||
},
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Information",
|
||||
"Microsoft.EntityFrameworkCore": "Error",
|
||||
"Microsoft.EntityFrameworkCore.Database.Command": "Warning",
|
||||
"Microsoft.AspNetCore": "Warning"
|
||||
}
|
||||
},
|
||||
"Serilog": {
|
||||
"MinimumLevel": {
|
||||
"Default": "Information",
|
||||
"Override": {
|
||||
"Microsoft": "Warning",
|
||||
"Microsoft.EntityFrameworkCore": "Error",
|
||||
"Microsoft.EntityFrameworkCore.Database.Command": "Warning",
|
||||
"Microsoft.Hosting.Lifetime": "Information"
|
||||
}
|
||||
}
|
||||
},
|
||||
"WorkerSettings": {
|
||||
"MqSettings": {
|
||||
"HostName": "parr-rabbitmq",
|
||||
"QueueName": "parr-task-workload",
|
||||
"User": "task_workload_reader",
|
||||
"Password": "sdjfgIYFdifyKJFDhvsad@!314d"
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user