From 6caab6b7298d306fcce354dea75c9e4fe051bab6 Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 12 Apr 2026 15:28:45 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D0=B7=D0=B0=D0=BC=D0=B5=D0=BD=D0=B0=20?= =?UTF-8?q?Kafka=20=E2=86=92=20shared-SQS=20=D0=B2=20IoT=20pipeline?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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: задокументировано решение --- .gitignore | 9 +- Dockerfile | 11 +- Makefile | 11 +- cmd/iot-operator/main.go | 143 ++++++++++ cmd/mqtt-bridge/main.go | 255 ++++++++++++++++++ cmd/sqs-consumer/main.go | 207 ++++++++++++++ deployments/k8s/iot-mqtt-bridge.yaml | 38 ++- deployments/k8s/iot-sqs-consumer.yaml | 62 +++++ .../2026-04-12-replace-kafka-with-sqs.md | 77 ++++++ go.mod | 9 +- go.sum | 24 +- internal/api/handler/handler.go | 2 - .../api/handler/iot_admin_stats_handler.go | 125 ++++----- 13 files changed, 848 insertions(+), 125 deletions(-) create mode 100644 cmd/iot-operator/main.go create mode 100644 cmd/mqtt-bridge/main.go create mode 100644 cmd/sqs-consumer/main.go create mode 100644 deployments/k8s/iot-sqs-consumer.yaml create mode 100644 doc/decisions/2026-04-12-replace-kafka-with-sqs.md diff --git a/.gitignore b/.gitignore index 38f0efb..fbe076c 100644 --- a/.gitignore +++ b/.gitignore @@ -1,7 +1,8 @@ -# Бинарники -iot-operator -mqtt-bridge -kafka-consumer +# Бинарники (только в корне проекта, не директории cmd/) +/iot-operator +/mqtt-bridge +/kafka-consumer +/sqs-consumer bin/ # Go diff --git a/Dockerfile b/Dockerfile index 7396a3a..4b5473a 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,6 +1,7 @@ # Создано: 2026-04-12 +# Изменено: 2026-04-12 (kafka-consumer → sqs-consumer) # Dockerfile для IoT managed service. -# Multi-stage build: 3 бинарника (iot-operator, mqtt-bridge, kafka-consumer). +# Multi-stage build: 3 бинарника (iot-operator, mqtt-bridge, sqs-consumer). # Образ: naeel/iot-operator (Docker Hub). FROM golang:1.25 AS builder @@ -19,11 +20,11 @@ COPY internal/ internal/ # iot-operator — controller-manager + REST API RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a -o iot-operator ./cmd/iot-operator/ -# mqtt-bridge — MQTT (EMQX) → Kafka +# mqtt-bridge — MQTT (EMQX) → shared-SQS RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a -o mqtt-bridge ./cmd/mqtt-bridge/ -# kafka-consumer — Kafka → Postgres -RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a -o kafka-consumer ./cmd/kafka-consumer/ +# sqs-consumer — shared-SQS → Postgres +RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a -o sqs-consumer ./cmd/sqs-consumer/ # --- Runtime --- FROM gcr.io/distroless/static:nonroot @@ -33,7 +34,7 @@ WORKDIR / # Какой запускать — определяется command в k8s Deployment. COPY --from=builder /workspace/iot-operator . COPY --from=builder /workspace/mqtt-bridge . -COPY --from=builder /workspace/kafka-consumer . +COPY --from=builder /workspace/sqs-consumer . USER 65532:65532 ENTRYPOINT ["/iot-operator"] diff --git a/Makefile b/Makefile index 485c879..d1f6746 100644 --- a/Makefile +++ b/Makefile @@ -15,16 +15,7 @@ build-bridge: CGO_ENABLED=0 go build -o bin/mqtt-bridge ./cmd/mqtt-bridge/ build-consumer: -CGO_ENABLED=0 go build -o bin/kafka-consumer ./cmd/kafka-consumer/ - -# Docker -docker-build: -docker build -t $(IMG) . - -docker-push: -docker push $(IMG) - -# Тесты и go mod + CGO_ENABLED=0 go build -o bin/sqs-consumer ./cmd/sqs-consumer/ test: go test ./... diff --git a/cmd/iot-operator/main.go b/cmd/iot-operator/main.go new file mode 100644 index 0000000..6fcf4c2 --- /dev/null +++ b/cmd/iot-operator/main.go @@ -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") +} diff --git a/cmd/mqtt-bridge/main.go b/cmd/mqtt-bridge/main.go new file mode 100644 index 0000000..19d682d --- /dev/null +++ b/cmd/mqtt-bridge/main.go @@ -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, + ) + } +} diff --git a/cmd/sqs-consumer/main.go b/cmd/sqs-consumer/main.go new file mode 100644 index 0000000..f9a0f7d --- /dev/null +++ b/cmd/sqs-consumer/main.go @@ -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 +} diff --git a/deployments/k8s/iot-mqtt-bridge.yaml b/deployments/k8s/iot-mqtt-bridge.yaml index 5495625..e0ddf7b 100644 --- a/deployments/k8s/iot-mqtt-bridge.yaml +++ b/deployments/k8s/iot-mqtt-bridge.yaml @@ -1,25 +1,18 @@ # Создано: 2026-04-04 -# Изменено: 2026-04-06 (MQTT→Kafka: убран RABBITMQ_URL, добавлен KAFKA_BROKERS, v0.1.67) -# Deployment iot-mqtt-bridge — MQTT→Kafka мост для IoT. +# Изменено: 2026-04-12 (замена Kafka → shared-SQS) +# Deployment iot-mqtt-bridge — MQTT→SQS мост для IoT. # # Получает MQTT сообщения от EMQX (подписка на "+/telemetry/+") -# и публикует в Kafka топик "iot.telemetry" (ключ = namespace). +# и отправляет в shared-SQS очередь "iot-telemetry" через AWS SDK. # # Credentials для MQTT подключения берутся из Secret iot-bridge-credentials. -# Этот Secret нужно создать вручную ДО деплоя: +# Credentials для SQS берутся из Secret iot-sqs-credentials. # -# # 1. Создать IoTDevice для bridge через API: -# curl -X POST .../v1/namespaces/sless-bridge/iot/devices \ -# -d '{"name":"bridge","device_id":"bridge","enabled":true}' -# -# # 2. Получить credentials: -# MQTT_USERNAME=$(kubectl get secret iot-bridge -n sless-bridge -o jsonpath='{.data.mqtt-username}' | base64 -d) -# MQTT_PASSWORD=$(kubectl get secret iot-bridge -n sless-bridge -o jsonpath='{.data.mqtt-password}' | base64 -d) -# -# # 3. Создать Secret для bridge Deployment (один раз): -# kubectl create secret generic iot-bridge-credentials -n sless \ -# --from-literal=MQTT_USERNAME="$MQTT_USERNAME" \ -# --from-literal=MQTT_PASSWORD="$MQTT_PASSWORD" +# Создание SQS Secret (один раз): +# kubectl create secret generic iot-sqs-credentials -n sless \ +# --from-literal=SQS_ENDPOINT="https://qu.kube5s.ru" \ +# --from-literal=SQS_ACCESS_KEY="" \ +# --from-literal=SQS_SECRET_KEY="" # # Применение: kubectl apply -f deployments/k8s/iot-mqtt-bridge.yaml @@ -47,19 +40,20 @@ spec: # При смене версии оператора — менять тег и здесь. image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.69 imagePullPolicy: Always - command: ["/iot-mqtt-bridge"] + command: ["/mqtt-bridge"] env: - name: MQTT_BROKER_URL value: "tcp://emqx.sless.svc:1883" - - name: KAFKA_BROKERS - value: "kafka.sless.svc.cluster.local:9092" + - name: SQS_QUEUE_NAME + value: "iot-telemetry" + - name: SQS_REGION + value: "us-east-1" envFrom: - secretRef: name: iot-bridge-credentials - # IOT_PG_DSN — сохранение телеметрии в Postgres (опционально) + # SQS_ENDPOINT, SQS_ACCESS_KEY, SQS_SECRET_KEY - secretRef: - name: iot-postgres-secret - optional: true + name: iot-sqs-credentials resources: requests: memory: "32Mi" diff --git a/deployments/k8s/iot-sqs-consumer.yaml b/deployments/k8s/iot-sqs-consumer.yaml new file mode 100644 index 0000000..c7219be --- /dev/null +++ b/deployments/k8s/iot-sqs-consumer.yaml @@ -0,0 +1,62 @@ +# Создано: 2026-04-12 +# Deployment iot-sqs-consumer — читает IoT телеметрию из shared-SQS → пишет в IoT Postgres. +# +# Long polling (WaitTimeSeconds=20) — минимизирует пустые запросы. +# Сообщение удаляется из SQS только после успешной записи в Postgres (at-least-once). +# +# Credentials для SQS берутся из Secret iot-sqs-credentials. +# IOT_PG_DSN берётся из Secret iot-postgres-secret. +# +# Создание SQS Secret (один раз, если ещё не создан): +# kubectl create secret generic iot-sqs-credentials -n sless \ +# --from-literal=SQS_ENDPOINT="https://qu.kube5s.ru" \ +# --from-literal=SQS_ACCESS_KEY="" \ +# --from-literal=SQS_SECRET_KEY="" +# +# Применение: kubectl apply -f deployments/k8s/iot-sqs-consumer.yaml + +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: iot-sqs-consumer + namespace: sless + labels: + app: iot-sqs-consumer +spec: + replicas: 1 + selector: + matchLabels: + app: iot-sqs-consumer + template: + metadata: + labels: + app: iot-sqs-consumer + spec: + containers: + - name: sqs-consumer + # Тот же образ что и оператор — все IoT бинари в одном образе. + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.69 + imagePullPolicy: Always + command: ["/sqs-consumer"] + env: + - name: SQS_QUEUE_NAME + value: "iot-telemetry" + - name: SQS_REGION + value: "us-east-1" + envFrom: + # SQS_ENDPOINT, SQS_ACCESS_KEY, SQS_SECRET_KEY + - secretRef: + name: iot-sqs-credentials + # IOT_PG_DSN — master DSN для IoT Postgres (per-tenant DB) + - secretRef: + name: iot-postgres-secret + resources: + requests: + memory: "32Mi" + cpu: "25m" + limits: + memory: "64Mi" + cpu: "100m" + imagePullSecrets: + - name: sless-registry-auth diff --git a/doc/decisions/2026-04-12-replace-kafka-with-sqs.md b/doc/decisions/2026-04-12-replace-kafka-with-sqs.md new file mode 100644 index 0000000..c5a74e0 --- /dev/null +++ b/doc/decisions/2026-04-12-replace-kafka-with-sqs.md @@ -0,0 +1,77 @@ +# Решение: замена Kafka на shared-SQS + +**Дата:** 2026-04-12 +**Статус:** в работе +**Ветка:** `feature/replace-kafka-with-sqs` + +--- + +## Контекст + +IoT-сервис использует Kafka (библиотека `segmentio/kafka-go`) как промежуточную очередь между MQTT-bridge и Postgres consumer. +Kafka — тяжёлый компонент: требует отдельный деплоймент в кластере (`deployments/k8s/kafka.yaml`), ZooKeeper/KRaft, настройку топиков, партиций. + +У нас есть собственный сервис **shared-SQS** — AWS SQS-совместимая очередь сообщений: +- Репа: https://gitea.services.ngcloud.ru/Nail/shared-SQS +- Endpoint: `https://qu.kube5s.ru` +- 17 SQS-операций, multi-tenant, Redis persistence, billing +- Работает через стандартные AWS SDK + +RabbitMQ в коде IoT **не используется** (только в старых документах как план MVP). + +--- + +## Решение + +Заменить Kafka → shared-SQS во всём IoT pipeline. + +--- + +## Что меняется + +### 1. mqtt-bridge (`cmd/mqtt-bridge/main.go`) +- **Было:** `kafka.Writer` → `WriteMessages()` в топик `iot.telemetry` +- **Стало:** AWS SDK Go v2 → `sqs.SendMessage()` в очередь `iot-telemetry` +- Env: `KAFKA_BROKERS` → `SQS_ENDPOINT`, `SQS_ACCESS_KEY`, `SQS_SECRET_KEY`, `SQS_QUEUE_NAME` + +### 2. kafka-consumer → sqs-consumer (`cmd/kafka-consumer/` → `cmd/sqs-consumer/`) +- **Было:** `kafka.NewReader` с consumer group, `FetchMessage()` + `CommitMessages()` +- **Стало:** polling loop: `sqs.ReceiveMessage(WaitTimeSeconds=20)` + `sqs.DeleteMessage()` +- Переименовать директорию и бинарник + +### 3. admin stats handler (`internal/api/handler/iot_admin_stats_handler.go`) +- **Было:** Kafka lag (ListOffsets + OffsetFetch) +- **Стало:** `sqs.GetQueueAttributes(ApproximateNumberOfMessages, ApproximateNumberOfMessagesNotVisible)` +- Env: `KAFKA_BROKERS` → SQS credentials в handler + +### 4. go.mod +- Убрать: `github.com/segmentio/kafka-go` +- Добавить: `github.com/aws/aws-sdk-go-v2`, `github.com/aws/aws-sdk-go-v2/service/sqs` + +### 5. Deployments +- Удалить: `deployments/k8s/kafka.yaml` +- Обновить: `deployments/k8s/iot-mqtt-bridge.yaml` — новые env vars +- Обновить: `deployments/k8s/iot-kafka-consumer.yaml` → `iot-sqs-consumer.yaml` + +### 6. Dockerfile / Makefile +- Переименовать бинарник `kafka-consumer` → `sqs-consumer` + +--- + +## Плюсы + +- Убираем Kafka из кластера (экономия ресурсов) +- Используем свой managed сервис (единая инфраструктура) +- AWS SDK — стандартная библиотека, код проще +- Billing и мониторинг из коробки в shared-SQS + +## Риски + +- SQS — pull-based (polling latency до 20s long poll vs Kafka push). Для IoT телеметрии приемлемо. +- At-least-once delivery — нужно учитывать idempotency (уже есть в текущем коде: INSERT ON CONFLICT) + +--- + +## SQS endpoint + +Пока используем публичный `https://qu.kube5s.ru`. Если есть внутрикластерный сервис — обновим. diff --git a/go.mod b/go.mod index e7160cc..c2f0599 100644 --- a/go.mod +++ b/go.mod @@ -3,11 +3,13 @@ module gitea.services.ngcloud.ru/Nail/IoT go 1.25 require ( + github.com/aws/aws-sdk-go-v2 v1.41.5 + github.com/aws/aws-sdk-go-v2/credentials v1.19.14 + github.com/aws/aws-sdk-go-v2/service/sqs v1.42.25 github.com/eclipse/paho.mqtt.golang v1.5.1 github.com/google/uuid v1.6.0 github.com/gorilla/mux v1.8.1 github.com/lib/pq v1.11.2 - github.com/segmentio/kafka-go v0.4.50 k8s.io/api v0.26.0 k8s.io/apimachinery v0.26.0 k8s.io/client-go v0.26.0 @@ -15,6 +17,9 @@ require ( ) require ( + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.1.2 // indirect github.com/davecgh/go-spew v1.1.1 // indirect @@ -36,13 +41,11 @@ require ( github.com/imdario/mergo v0.3.6 // indirect github.com/josharian/intern v1.0.0 // indirect github.com/json-iterator/go v1.1.12 // indirect - github.com/klauspost/compress v1.15.9 // indirect github.com/mailru/easyjson v0.7.6 // indirect github.com/matttproud/golang_protobuf_extensions v1.0.2 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.2 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect - github.com/pierrec/lz4/v4 v4.1.15 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.14.0 // indirect github.com/prometheus/client_model v0.3.0 // indirect diff --git a/go.sum b/go.sum index 695678f..6a6bf12 100644 --- a/go.sum +++ b/go.sum @@ -38,6 +38,18 @@ github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751/go.mod h1:LOuy github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= github.com/alecthomas/units v0.0.0-20190717042225-c3de453c63f4/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho= +github.com/aws/aws-sdk-go-v2 v1.41.5 h1:dj5kopbwUsVUVFgO4Fi5BIT3t4WyqIDjGKCangnV/yY= +github.com/aws/aws-sdk-go-v2 v1.41.5/go.mod h1:mwsPRE8ceUUpiTgF7QmQIJ7lgsKUPQOUl3o72QBrE1o= +github.com/aws/aws-sdk-go-v2/credentials v1.19.14 h1:n+UcGWAIZHkXzYt87uMFBv/l8THYELoX6gVcUvgl6fI= +github.com/aws/aws-sdk-go-v2/credentials v1.19.14/go.mod h1:cJKuyWB59Mqi0jM3nFYQRmnHVQIcgoxjEMAbLkpr62w= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 h1:Rgg6wvjjtX8bNHcvi9OnXWwcE0a2vGpbwmtICOsvcf4= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21/go.mod h1:A/kJFst/nm//cyqonihbdpQZwiUhhzpqTsdbhDdRF9c= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 h1:PEgGVtPoB6NTpPrBgqSE5hE/o47Ij9qk/SEZFbUOe9A= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21/go.mod h1:p+hz+PRAYlY3zcpJhPwXlLC4C+kqn70WIHwnzAfs6ps= +github.com/aws/aws-sdk-go-v2/service/sqs v1.42.25 h1:8Bv3TQ1Cob6HLlpUbAnWxeHhAkYScJO9RIHh2WPXaxw= +github.com/aws/aws-sdk-go-v2/service/sqs v1.42.25/go.mod h1:eDstEbM0OEnBUnNQxIA7j74Jy61cCU1S4EMlCtdMwzs= +github.com/aws/smithy-go v1.24.2 h1:FzA3bu/nt/vDvmnkg+R8Xl46gmzEDam6mZ1hzmwXFng= +github.com/aws/smithy-go v1.24.2/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/benbjohnson/clock v1.1.0 h1:Q92kusRqC1XV2MjkWETPvjJVqKetz1OzxZB7mHJLju8= github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q= @@ -188,8 +200,6 @@ github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7V github.com/julienschmidt/httprouter v1.3.0/go.mod h1:JR6WtHb+2LUe8TCKY3cZOxFyyO8IZAc4RVcycCCAKdM= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= -github.com/klauspost/compress v1.15.9 h1:wKRjX6JRtDdrE9qwa4b/Cip7ACOshUI4smpCQanqjSY= -github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/konsorten/go-windows-terminal-sequences v1.0.3/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= @@ -225,8 +235,6 @@ github.com/onsi/ginkgo/v2 v2.6.0 h1:9t9b9vRUbFq3C4qKFCGkVuq/fIHji802N1nrtkh1mNc= github.com/onsi/ginkgo/v2 v2.6.0/go.mod h1:63DOGlLAH8+REH8jUGdL3YpCpu7JODesutUjdENfUAc= github.com/onsi/gomega v1.24.1 h1:KORJXNNTzJXzu4ScJWssJfJMnJ+2QJqhoQSRwNlze9E= github.com/onsi/gomega v1.24.1/go.mod h1:3AOiACssS3/MajrniINInwbfOOtfZvplPzuRSmvt1jM= -github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= -github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -260,8 +268,6 @@ github.com/prometheus/procfs v0.7.3/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1 github.com/prometheus/procfs v0.8.0 h1:ODq8ZFEaYeCaZOJlZZdJA2AbQR98dSHSM1KW/You5mo= github.com/prometheus/procfs v0.8.0/go.mod h1:z7EfXMXOkbkqb9IINtpCn86r/to3BnA0uaxHdg830/4= github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= -github.com/segmentio/kafka-go v0.4.50 h1:mcyC3tT5WeyWzrFbd6O374t+hmcu1NKt2Pu1L3QaXmc= -github.com/segmentio/kafka-go v0.4.50/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E= github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88= @@ -278,12 +284,6 @@ github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/ github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0 h1:pSgiaMZlXftHpm5L7V1+rVB+AZJydKsMxsQBIJw4PKk= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= -github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= -github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= -github.com/xdg-go/scram v1.1.2 h1:FHX5I5B4i4hKRVRBCFRxq1iQRej7WO3hhBuJf+UUySY= -github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= -github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= -github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= github.com/yuin/goldmark v1.1.25/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.32/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= diff --git a/internal/api/handler/handler.go b/internal/api/handler/handler.go index c0f1165..1ba31e0 100644 --- a/internal/api/handler/handler.go +++ b/internal/api/handler/handler.go @@ -30,8 +30,6 @@ K8s client.Client Scheme *runtime.Scheme // IoTPG — хранилище IoT телеметрии (per-tenant Postgres). nil если IOT_PG_DSN не задан. IoTPG *iotpg.IoTPostgresStore -// KafkaBrokers — адреса Kafka брокеров для чтения consumer lag на admin странице. -KafkaBrokers string Log *slog.Logger } diff --git a/internal/api/handler/iot_admin_stats_handler.go b/internal/api/handler/iot_admin_stats_handler.go index 15072ad..006f843 100644 --- a/internal/api/handler/iot_admin_stats_handler.go +++ b/internal/api/handler/iot_admin_stats_handler.go @@ -1,4 +1,5 @@ // Создано: 2026-04-06 +// Изменено: 2026-04-12 (замена Kafka lag → SQS queue stats) // iot_admin_stats_handler.go — handler для страницы администратора IoT. // // Endpoints: @@ -6,8 +7,8 @@ // // Источники данных: // - PostgreSQL (IoTPG): counts per tenant, last 1h/24h, latest rows -// - Kafka: consumer lag (latest offset - committed offset для group iot-pg-consumer) -// - K8s: статус подов iot-mqtt-bridge и iot-kafka-consumer +// - SQS: approximate message count (ApproximateNumberOfMessages) +// - K8s: статус подов iot-mqtt-bridge и iot-sqs-consumer // // Авторизация: Bearer из env ADMIN_STATS_TOKEN. // Если ADMIN_STATS_TOKEN не задан — endpoint возвращает 503. @@ -19,10 +20,14 @@ import ( "fmt" "net/http" "os" + "strconv" "strings" "time" - kafka "github.com/segmentio/kafka-go" + "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" + sqstypes "github.com/aws/aws-sdk-go-v2/service/sqs/types" corev1 "k8s.io/api/core/v1" "sigs.k8s.io/controller-runtime/pkg/client" ) @@ -36,12 +41,11 @@ type iotAdminPodStatus struct { Age string `json:"age"` } -// iotAdminKafkaStats — информация о Kafka топике и consumer lag. -type iotAdminKafkaStats struct { - LatestOffset int64 `json:"latest_offset"` - CommittedOffset int64 `json:"committed_offset"` - ConsumerLag int64 `json:"consumer_lag"` - Error string `json:"error,omitempty"` +// iotAdminSQSStats — информация о SQS очереди для мониторинга. +type iotAdminSQSStats struct { + ApproximateMessages int64 `json:"approximate_messages"` + ApproximateMessagesNotVisible int64 `json:"approximate_messages_not_visible"` + Error string `json:"error,omitempty"` } // AdminStats обрабатывает GET /iot-admin/stats. @@ -78,8 +82,8 @@ func (h *Handler) AdminStats(w http.ResponseWriter, r *http.Request) { result["postgres"] = map[string]any{"reachable": false, "error": "IoTPG not configured"} } - // Kafka: consumer lag для топика iot.telemetry / группы iot-pg-consumer - result["kafka"] = h.collectIotKafkaLag(ctx) + // SQS: approximate message count для очереди iot-telemetry + result["sqs"] = h.collectIotSQSStats(ctx) // K8s: статус подов bridge и consumer result["pods"] = h.collectIotPodStatuses(ctx) @@ -87,72 +91,59 @@ func (h *Handler) AdminStats(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, result) } -// collectIotKafkaLag получает latest offset топика и committed offset consumer group, -// вычисляет lag = latest - committed. -// Topic: "iot.telemetry", Consumer Group: "iot-pg-consumer". -func (h *Handler) collectIotKafkaLag(ctx context.Context) iotAdminKafkaStats { - if h.KafkaBrokers == "" { - return iotAdminKafkaStats{Error: "KAFKA_BROKERS not configured"} +// collectIotSQSStats получает ApproximateNumberOfMessages из SQS очереди. +// Использует SQS_ENDPOINT, SQS_ACCESS_KEY, SQS_SECRET_KEY, SQS_QUEUE_NAME из env. +func (h *Handler) collectIotSQSStats(ctx context.Context) iotAdminSQSStats { + endpoint := os.Getenv("SQS_ENDPOINT") + accessKey := os.Getenv("SQS_ACCESS_KEY") + secretKey := os.Getenv("SQS_SECRET_KEY") + queueName := os.Getenv("SQS_QUEUE_NAME") + if queueName == "" { + queueName = "iot-telemetry" + } + region := os.Getenv("SQS_REGION") + if region == "" { + region = "us-east-1" } - brokers := strings.Split(h.KafkaBrokers, ",") - brokerAddr := kafka.TCP(brokers...) - - kc := &kafka.Client{ - Addr: brokerAddr, - Timeout: 5 * time.Second, + if endpoint == "" || accessKey == "" || secretKey == "" { + return iotAdminSQSStats{Error: "SQS credentials not configured (SQS_ENDPOINT, SQS_ACCESS_KEY, SQS_SECRET_KEY)"} } - const topic = "iot.telemetry" - const group = "iot-pg-consumer" + sqsClient := sqs.New(sqs.Options{ + Region: region, + Credentials: credentials.NewStaticCredentialsProvider( + accessKey, secretKey, "", + ), + BaseEndpoint: aws.String(endpoint), + }) - // Получаем latest offset (конец лога — сколько всего сообщений прошло) - offsetsResp, err := kc.ListOffsets(ctx, &kafka.ListOffsetsRequest{ - Addr: brokerAddr, - Topics: map[string][]kafka.OffsetRequest{ - topic: {kafka.LastOffsetOf(0)}, + // Получаем URL очереди + queueUrlOut, err := sqsClient.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{ + QueueName: aws.String(queueName), + }) + if err != nil { + return iotAdminSQSStats{Error: fmt.Sprintf("GetQueueUrl: %v", err)} + } + + // Запрашиваем атрибуты очереди — approximate message counts + attrsOut, err := sqsClient.GetQueueAttributes(ctx, &sqs.GetQueueAttributesInput{ + QueueUrl: queueUrlOut.QueueUrl, + AttributeNames: []sqstypes.QueueAttributeName{ + sqstypes.QueueAttributeNameApproximateNumberOfMessages, + sqstypes.QueueAttributeNameApproximateNumberOfMessagesNotVisible, }, }) if err != nil { - return iotAdminKafkaStats{Error: fmt.Sprintf("list offsets: %v", err)} + return iotAdminSQSStats{Error: fmt.Sprintf("GetQueueAttributes: %v", err)} } - var latestOffset int64 - if partitions, ok := offsetsResp.Topics[topic]; ok && len(partitions) > 0 { - if partitions[0].Error == nil { - latestOffset = partitions[0].LastOffset - } - } + approxMsg, _ := strconv.ParseInt(attrsOut.Attributes["ApproximateNumberOfMessages"], 10, 64) + approxNotVisible, _ := strconv.ParseInt(attrsOut.Attributes["ApproximateNumberOfMessagesNotVisible"], 10, 64) - // Получаем committed offset consumer group (что consumer уже обработал) - fetchResp, err := kc.OffsetFetch(ctx, &kafka.OffsetFetchRequest{ - Addr: brokerAddr, - GroupID: group, - Topics: map[string][]int{topic: {0}}, - }) - if err != nil { - return iotAdminKafkaStats{ - LatestOffset: latestOffset, - Error: fmt.Sprintf("offset fetch: %v", err), - } - } - - var committedOffset int64 - if partitions, ok := fetchResp.Topics[topic]; ok && len(partitions) > 0 { - if partitions[0].Error == nil { - committedOffset = partitions[0].CommittedOffset - } - } - - lag := latestOffset - committedOffset - if lag < 0 { - lag = 0 - } - - return iotAdminKafkaStats{ - LatestOffset: latestOffset, - CommittedOffset: committedOffset, - ConsumerLag: lag, + return iotAdminSQSStats{ + ApproximateMessages: approxMsg, + ApproximateMessagesNotVisible: approxNotVisible, } } @@ -160,7 +151,7 @@ func (h *Handler) collectIotKafkaLag(ctx context.Context) iotAdminKafkaStats { func (h *Handler) collectIotPodStatuses(ctx context.Context) map[string]any { result := map[string]any{} - for _, appLabel := range []string{"iot-mqtt-bridge", "iot-kafka-consumer"} { + for _, appLabel := range []string{"iot-mqtt-bridge", "iot-sqs-consumer"} { podList := &corev1.PodList{} if err := h.K8s.List(ctx, podList, client.InNamespace("sless"),