feat(api, core, dal, domain): В таблицу Tasks добавлено поле процент выполнения, так же добавлено это поле в response - статус формирования отчета. Создан интерфейс для отслеживания прогресса IProgressAsync. Из сервиса WorkloadService исключен ITaskManagementService.

This commit is contained in:
Mikhail Trubnikov
2026-05-20 14:26:25 +10:00
parent 8f61bbb73e
commit 68808e761b
17 changed files with 4160 additions and 41 deletions

View File

@@ -114,6 +114,10 @@ namespace PARR.Core.Services.TaskServices.Handlers
task.StatusCode = TaskItemStatusEnum.Processing;
task.DateModified = DateTimeOffset.UtcNow;
// Ставим процент 0
task.ProgressPercent = 0;
await taskRepository.UpdateProgressAsync(task.Id, 0);
return true;
}
@@ -129,7 +133,7 @@ namespace PARR.Core.Services.TaskServices.Handlers
/// <param name="task"></param>
/// <param name="result"></param>
/// <returns></returns>
private async System.Threading.Tasks.Task HandleSuccessAsync(TaskItem task, HandlerResult result)
private async Task HandleSuccessAsync(TaskItem task, HandlerResult result)
{
logger.LogInformation("Задача {TaskId} успешно выполнена", task.Id);
@@ -137,6 +141,9 @@ namespace PARR.Core.Services.TaskServices.Handlers
task.ProcessedAt = DateTimeOffset.UtcNow;
task.DateModified = DateTimeOffset.UtcNow;
// Ставим 100%
task.ProgressPercent = 100;
//todo: если есть результат, его можно сохранить, если надо
if (!string.IsNullOrEmpty(result.ResultData))
{
@@ -158,8 +165,19 @@ namespace PARR.Core.Services.TaskServices.Handlers
/// <param name="task"></param>
/// <param name="ex"></param>
/// <returns></returns>
private async System.Threading.Tasks.Task HandleErrorAsync(TaskItem task, Exception? ex)
private async Task HandleErrorAsync(TaskItem task, Exception? ex)
{
// задача в БД, у которой можно получить актуальный процент прогресса
var taskInDb = await taskRepository.Get()
.AsNoTracking()
.FirstOrDefaultAsync(t=>t.Id==task.Id);
if (taskInDb != null)
{
//так как в БД мы обновлям атомарно прогресс, то пришедший нам task не знает фактическое значение прогресса
task.ProgressPercent = taskInDb.ProgressPercent;
}
var attemptNumber = task.RetryCount + 1;
var errorMessage = ex?.Message ?? "Неизвестная ошибка";
var stackTrace = ex?.StackTrace;
@@ -219,7 +237,7 @@ namespace PARR.Core.Services.TaskServices.Handlers
/// </summary>
/// <param name="task"></param>
/// <returns></returns>
protected virtual async System.Threading.Tasks.Task ScheduleRetryAsync(TaskItem task)
protected virtual async Task ScheduleRetryAsync(TaskItem task)
{
// Здесь нужна интеграция с ITaskQueueService
// Для этого можно использовать событие или callback
@@ -254,7 +272,7 @@ namespace PARR.Core.Services.TaskServices.Handlers
/// <param name="task"></param>
/// <param name="result"></param>
/// <returns></returns>
protected virtual System.Threading.Tasks.Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask;
protected virtual Task OnSuccessAsync(TaskItem task, HandlerResult result) => System.Threading.Tasks.Task.CompletedTask;
/// <summary>
@@ -264,7 +282,43 @@ namespace PARR.Core.Services.TaskServices.Handlers
/// <param name="ex"></param>
/// <param name="attemptNumber"></param>
/// <returns></returns>
protected virtual System.Threading.Tasks.Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask;
protected virtual Task OnErrorAsync(TaskItem task, Exception? ex, int attemptNumber) => System.Threading.Tasks.Task.CompletedTask;
/// <summary>
/// Обновляет процент выполнения задачи.
/// Обновлять процент не чаще 1 раза в 5 сек.
/// Вызывается из конкретного хендлера при желании показать прогресс.
/// </summary>
/// <param name="taskId">Id задачи</param>
/// <param name="percent">Процент (0-100). Значения >100 будут ограничены до 100.</param>
/// <returns></returns>
protected async Task UpdateProgressAsync(Guid taskId, int percent)
{
if (percent < 0)
{
logger.LogWarning("Попытка установить отрицательный прогресс {Percent} для задачи {TaskId}", percent, taskId);
percent = 0;
}
else if (percent > 100)
{
logger.LogWarning("Прогресс {Percent} для задачи {TaskId} превысил 100%. Записываем 100.", percent, taskId);
percent = 100;
}
// Обновляем в БД
var affectedRows = await taskRepository.UpdateProgressAsync(taskId, percent);
if (affectedRows == 0)
{
// Задача не найдена или уже завершена — не считаем ошибкой
logger.LogDebug("Не удалось обновить прогресс задачи {TaskId} (возможно, уже завершена)", taskId);
}
else
{
logger.LogDebug("Прогресс задачи {TaskId} обновлен: {Percent}%", taskId, percent);
}
}
}

View File

@@ -1,4 +1,5 @@
using Microsoft.Extensions.Logging;
using PARR.Core.Common.Interfaces;
using PARR.Core.Repositories.Interfaces.TaskRepositories;
using PARR.Core.Services.Workload.Implementations;
using PARR.Domain.Entities.TaskEntities;
@@ -37,9 +38,18 @@ namespace PARR.Core.Services.TaskServices.Handlers
// тут логика построения отчета
// Создаём прогресс-репортер
var progress = new ProgressAsync<int>(async percent =>
{
await UpdateProgressAsync(task.Id, percent);
logger.LogDebug("Процент выполнения: {Percent}%", percent);
});
await workloadCacheService.DeleteWorkloadCacheAsync();
await workloadCacheService.DeleteTemplateCacheAsync();
await workloadCacheService.CreateTemplateCacheAsync();
await workloadCacheService.CreateTemplateCacheAsync(progress);
//await System.Threading.Tasks.Task.CompletedTask;

View File

@@ -0,0 +1,38 @@
using AutoMapper;
using Microsoft.EntityFrameworkCore;
using PARR.Core.Repositories.Interfaces.TaskRepositories;
using PARR.Domain.DTOs.TaskDTO;
using PARR.Domain.Enums;
namespace PARR.Core.Services.TaskServices.Helpers
{
/// <summary>
/// Хелпеы для TaskManagement
/// </summary>
internal static class TaskManagementHelpers
{
/// <summary>
/// Получить список активных задач
/// </summary>
/// <param name="taskRepository"></param>
/// <param name="mapper"></param>
/// <param name="typeCode"></param>
/// <returns></returns>
public static async Task<List<ActiveTask>> GetActiveTasksAsync(ITaskRepository taskRepository, IMapper mapper, 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);
}
}
}

View File

@@ -3,6 +3,7 @@ using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using PARR.Core.Common.Interfaces.RabbitServices;
using PARR.Core.Repositories.Interfaces.TaskRepositories;
using PARR.Core.Services.TaskServices.Helpers;
using PARR.Core.Services.TaskServices.Interfaces;
using PARR.Domain.Common.Rabbit.Messages;
using PARR.Domain.DTOs.TaskDTO;
@@ -112,18 +113,20 @@ namespace PARR.Core.Services.TaskServices.Implementations
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);
//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);
//if (typeCode.HasValue)
// query = query.Where(t => t.TypeCode == typeCode.Value);
var tasks = await query.OrderByDescending(t => t.DateCreated).ToListAsync();
//var tasks = await query.OrderByDescending(t => t.DateCreated).ToListAsync();
return mapper.Map<List<ActiveTask>>(tasks);
//return mapper.Map<List<ActiveTask>>(tasks);
return await TaskManagementHelpers.GetActiveTasksAsync(taskRepository, mapper, typeCode);
}

