Files
Deal/docs/superpowers/specs/2026-09-11-source-contract-design.md
T
Rustam Khalimov 27c7831910
ci / build-test (push) Canceled after 0s
Deal — единая кодовая база
SaaS-мониторинг Telegram: ядро (модули Cards/Kanban/Pipeline/Tenants/Settings/
Discovery, Api, Infrastructure), сервисы telegram/ai/ml/storage, фронт Vue,
контракты и grpc-hosting, деплой-конфиги (dev/prod/observability/CI-раннер),
Gitea Actions CI, документация (ТЗ, техдок, api-map, код-стайл, планы, бэклог).

Текущее состояние: все этапы роадмапа 0–12 закрыты, сборка 5 sln 0/0,
тесты 1340/130/52/38/9 зелёные.
2026-09-11 23:56:47 +03:00

14 KiB
Raw Blame History

Дизайн: единый контракт источника + общий Storage-сервис данных

Дата: 2026-09-11. Статус: реализовано в ядре (домен, Storage-сервис, персистентность, конвейер, wire, фронт); адаптер/провайдер telegram-сервиса и Storage-выгрузка — следующие шаги. Контракт не плодит типы вложений; файлы — в общем Storage.

1. Принцип

  1. Единый строго типизированный контракт. Любой источник (Telegram, WhatsApp, Avito, сайт, файл, Excel) через адаптер приводит данные к одному типу SourceItem. Ядро, AI и ML работают только с ним.
  2. Данные файлов — в общем Storage-сервисе. Каждый сервис-источник сам выгружает свои данные (картинки, видео, аудио, документы, любые файлы) в общий Storage с токеном валидации. Storage сам определяет тип и метаданные. В контракте хранится ссылка на файл, а не сам файл.
  3. Никаких подтипов вложений в контракте. Не плодим ImagePart/VideoPart/...; есть универсальный DataRef с полем Kind, которое заполняет Storage.
  4. Ссылки, контакты и прочее, что не является файлом, идут отдельными полями контента.
  5. В ядре нет Telegram-полей и слова Telegram (только в telegram-сервисе); в комментариях нет упоминаний задач/этапов/ТЗ.

2. Единый контракт (Deal.Modules.Cards)

public sealed record SourceItem
{
    public required SourceRef Source { get; init; }
    public required SourceContent Content { get; init; }
}

public sealed record SourceRef
{
    public required string Kind { get; init; }        // "telegram", "whatsapp", "avito", "file", "excel", ...
    public string? ExternalId { get; init; }           // id в источнике (сообщение/строка/файл)
    public string? DisplayName { get; init; }          // подпись в UI
    public string? OriginRef { get; init; }            // url / deep-link / путь
    public string? Author { get; init; }
    public DateTimeOffset ReceivedAt { get; init; }
    public IReadOnlyDictionary<string, string>? Extra { get; init; }
}

public sealed record SourceContent
{
    public string? Text { get; init; }                 // основной текст
    public string? Html { get; init; }                 // разметка (если есть)
    public string? Author { get; init; }               // отправитель
    public string? Subject { get; init; }              // тема/заголовок
    public IReadOnlyList<DataRef> Data { get; init; } = [];     // ссылки на файлы в Storage
    public IReadOnlyList<string>? Links { get; init; }          // ссылки (не файлы)
    public IReadOnlyList<ContactRef>? Contacts { get; init; }   // контакты
    public IReadOnlyDictionary<string, string>? Extra { get; init; } // прочее (не файл/не ссылка/не контакт)
}

DataRef — ссылка на объект в Storage; тип и метаданные определил Storage (nullable, чтобы не плодить типы):

public sealed record DataRef
{
    public required string Id { get; init; }        // идентификатор объекта в Storage
    public required string Ref { get; init; }       // ссылка (url/путь) для скачивания/отображения
    public string? Kind { get; init; }              // определил Storage: image/video/audio/document/archive/other
    public string? MimeType { get; init; }
    public string? FileName { get; init; }
    public long? Size { get; init; }
    public int? Width { get; init; }
    public int? Height { get; init; }
    public double? DurationSec { get; init; }
    public string? PreviewRef { get; init; }        // превью/thumbnail
    public string? Caption { get; init; }
    public int? Order { get; init; }
    public IReadOnlyDictionary<string, string>? Meta { get; init; } // прочие метаданные от Storage
}

ContactRef: Name?, Phone?, Email?, Url?, Kind? (контакт может быть квалифицирован).

3. Storage-сервис (общий)

Отдельный сервис (как ai/ml/telegram), владелец — данные. Источники и ядро только ссылаются на объекты.

  • Загрузка: Upload(stream, token, fileName?) → DataRef. Каждый сервис-источник выгружает свои данные сам, передавая токен валидации (сервисный токен/mTLS — уже есть в gRPC-обвязке).
  • Определение типа: Storage сам решает Kind/MimeType/размеры/длительность (контент-снифинг); контракт типы не задаёт.
  • Чтение: Get(id) → (stream, DataRef) либо выдача ссылки/временного URL.
  • Бэкенд: объектное хранилище (MinIO/S3). Путь/бакет — по тенанту.
  • Владение: единый общий сервис; каждый источник пишет в него со своим токеном, ядро/AI/ML читают по ссылке.

4. Адаптеры источников

public interface ISourceAdapter { string Kind { get; } SourceItem Normalize(object native); }

Владельцы: telegram → telegram-сервис; local → ручное создание (Cards); whatsapp/avito/web/file/ excel → соответствующий сервис. Файлы адаптер сам выгружает в Storage и кладёт в контракт DataRef.

