Skip to content

NATS Bridge

Назначение

Сервис-мост (bridge) между шиной событий NATS и сервером real-time доставки Centrifugo. Подписывается на NATS subjects платформы, оборачивает события в标准化 EventEnvelope и публикует в Centrifugo через HTTP API.

Порты

Нет HTTP-портов — чистый consumer/subscriber.

Архитектура потока данных

Бэкенд-сервисы → NATS subject → Bridge (валидация JSON → EventEnvelope) → HTTP POST → Centrifugo /api/publish → WebSocket → Клиенты

Ключевые возможности

  • Подписка на NATS subjects с поддержкой wildcards (pipeline.test.>)
  • Автоматическая обёртка событий в EventEnvelope (UUID, type, data, timestamp)
  • Умная маршрутизация: user_id → персональный канал, без user_id → broadcast
  • Ретрай подключения к NATS (10 попыток, 5с интервал)
  • Экспоненциальный бэк옩 для runtime reconnect (1s → 30s)
  • Graceful shutdown (SIGINT/SIGTERM → drain NATS)
  • Валидация JSON — отбрасывает невалидные сообщения

Маппинг subjects → channels

NATS SubjectNamespaceКанал
job.createdjobjob:user:<id> / job:broadcast
job.changedjobjob:user:<id> / job:broadcast
executor.statusexecutorexecutor:user:<id> / executor:broadcast
runner.changeexecutorexecutor:user:<id> / executor:broadcast
pipeline.test.>pipelinepipeline:user:<id> / pipeline:broadcast

EventEnvelope

json
{
  "id": "<UUID v4>",
  "type": "<NATS subject>",
  "data": { ... },
  "timestamp": "<RFC 3339 UTC>"
}

Исходящий вызов Centrifugo

POST {CENTRIFUGO_API_URL}/api/publish
Headers: X-API-Key: {CENTRIFUGO_API_KEY}, Content-Type: application/json
Body: {"channel": "<channel>", "data": <envelope_json>}

Конфигурация

ПеременнаяПо умолчаниюОписание
MESSAGING_URLnats://localhost:4222URL NATS-сервера
CENTRIFUGO_API_URLhttp://localhost:8100URL HTTP API Centrifugo
CENTRIFUGO_API_KEY""API ключ Centrifugo

Исходные файлы

ФайлНазначение
cmd/server/main.goТочка входа: конфиг, подключение NATS, bridge, graceful shutdown
internal/config/config.goЗагрузка конфигурации из env
internal/bridge/bridge.goОсновная логика: подписка, обработка, публикация
internal/bridge/models.goМодели: EventEnvelope, SubjectMapping, DefaultMappings()
internal/bridge/envelope.goСоздание EventEnvelope
internal/bridge/routing.goМаршрутизация: ResolveChannel()