Files
parr_api/PARR.AIHITMainSyncer/Services/SyncerService.cs

445 lines
20 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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<SyncerService> 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<SyncerService> 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<AihitMainDataMq>(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<string, string?>();
foreach (var kvp in objFromQuery.Properties)
normalizedProperties[kvp.Key.Trim()] = kvp.Value?.Trim();
// 2. Применяем замены значений
ApplyFieldValueReplacements(normalizedProperties);
var listAihitData = new List<KeyValuePair<string, string?>>();
// 3. Обрабатываем простые значения
var propertisWithSimpleValues = normalizedProperties
.Where(t => fieldsFromDB.Any(a => IsStringEqual(t.Key.Trim(), a.AihitName) && a.IsMultipleValue != true))
.Select(s => new KeyValuePair<string, string?>(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<string, string?>(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<KeyValuePair<string, string?>> SplitValues(KeyValuePair<string, string?> multiValueProperty, List<UnitField> fieldsFromDB)
{
var result = new List<KeyValuePair<string, string?>>();
// Обработка null или пустой строки
if (string.IsNullOrEmpty(multiValueProperty.Value))
{
result.Add(new KeyValuePair<string, string?>(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<string, string?>(multiValueProperty.Key, trimmedValue));
}
}
return result;
}
private void ApplyChange(Unit unit, List<UnitField> fieldsNameWithId, List<UnitFieldValue> valuesWithId, KeyValuePair<string, string?> 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<string?> 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<string?>();
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<string> fieldsFromAihit, List<UnitField> 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<Unit> 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<Unit> 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<Unit?> GetUnitByName(string name)
{
var normalized = name.Trim().ToUpperInvariant();
return await GetUnitWithFieldsAndValues()
.FirstOrDefaultAsync(u => u.Name == normalized);
}
private IQueryable<Unit> GetUnitWithFieldsAndValues()
{
return unitRepository.Get()
.Include(t => t.UnitValues)
.ThenInclude(uv => uv.Field)
.Include(t => t.UnitValues)
.ThenInclude(uf => uf.Value);
}
private void ApplyFieldValueReplacements(Dictionary<string, string?> properties)
{
foreach (var replacement in fieldValueReplacementsSettings.Replacements)
{
if (properties.ContainsKey(replacement.FieldName))
{
var currentValue = properties[replacement.FieldName];
// Это костыль, List<string?> не может хранить 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;
}
}
}
}
}
}