feat: замена Kafka → shared-SQS в IoT pipeline

- mqtt-bridge: kafka.Writer → SQS SendMessage (AWS SDK Go v2)
- kafka-consumer → sqs-consumer: polling loop с ReceiveMessage/DeleteMessage
- admin stats: Kafka lag → SQS GetQueueAttributes
- Удалён kafka-consumer, добавлен sqs-consumer
- go.mod: убран segmentio/kafka-go, добавлен aws-sdk-go-v2
- Dockerfile, Makefile: kafka-consumer → sqs-consumer
- .gitignore: исправлен чтобы не игнорировать cmd/ директории
- deployments: новый iot-sqs-consumer.yaml, обновлён mqtt-bridge
- doc/decisions: задокументировано решение
This commit is contained in:
Naeel
2026-04-12 15:28:45 +03:00
parent a5e52c9e1c
commit 6caab6b729
13 changed files with 848 additions and 125 deletions
+143
View File
@@ -0,0 +1,143 @@
// Создано: 2026-04-12
// main.go — точка входа IoT managed service.
// Запускает controller-manager (IoTDevice CRD) и REST API сервер (порт 9090) в одном процессе.
// Перенесено из sless/main.go — только IoT-специфичная часть.
//
// Компоненты:
// 1. Controller-Manager (controller-runtime) — reconcile IoTDevice CRD
// 2. REST API Server (gorilla/mux) — CRUD устройств, телеметрия, MQTT auth, admin
//
// Конфигурация через env vars:
// IOT_PG_DSN — DSN для IoT Postgres (опционально)
package main
import (
"context"
"fmt"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"k8s.io/apimachinery/pkg/runtime"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
iotv1alpha1 "gitea.services.ngcloud.ru/Nail/IoT/api/v1alpha1"
iotcontrollers "gitea.services.ngcloud.ru/Nail/IoT/controllers"
iotapi "gitea.services.ngcloud.ru/Nail/IoT/internal/api"
"gitea.services.ngcloud.ru/Nail/IoT/internal/api/handler"
"gitea.services.ngcloud.ru/Nail/IoT/internal/storage/iotpg"
)
var scheme = runtime.NewScheme()
func init() {
utilruntime.Must(clientgoscheme.AddToScheme(scheme))
// IoT API group iot.kube5s.ru/v1alpha1
utilruntime.Must(iotv1alpha1.AddToScheme(scheme))
}
func main() {
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
ctrl.SetLogger(zap.New(zap.UseDevMode(true)))
log.Info("starting IoT managed service operator")
// Controller Manager — управляет reconcile loop для IoTDevice CRD
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
Scheme: scheme,
HealthProbeBindAddress: ":8081",
LeaderElection: false,
})
if err != nil {
log.Error("unable to create controller manager", "err", err)
os.Exit(1)
}
// Регистрируем IoTDevice controller
if err = (&iotcontrollers.IoTDeviceReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
}).SetupWithManager(mgr); err != nil {
log.Error("unable to create IoTDevice controller", "err", err)
os.Exit(1)
}
// Health/ready пробы для k8s
if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
log.Error("unable to set up health check", "err", err)
os.Exit(1)
}
if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil {
log.Error("unable to set up ready check", "err", err)
os.Exit(1)
}
// IoT Postgres — опционален (если IOT_PG_DSN не задан — телеметрия отключена)
iotPGStore, err := iotpg.NewFromEnv(log)
if err != nil {
log.Error("failed to init IoT Postgres", "err", err)
os.Exit(1)
}
if iotPGStore != nil {
defer iotPGStore.Close()
}
// REST API handler
h := &handler.Handler{
K8s: mgr.GetClient(),
Scheme: mgr.GetScheme(),
IoTPG: iotPGStore,
Log: log,
}
router := iotapi.NewRouter(h, log)
// HTTP API сервер на порту 9090
apiServer := &http.Server{
Addr: ":9090",
Handler: router,
ReadTimeout: 30 * time.Second,
WriteTimeout: 60 * time.Second,
}
// Запускаем API сервер в горутине
go func() {
log.Info("starting REST API server", "addr", ":9090")
if err := apiServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Error("API server failed", "err", err)
os.Exit(1)
}
}()
// Graceful shutdown
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer cancel()
// Запускаем controller-manager (блокирует до ctx.Done)
go func() {
log.Info("starting controller manager")
if err := mgr.Start(ctx); err != nil {
log.Error("controller manager failed", "err", err)
os.Exit(1)
}
}()
<-ctx.Done()
log.Info("shutting down IoT operator")
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
defer shutdownCancel()
if err := apiServer.Shutdown(shutdownCtx); err != nil {
log.Error("API server shutdown error", "err", err)
}
fmt.Println("IoT operator stopped")
}
+255
View File
@@ -0,0 +1,255 @@
// Создано: 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,
)
}
}
+207
View File
@@ -0,0 +1,207 @@
// Создано: 2026-04-12 (замена kafka-consumer → sqs-consumer)
// sqs-consumer/main.go — iot-sqs-consumer: читает IoT телеметрию из shared-SQS → пишет в Postgres.
//
// Роль в архитектуре:
// SQS очередь "iot-telemetry" → iot-sqs-consumer → IoT Postgres (per-tenant DB)
//
// Long polling (WaitTimeSeconds=20) — минимизирует пустые запросы к SQS.
// При временной недоступности Postgres — сообщения остаются в очереди (visibility timeout).
// DeleteMessage вызывается только после успешной записи в Postgres.
//
// Конфигурация через env vars:
// 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)
// IOT_PG_DSN — postgres://user:pass@host:5432/iotdb (master DSN для IoT Postgres)
package main
import (
"context"
"encoding/json"
"log/slog"
"os"
"os/signal"
"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"
"gitea.services.ngcloud.ru/Nail/IoT/internal/storage/iotpg"
)
// sqsConsumerConfig — конфигурация из env vars.
type sqsConsumerConfig struct {
SQSEndpoint string
SQSAccessKey string
SQSSecretKey string
SQSQueueName string
SQSRegion string
}
// iotTelemetryMessage — envelope из SQS (идентичен bridge).
type iotTelemetryMessage struct {
Namespace string `json:"namespace"`
DeviceID string `json:"device_id"`
Topic string `json:"topic"`
Payload json.RawMessage `json:"payload"`
ReceivedAt string `json:"received_at"`
}
func main() {
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
cfg := loadConsumerConfig()
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer cancel()
log.Info("starting iot-sqs-consumer",
"sqs_endpoint", cfg.SQSEndpoint,
"sqs_queue", cfg.SQSQueueName,
)
// IoT Postgres — обязательный компонент для этого сервиса
iotStore, err := iotpg.NewFromEnv(log)
if err != nil || iotStore == nil {
log.Error("failed to connect to IoT Postgres — IOT_PG_DSN required", "err", err)
os.Exit(1)
}
defer iotStore.Close()
log.Info("connected to IoT Postgres")
// SQS клиент
sqsClient := sqs.New(sqs.Options{
Region: cfg.SQSRegion,
Credentials: credentials.NewStaticCredentialsProvider(
cfg.SQSAccessKey, cfg.SQSSecretKey, "",
),
BaseEndpoint: aws.String(cfg.SQSEndpoint),
})
// Получаем QueueUrl по имени один раз при старте
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)
log.Info("sqs consumer ready, starting polling loop...")
// Polling loop — long polling (WaitTimeSeconds=20) минимизирует пустые запросы.
// При ошибке SQS — backoff 5 сек и продолжаем.
for {
if ctx.Err() != nil {
break
}
resp, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: aws.String(queueURL),
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20, // long polling — SQS держит соединение до 20 сек
})
if err != nil {
if ctx.Err() != nil {
break // штатное завершение
}
log.Error("SQS ReceiveMessage failed", "err", err)
// Backoff при ошибках SQS — не спамим запросами
select {
case <-ctx.Done():
case <-time.After(5 * time.Second):
}
continue
}
for _, sqsMsg := range resp.Messages {
if sqsMsg.Body == nil {
continue
}
if err := processSQSTelemetry(ctx, *sqsMsg.Body, iotStore, log); err != nil {
log.Error("process telemetry message", "err", err, "message_id", derefStr(sqsMsg.MessageId))
// НЕ удаляем сообщение — оно вернётся в очередь после visibility timeout
continue
}
// Удаляем сообщение из SQS только после успешной записи в Postgres
_, err := sqsClient.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(queueURL),
ReceiptHandle: sqsMsg.ReceiptHandle,
})
if err != nil {
log.Error("SQS DeleteMessage failed", "err", err, "message_id", derefStr(sqsMsg.MessageId))
}
}
}
log.Info("shutting down iot-sqs-consumer")
}
// processSQSTelemetry десериализует envelope из SQS и записывает в Postgres.
func processSQSTelemetry(ctx context.Context, body string, store *iotpg.IoTPostgresStore, log *slog.Logger) error {
var envelope iotTelemetryMessage
if err := json.Unmarshal([]byte(body), &envelope); err != nil {
// Битое сообщение — логируем и пропускаем (не блокируем очередь)
log.Warn("failed to unmarshal telemetry envelope, skipping", "err", err, "raw", body)
return nil
}
// EnsureTenantDB идемпотентен — кэшируется после первого вызова
if err := store.EnsureTenantDB(ctx, envelope.Namespace); err != nil {
return err
}
if err := store.InsertTelemetry(ctx, envelope.Namespace, envelope.DeviceID, envelope.Payload); err != nil {
return err
}
log.Info("telemetry saved to Postgres",
"namespace", envelope.Namespace,
"device", envelope.DeviceID,
)
return nil
}
// loadConsumerConfig читает конфигурацию из env vars.
func loadConsumerConfig() sqsConsumerConfig {
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 sqsConsumerConfig{
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
}
// derefStr — безопасная разыменовка строкового указателя.
func derefStr(s *string) string {
if s == nil {
return ""
}
return *s
}