445 lines
20 KiB
C#
445 lines
20 KiB
C#
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;
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
} |