using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using PARR.AIHITMainLoader.Models; using PARR.AIHITMainLoader.Services; using PARR.AIHITMainLoader.Settings; using PARR.Core.Common.Interfaces; using PARR.Core.Common.Interfaces.RabbitServices; using System.Text.Encodings.Web; 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 IRabbitService mqService; private readonly IServiceProvider serviceProvider; private static readonly JsonSerializerOptions jsonOptions = new JsonSerializerOptions { Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping, }; public AihitMainLoader( ILogger logger, IIntervalService intervalService, WorkerSettings workerSettings, LoaderSettings loaderSettings, MqSettings mqSettings, IRabbitService 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 currentBatch = new List(); int packageSize = loaderSettings.PackageSize; if (packageSize <= 0) { logger.LogError("Некорректный размер пакета: {PackageSize}. Ожидалось > 0.", packageSize); return; } foreach (var service in services) { foreach (var responseArea in loaderSettings.ResponseAreas) { logger.LogDebug($"responseArea = {responseArea}"); var EKs = service.Invoke(responseArea); if (EKs?.Any() != true) continue; foreach (var item in EKs) { if (item == null) continue; var mainData = item.ToMainData(); if (mainData == null) continue; var json = JsonSerializer.Serialize(mainData, jsonOptions); currentBatch.Add(json); // Отправка, если набрали полный пакет if (currentBatch.Count >= packageSize) { await SendBatchAsync(currentBatch); currentBatch.Clear(); } } } } // Отправка остатка if (currentBatch.Count > 0) { await SendBatchAsync(currentBatch); } } } catch (Exception ex) { logger.LogError(ex, "Ошибка при выполнении загрузки данных из АИХ ИТ."); } } private async Task SendBatchAsync(List batch) { var sendResult = await mqService.SendAsync(mqSettings, batch.ToArray()); 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($"Ошибка при передаче данных в RabbitMQ. Не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}"); } } private List?>> GetServices(IAihitService service) { return new() { service.GetCvkData, service.GetPtkData, service.GetRegionalEKData, service.GetStoData, service.GetOrgData }; } } }