using AutoMapper;
using InfluxDB.Client.Api.Domain;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using PARR.BLL.Helpers;
using PARR.Core.Repositories.Interfaces;
using PARR.Core.Repositories.Interfaces.TemplateRepositories;
using PARR.Core.Services.NextRunServices;
using PARR.Core.Services.RobotTask.Interfaces;
using PARR.Core.Services.RobotTask.Models;
using PARR.Core.Services.Shortcodes;
using PARR.Domain.Common.Template;
using PARR.Domain.DTOs.RobotTask;
using PARR.Domain.Entities;
using PARR.Domain.Entities.Base.History;
using PARR.Domain.Entities.RobotEntities;
using PARR.Domain.Enums;
using PARR.Domain.Exceptions;
using PARR.Domain.Settings;
namespace PARR.Core.Services.RobotTask.Implementations
{
internal class RobotTaskService : IRobotTaskService
{
///
/// Количество заданий которые рассматриваем для взятия в работу.
/// Рекомендованное значение, кол-во роботов * 3
///
private readonly int TakeTasks = 15 * 3;
private readonly ILogger _logger;
private readonly IRobotConfigurationRepository _robotConfigurationRepository;
private readonly SettingsFromDb _settingsFromDb;
private readonly IRobotHistoryRepository _robotHistoryRepository;
private readonly IMapper _mapper;
private readonly IShortcodesService _shortcodesService;
private readonly INextRunService _nextRunService;
private readonly ITemplateRenamePendingRepository _templateRenamePendingRepository;
public RobotTaskService(
ILogger logger,
IRobotConfigurationRepository robotConfigurationRepository,
SettingsFromDb settingsFromDb,
IRobotHistoryRepository robotHistoryRepository,
IMapper mapper,
IShortcodesService shortcodesService,
INextRunService nextRunService,
ITemplateRenamePendingRepository templateRenamePendingRepository
)
{
_logger = logger;
_robotConfigurationRepository = robotConfigurationRepository;
_settingsFromDb = settingsFromDb;
_robotHistoryRepository = robotHistoryRepository;
_mapper = mapper;
_shortcodesService = shortcodesService;
_nextRunService = nextRunService;
_templateRenamePendingRepository = templateRenamePendingRepository;
}
public async Task GetTemplateTaskAsync(TaskStatusEnum taskStatusCode, bool acquireTask, string? robotIp, string? robotId)
{
var templateTask = await GetTaskAsync(RobotsEnum.TemplateOrder, taskStatusCode, acquireTask, robotIp, robotId, TimeSpan.Zero);
var task = _mapper.Map(templateTask);
task = task with { FullDescription = NormalizeLineEndingsToCrlf(await _shortcodesService.ApplyShortcodesAsync(task.FullDescription, templateTask.Template!)) };
task = task with { ShortDescription = await _shortcodesService.ApplyShortcodesAsync(task.ShortDescription, templateTask.Template!) };
task = task with { Solution = NormalizeLineEndingsToCrlf(await _shortcodesService.ApplyShortcodesAsync(task.Solution, templateTask.Template!)) };
task = task with { TnkName = await _shortcodesService.ApplyShortcodesAsync(task.TnkName, templateTask.Template!) };
task = task with { WorkName = await _shortcodesService.ApplyShortcodesAsync(task.WorkName, templateTask.Template!) };
task = task with { WorkGroup = await _shortcodesService.ApplyShortcodesAsync(task.WorkGroup, templateTask.Template!) };
task = task with { ResponseArea = await _shortcodesService.ApplyShortcodesAsync(task.ResponseArea, templateTask.Template!) };
task = task with { ClosingCode = _settingsFromDb.ClosingCode };
task = task with { Initiator = _settingsFromDb.Initiator };
task = task with { Category = _settingsFromDb.Category };
return task;
}
public async Task GetScheduleTaskAsync(TaskStatusEnum taskStatusCode, bool acquireTask, string? robotIp, string? robotId, IHistoryInitiator historyInitiator, TimeSpan scheduleCooldownDuration)
{
var scheduleTask = await GetTaskAsync(RobotsEnum.ScheduleOrder, taskStatusCode, acquireTask, robotIp, robotId, scheduleCooldownDuration);
// Проверяем nextRun, lastRun, обновляем их
var resultUpdateNextRun = await UpdateNextRunAsync(scheduleTask, historyInitiator);
if (!resultUpdateNextRun)
{
_logger.LogError("Ошибка при расчете NextRun для templateId: {templateId}", scheduleTask.TemplateId);
throw new NextRunException($"Ошибка при расчете NextRun для templateId: {scheduleTask.TemplateId}");
}
var task = _mapper.Map(scheduleTask);
task = task with { Timezone = _settingsFromDb.EsppScheduleTimezone };
task = task with { WorkGroup = await _shortcodesService.ApplyShortcodesAsync(task.WorkGroup, scheduleTask.Template!) };
task = task with { ResponseArea = await _shortcodesService.ApplyShortcodesAsync(task.ResponseArea, scheduleTask.Template!) };
//nextRun в часовой зоне УЗ Робота ЕСПП
var nextRunWithRobotTz = scheduleTask.Template!.NextRun.Add(_nextRunService.GetEsppAccountOffset());
//на всякий случай еще раз проверяем, что дата не устарела и отправляем задание
if (nextRunWithRobotTz < DateTimeOffset.UtcNow)
{
_logger.LogError("Ошибка при расчете NextRun для templateId: {templateId}, итоговое значение для робота, меньше чем сейчас {nextRunWithRobotTz}<{now}",
task.TemplateId, nextRunWithRobotTz, DateTimeOffset.UtcNow);
throw new NextRunException($"Ошибка при расчете NextRun для templateId: {scheduleTask.TemplateId}");
}
task = task with { NextStart = EsppScheduleHelpers.GetNextRun(nextRunWithRobotTz) };
task = task with { GenerationTime = EsppScheduleHelpers.GetGenerationTime(nextRunWithRobotTz) };
task = task with { RepeatRange = _settingsFromDb.ScheduleRepeatRange };
task = task with { };
return task;
}
///
/// Получить задачу для робота.
/// Метод может генерировать исключения.
///
///
///
/// Взять в работу
///
///
///
private async Task GetTaskAsync(RobotsEnum robotCode, TaskStatusEnum taskStatusCode, bool acquireTask, string? robotIp, string? robotId, TimeSpan scheduleCooldownDuration)
{
// 1. Ищем все задания с превышенным кол-вом попыток и просроченным временем, ставим им статус ошибки
await _robotConfigurationRepository.MarkExpiredTasksAsFailedAsync(_settingsFromDb.RobotAttemptsNumber, _settingsFromDb.RobotWaitTime);
// 2. Ищем доступные задания
var availableTasks = await GetAvailableTasksAsync(robotCode, taskStatusCode, scheduleCooldownDuration);
if (availableTasks.Count == 0)
throw new NotFoundException("Нет доступных заданий для робота");
Guid? acquiredTaskId = null;
if (acquireTask)
{
// Берем задание в работу, устанавливаем ему статус "В работе"
acquiredTaskId = await AcquireTaskAsync(availableTasks, robotIp, robotId);
if (acquiredTaskId == null)
throw new Exception($"Не удалось взять ни одну из доступных задач ({availableTasks.Count}) в работу");
}
else
{
// Берем первую задачу из списка доступных
acquiredTaskId = availableTasks.First();
_logger.LogDebug("Задача не требует захвата, взята первая из доступных: {TaskId}", acquiredTaskId);
}
//3. Получаем задачу со всеми нужными инклудами в зависимости от типа робота
var task = await GetTaskWithAllDataAsync(acquiredTaskId.Value, robotCode);
return task;
}
///
/// Получить список возможных заданий для взятия в работу.
/// Кол-во заданй ограничено переменной TakeTasks
///
///
///
///
private async Task> GetAvailableTasksAsync(RobotsEnum robotCode, TaskStatusEnum taskStatusCode, TimeSpan scheduleCooldownDuration)
{
var query = _robotConfigurationRepository.Get()
.AsNoTracking()
.Where(t => t.RobotCode == (int)robotCode);
// Если это задание для робота расписаний
if (robotCode == RobotsEnum.ScheduleOrder)
{
// Выбираем только записи с созданными шаблонами (у которых статус 30), а только потом ищем у них расписания
query = query.Where(t => t.Template!.RobotConfigurations.Any(x => x.RobotCode == (int)RobotsEnum.TemplateOrder && x.TaskStatusCode == (int)TaskStatusEnum.Ok));
// Не берем шаблоны, у которых lastRun + 3 часа < сейчас, и у них последний инициатор был или nextRun (10) или esppSchedule (5), это условие применяется только к активированным расписаниям
var cooldownThreshold = DateTimeOffset.UtcNow.Add(-scheduleCooldownDuration);
query = query.Where(t =>
// Условие кулдауна: проверяем, попадает ли шаблон под ЗАПРЕТ
!(
t.Template!.IsActiveSchedule
&& (t.Template.InitiatorParrComponentId == ParrComponentsEnum.EsppScheduleSync || t.Template.InitiatorParrComponentId == ParrComponentsEnum.NextRun)
&& t.Template.LastRun >= cooldownThreshold
)
);
}
// Сортируем по nextRun, чтобы те, у кого nextRun ближе к текущей, выполнились скорее
query = query.OrderBy(t => t.Template!.NextRun).ThenBy(t => t.Template!.IsActiveSchedule).ThenBy(t => t.Template!.IsActiveTemplate);
// Кандидаты заданий, Id задания и имя шаблона
//var tasks = new List();
var tasks = new List();
// Ищем первые TakeTasks заданий в статусе ОЖИДАНИЕ
tasks = await query
.Where(t =>
t.RobotStatusCode == (int)RobotStatusEnum.Wait
&& t.TaskStatusCode == (int)taskStatusCode
).Take(TakeTasks)
//.Select(t => t.Id)
.Select(t => new RobotTaskDetails(t.Id, t.Template!.Name, t.Template.NextRun))
.ToListAsync();
_logger.LogDebug("Найдено заданий в статусе 'Ожидание' {Count} шт. Робот '{Robot}'", tasks.Count, robotCode.ToString());
if (tasks.Count == 0)
{
// Ищем задания в статусе В РАБОТЕ, которые можно перезапустить
// Поиск по `RobotStatusCode` = 22.
// Далее проверяется `LastStatusUpdated`, что время последнего смены статуса не превышает допустимого(берется из настроек, поле `RobotWaitTime`)
// и что текущая попытка не больше разрешенной(берется из настроек, поле `RobotAttemptsNumber`) - если это так, берется эта запись.
var endDate = DateTimeOffset.UtcNow.Add(-_settingsFromDb.RobotWaitTime);
tasks = await query.Where(t => t.RobotStatusCode == (int)RobotStatusEnum.InProgress
&& t.TaskStatusCode == (int)taskStatusCode
&& t.AttemptsNumber < _settingsFromDb.RobotAttemptsNumber
&& t.LastRobotStatusUpdated < endDate)
.Take(TakeTasks)
//.Select(t => t.Id)
.Select(t => new RobotTaskDetails(t.Id, t.Template!.Name, t.Template.NextRun))
.ToListAsync();
_logger.LogDebug("Найдено заданий в статусе 'В работе' {Count} шт. Робот '{Robot}'", tasks.Count, robotCode.ToString());
}
if (robotCode == RobotsEnum.TemplateOrder)
{
// Если запрашиваем шаблоны, смотрим корректируем список заданий в зависимости от статуса переименования.
// Это не относится к расписаниям, потому что у переименованных расписаний статус Updating, а оно не возьмется в работу, пока не обновится шаблон
tasks = await ReplaceTemplateTasksForRenameAsync(tasks, robotCode);
}
return tasks.Select(t => t.TaskId).ToList();
}
///
/// Проверяет наличие шаблонов в процессе переименования и заменяет обычные задания на задания по переименованию.
/// Если связанный шаблон не переименован, и у него статус ошибки, целевому шаблону устанавливается статус ошибки.
///
///
///
private async Task> ReplaceTemplateTasksForRenameAsync(List tasks, RobotsEnum robotCode)
{
if (tasks.Count == 0 || robotCode != RobotsEnum.TemplateOrder)
return tasks;
// Ищем есть ли связанные шаблоны с таким имененм на переименование
var taskTemplateNames = tasks.Select(t => t.TemplateName).Distinct().ToList();
// Ищем записи в таблице переименований, где OldName совпадает с именами наших новых задач
var templatesToRename = await _templateRenamePendingRepository.Get().AsNoTracking()
.Where(t => taskTemplateNames.Contains(t.OldName))
.ToListAsync();
_logger.LogDebug("Найдено шаблонов в процессе переименования для текущих задач: {Count} шт.", templatesToRename.Count);
if (templatesToRename.Count == 0)
return tasks;
// Ищем конфигурации роботов для СТАРЫХ шаблонов (которые переименовываются) по ИД, смотрим, можем ли взять их в работу
var renameTemplateIds = templatesToRename.Select(t => t.TemplateId).ToList();
var renameTasks = await _robotConfigurationRepository.Get()
.AsNoTracking()
.Include(t => t.Template)
.Where(t =>
t.RobotCode == (int)robotCode
&& renameTemplateIds.Contains(t.TemplateId)
// Это может быть только обновление. Так как переименования для создаваемого шаблона быть не может
&& t.TaskStatusCode == (int)TaskStatusEnum.Updating
).ToListAsync();
// --- Блок обработки ошибок ---
// Если старый шаблон в ошибке и лимит попыток исчерпан, ставим ошибку и новому шаблону
var errorTasks = renameTasks
.Where(t =>
t.RobotStatusCode == (int)RobotStatusEnum.Error
&& t.AttemptsNumber >= _settingsFromDb.RobotAttemptsNumber
).ToList();
var tasksToSetErrorStatus = new List();
if (errorTasks.Count > 0)
{
_logger.LogDebug("Найдено старых заданий на переименование с ошибками: {ErrorCount}. Ставим ошибку целевым (новым) заданиям.", errorTasks.Count);
var errorTemplateNames = errorTasks.Select(t => t.Template!.Name).ToHashSet();
// Берем целевые таски, находим в них задания которым надо поставить ошибку
tasksToSetErrorStatus = tasks
.Where(t => errorTemplateNames.Contains(t.TemplateName))
.Select(t => t.TaskId)
.ToList();
if (tasksToSetErrorStatus.Count > 0)
{
// Устанавливаем ошибку целевым + пишем комментарий от робота + нажимаем комит
var logMessage = "[RobotTaskService] Установлен статус ошибки, так как не переименован связанный шаблон";
await SetErrorStatusAsync(tasksToSetErrorStatus, logMessage);
}
}
// --- Блок подмены задач ---
var endDate = DateTimeOffset.UtcNow.Add(-_settingsFromDb.RobotWaitTime);
// Фильтруем старые задачи, которые МОЖНО взять в работу. Смотрим статусы роботов, можно взять в работу, только если (RobotStatus == Wait) или (InpRogress но которые еще не просрочены)
var allowedRenameTasks = renameTasks.Where(t =>
t.RobotStatusCode == (int)RobotStatusEnum.Wait
|| (t.RobotStatusCode == (int)RobotStatusEnum.InProgress
&& t.AttemptsNumber < _settingsFromDb.RobotAttemptsNumber
&& t.LastRobotStatusUpdated < endDate)
).ToList();
// Словарь для поиска подменной задачи по имени шаблона.
// GroupBy + First на случай, если в бд есть дубликаты, но такого быть не может
var renameTasksToDictionary = allowedRenameTasks
.GroupBy(t => t.Template!.Name)
.ToDictionary(
t => t.Key,
t => new RobotTaskDetails(t.First().Id, t.First().Template!.Name, t.First().Template!.NextRun)
);
var errorTaskIdsSet = tasksToSetErrorStatus.ToHashSet();
// Создаем итоговый список
var finalTasks = new List(tasks.Count);
int replacedCount = 0;
int errorCount = errorTaskIdsSet.Count;
// Проходим по ИСХОДНОМУ списку, чтобы сохранить его порядок сортировки
foreach (var task in tasks)
{
// Если задаче нужно поставить ошибку, просто пропускаем ее (она не попадет в итоговый список)
if (errorTaskIdsSet.Contains(task.TaskId))
{
continue;
}
// Если для этого имени шаблона есть разрешенная задача на переименование - вставляем ее на место текущей
if (renameTasksToDictionary.TryGetValue(task.TemplateName, out var renameTask))
{
finalTasks.Add(renameTask);
replacedCount++;
}
else
{
// Иначе оставляем исходную задачу на месте
finalTasks.Add(task);
}
}
_logger.LogInformation(
"Трансформация пула задач (Rename). Исходных: {OriginalCount}. Отклонено (Error): {ErrorCount}. " +
"Заменено на старые: {ReplacedCount}. Итого к выдаче: {FinalCount}",
tasks.Count, errorCount, replacedCount, finalTasks.Count
);
// Возвращаем без дополнительной сортировки по NextRun. Порядок сохранен начального списка
return finalTasks;
}
///
/// Установить статус задания - ошибка
///
///
///
private async Task SetErrorStatusAsync(List taskIds, string logMessage)
{
if (taskIds == null || taskIds.Count == 0)
return;
var tasks = await _robotConfigurationRepository.Get()
.Include(t => t.Template)
.Where(t => taskIds.Contains(t.Id))
.ToListAsync();
if (tasks.Count == 0)
return;
foreach (var task in tasks)
{
// Так как это целевой шаблон, то ставим ему сразу максимальное кол-во попыток и ошибку, чтоб больше он не выдавался в заданиях, пока не исправим связанный
// Устанавливаем статус ошибки
_robotConfigurationRepository.SetErrorRobotStatusAndMaxAttempts(task);
// Пишем в лог роботу
var history = new RobotHistory
{
Id = Guid.NewGuid(),
HistoryLevel = (int)RobotStatusEnum.Error,
TaskStatusCode = task.TaskStatusCode,
RobotConfigurationId = task.Id,
RobotIp = null,
RobotId = ParrComponentsEnum.Api.ToString(),
RobotMessage = logMessage
};
await _robotHistoryRepository.CreateAsync(history);
_logger.LogInformation("Для целевого задания {TaskId} (шаблон '{TemplateName}') установлен статус ошибки, " +
"так как связанное задание со старым шаблоном не было успешно выполнено.",
task.Id, task.Template!.Name);
}
if (await _robotHistoryRepository.CommitAsync())
{
_logger.LogDebug("Установлен статус 'Ошибка', для заданий {TaskCount} шт.", tasks.Count);
}
else
{
_logger.LogError("Ошибка при установке статуса задания 'Ошибка', для заданий {TaskCount} шт. Транзакция отменена", tasks.Count);
throw new DbErrorException("Не удалось сохранить изменения статусов заданий при обработке переименования шаблона.");
}
}
///
/// Взять задачу в работу
///
///
///
private async Task AcquireTaskAsync(List tasks, string? robotIp, string? robotId)
{
foreach (var taskId in tasks)
{
var isChangedStatus = await _robotConfigurationRepository.SetInProgressStatusAsync(taskId);
if (isChangedStatus)
{
_logger.LogDebug("Захвачена задача {TaskId}", taskId);
var task = await _robotConfigurationRepository.Get()
.AsNoTracking()
.FirstAsync(t => t.Id == taskId);
// пишем в историю робота
var history = new RobotHistory
{
Id = Guid.NewGuid(),
HistoryLevel = (int)RobotHistoryLevelEnum.Start,
TaskStatusCode = task.TaskStatusCode,
RobotConfigurationId = taskId,
RobotIp = robotIp,
RobotId = robotId
};
if (!await _robotHistoryRepository.CreateAsync(history) || !await _robotHistoryRepository.CommitAsync())
throw new DbErrorException("Ошибка при добавлении истории робота, при взятии задания в работу.");
return taskId;
}
else
{
_logger.LogDebug("Не удалось захватить задачу {TaskId}", taskId);
}
}
_logger.LogDebug("Не удалось захватить ни одну из доступных задач для робота");
return null;
}
///
/// Получить задачу со всем необходимыми полями
///
///
///
///
private async Task GetTaskWithAllDataAsync(Guid taskId, RobotsEnum robotCode)
{
IQueryable query = _robotConfigurationRepository.Get()
//.AsNoTracking() // нужно обязательно трекать, так как может измениться nextRun и его нужно будет сохранить
.AsSingleQuery()
// Общие инклуды для шаблонов и расписаний
// Units
.Include(t => t.Template)
.ThenInclude(t => t!.Unit)
.ThenInclude(t => t!.UnitValues)
.ThenInclude(t => t!.Field)
.Include(t => t.Template)
.ThenInclude(t => t!.Unit)
.ThenInclude(t => t!.UnitValues)
.ThenInclude(t => t!.Value)
.Include(t => t.Template)
.ThenInclude(t => t!.UnitsInTemplate)
// Группа с типом
.Include(t => t.Template)
.ThenInclude(t => t!.Job)
.ThenInclude(t => t!.Group)
.ThenInclude(t => t!.GroupType)
// ТНК
.Include(t => t.Template)
.ThenInclude(t => t!.Job)
.ThenInclude(t => t!.Tnk)
.ThenInclude(t => t!.Subprocess)
.ThenInclude(t => t!.Process);
if (robotCode == RobotsEnum.ScheduleOrder)
{
// Инклуды только для расписаний
query = query
// Расписание
.Include(t => t.Template)
.ThenInclude(t => t!.Job)
.ThenInclude(t => t!.Group)
.ThenInclude(t => t!.EsppSchValues)
.ThenInclude(t => t!.EsppSchTypeConfig)
.ThenInclude(t => t!.EsppSchTypeSchedule)
// Исключения по типу
.Include(t => t.Template)
.ThenInclude(t => t!.Job)
.ThenInclude(t => t!.Group)
.ThenInclude(t => t!.ScheduleExcludeType)
// Исключения по календарю
.Include(t => t.Template)
.ThenInclude(t => t!.Job)
.ThenInclude(t => t!.Group)
.ThenInclude(t => t!.ScheduleExcludeTypeCalendar);
}
return await query.FirstAsync(t => t.Id == taskId);
}
///
/// Приводит переносы строк в тексте к формату CRLF (\r\n)
///
/// Исходный текст
/// Текст с унифицированными переносами строк
private string NormalizeLineEndingsToCrlf(string? text)
{
if (string.IsNullOrEmpty(text))
return string.Empty;
// Заменяем любые варианты переносов (\r\n, \r, \n) на единый \r\n
return System.Text.RegularExpressions.Regex.Replace(text, @"\r\n|\r|\n", "\r\n");
}
///
/// Обоновить NextRun если он устарел
///
///
///
private async Task UpdateNextRunAsync(RobotConfiguration task, IHistoryInitiator historyInitiator)
{
var template = task.Template!;
//var nextRun = await esppScheduleTransformService.GetNextDateAsync(template.Job!.GroupId, template!.Job!.Group!.ReferenceDate);
var nextRun = await _nextRunService.GetNextRunForTemplateAsync(template.Id, false);
if (!nextRun.HasValue)
{
_logger.LogError("При обновлении nextRun для шаблона {TemplateId}, расчитанный nextRun=null, ошибка в расчетах.", template.Id);
return false;
}
if (nextRun.Value < DateTimeOffset.UtcNow)
{
_logger.LogError("При обновлении nextRun для шаблона {TemplateId}, расчитанный nextRun