View File

@@ -177,6 +177,7 @@ namespace PARR.Core.Services.TaskServices.Implementations
task.RetryCount++;
task.DateModified = DateTimeOffset.UtcNow;
task.ProgressPercent = 0;
await taskRepository.CommitAsync();

View File

@@ -82,7 +82,7 @@ namespace PARR.Core.Services.Workload.Implementations
/// КЭШ рабочих групп и зон ответственности
/// </summary>
/// <returns></returns>
public async Task CreateTemplateCacheAsync()
public async Task CreateTemplateCacheAsync(IProgressAsync<int> progress)
{
var dtStart = DateTimeOffset.UtcNow;
@@ -97,10 +97,25 @@ namespace PARR.Core.Services.Workload.Implementations
logger.LogInformation("Получил из БД шаблонов {TemplateCount} шт.", templates.Count);
//Ставим сразу 5% - получили из БД
await progress.ReportAsync(5);
if (templates.Count == 0)
{
logger.LogWarning("Список шаблонов пуст, кэш не формируется");
await progress.ReportAsync(100);
return;
}
//2. Получить по каждому шаблону ЗО и РГ
var cacheData = new List<TemplateCache>();
int templateCount = 0;
int totalTemplates = templates.Count;
// последний отправленный процент
int lastReportedPercent = 5;
foreach (var template in templates)
{
var (workGroup, responseArea) = await ApplyShortcodesAsync(template);
@@ -116,8 +131,26 @@ namespace PARR.Core.Services.Workload.Implementations
cacheData.Add(cacheItem);
templateCount++;
#region Расчет и отрпавка прогресса
// Считаем текущий процент(0 - 100). Этап идёт от 5% до 95%
var currentPercent = 5 + (int)((templateCount * 90.0) / totalTemplates);
// Округляем вниз до кратного 5 (33→30, 37→35, 99→95)
var roundedPercent = (currentPercent / 5) * 5;
// Отправляем, только если прогресс реально изменился (защита от дублей)
if (roundedPercent > lastReportedPercent && roundedPercent <= 95)
{
await progress.ReportAsync(roundedPercent);
lastReportedPercent = roundedPercent;
}
#endregion
}
//95%
await progress.ReportAsync(95);
//3. Сохранить в Redis
logger.LogDebug("Для сохранения в Redis подготовлен массив из {Count} строк", cacheData.Count);
var hashKey = redisCacheService.GetKey(cacheTemplateKey);
@@ -127,6 +160,8 @@ namespace PARR.Core.Services.Workload.Implementations
await redisCacheService.SetHashFieldAsync(hashKey, item.Data.TemplateId.ToString(), item, cacheTemplateTtl);
}
await progress.ReportAsync(100);
var duration = DateTimeOffset.UtcNow - dtStart;
logger.LogInformation("Завершено формирование кэша. Key: {HashKey}, ttl: {Ttl}, кол-во записей: {CacheCount}. Продолжительность формирования КЭШ {Duration}", hashKey, cacheTemplateTtl, cacheData.Count, duration);
}

