feat(dal): TaskManagementService - создать задачу в БД и отправить в очередь

This commit is contained in:
Mikhail Trubnikov
2026-03-23 16:56:31 +10:00
parent eb116bc16e
commit 583da7e433
14 changed files with 220 additions and 116 deletions

View File

@@ -4,6 +4,7 @@ using PARR.API.Contracts.V1;
using PARR.API.Contracts.V1.Responses.Base; using PARR.API.Contracts.V1.Responses.Base;
using PARR.API.Controllers.V1.Base; using PARR.API.Controllers.V1.Base;
using PARR.API.Services.Interfaces; using PARR.API.Services.Interfaces;
using PARR.API.Settings;
using PARR.DAL.Cache.Services.Base; using PARR.DAL.Cache.Services.Base;
using PARR.DAL.NextRunServices; using PARR.DAL.NextRunServices;
using PARR.DAL.Services.Interfaces; using PARR.DAL.Services.Interfaces;
@@ -23,7 +24,8 @@ namespace PARR.API.Controllers.V1
IClientService clientService, IClientService clientService,
IRedisCacheService redisCacheService, IRedisCacheService redisCacheService,
INextRunService nextRunService, INextRunService nextRunService,
ITemplateService templateService ITemplateService templateService,
MqSettings mqSettings
) )
{ {
this.clientService = clientService; this.clientService = clientService;

View File

@@ -11,6 +11,9 @@
{ {
public required string Name { get; set; } public required string Name { get; set; }
/// <summary>
/// Промежуток времени, если робот не доступен, не считаем это ошибкой
/// </summary>
public TimeSpan InaccessibilityTime { get; set; } public TimeSpan InaccessibilityTime { get; set; }
} }

View File

@@ -1,43 +1,22 @@
using PARR.BLL.Contracts.Interfaces; using PARR.BLL.Contracts.Interfaces;
using PARR.DAL.Contracts;
namespace PARR.API.Settings namespace PARR.API.Settings
{ {
public class MqSettings public class MqSettings
{ {
public MqGenerateTemplates GenerateTemplates { get; set; } = new MqGenerateTemplates(); public MqSettingsBase GenerateTemplates { get; set; } = new();
public MqTemplateDistributor TemplateDistributor { get; set; } = new MqTemplateDistributor(); public MqSettingsBase TemplateDistributor { get; set; } = new();
public MqTemplateActivator TemplateActivator { get; set; } = new MqTemplateActivator(); public MqSettingsBase TemplateActivator { get; set; } = new();
public MqStatistics Statistics { get; set; } = new MqStatistics(); public MqStatistics Statistics { get; set; } = new MqStatistics();
public MqTemplatesMatcher TemplatesMatcher { get; set; } = new MqTemplatesMatcher(); public MqSettingsBase TemplatesMatcher { get; set; } = new();
public MqTemplatesUpdater TemplatesUpdater { get; set; } = new MqTemplatesUpdater(); public MqSettingsBase TemplatesUpdater { get; set; } = new();
public MqNextRun NextRun { get; set; } = new MqNextRun(); public MqSettingsBase NextRun { get; set; } = new();
}
public class MqGenerateTemplates : IMqSettings /// <summary>
{ /// Настройки MQ для задач
public string HostName { get; set; } = string.Empty; /// </summary>
public string QueueName { get; set; } = string.Empty; public Dictionary<TaskTypeEnum, MqSettingsBase> Tasks { get; set; } = new();
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
public class MqTemplateDistributor : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
public class MqTemplateActivator : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
} }
public class MqStatistics public class MqStatistics
@@ -54,30 +33,4 @@ namespace PARR.API.Settings
public ushort? PrefetchCount { get; set; } = 0; public ushort? PrefetchCount { get; set; } = 0;
} }
public class MqTemplatesMatcher : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
public class MqTemplatesUpdater : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
public class MqNextRun : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
} }

View File

