2026-03-23 16:56:31 +10:00
using Microsoft.EntityFrameworkCore ;
using Microsoft.Extensions.Logging ;
2026-04-15 16:48:23 +10:00
using PARR.Core.Common.Interfaces.RabbitServices ;
2026-04-14 09:56:31 +10:00
using PARR.Core.Repositories.Interfaces.TaskRepositories ;
2026-04-30 15:10:35 +10:00
using PARR.Core.Services.TaskServices.Interfaces ;
2026-04-15 16:48:23 +10:00
using PARR.Domain.Common.Rabbit.Messages ;
2026-04-13 16:58:30 +10:00
using PARR.Domain.Entities.Base.History ;
using PARR.Domain.Entities.TaskEntities ;
using PARR.Domain.Enums ;
2026-04-14 12:01:49 +10:00
using PARR.Domain.Settings ;
2026-03-23 16:56:31 +10:00
using System.Text.Encodings.Web ;
using System.Text.Json ;
2026-04-30 15:10:35 +10:00
namespace PARR.Core.Services.TaskServices.Implementations
2026-03-23 16:56:31 +10:00
{
internal class TaskManagementService : ITaskManagementService
{
private readonly ILogger < TaskManagementService > logger ;
2026-04-15 16:48:23 +10:00
private readonly ITaskTypeRepository taskTypeRepository ;
private readonly ITaskRepository taskRepository ;
private readonly IRabbitService rabbitService ;
2026-03-23 16:56:31 +10:00
public TaskManagementService (
ILogger < TaskManagementService > logger ,
2026-04-15 16:48:23 +10:00
ITaskTypeRepository taskTypeRepository ,
ITaskRepository taskRepository ,
IRabbitService rabbitService
2026-03-23 16:56:31 +10:00
)
{
this . logger = logger ;
2026-04-15 16:48:23 +10:00
this . taskTypeRepository = taskTypeRepository ;
this . taskRepository = taskRepository ;
this . rabbitService = rabbitService ;
2026-03-23 16:56:31 +10:00
}
public async Task < Guid > CreateTaskAsync < T > ( TaskTypeEnum typeCode , T payload , IHistoryInitiator initiator , IMqSettings mqSettings )
{
2026-04-15 16:48:23 +10:00
var taskType = await taskTypeRepository . Get ( ) . AsNoTracking ( ) . FirstOrDefaultAsync ( t = > t . Code = = typeCode ) ;
2026-03-23 16:56:31 +10:00
if ( taskType = = null )
{
logger . LogError ( "Тип задачи {TypeCode} не найден в БД" , typeCode ) ;
throw new InvalidOperationException ( $"Тип задачи typeCode не найден в БД" ) ;
}
// Проверка IsSingleton: если задача уже активна — возвращаем её
if ( taskType . IsSingleton )
{
var existingTask = await GetActiveSingletonTaskAsync ( typeCode ) ;
if ( existingTask ! = null )
{
2026-04-15 16:48:23 +10:00
logger . LogInformation ( "Задача типа {TypeCode} уже активна (id: {ExistingId}). Возвращаем существующую." , typeCode , existingTask . Id ) ;
2026-03-23 16:56:31 +10:00
return existingTask . Id ;
}
}
var payloadStr = PayloadToString ( payload ) ;
// Создаем новую запись задачи
var task = new TaskItem
{
Id = Guid . NewGuid ( ) ,
TypeCode = typeCode ,
StatusCode = TaskItemStatusEnum . Pending ,
Payload = payloadStr ,
RetryCount = 0 ,
2026-04-15 16:48:23 +10:00
ProcessedAt = null
2026-03-23 16:56:31 +10:00
} ;
2026-04-15 16:48:23 +10:00
if ( ! await taskRepository . CreateAsync ( task ) | | ! await taskRepository . CommitAsync ( initiator ) )
2026-03-23 16:56:31 +10:00
{
logger . LogError ( "Ошибка при сохранении задачи в БД" ) ;
2026-04-15 16:48:23 +10:00
throw new DbUpdateException ( "Ошибка при сохранении задачи в БД" ) ;
2026-03-23 16:56:31 +10:00
}
logger . LogInformation ( "Задача {TaskId} типа {TypeCode} сохранена в БД с о статусом Pending" , task . Id , typeCode ) ;
2026-04-15 16:48:23 +10:00
var message = new TaskMessage { TaskId = task . Id , TypeCode = task . TypeCode } ;
2026-03-23 16:56:31 +10:00
2026-04-15 16:48:23 +10:00
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 ) ;
2026-04-23 16:08:16 +10:00
// так как задачу не смогли отправить в очередь повторно, обновим ей DateModified, относительно нее считается время обработки, и задача встанет на повтор позже сама
task . DateModified = DateTimeOffset . UtcNow ;
await taskRepository . CommitAsync ( ) ;
2026-04-15 16:48:23 +10:00
}
2026-03-23 16:56:31 +10:00
return task . Id ;
}
/// <summary>
/// Получить активную singleton задачу указанного типа
/// </summary>
/// <param name="typeCode"></param>
/// <returns></returns>
private async Task < TaskItem ? > GetActiveSingletonTaskAsync ( TaskTypeEnum typeCode )
{
2026-04-15 16:48:23 +10:00
return await taskRepository . Get ( ) . AsNoTracking ( )
2026-03-23 16:56:31 +10:00
. 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>
2026-04-15 16:48:23 +10:00
private string? PayloadToString < T > ( T payload )
2026-03-23 16:56:31 +10:00
{
var jsonOptions = new JsonSerializerOptions
{
Encoder = JavaScriptEncoder . UnsafeRelaxedJsonEscaping
} ;
2026-04-15 16:48:23 +10:00
if ( payload = = null )
return null ;
2026-03-23 16:56:31 +10:00
return JsonSerializer . Serialize ( payload , jsonOptions ) ;
}
}
}