95 lines
8.2 KiB
Markdown
95 lines
8.2 KiB
Markdown
# Интерфейсы сервиса 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 = `<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) | Файловое хранилище |
|