// Создано: 2026-04-04 // Изменено: 2026-04-12 (замена Kafka → shared-SQS) // mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → shared-SQS. // // Роль в архитектуре: // IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] // → SQS очередь "iot-telemetry" // → iot-sqs-consumer → Postgres (история телеметрии) // → event-dispatcher → Serverless Functions (триггеры) // // Логика: // 1. Подключиться к EMQX как MQTT клиент (credentials из env) // 2. Подписаться на топик "+/telemetry/+" (any namespace / telemetry / any device) // 3. При получении сообщения — отправить в SQS очередь "iot-telemetry" // 4. Payload оборачивается в envelope с метаданными (namespace, device_id, ts) // // Конфигурация через env vars: // MQTT_BROKER_URL — tcp://emqx.sless.svc:1883 // MQTT_USERNAME — username для подключения bridge к EMQX // MQTT_PASSWORD — пароль bridge клиента // SQS_ENDPOINT — https://qu.kube5s.ru (или внутрикластерный endpoint) // SQS_ACCESS_KEY — Access Key для shared-SQS tenant // SQS_SECRET_KEY — Secret Key для shared-SQS tenant // SQS_QUEUE_NAME — имя очереди (default: iot-telemetry) // SQS_REGION — регион (default: us-east-1) package main import ( "context" "encoding/json" "fmt" "log/slog" "os" "os/signal" "strings" "syscall" "time" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/sqs" mqtt "github.com/eclipse/paho.mqtt.golang" ) // mqttBridgeConfig — конфигурация сервиса из env vars. type mqttBridgeConfig struct { MQTTBrokerURL string MQTTUsername string MQTTPassword string SQSEndpoint string SQSAccessKey string SQSSecretKey string SQSQueueName string SQSRegion string } // iotTelemetryMessage — envelope сообщения публикуемого в SQS. // Потребители (iot-sqs-consumer, event-dispatcher) читают этот формат. 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, RFC3339) 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, "sqs_endpoint", cfg.SQSEndpoint, "sqs_queue", cfg.SQSQueueName, ) // SQS клиент через AWS SDK Go v2. Endpoint переопределяем на shared-SQS. sqsClient := sqs.New(sqs.Options{ Region: cfg.SQSRegion, Credentials: credentials.NewStaticCredentialsProvider( cfg.SQSAccessKey, cfg.SQSSecretKey, "", ), BaseEndpoint: aws.String(cfg.SQSEndpoint), }) // Получаем QueueUrl по имени — SQS API требует URL, а не имя очереди. // Делаем это один раз при старте. queueUrlOut, err := sqsClient.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{ QueueName: aws.String(cfg.SQSQueueName), }) if err != nil { log.Error("failed to get SQS queue URL — очередь должна существовать", "queue", cfg.SQSQueueName, "err", err) os.Exit(1) } queueURL := *queueUrlOut.QueueUrl log.Info("SQS queue resolved", "queue_url", queueURL) // Создаём 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 сообщений — отправляет envelope в SQS messageHandler := buildMQTTMessageHandler(ctx, sqsClient, queueURL, 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"), SQSEndpoint: required("SQS_ENDPOINT"), SQSAccessKey: required("SQS_ACCESS_KEY"), SQSSecretKey: required("SQS_SECRET_KEY"), SQSQueueName: getEnvOrDefault("SQS_QUEUE_NAME", "iot-telemetry"), SQSRegion: getEnvOrDefault("SQS_REGION", "us-east-1"), } } func getEnvOrDefault(key, defaultVal string) string { if v := os.Getenv(key); v != "" { return v } return defaultVal } // connectMQTT устанавливает подключение к EMQX брокеру. 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() 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 } // buildMQTTMessageHandler возвращает обработчик MQTT сообщений. // При получении — формирует envelope и отправляет в SQS очередь. // MessageGroupId = namespace (для FIFO очередей — группировка по тенанту). func buildMQTTMessageHandler(ctx context.Context, sqsClient *sqs.Client, queueURL string, log *slog.Logger) mqtt.MessageHandler { return func(_ mqtt.Client, msg mqtt.Message) { topic := msg.Topic() payload := msg.Payload() // Топик: "{namespace}/telemetry/{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) } 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 } // Отправляем в SQS. SendMessage — синхронный, но HTTP round-trip быстрый. // При ошибке логируем и продолжаем — не блокируем MQTT callback надолго. _, err = sqsClient.SendMessage(ctx, &sqs.SendMessageInput{ QueueUrl: aws.String(queueURL), MessageBody: aws.String(string(body)), }) if err != nil { log.Error("SQS SendMessage failed", "queue_url", queueURL, "err", err) return } log.Info("forwarded IoT telemetry to SQS", "mqtt_topic", topic, "namespace", ns, "device", deviceID, "sqs_queue", queueURL, ) } }