diff --git a/backlog.md b/backlog.md index 3d8f0fc..7881bee 100644 --- a/backlog.md +++ b/backlog.md @@ -42,7 +42,7 @@ | ID | Пункт | Источник | Приоритет | Статус | |---|---|---|---|---| -| BL-ALERT-BUDGET | Метрика расхода токенов для алерта бюджета (сейчас алерт не настроен — нет метрики) | этап 12, A | P2 | TECHDEBT | +| BL-ALERT-BUDGET | **Сделано (2026-09-11):** метрика `deal.ai.budget.used.ratio{tenant}` (доля израсходованного ИИ-бюджета периода, 0..1) в `DealMetrics` + сбор в `RuntimeDepthsCollector`/`DealMetricsCollector`; на её основе оператор настраивает алерт в Prometheus/Grafana | этап 12, A | P2 | DONE | | BL-LOG-ACTOR | Идентичность актора в access-логах (сейчас login/tenant только в `audit_log`) | этап 12, T6 | P3 | TECHDEBT | | BL-GRACEFUL | Дополнительные проверки устойчивости/ретраев (по результатам нагрузочного прогона) | этап 12, C | P2 | BACKLOG | | BL-SUSPICIOUS | Расширение детектора подозрительной активности (правила/пороги по логам безопасности) | ТЗ §10.5, этап 12 | P3 | BACKLOG | diff --git a/docs/superpowers/STATUS.md b/docs/superpowers/STATUS.md index 4b846eb..5325f9c 100644 --- a/docs/superpowers/STATUS.md +++ b/docs/superpowers/STATUS.md @@ -12,8 +12,9 @@ > `SourceIngressGrpcService`), `PushMessage` из telegram.proto удалён. Сухой прогон текста по конвейеру > (стоп-правила → ML → ИИ) без записи: `POST /api/admin/check-message` + UI настроек. Remote-просмотр > исходника: `TelegramService.ReadSource` + `TelegramSourceContentProvider` + UI «Обновить из источника». -> Ядро: build 5 sln 0/0, `Deal.Tests.Unit` **1307/1307 PASS**, telegram **130/130**, фронт `build` + -> `lint:i18n` зелёные. Детали — `docs/superpowers/specs/2026-09-11-source-contract-design.md`. +> Метрика алертинга `deal.ai.budget.used.ratio{tenant}`. Ядро: build 5 sln 0/0, +> `Deal.Tests.Unit` **1310/1310 PASS**, telegram **130/130**, фронт `build` + `lint:i18n` зелёные. +> Детали — `docs/superpowers/specs/2026-09-11-source-contract-design.md`. > Осталось (в backlog): `GET /api/cards/{id}/source` + `ISourceContentProvider`, выгрузка вложений > telegram-адаптером в Storage, `TelegramSourceContentProvider`, перенос оставшейся Telegram-специфики > (`TelegramStore`, `Dialogs`/`TgMessages`, Discovery) в telegram-сервис. diff --git a/docs/technical/Техническая-документация-Дейл.md b/docs/technical/Техническая-документация-Дейл.md index fb3aa50..ed43734 100644 --- a/docs/technical/Техническая-документация-Дейл.md +++ b/docs/technical/Техническая-документация-Дейл.md @@ -272,7 +272,9 @@ settings(Key varchar(200) PK, ValueJson text, UpdatedAt timestamptz) -- - `deal_ml_calls_total` / `deal_ml_tokens_total` — вызовы и оценка токенов локального ML; - `deal_audit_events_total{event,actor}` — события аудита по типу/актору; - `deal_pipeline_queue_depth`, `deal_ml_outbox_depth` — суммарные глубины очередей (пайплайн, MlOutbox) - по всем тенантам; `deal_sessions_active` — активные непросроченные сессии пользователей и операторов. + по всем тенантам; `deal_sessions_active` — активные непросроченные сессии пользователей и операторов; + - `deal_ai_budget_used_ratio{tenant}` — доля израсходованного ИИ-бюджета периода (0..1) по тенантам + (осознанное исключение из низкокардинального правила: бюджеты пер-тенантные, алерт должен знать тенанта). Gauge-значения собирает фоновый `DealMetricsCollector` ядра (каждые 15 с) через существующие сервисы/хранилища (`PipelineProcessingService.QueueCountsAsync`, `IMlLearningStore.CountOutboxAsync`, `public.sessions`/`operator_sessions`); инкремент счётчиков токенов/аудита — там же, где пишутся diff --git a/src/core/Deal.Api/Observability/DealMetricsCollector.cs b/src/core/Deal.Api/Observability/DealMetricsCollector.cs index 3ac873a..15fcf8a 100644 --- a/src/core/Deal.Api/Observability/DealMetricsCollector.cs +++ b/src/core/Deal.Api/Observability/DealMetricsCollector.cs @@ -100,6 +100,8 @@ public sealed class DealMetricsCollector : IHostedService DealMetrics.SetPipelineQueueDepth(depths.PipelineQueue); DealMetrics.SetMlOutboxDepth(depths.MlOutbox); DealMetrics.SetActiveSessions(depths.ActiveSessions); + DealMetrics.ReplaceBudgetRatios( + (depths.Budgets ?? []).Select(budget => (budget.TenantId, budget.UsedRatio))); } catch (OperationCanceledException) { diff --git a/src/core/Deal.Api/Observability/RuntimeDepthsCollector.cs b/src/core/Deal.Api/Observability/RuntimeDepthsCollector.cs index 06d9e8f..efc4184 100644 --- a/src/core/Deal.Api/Observability/RuntimeDepthsCollector.cs +++ b/src/core/Deal.Api/Observability/RuntimeDepthsCollector.cs @@ -39,8 +39,10 @@ public sealed class RuntimeDepthsCollector { await using AsyncServiceScope cycleScope = _scopeFactory.CreateAsyncScope(); int sessions = await CountActiveSessionsAsync(cycleScope, ct); - (long queue, long outbox) = await SumTenantDepthsAsync(cycleScope, ct); - return new RuntimeDepthsDto(queue, outbox, sessions); + IReadOnlyList tenants = await ListTenantsAsync(cycleScope, ct); + (long queue, long outbox) = await SumTenantDepthsAsync(tenants, ct); + IReadOnlyList budgets = await CollectBudgetsAsync(cycleScope, tenants, ct); + return new RuntimeDepthsDto(queue, outbox, sessions, budgets); } // Считает активные непросроченные сессии пользователей и операторов (public-схема). @@ -68,24 +70,16 @@ public sealed class RuntimeDepthsCollector } } - // Суммирует глубины очередей по всем тенантам реестра. - // cycleScope: Scope прохода (реестр тенантов читается без tenant-контекста). + // Список тенантов реестра (без tenant-контекста — public-схема). + // cycleScope: Scope прохода. // ct: Токен отмены. - // Возвращает: Пара (сумма очереди пайплайна, сумма MlOutbox); недоступность реестра — (0, 0). - private async Task<(long Queue, long Outbox)> SumTenantDepthsAsync(AsyncServiceScope cycleScope, CancellationToken ct) + // Возвращает: Тенанты реестра; недоступность реестра — пустой список. + private async Task> ListTenantsAsync(AsyncServiceScope cycleScope, CancellationToken ct) { - long queueDepth = 0; - long outboxDepth = 0; try { ITenantRepository tenantRepository = cycleScope.ServiceProvider.GetRequiredService(); - IReadOnlyList tenants = await tenantRepository.ListAsync(ct); - foreach (TenantRecordDto tenant in tenants) - { - (int queue, int outbox) = await CollectTenantAsync(tenant, ct); - queueDepth += queue; - outboxDepth += outbox; - } + return await tenantRepository.ListAsync(ct); } catch (OperationCanceledException) { @@ -94,11 +88,76 @@ public sealed class RuntimeDepthsCollector catch (Exception exception) { _logger.LogWarning(exception, "Сборщик глубин: обход реестра тенантов не удался"); + return []; + } + } + + // Суммирует глубины очередей переданных тенантов. + // tenants: Тенанты реестра. + // ct: Токен отмены. + // Возвращает: Пара (сумма очереди пайплайна, сумма MlOutbox). + private async Task<(long Queue, long Outbox)> SumTenantDepthsAsync( + IReadOnlyList tenants, + CancellationToken ct) + { + long queueDepth = 0; + long outboxDepth = 0; + foreach (TenantRecordDto tenant in tenants) + { + (int queue, int outbox) = await CollectTenantAsync(tenant, ct); + queueDepth += queue; + outboxDepth += outbox; } return (queueDepth, outboxDepth); } + // Собирает долю израсходованного ИИ-бюджета по тенантам. + // cycleScope: Scope прохода (public-схема: tenant_limits — без tenant-контекста). + // tenants: Тенанты реестра. + // ct: Токен отмены. + // Возвращает: Доли 0..1 по тенантам; сбой одного тенанта — он пропускается. + private async Task> CollectBudgetsAsync( + AsyncServiceScope cycleScope, + IReadOnlyList tenants, + CancellationToken ct) + { + var budgets = new List(tenants.Count); + try + { + ITenantLimitStore limits = cycleScope.ServiceProvider.GetRequiredService(); + foreach (TenantRecordDto tenant in tenants) + { + try + { + BudgetStateDto state = await limits.GetStateAsync(tenant.Id, ct); + double ratio = state.BudgetTokens > 0 + ? Math.Clamp(state.UsedTokens / (double)state.BudgetTokens, 0.0, 1.0) + : 0.0; + budgets.Add(new TenantBudgetUsageDto(tenant.Id.ToString("N"), ratio)); + } + catch (OperationCanceledException) + { + throw; + } + catch (Exception exception) + { + _logger.LogWarning(exception, "Сборщик глубин: бюджет тенанта {TenantId} не прочитан", tenant.Id); + } + } + } + catch (OperationCanceledException) + { + throw; + } + catch (Exception exception) + { + _logger.LogWarning(exception, "Сборщик глубин: секция бюджета не удалась"); + } + + return budgets; + } + // Считает глубины очередей одного тенанта в собственном scope (SetTenant → сервисы → Reset). // tenant: Тенант реестра (Id в формате Guid; схема — tenant_<N>). // ct: Токен отмены прохода. diff --git a/src/core/Deal.Api/Observability/RuntimeDepthsDto.cs b/src/core/Deal.Api/Observability/RuntimeDepthsDto.cs index 2c31660..1df166f 100644 --- a/src/core/Deal.Api/Observability/RuntimeDepthsDto.cs +++ b/src/core/Deal.Api/Observability/RuntimeDepthsDto.cs @@ -6,4 +6,9 @@ namespace Deal.Api.Observability; /// Суммарная глубина очереди пайплайна (new+filtered) по всем тенантам. /// Суммарная глубина очереди обучения ML (MlOutbox) по всем тенантам. /// Число активных непросроченных сессий пользователей и операторов. -public sealed record RuntimeDepthsDto(long PipelineQueue, long MlOutbox, int ActiveSessions); +/// Использование ИИ-бюджета по тенантам (пусто — данных нет). +public sealed record RuntimeDepthsDto( + long PipelineQueue, + long MlOutbox, + int ActiveSessions, + IReadOnlyList? Budgets = null); diff --git a/src/core/Deal.Api/Observability/TenantBudgetUsageDto.cs b/src/core/Deal.Api/Observability/TenantBudgetUsageDto.cs new file mode 100644 index 0000000..9aef560 --- /dev/null +++ b/src/core/Deal.Api/Observability/TenantBudgetUsageDto.cs @@ -0,0 +1,8 @@ +namespace Deal.Api.Observability; + +/// +/// Использование ИИ-бюджета тенанта — для метрики алертинга. +/// +/// Id тенанта (Guid в формате «N»). +/// Доля израсходованного бюджета периода 0..1; бюджет не задан — 0. +public sealed record TenantBudgetUsageDto(string TenantId, double UsedRatio); diff --git a/src/core/Deal.SharedKernel/Observability/DealMetrics.cs b/src/core/Deal.SharedKernel/Observability/DealMetrics.cs index 0d480eb..ef221ee 100644 --- a/src/core/Deal.SharedKernel/Observability/DealMetrics.cs +++ b/src/core/Deal.SharedKernel/Observability/DealMetrics.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using System.Diagnostics; using System.Diagnostics.Metrics; @@ -38,6 +39,16 @@ public static class DealMetrics /// public const string AuditActorTagName = "actor"; + /// + /// Имя метки тенанта + /// + public const string TenantTagName = "tenant"; + + /// + /// Имя метрики доли израсходованного ИИ-бюджета + /// + public const string AiBudgetUsedRatioName = "deal.ai.budget.used.ratio"; + // Meter прикладных метрик (static — один на процесс, как того требует System.Diagnostics.Metrics). private static readonly Meter Meter = new(MeterName); @@ -98,6 +109,17 @@ public static class DealMetrics (Func)(() => (long)Interlocked.Read(ref DealMetrics._activeSessions)), description: "Активные непросроченные сессии пользователей и операторов."); + /// + /// Доля израсходованного ИИ-бюджета по тенантам + /// + public static readonly ObservableGauge AiBudgetUsedRatio = + Meter.CreateObservableGauge( + AiBudgetUsedRatioName, + ObserveBudgetRatios, + description: "Доля израсходованного ИИ-бюджета периода по тенантам (0..1)."); + + private static readonly ConcurrentDictionary BudgetRatios = new(StringComparer.Ordinal); + private static long _pipelineQueueDepth; private static long _mlOutboxDepth; private static long _activeSessions; @@ -163,4 +185,23 @@ public static class DealMetrics /// /// Активные непросроченные сессии пользователей и операторов. public static void SetActiveSessions(long count) => Interlocked.Exchange(ref _activeSessions, count); + + /// + /// Публикует доли расхода ИИ-бюджета по тенантам (заменяет предыдущий снимок) + /// + /// Пары (id тенанта, доля 0..1). + public static void ReplaceBudgetRatios(IEnumerable<(string TenantId, double Ratio)> ratios) + { + BudgetRatios.Clear(); + foreach ((string tenantId, double ratio) in ratios) + { + BudgetRatios[tenantId] = Math.Clamp(ratio, 0.0, 1.0); + } + } + + // Снимок долей расхода бюджета для observable-gauge (по измерению на тенант). + private static IEnumerable> ObserveBudgetRatios() + => BudgetRatios.Select(pair => new Measurement( + pair.Value, + new TagList { { TenantTagName, pair.Key } })); } diff --git a/src/core/tests/Deal.Tests.Unit/Api/RuntimeDepthsCollectorTests.cs b/src/core/tests/Deal.Tests.Unit/Api/RuntimeDepthsCollectorTests.cs index a3bde6b..3f315b1 100644 --- a/src/core/tests/Deal.Tests.Unit/Api/RuntimeDepthsCollectorTests.cs +++ b/src/core/tests/Deal.Tests.Unit/Api/RuntimeDepthsCollectorTests.cs @@ -23,6 +23,7 @@ public sealed class RuntimeDepthsCollectorTests { private static readonly Guid TenantA = Guid.Parse("11111111-1111-1111-1111-111111111111"); private static readonly Guid TenantB = Guid.Parse("22222222-2222-2222-2222-222222222222"); + private static readonly Guid TenantC = Guid.Parse("33333333-3333-3333-3333-333333333333"); [Fact] public async Task CollectAsync_SumsQueueAndOutboxAcrossTenants() @@ -50,6 +51,53 @@ public sealed class RuntimeDepthsCollectorTests Assert.Equal(3, depths.PipelineQueue); Assert.Equal(1, depths.MlOutbox); Assert.Equal(0, depths.ActiveSessions); + Assert.NotNull(depths.Budgets); + Assert.Equal(2, depths.Budgets!.Count); + } + + [Fact] + public async Task CollectAsync_CollectsBudgetRatiosPerTenant() + { + var limits = new FakeTenantLimitStore(); + DateTimeOffset periodStart = DateTimeOffset.UtcNow; + limits.Preload(TenantA, 100, TenantLimitPeriods.Month, periodStart, 80); + limits.Preload(TenantB, 0, TenantLimitPeriods.Month, periodStart, 55); + limits.Preload(TenantC, 100, TenantLimitPeriods.Month, periodStart, 150); + + RuntimeDepthsCollector collector = Build( + PipelineStores(TenantA, TenantB, TenantC), + OutboxStores(TenantA, TenantB, TenantC), + TenantRecords(TenantA, TenantB, TenantC), + limits); + + RuntimeDepthsDto depths = await collector.CollectAsync(CancellationToken.None); + + Assert.NotNull(depths.Budgets); + Assert.Equal(3, depths.Budgets!.Count); + Assert.Equal(0.8, RatioFor(depths.Budgets, TenantA), 3); + Assert.Equal(0.0, RatioFor(depths.Budgets, TenantB), 3); + Assert.Equal(1.0, RatioFor(depths.Budgets, TenantC), 3); + } + + [Fact] + public async Task CollectAsync_BudgetReadFailure_SkipsFailedTenant() + { + var limits = new FakeTenantLimitStore(); + limits.FailStateReads.Add(TenantA); + limits.Preload(TenantB, 100, TenantLimitPeriods.Month, DateTimeOffset.UtcNow, 80); + + RuntimeDepthsCollector collector = Build( + PipelineStores(TenantA, TenantB), + OutboxStores(TenantA, TenantB), + TenantRecords(TenantA, TenantB), + limits); + + RuntimeDepthsDto depths = await collector.CollectAsync(CancellationToken.None); + + Assert.NotNull(depths.Budgets); + TenantBudgetUsageDto budget = Assert.Single(depths.Budgets!); + Assert.Equal(TenantB.ToString("N"), budget.TenantId); + Assert.Equal(0.8, budget.UsedRatio, 3); } [Fact] @@ -74,11 +122,26 @@ public sealed class RuntimeDepthsCollectorTests Status = PipelineQueueStatuses.New, }; + private static Dictionary PipelineStores(params Guid[] tenants) + => tenants.ToDictionary(tenant => tenant.ToString("N"), _ => new FakePipelineStore()); + + private static Dictionary OutboxStores(params Guid[] tenants) + => tenants.ToDictionary(tenant => tenant.ToString("N"), _ => new FakeMlLearningStore()); + + private static TenantRecordDto[] TenantRecords(params Guid[] tenants) + => tenants + .Select((tenant, index) => new TenantRecordDto(tenant, $"T{index}", "active", DateTimeOffset.UtcNow)) + .ToArray(); + + private static double RatioFor(IReadOnlyList budgets, Guid tenant) + => budgets.Single(budget => budget.TenantId == tenant.ToString("N")).UsedRatio; + // Собирает коллектор поверх tenant-scoped фейков (как реальные адаптеры по ITenantContext). private static RuntimeDepthsCollector Build( IReadOnlyDictionary pipelineByTenant, IReadOnlyDictionary outboxByTenant, - IReadOnlyList? tenants = null) + IReadOnlyList? tenants = null, + ITenantLimitStore? limitStore = null) { var services = new ServiceCollection(); services.AddSingleton(); @@ -94,6 +157,7 @@ public sealed class RuntimeDepthsCollectorTests services.AddScoped(); services.AddScoped(); services.AddScoped(provider => outboxByTenant[CurrentTenant(provider)]); + services.AddSingleton(limitStore ?? new FakeTenantLimitStore()); ServiceProvider provider = services.BuildServiceProvider(); return new RuntimeDepthsCollector( diff --git a/src/core/tests/Deal.Tests.Unit/Modules/Tenants/FakeTenantLimitStore.cs b/src/core/tests/Deal.Tests.Unit/Modules/Tenants/FakeTenantLimitStore.cs index 2be5a27..cb1cb2c 100644 --- a/src/core/tests/Deal.Tests.Unit/Modules/Tenants/FakeTenantLimitStore.cs +++ b/src/core/tests/Deal.Tests.Unit/Modules/Tenants/FakeTenantLimitStore.cs @@ -55,6 +55,11 @@ public sealed class FakeTenantLimitStore : ITenantLimitStore /// public string TenantStatus { get; set; } + /// + /// Тенанты, для которых бросает исключение (устойчивость сборщиков) + /// + public HashSet FailStateReads { get; } = new(); + /// /// Кладёт готовую строку лимита /// @@ -124,6 +129,11 @@ public sealed class FakeTenantLimitStore : ITenantLimitStore /// public Task GetStateAsync(Guid tenantId, CancellationToken ct) { + if (FailStateReads.Contains(tenantId)) + { + throw new InvalidOperationException($"Фейковый сбой чтения состояния лимита тенанта {tenantId}."); + } + Row row = Ensure(tenantId, TokenBudgetDefaults.Default); ResetIfPeriodExpired(row); return Task.FromResult(ToStateDto(row, tenantId)); diff --git a/src/core/tests/Deal.Tests.Unit/SharedKernel/Observability/DealMetricsBudgetRatioTests.cs b/src/core/tests/Deal.Tests.Unit/SharedKernel/Observability/DealMetricsBudgetRatioTests.cs new file mode 100644 index 0000000..a1870b4 --- /dev/null +++ b/src/core/tests/Deal.Tests.Unit/SharedKernel/Observability/DealMetricsBudgetRatioTests.cs @@ -0,0 +1,59 @@ +using System.Diagnostics.Metrics; +using Deal.SharedKernel.Observability; + +namespace Deal.Tests.Unit.SharedKernel.Observability; + +public sealed class DealMetricsBudgetRatioTests +{ + [Fact] + public void ReplaceBudgetRatios_ClampsAndReplacesSnapshot() + { + var measurements = new Dictionary(StringComparer.Ordinal); + using var listener = new MeterListener(); + listener.InstrumentPublished = (instrument, current) => + { + if (instrument.Meter.Name == DealMetrics.MeterName + && instrument.Name == DealMetrics.AiBudgetUsedRatioName) + { + current.EnableMeasurementEvents(instrument); + } + }; + listener.SetMeasurementEventCallback((_, value, tags, _) => + { + string? tenant = null; + foreach (KeyValuePair tag in tags) + { + if (tag.Key == DealMetrics.TenantTagName) + { + tenant = tag.Value?.ToString(); + } + } + + if (tenant is not null) + { + measurements[tenant] = value; + } + }); + listener.Start(); + + DealMetrics.ReplaceBudgetRatios(new[] + { + ("tenant-a", 1.5), + ("tenant-b", -0.2), + }); + measurements.Clear(); + listener.RecordObservableInstruments(); + + Assert.Equal(1.0, measurements["tenant-a"], 3); + Assert.Equal(0.0, measurements["tenant-b"], 3); + Assert.Equal(2, measurements.Count); + + DealMetrics.ReplaceBudgetRatios(new[] { ("tenant-c", 0.4) }); + measurements.Clear(); + listener.RecordObservableInstruments(); + + Assert.Equal(0.4, measurements["tenant-c"], 3); + Assert.False(measurements.ContainsKey("tenant-a")); + Assert.Single(measurements); + } +}