Files
parr_api/PARR.TemplateMatcher/Services/Implementations/SimpleTemplateSynchronizer.cs

350 lines
16 KiB
C#
Raw Permalink Normal View History

using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using PARR.Core.Repositories.Interfaces;
using PARR.Core.Repositories.Interfaces.JobRepositories;
using PARR.Core.Services.MatchingStatusService;
using PARR.Core.Services.UnitFilterService;
using PARR.Domain.Cache.Models;
using PARR.Domain.Common.Rabbit.Messages.TemplateMatching;
using PARR.Domain.Entities.Base.History;
using PARR.Domain.Enums;
using PARR.Domain.Settings;
using PARR.TemplateMatcher.Exceptions;
using PARR.TemplateMatcher.Services.Interfaces;
using PARR.TemplateMatcher.Services.SimpleSync;
using PARR.TemplateMatcher.Settings;
using System.Diagnostics;
2025-12-26 21:47:53 +10:00
namespace PARR.TemplateMatcher.Services.Implementations;
internal class SimpleTemplateSynchronizer : ITemplateSynchronizer
{
2025-12-26 21:47:53 +10:00
#if DEBUG
private readonly Guid _targetUnitId = Guid.Parse("358437ac-1eeb-4c00-840c-998326f657ac");
2025-12-26 21:47:53 +10:00
#endif
private readonly IEnumerable<ISimpleSyncStage> _readStages;
private readonly IEnumerable<ISimpleSyncWriteStage> _writeStages;
private readonly ILogger<SimpleTemplateSynchronizer> _logger;
private readonly IUnitFilterService _unitFilterService;
private readonly ITemplateRepository _templateService;
private readonly IJobRepository _jobService;
private readonly ITemplateNameNormalizer _templateNameNormalizer;
private readonly ITemplateMqPublisher _templateMqPublisher;
private readonly IMatchingStatusService _matchingStatusService;
private readonly SettingsFromDb _settingsFromDb;
private readonly IUnusedTemplatesSyncService _unusedTemplatesSyncService;
public SimpleTemplateSynchronizer(
IEnumerable<ISimpleSyncStage> readStages,
IEnumerable<ISimpleSyncWriteStage> writeStages,
ILogger<SimpleTemplateSynchronizer> logger,
IUnitFilterService unitFilterService,
ITemplateRepository templateService,
IJobRepository jobService,
2025-12-26 21:47:53 +10:00
ITemplateNameNormalizer templateNameNormalizer,
ITemplateMqPublisher templateMqPublisher,
IMatchingStatusService matchingStatusService,
SettingsFromDb settingsFromDb,
IOptions<TemplateSettings> templateSettings,
IUnusedTemplatesSyncService unusedTemplatesSyncService
2025-12-26 21:47:53 +10:00
)
{
_readStages = readStages;
_writeStages = writeStages;
_logger = logger;
_unitFilterService = unitFilterService;
_templateService = templateService;
_jobService = jobService;
_templateNameNormalizer = templateNameNormalizer;
_templateMqPublisher = templateMqPublisher;
_matchingStatusService = matchingStatusService;
_settingsFromDb = settingsFromDb;
_unusedTemplatesSyncService = unusedTemplatesSyncService;
}
public async Task SyncTemplatesForJobAsync(Guid jobId, HistoryInitiator initiator, bool dryRun = false, CancellationToken ct = default)
{
if (jobId == _settingsFromDb.JobIdForUnusedTemplates)
{
await _unusedTemplatesSyncService.SyncAsync(jobId, initiator, dryRun);
return;
}
_logger.LogInformation("Начало синхронизации шаблонов для Job {JobId}", jobId);
var existingStatus = await _matchingStatusService.GetStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
if (existingStatus.DetailsJobs?.Count > 0)
{
_logger.LogWarning("Синхронизация для Job {JobId} уже запущена. Пропускаем.", jobId);
return;
}
var initialStatus = new MatchingStatusItemDto
{
DateStart = DateTimeOffset.UtcNow,
Action = TemplateMatcherActionEnum.Sync,
Comment = "Начало синхронизации"
};
await _matchingStatusService.SetMatchingStatusAsync(
jobId, SyncTaskEntityTypeEnum.Job,
new MatchingStatusItem { Data = initialStatus, Timestamp = DateTimeOffset.UtcNow, Source = nameof(SimpleTemplateSynchronizer) },
TimeSpan.FromMinutes(35));
// Таймер запускается ПОСЛЕ инфраструктурных операций (статус, проверка блокировки)
var totalSw = Stopwatch.StartNew();
var context = new SimpleSyncContext { JobId = jobId, Initiator = initiator };
try
{
foreach (var stage in _readStages)
{
var stageSw = Stopwatch.StartNew();
await stage.ExecuteAsync(context);
stageSw.Stop();
_logger.LogDebug("[Perf] Job '{JobName}' ({JobId}) | Этап: {Stage} | Время: {Ms} мс",
context.JobName, jobId, stage.StageName, stageSw.ElapsedMilliseconds);
}
// Подробный отчёт по результатам read-этапов
PrintDryRunReport(context);
if (!dryRun)
{
foreach (var stage in _writeStages)
{
var stageSw = Stopwatch.StartNew();
await stage.ExecuteAsync(context);
stageSw.Stop();
_logger.LogDebug("[Perf] Job '{JobName}' ({JobId}) | Этап: {Stage} | Время: {Ms} мс",
context.JobName, jobId, stage.StageName, stageSw.ElapsedMilliseconds);
}
await SetStatusAsync(jobId, "Синхронизация завершена успешно");
await _matchingStatusService.DeleteMatchingStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
_logger.LogInformation("Синхронизация шаблонов завершена для Job '{JobName}' ({JobId})", context.JobName, jobId);
}
else
{
_logger.LogInformation("[DryRun] Write-этапы пропущены. Изменения в БД и MQ не выполнены.");
await SetStatusAsync(jobId, "DryRun: Анализ завершён без изменений");
await _matchingStatusService.DeleteMatchingStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
}
totalSw.Stop();
_logger.LogInformation("[Perf] Job '{JobName}' ({JobId}) | ИТОГО: {TotalMs} мс",
context.JobName, jobId, totalSw.ElapsedMilliseconds);
}
catch (SyncEarlyExitException ex)
{
totalSw.Stop();
_logger.LogInformation("Job {JobId}: {Reason} ({ElapsedMs} мс)",
jobId, ex.Reason, totalSw.ElapsedMilliseconds);
await SetStatusAsync(jobId, ex.Reason);
await _matchingStatusService.DeleteMatchingStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
}
catch (Exception ex)
{
totalSw.Stop();
_logger.LogError(ex, "Ошибка при синхронизации Job '{JobName}' ({JobId}) через {ElapsedMs} мс",
context.JobName ?? string.Empty, jobId, totalSw.ElapsedMilliseconds);
await SetStatusAsync(jobId, $"Ошибка: {ex.Message}");
throw;
}
2025-12-26 21:47:53 +10:00
}
private void PrintDryRunReport(SimpleSyncContext context)
2025-12-26 21:47:53 +10:00
{
_logger.LogInformation("========== [DryRun] ОТЧЁТ по Job '{JobName}' ({JobId}) ==========",
context.JobName, context.JobId);
_logger.LogInformation("[DryRun] Отфильтровано юнитов: {Count}", context.FilteredUnitIds.Count);
_logger.LogInformation("[DryRun] Существующих активных шаблонов: {Count}", context.ExistingUsedTemplates.Count);
_logger.LogInformation("[DryRun] Новых юнитов (требуют аллокации): {Count}", context.NewUnitIds.Count);
_logger.LogInformation("[DryRun] Шаблонов для деактивации (Unused): {Count}", context.UnusedTemplates.Count);
_logger.LogInformation("[DryRun] Шаблонов для переименования: {Count}", context.TemplatesToRename.Count);
if (context.NewUnitIds.Any())
{
_logger.LogDebug("[DryRun] --- Новые юниты (будут созданы шаблоны) ---");
foreach (var unitId in context.NewUnitIds)
{
_logger.LogDebug("[DryRun] + {Unit}", context.FormatUnit(unitId));
}
}
if (context.UnusedTemplates.Any())
{
_logger.LogDebug("[DryRun] --- Шаблоны для деактивации ---");
foreach (var template in context.UnusedTemplates)
{
_logger.LogDebug("[DryRun] - {TemplateId} (Unit: {Unit}, Name: '{Name}')",
template.Id, context.FormatUnit(template.UnitId), template.Name);
}
}
if (context.TemplatesToRename.Any())
{
_logger.LogDebug("[DryRun] --- Шаблоны для переименования ---");
foreach (var (template, expectedName) in context.TemplatesToRename)
{
_logger.LogDebug("[DryRun] ~ {TemplateId} (Unit: {Unit}): '{OldName}' -> '{NewName}'",
template.Id, context.FormatUnit(template.UnitId), template.Name, expectedName);
}
}
_logger.LogInformation("========== [DryRun] КОНЕЦ ОТЧЁТА ==========");
}
public Task SyncTemplatesForJobGroupAsync(Guid jobGroupId, HistoryInitiator initiator, bool dryRun = false, CancellationToken ct = default)
{
_logger.LogWarning(
"SimpleTemplateSynchronizer: SyncTemplatesForJobGroup вызван для JobGroup {JobGroupId} (DryRun={DryRun}). Это не поддерживаемая операция.",
jobGroupId, dryRun);
return Task.CompletedTask;
}
public async Task UpdateTemplatesForJobAsync(Guid jobId, HistoryInitiator initiator, CancellationToken ct = default)
{
_logger.LogDebug("Обновление шаблонов для Job {JobId}", jobId);
// === Проверка: уже запущена? ===
var existingStatus = await _matchingStatusService.GetStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
if (existingStatus.DetailsJobs?.Count > 0)
{
_logger.LogWarning("Обновление для Job {JobId} уже запущено. Пропускаем.", jobId);
2025-12-26 21:47:53 +10:00
return;
}
var initialStatus = new MatchingStatusItemDto
{
DateStart = DateTimeOffset.UtcNow,
Action = TemplateMatcherActionEnum.Update,
Comment = "Начало обновления имён"
};
await _matchingStatusService.SetMatchingStatusAsync(
jobId,
SyncTaskEntityTypeEnum.Job,
new MatchingStatusItem { Data = initialStatus, Timestamp = DateTimeOffset.UtcNow, Source = nameof(SimpleTemplateSynchronizer) },
TimeSpan.FromMinutes(30)
);
try
{
var job = await _jobService.Get()
.AsNoTracking()
.Include(j => j.AutoControl)
.Include(j => j.Tnk)
.Include(j => j.Group)
.ThenInclude(g => g!.GroupType)
.Include(j => j.UnitFilters)
.ThenInclude(uf => uf.RelationshipFilters)
.FirstOrDefaultAsync(j => j.Id == jobId);
if (job == null)
{
_logger.LogWarning("Job {JobId} не найден.", jobId);
await SetStatusAsync(jobId, "Job не найден");
return;
}
// === Получение отфильтрованных юнитов с полной информацией ===
var filteredUnits = await _unitFilterService.GetUnitsByJobFilterAsync(jobId);
if (filteredUnits == null || !filteredUnits.Any())
{
_logger.LogInformation("Для Job {JobId} фильтры не дали Unit'ов.", jobId);
await SetStatusAsync(jobId, "Нет Unit'ов — обновление не требуется");
await _matchingStatusService.DeleteMatchingStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
return;
}
// Извлекаем ID юнитов для последующих операций
var unitIds = filteredUnits.Select(u => u.Id).ToList();
2025-12-26 21:47:53 +10:00
#if DEBUG
// Отладка: проверить, есть ли юнит в unitIds
if (unitIds.Contains(_targetUnitId))
{
_logger.LogDebug("Юнит {TargetUnitId} найден в unitIds.", _targetUnitId);
}
else
{
_logger.LogDebug("Юнит {TargetUnitId} НЕ найден в unitIds.", _targetUnitId);
}
2025-12-26 21:47:53 +10:00
#endif
var existingTemplates = await _templateService.Get()
.AsNoTracking()
.Include(t => t.Unit)
.Include(t => t.UnitsInTemplate)
.Include(t => t.Job)
.ThenInclude(t => t!.Group)
.ThenInclude(t => t!.GroupType)
.Include(t => t.Job)
.ThenInclude(t => t!.Tnk)
.Where(t => t.JobId == jobId && t.StatusTypeId == TemplateStatusTypeEnum.Used)
.ToListAsync();
foreach (var template in existingTemplates)
{
if (unitIds.Contains(template.UnitId))
{
var expectedName = await _templateNameNormalizer.GetNormalizedTemplateNameAsync(template);
if (!string.Equals(template.Name, expectedName, StringComparison.OrdinalIgnoreCase))
2025-12-26 21:47:53 +10:00
{
_logger.LogDebug("Шаблон {TemplateId} требует обновления имени: старое = '{OldName}', новое = '{NewName}'", template.Id, template.Name, expectedName);
var updateRequest = new TemplateUpdaterMessage
{
TemplateId = template.Id,
JobId = jobId,
UnitId = template.UnitId,
Name = expectedName,
IsActiveTemplate = template.IsActiveTemplate,
IsActiveSchedule = template.IsActiveSchedule,
IsNew = false,
Index = template.Index,
StatusTypeId = TemplateStatusTypeEnum.Used,
Initiator = initiator,
UnitsInTemplate = new List<UnitInTemplateMessage>() // для простого шаблона
};
await _templateMqPublisher.PublishUpdateAsync(updateRequest);
}
2025-12-26 21:47:53 +10:00
}
}
await SetStatusAsync(jobId, "Обновление завершено");
await _matchingStatusService.DeleteMatchingStatusAsync(jobId, SyncTaskEntityTypeEnum.Job);
_logger.LogInformation("Обновление шаблонов завершено для Job {JobId}.", jobId);
}
catch (Exception ex)
{
_logger.LogError(ex, "Ошибка при обновлении Job {JobId}", jobId);
await SetStatusAsync(jobId, $"Ошибка: {ex.Message}");
throw;
}
}
private async Task SetStatusAsync(Guid jobId, string comment)
{
var status = new MatchingStatusItemDto
{
DateStart = DateTimeOffset.UtcNow,
Action = TemplateMatcherActionEnum.Sync,
Comment = comment
};
await _matchingStatusService.SetMatchingStatusAsync(
jobId,
SyncTaskEntityTypeEnum.Job,
new MatchingStatusItem { Data = status, Timestamp = DateTimeOffset.UtcNow, Source = nameof(SimpleTemplateSynchronizer) },
TimeSpan.FromMinutes(30)
);
}
}