167 lines
7.4 KiB
C#
167 lines
7.4 KiB
C#
using AutoMapper;
|
||
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.DTOs.TaskDTO;
|
||
using PARR.Domain.Entities.Base.History;
|
||
using PARR.Domain.Enums;
|
||
using PARR.Domain.Exceptions;
|
||
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;
|
||
private readonly IMapper mapper;
|
||
|
||
public TaskManagementService(
|
||
ILogger<TaskManagementService> logger,
|
||
ITaskTypeRepository taskTypeRepository,
|
||
ITaskRepository taskRepository,
|
||
IRabbitService rabbitService,
|
||
IMapper mapper
|
||
)
|
||
{
|
||
this.logger = logger;
|
||
this.taskTypeRepository = taskTypeRepository;
|
||
this.taskRepository = taskRepository;
|
||
this.rabbitService = rabbitService;
|
||
this.mapper = mapper;
|
||
}
|
||
|
||
|
||
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;
|
||
|
||
logger.LogInformation("Задача типа {TypeCode} уже активна (id: {ExistingId}).", typeCode, existingTask.Id);
|
||
|
||
throw new AlreadyExistsException($"Уже существует активная задача (id: {existingTask.Id}). Нельзя создать новую пока не выполнится существующая.");
|
||
}
|
||
}
|
||
|
||
var payloadStr = PayloadToString(payload);
|
||
|
||
// Создаем новую запись задачи
|
||
var task = new Domain.Entities.TaskEntities.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 DbErrorException("Ошибка при сохранении задачи в БД");
|
||
}
|
||
|
||
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;
|
||
}
|
||
|
||
|
||
public async Task<List<ActiveTask>> GetActiveTasksAsync(TaskTypeEnum? typeCode = null)
|
||
{
|
||
var query = taskRepository.Get()
|
||
.AsNoTracking()
|
||
.Include(t => t.TaskStatus)
|
||
.Include(t => t.TaskType)
|
||
.Where(t => t.StatusCode == TaskItemStatusEnum.Pending || t.StatusCode == TaskItemStatusEnum.Processing);
|
||
|
||
if (typeCode.HasValue)
|
||
query = query.Where(t => t.TypeCode == typeCode.Value);
|
||
|
||
var tasks = await query.OrderByDescending(t => t.DateCreated).ToListAsync();
|
||
|
||
return mapper.Map<List<ActiveTask>>(tasks);
|
||
}
|
||
|
||
|
||
/// <summary>
|
||
/// Получить активную singleton задачу указанного типа
|
||
/// </summary>
|
||
/// <param name="typeCode"></param>
|
||
/// <returns></returns>
|
||
private async Task<Domain.Entities.TaskEntities.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);
|
||
}
|
||
|
||
|
||
}
|
||
}
|