using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Logging; using PARR.AIHITMainSyncer.Settings; using PARR.Core.Common.Interfaces; using PARR.Core.Repositories.Interfaces.Unit; using PARR.Core.Services.UnitService.Interfaces; using PARR.Domain.Common.Rabbit.Messages; using PARR.Domain.Entities.Unit; namespace PARR.AIHITMainSyncer.Services { internal class SyncerService : ISyncerService { private readonly ILogger logger; private readonly ITransformService transformService; private readonly IUnitRepository unitRepository; private readonly IUnitFieldRepository unitFieldRepository; private readonly IUnitFieldValueRepository unitFieldValueRepository; private readonly FieldValueReplacementsSettings fieldValueReplacementsSettings; private readonly IUnitService unitService; public SyncerService( ILogger logger, ITransformService transformService, IUnitRepository unitRepository, IUnitFieldRepository unitFieldRepository, IUnitFieldValueRepository unitFieldValueRepository, FieldValueReplacementsSettings fieldValueReplacementsSettings, IUnitService unitService ) { this.logger = logger; this.transformService = transformService; this.unitRepository = unitRepository; this.unitFieldRepository = unitFieldRepository; this.unitFieldValueRepository = unitFieldValueRepository; this.fieldValueReplacementsSettings = fieldValueReplacementsSettings; this.unitService = unitService; } public async Task SyncAsync(string msg) { logger.LogInformation("Запуск синхронизации данных из очереди сообщений в ПАРР."); var objFromQuery = transformService.GetModelFromJson(msg); if (objFromQuery == null) { logger.LogWarning("Не удалось десериализовать сообщение в объект AihitMainDataMq. Сообщение: {Message}", msg); return; } await SyncUnitAsync(objFromQuery); } private async Task SyncUnitAsync(AihitMainDataMq objFromQuery) { logger.LogDebug("Начинаем синхронизацию Unit для объекта с именем: {UnitName}", objFromQuery.Name); //Получаем сразу словарь атрибутов, чтобы не ходить каждый раз в базу var fieldsFromDB = await unitFieldRepository.Get().AsNoTracking().ToListAsync(); if (!objFromQuery.Properties.Any()) { logger.LogWarning("Объект {UnitName} не содержит свойств для синхронизации", objFromQuery.Name); return; } // 1. Создаем копию и нормализуем свойства var normalizedProperties = new Dictionary(); foreach (var kvp in objFromQuery.Properties) normalizedProperties[kvp.Key.Trim()] = kvp.Value?.Trim(); // 2. Применяем замены значений ApplyFieldValueReplacements(normalizedProperties); var listAihitData = new List>(); // 3. Обрабатываем простые значения var propertisWithSimpleValues = normalizedProperties .Where(t => fieldsFromDB.Any(a => IsStringEqual(t.Key.Trim(), a.AihitName) && a.IsMultipleValue != true)) .Select(s => new KeyValuePair(s.Key.Trim(), s.Value?.Trim())).ToList(); listAihitData.AddRange(propertisWithSimpleValues); // 4. Обрабатываем множественные значения var propertiesWithMultipleValues = objFromQuery.Properties.Where(t => !propertisWithSimpleValues.Any(a => t.Key.Trim() == a.Key)).ToList(); propertiesWithMultipleValues.ForEach(p => listAihitData.AddRange( SplitValues(p, fieldsFromDB).Distinct().Select(s => new KeyValuePair(s.Key.Trim(), s.Value?.Trim())) )); var unit = await GetUnit(objFromQuery.Name); var incomingPairs = listAihitData .Select(kv => new { kv.Key, kv.Value }) .ToList(); // Определяем, какие ПОЛЯ обновляются этим сообщением чтобы случайно не удалить не наши var fieldsInMessage = listAihitData .Select(kv => kv.Key) .Distinct(StringComparer.OrdinalIgnoreCase) .ToHashSet(); var unitValuesToRemove = unit.UnitValues .Where(uv => uv.Field != null && fieldsInMessage.Contains(uv.Field.AihitName) && !listAihitData.Any(kv => IsStringEqual(uv.Field.AihitName, kv.Key) && IsStringEqual(uv.Value?.Value, kv.Value) ) ) .ToList(); bool hasChanges = false; if (unitValuesToRemove.Any()) { logger.LogDebug("Удаляем {Count} устаревших значений для Unit {UnitName}", unitValuesToRemove.Count, unit.Name); unitValuesToRemove.ForEach(uv => unit.UnitValues.Remove(uv)); hasChanges = true; } var unitValueMissing = listAihitData.Where(t => !unit.UnitValues.Any(a => IsStringEqual(t.Key, a.Field!.AihitName) && IsStringEqual(t.Value, a.Value!.Value) )).ToList(); if (unitValueMissing.Any()) { logger.LogDebug("Найдено {Count} отсутствующих пар значений для Unit {UnitName}", unitValueMissing.Count, unit.Name); var fieldsInChanges = unitValueMissing.Select(s => s.Key).Distinct().ToList(); var valuesInChanges = unitValueMissing.Select(s => s.Value).Distinct().ToList(); await SyncFieldsAsync(fieldsInChanges!, fieldsFromDB); await SyncValuesAsync(valuesInChanges); var normalizedFieldNames = fieldsInChanges .Select(n => n?.Trim()) .Where(n => n != null) .ToHashSet(StringComparer.OrdinalIgnoreCase); var fieldsNameWithId = await unitFieldRepository.Get().AsNoTracking() .Where(f => normalizedFieldNames.Contains(f.AihitName)) .ToListAsync(); var normalizedValues = valuesInChanges .Where(v => v != null) .Select(v => v!.Trim()) .ToHashSet(StringComparer.OrdinalIgnoreCase); var valuesWithId = await unitFieldValueRepository.Get().AsNoTracking() .Where(v => (v.Value == null && valuesInChanges.Contains(null)) || (v.Value != null && normalizedValues.Contains(v.Value))) .ToListAsync(); foreach (var change in unitValueMissing) ApplyChange(unit, fieldsNameWithId, valuesWithId, change); hasChanges = true; } else logger.LogDebug("Нет новых значений для синхронизации для Unit {UnitName}", unit.Name); // Сохраняем изменения и инвалидируем кэш только если были реальные изменения if (hasChanges) { if (!await unitRepository.CommitAsync()) { logger.LogError("Не удалось применить изменения по актуализации атрибутов Unit {UnitName} в базе данных", unit.Name); return; } // Инвалидация кэша после успешного сохранения await unitService.RemoveFromCacheAsync(unit.Id); logger.LogInformation("Успешно обновлены атрибуты для Unit {UnitName}. Кэш инвалидирован.", unit.Name); } logger.LogInformation("Завершена синхронизация Unit {UnitName}", unit.Name); return; } private List> SplitValues(KeyValuePair multiValueProperty, List fieldsFromDB) { var result = new List>(); // Обработка null или пустой строки if (string.IsNullOrEmpty(multiValueProperty.Value)) { result.Add(new KeyValuePair(multiValueProperty.Key, multiValueProperty.Value)); return result; } // Поиск конфигурации поля для тега var tagUnitField = fieldsFromDB.FirstOrDefault(f => f.Code == "tag"); bool isTagProperty = tagUnitField != null && multiValueProperty.Key == tagUnitField.AihitName; string[] splitValues; if (isTagProperty) { // Для тегов нормализуем разделители до пробела var normalizedValue = multiValueProperty.Value .Replace(';', ' ') .Replace('#', ' ') .Replace(',', ' '); splitValues = normalizedValue.Split(' ', StringSplitOptions.RemoveEmptyEntries); } else { // Для остальных полей стандартный разделитель - запятая splitValues = multiValueProperty.Value.Split(',', StringSplitOptions.RemoveEmptyEntries); } foreach (var value in splitValues) { var trimmedValue = value.Trim(); // Добавляем только непустые значения после обрезки пробелов if (!string.IsNullOrEmpty(trimmedValue)) { result.Add(new KeyValuePair(multiValueProperty.Key, trimmedValue)); } } return result; } private void ApplyChange(Unit unit, List fieldsNameWithId, List valuesWithId, KeyValuePair change) { var field = fieldsNameWithId.FirstOrDefault(t => IsStringEqual(t.AihitName, change.Key)); var value = valuesWithId.FirstOrDefault(t => IsStringEqual(t.Value, change.Value)); if (field == null || value == null) { logger.LogError("Ошибка присвоения значения({Value}) аттрибуту({Key}). FieldId - {FieldId}, ValueId - {ValueId}", change.Value, change.Key, field?.Id, value?.Id); } else { var unitInValue = new UnitInValue { UnitId = unit.Id, FieldId = field.Id, ValueId = value.Id, DateCreated = DateTimeOffset.UtcNow }; unit.UnitValues.Add(unitInValue); } } private async Task SyncValuesAsync(List values) { bool hasNull = values.Any(v => v is null); var normalizedNoneNullValues = values .Where(v => v != null) .Select(v => v!.Trim()) .Distinct(StringComparer.OrdinalIgnoreCase) .ToList(); logger.LogDebug("Синхронизируем {Count} значений (включая null: {HasNull})", normalizedNoneNullValues.Count, hasNull); var existingValuesInDb = await unitFieldValueRepository.Get() .Where(v => normalizedNoneNullValues.Contains(v.Value!) || v.Value == null) .Select(v => v.Value) .ToListAsync(); if (!hasNull) existingValuesInDb.RemoveAll(v => v is null); var newNonNullValues = normalizedNoneNullValues .Where(v => !existingValuesInDb.Any(ev => string.Equals(ev, v, StringComparison.OrdinalIgnoreCase))) .ToList(); var newValuesToInsert = new List(); newValuesToInsert.AddRange(newNonNullValues); if (hasNull && !existingValuesInDb.Contains(null)) newValuesToInsert.Add(null); // Если нет новых значений, выходим БЕЗ вызова CommitAsync if (newValuesToInsert.Count == 0) { logger.LogDebug("Нет новых значений для добавления в базу данных"); return; } var newFieldValues = newValuesToInsert.Select(v => new UnitFieldValue { Id = Guid.NewGuid(), Value = v }) .ToList(); if (!await unitFieldValueRepository.AddRangeAsync(newFieldValues)) { logger.LogError("Не удалось добавить {Count} новых значений FieldValues в контекст", newValuesToInsert.Count); return; } // Коммит выполняется ТОЛЬКО если были добавлены новые сущности if (!await unitFieldValueRepository.CommitAsync()) { logger.LogError("Не удалось сохранить {Count} новых значений FieldValues в базу данных", newValuesToInsert.Count); return; } logger.LogDebug("Успешно сохранено {Count} новых значений FieldValues", newValuesToInsert.Count); } private async Task SyncFieldsAsync(List fieldsFromAihit, List fieldsFromDB) { var newFields = fieldsFromAihit.Where(t => !fieldsFromDB.Any(f => IsStringEqual(f.AihitName, t))).ToList(); // Если нет новых полей, выходим БЕЗ вызова CommitAsync if (!newFields.Any()) { logger.LogDebug("Нет новых полей для добавления в базу данных"); return; } logger.LogDebug("Найдено {Count} новых полей для добавления", newFields.Count); foreach (var item in newFields) { var field = new UnitField { Id = Guid.NewGuid(), AihitName = item, EsppName = null, }; if (!await unitFieldRepository.CreateAsync(field)) { logger.LogError("Не удалось добавить поле в контекст: {FieldName}", item); continue; } // Коммит выполняется ТОЛЬКО после успешного добавления конкретного поля if (!await unitFieldRepository.CommitAsync()) { logger.LogError("Не удалось сохранить поле в базу данных: {FieldName}", item); } else { logger.LogDebug("Успешно сохранено поле: {FieldName}", item); } } } private bool IsStringEqual(string? value1, string? value2) { return string.Equals( value1?.Trim(), value2?.Trim(), StringComparison.OrdinalIgnoreCase); } private async Task GetUnit(string name) { var unit = await GetUnitByName(name.Trim()); if (unit == null) { logger.LogDebug("Unit с именем {UnitName} не найден, создаем новый", name); unit = await CreateUnitAsync(name); } else { logger.LogDebug("Найден существующий Unit с именем {UnitName}", name); } return unit; } private async Task CreateUnitAsync(string name) { var unit = new Unit { Name = name.Trim().ToUpperInvariant() }; if (!await unitRepository.CreateAsync(unit) || !await unitRepository.CommitAsync()) { logger.LogError("Не удалось создать Unit {UnitName}", unit.Name); throw new InvalidOperationException($"Не удалось создать Unit {unit.Name}"); } else { logger.LogInformation("Создан Unit: {UnitName}", unit.Name); } return unit; } public async Task GetUnitByName(string name) { var normalized = name.Trim().ToUpperInvariant(); return await GetUnitWithFieldsAndValues() .FirstOrDefaultAsync(u => u.Name == normalized); } private IQueryable GetUnitWithFieldsAndValues() { return unitRepository.Get() .Include(t => t.UnitValues) .ThenInclude(uv => uv.Field) .Include(t => t.UnitValues) .ThenInclude(uf => uf.Value); } private void ApplyFieldValueReplacements(Dictionary properties) { foreach (var replacement in fieldValueReplacementsSettings.Replacements) { if (properties.ContainsKey(replacement.FieldName)) { var currentValue = properties[replacement.FieldName]; // Это костыль, List не может хранить null, вместо null он хранит "" var stringNull = "null"; // Проверяем, содержится ли текущее значение в списке OldValues (с учетом регистра) if (replacement.OldValues.Contains(currentValue ?? stringNull, StringComparer.OrdinalIgnoreCase)) { logger.LogDebug("Заменяем значение поля '{FieldName}' с '{OldValue}' на '{NewValue}'", replacement.FieldName, currentValue, replacement.NewValue); properties[replacement.FieldName] = replacement.NewValue; } } } } } }