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 ) { this._logger = logger; this._robotConfigurationRepository = robotConfigurationRepository; this._settingsFromDb = settingsFromDb; this._robotHistoryRepository = robotHistoryRepository; this._mapper = mapper; this._shortcodesService = shortcodesService; this._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); 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).ToList(); // Берем целевые таски, находим в них задания которым надо поставить ошибку tasksToSetErrorStatus = tasks.Where(t => errorTemplateNames.Contains(t.TemplateName)).Select(t => t.TaskId).ToList(); // Устанавливаем ошибку целевым + пишем комментарий от робота + нажимаем комит var logMessage = "[RobotTaskService] Установлен статус ошибки, так как не переименован связанный шаблон"; await SetErrorStatusAsync(tasksToSetErrorStatus, logMessage); } var endDate = DateTimeOffset.UtcNow.Add(-_settingsFromDb.RobotWaitTime); // Смотрим статусы роботов, можно взять в работу, только если (RobotStatus == Wait) или (InpRogress но которые еще не просрочены) var allowedTasks = renameTasks.Where(t => t.RobotStatusCode == (int)RobotStatusEnum.Wait || (t.RobotStatusCode == (int)RobotStatusEnum.InProgress && t.AttemptsNumber < _settingsFromDb.RobotAttemptsNumber && t.LastRobotStatusUpdated < endDate) ); // Формируем список заданий var originalCount = tasks.Count; // Запоминаем сколько было изначально // Из исходных тасков удалить те которым установлен статус ошибки var filteredOriginalTasks = tasks.Where(t => !tasksToSetErrorStatus.Contains(t.TaskId)).ToList(); var errorCount = tasksToSetErrorStatus.Count; // Столько ушло в ошибку // Определяем, какие имена шаблонов мы БУДЕМ подменять на задачи переименования var allowedRenameTemplateNames = allowedTasks.Select(t => t.Template!.Name).ToList(); // Из исходных удаляем те, которые мы сейчас заменим (подменим) var resultsTasks = filteredOriginalTasks.Where(t => !allowedRenameTemplateNames.Contains(t.TemplateName)).ToList(); // Считаем сколько именно задач мы ВЫКИНУЛИ из исходного списка ради подмены var replacedCount = filteredOriginalTasks.Count - resultsTasks.Count; // В исходные добавляем подменные таски для переименования (мапим их в RobotTaskDetails) var mappedRenameTasks = allowedTasks.Select(t => new RobotTaskDetails(t.Id, t.Template!.Name, t.Template.NextRun)).ToList(); resultsTasks.AddRange(mappedRenameTasks); var injectedCount = mappedRenameTasks.Count; // Столько задач переименования добавили взамен // Логируем итоговую статистику трансформации пула задач // TODO: Позже этот лог можно понизить до Debug!!! _logger.LogInformation( "Трансформация пула задач завершена. Исходных: {OriginalCount} шт. " + "Отклонено (ошибка): {ErrorCount} шт. Удалено обычных для подмены: {ReplacedCount} шт. " + "Внедрено задач переименования: {InjectedCount} шт. Итого к выдаче: {FinalCount} шт.", originalCount, errorCount, replacedCount, injectedCount, resultsTasks.Count); // Затем отсортируем по NextRun, чтоб ближайшие были выше return resultsTasks.OrderBy(t => t.NextRun).ToList(); } /// /// Установить статус задания - ошибка /// /// /// private async Task SetErrorStatusAsync(List taskIds, string logMessage) { 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); } } /// /// Взять задачу в работу /// /// /// 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