diff --git a/PARR.API/Controllers/V1/TestController.cs b/PARR.API/Controllers/V1/TestController.cs
index ccc11361..04515187 100644
--- a/PARR.API/Controllers/V1/TestController.cs
+++ b/PARR.API/Controllers/V1/TestController.cs
@@ -4,6 +4,7 @@ using PARR.API.Contracts.V1;
using PARR.API.Contracts.V1.Responses.Base;
using PARR.API.Controllers.V1.Base;
using PARR.API.Services.Interfaces;
+using PARR.API.Settings;
using PARR.DAL.Cache.Services.Base;
using PARR.DAL.NextRunServices;
using PARR.DAL.Services.Interfaces;
@@ -23,7 +24,8 @@ namespace PARR.API.Controllers.V1
IClientService clientService,
IRedisCacheService redisCacheService,
INextRunService nextRunService,
- ITemplateService templateService
+ ITemplateService templateService,
+ MqSettings mqSettings
)
{
this.clientService = clientService;
diff --git a/PARR.API/Settings/MonitoringSettings.cs b/PARR.API/Settings/MonitoringSettings.cs
index 51e268ab..1f34046f 100644
--- a/PARR.API/Settings/MonitoringSettings.cs
+++ b/PARR.API/Settings/MonitoringSettings.cs
@@ -11,6 +11,9 @@
{
public required string Name { get; set; }
+ ///
+ /// Промежуток времени, если робот не доступен, не считаем это ошибкой
+ ///
public TimeSpan InaccessibilityTime { get; set; }
}
diff --git a/PARR.API/Settings/MqSettings.cs b/PARR.API/Settings/MqSettings.cs
index c319aa22..1211d6d2 100644
--- a/PARR.API/Settings/MqSettings.cs
+++ b/PARR.API/Settings/MqSettings.cs
@@ -1,43 +1,22 @@
using PARR.BLL.Contracts.Interfaces;
+using PARR.DAL.Contracts;
namespace PARR.API.Settings
{
public class MqSettings
{
- public MqGenerateTemplates GenerateTemplates { get; set; } = new MqGenerateTemplates();
- public MqTemplateDistributor TemplateDistributor { get; set; } = new MqTemplateDistributor();
- public MqTemplateActivator TemplateActivator { get; set; } = new MqTemplateActivator();
+ public MqSettingsBase GenerateTemplates { get; set; } = new();
+ public MqSettingsBase TemplateDistributor { get; set; } = new();
+ public MqSettingsBase TemplateActivator { get; set; } = new();
public MqStatistics Statistics { get; set; } = new MqStatistics();
- public MqTemplatesMatcher TemplatesMatcher { get; set; } = new MqTemplatesMatcher();
- public MqTemplatesUpdater TemplatesUpdater { get; set; } = new MqTemplatesUpdater();
- public MqNextRun NextRun { get; set; } = new MqNextRun();
- }
+ public MqSettingsBase TemplatesMatcher { get; set; } = new();
+ public MqSettingsBase TemplatesUpdater { get; set; } = new();
+ public MqSettingsBase NextRun { get; set; } = new();
- public class MqGenerateTemplates : 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 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;
+ ///
+ /// Настройки MQ для задач
+ ///
+ public Dictionary Tasks { get; set; } = new();
}
public class MqStatistics
@@ -54,30 +33,4 @@ namespace PARR.API.Settings
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;
- }
}
diff --git a/PARR.API/appsettings.Development.json b/PARR.API/appsettings.Development.json
index 7064c31d..ecbe68f1 100644
--- a/PARR.API/appsettings.Development.json
+++ b/PARR.API/appsettings.Development.json
@@ -48,6 +48,11 @@
"NextRun": {
"HostName": "10.99.253.216"
},
+ "Tasks": {
+ "Workload": {
+ "HostName": "10.99.253.216"
+ }
+ },
"Statistics": {
"StatMqAuth": {
"HostName": "10.99.253.216"
diff --git a/PARR.API/appsettings.json b/PARR.API/appsettings.json
index 7f1857f1..0d42f7d2 100644
--- a/PARR.API/appsettings.json
+++ b/PARR.API/appsettings.json
@@ -66,6 +66,14 @@
"User": "next_run_updater_writer",
"Password": "KJHGfdsjfhg$976412!kay"
},
+ "Tasks": {
+ "Workload": {
+ "HostName": "parr-rabbitmq",
+ "QueueName": "parr-task-workload",
+ "User": "task_workload_writer",
+ "Password": "UFtsduifytIUR$^85382Wt1few"
+ }
+ },
"Statistics": {
"StatMqAuth": {
"HostName": "parr-rabbitmq",
diff --git a/PARR.BLL/Contracts/Interfaces/IMqSettings.cs b/PARR.BLL/Contracts/Interfaces/IMqSettings.cs
index aa6f9aa1..a5a72023 100644
--- a/PARR.BLL/Contracts/Interfaces/IMqSettings.cs
+++ b/PARR.BLL/Contracts/Interfaces/IMqSettings.cs
@@ -14,4 +14,13 @@
///
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;
+ }
}
diff --git a/PARR.DAL/Cache/Services/Base/RedisCacheService.cs b/PARR.DAL/Cache/Services/Base/RedisCacheService.cs
index 064ce02d..42f89183 100644
--- a/PARR.DAL/Cache/Services/Base/RedisCacheService.cs
+++ b/PARR.DAL/Cache/Services/Base/RedisCacheService.cs
@@ -1,7 +1,6 @@
using Microsoft.Extensions.Caching.Distributed;
using StackExchange.Redis;
using System.Text.Json;
-using static System.Runtime.InteropServices.JavaScript.JSType;
namespace PARR.DAL.Cache.Services.Base
{
diff --git a/PARR.DAL/ParrDalInstaller.cs b/PARR.DAL/ParrDalInstaller.cs
index 0307b397..d012683c 100644
--- a/PARR.DAL/ParrDalInstaller.cs
+++ b/PARR.DAL/ParrDalInstaller.cs
@@ -26,6 +26,7 @@ using PARR.DAL.Services.Interfaces.Schedule;
using PARR.DAL.Services.Interfaces.TaskServices;
using PARR.DAL.Services.Interfaces.Unit;
using PARR.DAL.Settings;
+using PARR.DAL.TaskServices;
using StackExchange.Redis;
namespace PARR.DAL
@@ -150,6 +151,9 @@ namespace PARR.DAL
services.AddTransient();
services.AddTransient();
+ services.AddTransient();
+
+ services.AddTransient();
#endregion
diff --git a/PARR.DAL/Services/Implementations/TaskServices/TaskTypeService.cs b/PARR.DAL/Services/Implementations/TaskServices/TaskTypeService.cs
new file mode 100644
index 00000000..818556a0
--- /dev/null
+++ b/PARR.DAL/Services/Implementations/TaskServices/TaskTypeService.cs
@@ -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 Get()
+ {
+ return dataContext.TaskTypes;
+ }
+ }
+}
diff --git a/PARR.DAL/Services/Interfaces/TaskServices/ITaskTypeService.cs b/PARR.DAL/Services/Interfaces/TaskServices/ITaskTypeService.cs
new file mode 100644
index 00000000..780910e7
--- /dev/null
+++ b/PARR.DAL/Services/Interfaces/TaskServices/ITaskTypeService.cs
@@ -0,0 +1,9 @@
+using PARR.DAL.Models.TaskModels;
+
+namespace PARR.DAL.Services.Interfaces.TaskServices
+{
+ public interface ITaskTypeService
+ {
+ IQueryable Get();
+ }
+}
diff --git a/PARR.DAL/TaskServices/ITaskManagementService.cs b/PARR.DAL/TaskServices/ITaskManagementService.cs
new file mode 100644
index 00000000..712d5ea6
--- /dev/null
+++ b/PARR.DAL/TaskServices/ITaskManagementService.cs
@@ -0,0 +1,23 @@
+using PARR.BLL.Contracts.Interfaces;
+using PARR.Common.Domain;
+using PARR.DAL.Contracts;
+
+namespace PARR.DAL.TaskServices
+{
+ ///
+ /// Управления задачами (очередями)
+ ///
+ public interface ITaskManagementService
+ {
+ ///
+ /// Создать новую задачу в БД и опубликовать в очередь
+ ///
+ ///
+ /// Тип задачи
+ /// Payload
+ /// Инициатор
+ /// Настройки очереди
+ ///
+ Task CreateTaskAsync(TaskTypeEnum typeCode, T payload, IHistoryInitiator initiator, IMqSettings mqSettings);
+ }
+}
diff --git a/PARR.DAL/TaskServices/TaskManagementService.cs b/PARR.DAL/TaskServices/TaskManagementService.cs
new file mode 100644
index 00000000..806bc345
--- /dev/null
+++ b/PARR.DAL/TaskServices/TaskManagementService.cs
@@ -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 logger;
+ private readonly ITaskTypeService taskTypeService;
+ private readonly ITaskService taskService;
+
+ public TaskManagementService(
+ ILogger logger,
+ ITaskTypeService taskTypeService,
+ ITaskService taskService
+ )
+ {
+ this.logger = logger;
+ this.taskTypeService = taskTypeService;
+ this.taskService = taskService;
+ }
+
+
+ public async Task CreateTaskAsync(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;
+ }
+
+
+ ///
+ /// Получить активную singleton задачу указанного типа
+ ///
+ ///
+ ///
+ private async Task GetActiveSingletonTaskAsync(TaskTypeEnum typeCode)
+ {
+ return await taskService.Get().AsNoTracking()
+ .FirstOrDefaultAsync(t => t.TypeCode == typeCode && (
+ t.StatusCode == TaskItemStatusEnum.Pending
+ || t.StatusCode == TaskItemStatusEnum.Processing
+ ));
+ }
+
+
+ ///
+ /// Payload конвертировать в string
+ ///
+ ///
+ ///
+ ///
+ private string PayloadToString(T payload)
+ {
+ var jsonOptions = new JsonSerializerOptions
+ {
+ Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping
+ };
+
+ return JsonSerializer.Serialize(payload, jsonOptions);
+ }
+
+
+ }
+}
diff --git a/PARR.DAL/TransformServices/INextRunModifierService.cs b/PARR.DAL/TransformServices/INextRunModifierService.cs
deleted file mode 100644
index 1b19d938..00000000
--- a/PARR.DAL/TransformServices/INextRunModifierService.cs
+++ /dev/null
@@ -1,24 +0,0 @@
-
-
-namespace PARR.DAL.TransformServices
-{
- /////
- ///// Сервис по изменению NextRun
- /////
- //public interface INextRunModifierService
- //{
- // ///
- // /// Получить NextRun с учетом тайм зоны УЗ робота ЕСПП
- // ///
- // ///
- // ///
- // DateTimeOffset GetNextRunByAccountRobotTimeZone(DateTimeOffset nextRun);
-
- // /////
- // ///// Получить РАБОЧИЙ день следующего срабатывания
- // /////
- // /////
- // /////
- // //Task GetWorkDayAsync(DateTimeOffset nextRun);
- //}
-}
diff --git a/PARR.DAL/TransformServices/NextRunModifierService.cs b/PARR.DAL/TransformServices/NextRunModifierService.cs
deleted file mode 100644
index 00c49db7..00000000
--- a/PARR.DAL/TransformServices/NextRunModifierService.cs
+++ /dev/null
@@ -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 logger;
- // private readonly IWeekendDayService weekendDayService;
-
- // public NextRunModifierService(
- // SettingsFromDb settingsFromDb,
- // ILogger logger,
- // IWeekendDayService weekendDayService
- // )
- // {
- // this.settingsFromDb = settingsFromDb;
- // this.logger = logger;
- // this.weekendDayService = weekendDayService;
- // }
-
-
- // public DateTimeOffset GetNextRunByAccountRobotTimeZone(DateTimeOffset nextRun)
- // {
- // return nextRun.AddHours(settingsFromDb.EsppRobotAccountTimeZoneHour);
- // }
-
-
- //}
-}