# Интерфейсы сервиса dps-message-hub # Версия: 0.1.0 Документ описывает интерфейсную поверхность сервиса: потребляемые топики Kafka и обращения к внешним зависимостям (PostgreSQL). > `dps-message-hub` — это чистый Kafka-воркер на [FastStream](https://faststream.airt.ai/). У него **нет HTTP/REST API** (в коде нет FastAPI/Flask/aiohttp, health-роутов и т.п.), поэтому файла `openapi.yaml` для сервиса нет. Исходящих HTTP-запросов к другим сервисам он тоже не выполняет — единственный получатель данных — база PostgreSQL. ## Как устроено взаимодействие Единое FastStream-приложение (`dps_message_hub.infra.app:get_app`, `src/dps_message_hub/infra/app.py`) объединяет: - **Kafka** — потребитель сообщений (`KafkaBroker` + `KafkaRouter`), топик `assets` (`src/dps_message_hub/features/assets/interface/kafka.py`); - **PostgreSQL** — пул соединений `psycopg` (`AsyncConnectionPool`), открывается/закрывается в lifespan (`src/dps_message_hub/infra/lifespan.py`, `infra/database.py`). Каждое сообщение проходит через middleware `Retry` (`src/dps_message_hub/interface/middleware.py`): при исключении обработка повторяется бесконечно с экспоненциальной задержкой (`1 << min(10, retry_count)` секунд). Коммит оффсета ручной — у потребителя `auto_commit=False`. ## Kafka-потребители (входящие сообщения) Реальное имя топика задаётся переменной `DPS_MESSAGE_HUB_KAFKA__TOPICS` (маппинг логического имени `assets` в имя топика Kafka). Формат сообщения — `BrokerMessageDto` (`src/dps_message_hub/infra/dto.py`): поля `schema_version`, `model`, `sender`, `type`, `body`, `diff`, `timestamp`, `trace_id` (alias `xtraceId`), `user_id`, `tenants`, `tags`. | Топик (логич.) | `group_id` | Offset reset | `auto_commit` | Обработчик | Назначение | | --- | --- | --- | --- | --- | --- | | `assets` | `dps_assets_consumer` | earliest | `false` | `assets` | Обновление событий разметки (`markup_event`) по изменению ассета | Диспетчеризация внутри обработчика `assets` по полю `type` (`src/dps_message_hub/features/assets/interface/kafka.py`); `body` разбирается в модель `AssetUpdate` (поля `id`, `attributes[]`), `diff` — произвольный `dict`: | `type` сообщения | Условие (`diff`) | Действие | Назначение | | --- | --- | --- | --- | | `model_updated` | `diff is None` | — | Пропуск (нет изменений) | | `model_updated` | в `diff` есть ключ `resource_id` | `unlink_asset(asset_id)` | Отвязка ассета: деактивация его событий разметки | | `model_updated` | в `diff` есть ключ `attributes` | `update_markup_events(asset_id, attributes)` | Пересоздание событий разметки по обновлённым атрибутам | | `model_deleted` | — | `unlink_asset(asset_id)` | Отвязка ассета при удалении модели | Бизнес-логика (`src/dps_message_hub/features/assets/application/asset.py`): для обновляемых атрибутов существующие активные `markup_event` деактивируются (`inactivation_time`), затем в транзакции вставляются новые записи. Тип атрибута (поле `type`) определяет целевую колонку значения (`value_int`/`value_float`/`value_string`/`value_option_id`/`unit_option_id`/`value_dt`/`value_date`) — маппинг в `features/assets/domain/entities.py`. ## Внешние зависимости (инфраструктура) | Зависимость | Назначение | | --- | --- | | Kafka | Источник сообщений (топик `assets`). Подключение `SASLScram512`, при заданном `KAFKA__SSL_CAFILE` — по `SASL_SSL` | | PostgreSQL | Хранилище данных (`psycopg` async, таблица `markup_event`). Операции: `SELECT`, `UPDATE`, массовая вставка `COPY ... FROM STDIN` (`features/assets/infra/database.py`) | ## Исходящие HTTP-запросы к внешним сервисам Отсутствуют. Сервис не содержит HTTP-клиентов (`httpx`/`aiohttp`/`requests`) и не обращается к другим сервисам по HTTP; все побочные эффекты — запись в PostgreSQL.