From d970958e5cc43dc565774cb58a5483e840b72666 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Sun, 16 Aug 2026 15:24:26 +0400 Subject: [PATCH] fix(bridge): resubscribe on every (re)connect; emqx v0.2.1 dashboard off; iot-service v0.1.1 --- HISTORY/2026-08-16-session-log.md | 58 +++++++++++++++++++++++++++++++ Makefile | 2 +- emqx/Makefile | 2 +- emqx/emqx.conf | 5 +-- internal/service/bridge/bridge.go | 33 +++++++++++------- 5 files changed, 83 insertions(+), 17 deletions(-) diff --git a/HISTORY/2026-08-16-session-log.md b/HISTORY/2026-08-16-session-log.md index 2b8d31f..8580fea 100644 --- a/HISTORY/2026-08-16-session-log.md +++ b/HISTORY/2026-08-16-session-log.md @@ -847,3 +847,61 @@ EMQX_LOG__CONSOLE_HANDLER__LEVEL=debug (диагностика) (devices + bridge) + authorization postgres + http-фолбэк; монолит создаёт bridge-строку в iot_devices (EnsureBridgeDevice). 4. Процедура деплоя для чужого кластера (документ). + +--- + +## 26. Tail 2 выполнен + найден и исправлен баг реконнекта бриджа (14:30 GMT+03) + +### 26.1 EMQX v0.2.1 — dashboard выключен +- `emqx/emqx.conf`: `dashboard { listeners.http { enable = false } }`. + Причина: на CPU-квоте 500m swagger-генерация dashboard падает с таймаутами + (crash-loop emqx_dashboard_listener), засоряет логи. +- `emqx/Makefile`: VERSION → v0.2.1. Сборка локально, smoke-тест + (docker run -p 18083:8083 + wscat -s mqtt — соединение держится), пуш + `naeel/iot-emqx:v0.2.1` + latest (digest 16974eac84e0). +- Деплой на платформу: `kubectl -n 2fdd9658-… set image deployment/containerk8s + app=naeel/iot-emqx:v0.2.1` + patch CPU: limits cpu=2, memory=1Gi + (deck-рекомендация: задавать CPU ≥1000m при создании). +- Проверка v0.2.1: под Running, bridge-подключение живо (PINGREQ/PINGRESP), + внешний wss: wscat держит соединение, curl без subprotocol → 400 (норма), + dashboard-падений в логах нет. + +### 26.2 БАГ: бридж не переподписывался после реконнекта (найден в 26.1) +- Симптом: живая публикация (paho, rc=0) → в EMQX PUBLISH принят, + authorization_permission_allowed, publish_to — но до PG не дошёл; + в логах монолита после реконнекта нет второго "bridge: subscribed". +- Хронология: монолит v0.1.0 рестартован в 11:11 (AUTH_TEST_MODE=false), + "bridge: subscribed" в 11:11:18; EMQX перезапущен (v0.2.1), бридж + реконнект в 11:14:08 — БЕЗ повторной подписки; сессия EMQX потеряна → + подписчиков на +/telemetry/+ нет → сообщения молча теряются. +- Причина в коде: `SetCleanSession(false)` + подписка один раз в Run(); + paho при реконнекте полагается на возобновление сессии сервером, + resubscribe не выполнялся. +- ФИКС `internal/service/bridge/bridge.go`: подписка перенесена в + OnConnectHandler (выполняется при КАЖДОМ подключении); handler передаётся + в connectMQTT; Run() больше не подписывается вручную. +- Монолит: Makefile VERSION → v0.1.1; go build/vet OK; образ + `naeel/iot-service:v0.1.1` + latest (digest ba120dff63d7) собраны и запушены. +- Деплой: `kubectl -n 4504e05b-… set image deployment/containerk8s + app=naeel/iot-service:v0.1.1`. В логах нового пода: connected + subscribed + (11:22:25), после вытеснения старого пода (close 1000 — EMQX кикает + дубль clientid) реконнект 11:22:25.109 + ПОВТОРНЫЙ subscribed — фикс работает. + +### 26.3 Живой e2e после фикса (PASS) +- paho-mqtt (websockets, path=/mqtt) → wss://exqx.containerk8s.dev.nubes.ru, + client e2e-paho-2, username test_dev-001 → publish test/telemetry/dev-001 + {"e2e_v011":true,"temp":42.0}. +- API (JWT структурный, Bearer): GET /v1/namespaces/test/iot/telemetry → + count=2: id=2 ts=2026-08-16T14:23:22+03:00 payload {"temp":42.0,"e2e_v011":true}. +- Цепочка подтверждена: устройство → wss → EMQX → bridge → SQS iot-telemetry → + consumer → PG (tenant_test) → API. + +### 26.4 Инструментальные заметки +- Сырые MQTT-пакеты через python websocket-client снаружи → EMQX рвёт с + {error,bad_encoding} (оба теста); wscat и paho работают штатно → + для e2e-тестов использовать paho или wscat, не велосипеды на сырых кадрах. +- JWT для тестов API: структурный (3 части, payload с sub+exp), подпись не + проверяется (validateJWT) — минтить python-ом. +- EMQX-логика e2e-диагностики: grep EMQX-логов по clientid → CONNECT/auth → + authorization_permission_allowed → publish_to; монолит: auth/acl 200, + bridge subscribed; PG: API telemetry. diff --git a/Makefile b/Makefile index 8a1422a..79eaa11 100644 --- a/Makefile +++ b/Makefile @@ -1,7 +1,7 @@ # Makefile — монолит iot-service (образ naeel/iot-service). # Старые k8s-цели — в legacy/Makefile.old. -VERSION ?= v0.1.0 +VERSION ?= v0.1.1 IMAGE ?= naeel/iot-service LDFLAGS = -X main.version=$(VERSION) diff --git a/emqx/Makefile b/emqx/Makefile index bac94f8..8418662 100644 --- a/emqx/Makefile +++ b/emqx/Makefile @@ -1,6 +1,6 @@ # EMQX-контейнер для IoT (образ naeel/iot-emqx). -VERSION ?= v0.2.0 +VERSION ?= v0.2.1 IMAGE ?= naeel/iot-emqx .PHONY: docker-build docker-push diff --git a/emqx/emqx.conf b/emqx/emqx.conf index 44d6950..2ef72bd 100644 --- a/emqx/emqx.conf +++ b/emqx/emqx.conf @@ -96,9 +96,10 @@ listeners.ws.default { max_connections = 512 } -## Dashboard — для диагностики через port-forward (наружу не экспонируется). +## Dashboard ОТКЛЮЧЁН (enable=false): на CPU-квоте 500m swagger-генерация падает +## с таймаутами и засоряет логи. Диагностика — через порт 1883/ctl. dashboard { listeners.http { - bind = 18083 + enable = false } } diff --git a/internal/service/bridge/bridge.go b/internal/service/bridge/bridge.go index 41bf84e..4d06aa2 100644 --- a/internal/service/bridge/bridge.go +++ b/internal/service/bridge/bridge.go @@ -30,22 +30,18 @@ func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, log *sl } log.Info("bridge: SQS queue resolved", "queue", cfg.SQSQueueName) - client, err := connectMQTT(ctx, cfg, log) + // Обработчик сообщений передаём в connectMQTT: подписка выполняется в + // OnConnectHandler, т.е. повторяется при каждом (ре)подключении. + // Иначе после потери сессии EMQX бридж оставался бы без подписки + // (прецедент 2026-08-16: реконнект без resubscribe → телеметрия терялась). + handler := newMessageHandler(ctx, sqsClient, queueURL, log) + + client, err := connectMQTT(ctx, cfg, log, handler) if err != nil { return err } defer client.Disconnect(250) - handler := newMessageHandler(ctx, sqsClient, queueURL, log) - token := client.Subscribe(telemetryTopicFilter, 1, handler) - if !token.WaitTimeout(10 * time.Second) { - return errSubscribeTimeout - } - if token.Error() != nil { - return token.Error() - } - log.Info("bridge: subscribed", "filter", telemetryTopicFilter) - <-ctx.Done() log.Info("bridge: shutting down") return nil @@ -60,7 +56,8 @@ type errBridge string func (e errBridge) Error() string { return string(e) } // connectMQTT устанавливает подключение к EMQX с автореконнектом. -func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger) (mqtt.Client, error) { +// handler передаётся в OnConnectHandler: подписка повторяется при каждом подключении. +func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger, handler mqtt.MessageHandler) (mqtt.Client, error) { opts := mqtt.NewClientOptions() opts.AddBroker(cfg.MQTTBrokerURL) opts.SetClientID(cfg.MQTTClientID) @@ -78,8 +75,18 @@ func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger) (mqt opts.SetReconnectingHandler(func(_ mqtt.Client, _ *mqtt.ClientOptions) { log.Info("bridge: MQTT reconnecting...") }) - opts.SetOnConnectHandler(func(_ mqtt.Client) { + opts.SetOnConnectHandler(func(c mqtt.Client) { log.Info("bridge: MQTT connected") + token := c.Subscribe(telemetryTopicFilter, 1, handler) + if !token.WaitTimeout(10 * time.Second) { + log.Error("bridge: resubscribe timeout", "filter", telemetryTopicFilter) + return + } + if token.Error() != nil { + log.Error("bridge: resubscribe failed", "filter", telemetryTopicFilter, "err", token.Error()) + return + } + log.Info("bridge: subscribed", "filter", telemetryTopicFilter) }) client := mqtt.NewClient(opts)