@@ -48,6 +48,11 @@
"NextRun": { "NextRun": {
"HostName": "10.99.253.216" "HostName": "10.99.253.216"
}, },
"Tasks": {
"Workload": {
"HostName": "10.99.253.216"
}
},
"Statistics": { "Statistics": {
"StatMqAuth": { "StatMqAuth": {
"HostName": "10.99.253.216" "HostName": "10.99.253.216"

View File

@@ -66,6 +66,14 @@
"User": "next_run_updater_writer", "User": "next_run_updater_writer",
"Password": "KJHGfdsjfhg$976412!kay" "Password": "KJHGfdsjfhg$976412!kay"
}, },
"Tasks": {
"Workload": {
"HostName": "parr-rabbitmq",
"QueueName": "parr-task-workload",
"User": "task_workload_writer",
"Password": "UFtsduifytIUR$^85382Wt1few"
}
},
"Statistics": { "Statistics": {
"StatMqAuth": { "StatMqAuth": {
"HostName": "parr-rabbitmq", "HostName": "parr-rabbitmq",

View File

@@ -14,4 +14,13 @@
/// </summary> /// </summary>
ushort? PrefetchCount { get; set; } ushort? PrefetchCount { get; set; }
} }
public class MqSettingsBase : IMqSettings
{
public string HostName { get; set; } = string.Empty;
public string QueueName { get; set; } = string.Empty;
public string User { get; set; } = string.Empty;
public string Password { get; set; } = string.Empty;
public ushort? PrefetchCount { get; set; } = 0;
}
} }

View File

@@ -1,7 +1,6 @@
using Microsoft.Extensions.Caching.Distributed; using Microsoft.Extensions.Caching.Distributed;
using StackExchange.Redis; using StackExchange.Redis;
using System.Text.Json; using System.Text.Json;
using static System.Runtime.InteropServices.JavaScript.JSType;
namespace PARR.DAL.Cache.Services.Base namespace PARR.DAL.Cache.Services.Base
{ {

View File

@@ -26,6 +26,7 @@ using PARR.DAL.Services.Interfaces.Schedule;
using PARR.DAL.Services.Interfaces.TaskServices; using PARR.DAL.Services.Interfaces.TaskServices;
using PARR.DAL.Services.Interfaces.Unit; using PARR.DAL.Services.Interfaces.Unit;
using PARR.DAL.Settings; using PARR.DAL.Settings;
using PARR.DAL.TaskServices;
using StackExchange.Redis; using StackExchange.Redis;
namespace PARR.DAL namespace PARR.DAL
@@ -150,6 +151,9 @@ namespace PARR.DAL
services.AddTransient<ITaskErrorService, TaskErrorService>(); services.AddTransient<ITaskErrorService, TaskErrorService>();
services.AddTransient<ITaskService, TaskService>(); services.AddTransient<ITaskService, TaskService>();
services.AddTransient<ITaskTypeService, TaskTypeService>();
services.AddTransient<ITaskManagementService, TaskManagementService>();
#endregion #endregion

View File

@@ -0,0 +1,21 @@
using PARR.DAL.Context;
using PARR.DAL.Models.TaskModels;
using PARR.DAL.Services.Interfaces.TaskServices;
namespace PARR.DAL.Services.Implementations.TaskServices
{
internal class TaskTypeService : ITaskTypeService
{
private readonly DataContext dataContext;
public TaskTypeService(DataContext dataContext)
{
this.dataContext = dataContext;
}
public IQueryable<TaskType> Get()
{
return dataContext.TaskTypes;
}
}
}

View File

@@ -0,0 +1,9 @@
using PARR.DAL.Models.TaskModels;
namespace PARR.DAL.Services.Interfaces.TaskServices
{
public interface ITaskTypeService
{
IQueryable<TaskType> Get();
}
}

View File

@@ -0,0 +1,23 @@
using PARR.BLL.Contracts.Interfaces;
using PARR.Common.Domain;
using PARR.DAL.Contracts;
namespace PARR.DAL.TaskServices
{
/// <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);
}
}

View File

@@ -0,0 +1,124 @@
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using PARR.BLL.Contracts.Interfaces;
using PARR.Common.Domain;
using PARR.DAL.Contracts;
using PARR.DAL.Models.TaskModels;
using PARR.DAL.Services.Interfaces.TaskServices;
using System.Text.Encodings.Web;
using System.Text.Json;
namespace PARR.DAL.TaskServices
{
internal class TaskManagementService : ITaskManagementService
{
private readonly ILogger<TaskManagementService> logger;
private readonly ITaskTypeService taskTypeService;
private readonly ITaskService taskService;
public TaskManagementService(
ILogger<TaskManagementService> logger,
ITaskTypeService taskTypeService,
ITaskService taskService
)
{
this.logger = logger;
this.taskTypeService = taskTypeService;
this.taskService = taskService;
}
public async Task<Guid> CreateTaskAsync<T>(TaskTypeEnum typeCode, T payload, IHistoryInitiator initiator, IMqSettings mqSettings)
{
var taskType = await taskTypeService.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,
InitiatorIp = initiator.InitiatorIp,
InitiatorParrComponentId = initiator.InitiatorParrComponentId,
InitiatorComment = initiator.InitiatorComment
};
if (!await taskService.CreateAsync(task) || !await taskService.CommitAsync())
{
logger.LogError("Ошибка при сохранении задачи в БД");
return default;
}
logger.LogInformation("Задача {TaskId} типа {TypeCode} сохранена в БД со статусом Pending", task.Id, typeCode);
// Публикуме задачу в очередь. Если вдруг даже не получится,
// задача останется в БД со статусом pending, и Reconciliation Task позже ее возмет в работу
// todo: опубликовать в очередь
// вопросы, зачем scoped? может можно AddTransient?
// нужно придумать шаблонную модель для очереди
return task.Id;
}
/// <summary>
/// Получить активную singleton задачу указанного типа
/// </summary>
/// <param name="typeCode"></param>
/// <returns></returns>
private async Task<TaskItem?> GetActiveSingletonTaskAsync(TaskTypeEnum typeCode)
{
return await taskService.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
};
return JsonSerializer.Serialize(payload, jsonOptions);
}
}
}

View File

@@ -1,24 +0,0 @@

namespace PARR.DAL.TransformServices
{
///// <summary>
///// Сервис по изменению NextRun
///// </summary>
//public interface INextRunModifierService
//{
// /// <summary>
// /// Получить NextRun с учетом тайм зоны УЗ робота ЕСПП
// /// </summary>
// /// <param name="nextRun"></param>
// /// <returns></returns>
// DateTimeOffset GetNextRunByAccountRobotTimeZone(DateTimeOffset nextRun);
// ///// <summary>
// ///// Получить РАБОЧИЙ день следующего срабатывания
// ///// </summary>
// ///// <param name="nextRun"></param>
// ///// <returns></returns>
// //Task<DateTimeOffset> GetWorkDayAsync(DateTimeOffset nextRun);
//}
}

View File

@@ -1,32 +0,0 @@
using Microsoft.Extensions.Logging;
using PARR.DAL.Contracts;
using PARR.DAL.Services.Interfaces;
namespace PARR.DAL.TransformServices
{
//internal class NextRunModifierService : INextRunModifierService
//{
// private readonly SettingsFromDb settingsFromDb;
// private readonly ILogger<NextRunModifierService> logger;
// private readonly IWeekendDayService weekendDayService;
// public NextRunModifierService(
// SettingsFromDb settingsFromDb,
// ILogger<NextRunModifierService> logger,
// IWeekendDayService weekendDayService
// )
// {
// this.settingsFromDb = settingsFromDb;
// this.logger = logger;
// this.weekendDayService = weekendDayService;
// }
// public DateTimeOffset GetNextRunByAccountRobotTimeZone(DateTimeOffset nextRun)
// {
// return nextRun.AddHours(settingsFromDb.EsppRobotAccountTimeZoneHour);
// }
//}
}