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 _readStages; private readonly IEnumerable _writeStages; private readonly ILogger _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 readStages, IEnumerable writeStages, ILogger logger, IUnitFilterService unitFilterService, ITemplateRepository templateService, IJobRepository jobService, ITemplateNameNormalizer templateNameNormalizer, ITemplateMqPublisher templateMqPublisher, IMatchingStatusService matchingStatusService, SettingsFromDb settingsFromDb, IOptions 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() // для простого шаблона }; 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) ); } }