using InfluxDB.Client.Api.Domain; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Logging; using PARR.Core.Repositories.Interfaces; using PARR.Core.Services.RobotTask.Interfaces; using PARR.Domain.Entities; using PARR.Domain.Enums; using PARR.Domain.Exceptions; using PARR.Domain.Settings; namespace PARR.Core.Services.RobotTask.Implementations { internal class RobotTaskService : IRobotTaskService { /// /// Количество заданий которые рассматриваем для взятия в работу. /// private readonly int TakeTasks = 10; private readonly ILogger logger; private readonly IRobotConfigurationRepository robotConfigurationRepository; private readonly SettingsFromDb settingsFromDb; private readonly IRobotHistoryRepository robotHistoryRepository; public RobotTaskService( ILogger logger, IRobotConfigurationRepository robotConfigurationRepository, SettingsFromDb settingsFromDb, IRobotHistoryRepository robotHistoryRepository ) { this.logger = logger; this.robotConfigurationRepository = robotConfigurationRepository; this.settingsFromDb = settingsFromDb; this.robotHistoryRepository = robotHistoryRepository; } public async Task GetTaskAsync(RobotsEnum robotCode, TaskStatusEnum taskStatusCode, bool acquireTask, string? robotIp) { // 1. Ищем все задания с превышенным кол-вом попыток и просроченным временем, ставим им статус ошибки await robotConfigurationRepository.MarkExpiredTasksAsFailedAsync(settingsFromDb.RobotAttemptsNumber, settingsFromDb.RobotWaitTime); // 2. Ищем доступные задания var availableTasks = await GetAvailableTasksAsync(robotCode, taskStatusCode); if (availableTasks.Count == 0) throw new NotFoundException("Нет доступных заданий для робота"); Guid? acquiredTaskId = null; if (acquireTask) { // Берем задание в работу, устанавливаем ей статус "В работе" acquiredTaskId = await AcquireTaskAsync(availableTasks, robotIp); if (acquiredTaskId == null) throw new NotFoundException($"Не удалось взять ни одну из доступных задач ({availableTasks.Count}) в работу"); } else { // Берем первую задачу из списка доступных acquiredTaskId = availableTasks.First(); logger.LogDebug("Задача не требует захвата, взята первая из доступных: {TaskId}", acquiredTaskId); } //3. Готовим модель ответа var task = await GetTaskWithAllDataAsync(acquiredTaskId.Value, robotCode); //todo: } /// /// Получить список возможных заданий для взятия в работу. /// Кол-во заданй ограничено переменной TakeTasks /// /// /// /// private async Task> GetAvailableTasksAsync(RobotsEnum robotCode, TaskStatusEnum taskStatusCode) { var query = robotConfigurationRepository.Get() .AsNoTracking() .Where(t => t.RobotCode == (int)robotCode && t.TaskStatusCode == (int)taskStatusCode); // Если это задание для робота расписаний if (robotCode == RobotsEnum.ScheduleOrder) { // Выбираем только записи с созданными шаблонами (у которых статус 30), а только потом ищем у них расписания var createdTemplates = robotConfigurationRepository.Get() .Where(t => t.RobotCode == (int)RobotsEnum.TemplateOrder && t.TaskStatusCode == (int)TaskStatusEnum.Ok) .Select(t => t.TemplateId); query = query.Where(t => createdTemplates.Contains(t.TemplateId)); } // Сортируем по nextRun, чтобы те, у кого nextRun ближе к текущей, выполнились скорее query = query.OrderBy(t => t.Template!.NextRun).ThenBy(t => t.Template!.IsActiveSchedule).ThenBy(t => t.Template!.IsActiveTemplate); // Кандидаты заданий var tasks = new List(); // Ещем первые 10 заданий в статусе ОЖИДАНИЕ tasks = await query.Where(t => t.RobotStatusCode == (int)RobotStatusEnum.Wait).Take(TakeTasks).Select(t => t.Id).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.AttemptsNumber < settingsFromDb.RobotAttemptsNumber && t.LastRobotStatusUpdated < endDate) .Take(TakeTasks) .Select(t => t.Id) .ToListAsync(); logger.LogDebug("Найдено заданий в статусе 'В работе' {Count} шт. Робот '{Robot}'", tasks.Count, robotCode.ToString()); } return tasks; } /// /// Взять задачу в работу /// /// /// private async Task AcquireTaskAsync(List tasks, string? robotIp) { 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 }; 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) { var query = robotConfigurationRepository.Get() .AsNoTracking() .AsSingleQuery() //todo: тут общие инклуды ; if(robotCode== RobotsEnum.ScheduleOrder) { // тут инклуды только для расписаний } return await query.FirstAsync(t => t.Id == taskId); } } }