305 lines
12 KiB
Go
305 lines
12 KiB
Go
// Создано: 2026-04-04
|
||
// Изменено: 2026-04-05 (добавлен INSERT в IoT Postgres)
|
||
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → RabbitMQ.
|
||
//
|
||
// Роль в архитектуре:
|
||
// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] → RabbitMQ → event-dispatcher → function
|
||
//
|
||
// Логика:
|
||
// 1. Подключиться к EMQX как MQTT клиент (credentials из env)
|
||
// 2. Подписаться на топик "+/telemetry/+" (any namespace / telemetry / any device)
|
||
// 3. При получении сообщения:
|
||
// - Извлечь namespace из топика — первый сегмент до "/"
|
||
// - Опубликовать в RabbitMQ queue "iot.{namespace}.telemetry"
|
||
// - Payload передаётся as-is (JSON от устройства)
|
||
// 4. Переподключаться к RabbitMQ при разрыве (reconnect loop)
|
||
//
|
||
// Конфигурация через env vars:
|
||
// MQTT_BROKER_URL — tcp://emqx.sless.svc:1883
|
||
// MQTT_USERNAME — username для подключения bridge к EMQX
|
||
// MQTT_PASSWORD — пароль bridge клиента
|
||
// RABBITMQ_URL — amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/
|
||
//
|
||
// ВАЖНО: bridge клиент должен проходить EMQX auth — нужен IoTDevice "iot-bridge" в namespace "sless-bridge".
|
||
// Для MVP: выделить специальный namespace "sless-bridge" с устройством "bridge",
|
||
// и использовать его credentials для подключения bridge сервиса.
|
||
// Или: зарегистрировать bridge устройство через API и записать credentials в Secret.
|
||
|
||
package main
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log/slog"
|
||
"os"
|
||
"os/signal"
|
||
"strings"
|
||
"syscall"
|
||
"time"
|
||
|
||
mqtt "github.com/eclipse/paho.mqtt.golang"
|
||
amqp "github.com/rabbitmq/amqp091-go"
|
||
|
||
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/iotpg"
|
||
)
|
||
|
||
// mqttBridgeConfig — конфигурация сервиса из env vars.
|
||
type mqttBridgeConfig struct {
|
||
MQTTBrokerURL string
|
||
MQTTUsername string
|
||
MQTTPassword string
|
||
RabbitMQURL string
|
||
}
|
||
|
||
// iotTelemetryMessage — структура сообщения публикуемого в RabbitMQ.
|
||
// Оборачивает MQTT payload в envelope с метаданными.
|
||
type iotTelemetryMessage struct {
|
||
// Namespace — k8s namespace пользователя (из MQTT topic)
|
||
Namespace string `json:"namespace"`
|
||
// DeviceID — идентификатор устройства (из MQTT topic, последний сегмент)
|
||
DeviceID string `json:"device_id"`
|
||
// Topic — оригинальный MQTT topic
|
||
Topic string `json:"topic"`
|
||
// Payload — данные от устройства (JSON передаётся as-is / строка если не JSON)
|
||
Payload json.RawMessage `json:"payload"`
|
||
// ReceivedAt — время получения сообщения мостом (UTC)
|
||
ReceivedAt string `json:"received_at"`
|
||
}
|
||
|
||
func main() {
|
||
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
|
||
|
||
cfg := loadBridgeConfig()
|
||
|
||
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
|
||
defer cancel()
|
||
|
||
log.Info("starting iot-mqtt-bridge",
|
||
"mqtt_broker", cfg.MQTTBrokerURL,
|
||
"mqtt_username", cfg.MQTTUsername,
|
||
)
|
||
|
||
// RabbitMQ connection с reconnect loop
|
||
rabbitConn, err := connectRabbitMQWithRetry(ctx, cfg.RabbitMQURL, log)
|
||
if err != nil {
|
||
log.Error("failed to connect to RabbitMQ", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
defer rabbitConn.Close()
|
||
|
||
rabbitCh, err := rabbitConn.Channel()
|
||
if err != nil {
|
||
log.Error("failed to open RabbitMQ channel", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
defer rabbitCh.Close()
|
||
|
||
// IoT Postgres — сохранение телеметрии (per-tenant DB).
|
||
// Опционально: если IOT_PG_DSN не задан — продолжаем работать без Postgres (только RabbitMQ)
|
||
iotPGStore, err := iotpg.NewFromEnv(log)
|
||
if err != nil {
|
||
log.Error("failed to connect to IoT Postgres", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
if iotPGStore != nil {
|
||
defer iotPGStore.Close()
|
||
log.Info("connected to IoT Postgres for telemetry storage")
|
||
}
|
||
|
||
// Создаём MQTT клиент
|
||
mqttClient, err := connectMQTT(cfg, log)
|
||
if err != nil {
|
||
log.Error("failed to connect to MQTT broker", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
defer mqttClient.Disconnect(500)
|
||
|
||
// Функция-обработчик MQTT сообщений
|
||
// Вызывается в goroutine paho при каждом сообщении
|
||
messageHandler := buildMQTTMessageHandler(ctx, rabbitCh, iotPGStore, log)
|
||
|
||
// Подписываемся на все telemetry топики всех namespace
|
||
// "+/telemetry/+" = {любой namespace}/telemetry/{любой deviceId}
|
||
const telemetryTopicFilter = "+/telemetry/+"
|
||
token := mqttClient.Subscribe(telemetryTopicFilter, 1, messageHandler)
|
||
token.Wait()
|
||
if token.Error() != nil {
|
||
log.Error("mqtt subscribe failed", "topic", telemetryTopicFilter, "err", token.Error())
|
||
os.Exit(1)
|
||
}
|
||
log.Info("subscribed to MQTT topic", "filter", telemetryTopicFilter)
|
||
|
||
<-ctx.Done()
|
||
log.Info("shutting down iot-mqtt-bridge")
|
||
}
|
||
|
||
// loadBridgeConfig читает конфигурацию из env vars.
|
||
// Завершает процесс если обязательные переменные отсутствуют.
|
||
func loadBridgeConfig() mqttBridgeConfig {
|
||
required := func(key string) string {
|
||
v := os.Getenv(key)
|
||
if v == "" {
|
||
slog.Error("required env var not set", "key", key)
|
||
os.Exit(1)
|
||
}
|
||
return v
|
||
}
|
||
|
||
return mqttBridgeConfig{
|
||
MQTTBrokerURL: getEnvOrDefault("MQTT_BROKER_URL", "tcp://emqx.sless.svc:1883"),
|
||
MQTTUsername: required("MQTT_USERNAME"),
|
||
MQTTPassword: required("MQTT_PASSWORD"),
|
||
RabbitMQURL: required("RABBITMQ_URL"),
|
||
}
|
||
}
|
||
|
||
func getEnvOrDefault(key, defaultVal string) string {
|
||
if v := os.Getenv(key); v != "" {
|
||
return v
|
||
}
|
||
return defaultVal
|
||
}
|
||
|
||
// connectMQTT устанавливает подключение к EMQX брокеру.
|
||
// AutoReconnect=true — paho сам переподключается при разрыве.
|
||
func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
|
||
opts := mqtt.NewClientOptions()
|
||
opts.AddBroker(cfg.MQTTBrokerURL)
|
||
opts.SetClientID("sless-iot-bridge")
|
||
opts.SetUsername(cfg.MQTTUsername)
|
||
opts.SetPassword(cfg.MQTTPassword)
|
||
opts.SetAutoReconnect(true)
|
||
opts.SetConnectRetry(true)
|
||
opts.SetConnectRetryInterval(5 * time.Second)
|
||
opts.SetKeepAlive(30 * time.Second)
|
||
opts.SetCleanSession(false) // сохраняем подписки при реконнекте
|
||
|
||
opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) {
|
||
log.Warn("MQTT connection lost, reconnecting...", "err", err)
|
||
})
|
||
opts.SetReconnectingHandler(func(_ mqtt.Client, _ *mqtt.ClientOptions) {
|
||
log.Info("MQTT reconnecting...")
|
||
})
|
||
opts.SetOnConnectHandler(func(_ mqtt.Client) {
|
||
log.Info("MQTT connected to broker")
|
||
})
|
||
|
||
client := mqtt.NewClient(opts)
|
||
token := client.Connect()
|
||
// Ждём максимум 30 секунд
|
||
if !token.WaitTimeout(30 * time.Second) {
|
||
return nil, fmt.Errorf("MQTT connect timeout")
|
||
}
|
||
if token.Error() != nil {
|
||
return nil, fmt.Errorf("MQTT connect: %w", token.Error())
|
||
}
|
||
return client, nil
|
||
}
|
||
|
||
// connectRabbitMQWithRetry подключается к RabbitMQ с повторными попытками.
|
||
// Retry нужен потому что RabbitMQ может стартовать позже bridge сервиса.
|
||
func connectRabbitMQWithRetry(ctx context.Context, url string, log *slog.Logger) (*amqp.Connection, error) {
|
||
const maxAttempts = 10
|
||
for attempt := 1; attempt <= maxAttempts; attempt++ {
|
||
conn, err := amqp.Dial(url)
|
||
if err == nil {
|
||
log.Info("connected to RabbitMQ", "attempt", attempt)
|
||
return conn, nil
|
||
}
|
||
log.Warn("RabbitMQ connection failed, retrying...", "attempt", attempt, "err", err)
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil, ctx.Err()
|
||
case <-time.After(5 * time.Second):
|
||
}
|
||
}
|
||
return nil, fmt.Errorf("exhausted %d RabbitMQ connection attempts", maxAttempts)
|
||
}
|
||
|
||
// buildMQTTMessageHandler возвращает функцию-обработчик MQTT сообщений.
|
||
// Замыкание над rabbitCh (RabbitMQ channel), iotStore (может быть nil) и logger.
|
||
// Порядок действий при получении сообщения:
|
||
// 1. INSERT в IoT Postgres (tenant DB) — если iotStore != nil
|
||
// 2. Publish в RabbitMQ — всегда (для event-dispatcher → function triggers)
|
||
//
|
||
// Ошибка INSERT не блокирует RabbitMQ publish — разные failure domain.
|
||
func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotStore *iotpg.IoTPostgresStore, log *slog.Logger) mqtt.MessageHandler {
|
||
return func(_ mqtt.Client, msg mqtt.Message) {
|
||
topic := msg.Topic()
|
||
payload := msg.Payload()
|
||
|
||
// Топик: "{namespace}/telemetry/{deviceId}"
|
||
// Извлекаем namespace (первый сегмент) и deviceId (третий сегмент)
|
||
parts := strings.SplitN(topic, "/", 3)
|
||
if len(parts) != 3 {
|
||
log.Warn("unexpected MQTT topic format, skipping", "topic", topic)
|
||
return
|
||
}
|
||
ns := parts[0]
|
||
deviceID := parts[2]
|
||
|
||
// Нормализуем payload: если это не JSON — оборачиваем в строку
|
||
rawPayload := json.RawMessage(payload)
|
||
if !json.Valid(payload) {
|
||
quotedBytes, _ := json.Marshal(string(payload))
|
||
rawPayload = json.RawMessage(quotedBytes)
|
||
}
|
||
|
||
// ШАГ 1: INSERT в IoT Postgres — сохраняем телеметрию в per-tenant DB
|
||
// EnsureTenantDB идемпотентен: кэшируется после первого вызова
|
||
if iotStore != nil {
|
||
if err := iotStore.EnsureTenantDB(ctx, ns); err != nil {
|
||
log.Error("ensure tenant DB", "namespace", ns, "err", err)
|
||
// НЕ возвращаемся — продолжаем RabbitMQ publish
|
||
} else if err := iotStore.InsertTelemetry(ctx, ns, deviceID, rawPayload); err != nil {
|
||
log.Error("insert telemetry", "topic", topic, "err", err)
|
||
// НЕ возвращаемся — RabbitMQ не должен зависеть от Postgres
|
||
} else {
|
||
log.Debug("telemetry saved to Postgres", "namespace", ns, "device", deviceID)
|
||
}
|
||
}
|
||
|
||
// ШАГ 2: Publish в RabbitMQ (для event-dispatcher → function triggers)
|
||
envelope := iotTelemetryMessage{
|
||
Namespace: ns,
|
||
DeviceID: deviceID,
|
||
Topic: topic,
|
||
Payload: rawPayload,
|
||
ReceivedAt: time.Now().UTC().Format(time.RFC3339),
|
||
}
|
||
|
||
body, err := json.Marshal(envelope)
|
||
if err != nil {
|
||
log.Error("marshal telemetry message", "topic", topic, "err", err)
|
||
return
|
||
}
|
||
|
||
// Queue name: "iot.{namespace}.telemetry"
|
||
// Declare-on-publish: если queue не существует — создаём
|
||
queueName := fmt.Sprintf("iot.%s.telemetry", ns)
|
||
if _, err := rabbitCh.QueueDeclare(queueName, true, false, false, false, nil); err != nil {
|
||
log.Error("declare RabbitMQ queue", "queue", queueName, "err", err)
|
||
return
|
||
}
|
||
|
||
err = rabbitCh.Publish(
|
||
"", // exchange — default exchange
|
||
queueName, // routing key = queue name для default exchange
|
||
false, // mandatory
|
||
false, // immediate
|
||
amqp.Publishing{
|
||
ContentType: "application/json",
|
||
Body: body,
|
||
DeliveryMode: amqp.Persistent, // сохранять при рестарте RabbitMQ
|
||
},
|
||
)
|
||
if err != nil {
|
||
log.Error("publish to RabbitMQ", "queue", queueName, "err", err)
|
||
return
|
||
}
|
||
|
||
log.Info("forwarded IoT telemetry", "topic", topic, "namespace", ns, "device", deviceID, "queue", queueName)
|
||
}
|
||
}
|