# Интерфейсы сервиса message-hub # Версия: 0.1.0 Документ описывает интерфейсную поверхность сервиса: HTTP-эндпоинты, WebSocket (Socket.IO), потребляемые топики Kafka и исходящие HTTP-запросы к внешним сервисам. > В отличие от классических backend-сервисов, у `message-hub` нет публичного REST API и, соответственно, нет OpenAPI-схемы (health-роуты объявлены с `include_in_schema=False`). Основные интерфейсы — это Kafka-потребители и Socket.IO. Поэтому файла `openapi.yaml` для сервиса нет. ## Как устроено взаимодействие Единое ASGI-приложение (`main:app`, `src/main.py`) объединяет три поверхности: - **HTTP** — health-проверки (`src/health.py`), обслуживаются FastStream ASGI; - **WebSocket** — Socket.IO-сервер (`socketio.AsyncServer` + `AsyncRedisManager`), пространство имён `/project` (`src/ws/namespaces.py`); - **Kafka** — потребители сообщений (`src/consumers/*.py`) на базе FastStream `KafkaRouter`. Исходящие вызовы к внешним сервисам выполняются через `httpx.AsyncClient` (`src/config/sarex.py`, `mailer.py`), базовый хост берётся из соответствующего `*_HOST` (см. `CONFIGURATION.md`). ## HTTP-эндпоинты (входящие) | Метод | Путь | Ответ | Назначение | | --- | --- | --- | --- | | GET | `/health/live` | `204 No Content` | Liveness-проба (всегда 204, если процесс жив) | | GET | `/health/ready` | `204` / `500` | Readiness-проба: проверяет доступность Kafka (`broker.ping`) и Redis (`ping`); `500`, если хотя бы один недоступен | Слушает `0.0.0.0:8000` (gunicorn + UvicornWorker). В Helm пробы настроены на `/health/live` и `/health/ready`, порт `8000`. ## WebSocket (Socket.IO) Пространство имён: **`/project`** (`ProjectNamespace`). Менеджер состояния — Redis (`AsyncRedisManager`), CORS — `*`. Параметры подключения (query string при `connect`): `project_id` (int), `user_id` (int). Клиент помещается в комнату `project:{project_id}`. | Направление | Событие | Данные | Назначение | | --- | --- | --- | --- | | client → server | `connect` | query: `project_id`, `user_id` | Подключение; вход в комнату проекта, регистрация присутствия в Redis | | client → server | `heartbeat` | — | Продление TTL присутствия пользователя в проекте | | client → server | `disconnect` | — | Отключение; выход из комнаты, снятие присутствия | | server → client | `connected_users` | `list[int]` (user_id) | Актуальный список пользователей, подключённых к проекту (рассылается в комнату `project:{project_id}` при connect/disconnect) | Ключи присутствия в Redis (TTL = `SETTINGS_CACHE_EXPIRATION`): `project:{project_id}:{user_id}`, `user_sids:{project_id}:{user_id}:{sid}`. ## Kafka-потребители (входящие сообщения) Реальные имена топиков задаются переменной `SETTINGS_TOPICS` (маппинг логических имён `planning`/`assets`/`issues` в имена топиков). Формат сообщения — `MessageSchema` (`src/schemas/message.py`): поля `schema_version`, `model`, `sender`, `type`, `body`, `timestamp`, `xtraceId`, `user_id`, `tenants`, `tags`. | Топик (логич.) | `group_id` | Offset reset | Обработчик | Назначение | | --- | --- | --- | --- | --- | | `assets` | `assets_consumer` | earliest | `update_attributes_with_assets` | Обновление атрибутов по ассетам | | `planning` | `planning` | earliest | диспетчер по `type` (см. ниже) | Обработка событий планирования | | `issues` | `project_entity` | earliest | `create_or_update_entity` | Создание/обновление сущности проекта из issue | | `issues` | `analytic_values` | latest | `proceed_entity_value` | Обработка значений аналитики по сущности | Диспетчеризация топика `planning` по полю `type` (`src/consumers/planning.py`): | `type` сообщения | Обработчик | Назначение | | --- | --- | --- | | `auto_scheduling` | `handle_auto_scheduling` | Автопланирование | | `get_converted_file` | `handle_file_export` | Экспорт/конвертация файла | | `system_log` | `handle_system_log` | Системный журнал изменений | | `email_notifications` | `handle_email_notification` | Email-уведомления | | `sync_entity_to_project` | `handle_project_entity_sync` | Синхронизация сущности в проект | | `sync_detailed_tasks_attributes` | `handle_task_attributes_sync` | Синхронизация атрибутов детальных задач | | `sync_tasks` | `handle_task_sync` | Синхронизация задач | | `detailed_tasks_analytics` | `handle_task_analytics` | Аналитика по детальным задачам | | `update_project` | — | Пропускается (в списке `PLANNING_SKIP_TYPES`) | Обработка обёрнута в `retry_handler` (`src/infrastructure/kafka/retry.py`): до `SETTINGS_MAX_RETRIES` попыток с экспоненциальной задержкой (старт `SETTINGS_RETRY_DELAY`), ручной `ack` после успеха либо исчерпания попыток. ## Исходящие HTTP-запросы к внешним сервисам Базовый хост каждого сервиса — из соответствующего `*_HOST` (см. `CONFIGURATION.md`). Итоговый URL = `` + путь из таблицы. | Сервис (`config`) | Метод | Путь | Назначение | | --- | --- | --- | --- | | `bi` (`BI_HOST`) | POST | `/internal/values/sync_value/` | Синхронизация значений аналитики | | `pm` (`PM_HOST`) | POST | `/internal/pm/detailed_tasks/` | Синхронизация атрибутов детальных задач | | `pm` (`PM_HOST`) | POST | `/internal/pm/{endpoint}/` | Автопланирование (endpoint из тела сообщения) | | `pm` (`PM_HOST`) | POST | `/internal/pm/sync_tasks/` | Синхронизация задач | | `pdf` (`PDF_CONVERTER_HOST`) | POST | `/convert_to_pdf/` | Конвертация HTML → PDF | | `mailer` (`MAILER_HOST`) | POST | `{MAILER_PREFIX}/emails/bulk` | Массовая отправка email | | `issues` (`ISSUES_HOST`) | GET | `/api/issue-types/?company_id={id}` | Список типов issue компании | | `issues` (`ISSUES_HOST`) | GET | `/api/companies/{company_id}/status-model/v2/?issue_type_id={id}` | Модель статусов по типу issue | | `eav` (`EAV_HOST`) | POST | `/api/v4/assets/search/` | Поиск ассетов по идентификаторам | | `eav` (`EAV_HOST`) | GET | `/api/v4/attribute/` | Список атрибутов | | `sarex` (`SAREX_HOST`) | GET | `internal/client/token/{user_id}/` | Получение токена клиента | ## Внешние зависимости (инфраструктура) | Зависимость | Назначение | | --- | --- | | Kafka | Источник сообщений (топики `planning`/`assets`/`issues`) | | PostgreSQL | Хранилище данных (SQLAlchemy + psycopg) | | Redis | Менеджер состояния Socket.IO и хранилище присутствия пользователей | | S3 (Yandex Object Storage) | Файловое хранилище |