// Создано: 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) } }