From e46a8bb3e38c3c28f28828ff00c5188e8ffca56a Mon Sep 17 00:00:00 2001 From: Naeel Date: Mon, 6 Apr 2026 16:03:57 +0300 Subject: [PATCH] doc: 2026-04-06 thinking log + kafka integration plan --- doc/progress.md | 36 ++++++++++++++- doc/thinking/2026-04-06.md | 93 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 128 insertions(+), 1 deletion(-) create mode 100644 doc/thinking/2026-04-06.md diff --git a/doc/progress.md b/doc/progress.md index 39716bc..20ab09b 100644 --- a/doc/progress.md +++ b/doc/progress.md @@ -1,6 +1,40 @@ # Прогресс разработки -Последнее обновление: 2026-04-05 +Последнее обновление: 2026-04-06 + +--- + +## 2026-04-06 — Kafka интеграция (ветка iot-kafka, в процессе) + +### Цель +Заменить прямой INSERT в Postgres из bridge на Kafka pipeline: +``` +MQTT → bridge → Kafka → iot-kafka-consumer → Postgres + → event-dispatcher → Functions (будущее) +``` + +### Обоснование +- RabbitMQ для IoT был подключён в bridge но бесполезен — никто не читал очередь +- Kafka даёт буферизацию, retention 7 дней, множество потребителей +- При переходе на managed Kafka в prod — только меняется KAFKA_BROKERS в Secret + +### План +1. ✅ Документация + план +2. ✅ Ветка `iot-kafka` +3. ⏳ Helm: установить Kafka (bitnami, KRaft, 1 нод, PVC) в namespace `sless` +4. ⏳ bridge: убрать RabbitMQ, добавить Kafka producer (`segmentio/kafka-go`) +5. ⏳ Новый сервис `iot/cmd/kafka-consumer/main.go` +6. ⏳ Dockerfile + deployment манифесты +7. ⏳ Сборка v0.1.67, деплой, тест E2E + +### Что НЕ меняется +- EMQX, operator, REST API, IoT Console +- `iotpg` storage package +- ACL, auth, namespace-изоляция + +--- + +## 2026-04-05 (вечер) — v0.1.66: UX-правки + деструктивный инцидент --- diff --git a/doc/thinking/2026-04-06.md b/doc/thinking/2026-04-06.md new file mode 100644 index 0000000..9abc325 --- /dev/null +++ b/doc/thinking/2026-04-06.md @@ -0,0 +1,93 @@ +# Thinking Log — 2026-04-06 +## Агент: GitHub Copilot (Claude Sonnet 4.6) + +--- + +## Архитектурные обсуждения перед началом Kafka + +### Контекст +Пользователь обсуждал будущую prod-архитектуру IoT сервиса. +Никакого кода не менялось — чистое планирование. + +### Итоги обсуждений + +**Три отдельных кластера (принято):** +1. IoT кластер — EMQX, bridge, Kafka, iot-consumer, Postgres, REST API +2. Serverless кластер — operator, builder, event-dispatcher, Functions +3. Infra/Control кластер — Terraform для provisioning кластеров 1 и 2, DNS, TLS, auth, billing + +Это классическая схема "control plane отдельно от data plane". + +**Kafka — выбор подтверждён:** +- Сейчас: bridge → Postgres напрямую (синхронно, без буфера) +- Prod: bridge → Kafka → {consumer → Postgres, event-dispatcher → Functions} +- Dev/test: Kafka через Helm (bitnami, KRaft mode, 1 нод, PVC) +- Prod: managed Kafka (Confluent/Aiven) — только меняется KAFKA_BROKERS в Secret + +**Postgres → managed облачный: легко** +- bridge и API используют DATABASE_URL из env +- Для переключения: только заменить Secret в кластере +- Код не трогается + +**Состояние RabbitMQ для IoT (важное открытие):** +- Bridge сейчас пишет в RabbitMQ очередь `iot.{namespace}.telemetry` +- НО event-dispatcher эту очередь не читает — он настроен на serverless functions triggers +- То есть IoT-сообщения в RabbitMQ лежат мёртвым грузом — никто не читает +- Kafka заменяет RabbitMQ для IoT-части полностью + +**Что проверяли в кластере:** +- 2026-04-05: только один активный тенант `sless-16367aacb67a4a01` (созданный после инцидента) +- Устройство `device2`, одно сообщение: `{"msg":"hello1dddd1777"}` от 14:34 UTC +- 2026-04-06: kubeconfig истёк → обновил → тот же один тенант, никто новый не входил + +--- + +## План интеграции Kafka + +### Анализ текущего bridge + +Читал `iot/cmd/mqtt-bridge/main.go`. Текущая логика в `buildMQTTMessageHandler`: +1. Получает MQTT сообщение +2. Публикует в RabbitMQ (бесполезно — никто не читает) +3. Пишет напрямую в Postgres через iotpg.Store + +С Kafka нужно: +1. Получает MQTT сообщение +2. Публикует в Kafka топик `iot.telemetry` (единый топик, namespace в payload) +3. Убрать прямой INSERT в Postgres из bridge + +### Что создаётся заново + +**`iot/cmd/kafka-consumer/main.go`** — новый сервис: +- Читает из Kafka топика `iot.telemetry` +- Пишет в Postgres (та же логика что сейчас в bridge) +- Consumer group: `iot-pg-consumer` + +**Изменения в bridge:** +- Убрать RabbitMQ +- Добавить Kafka producer (библиотека `github.com/segmentio/kafka-go`) +- Env var: `KAFKA_BROKERS` вместо `RABBITMQ_URL` + +**Новые env vars:** +- bridge: `KAFKA_BROKERS=kafka.sless.svc.cluster.local:9092` +- consumer: `KAFKA_BROKERS=...`, `IOT_PG_DSN=...` + +### Что НЕ меняется +- EMQX, operator, REST API, IoT Console — не трогаются +- `iotpg` storage package — используется consumer-ом напрямую +- ACL, auth, namespace-изоляция — не меняются + +### Порядок работы +1. Документация + коммит (сейчас) +2. Ветка `iot-kafka` +3. Helm: установить Kafka в namespace `sless` +4. Переписать bridge: убрать RabbitMQ, добавить Kafka producer +5. Создать `iot/cmd/kafka-consumer/main.go` +6. Обновить Dockerfile (добавить сборку consumer) +7. Обновить deployment манифесты +8. Сборка v0.1.67, деплой, тест + +### Риски +- `kafka-go` vs `confluent-kafka-go` — выбираем `segmentio/kafka-go` (pure Go, без CGO, совместим с alpine) +- KRaft mode в Helm bitnami — убедиться что включён (без Zookeeper) +- Topic `iot.telemetry` — создаётся автоматически при первой публикации (auto.create.topics.enable=true по умолчанию)