View File

@@ -1,9 +1,10 @@
using Microsoft.Extensions.Logging;
using AutoMapper;
using Microsoft.Extensions.Logging;
using PARR.Core.Repositories.Interfaces.TaskRepositories;
using PARR.Core.Services.NextRunServices;
using PARR.Core.Services.TaskServices.Interfaces;
using PARR.Core.Services.TaskServices.Helpers;
using PARR.Core.Services.Workload.Interfaces;
using PARR.Core.Services.Workload.Models;
using PARR.Domain.DTOs.TaskDTO;
using PARR.Domain.DTOs.Workload;
using PARR.Domain.Enums;
@@ -14,19 +15,26 @@ namespace PARR.Core.Services.Workload.Implementations
private readonly ILogger<WorkloadService> logger;
private readonly WorkloadCacheService workloadCacheService;
private readonly INextRunService nextRunService;
private readonly ITaskManagementService taskManagementService;
private readonly ITaskRepository taskRepository;
private readonly IMapper mapper;
//private readonly ITaskManagementService taskManagementService;
public WorkloadService(
ILogger<WorkloadService> logger,
WorkloadCacheService workloadCacheService,
INextRunService nextRunService,
ITaskManagementService taskManagementService
ITaskRepository taskRepository,
IMapper mapper
//ITaskManagementService taskManagementService
)
{
this.logger = logger;
this.workloadCacheService = workloadCacheService;
this.nextRunService = nextRunService;
this.taskManagementService = taskManagementService;
this.taskRepository = taskRepository;
this.mapper = mapper;
//this.taskManagementService = taskManagementService;
}
@@ -122,14 +130,15 @@ namespace PARR.Core.Services.Workload.Implementations
// 1. проверить что выполнились все задания на обновление кэш
var activeTasks = await taskManagementService.GetActiveTasksAsync(TaskTypeEnum.Workload);
//var activeTasks = await taskManagementService.GetActiveTasksAsync(TaskTypeEnum.Workload);
var activeTasks = await TaskManagementHelpers.GetActiveTasksAsync(taskRepository, mapper, TaskTypeEnum.Workload);
if (activeTasks.Count > 0)
{
// идет формирование отчета
// задача может быт только одна, так как у нее такой тип
var task = activeTasks.First();
return new WorkloadCacheInfo { Status = WorkloadCacheInfo.StatusEnum.Processing, Date = task.DateCreated };
return new WorkloadCacheInfo { Status = WorkloadCacheInfo.StatusEnum.Processing, Date = task.DateCreated, ProgressPercent = task.ProgressPercent };
}
// 2. получить дату кэш