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

350 lines
16 KiB
C#
Raw Permalink 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 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;
namespace PARR.TemplateMatcher.Services.Implementations;
internal class SimpleTemplateSynchronizer : ITemplateSynchronizer
{
#if DEBUG
private readonly Guid _targetUnitId = Guid.Parse("358437ac-1eeb-4c00-840c-998326f657ac");
#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,
ITemplateNameNormalizer templateNameNormalizer,
ITemplateMqPublisher templateMqPublisher,
IMatchingStatusService matchingStatusService,
SettingsFromDb settingsFromDb,
IOptions<TemplateSettings> templateSettings,
IUnusedTemplatesSyncService unusedTemplatesSyncService
)
{
_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;
}
}
private void PrintDryRunReport(SimpleSyncContext context)
{
_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);
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();
#if DEBUG
// Отладка: проверить, есть ли юнит в unitIds
if (unitIds.Contains(_targetUnitId))
{
_logger.LogDebug("Юнит {TargetUnitId} найден в unitIds.", _targetUnitId);
}
else
{
_logger.LogDebug("Юнит {TargetUnitId} НЕ найден в unitIds.", _targetUnitId);
}
#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))
{
_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);
}
}
}
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)
);
}
}