using Deal.SharedKernel.Observability; namespace Deal.Api.Observability; /// /// Фоновый сборщик gauge-метрик ядра /// public sealed class DealMetricsCollector : IHostedService { /// /// Период сбора gauge-метрик — 15 с /// public const int CollectionPeriodSeconds = 15; private readonly RuntimeDepthsCollector _depthsCollector; private readonly ILogger _logger; // Отмена при остановке хоста: прерывает текущий проход (EF-запросы наблюдают токен). private readonly CancellationTokenSource _shutdownCts = new(); private Timer? _timer; private Task? _currentIteration; private int _iterationInProgress; /// /// Создаёт фоновый сборщик gauge-метрик. /// /// Сборщик глубин очередей/сессий (общий с операторским health). /// Логгер ошибок цикла. public DealMetricsCollector(RuntimeDepthsCollector depthsCollector, ILogger logger) { ArgumentNullException.ThrowIfNull(depthsCollector); ArgumentNullException.ThrowIfNull(logger); _depthsCollector = depthsCollector; _logger = logger; } /// public Task StartAsync(CancellationToken ct) { _timer = new Timer( static state => ((DealMetricsCollector)state!).RunIteration(), this, TimeSpan.Zero, TimeSpan.FromSeconds(CollectionPeriodSeconds)); return Task.CompletedTask; } /// public async Task StopAsync(CancellationToken ct) { _timer?.Change(Timeout.InfiniteTimeSpan, Timeout.InfiniteTimeSpan); _timer?.Dispose(); _timer = null; _shutdownCts.Cancel(); Task? iteration = Volatile.Read(ref _currentIteration); if (iteration is not null) { try { await iteration.WaitAsync(ct); } catch (OperationCanceledException) { // Лимит остановки истёк — ошибки прохода проглочены внутри CollectCycleAsync. } } } /// /// Один проход сбора /// /// Задача прохода (завершается без исключений). public Task RunCycleAsync(CancellationToken ct) { if (Interlocked.CompareExchange(ref _iterationInProgress, 1, 0) != 0) { return Task.CompletedTask; } Task iteration = CollectCycleAsync(ct); Volatile.Write(ref _currentIteration, iteration); return iteration; } // Запускает проход из callback таймера (guard — внутри RunCycleAsync). private void RunIteration() { _ = RunCycleAsync(_shutdownCts.Token); } // Тело прохода: снимок глубин/сессий и публикация в meter Deal; guard сбрасывается в finally. // ct: Токен отмены (остановка хоста). private async Task CollectCycleAsync(CancellationToken ct) { try { RuntimeDepthsDto depths = await _depthsCollector.CollectAsync(ct); DealMetrics.SetPipelineQueueDepth(depths.PipelineQueue); DealMetrics.SetMlOutboxDepth(depths.MlOutbox); DealMetrics.SetActiveSessions(depths.ActiveSessions); DealMetrics.ReplaceBudgetRatios( (depths.Budgets ?? []).Select(budget => (budget.TenantId, budget.UsedRatio))); } catch (OperationCanceledException) { // Остановка хоста: проход прерван по токену — штатный выход, не ошибка. } catch (Exception exception) { _logger.LogError(exception, "Сборщик метрик: проход по тенантам не удался"); } finally { Interlocked.Exchange(ref _iterationInProgress, 0); } } }