8.2 KiB
Интерфейсы сервиса 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) на базе FastStreamKafkaRouter.
Исходящие вызовы к внешним сервисам выполняются через 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 = <HOST> + путь из таблицы.
Сервис (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) | Файловое хранилище |