From 6d074834a7d617ea532cce69888ea3128351ff8f Mon Sep 17 00:00:00 2001 From: Rustam Khalimov Date: Fri, 11 Sep 2026 16:32:49 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B5=D0=B2=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=B8=20=D0=B2=D1=85=D0=BE=D0=B4=D1=8F=D1=89=D0=B8=D0=B9?= =?UTF-8?q?=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=20=D0=B8=D1=81=D1=82=D0=BE?= =?UTF-8?q?=D1=87=D0=BD=D0=B8=D0=BA=D0=BE=D0=B2=20=D0=BD=D0=B0=20generic-?= =?UTF-8?q?=D0=BA=D0=BE=D0=BD=D1=82=D1=80=D0=B0=D0=BA=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Добавлен sources.proto с PushSource; приём в ядре вынесен в SourceIngressGrpcService с SourceProtoMapper и ISourceIngestObserver, тенант определяется IngressTenantResolver. Из telegram.proto удалён PushMessage, telegram-сервис шлёт generic-записи, превью сохраняет TelegramSourceIngestObserver. --- backlog.md | 2 +- docs/superpowers/STATUS.md | 6 +- .../2026-09-11-source-contract-design.md | 14 +- .../Техническая-документация-Дейл.md | 10 +- src/contracts/Deal.Proto.csproj | 1 + src/contracts/sources.proto | 91 ++++++ src/contracts/telegram.proto | 35 +-- src/core/Deal.Api/Program.cs | 4 + .../Deal.Api/Sources/IngressTenantResolver.cs | 64 ++++ .../Sources/SourceIngressGrpcService.cs | 103 +++++++ .../Deal.Api/Sources/SourceProtoMapper.cs | 85 ++++++ .../Telegram/TelegramIngressService.cs | 161 +--------- .../Sources/ISourceIngestObserver.cs | 15 + .../Application/TelegramModuleRegistrar.cs | 2 + .../TelegramSourceIngestObserver.cs | 47 +++ .../Deal.Modules.Telegram.csproj | 1 + .../Api/IngressRateLimitInterceptorTests.cs | 39 +-- .../Support/SourceIngressGrpcServiceTests.cs | 288 ++++++++++++++++++ .../Support/SourceProtoMapperTests.cs | 211 +++++++++++++ .../Support/TelegramIngressServiceTests.cs | 205 +------------ .../Support/TelegramIngressTestHost.cs | 4 + .../TelegramSourceIngestObserverTests.cs | 109 +++++++ .../Grpc/CoreIngressClientTests.cs | 38 ++- .../Grpc/DialogRpcTests.cs | 2 +- .../Grpc/FakeIngressServer.cs | 111 +++++-- .../Grpc/RealtimeListenerTests.cs | 10 +- .../Deal.Telegram.Tests/Grpc/TestDoubles.cs | 21 +- .../Telegram/BackfillServiceTests.cs | 13 +- .../Telegram/RealtimeSweepTests.cs | 4 +- .../Deal.Telegram/Core/CoreIngressClient.cs | 67 ++-- .../Deal.Telegram/Core/ICoreIngressClient.cs | 11 +- .../Deal.Telegram/Dialogs/BackfillService.cs | 8 +- .../Dialogs/DialogProtoMapper.cs | 33 +- .../Deal.Telegram/Dialogs/RealtimeListener.cs | 10 +- .../Deal.Telegram/Dialogs/RealtimeSweep.cs | 7 +- .../Hosting/RealtimeMonitorService.cs | 2 +- .../Sessions/SessionErrorMessages.cs | 2 +- .../Deal.Telegram/Telegram/ISessionClient.cs | 2 +- .../Deal.Telegram/Telegram/TelegramMessage.cs | 6 +- 39 files changed, 1312 insertions(+), 532 deletions(-) create mode 100644 src/contracts/sources.proto create mode 100644 src/core/Deal.Api/Sources/IngressTenantResolver.cs create mode 100644 src/core/Deal.Api/Sources/SourceIngressGrpcService.cs create mode 100644 src/core/Deal.Api/Sources/SourceProtoMapper.cs create mode 100644 src/core/Deal.Modules.Cards/Application/Sources/ISourceIngestObserver.cs create mode 100644 src/core/Deal.Modules.Telegram/Application/TelegramSourceIngestObserver.cs create mode 100644 src/core/tests/Deal.Tests.Unit/Support/SourceIngressGrpcServiceTests.cs create mode 100644 src/core/tests/Deal.Tests.Unit/Support/SourceProtoMapperTests.cs create mode 100644 src/core/tests/Deal.Tests.Unit/Support/TelegramSourceIngestObserverTests.cs diff --git a/backlog.md b/backlog.md index a2807b6..f10a75f 100644 --- a/backlog.md +++ b/backlog.md @@ -59,7 +59,7 @@ | TD-OLD-DOCS | Исторические доки несут старые термины под пометками (переписывать не нужно) | docs sweep | P3 | TECHDEBT | | TD-SOURCE-PROVIDER | `ISourceContentProvider`/`SourceContentResolver` и `GET /api/cards/{id}/source` добавлены; осталось — реализовать провайдеры источников (telegram/local/file) с ленивой догрузкой | generic source 2026-09-11 | P1 | BACKLOG | | TD-STORE-ATTACH | Выгрузка вложений источника в Storage-сервис адаптером + `TelegramSourceContentProvider` (только telegram-сервис) + рендер `DataRef` в UI | generic source 2026-09-11 | P1 | BACKLOG | -| TD-TG-CORE-SPLIT | Перенести оставшуюся Telegram-специфику ядра (`TelegramStore`, таблицы `Dialogs`/`TgMessages`, Discovery, ingress-proto) в telegram-сервис; generic ingress-контракт источников | generic source 2026-09-11 | P2 | TECHDEBT | +| TD-TG-CORE-SPLIT | Перенести оставшуюся Telegram-специфику ядра (`TelegramStore`, таблицы `Dialogs`/`TgMessages`, Discovery) в telegram-сервис. **Сделано (2026-09-11):** входящий поток переведён на generic `sources.proto`/`PushSource`, `PushMessage` удалён, приём в ядре generic (`SourceIngressGrpcService`). Осталось: каталог (`SyncDialogs`/`ReportStatus`), `TelegramStore`, Discovery | generic source 2026-09-11 | P2 | TECHDEBT | | TD-SOURCE-CONTACTS | Квалификатор контактов знает форматы профилей (t.me/`@handle`) — вынести в расширяемые правила источников | generic source 2026-09-11 | P3 | BACKLOG | | TD-APIMAP-COUNT | Ручной подсчёт числа операторских ручек в `api-map` (может расходиться с группировкой) | docs sweep | P3 | TECHDEBT | diff --git a/docs/superpowers/STATUS.md b/docs/superpowers/STATUS.md index 22aab98..c495a5e 100644 --- a/docs/superpowers/STATUS.md +++ b/docs/superpowers/STATUS.md @@ -8,8 +8,10 @@ > `SourceItem` (`SourceRef` + `SourceContent`, вложения — `DataRef` → общий Storage). Дашборды/канбан/конвейер > больше не знают о Telegram; Telegram-специфика — только в тонком адаптере приёма. Tenant-миграции > пересозданы с нуля (init). Добавлены extension-point `ISourceContentProvider`/`SourceContentResolver` и -> `GET /api/cards/{id}/source`. Ядро: build 4 sln 0/0, `Deal.Tests.Unit` **1280/1280 PASS**, фронт `build` + -> `lint:i18n` зелёные. Детали — `docs/superpowers/specs/2026-09-11-source-contract-design.md`. +> `GET /api/cards/{id}/source`. Входящий поток источников — generic (`sources.proto`/`PushSource`, +> `SourceIngressGrpcService`), `PushMessage` из telegram.proto удалён. Ядро: build 5 sln 0/0, +> `Deal.Tests.Unit` **1289/1289 PASS**, telegram **125/125**, фронт `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/superpowers/specs/2026-09-11-source-contract-design.md b/docs/superpowers/specs/2026-09-11-source-contract-design.md index d0eab20..d3a4989 100644 --- a/docs/superpowers/specs/2026-09-11-source-contract-design.md +++ b/docs/superpowers/specs/2026-09-11-source-contract-design.md @@ -165,7 +165,19 @@ API ядра: `GET /api/cards/{id}/source` → generic контент. ### Маппинг источника -- Адаптер Telegram (в ядре — тонкий край приёма): `Kind="telegram"`, `ExternalId=MsgId`, +- Telegram-адаптер (в ядре — тонкий край приёма): `Kind="telegram"`, `ExternalId=MsgId`, `OriginRef=DialogId`, `DisplayName=ChannelName`, `Extra["hue"]=ChannelHue` (иначе дефолт), `ReceivedAt=msgAt`, `Content.Text=Text`, `Content.Author=ChannelName`. - Дашборды/карточки/конвейер работают только с `SourceRef`/`SourceContent`; Telegram-поля не проходят дальше адаптера. + +### Входящий поток (generic, 2026-09-11) + +- `src/contracts/sources.proto` → сервис `SourceIngressService.PushSource` с generic-типами + `SourceRefProto`/`SourceContentProto`/`DataRefProto`/`ContactRefProto`; вложения — ссылки на Storage. +- Ядро: `Deal.Api/Sources/SourceIngressGrpcService` (приём) + `SourceProtoMapper` (proto → домен) + + `ISourceIngestObserver` (вторичная обработка принятой записи, сбой наблюдателя не влияет на приём) + + `IngressTenantResolver` (тенант по metadata). +- Из `telegram.proto` удалён `IngressService.PushMessage` (остались `SyncDialogs`/`ReportStatus`); + telegram-сервис шлёт записи через `PushSource` (`kind="telegram"`). Превью каталога/TgMessages сохраняет + `TelegramSourceIngestObserver` (ядро, telegram-модуль — единственное место с telegram-спецификой приёма). +- Любой другой источник (whatsapp/avito/файл/excel) шлёт тот же `PushSource` со своим `source.kind`. diff --git a/docs/technical/Техническая-документация-Дейл.md b/docs/technical/Техническая-документация-Дейл.md index 89641c9..6c7aa3a 100644 --- a/docs/technical/Техническая-документация-Дейл.md +++ b/docs/technical/Техническая-документация-Дейл.md @@ -163,7 +163,8 @@ settings(Key varchar(200) PK, ValueJson text, UpdatedAt timestamptz) -- - `telegram.proto` (пакет `deal.telegram.v1`) — два сервиса: `TelegramService` — команды ядра к telegram-service (GetStatus, StartPhone, StartQr, SendCode, SendPassword, Logout, RefreshDialogs, SetMonitor, SetMonitorAll, Backfill, ReadRecent, Search, GetInfo, ReadForEval, Join, Leave); - `IngressService` — исходящий поток telegram-service → core (PushMessage, SyncDialogs, ReportStatus; + `SourceIngressService.PushSource` (generic-контракт источников, `sources.proto`) и `IngressService` + (SyncDialogs, ReportStatus) — исходящий поток telegram-service → core; сервер — gRPC-ингресс core :5082). - `ml.proto` (пакет `deal.ml.v1`) — `MlService`: Predict (text → {take,label,scores,margin,type}), Status, Reset, TrainBatch; всё с metadata `tenant-id`. @@ -959,7 +960,7 @@ ml 36/36 PASS; build 0/0 всех четырёх sln; `docker compose -f deploy/ |---|---|---| | `deal-postgres` | **5433** | БД (host-порт; внутри — 5432) | | `deal-minio` | **9000/9001** | S3-API / консоль (файлы вложений; в Local-режиме необязателен) | -| `deal-core` (Deal.Api) | HTTP **5080**, gRPC-ингресс **5082** | портал `/api` + приём PushMessage/SyncDialogs/ReportStatus | +| `deal-core` (Deal.Api) | HTTP **5080**, gRPC-ингресс **5082** | портал `/api` + приём PushSource/SyncDialogs/ReportStatus | | `deal-telegram-service` | **5101** | Telegram: сессии/QR/диалоги/мониторинг/backfill/discovery-операции | | `deal-ai-service` | **5102** | LLM-фасад: Filter/Classify/GenerateKeywords/EvaluateFit | | `deal-ml-service` | **5103** | инкрементальная модель per-tenant: predict/train/status/reset | @@ -1032,6 +1033,11 @@ docker compose -f deploy/compose.dev.yml down # погасить ст `PipelineIngestService.EnqueueAsync` (тот же контракт, что demo-ingest) + превью в `TgMessages`; `SyncDialogs` → синхронизация каталога/мониторинга; `ReportStatus` → KV `tgStatus`/`tgAccount` + SSE `system_status`/тосты переходов. + + > **Актуально с 2026-09-11:** приём записей вынесен из `TelegramIngressService` в generic + > `Deal.Api/Sources/SourceIngressGrpcService` (`sources.proto` → `PushSource`, proto → домен через + > `SourceProtoMapper`, тенант — `IngressTenantResolver`, превью — `TelegramSourceIngestObserver`). + > `TelegramIngressService` обслуживает только `SyncDialogs`/`ReportStatus`. - Эндпоинты 1:1 api-map §3.3 (14 шт.): статус (§4.9 — live-поля фазы, `account` из KV, `monitored` из `count(Dialogs WHERE Monitor)`, `keysSet`), start-phone/start-qr/send-code/send-password/logout, QR-image (SVG, Net.Codecrete.QrCodeGenerator; 404 «QR не активен — начните вход по QR»), dialogs/refresh/ diff --git a/src/contracts/Deal.Proto.csproj b/src/contracts/Deal.Proto.csproj index a1535e8..b70c81a 100644 --- a/src/contracts/Deal.Proto.csproj +++ b/src/contracts/Deal.Proto.csproj @@ -39,6 +39,7 @@ + diff --git a/src/contracts/sources.proto b/src/contracts/sources.proto new file mode 100644 index 0000000..fe9f647 --- /dev/null +++ b/src/contracts/sources.proto @@ -0,0 +1,91 @@ +// +// Контракт входящего потока источников: сервис-источник → ядро. +// +// Любой источник (telegram, whatsapp, avito, сайт, файл, excel) приводит свои +// данные к единому контракту SourceItem и шлёт их одним вызовом PushSource. +// Ядро не знает о природе источника: вид задаётся полем source.kind. +// +// tenant-id — id тенанта (metadata; единственный источник принадлежности); +// service-token — общий токен сервисов (env DEAL_SERVICE_TOKEN); неверный/ +// пустой → UNAUTHENTICATED. +// +// Файлы вложений источник сам выгружает в сервис данных (storage.proto) и +// передаёт здесь ссылкой в content.data. +syntax = "proto3"; + +package deal.sources.v1; + +option csharp_namespace = "Deal.Grpc.Sources"; + +service SourceIngressService { + // Запись источника → очередь конвейера ядра. Дубль (вид+оригинал+внешний id) + // уже в очереди — очередь не растёт (duplicate=true). + rpc PushSource(PushSourceRequest) returns (PushSourceReply); +} + +message PushSourceRequest { + SourceRefProto source = 1; + SourceContentProto content = 2; + // True — вернуть из отсева: правила/устарелость/ИИ-фильтр пропускаются. + bool force = 3; +} + +message SourceRefProto { + // Дискриминатор источника, задаёт владелец (telegram/file/local/...). + string kind = 1; + // Идентификатор записи в источнике (сообщение/строка/файл). + optional string external_id = 2; + // Подпись источника для интерфейса. + optional string display_name = 3; + // Ссылка на оригинал (url, deep-link, путь). + optional string origin_ref = 4; + optional string author = 5; + // Время получения записи, epoch-ms. + int64 received_at = 6; + // Прочие метаданные источника (в т.ч. цвет интерфейса — hue). + map extra = 7; +} + +message SourceContentProto { + optional string text = 1; + optional string html = 2; + optional string author = 3; + optional string subject = 4; + // Вложения: ссылки на объекты сервиса данных. + repeated DataRefProto data = 5; + // Ссылки, не являющиеся файлами. + repeated string links = 6; + repeated ContactRefProto contacts = 7; + map extra = 8; +} + +message DataRefProto { + string id = 1; + string ref = 2; + optional string kind = 3; + optional string mime_type = 4; + optional string file_name = 5; + optional int64 size = 6; + optional int32 width = 7; + optional int32 height = 8; + optional double duration_sec = 9; + optional string preview_ref = 10; + optional string caption = 11; + optional int32 order = 12; + map meta = 13; +} + +message ContactRefProto { + optional string kind = 1; + optional string name = 2; + optional string phone = 3; + optional string email = 4; + optional string url = 5; +} + +message PushSourceReply { + // True — запись принята (пустой текст и пустой контент — accepted=false). + bool accepted = 1; + // True — дубль уже в очереди (очередь не выросла). + bool duplicate = 2; +} diff --git a/src/contracts/telegram.proto b/src/contracts/telegram.proto index 4d8cea7..ca1de0f 100644 --- a/src/contracts/telegram.proto +++ b/src/contracts/telegram.proto @@ -2,8 +2,10 @@ // * TelegramService — команды ядра к telegram-service (порт-гейт ITelegramGateway): // подключение/отключение аккаунта, каталог диалогов, мониторинг, backfill, // превью, discovery-операции (поиск/инфо/чтение/вступление/выход); -// * IngressService — исходящий поток telegram-service → ядро: сырые сообщения -// (PushMessage), синхронизация каталога (SyncDialogs), статус аккаунта +// * IngressService — исходящий поток telegram-service → ядро: синхронизация +// каталога (SyncDialogs), статус аккаунта (ReportStatus). +// * Сообщения мониторящихся диалогов приходят в ядро generic-контрактом +// sources.proto (PushSource) — ядро не привязано к природе источника. // // // tenant-id — id тенанта (строка; единственный источник принадлежности, @@ -328,11 +330,6 @@ message LeaveReply { // --------------------------------------------------------------------------- service IngressService { - // Новое/догоняющее сообщение мониторящегося диалога → очередь пайплайна - // ядра (PipelineIngestService.EnqueueAsync, контракт demo-ingest; + превью в - // TgMessages). Дубль dialog+msgId уже в очереди — не растёт (duplicate=true). - rpc PushMessage(PushMessageRequest) returns (PushMessageReply); - // Синхронизация каталога диалогов: ядро применяет entries (SyncFromTelegram: // авто-мониторинг по autoMonitorNew, обновление, удаление отсутствующих) и // отвечает актуальным списком monitored id — сервис держит зеркало @@ -342,30 +339,6 @@ service IngressService { rpc ReportStatus(ReportStatusRequest) returns (ReportStatusReply); } -// дубль-гвард; msg_at — время исходного сообщения, без него ядро подставит now. -message PushMessageRequest { - // Id диалога-источника (подписанный; пуст — приём no-op). - string dialog_id = 1; - // Имя канала/диалога (title/first_name или id). - string channel_name = 2; - // Username канала/диалога (пуст, если нет). - string channel_handle = 3; - string channel_hue = 4; - // Id исходного сообщения в Telegram (дубль-гвард dialog+msgId). - optional int64 msg_id = 5; - // Текст сообщения (сервис шлёт как есть; приём обрежет до 6000). - string text = 6; - // Время исходного сообщения, epoch-ms; пусто — ядро подставит now. - optional int64 msg_at = 7; -} - -message PushMessageReply { - // True — сообщение принято (no-op с пустым текстом/диалогом — accepted=false). - bool accepted = 1; - // True — дубль dialog_id+msg_id уже в очереди (очередь не выросла). - bool duplicate = 2; -} - message SyncDialogsRequest { // Актуальный каталог диалогов (собирает сервис, как refresh_dialogs). repeated DialogEntry entries = 1; diff --git a/src/core/Deal.Api/Program.cs b/src/core/Deal.Api/Program.cs index 7182d58..9a2b0c3 100644 --- a/src/core/Deal.Api/Program.cs +++ b/src/core/Deal.Api/Program.cs @@ -9,6 +9,7 @@ using Deal.Api.Logging; using Deal.Api.Middleware; using Deal.Api.Observability; using Deal.Api.Services; +using Deal.Api.Sources; using Deal.Api.Telegram; using Deal.Contracts.Integrations.Abstractions; using Deal.Infrastructure; @@ -199,6 +200,8 @@ if (rateLimitOptions.Enabled) } builder.Services.AddScoped(); +builder.Services.AddScoped(); +builder.Services.AddScoped(); builder.Services .AddGrpcHealthChecks() @@ -361,6 +364,7 @@ app.MapTelegramEndpoints(); app.MapTelegramQrImageEndpoint(); app.MapDiscoveryEndpoints(); app.MapGrpcService().DisableRateLimiting(); +app.MapGrpcService().DisableRateLimiting(); app.MapGrpcHealthChecksService().DisableRateLimiting(); app.Run(); diff --git a/src/core/Deal.Api/Sources/IngressTenantResolver.cs b/src/core/Deal.Api/Sources/IngressTenantResolver.cs new file mode 100644 index 0000000..6cd3a4c --- /dev/null +++ b/src/core/Deal.Api/Sources/IngressTenantResolver.cs @@ -0,0 +1,64 @@ +using Deal.Modules.Tenants.Application.Abstractions; +using Deal.Modules.Tenants.Application.Models; +using Grpc.Core; + +namespace Deal.Api.Sources; + +/// +/// Определение тенанта gRPC-вызова входящего потока по metadata. +/// +/// Фабрика scope для чтения реестра тенантов. +/// Логгер отказов. +public sealed class IngressTenantResolver( + IServiceScopeFactory scopeFactory, + ILogger logger) +{ + /// + /// Ключ metadata с id тенанта. + /// + public const string TenantIdMetadataKey = "tenant-id"; + + private const string MissingTenantIdDetail = "tenant-id отсутствует в metadata"; + + /// + /// Читает обязательный id тенанта из metadata. + /// + /// Контекст вызова. + /// Значение tenant-id. + /// tenant-id отсутствует (UNAUTHENTICATED). + public static string RequireTenantIdMetadata(ServerCallContext context) + { + string? tenantId = context.RequestHeaders.GetValue(TenantIdMetadataKey); + if (string.IsNullOrWhiteSpace(tenantId)) + { + throw new RpcException(new Status(StatusCode.Unauthenticated, MissingTenantIdDetail)); + } + + return tenantId; + } + + /// + /// Находит тенанта вызова в реестре. + /// + /// Контекст вызова. + /// Запись тенанта либо null — id неизвестен/не Guid/нет в реестре. + public async Task ResolveAsync(ServerCallContext context) + { + string tenantId = RequireTenantIdMetadata(context); + if (!Guid.TryParse(tenantId, out Guid tenantGuid)) + { + logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (id не Guid)", context.Method, tenantId); + return null; + } + + await using AsyncServiceScope registryScope = scopeFactory.CreateAsyncScope(); + ITenantRepository repository = registryScope.ServiceProvider.GetRequiredService(); + TenantRecordDto? tenant = await repository.FindByIdAsync(tenantGuid, context.CancellationToken).ConfigureAwait(false); + if (tenant is null) + { + logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (нет в реестре)", context.Method, tenantId); + } + + return tenant; + } +} diff --git a/src/core/Deal.Api/Sources/SourceIngressGrpcService.cs b/src/core/Deal.Api/Sources/SourceIngressGrpcService.cs new file mode 100644 index 0000000..061e42f --- /dev/null +++ b/src/core/Deal.Api/Sources/SourceIngressGrpcService.cs @@ -0,0 +1,103 @@ +using Deal.Grpc.Sources; +using Deal.Modules.Cards.Application.Sources; +using Deal.Modules.Pipeline.Application.Abstractions; +using Deal.Modules.Pipeline.Application.Models; +using Deal.Modules.Pipeline.Application.Services; +using Deal.Modules.Tenants.Application.Models; +using Deal.SharedKernel.Tenants.Abstractions; +using Deal.SharedKernel.Tenants.Models; +using Grpc.Core; + +namespace Deal.Api.Sources; + +/// +/// gRPC-сервер входящего потока источников (generic-контракт). +/// +/// Фабрика tenant-scope приёма. +/// Определение тенанта вызова. +/// Логгер приёма. +public sealed class SourceIngressGrpcService( + IServiceScopeFactory scopeFactory, + IngressTenantResolver tenants, + ILogger logger) : SourceIngressService.SourceIngressServiceBase +{ + /// + /// PushSource — запись источника в очередь пайплайна тенанта. + /// + /// Запись источника (source + content + force). + /// Контекст вызова (metadata tenant-id + service-token). + /// accepted — запись принята (либо дубль), duplicate — уже была в очереди. + public override async Task PushSource(PushSourceRequest request, ServerCallContext context) + { + TenantRecordDto? tenant = await tenants.ResolveAsync(context).ConfigureAwait(false); + if (tenant is null) + { + return new PushSourceReply(); + } + + await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope(); + ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService(); + try + { + tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N"))); + SourceItem item = SourceProtoMapper.ToItem(request); + PipelineIngestService ingest = tenantScope.ServiceProvider.GetRequiredService(); + + PipelineIngestResultDto result = await ingest.EnqueueAsync( + new QueuedMessage { Item = item, Force = request.Force }, + context.CancellationToken).ConfigureAwait(false); + + foreach (ISourceIngestObserver observer in tenantScope.ServiceProvider.GetServices()) + { + await NotifySafelyAsync(observer, item, tenant.Id, context.CancellationToken).ConfigureAwait(false); + } + + logger.LogInformation( + "Аудит: PushSource {TenantId} {Kind} {ExternalId} → {Outcome}", + tenant.Id, + item.Source.Kind, + item.Source.ExternalId ?? "-", + result.Duplicate ? "duplicate" : result.Id is null ? "no-op" : "queued"); + + return new PushSourceReply + { + Accepted = result.Id is not null || result.Duplicate, + Duplicate = result.Duplicate, + }; + } + catch (OperationCanceledException) + { + throw; + } + catch (Exception exception) + { + logger.LogWarning(exception, "Аудит: PushSource {TenantId} → не принято (сбой схемы/БД)", tenant.Id); + return new PushSourceReply(); + } + finally + { + tenantContext.Reset(); + } + } + + // Наблюдатель — вторичная обработка: его сбой не влияет на приём записи. + private async Task NotifySafelyAsync( + ISourceIngestObserver observer, + SourceItem item, + Guid tenantId, + CancellationToken ct) + { + try + { + await observer.OnIngestedAsync(item, ct).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + throw; + } + catch (Exception exception) + { + logger.LogDebug(exception, "PushSource {TenantId}: наблюдатель {Observer} не отработал (приём не затронут)", tenantId, observer.GetType().Name); + } + } +} diff --git a/src/core/Deal.Api/Sources/SourceProtoMapper.cs b/src/core/Deal.Api/Sources/SourceProtoMapper.cs new file mode 100644 index 0000000..43d4f87 --- /dev/null +++ b/src/core/Deal.Api/Sources/SourceProtoMapper.cs @@ -0,0 +1,85 @@ +using Deal.Grpc.Sources; +using Deal.Modules.Cards.Application.Sources; + +namespace Deal.Api.Sources; + +/// +/// Маппинг generic-контракта источника в доменную модель. +/// +public static class SourceProtoMapper +{ + /// + /// Преобразует protobuf-запись источника в доменную. + /// + /// Запрос PushSource. + /// Запись источника домена. + public static SourceItem ToItem(PushSourceRequest request) => new() + { + Source = ToRef(request.Source), + Content = ToContent(request.Content), + }; + + private static SourceRef ToRef(SourceRefProto proto) => new() + { + Kind = proto.Kind, + ExternalId = proto.HasExternalId ? proto.ExternalId : null, + DisplayName = proto.HasDisplayName ? proto.DisplayName : null, + OriginRef = proto.HasOriginRef ? proto.OriginRef : null, + Author = proto.HasAuthor ? proto.Author : null, + ReceivedAt = proto.ReceivedAt != 0 ? DateTimeOffset.FromUnixTimeMilliseconds(proto.ReceivedAt) : DateTimeOffset.UtcNow, + Extra = proto.Extra.Count > 0 ? new Dictionary(proto.Extra) : null, + }; + + private static SourceContent ToContent(SourceContentProto proto) + { + var data = new List(proto.Data.Count); + foreach (DataRefProto item in proto.Data) + { + data.Add(ToData(item)); + } + + var contacts = new List(proto.Contacts.Count); + foreach (ContactRefProto contact in proto.Contacts) + { + contacts.Add(ToContact(contact)); + } + + return new SourceContent + { + Text = proto.HasText ? proto.Text : null, + Html = proto.HasHtml ? proto.Html : null, + Author = proto.HasAuthor ? proto.Author : null, + Subject = proto.HasSubject ? proto.Subject : null, + Data = data, + Links = proto.Links.Count > 0 ? proto.Links.ToList() : null, + Contacts = contacts.Count > 0 ? contacts : null, + Extra = proto.Extra.Count > 0 ? new Dictionary(proto.Extra) : null, + }; + } + + private static DataRef ToData(DataRefProto proto) => new() + { + Id = proto.Id, + Ref = proto.Ref, + Kind = proto.HasKind ? proto.Kind : null, + MimeType = proto.HasMimeType ? proto.MimeType : null, + FileName = proto.HasFileName ? proto.FileName : null, + Size = proto.HasSize ? proto.Size : null, + Width = proto.HasWidth ? proto.Width : null, + Height = proto.HasHeight ? proto.Height : null, + DurationSec = proto.HasDurationSec ? proto.DurationSec : null, + PreviewRef = proto.HasPreviewRef ? proto.PreviewRef : null, + Caption = proto.HasCaption ? proto.Caption : null, + Order = proto.HasOrder ? proto.Order : null, + Meta = proto.Meta.Count > 0 ? new Dictionary(proto.Meta) : null, + }; + + private static ContactRef ToContact(ContactRefProto proto) => new() + { + Kind = proto.HasKind ? proto.Kind : null, + Name = proto.HasName ? proto.Name : null, + Phone = proto.HasPhone ? proto.Phone : null, + Email = proto.HasEmail ? proto.Email : null, + Url = proto.HasUrl ? proto.Url : null, + }; +} diff --git a/src/core/Deal.Api/Telegram/TelegramIngressService.cs b/src/core/Deal.Api/Telegram/TelegramIngressService.cs index ef5f9c3..730a607 100644 --- a/src/core/Deal.Api/Telegram/TelegramIngressService.cs +++ b/src/core/Deal.Api/Telegram/TelegramIngressService.cs @@ -1,15 +1,11 @@ -using System.Globalization; using System.Text.Json; using Deal.Api.Events; +using Deal.Api.Sources; using Deal.Contracts.Integrations.Models; using Deal.Grpc.Telegram; -using Deal.Modules.Cards.Application.Sources; -using Deal.Modules.Pipeline.Application.Models; -using Deal.Modules.Pipeline.Application.Services; using Deal.Modules.Settings.Application.Abstractions; using Deal.Modules.Settings.Application.Models; using Deal.Modules.Telegram.Application; -using Deal.Modules.Tenants.Application.Abstractions; using Deal.Modules.Tenants.Application.Models; using Deal.SharedKernel.Tenants.Abstractions; using Deal.SharedKernel.Tenants.Models; @@ -18,14 +14,18 @@ using Grpc.Core; namespace Deal.Api.Telegram; /// -/// gRPC-сервер входящего потока telegram-service → ядро. +/// gRPC-сервер служебного потока telegram-service → ядро: каталог диалогов и статус аккаунта. /// public sealed class TelegramIngressService( IServiceScopeFactory scopeFactory, SseBroker broker, + IngressTenantResolver tenants, ILogger logger) : IngressService.IngressServiceBase { - public const string TenantIdMetadataKey = "tenant-id"; + /// + /// Ключ metadata с id тенанта. + /// + public const string TenantIdMetadataKey = IngressTenantResolver.TenantIdMetadataKey; private const string SystemStatusEventType = "system_status"; @@ -41,93 +41,12 @@ public sealed class TelegramIngressService( // Иконка тоста отключения (из набора Icon.vue фронта). private const string DisconnectedToastIcon = "logout"; - // Деталь отказа: metadata tenant-id отсутствует (UNAUTHENTICATED, README src/contracts). - private const string MissingTenantIdDetail = "tenant-id отсутствует в metadata"; - private static readonly JsonSerializerOptions StatusJsonOptions = new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase, PropertyNameCaseInsensitive = true, }; - /// - /// PushMessage — сообщение диалога в очередь пайплайна тенанта + превью. - /// - /// Сообщение из потока telegram-service. - /// Контекст вызова (metadata tenant-id + service-token). - /// accepted — сообщение принято (либо дубль), duplicate — уже было в очереди. - public override async Task PushMessage(PushMessageRequest request, ServerCallContext context) - { - TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false); - if (tenant is null) - { - return new PushMessageReply(); - } - - await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope(); - ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService(); - try - { - // Resolve ПОСЛЕ SetTenant: TenantDbContext (и его адаптеры) строятся от схемы текущего тенанта. - tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N"))); - PipelineIngestService ingest = tenantScope.ServiceProvider.GetRequiredService(); - - string hue = string.IsNullOrWhiteSpace(request.ChannelHue) ? SourceRefs.DefaultHue : request.ChannelHue; - DateTimeOffset receivedAt = request.HasMsgAt - ? DateTimeOffset.FromUnixTimeMilliseconds(request.MsgAt) - : DateTimeOffset.UtcNow; - var item = new SourceItem - { - Source = new SourceRef - { - Kind = "telegram", - ExternalId = request.HasMsgId ? request.MsgId.ToString(CultureInfo.InvariantCulture) : null, - OriginRef = request.DialogId, - DisplayName = request.ChannelName, - ReceivedAt = receivedAt, - Extra = new Dictionary { [SourceRefs.HueKey] = hue }, - }, - Content = new SourceContent - { - Text = request.Text, - Author = request.ChannelName, - }, - }; - - PipelineIngestResultDto result = await ingest.EnqueueAsync( - new QueuedMessage { Item = item }, - context.CancellationToken).ConfigureAwait(false); - - await SavePreviewSafelyAsync(tenantScope, tenant, request, context.CancellationToken).ConfigureAwait(false); - - logger.LogInformation( - "Аудит: PushMessage {TenantId} диалог {DialogId} msg {MsgId} → {Outcome}", - tenant.Id, - request.DialogId, - request.HasMsgId ? request.MsgId.ToString() : "-", - result.Duplicate ? "duplicate" : result.Id is null ? "no-op" : "queued"); - - return new PushMessageReply - { - Accepted = result.Id is not null || result.Duplicate, - Duplicate = result.Duplicate, - }; - } - catch (OperationCanceledException) - { - throw; - } - catch (Exception exception) - { - logger.LogWarning(exception, "Аудит: PushMessage {TenantId} → не принято (сбой схемы/БД)", tenant.Id); - return new PushMessageReply(); - } - finally - { - tenantContext.Reset(); - } - } - /// /// SyncDialogs — синхронизация каталога диалогов аккаунта. /// @@ -136,7 +55,7 @@ public sealed class TelegramIngressService( /// monitored_ids — диалоги с включённым мониторингом после применения каталога. public override async Task SyncDialogs(SyncDialogsRequest request, ServerCallContext context) { - TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false); + TenantRecordDto? tenant = await tenants.ResolveAsync(context).ConfigureAwait(false); if (tenant is null) { return new SyncDialogsReply(); @@ -193,7 +112,7 @@ public sealed class TelegramIngressService( /// ok — статус принят и сохранён. public override async Task ReportStatus(ReportStatusRequest request, ServerCallContext context) { - TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false); + TenantRecordDto? tenant = await tenants.ResolveAsync(context).ConfigureAwait(false); if (tenant is null) { return new ReportStatusReply(); @@ -237,68 +156,6 @@ public sealed class TelegramIngressService( } } - private async Task ResolveTenantAsync(ServerCallContext context) - { - string tenantId = RequireTenantIdMetadata(context); - if (!Guid.TryParse(tenantId, out Guid tenantGuid)) - { - logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (id не Guid)", context.Method, tenantId); - return null; - } - - await using AsyncServiceScope registryScope = scopeFactory.CreateAsyncScope(); - ITenantRepository repository = registryScope.ServiceProvider.GetRequiredService(); - TenantRecordDto? tenant = await repository.FindByIdAsync(tenantGuid, context.CancellationToken).ConfigureAwait(false); - if (tenant is null) - { - logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (нет в реестре)", context.Method, tenantId); - } - - return tenant; - } - - // Читает tenant-id из metadata (обязателен; отсутствие — UNAUTHENTICATED, README). - // context: Контекст вызова. - // Возвращает: Значение tenant-id. - private static string RequireTenantIdMetadata(ServerCallContext context) - { - string? tenantId = context.RequestHeaders.GetValue(TenantIdMetadataKey); - if (string.IsNullOrWhiteSpace(tenantId)) - { - throw new RpcException(new Status(StatusCode.Unauthenticated, MissingTenantIdDetail)); - } - - return tenantId; - } - - private async Task SavePreviewSafelyAsync( - AsyncServiceScope tenantScope, - TenantRecordDto tenant, - PushMessageRequest request, - CancellationToken ct) - { - try - { - DialogsService dialogs = tenantScope.ServiceProvider.GetRequiredService(); - DateTimeOffset? msgAt = request.HasMsgAt ? DateTimeOffset.FromUnixTimeMilliseconds(request.MsgAt) : null; - await dialogs.SavePreviewAsync( - request.DialogId, - request.HasMsgId ? request.MsgId : null, - request.Text, - msgAt, - ct).ConfigureAwait(false); - } - catch (OperationCanceledException) - { - throw; - } - catch (Exception exception) - { - // Превью — вторичная запись: приём не затронут (лог дебага, не ошибка RPC). - logger.LogDebug(exception, "PushMessage {TenantId}: превью не сохранено (приём не затронут)", tenant.Id); - } - } - // Публикует SSE system_status (каждый репорт) и тосты на переходах connected. // tenantId: Тенант канала (реестровый Guid). // previous: Предыдущий снимок из KV (null — первый репорт). diff --git a/src/core/Deal.Modules.Cards/Application/Sources/ISourceIngestObserver.cs b/src/core/Deal.Modules.Cards/Application/Sources/ISourceIngestObserver.cs new file mode 100644 index 0000000..af31efe --- /dev/null +++ b/src/core/Deal.Modules.Cards/Application/Sources/ISourceIngestObserver.cs @@ -0,0 +1,15 @@ +namespace Deal.Modules.Cards.Application.Sources; + +/// +/// Уведомление о принятой записи источника. +/// +public interface ISourceIngestObserver +{ + /// + /// Обрабатывает принятую запись источника. + /// + /// Принятая запись. + /// Токен отмены. + /// Завершается после обработки; сбой наблюдателя не влияет на приём. + public Task OnIngestedAsync(SourceItem item, CancellationToken ct); +} diff --git a/src/core/Deal.Modules.Telegram/Application/TelegramModuleRegistrar.cs b/src/core/Deal.Modules.Telegram/Application/TelegramModuleRegistrar.cs index 2e378e7..c0f4a5b 100644 --- a/src/core/Deal.Modules.Telegram/Application/TelegramModuleRegistrar.cs +++ b/src/core/Deal.Modules.Telegram/Application/TelegramModuleRegistrar.cs @@ -1,3 +1,4 @@ +using Deal.Modules.Cards.Application.Sources; using Microsoft.Extensions.DependencyInjection; namespace Deal.Modules.Telegram.Application; @@ -13,6 +14,7 @@ public static class TelegramModuleRegistrar public static IServiceCollection AddTelegramModule(this IServiceCollection services) { services.AddScoped(); + services.AddScoped(); return services; } } diff --git a/src/core/Deal.Modules.Telegram/Application/TelegramSourceIngestObserver.cs b/src/core/Deal.Modules.Telegram/Application/TelegramSourceIngestObserver.cs new file mode 100644 index 0000000..b630016 --- /dev/null +++ b/src/core/Deal.Modules.Telegram/Application/TelegramSourceIngestObserver.cs @@ -0,0 +1,47 @@ +using Deal.Modules.Cards.Application.Sources; +using Microsoft.Extensions.Logging; + +namespace Deal.Modules.Telegram.Application; + +/// +/// Превью принятых сообщений источника telegram: последнее сообщение каталога и TgMessages. +/// +/// Каталог диалогов тенанта. +/// Логгер сбоев записи. +public sealed class TelegramSourceIngestObserver( + DialogsService dialogs, + ILogger logger) : ISourceIngestObserver +{ + private const string TelegramKind = "telegram"; + + /// + public async Task OnIngestedAsync(SourceItem item, CancellationToken ct) + { + if (!string.Equals(item.Source.Kind, TelegramKind, StringComparison.OrdinalIgnoreCase)) + { + return; + } + + string dialogId = item.Source.OriginRef ?? string.Empty; + if (dialogId.Length == 0) + { + return; + } + + long? msgId = long.TryParse(item.Source.ExternalId, out long parsed) ? parsed : null; + try + { + await dialogs + .SavePreviewAsync(dialogId, msgId, item.Content.Text ?? string.Empty, item.Source.ReceivedAt, ct) + .ConfigureAwait(false); + } + catch (OperationCanceledException) + { + throw; + } + catch (Exception exception) + { + logger.LogDebug(exception, "Превью {DialogId} не сохранено (приём не затронут)", dialogId); + } + } +} diff --git a/src/core/Deal.Modules.Telegram/Deal.Modules.Telegram.csproj b/src/core/Deal.Modules.Telegram/Deal.Modules.Telegram.csproj index 9fc6259..3a1cd3b 100644 --- a/src/core/Deal.Modules.Telegram/Deal.Modules.Telegram.csproj +++ b/src/core/Deal.Modules.Telegram/Deal.Modules.Telegram.csproj @@ -4,6 +4,7 @@ + diff --git a/src/core/tests/Deal.Tests.Unit/Api/IngressRateLimitInterceptorTests.cs b/src/core/tests/Deal.Tests.Unit/Api/IngressRateLimitInterceptorTests.cs index ec88020..900ed26 100644 --- a/src/core/tests/Deal.Tests.Unit/Api/IngressRateLimitInterceptorTests.cs +++ b/src/core/tests/Deal.Tests.Unit/Api/IngressRateLimitInterceptorTests.cs @@ -1,6 +1,6 @@ using Deal.Api.Configuration; using Deal.Api.Middleware; -using Deal.Grpc.Telegram; +using Deal.Grpc.Sources; using Deal.Modules.Pipeline.Application.Abstractions; using Deal.Modules.Tenants.Application.Abstractions; using Deal.Modules.Tenants.Application.Models; @@ -28,7 +28,7 @@ public sealed class IngressRateLimitInterceptorTests // Тенант B сценариев (в реестре) — проверка изоляции окон. private static readonly Guid TenantB = Guid.NewGuid(); - // Id диалога сценариев PushMessage. + // Id диалога-оригинала сценариев PushSource. private const string DialogId = "d_channel_rate_limit"; // Окно сценария по умолчанию: 2/мин — быстрый тест без ожидания реального окна. @@ -43,7 +43,7 @@ public sealed class IngressRateLimitInterceptorTests /// 3-й вызов тенанта в минуту /// [Fact] - public async Task PushMessage_ExceedingTenantWindow_ThirdRejectedOtherTenantPasses() + public async Task PushSource_ExceedingTenantWindow_ThirdRejectedOtherTenantPasses() { var registry = new FakeTenantRegistry(Tenant(TenantA), Tenant(TenantB)); @@ -52,8 +52,8 @@ public sealed class IngressRateLimitInterceptorTests services => services.AddScoped(_ => new FakePipelineStore()), async channel => { - PushMessageReply first = await PushAsync(channel, TenantA); - PushMessageReply second = await PushAsync(channel, TenantA); + PushSourceReply first = await PushAsync(channel, TenantA); + PushSourceReply second = await PushAsync(channel, TenantA); Assert.True(first.Accepted); Assert.True(second.Accepted); @@ -62,7 +62,7 @@ public sealed class IngressRateLimitInterceptorTests Assert.Equal(RateLimitPolicies.RejectedDetail, rejection.Status.Detail); // Окно TenantB — отдельная партиция: вызов проходит. - PushMessageReply tenantB = await PushAsync(channel, TenantB); + PushSourceReply tenantB = await PushAsync(channel, TenantB); Assert.True(tenantB.Accepted); }); } @@ -81,7 +81,7 @@ public sealed class IngressRateLimitInterceptorTests async channel => { // Окно ингресса — 1/мин: единственный разрешённый вызов исчерпывает лимит. - PushMessageReply push = await PushAsync(channel, TenantA); + PushSourceReply push = await PushAsync(channel, TenantA); Assert.True(push.Accepted); var healthClient = new Health.HealthClient(channel); @@ -128,22 +128,25 @@ public sealed class IngressRateLimitInterceptorTests GrpcIngressPerMinute = permitsPerMinute, }); - // Вызывает PushMessage с metadata сценария (service-token + tenant-id, deadline 10 с). + // Вызывает PushSource с metadata сценария (service-token + tenant-id, deadline 10 с). // channel: Канал к хосту. // tenantId: Id тенанта в metadata. - private static Task PushAsync(GrpcChannel channel, Guid tenantId) + private static Task PushAsync(GrpcChannel channel, Guid tenantId) { - var client = new IngressService.IngressServiceClient(channel); - var request = new PushMessageRequest + var client = new SourceIngressService.SourceIngressServiceClient(channel); + var request = new PushSourceRequest { - DialogId = DialogId, - ChannelName = "Канал", - ChannelHandle = "kanal_handle", - ChannelHue = "#a33", - Text = "Сообщение сценария лимита", - MsgId = 1, + Source = new SourceRefProto + { + Kind = "telegram", + OriginRef = DialogId, + ExternalId = "1", + DisplayName = "Канал", + }, + Content = new SourceContentProto { Text = "Сообщение сценария лимита" }, }; - AsyncUnaryCall call = client.PushMessageAsync( + request.Source.Extra["hue"] = "#a33"; + AsyncUnaryCall call = client.PushSourceAsync( request, new CallOptions(TelegramIngressTestHost.CallMetadata(ValidToken, tenantId), deadline: Deadline())); return call.ResponseAsync; diff --git a/src/core/tests/Deal.Tests.Unit/Support/SourceIngressGrpcServiceTests.cs b/src/core/tests/Deal.Tests.Unit/Support/SourceIngressGrpcServiceTests.cs new file mode 100644 index 0000000..0a3b07f --- /dev/null +++ b/src/core/tests/Deal.Tests.Unit/Support/SourceIngressGrpcServiceTests.cs @@ -0,0 +1,288 @@ +using Deal.Grpc.Sources; +using Deal.Modules.Cards.Application.Sources; +using Deal.Modules.Pipeline.Application.Abstractions; +using Deal.Modules.Pipeline.Application.Models; +using Deal.Modules.Tenants.Application.Abstractions; +using Deal.Modules.Tenants.Application.Models; +using Deal.Tests.Unit.Modules.Kanban; +using Deal.Tests.Unit.Modules.Tenants; +using Grpc.Core; +using Grpc.Net.Client; +using Microsoft.Extensions.DependencyInjection; + +namespace Deal.Tests.Unit.Support; + +/// +/// In-proc gRPC-тесты приёма источников (PushSource). +/// +public sealed class SourceIngressGrpcServiceTests +{ + // Токен сценариев теста. + private const string ValidToken = TelegramIngressTestHost.DefaultToken; + + // Тенант A сценариев (в реестре). + private static readonly Guid TenantA = Guid.NewGuid(); + + // Вид источника сценариев. + private const string Kind = "telegram"; + + // Id диалога-оригинала сценариев. + private const string OriginRef = "d_channel_100"; + + // Подпись источника сценариев. + private const string DisplayName = "Канал"; + + // Цвет источника сценариев. + private const string Hue = "#a33"; + + // Время получения записи сценариев (epoch-ms). + private const long ReceivedAtMs = 1_700_000_000_000; + + // ─── PushSource: приём, дубль, несуществующий тенант ────────────────── + + /// + /// PushSource кладёт строку очереди тенанта и принимает запись + /// + [Fact] + public async Task PushSource_ValidTenant_EnqueuesQueueRowAndAccepts() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync( + registry, + services => services.AddScoped(_ => store), + async channel => + { + PushSourceReply reply = await PushAsync( + channel, + Request("Вакансия Python-разработчика", externalId: "7"), + TenantA, + ValidToken); + + Assert.True(reply.Accepted); + Assert.False(reply.Duplicate); + + QueueItemDto row = Assert.Single(store.Queue); + Assert.Equal(Kind, row.Source.Kind); + Assert.Equal(OriginRef, row.Source.OriginRef); + Assert.Equal("7", row.Source.ExternalId); + Assert.Equal(DisplayName, row.Source.DisplayName); + Assert.Equal(DateTimeOffset.FromUnixTimeMilliseconds(ReceivedAtMs), row.Source.ReceivedAt); + Assert.Equal(Hue, row.Source.ResolveHue()); + Assert.Equal("Вакансия Python-разработчика", row.Text); + }); + } + + /// + /// Повтор PushSource того же вида+оригинала+внешнего id — duplicate=true, очередь не растёт + /// + [Fact] + public async Task PushSource_SameSourceKeyTwice_SecondIsDuplicateAndQueueNotGrown() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync( + registry, + services => services.AddScoped(_ => store), + async channel => + { + PushSourceReply first = await PushAsync(channel, Request("Сообщение канала", externalId: "42"), TenantA, ValidToken); + PushSourceReply second = await PushAsync(channel, Request("Сообщение канала", externalId: "42"), TenantA, ValidToken); + + Assert.True(first.Accepted); + Assert.True(second.Accepted); + Assert.True(second.Duplicate); + Assert.Single(store.Queue); + }); + } + + /// + /// PushSource без текста не принимается и очередь не растёт + /// + [Fact] + public async Task PushSource_EmptyText_NotAcceptedAndQueueNotGrown() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync( + registry, + services => services.AddScoped(_ => store), + async channel => + { + PushSourceReply reply = await PushAsync(channel, Request(" ", externalId: "9"), TenantA, ValidToken); + + Assert.False(reply.Accepted); + Assert.False(reply.Duplicate); + Assert.Empty(store.Queue); + }); + } + + /// + /// PushSource для несуществующего тенанта не падает + /// + [Fact] + public async Task PushSource_UnknownTenant_NotAcceptedWithoutRpcError() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync( + registry, + services => services.AddScoped(_ => store), + async channel => + { + PushSourceReply reply = await PushAsync(channel, Request("Сообщение чужого тенанта", externalId: "5"), Guid.NewGuid(), ValidToken); + + Assert.False(reply.Accepted); + Assert.False(reply.Duplicate); + Assert.Empty(store.Queue); + }); + } + + // ─── Интерцептор service-token и metadata tenant-id ───────────────────── + + /// + /// Запрос без metadata «service-token» → UNAUTHENTICATED. + /// + [Fact] + public async Task PushSource_WithoutToken_IsUnauthenticated() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync(registry, services => services.AddScoped(_ => store), async channel => + { + RpcException exception = await Assert.ThrowsAsync( + () => PushAsync(channel, Request("текст", externalId: "1"), TenantA, null)); + Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); + }); + } + + /// + /// Запрос с неверным токеном → UNAUTHENTICATED. + /// + [Fact] + public async Task PushSource_WithWrongToken_IsUnauthenticated() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync(registry, services => services.AddScoped(_ => store), async channel => + { + RpcException exception = await Assert.ThrowsAsync( + () => PushAsync(channel, Request("текст", externalId: "2"), TenantA, "wrong-token")); + Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); + }); + } + + /// + /// Fail-closed: DEAL_SERVICE_TOKEN не задан — RPC ингресса отклоняется даже с «каким-то» токеном. + /// + [Fact] + public async Task PushSource_UnsetEnvToken_FailsClosed() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await TelegramIngressTestHost.RunAsync( + serviceToken: null, + configureServices: services => + { + services.AddSingleton(registry); + services.AddScoped(_ => store); + }, + scenario: async channel => + { + RpcException exception = await Assert.ThrowsAsync( + () => PushAsync(channel, Request("текст", externalId: "3"), TenantA, ValidToken)); + Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); + }); + } + + /// + /// Metadata tenant-id отсутствует → UNAUTHENTICATED + /// + [Fact] + public async Task PushSource_WithoutTenantIdMetadata_IsUnauthenticated() + { + var store = new FakePipelineStore(); + var registry = new FakeTenantRegistry(Tenant(TenantA)); + + await RunAsync(registry, services => services.AddScoped(_ => store), async channel => + { + RpcException exception = await Assert.ThrowsAsync( + () => PushAsync(channel, Request("текст", externalId: "4"), null, ValidToken)); + Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); + }); + } + + // ─── Контекст и хелперы ───────────────────────────────────────────────── + + // Запись реестра тенанта (как строка public.tenants). + // id: Идентификатор тенанта. + private static TenantRecordDto Tenant(Guid id) => + new(id, Name: "tenant", Status: "active", CreatedAt: DateTimeOffset.UtcNow); + + // Поднимает хост ингресса: реестр тенантов + tenant-адаптеры сценария (общий харнесс). + // registry: Реестр тенантов сценария. + // registerTenantServices: Дополнительные tenant-scoped адаптеры сценария (IPipelineStore). + // scenario: Сценарий с gRPC-каналом. + private static Task RunAsync( + FakeTenantRegistry registry, + Action registerTenantServices, + Func scenario) + => TelegramIngressTestHost.RunAsync( + ValidToken, + services => + { + services.AddSingleton(registry); + registerTenantServices(services); + }, + scenario); + + // Вызывает PushSource с metadata сценария (deadline 10 с — контракт README). + // channel: Канал к хосту. + // request: Запрос PushSource. + // tenantId: Id тенанта в metadata (null — без заголовка tenant-id). + // tokenHeader: Значение metadata «service-token» (null — без заголовка). + private static Task PushAsync( + GrpcChannel channel, + PushSourceRequest request, + Guid? tenantId, + string? tokenHeader) + { + var client = new SourceIngressService.SourceIngressServiceClient(channel); + AsyncUnaryCall call = client.PushSourceAsync( + request, + new CallOptions(TelegramIngressTestHost.CallMetadata(tokenHeader, tenantId), deadline: Deadline())); + return call.ResponseAsync; + } + + // Запрос PushSource сценария (вид/оригинал/подпись/hue фиксированы). + // text: Текст сообщения. + // externalId: Id записи в источнике. + private static PushSourceRequest Request(string text, string externalId) + { + var request = new PushSourceRequest + { + Source = new SourceRefProto + { + Kind = Kind, + ExternalId = externalId, + DisplayName = DisplayName, + OriginRef = OriginRef, + ReceivedAt = ReceivedAtMs, + }, + Content = new SourceContentProto { Text = text }, + }; + request.Source.Extra["hue"] = Hue; + return request; + } + + // Deadline вызовов теста (контракт ингресса — 10 с). + private static DateTime Deadline() + => DateTime.UtcNow.AddSeconds(TelegramIngressTestHost.RpcDeadlineSeconds); +} diff --git a/src/core/tests/Deal.Tests.Unit/Support/SourceProtoMapperTests.cs b/src/core/tests/Deal.Tests.Unit/Support/SourceProtoMapperTests.cs new file mode 100644 index 0000000..2f0339c --- /dev/null +++ b/src/core/tests/Deal.Tests.Unit/Support/SourceProtoMapperTests.cs @@ -0,0 +1,211 @@ +using Deal.Api.Sources; +using Deal.Grpc.Sources; +using Deal.Modules.Cards.Application.Sources; + +namespace Deal.Tests.Unit.Support; + +/// +/// Тесты маппинга generic-контракта источника в доменную модель. +/// +public sealed class SourceProtoMapperTests +{ + /// + /// Заполненный запрос переносится в источник и содержимое целиком + /// + [Fact] + public void ToItem_FullyPopulatedRequest_MapsSourceAndContent() + { + var request = new PushSourceRequest + { + Source = new SourceRefProto + { + Kind = "telegram", + ExternalId = "42", + DisplayName = "Канал", + OriginRef = "d_1", + Author = "@author", + ReceivedAt = 1_700_000_000_000, + }, + Content = new SourceContentProto + { + Text = "Вакансия", + Html = "Вакансия", + Author = "Иван", + Subject = "Тема", + Links = { "https://a", "https://b" }, + }, + }; + request.Source.Extra["hue"] = "#a33"; + request.Content.Extra["lang"] = "ru"; + + SourceItem item = SourceProtoMapper.ToItem(request); + + Assert.Equal("telegram", item.Source.Kind); + Assert.Equal("42", item.Source.ExternalId); + Assert.Equal("Канал", item.Source.DisplayName); + Assert.Equal("d_1", item.Source.OriginRef); + Assert.Equal("@author", item.Source.Author); + Assert.Equal(DateTimeOffset.FromUnixTimeMilliseconds(1_700_000_000_000), item.Source.ReceivedAt); + Assert.Equal("#a33", item.Source.Extra!["hue"]); + + Assert.Equal("Вакансия", item.Content.Text); + Assert.Equal("Вакансия", item.Content.Html); + Assert.Equal("Иван", item.Content.Author); + Assert.Equal("Тема", item.Content.Subject); + Assert.Equal(["https://a", "https://b"], item.Content.Links); + Assert.Equal("ru", item.Content.Extra!["lang"]); + Assert.Empty(item.Content.Data); + } + + /// + /// Вложения и контакты переносятся со всеми необязательными полями + /// + [Fact] + public void ToItem_DataAndContacts_MapsAllFields() + { + var request = new PushSourceRequest + { + Source = new SourceRefProto { Kind = "file" }, + Content = new SourceContentProto + { + Data = + { + new DataRefProto + { + Id = "obj_1", + Ref = "storage://obj_1", + Kind = "image", + MimeType = "image/png", + FileName = "screen.png", + Size = 2048, + Width = 800, + Height = 600, + DurationSec = 1.5, + PreviewRef = "storage://obj_1/thumb", + Caption = "Подпись", + Order = 2, + }, + }, + Contacts = + { + new ContactRefProto + { + Kind = "phone", + Name = "Иван", + Phone = "+79990000000", + Email = "ivan@example.com", + Url = "https://example.com", + }, + }, + }, + }; + request.Content.Data[0].Meta["source"] = "upload"; + + SourceItem item = SourceProtoMapper.ToItem(request); + + DataRef data = Assert.Single(item.Content.Data); + Assert.Equal("obj_1", data.Id); + Assert.Equal("storage://obj_1", data.Ref); + Assert.Equal("image", data.Kind); + Assert.Equal("image/png", data.MimeType); + Assert.Equal("screen.png", data.FileName); + Assert.Equal(2048, data.Size); + Assert.Equal(800, data.Width); + Assert.Equal(600, data.Height); + Assert.Equal(1.5, data.DurationSec); + Assert.Equal("storage://obj_1/thumb", data.PreviewRef); + Assert.Equal("Подпись", data.Caption); + Assert.Equal(2, data.Order); + Assert.Equal("upload", data.Meta!["source"]); + + ContactRef contact = Assert.Single(item.Content.Contacts!); + Assert.Equal("phone", contact.Kind); + Assert.Equal("Иван", contact.Name); + Assert.Equal("+79990000000", contact.Phone); + Assert.Equal("ivan@example.com", contact.Email); + Assert.Equal("https://example.com", contact.Url); + } + + /// + /// Незаполненные необязательные поля источника/содержимого дают null + /// + [Fact] + public void ToItem_MissingOptionals_MapToNull() + { + var request = new PushSourceRequest + { + Source = new SourceRefProto { Kind = "local" }, + Content = new SourceContentProto { Text = "x" }, + }; + + SourceItem item = SourceProtoMapper.ToItem(request); + + Assert.Null(item.Source.ExternalId); + Assert.Null(item.Source.DisplayName); + Assert.Null(item.Source.OriginRef); + Assert.Null(item.Source.Author); + Assert.Null(item.Source.Extra); + Assert.Null(item.Content.Html); + Assert.Null(item.Content.Author); + Assert.Null(item.Content.Subject); + Assert.Null(item.Content.Links); + Assert.Null(item.Content.Contacts); + Assert.Null(item.Content.Extra); + } + + /// + /// Незаполненные необязательные поля вложения/контакта дают null + /// + [Fact] + public void ToItem_DataAndContactsWithoutOptionals_MapToNull() + { + var request = new PushSourceRequest + { + Source = new SourceRefProto { Kind = "file" }, + Content = new SourceContentProto + { + Data = { new DataRefProto { Id = "obj_1", Ref = "storage://obj_1" } }, + Contacts = { new ContactRefProto() }, + }, + }; + + SourceItem item = SourceProtoMapper.ToItem(request); + + DataRef data = Assert.Single(item.Content.Data); + Assert.Null(data.Kind); + Assert.Null(data.MimeType); + Assert.Null(data.FileName); + Assert.Null(data.Width); + Assert.Null(data.Height); + Assert.Null(data.DurationSec); + Assert.Null(data.PreviewRef); + Assert.Null(data.Caption); + Assert.Null(data.Order); + Assert.Null(data.Meta); + + ContactRef contact = Assert.Single(item.Content.Contacts!); + Assert.Null(contact.Kind); + Assert.Null(contact.Name); + Assert.Null(contact.Phone); + Assert.Null(contact.Email); + Assert.Null(contact.Url); + } + + /// + /// received_at=0 подменяется текущим временем получения + /// + [Fact] + public void ToItem_ZeroReceivedAt_UsesCurrentTime() + { + DateTimeOffset before = DateTimeOffset.UtcNow; + var request = new PushSourceRequest + { + Source = new SourceRefProto { Kind = "local" }, + Content = new SourceContentProto { Text = "x" }, + }; + + SourceItem item = SourceProtoMapper.ToItem(request); + + Assert.InRange(item.Source.ReceivedAt, before.AddSeconds(-1), DateTimeOffset.UtcNow.AddSeconds(1)); + } +} diff --git a/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressServiceTests.cs b/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressServiceTests.cs index 1e4a4be..fbf73ef 100644 --- a/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressServiceTests.cs +++ b/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressServiceTests.cs @@ -1,16 +1,12 @@ using System.Text.Json; using Deal.Api.Events; using Deal.Grpc.Telegram; -using Deal.Modules.Cards.Application.Sources; -using Deal.Modules.Pipeline.Application.Abstractions; -using Deal.Modules.Pipeline.Application.Models; using Deal.Modules.Settings.Application.Abstractions; using Deal.Modules.Settings.Application.Models; using Deal.Modules.Telegram.Application; using Deal.Modules.Tenants.Application.Abstractions; using Deal.Modules.Tenants.Application.Models; using Deal.Tests.Unit.Contracts; -using Deal.Tests.Unit.Modules.Kanban; using Deal.Tests.Unit.Modules.Settings; using Deal.Tests.Unit.Modules.Tenants; using Grpc.Core; @@ -20,7 +16,7 @@ using Microsoft.Extensions.DependencyInjection; namespace Deal.Tests.Unit.Support; /// -/// In-proc gRPC-тесты входящего потока telegram-service → ядро. +/// In-proc gRPC-тесты служебного потока telegram-service → ядро (SyncDialogs/ReportStatus). /// public sealed class TelegramIngressServiceTests { @@ -30,164 +26,6 @@ public sealed class TelegramIngressServiceTests // Тенант A сценариев (в реестре). private static readonly Guid TenantA = Guid.NewGuid(); - // Id диалога сценариев PushMessage. - private const string DialogId = "d_channel_100"; - - // ─── PushMessage: приём, дубль, несуществующий тенант ────────────────── - - /// - /// PushMessage кладёт строку очереди тенанта - /// - [Fact] - public async Task PushMessage_ValidTenant_EnqueuesQueueRowAndAccepts() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync( - registry, - services => services.AddScoped(_ => store), - async channel => - { - PushMessageReply reply = await PushAsync(channel, Request("Вакансия Python-разработчика", msgId: 7), TenantA, ValidToken); - - Assert.True(reply.Accepted); - Assert.False(reply.Duplicate); - - QueueItemDto row = Assert.Single(store.Queue); - Assert.Equal("telegram", row.Source.Kind); - Assert.Equal(DialogId, row.Source.OriginRef); - Assert.Equal("7", row.Source.ExternalId); - Assert.Equal("Вакансия Python-разработчика", row.Text); - Assert.Equal("Канал", row.Source.DisplayName); - Assert.Equal("#a33", row.Source.ResolveHue()); - }); - } - - /// - /// Повтор PushMessage того же dialogId+msgId — duplicate=true, очередь не растёт - /// - [Fact] - public async Task PushMessage_SameDialogAndMsgIdTwice_SecondIsDuplicateAndQueueNotGrown() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync( - registry, - services => services.AddScoped(_ => store), - async channel => - { - PushMessageReply first = await PushAsync(channel, Request("Сообщение канала", msgId: 42), TenantA, ValidToken); - PushMessageReply second = await PushAsync(channel, Request("Сообщение канала", msgId: 42), TenantA, ValidToken); - - Assert.True(first.Accepted); - Assert.True(second.Accepted); - Assert.True(second.Duplicate); - Assert.Single(store.Queue); - }); - } - - /// - /// PushMessage для несуществующего тенанта не падает - /// - [Fact] - public async Task PushMessage_UnknownTenant_NotAcceptedWithoutRpcError() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync( - registry, - services => services.AddScoped(_ => store), - async channel => - { - PushMessageReply reply = await PushAsync(channel, Request("Сообщение чужого тенанта"), Guid.NewGuid(), ValidToken); - - Assert.False(reply.Accepted); - Assert.False(reply.Duplicate); - Assert.Empty(store.Queue); - }); - } - - // ─── Интерцептор service-token и metadata tenant-id ───────────────────── - - /// - /// Запрос без metadata «service-token» → UNAUTHENTICATED. - /// - [Fact] - public async Task PushMessage_WithoutToken_IsUnauthenticated() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync(registry, services => services.AddScoped(_ => store), async channel => - { - RpcException exception = await Assert.ThrowsAsync( - () => PushAsync(channel, Request("текст"), TenantA, null)); - Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); - }); - } - - /// - /// Запрос с неверным токеном → UNAUTHENTICATED. - /// - [Fact] - public async Task PushMessage_WithWrongToken_IsUnauthenticated() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync(registry, services => services.AddScoped(_ => store), async channel => - { - RpcException exception = await Assert.ThrowsAsync( - () => PushAsync(channel, Request("текст"), TenantA, "wrong-token")); - Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); - }); - } - - /// - /// Fail-closed: DEAL_SERVICE_TOKEN не задан — RPC ингресса отклоняется даже с «каким-то» токеном. - /// - [Fact] - public async Task PushMessage_UnsetEnvToken_FailsClosed() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await TelegramIngressTestHost.RunAsync( - serviceToken: null, - configureServices: services => - { - services.AddSingleton(registry); - services.AddScoped(_ => new FakeSettingsStore()); - services.AddScoped(_ => store); - }, - scenario: async channel => - { - RpcException exception = await Assert.ThrowsAsync( - () => PushAsync(channel, Request("текст"), TenantA, ValidToken)); - Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); - }); - } - - /// - /// Metadata tenant-id отсутствует → UNAUTHENTICATED - /// - [Fact] - public async Task PushMessage_WithoutTenantIdMetadata_IsUnauthenticated() - { - var store = new FakePipelineStore(); - var registry = new FakeTenantRegistry(Tenant(TenantA)); - - await RunAsync(registry, services => services.AddScoped(_ => store), async channel => - { - RpcException exception = await Assert.ThrowsAsync( - () => PushAsync(channel, Request("текст"), null, ValidToken)); - Assert.Equal(StatusCode.Unauthenticated, exception.StatusCode); - }); - } - // ─── SyncDialogs и ReportStatus ───────────────────────────────────────── /// @@ -337,7 +175,7 @@ public sealed class TelegramIngressServiceTests // Поднимает хост ингресса: реестр тенантов + tenant-адаптеры сценария (общий харнесс). // registry: Реестр тенантов сценария. // registerTenantServices: Дополнительные tenant-scoped адаптеры сценария - // (IPipelineStore/ISettingsStore/брокер). + // (ISettingsStore/ITelegramStore/брокер). // scenario: Сценарий с gRPC-каналом. private static Task RunAsync( FakeTenantRegistry registry, @@ -352,24 +190,6 @@ public sealed class TelegramIngressServiceTests }, scenario); - // Вызывает PushMessage с metadata сценария (deadline 10 с — контракт README). - // channel: Канал к хосту. - // request: Запрос PushMessage. - // tenantId: Id тенанта в metadata (null — без заголовка tenant-id). - // tokenHeader: Значение metadata «service-token» (null — без заголовка). - private static Task PushAsync( - GrpcChannel channel, - PushMessageRequest request, - Guid? tenantId, - string? tokenHeader) - { - var client = new IngressService.IngressServiceClient(channel); - AsyncUnaryCall call = client.PushMessageAsync( - request, - new CallOptions(TelegramIngressTestHost.CallMetadata(tokenHeader, tenantId), deadline: Deadline())); - return call.ResponseAsync; - } - // Вызывает SyncDialogs с metadata сценария. // channel: Канал к хосту. // request: Запрос SyncDialogs. @@ -412,27 +232,6 @@ public sealed class TelegramIngressServiceTests return call.ResponseAsync; } - // Запрос PushMessage сценария (канальные поля фиксированы). - // text: Текст сообщения. - // msgId: Id сообщения в Telegram (null — без msg_id). - private static PushMessageRequest Request(string text, long? msgId = null) - { - var request = new PushMessageRequest - { - DialogId = DialogId, - ChannelName = "Канал", - ChannelHandle = "kanal_handle", - ChannelHue = "#a33", - Text = text, - }; - if (msgId is not null) - { - request.MsgId = msgId.Value; - } - - return request; - } - // Deadline вызовов теста (контракт ингресса — 10 с). private static DateTime Deadline() => DateTime.UtcNow.AddSeconds(TelegramIngressTestHost.RpcDeadlineSeconds); diff --git a/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressTestHost.cs b/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressTestHost.cs index f90b691..804cdaf 100644 --- a/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressTestHost.cs +++ b/src/core/tests/Deal.Tests.Unit/Support/TelegramIngressTestHost.cs @@ -1,6 +1,7 @@ using System.Net; using Deal.Api.Configuration; using Deal.Api.Events; +using Deal.Api.Sources; using Deal.Api.Telegram; using Deal.Contracts.Integrations.Abstractions; using Deal.Infrastructure.Data; @@ -75,6 +76,8 @@ internal static class TelegramIngressTestHost builder.Services.AddSingleton(); builder.Services.AddScoped(); builder.Services.AddScoped(); + builder.Services.AddScoped(); + builder.Services.AddScoped(); builder.Services.AddGrpc(grpc => { grpc.Interceptors.Add(); @@ -105,6 +108,7 @@ internal static class TelegramIngressTestHost app = builder.Build(); app.MapGrpcService(); + app.MapGrpcService(); if (rateLimitOptions is { Enabled: true }) { app.MapGrpcHealthChecksService(); diff --git a/src/core/tests/Deal.Tests.Unit/Support/TelegramSourceIngestObserverTests.cs b/src/core/tests/Deal.Tests.Unit/Support/TelegramSourceIngestObserverTests.cs new file mode 100644 index 0000000..cb3165f --- /dev/null +++ b/src/core/tests/Deal.Tests.Unit/Support/TelegramSourceIngestObserverTests.cs @@ -0,0 +1,109 @@ +using Deal.Modules.Cards.Application.Sources; +using Deal.Modules.Telegram.Application; +using Deal.Tests.Unit.Contracts; +using Deal.Tests.Unit.Modules.Settings; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Deal.Tests.Unit.Support; + +/// +/// Тесты наблюдателя, сохраняющего превью принятых сообщений. +/// +public sealed class TelegramSourceIngestObserverTests +{ + // Id диалога-оригинала сценариев. + private const string DialogId = "d_1"; + + /// + /// Принятая запись известного вида сохраняет превью сообщения + /// + [Fact] + public async Task OnIngestedAsync_KnownKind_SavesPreview() + { + (TelegramSourceIngestObserver observer, FakeTelegramStore store) = Create(); + store.Seed(Dialog()); + DateTimeOffset receivedAt = DateTimeOffset.UtcNow.AddMinutes(-1); + + await observer.OnIngestedAsync(Item("telegram", DialogId, "7", "Нужен Python", receivedAt), CancellationToken.None); + + FakeTelegramMessageRow preview = Assert.Single(store.Messages); + Assert.Equal("m_d_1_7", preview.Id); + Assert.Equal("Нужен Python", preview.Text); + Assert.Equal(receivedAt, preview.MsgAt); + Assert.Equal("Нужен Python", store.Dialogs[0].LastText); + Assert.Equal(receivedAt, store.Dialogs[0].LastAt); + } + + /// + /// Чужой вид источника не трогает хранилище + /// + [Fact] + public async Task OnIngestedAsync_ForeignKind_DoesNotTouchStore() + { + (TelegramSourceIngestObserver observer, FakeTelegramStore store) = Create(); + store.Seed(Dialog()); + + await observer.OnIngestedAsync(Item("avito", DialogId, "7", "Текст", DateTimeOffset.UtcNow), CancellationToken.None); + + Assert.Empty(store.Messages); + Assert.Equal(string.Empty, store.Dialogs[0].LastText); + Assert.Null(store.Dialogs[0].LastAt); + } + + /// + /// Запись без оригинала не сохраняет превью + /// + [Fact] + public async Task OnIngestedAsync_NoOriginRef_DoesNotTouchStore() + { + (TelegramSourceIngestObserver observer, FakeTelegramStore store) = Create(); + store.Seed(Dialog()); + + await observer.OnIngestedAsync(Item("telegram", null, "7", "Текст", DateTimeOffset.UtcNow), CancellationToken.None); + + Assert.Empty(store.Messages); + Assert.Equal(string.Empty, store.Dialogs[0].LastText); + } + + // Собирает наблюдателя на фейковом хранилище диалогов. + private static (TelegramSourceIngestObserver Observer, FakeTelegramStore Store) Create() + { + var store = new FakeTelegramStore(); + var dialogs = new DialogsService( + store, + new FakeSettingsStore(), + new FakeTelegramGateway(), + NullLogger.Instance); + var observer = new TelegramSourceIngestObserver(dialogs, NullLogger.Instance); + return (observer, store); + } + + // Строка каталога для сценариев наблюдателя. + private static FakeTelegramDialogRow Dialog() => new( + DialogId, + Name: "Канал", + Handle: "kanal_handle", + Kind: "channel", + Hue: "#a33", + Monitor: true, + LastText: string.Empty, + LastAt: null, + Backfilled: false, + UpdatedAt: DateTimeOffset.UtcNow); + + // Запись источника сценариев. + // kind: Вид источника. + // originRef: Оригинал (диалог); null — без оригинала. + // externalId: Id сообщения. + private static SourceItem Item(string kind, string? originRef, string externalId, string text, DateTimeOffset receivedAt) => new() + { + Source = new SourceRef + { + Kind = kind, + OriginRef = originRef, + ExternalId = externalId, + ReceivedAt = receivedAt, + }, + Content = new SourceContent { Text = text }, + }; +} diff --git a/src/telegram-service/Deal.Telegram.Tests/Grpc/CoreIngressClientTests.cs b/src/telegram-service/Deal.Telegram.Tests/Grpc/CoreIngressClientTests.cs index e3c9f39..cca9d0d 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Grpc/CoreIngressClientTests.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Grpc/CoreIngressClientTests.cs @@ -1,3 +1,4 @@ +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Deal.Telegram.Core; using Deal.Telegram.Sessions; @@ -16,31 +17,35 @@ public sealed class CoreIngressClientTests private const string ServiceToken = "deal-test-token"; /// - /// PushMessage: тело уходит, metadata tenant-id/service-token на месте; ответ accepted. + /// PushSource: тело уходит, metadata tenant-id/service-token на месте; ответ accepted. /// [Fact] - public async Task PushMessage_SendsMessageWithMetadata_AndReturnsReply() + public async Task PushSource_SendsSourceWithMetadata_AndReturnsReply() { (string endpoint, RecordingIngressService server, WebApplication app) = await FakeIngressServer.StartAsync(); try { CoreIngressClient client = CreateClient(endpoint); - PushMessageReply reply = await client.PushMessageAsync( + PushSourceReply reply = await client.PushSourceAsync( TenantId, - new PushMessageRequest + new PushSourceRequest { - DialogId = "-1001234567890", - ChannelName = "IT Канал", - Text = "сообщение", - MsgId = 42, + Source = new SourceRefProto + { + Kind = "telegram", + ExternalId = "42", + OriginRef = "-1001234567890", + DisplayName = "IT Канал", + }, + Content = new SourceContentProto { Text = "сообщение" }, }, CancellationToken.None); Assert.True(reply.Accepted); - PushMessageRequest sent = Assert.Single(server.Pushes); - Assert.Equal("-1001234567890", sent.DialogId); - Assert.Equal("сообщение", sent.Text); + PushSourceRequest sent = Assert.Single(server.Pushes); + Assert.Equal("-1001234567890", sent.Source.OriginRef); + Assert.Equal("сообщение", sent.Content.Text); Assert.Equal(TenantId, Assert.Single(server.Tenants)); Assert.Equal(ServiceToken, Assert.Single(server.Tokens)); } @@ -83,13 +88,20 @@ public sealed class CoreIngressClientTests /// Ядро недоступно → SessionException UNAVAILABLE «Ядро недоступно…» /// [Fact] - public async Task PushMessage_UnreachableCore_ThrowsSessionException() + public async Task PushSource_UnreachableCore_ThrowsSessionException() { // Порт без сервера: соединение отклоняется — вызов падает до deadline (RpcTimeout 15 с). CoreIngressClient client = CreateClient("http://127.0.0.1:1"); SessionException exception = await Assert.ThrowsAsync( - () => client.PushMessageAsync(TenantId, new PushMessageRequest { DialogId = "-1001", Text = "x" }, CancellationToken.None)); + () => client.PushSourceAsync( + TenantId, + new PushSourceRequest + { + Source = new SourceRefProto { Kind = "telegram", OriginRef = "-1001" }, + Content = new SourceContentProto { Text = "x" }, + }, + CancellationToken.None)); Assert.Equal(StatusCode.Unavailable, exception.Code); Assert.Equal(SessionErrorMessages.IngressUnavailable, exception.Message); diff --git a/src/telegram-service/Deal.Telegram.Tests/Grpc/DialogRpcTests.cs b/src/telegram-service/Deal.Telegram.Tests/Grpc/DialogRpcTests.cs index 531f8d7..b342aea 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Grpc/DialogRpcTests.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Grpc/DialogRpcTests.cs @@ -143,7 +143,7 @@ public sealed class DialogRpcTests Assert.Equal(2, reply.Processed); Assert.Equal(2, ingressServer.Pushes.Count); - Assert.Equal([9, 10], ingressServer.Pushes.Select(push => push.MsgId)); + Assert.Equal(["9", "10"], ingressServer.Pushes.Select(push => push.Source.ExternalId)); Assert.All(ingressServer.Tenants, tenant => Assert.Equal(TenantId, tenant)); Assert.Equal([DialogA], factory.CreatedClients.Single().MarkedReadDialogs); }); diff --git a/src/telegram-service/Deal.Telegram.Tests/Grpc/FakeIngressServer.cs b/src/telegram-service/Deal.Telegram.Tests/Grpc/FakeIngressServer.cs index 1b330d4..5ad0d71 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Grpc/FakeIngressServer.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Grpc/FakeIngressServer.cs @@ -1,5 +1,6 @@ using System.Net; using System.Net.Sockets; +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Grpc.Core; using Microsoft.AspNetCore.Builder; @@ -9,12 +10,14 @@ using Microsoft.Extensions.DependencyInjection; namespace Deal.Telegram.Tests.Grpc; -internal sealed class RecordingIngressService : IngressService.IngressServiceBase +// Общее состояние записанных вызовов: PushSource (SourceIngressService) и SyncDialogs (IngressService) +// хостятся на одном фейк-сервере, но лежат в разных protobuf-сервисах. +internal sealed class IngressRecorder { /// - /// Полученные PushMessage, в порядке вызовов. + /// Полученные PushSource, в порядке вызовов. /// - public List Pushes { get; } = []; + public List Pushes { get; } = []; /// /// Полученные SyncDialogs, в порядке вызовов. @@ -36,33 +39,87 @@ internal sealed class RecordingIngressService : IngressService.IngressServiceBas /// public List MonitoredIds { get; set; } = []; - /// - public override Task PushMessage(PushMessageRequest request, ServerCallContext context) - { - Record(context); - Pushes.Add(request); - return Task.FromResult(new PushMessageReply { Accepted = true, Duplicate = false }); - } - - /// - public override Task SyncDialogs(SyncDialogsRequest request, ServerCallContext context) - { - Record(context); - Syncs.Add(request); - var reply = new SyncDialogsReply(); - reply.MonitoredIds.AddRange(MonitoredIds); - return Task.FromResult(reply); - } - - private void Record(ServerCallContext context) + /// + /// Записывает metadata вызова + /// + /// Контекст серверного вызова. + public void Record(ServerCallContext context) { Tenants.Add(context.RequestHeaders.GetValue("tenant-id") ?? string.Empty); Tokens.Add(context.RequestHeaders.GetValue("service-token") ?? string.Empty); } } -// Поднимает in-proc gRPC-сервер IngressService на эфемерном порту (без service-token-интерцептора: -// тест проверяет, что клиент шлёт токен, а не что сервер его принимает). +internal sealed class RecordingIngressService : SourceIngressService.SourceIngressServiceBase +{ + private readonly IngressRecorder _recorder; + + public RecordingIngressService(IngressRecorder recorder) + { + _recorder = recorder; + } + + /// + /// Полученные PushSource, в порядке вызовов. + /// + public List Pushes => _recorder.Pushes; + + /// + /// Полученные SyncDialogs, в порядке вызовов. + /// + public List Syncs => _recorder.Syncs; + + /// + /// Значения tenant-id полученных вызовов + /// + public List Tenants => _recorder.Tenants; + + /// + /// Значения service-token полученных вызовов + /// + public List Tokens => _recorder.Tokens; + + /// + /// Monitored-набор ответа SyncDialogs + /// + public List MonitoredIds + { + get => _recorder.MonitoredIds; + set => _recorder.MonitoredIds = value; + } + + /// + public override Task PushSource(PushSourceRequest request, ServerCallContext context) + { + _recorder.Record(context); + _recorder.Pushes.Add(request); + return Task.FromResult(new PushSourceReply { Accepted = true, Duplicate = false }); + } +} + +internal sealed class RecordingSyncService : IngressService.IngressServiceBase +{ + private readonly IngressRecorder _recorder; + + public RecordingSyncService(IngressRecorder recorder) + { + _recorder = recorder; + } + + /// + public override Task SyncDialogs(SyncDialogsRequest request, ServerCallContext context) + { + _recorder.Record(context); + _recorder.Syncs.Add(request); + var reply = new SyncDialogsReply(); + reply.MonitoredIds.AddRange(_recorder.MonitoredIds); + return Task.FromResult(reply); + } +} + +// Поднимает in-proc gRPC-сервер исходящих контрактов (SourceIngressService/IngressService) на +// эфемерном порту (без service-token-интерцептора: тест проверяет, что клиент шлёт токен, +// а не что сервер его принимает). internal static class FakeIngressServer { /// @@ -70,17 +127,21 @@ internal static class FakeIngressServer /// public static async Task<(string Endpoint, RecordingIngressService Server, WebApplication App)> StartAsync() { - var server = new RecordingIngressService(); + var recorder = new IngressRecorder(); + var server = new RecordingIngressService(recorder); + var syncServer = new RecordingSyncService(recorder); int port = FreeTcpPort(); WebApplicationBuilder builder = WebApplication.CreateBuilder(); builder.WebHost.ConfigureKestrel(kestrel => kestrel.Listen(IPAddress.Loopback, port, listen => listen.Protocols = HttpProtocols.Http2)); builder.Services.AddSingleton(server); + builder.Services.AddSingleton(syncServer); builder.Services.AddGrpc(); WebApplication app = builder.Build(); app.MapGrpcService(); + app.MapGrpcService(); await app.StartAsync(); return ($"http://127.0.0.1:{port}", server, app); diff --git a/src/telegram-service/Deal.Telegram.Tests/Grpc/RealtimeListenerTests.cs b/src/telegram-service/Deal.Telegram.Tests/Grpc/RealtimeListenerTests.cs index d410484..47a92d7 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Grpc/RealtimeListenerTests.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Grpc/RealtimeListenerTests.cs @@ -1,4 +1,4 @@ -using Deal.Grpc.Telegram; +using Deal.Grpc.Sources; using Deal.Telegram.Dialogs; using Deal.Telegram.Sessions; using Deal.Telegram.Telegram; @@ -35,9 +35,9 @@ public sealed class RealtimeListenerTests await harness.Client.RaiseMessageAsync( new TelegramMessage(DialogId, 101, "свежее сообщение", 1_700_000_000_101, "IT Канал", "it_channel")); - PushMessageRequest push = Assert.Single(ingress.Pushes).Message; - Assert.Equal(DialogId, push.DialogId); - Assert.Equal("свежее сообщение", push.Text); + PushSourceRequest push = Assert.Single(ingress.Pushes).Request; + Assert.Equal(DialogId, push.Source.OriginRef); + Assert.Equal("свежее сообщение", push.Content.Text); Assert.Equal(DialogId, Assert.Single(harness.Client.MarkedReadDialogs)); } finally @@ -73,7 +73,7 @@ public sealed class RealtimeListenerTests } /// - /// Сбой PushMessage не роняет realtime + /// Сбой PushSource не роняет realtime /// [Fact] public async Task Listener_PushFails_NoMarkRead() diff --git a/src/telegram-service/Deal.Telegram.Tests/Grpc/TestDoubles.cs b/src/telegram-service/Deal.Telegram.Tests/Grpc/TestDoubles.cs index 73aaf73..5b0b514 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Grpc/TestDoubles.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Grpc/TestDoubles.cs @@ -1,3 +1,4 @@ +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Deal.Telegram.Core; using Deal.Telegram.Dialogs; @@ -7,9 +8,9 @@ namespace Deal.Telegram.Tests.Grpc; internal sealed class FakeIngress : ICoreIngressClient { /// - /// Отправленные PushMessage + /// Отправленные PushSource /// - public List<(string TenantId, PushMessageRequest Message)> Pushes { get; } = []; + public List<(string TenantId, PushSourceRequest Request)> Pushes { get; } = []; /// /// Синхронизации каталога @@ -22,7 +23,7 @@ internal sealed class FakeIngress : ICoreIngressClient public List MonitoredIdsToReturn { get; set; } = []; /// - /// Ошибка PushMessageAsync + /// Ошибка PushSourceAsync /// public Exception? PushError { get; set; } @@ -32,19 +33,19 @@ internal sealed class FakeIngress : ICoreIngressClient public Exception? SyncError { get; set; } /// - /// Текст последнего PushMessage + /// Текст последнего PushSource /// - public string? LastPushText => Pushes.Count == 0 ? null : Pushes[^1].Message.Text; + public string? LastPushText => Pushes.Count == 0 ? null : Pushes[^1].Request.Content.Text; /// /// Id последнего отправленного диалога. /// - public string? LastPushDialogId => Pushes.Count == 0 ? null : Pushes[^1].Message.DialogId; + public string? LastPushDialogId => Pushes.Count == 0 ? null : Pushes[^1].Request.Source.OriginRef; /// - public Task PushMessageAsync( + public Task PushSourceAsync( string tenantId, - PushMessageRequest message, + PushSourceRequest request, CancellationToken cancellationToken) { if (PushError is not null) @@ -52,8 +53,8 @@ internal sealed class FakeIngress : ICoreIngressClient throw PushError; } - Pushes.Add((tenantId, message)); - return Task.FromResult(new PushMessageReply { Accepted = true, Duplicate = false }); + Pushes.Add((tenantId, request)); + return Task.FromResult(new PushSourceReply { Accepted = true, Duplicate = false }); } /// diff --git a/src/telegram-service/Deal.Telegram.Tests/Telegram/BackfillServiceTests.cs b/src/telegram-service/Deal.Telegram.Tests/Telegram/BackfillServiceTests.cs index 6ff054e..e602d4c 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Telegram/BackfillServiceTests.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Telegram/BackfillServiceTests.cs @@ -35,13 +35,12 @@ public sealed class BackfillServiceTests Assert.Equal(3, processed); // От старых к новым (в fake список хранится от новых к старым, как get_messages). - Assert.Equal([28, 29, 30], ingress.Pushes.Select(push => push.Message.MsgId)); + Assert.Equal(["28", "29", "30"], ingress.Pushes.Select(push => push.Request.Source.ExternalId)); Assert.Equal(TenantId, ingress.Pushes[0].TenantId); - Assert.Equal(DialogId, ingress.Pushes[0].Message.DialogId); - Assert.Equal(ChannelName, ingress.Pushes[0].Message.ChannelName); - Assert.Equal(ChannelHandle, ingress.Pushes[0].Message.ChannelHandle); - Assert.Equal(DialogHue.Compute(DialogId, ChannelName), ingress.Pushes[0].Message.ChannelHue); - Assert.Equal("новое сообщение", ingress.Pushes[2].Message.Text); + Assert.Equal(DialogId, ingress.Pushes[0].Request.Source.OriginRef); + Assert.Equal(ChannelName, ingress.Pushes[0].Request.Source.DisplayName); + Assert.Equal(DialogHue.Compute(DialogId, ChannelName), ingress.Pushes[0].Request.Source.Extra["hue"]); + Assert.Equal("новое сообщение", ingress.Pushes[2].Request.Content.Text); Assert.Equal(3, pacer.Waits.Count); Assert.All(pacer.Waits, wait => { @@ -114,7 +113,7 @@ public sealed class BackfillServiceTests } /// - /// Сбой PushMessage прерывает backfill без read-ack + /// Сбой PushSource прерывает backfill без read-ack /// [Fact] public async Task Backfill_PushFails_Throws_WithoutMarkRead() diff --git a/src/telegram-service/Deal.Telegram.Tests/Telegram/RealtimeSweepTests.cs b/src/telegram-service/Deal.Telegram.Tests/Telegram/RealtimeSweepTests.cs index 5f27439..4c5d299 100644 --- a/src/telegram-service/Deal.Telegram.Tests/Telegram/RealtimeSweepTests.cs +++ b/src/telegram-service/Deal.Telegram.Tests/Telegram/RealtimeSweepTests.cs @@ -47,8 +47,8 @@ public sealed class RealtimeSweepTests Assert.True(catalog.IsMonitored(TenantId, DialogId)); Assert.False(catalog.IsMonitored(TenantId, OtherDialogId)); // Дочитаны только сообщения мониторящегося диалога, от старых к новым. - Assert.Equal([28, 29, 30], ingress.Pushes.Select(push => push.Message.MsgId)); - Assert.Equal(DialogId, ingress.Pushes[0].Message.DialogId); + Assert.Equal(["28", "29", "30"], ingress.Pushes.Select(push => push.Request.Source.ExternalId)); + Assert.Equal(DialogId, ingress.Pushes[0].Request.Source.OriginRef); Assert.Equal([DialogId], harness.Client.MarkedReadDialogs); } finally diff --git a/src/telegram-service/Deal.Telegram/Core/CoreIngressClient.cs b/src/telegram-service/Deal.Telegram/Core/CoreIngressClient.cs index 4fd81f4..6b5f350 100644 --- a/src/telegram-service/Deal.Telegram/Core/CoreIngressClient.cs +++ b/src/telegram-service/Deal.Telegram/Core/CoreIngressClient.cs @@ -2,6 +2,7 @@ using Deal.Grpc.Hosting.Interceptors; using Deal.Grpc.Hosting.Models; using Deal.Grpc.Hosting.Options; using Deal.Grpc.Hosting.Services; +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Deal.Telegram.Sessions; using Grpc.Core; @@ -22,7 +23,8 @@ public sealed class CoreIngressClient : ICoreIngressClient private readonly ILogger _logger; private readonly MtlsCertificates? _mtlsCertificates; private readonly object _channelGate = new(); - private IngressService.IngressServiceClient? _client; + private IngressService.IngressServiceClient? _ingressClient; + private SourceIngressService.SourceIngressServiceClient? _sourceClient; private GrpcChannel? _channel; /// @@ -42,19 +44,19 @@ public sealed class CoreIngressClient : ICoreIngressClient } /// - public async Task PushMessageAsync( + public async Task PushSourceAsync( string tenantId, - PushMessageRequest message, + PushSourceRequest request, CancellationToken cancellationToken) { - IngressService.IngressServiceClient client = GetClient(); + SourceIngressService.SourceIngressServiceClient client = GetSourceClient(); try { - return await client.PushMessageAsync(message, CallOptions(tenantId)).ResponseAsync.WaitAsync(cancellationToken).ConfigureAwait(false); + return await client.PushSourceAsync(request, CallOptions(tenantId)).ResponseAsync.WaitAsync(cancellationToken).ConfigureAwait(false); } catch (Exception exception) when (exception is not OperationCanceledException) { - throw Fail("PushMessage", exception); + throw Fail("PushSource", exception); } } @@ -64,7 +66,7 @@ public sealed class CoreIngressClient : ICoreIngressClient IReadOnlyList entries, CancellationToken cancellationToken) { - IngressService.IngressServiceClient client = GetClient(); + IngressService.IngressServiceClient client = GetIngressClient(); var request = new SyncDialogsRequest(); request.Entries.AddRange(entries); try @@ -78,31 +80,44 @@ public sealed class CoreIngressClient : ICoreIngressClient } } - // Лениво создаёт gRPC-канал к ядру (адрес неизменен на время жизни процесса). - private IngressService.IngressServiceClient GetClient() + // Клиент SourceIngress (PushSource); канал к ядру создаётся лениво. + private SourceIngressService.SourceIngressServiceClient GetSourceClient() { lock (_channelGate) { - if (_client is null) - { - if (_mtlsCertificates is not null) - { - _channel = GrpcChannel.ForAddress( - _options.IngressEndpoint, - new GrpcChannelOptions { HttpHandler = _mtlsCertificates.CreateClientHttpHandler() }); - } - else - { - _channel = GrpcChannel.ForAddress(_options.IngressEndpoint); - } - - _client = new IngressService.IngressServiceClient(_channel); - } - - return _client; + EnsureChannelLocked(); + return _sourceClient!; } } + // Клиент Ingress (SyncDialogs); канал к ядру создаётся лениво. + private IngressService.IngressServiceClient GetIngressClient() + { + lock (_channelGate) + { + EnsureChannelLocked(); + return _ingressClient!; + } + } + + // Лениво создаёт gRPC-канал к ядру (адрес неизменен на время жизни процесса). + private void EnsureChannelLocked() + { + if (_channel is not null) + { + return; + } + + _channel = _mtlsCertificates is not null + ? GrpcChannel.ForAddress( + _options.IngressEndpoint, + new GrpcChannelOptions { HttpHandler = _mtlsCertificates.CreateClientHttpHandler() }) + : GrpcChannel.ForAddress(_options.IngressEndpoint); + + _ingressClient = new IngressService.IngressServiceClient(_channel); + _sourceClient = new SourceIngressService.SourceIngressServiceClient(_channel); + } + private CallOptions CallOptions(string tenantId) { var metadata = new Metadata diff --git a/src/telegram-service/Deal.Telegram/Core/ICoreIngressClient.cs b/src/telegram-service/Deal.Telegram/Core/ICoreIngressClient.cs index 76207a4..8d586db 100644 --- a/src/telegram-service/Deal.Telegram/Core/ICoreIngressClient.cs +++ b/src/telegram-service/Deal.Telegram/Core/ICoreIngressClient.cs @@ -1,3 +1,4 @@ +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; namespace Deal.Telegram.Core; @@ -8,15 +9,15 @@ namespace Deal.Telegram.Core; public interface ICoreIngressClient { /// - /// Отправляет сообщение диалога в ядро. + /// Отправляет запись источника в ядро. /// /// Id тенанта. - /// Сообщение в контракте PushMessageRequest. + /// Запись в контракте PushSourceRequest. /// Отмена вызова. - /// Ответ ядра (accepted/duplicate — дубль dialog+msgId в очереди не растёт). - public Task PushMessageAsync( + /// Ответ ядра (accepted/duplicate — дубль источника в очереди не растёт). + public Task PushSourceAsync( string tenantId, - PushMessageRequest message, + PushSourceRequest request, CancellationToken cancellationToken); /// diff --git a/src/telegram-service/Deal.Telegram/Dialogs/BackfillService.cs b/src/telegram-service/Deal.Telegram/Dialogs/BackfillService.cs index c0ae0fd..1fe0758 100644 --- a/src/telegram-service/Deal.Telegram/Dialogs/BackfillService.cs +++ b/src/telegram-service/Deal.Telegram/Dialogs/BackfillService.cs @@ -1,5 +1,5 @@ using System.Collections.Concurrent; -using Deal.Grpc.Telegram; +using Deal.Grpc.Sources; using Deal.Telegram.Core; using Deal.Telegram.Sessions; using Deal.Telegram.Telegram; @@ -54,7 +54,7 @@ public sealed class BackfillService /// Создаёт службу backfill'а. /// /// Пул сессий тенантов (read-операции диалогов). - /// Канал в ядро (PushMessage). + /// Канал в ядро (PushSource). /// Анти-бан-паузы (реальный — случайные, тесты — фейк). /// Логгер. public BackfillService( @@ -168,8 +168,8 @@ public sealed class BackfillService int processed = 0; foreach (TelegramMessage message in messages.Reverse()) { - PushMessageReply reply = await _ingress - .PushMessageAsync(tenantId, DialogProtoMapper.ToPushRequest(message), cancellationToken) + PushSourceReply reply = await _ingress + .PushSourceAsync(tenantId, DialogProtoMapper.ToSourceRequest(message), cancellationToken) .ConfigureAwait(false); if (reply.Accepted && !reply.Duplicate) { diff --git a/src/telegram-service/Deal.Telegram/Dialogs/DialogProtoMapper.cs b/src/telegram-service/Deal.Telegram/Dialogs/DialogProtoMapper.cs index bd1d0b2..a99a971 100644 --- a/src/telegram-service/Deal.Telegram/Dialogs/DialogProtoMapper.cs +++ b/src/telegram-service/Deal.Telegram/Dialogs/DialogProtoMapper.cs @@ -1,4 +1,5 @@ using System.Globalization; +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Deal.Telegram.Telegram; @@ -9,6 +10,10 @@ namespace Deal.Telegram.Dialogs; /// public static class DialogProtoMapper { + private const string TelegramKind = "telegram"; + + private const string HueKey = "hue"; + /// /// Нейтральный диалог → DialogEntry контракта /// @@ -38,19 +43,27 @@ public static class DialogProtoMapper }; /// - /// Нейтральное сообщение → PushMessageRequest ингресса. + /// Нейтральное сообщение → PushSourceRequest ингресса. /// /// Сообщение диалога (непустой текст). - /// Запрос Ingress.PushMessage с канальными полями и дубль-гвардом msg_id. - public static PushMessageRequest ToPushRequest(TelegramMessage message) + /// Запрос SourceIngress.PushSource с источником telegram и содержимым. + public static PushSourceRequest ToSourceRequest(TelegramMessage message) => new() { - DialogId = message.DialogId, - ChannelName = message.DialogName, - ChannelHandle = message.DialogHandle, - ChannelHue = DialogHue.Compute(message.DialogId, message.DialogName), - MsgId = message.Id, - Text = message.Text, - MsgAt = message.DateMs, + Source = new SourceRefProto + { + Kind = TelegramKind, + ExternalId = message.Id.ToString(CultureInfo.InvariantCulture), + OriginRef = message.DialogId, + DisplayName = message.DialogName, + Author = message.DialogName, + ReceivedAt = message.DateMs, + Extra = { [HueKey] = DialogHue.Compute(message.DialogId, message.DialogName) }, + }, + Content = new SourceContentProto + { + Text = message.Text, + Author = message.DialogName, + }, }; } diff --git a/src/telegram-service/Deal.Telegram/Dialogs/RealtimeListener.cs b/src/telegram-service/Deal.Telegram/Dialogs/RealtimeListener.cs index ee17863..cbf252e 100644 --- a/src/telegram-service/Deal.Telegram/Dialogs/RealtimeListener.cs +++ b/src/telegram-service/Deal.Telegram/Dialogs/RealtimeListener.cs @@ -1,4 +1,4 @@ -using Deal.Grpc.Telegram; +using Deal.Grpc.Sources; using Deal.Telegram.Core; using Deal.Telegram.Sessions; using Deal.Telegram.Telegram; @@ -20,7 +20,7 @@ public sealed class RealtimeListener /// /// Ready-сессия тенанта (события сообщений её клиента). /// Зеркало мониторинга. - /// Канал в ядро (PushMessage). + /// Канал в ядро (PushSource). /// Логгер. public RealtimeListener( TenantSession session, @@ -52,7 +52,7 @@ public sealed class RealtimeListener _session.SetListenerActive(false); } - // Обработчик входящего сообщения: фильтр по зеркалу → PushMessage в ядро → mark-as-read. + // Обработчик входящего сообщения: фильтр по зеркалу → PushSource в ядро → mark-as-read. // Сбой отправки не роняет realtime: сообщение остаётся непрочитанным и догоняется realtime_sweep. // message: Входящее сообщение аккаунта. private async Task OnMessageReceivedAsync(TelegramMessage message) @@ -64,8 +64,8 @@ public sealed class RealtimeListener return; } - PushMessageReply reply = await _ingress - .PushMessageAsync(_session.TenantId, DialogProtoMapper.ToPushRequest(message), CancellationToken.None) + PushSourceReply reply = await _ingress + .PushSourceAsync(_session.TenantId, DialogProtoMapper.ToSourceRequest(message), CancellationToken.None) .ConfigureAwait(false); if (reply.Accepted) { diff --git a/src/telegram-service/Deal.Telegram/Dialogs/RealtimeSweep.cs b/src/telegram-service/Deal.Telegram/Dialogs/RealtimeSweep.cs index db8b4b0..616a974 100644 --- a/src/telegram-service/Deal.Telegram/Dialogs/RealtimeSweep.cs +++ b/src/telegram-service/Deal.Telegram/Dialogs/RealtimeSweep.cs @@ -1,3 +1,4 @@ +using Deal.Grpc.Sources; using Deal.Grpc.Telegram; using Deal.Telegram.Core; using Deal.Telegram.Sessions; @@ -40,7 +41,7 @@ public sealed class RealtimeSweep /// /// Пул сессий тенантов. /// Зеркало каталога/мониторинга (актуализируется ответом SyncDialogs). - /// Канал в ядро (PushMessage/SyncDialogs). + /// Канал в ядро (PushSource/SyncDialogs). /// Логгер. public RealtimeSweep( SessionFarm sessionFarm, @@ -145,8 +146,8 @@ public sealed class RealtimeSweep int added = 0; foreach (TelegramMessage message in messages.Reverse()) { - PushMessageReply reply = await _ingress - .PushMessageAsync(tenantId, DialogProtoMapper.ToPushRequest(message), cancellationToken) + PushSourceReply reply = await _ingress + .PushSourceAsync(tenantId, DialogProtoMapper.ToSourceRequest(message), cancellationToken) .ConfigureAwait(false); if (reply.Accepted && !reply.Duplicate) { diff --git a/src/telegram-service/Deal.Telegram/Hosting/RealtimeMonitorService.cs b/src/telegram-service/Deal.Telegram/Hosting/RealtimeMonitorService.cs index f428678..bbfbf09 100644 --- a/src/telegram-service/Deal.Telegram/Hosting/RealtimeMonitorService.cs +++ b/src/telegram-service/Deal.Telegram/Hosting/RealtimeMonitorService.cs @@ -27,7 +27,7 @@ public sealed class RealtimeMonitorService : BackgroundService /// /// Пул сессий тенантов. /// Зеркало мониторинга (фильтр сообщений). - /// Канал в ядро (PushMessage). + /// Канал в ядро (PushSource). /// Фабрика логгеров (логгеры listener'ов). public RealtimeMonitorService( SessionFarm sessionFarm, diff --git a/src/telegram-service/Deal.Telegram/Sessions/SessionErrorMessages.cs b/src/telegram-service/Deal.Telegram/Sessions/SessionErrorMessages.cs index 792cf34..7e1686a 100644 --- a/src/telegram-service/Deal.Telegram/Sessions/SessionErrorMessages.cs +++ b/src/telegram-service/Deal.Telegram/Sessions/SessionErrorMessages.cs @@ -53,7 +53,7 @@ public static class SessionErrorMessages public const string TelegramUnavailable = "Telegram недоступен — повторите попытку позже"; /// - /// Ядро (gRPC-ингресс) недоступно — PushMessage/SyncDialogs не доставлены (UNAVAILABLE). + /// Ядро (gRPC-ингресс) недоступно — PushSource/SyncDialogs не доставлены (UNAVAILABLE). /// public const string IngressUnavailable = "Ядро недоступно — повторите попытку позже"; diff --git a/src/telegram-service/Deal.Telegram/Telegram/ISessionClient.cs b/src/telegram-service/Deal.Telegram/Telegram/ISessionClient.cs index 14a76db..09d2082 100644 --- a/src/telegram-service/Deal.Telegram/Telegram/ISessionClient.cs +++ b/src/telegram-service/Deal.Telegram/Telegram/ISessionClient.cs @@ -92,7 +92,7 @@ public interface ISessionClient : IAsyncDisposable /// Подписанный id диалога (каналы "-100…", группы "-…", личные "+…"). /// Сколько последних сообщений запросить. /// Отмена операции. - /// Сообщения диалога (с канальными полями для PushMessage). + /// Сообщения диалога (с канальными полями для PushSource). public Task> GetMessagesAsync( string dialogId, int limit, diff --git a/src/telegram-service/Deal.Telegram/Telegram/TelegramMessage.cs b/src/telegram-service/Deal.Telegram/Telegram/TelegramMessage.cs index 2d4e009..33e6880 100644 --- a/src/telegram-service/Deal.Telegram/Telegram/TelegramMessage.cs +++ b/src/telegram-service/Deal.Telegram/Telegram/TelegramMessage.cs @@ -12,8 +12,8 @@ public sealed record TelegramMessage /// Id сообщения в Telegram (дубль-гвард dialog+msgId ядра). /// Текст сообщения (непустой). /// Время сообщения, epoch-ms. - /// Имя диалога (title/first_name) для PushMessage.channel_name. - /// Username диалога для PushMessage.channel_handle. + /// Имя диалога (title/first_name) для source.display_name. + /// Username диалога. public TelegramMessage( string dialogId, int id, @@ -51,7 +51,7 @@ public sealed record TelegramMessage public long DateMs { get; } /// - /// Имя диалога (title/first_name) — для PushMessage.channel_name. + /// Имя диалога (title/first_name) — для source.display_name. /// public string DialogName { get; }