using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using PARR.AIHITMainLoader.Models; using PARR.AIHITMainLoader.Services; using PARR.AIHITMainLoader.Settings; using PARR.BLL.Domain.Mq; using PARR.BLL.Services.Interfaces; using System.Text.Json; namespace PARR.AIHITMainLoader { internal class AihitMainLoader : IAihitMainLoader { private readonly ILogger logger; private readonly IIntervalService intervalService; private readonly WorkerSettings workerSettings; private readonly LoaderSettings loaderSettings; private readonly MqSettings mqSettings; private readonly IMqService mqService; private readonly IServiceProvider serviceProvider; public AihitMainLoader( ILogger logger, IIntervalService intervalService, WorkerSettings workerSettings, LoaderSettings loaderSettings, MqSettings mqSettings, IMqService mqService, IServiceProvider serviceProvider ) { this.logger = logger; this.intervalService = intervalService; this.workerSettings = workerSettings; this.loaderSettings = loaderSettings; this.mqSettings = mqSettings; this.mqService = mqService; this.serviceProvider = serviceProvider; } public async Task StartAsync() { logger.LogInformation("Запуск сервиса загрузки данных из АИХ ИТ."); await intervalService.IntervalInitAsync(LoadDataAsync, workerSettings.RepeatEvery); } private async Task LoadDataAsync() { try { logger.LogInformation("Запуск загрузки данных из АИХ ИТ."); using (var scope = serviceProvider.CreateScope()) { var aihitService = scope.ServiceProvider.GetRequiredService(); var services = GetServices(aihitService); var listToMq = new List(); foreach (var service in services) { foreach (var responseArea in loaderSettings.ResponseAreas) { logger.LogDebug($"responseArea = {responseArea}"); var EKs = service.Invoke(responseArea); if (EKs?.Any() == true) { foreach (var item in EKs) { if (item != null) { var mainData = item.ToMainData(); if (mainData != null) listToMq.Add(mainData); } } } } } var preparedDataList = listToMq.Select(d => JsonSerializer.Serialize(d)).ToList(); int count = preparedDataList.Count; int packageSize = loaderSettings.PackageSize; if (packageSize <= 0) { logger.LogError("Некорректный размер пакета: {PackageSize}. Ожидалось > 0.", packageSize); return; } var parts = count == 0 ? 0 : (count + packageSize - 1) / packageSize; for (var i = 0; i < parts; i++) { int start = i * packageSize; int length = Math.Min(packageSize, count - start); var batch = preparedDataList.GetRange(start, length).ToArray(); var sendResult = await mqService.SendAsync(mqSettings, batch); if (sendResult.IsSuccess) { logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count()}"); logger.LogDebug($"Отправлено сообщений: {batch.Count()}, пример первого: {batch.FirstOrDefault()?.Substring(0, Math.Min(200, batch.FirstOrDefault()?.Length ?? 0))}..."); } else logger.LogError($"Ошибка при передаче данных в Rabbit, не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}"); } } } catch (Exception ex) { logger.LogError(ex, "Ошибка при выполнении загрузки данных из АИХ ИТ."); } } private List?>> GetServices(IAihitService service) { return new() { service.GetCvkData, service.GetPtkData, service.GetRegionalEKData, service.GetStoData, service.GetOrgData }; } } }