5. Загрузка исходника карточки

Единый способ: по SourceRef.Kind — провайдер, возвращающий SourceContent (для файла — через Storage по DataRef.Ref, для сообщения — у источника). ISourceContentProvider { Kind; LoadAsync(SourceRef) } + реестр. API ядра: GET /api/cards/{id}/source → generic контент.

6. Персистентность

В карточках вместо плоских Telegram-колонок:

  • SourceKind (text); SourceJson (jsonb, SourceRef);
  • ContentJson (jsonb, SourceContent — текст + DataRef-ссылки + прочее);
  • SourceText (text, FTS);
  • SourceRefUrl (text?, OriginRef).

Конвертер контента общий (без per-source сериализаторов). Миграции: старые удаляем → новый init с нуля.

7. Wire и фронт

  • CardDto.Source = { kind, displayName?, originRef?, receivedAt }.
  • GET /api/cards/{id}/source{ text?, html?, author?, subject?, data[], links[], contacts[], extra? }.
  • Фронт: generic блок источника + универсальный просмотрщик (по DataRef.Kind — картинка/видео/аудио/файл; ссылки/контакты — списками).

8. Этапы

  1. Домен: SourceItem/SourceRef/SourceContent/DataRef/ContactRef; удалить Telegram-маркеры из Cards.
  2. Storage-сервис: контракт gRPC, определение типа, токен валидации, бэкенд MinIO; регистрация.
  3. Персистентность: SourceKind/SourceJson/ContentJson/SourceText/SourceRefUrl, общий конвертер, новый init, маппинг KanbanStore.
  4. Pipeline: приём SourceItem, загрузка вложений в Storage адаптером, без Telegram-полей.
  5. Wire/API: generic Source в CardDto, GET /api/cards/{id}/source, провайдеры.
  6. Frontend: generic источник + универсальный просмотрщик.
  7. Telegram: адаптер + провайдер исходника (только в telegram-сервисе) + выгрузка в Storage.
  8. Комментарии: убрать упоминания Telegram из ядра и задачи/этапы — везде.

9. Реализация: зафиксированные сигнатуры

Ядро: домен

  • SourceRefs (Deal.Modules.Cards/Application/Sources): Empty, DefaultHue = "#666", HueKey = "hue", расширения DedupeKey() (вид|оригинал|внешний id), ResolveHue().
  • CardSnapshot: вместо ChannelName/ChannelHandle/ChannelHue/SourceMsg/SourceDialogId/SourceMsgIdSourceRef Source + SourceContent Content; ReceivedAt остаётся.
  • CardDto: вместо Channel/SourceMsg/SourceDialogId/SourceMsgId/прежнего SourceSourceRef Source + SourceContent Content; ReceivedAtMs остаётся. CardChannelDto/CardSourceDto удалены.
  • ICardStore.GetCardBySourceAsync(SourceRef source, CancellationToken ct).

Ядро: конвейер

  • QueuedMessage { required SourceItem Item; bool Force; }.
  • QueueItemDto { string Id; SourceRef Source; SourceContent Content; string Text; string Status; long MsgAtMs; long QueuedAtMs; bool Force; } (JsonIgnore на Force).
  • RejectRecord { SourceRef Source; SourceContent Content; string Text; long MsgAtMs; string DecidedBy; string Stage; string Reason; string Kw; string? DeterministicId; } (DeterministicId = r_{Kind}_{OriginRef}_{ExternalId}).
  • RejectedItemDto: Source/Content, DecidedBy/DecidedByLabel (решение), остальное как было.
  • IPipelineStore.ExistsDuplicateAsync(SourceRef source, CancellationToken ct).
  • PipelineChannelDto удалён.

Схема БД (схема тенанта)

  • Cards: удалить ChannelName/ChannelHandle/ChannelHue/SourceMsg/SourceDialogId/SourceMsgId; добавить SourceKind, SourceExternalId, SourceOriginRef (text, для запросов), SourceJson (text), ContentJson (text), SourceText (text). FTS: Title+Summary+SourceText+Contact.
  • QueueItems: удалить DialogId/ChannelName/ChannelHandle/ChannelHue/MsgId; добавить SourceKey (text, уникальный ключ дедупа), SourceJson, ContentJson. Text/MsgAt/Status/Force/CreatedAt/UpdatedAt остаются.
  • RejectedItems: удалить DialogId/MsgId/ChannelName/ChannelHandle/ChannelHue; добавить SourceKey, SourceJson, ContentJson. FTS — по Text.
  • Миграции tenant: старые удалить, сгенерировать новый init с нуля (данных нет).

Маппинг источника

  • 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.

Remote-просмотр исходника (2026-09-11)

  • TelegramService.ReadSource(ReadSourceRequest{dialog_id, msg_id})ReadSourceReply{found, text?, time?} (src/contracts/telegram.proto); telegram-сервис достаёт конкретное сообщение (ISessionClient.GetMessageAsync → TL Messages_GetMessages). Медиа без текста → found=false.
  • Ядро: ITelegramGateway.ReadSourceAsync + TelegramSourceContentProvider (ISourceContentProvider, Kind="telegram", Deal.Infrastructure/Integrations/Sources) — резолвится SourceContentResolver.
  • GET /api/cards/{id}/source отдаёт результат провайдера либо сохранённое содержимое карточки. Фронт: кнопка «Обновить из источника» в подробной карточке (CardDrawer.vueloadCardSource).