Merge branch 'report' into dev
This commit is contained in:
@@ -22,7 +22,6 @@ using PARR.DAL.Services.Interfaces;
|
||||
using PARR.DAL.Services.Interfaces.Job;
|
||||
using PARR.DAL.Services.Interfaces.Schedule;
|
||||
using PARR.DAL.Services.Interfaces.Unit;
|
||||
using PARR.DAL.TaskServices;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
namespace PARR.DAL
|
||||
@@ -154,8 +153,6 @@ namespace PARR.DAL
|
||||
services.AddScoped<ITaskRepository, TaskRepository>();
|
||||
services.AddScoped<ITaskTypeRepository, TaskTypeRepository>();
|
||||
|
||||
services.AddTransient<ITaskManagementService, TaskManagementService>();
|
||||
|
||||
#endregion
|
||||
|
||||
//services.AddTransient<INextRunModifierService, NextRunModifierService>();
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.EntityFrameworkCore.ChangeTracking;
|
||||
using Microsoft.EntityFrameworkCore.Metadata.Internal;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Repositories.Base;
|
||||
using PARR.DAL.Context;
|
||||
using PARR.Domain.Common.Pagination;
|
||||
using PARR.Domain.Entities.Attributes;
|
||||
using PARR.Domain.Entities.Base;
|
||||
using PARR.Domain.Entities.Base.History;
|
||||
using PARR.Domain.Entities.Base.History.Base;
|
||||
@@ -99,6 +101,9 @@ namespace PARR.DAL.Repositories.Base
|
||||
/// <param name="initiator"></param>
|
||||
private void SetInitiator(IHistoryInitiator? initiator)
|
||||
{
|
||||
if (initiator == null)
|
||||
return;
|
||||
|
||||
logger.LogDebug("Устанавливаю инициатора для изменений");
|
||||
|
||||
// Задаем инициатора только для новых и измененных записей
|
||||
@@ -130,10 +135,28 @@ namespace PARR.DAL.Repositories.Base
|
||||
/// <param name="obj"></param>
|
||||
private void DateModifiedResolver(EntityEntry obj)
|
||||
{
|
||||
if (obj.Entity is IBaseEntityDateModified)
|
||||
if (obj.Entity is IBaseEntityDateModified entity)
|
||||
{
|
||||
logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name);
|
||||
(obj.Entity as IBaseEntityDateModified)!.DateModified = DateTimeOffset.UtcNow;
|
||||
//logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name);
|
||||
//(obj.Entity as IBaseEntityDateModified)!.DateModified = DateTimeOffset.UtcNow;
|
||||
|
||||
var entityType = entity.GetType();
|
||||
|
||||
//Ищем DateModified
|
||||
var properyInfo = entityType.GetProperty(nameof(IBaseEntityDateModified.DateModified));
|
||||
|
||||
// Проверяем наличие атрибута ManualControlAttribute на этом свойстве
|
||||
bool isManual = properyInfo?.GetCustomAttributes(typeof(ManualControlAttribute), false).Any() ?? false;
|
||||
|
||||
if (!isManual)
|
||||
{
|
||||
logger.LogDebug("Обновляю DateModified для сущности типа {EntityType}", obj.Entity.GetType().Name);
|
||||
entity.DateModified = DateTimeOffset.UtcNow;
|
||||
}
|
||||
else
|
||||
{
|
||||
logger.LogDebug("Пропуск обновления DateModified (ManualControl) для {EntityType}", entityType.Name);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -319,7 +342,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;
|
||||
}
|
||||
|
||||
// метод атомарного взятия в работу
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,23 +0,0 @@
|
||||
using PARR.Domain.Entities.Base.History;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -1,124 +0,0 @@
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using PARR.Core.Repositories.Interfaces.TaskRepositories;
|
||||
using PARR.Domain.Entities.Base.History;
|
||||
using PARR.Domain.Entities.TaskEntities;
|
||||
using PARR.Domain.Enums;
|
||||
using PARR.Domain.Settings;
|
||||
using System.Text.Encodings.Web;
|
||||
using System.Text.Json;
|
||||
|
||||
namespace PARR.DAL.TaskServices
|
||||
{
|
||||
internal class TaskManagementService : ITaskManagementService
|
||||
{
|
||||
private readonly ILogger<TaskManagementService> logger;
|
||||
private readonly ITaskTypeRepository taskTypeService;
|
||||
private readonly ITaskRepository taskService;
|
||||
|
||||
public TaskManagementService(
|
||||
ILogger<TaskManagementService> logger,
|
||||
ITaskTypeRepository taskTypeService,
|
||||
ITaskRepository 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);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user