Files
sless/doc/thinking/2026-04-04.md
T
Naeel 1e53766c46 feat(iot): Этапы 2-7 — MQTT auth, IoT API, EMQX, mqtt-bridge, E2E demo
Этап 2+4: internal/api/handler/iot_device_handler.go
  - MQTTAuth: POST /internal/mqtt/auth (без JWT, для EMQX)
  - CreateIoTDevice, ListIoTDevices, GetIoTDevice (c password), DeleteIoTDevice, UpdateIoTDevice
  - crypto/subtle.ConstantTimeCompare против timing attacks

Этап 4: internal/api/router.go
  - /v1/namespaces/{ns}/iot/devices CRUD
  - /internal/mqtt/auth (без JWT middleware)

Этап 3: deployments/k8s/emqx.yaml
  - EMQX 5.5.1, emqx.conf (HOCON) с HTTP auth backend
  - Сервис exposure: 1883 (MQTT), 8083 (WS), 18083 (Dashboard)

Этап 3: iot/cmd/mqtt-bridge/main.go
  - paho.mqtt.golang: подписка на +/telemetry/+
  - amqp091-go: publish в iot.{namespace}.telemetry
  - deployments/k8s/iot-mqtt-bridge.yaml

Этап 7: examples/IOT/ — E2E demo (main.tf, handler.py, README.md)

go.mod: добавлен github.com/eclipse/paho.mqtt.golang v1.5.1
go build ./... — ошибок нет
2026-04-04 09:45:23 +03:00

33 KiB
Raw Blame History

Лог мышления — 2026-04-04

Агент: GitHub Copilot (Claude Opus 4.6)


Задача: Архитектура Managed IoT Service

Что имеем

Изучил текущую архитектуру sless:

  • Каждому пользователю — свой namespace sless-{hash} (CRD объекты) + sless-fn-{hash} (рабочие нагрузки)
  • Есть 3 типа триггеров: HTTP, Cron, Event (RabbitMQ)
  • Event-dispatcher уже умеет: подписка на RabbitMQ queue → POST в функцию
  • Сборка через kaniko, образы в registry, S3 для кода

Вопрос пользователя

Нужен managed IoT сервис. Вопрос: каждому юзеру свой брокер (Rabbit/Kafka), свой Postgres?

Мои рассуждения

Вариант A: Всё изолированно (per-user)

  • Каждому юзеру: свой MQTT-брокер (EMQX/VerneMQ), свой RabbitMQ, свой Postgres
  • Плюсы: полная изоляция, нет noisy neighbor, простая модель безопасности
  • Минусы: огромный расход ресурсов. 100 юзеров = 100 MQTT-брокеров + 100 Postgres + 100 RabbitMQ. Это нереально на одном кластере

Вариант B: Shared инфраструктура с логической изоляцией

  • Один MQTT-брокер (EMQX) — multi-tenant через vhost/namespace prefix в топиках
  • Один RabbitMQ (уже есть!) — vhost per user
  • Один Postgres — schema per user или row-level security
  • Плюсы: экономия ресурсов, управляемость
  • Минусы: сложнее изоляция, risk noisy neighbor

Вариант C: Гибридный (мой выбор)

  • Shared: MQTT-брокер (EMQX с multi-tenancy), PostgreSQL (schema per user)
  • Per-user в namespace: только легковесные компоненты — bridge/adapter pod
  • Существующий RabbitMQ: использовать как есть, vhost per user
  • Reason: IoT-устройства общаются через MQTT → сообщения попадают в RabbitMQ через bridge → event-dispatcher уже умеет доставлять в функции

Архитектурная цепочка (Вариант C)

IoT Device → MQTT (topic: {user-prefix}/device/telemetry)
  → EMQX Rule Engine / Bridge → RabbitMQ vhost={user} queue={trigger-queue}
    → event-dispatcher (уже есть!) → POST → serverless function
      → function пишет в Postgres (per-user schema) / отправляет команду обратно
        → MQTT publish → device

Что нового нужно создать

  1. MQTT-брокер — EMQX (есть multi-tenancy, WebSocket, rule engine, k8s operator)
  2. IoT Device Registry — CRD IoTDevice (имя, credentials, metadata, user namespace)
  3. IoT Bridge — сервис MQTT→RabbitMQ per user-namespace (или EMQX rule engine)
  4. Device Auth — MQTT username/password или client certificates, хранение в k8s secrets per namespace
  5. Terraform provider — ресурсы: sless_iot_device, sless_iot_rule

Postgres

НЕ нужен per-user Postgres. Достаточно:

  • Shared Postgres (уже есть для логов)
  • Schema per user для IoT data (telemetry, device state)
  • Или TimescaleDB extension для time-series данных

RabbitMQ / Kafka

  • RabbitMQ уже есть и используется event-dispatcher'ом
  • Vhost per user — достаточная изоляция
  • Kafka избыточен на данном этапе (нужен только при миллионах сообщений/сек)

Вывод

Поднимать per-user брокер/БД — НЕ нужно. Нужна shared инфраструктура с логической изоляцией. Основная работа: EMQX + мост в RabbitMQ + device registry CRD.


Агент: GitHub Copilot (Claude Sonnet 4.6)

Дата сессии: 2026-04-04

Задача: Реализация IoT MVP по плану из doc/iot-mvp-plan.md


Шаг 0: Изучение кодовой базы (перед написанием кода)

Прочитал следующие файлы:

  • api/v1alpha1/function_types.go — паттерн написания CRD типов
  • api/v1alpha1/trigger_types.go — использование TriggerType string const
  • api/v1alpha1/groupversion_info.go — паттерн регистрации API group
  • controllers/function_controller.go — паттерн контроллера, finalizer, Reconcile
  • controllers/service_controller.go — второй контроллер для сравнения
  • internal/api/router.go — регистрация маршрутов, auth middleware
  • main.go — как регистрируются контроллеры, как подключается API
  • go.mod — версии зависимостей (k8s v0.26.0, controller-runtime v0.14.1)

Что понял

Go module: gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless API сервер: gorilla/mux, auth через middleware на /v1/ subrouter Паттерн контроллера: r.Get → finalizer → business logic → r.Status().Update() Finalizer: у function = sless.kube5s.ru/finalizer, у service = sless.kube5s.ru/service-finalizer

  • Используют ручной containsString/removeString, НЕ controllerutil
  • Я решил использовать controllerutil.ContainsFinalizer/AddFinalizer/RemoveFinalizer — это более современный и рекомендуемый подход из controller-runtime OwnerReference: нигде не вижу в существующих контроллерах. Для IoT Secret нужна — чтобы Secret автоматически удалялся при удалении IoTDevice через cascade deletion. Использую controllerutil.SetControllerReference.

Почему controllerutil вместо ручного containsString

Существующий код использует устаревший подход (хелперы containsString/removeString). controllerutil — официальный API controller-runtime для работы с finalizers. Это не рефакторинг старого кода — это новый контроллер, пишу его правильно сразу.


Шаг 1: CRD файлы и контроллер

Решение по структуре директорий

По плану: iot/api/v1alpha1/ и iot/controllers/. Это отдельный package path от api/v1alpha1.

  • IoT CRD package: gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/api/v1alpha1
  • IoT контроллер package: gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/controllers
  • Тот же go.mod — не нужно создавать отдельный модуль

Решение по API Group

По плану: iot.kube5s.ru — отдельная от sless.kube5s.ru. Причина: при выносе в отдельную репу CRD не будет конфликтовать. Принимаю.

MQTTUsername формат

По плану: {namespace}_{deviceId}. Пример: sless-abc123def456_sensor-01 Причина: EMQX требует глобально уникальный username. Namespace даёт изоляцию.

Secret name

По плану: iot-{deviceId}. Возможная проблема: deviceId может содержать символы недопустимые в k8s Secret именах (только [a-z0-9-]). Решение: в kubebuilder validation на DeviceID добавить regex [a-z0-9-]+. Если deviceId уже проходит валидацию — проблемы нет. В плане валидация не упомянута, но это необходимо чтобы имя Secret было валидным. Добавлю +kubebuilder:validation:Pattern.

Генерация пароля

32 байта через crypto/rand.Readhex.EncodeToString = 64 символа. Это достаточно энтропии (256 бит).

OwnerReference у Secret

С OwnerReference Secret автоматически удалится при удалении IoTDevice (cascade GC в k8s). Поэтому в finalizer обработчике нет нужды явно удалять Secret — просто убираем finalizer.

Но есть нюанс: если IoTDevice и Secret находятся в одном namespace — cascade deletion работает. В нашем случае оба в sless-{hash} — OK.

Обработка статуса

r.Status().Update() — только subresource. Не трогает spec или metadata. Это важно чтобы не вызвать лишний reconcile цикл (обновление spec → новый reconcile → loop).

Disabled устройство

Если spec.enabled == false:

  • Secret НЕ создаём (устройство не должно подключаться)
  • Если Secret уже существует — НЕ удаляем (при re-enable пароль не изменится)
  • Status: phase = "Disabled" Это соответствует плану.

Стоп — перечитал план: "Установить status.phase = 'Disabled' — НЕ удалять Secret". Значит если disabled — просто обновить статус, Secret остаётся. Принимаю.


Что создаю (Этап 1)

  1. iot/api/v1alpha1/device_types.go — CRD IoTDevice
  2. iot/api/v1alpha1/groupversion_info.go — API group iot.kube5s.ru/v1alpha1
  3. iot/controllers/iotdevice_controller.go — контроллер
  4. iot/config/crd/bases/ — директория для CRD YAML (создаётся controller-gen через SSH)
  5. Обновление main.go — регистрация IoT схемы и контроллера

Этапы 2+ (MQTT auth, EMQX, API routes, Terraform) — отдельно после одобрения Этапа 1.


Этап 2-7: план перед реализацией

API port

Из deployments/k8s/operator.yaml: API_PORT: "9090", сервис sless-operator.sless.svc:9090. RabbitMQ: amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/

EMQX версия — проблема

В плане указан emqx/emqx:5.5.1. Изучил вопрос:

  • EMQX 5.x open source НЕ имеет встроенного RabbitMQ bridge (только в Enterprise)
  • EMQX 4.x имеет RabbitMQ bridge через plugin, конфигурируется env vars

Рассматривал варианты: A. EMQX 4.4 — встроенный bridge, но env vars другого формата чем в плане B. EMQX 5.x + HTTP Webhook rule → наш bridge HTTP сервер C. EMQX 5.x + MQTT client (paho) в bridge сервисе

Выбрал вариант C: mqtt-bridge Go сервис с github.com/eclipse/paho.mqtt.golang

  • Не зависит от версии EMQX (работает с любым MQTT брокером)
  • amqp091-go уже в go.mod
  • paho.mqtt.golang добавляется через go get по SSH
  • Самый надёжный и тестируемый подход

EMQX 5.5.1: используем только для HTTP auth (через emqx.conf HOCON). Bridge service подключается к EMQX как обычный MQTT клиент.

MQTT Auth

Константы из существующего кода и CRD:

  • username format: {namespace}_{deviceId}_ разделитель безопасен (namespace не содержит _)
  • Secret name: iot-{deviceId}
  • Always return HTTP 200, body {"result": "allow"|"deny"} (безопасно для обеих версий EMQX)
  • crypto/subtle.ConstantTimeCompare для сравнения паролей

Структура файлов Этапов 2-7

  • internal/api/handler/iot_device_handler.go — MQTT auth + IoT CRUD handlers
  • internal/api/router.go — добавить IoT routes
  • deployments/k8s/emqx.yaml — EMQX deployment c emqx.conf ConfigMap (только HTTP auth)
  • iot/cmd/mqtt-bridge/main.go — MQTT subscriber → RabbitMQ publisher
  • deployments/k8s/iot-mqtt-bridge.yaml — Deployment mqtt-bridge
  • examples/IOT/ — E2E demo

Результат выполнения Этапов 2-7

Создано:

  • internal/api/handler/iot_device_handler.go — MQTT auth + IoT CRUD handlers
  • internal/api/router.go — IoT routes + /internal/mqtt/auth
  • deployments/k8s/emqx.yaml — EMQX 5.5.1 deployment с emqx.conf (HTTP auth)
  • iot/cmd/mqtt-bridge/main.go — MQTT subscriber → RabbitMQ publisher (paho + amqp091-go)
  • deployments/k8s/iot-mqtt-bridge.yaml — Deployment mqtt-bridge
  • examples/IOT/ — E2E demo (main.tf, handler.py, README.md)

go.mod: добавлен github.com/eclipse/paho.mqtt.golang v1.5.1

go build ./... — ошибок нет.

Не реализовано (отложено):

  • Этап 6 (Terraform Provider) — находится в отдельном репозитории, путь неизвестен
  • Terraform ресурс sless_iot_device — реализуется отдельно в provider репо

Ключевые архитектурные решения:

  • EMQX 5.5.1 (как в плане) — HTTP auth через emqx.conf HOCON
  • mqtt-bridge использует paho.mqtt.golang (MQTT subscriber), а не EMQX webhook — версионно-независимо
  • MQTTAuth всегда возвращает HTTP 200 (совместимо с EMQX 4.x и 5.x)
  • crypto/subtle.ConstantTimeCompare для защиты от timing attacks
  • GetIoTDevice — единственный endpoint с mqtt_password (security by design)

Создано:

  • iot/api/v1alpha1/device_types.go — CRD IoTDevice с IoTDevicePhase константами
  • iot/api/v1alpha1/groupversion_info.go — API group iot.kube5s.ru/v1alpha1
  • iot/api/v1alpha1/zz_generated.deepcopy.go — сгенерирован controller-gen
  • iot/config/crd/bases/iot.kube5s.ru_iotdevices.yaml — CRD манифест
  • iot/controllers/iotdevice_controller.go — контроллер с полным Reconcile loop
  • main.go — добавлена регистрация IoT схемы и контроллера

Проверка: go build ./... — ошибок нет.

Отклонения от плана (с обоснованием):

  • Добавлен +kubebuilder:validation:Pattern на DeviceID — необходимо чтобы имя Secret iot-{deviceId} было валидным k8s именем. В плане не указано, но без валидации контроллер может сломаться.
  • Использован controllerutil.ContainsFinalizer/AddFinalizer/RemoveFinalizer вместо ручных хелперов — более современный подход из controller-runtime. Старый код не трогал.
  • Добавлена константа IoTDevicePhase (типизированный string) вместо plain string в плане — для type-safety и читаемости.

Задача: Решения по архитектуре IoT MVP

Вопрос 1: Репозиторий — та же или новая?

Анализирую связи:

  • IoT генерирует события → RabbitMQ → event-dispatcher → serverless функции
  • Связь между IoT и sless = на уровне message bus (RabbitMQ), НЕ на уровне кода
  • Общее: концепция user namespace (sless-{hash}), аутентификация (JWT→namespace)
  • Разное: домен (устройства vs функции), протоколы (MQTT vs HTTP), CRD-типы

Вариант A: Та же репа

  • Плюс: общий go.mod, общие утилиты namespace, быстрый старт
  • Плюс: один оператор — проще деплоить для демо
  • Минус: два домена в одной репе — запутает
  • Минус: разные циклы релизов в будущем

Вариант B: Новая репа

  • Плюс: чистое разделение, независимые релизы
  • Минус: дублирование namespace-логики или общая библиотека
  • Минус: overhead для демо слишком большой

Вариант C (мой выбор): Та же репа, изолированная структура

  • Весь IoT-код в директории iot/ на верхнем уровне
  • Свои контроллеры: iot/controllers/
  • Свои CRD: iot/api/v1alpha1/
  • Свой деплоймент (отдельный binary или часть того же оператора)
  • Легко вынести в отдельную репу позже — просто перемещаем iot/
  • Для демо: контроллеры IoT встраиваются в тот же operator binary (один pod)

Reason: связь IoT↔sless через RabbitMQ — слабая. Код не зависит друг от друга. Но для демо удобнее держать вместе. Структура iot/ позволяет легко разделить.

Вопрос 2: Terraform provider — расширять или новый?

Факты:

  • Текущий провайдер: sless (terraform-provider-sless)
  • Ресурсы: sless_function, sless_trigger, sless_service
  • Auth: JWT → namespace

Анализ:

  • Имя "sless" не подходит для IoT-ресурсов (sless_iot_device — странно)
  • Но auth/namespace логика идентична
  • Для демо: расширение существующего — быстрее всего
  • Для прода: нужен единый провайдер nubes (бренд облака) с подресурсами, или отдельный nubes-iot

Мой выбор: расширить текущий для демо

  • Добавить sless_iot_device, sless_iot_rule
  • Имя неидеальное, но для демо ОК
  • Для прода: переименование в nubes — отдельная задача (breaking change)
  • Альтернатива: сразу назвать новый провайдер nubes-iot, но это overhead для демо

Рекомендация пользователю: решить позже, когда IoT станет полноценным сервисом. Для демо — расширяем sless.

Вопрос 3: Scope MVP — что включаем?

Полный IoT-сервис (для справки):

  1. MQTT-брокер ✓
  2. Device Registry ✓
  3. Device Auth ✓
  4. Rules Engine (маршрутизация)
  5. Time-series storage (телеметрия)
  6. Device Shadow/Twin (состояние)
  7. Command Channel (cloud→device)
  8. Dashboard/мониторинг

MVP (демо с возможностью усложнения):

ДА, включаем:

  1. EMQX — деплой через YAML/Helm в кластер
  2. CRD IoTDevice — имя, namespace, credentials (username/password), metadata
  3. IoT-контроллер — reconcile IoTDevice → создаёт MQTT credentials в EMQX через HTTP API
  4. EMQX → RabbitMQ bridge — маршрутизация: MQTT topic → RabbitMQ queue
  5. Включение event-триггеров в sless API (снятие блокировки)
  6. Terraform: sless_iot_device (CRUD)
  7. E2E демо: device → MQTT → function вызывается

НЕТ, откладываем:

  • Device Shadow — усложнение, не нужно для демо
  • Rules Engine — для демо хватит простой маршрутизации topic→queue
  • Time-series storage — функция сама может писать в Postgres
  • Command channel (cloud→device) — второй этап
  • Client certificates — для демо username/password
  • Dashboard — Grafana + метрики EMQX потом

Вопрос 4: Архитектура MVP — как именно работает

Цепочка данных:

IoT Device
  → MQTT connect (username=deviceId, password=deviceSecret)
  → EMQX (topic: {namespace}/telemetry/{deviceId})
    → EMQX Rule + Bridge → RabbitMQ (queue: iot.{namespace}.{topic-pattern})
      → sless event-dispatcher (существующий) → POST body → serverless function
        → function обрабатывает данные

Аутентификация устройств:

  • EMQX HTTP Auth Backend → наш API: POST /internal/mqtt/auth
  • Контроллер при создании IoTDevice → генерирует credentials → хранит в k8s Secret
  • EMQX проверяет при MQTT CONNECT: запрос к нашему API → проверка credentials → ACL (device видит только свой namespace)

Почему EMQX HTTP Auth, а не встроенная БД:

  • При добавлении/удалении устройства не нужно перезагружать EMQX
  • ACL динамический — привязан к namespace
  • Возможность усложнения (certificates, OAuth) без изменения EMQX

Структура файлов (план):

iot/
  api/v1alpha1/
    device_types.go          # CRD IoTDevice
    groupversion_info.go
    zz_generated.deepcopy.go
  controllers/
    device_controller.go     # Reconcile: создаёт credentials, Secret
  internal/
    emqx/
      client.go              # HTTP-клиент к EMQX Management API
    mqtt_auth/
      handler.go             # HTTP Auth Backend для EMQX
  deployments/
    emqx.yaml                # Деплой EMQX в кластер

Что НЕ нужно создавать с нуля:

  • RabbitMQ — есть
  • Event-dispatcher — есть (только включить event triggers)
  • Namespace-логика — есть (переиспользуем)
  • API-сервер (JWT auth, routing) — есть, добавляем IoT-эндпоинты

Вопрос: RabbitMQ vs Kafka для IoT

Контекст

  • RabbitMQ уже развёрнут, event-dispatcher написан под AMQP
  • Пользователь хочет "с прицелом на будущее, без переделок"
  • IoT = потенциально тысячи устройств, миллионы сообщений

Сравнение для IoT

Критерий RabbitMQ Kafka
Модель Push (broker → consumer) Pull (consumer → offset)
Хранение Сообщение удаляется после ack Лог хранится N дней (replay!)
Масштаб до ~50K msg/sec (один node) миллионы msg/sec
Multi-consumer нет (сообщение потреблено = удалено) да (разные consumer groups)
IoT replay невозможен ключевая фича
Операционные затраты проще сложнее (KRaft, partitions)
Per-user изоляция vhost topic prefix, ACL
Уже есть да нет

Для IoT Kafka объективно лучше потому что:

  1. Replay — устройство отправило данные, функция упала → с Kafka можно перечитать. С RabbitMQ — данные потеряны (nack+requeue — не то же самое)
  2. Multi-consumer — одни и те же данные от датчика читают: функция-алерт, функция-аналитика, Trino (!) для запросов. С RabbitMQ нужны exchange fanout + N очередей (костыль)
  3. Time-series природа — IoT телеметрия = append-only лог. Kafka создан именно для этого
  4. Масштаб — 1000 устройств × 1 msg/sec = 1000 msg/sec (RabbitMQ справится). 100K устройств × 10 msg/sec = 1M msg/sec (RabbitMQ не справится)
  5. Trino — Kafka connector для Trino → SQL-запросы по потоку данных в реальном времени. С RabbitMQ этого нет

НО: event-dispatcher сейчас написан под AMQP

Переход = переписать event-dispatcher. Это ~300 строк кода. Не страшно.

Мой вывод

Для IoT — Kafka. Для существующих serverless event-триггеров — RabbitMQ остаётся (менять рабочее не нужно).

Два брокера — нормально:

  • RabbitMQ — существующие event-триггеры sless (task queue: одно сообщение → одна функция)
  • Kafka — IoT data pipeline (stream: данные хранятся, читаются многими)

Это разные паттерны использования. Не нужно одним инструментом решать обе задачи.


Вопрос: может перевести sless тоже на Kafka — пока не обросли?

Думаю вслух

Пользователь прав — переделывать потом больнее. Сейчас event-dispatcher = ~300 строк, event-триггеры ещё даже не включены в API. Идеальный момент для смены.

Для sless event-триггеров: что именно нужно?

Паттерн: сообщение пришло → вызвать ОДНУ функцию → подтвердить/повторить.

Нужно для sless RabbitMQ Kafka
Доставка 1 сообщение → 1 функция нативно (queue) consumer group (работает)
Retry при ошибке nack+requeue / dead letter — нативно нужна логика retry-topic (код)
Dead letter queue встроен нужен отдельный topic + код
Приоритеты сообщений да нет
Задержка доставки (delay) плагин, просто нет нативно

RabbitMQ для task queue объективно удобнее. Kafka для этого работает, но требует больше кода.

НО: два брокера в проде — это боль

  • Два кластера мониторить
  • Два набора алертов
  • Два набора бэкапов
  • Две точки отказа
  • Двойное потребление ресурсов

Варианты

Вариант A: Два брокера (RabbitMQ для sless, Kafka для IoT)

  • Плюс: каждый инструмент для своей задачи
  • Минус: операционная сложность × 2

Вариант B: Kafka для всего

  • Плюс: один брокер, одна инфраструктура
  • Плюс: sless event-dispatcher переписать СЕЙЧАС — пока маленький
  • Минус: retry/DLQ для sless придётся писать руками (~50 строк)
  • Минус: Kafka тяжелее (3 ноды KRaft минимум для прода)

Вариант C: Redpanda вместо Kafka

  • Kafka-совместимый API, но single-binary, легче в ops
  • Но менее зрелый, меньше community

Мой вывод

Kafka для всего. Причины:

  1. Event-триггеры в sless ещё не запущены — переписать сейчас = 0 стоимости миграции
  2. Один брокер вместо двух — проще эксплуатация
  3. Retry через retry-topic — стандартный паттерн, ~50 строк кода
  4. Kafka для sless event-триггеров работает нормально (consumer group, offset commit = тот же ack)
  5. С прицелом: если sless и IoT оба на Kafka — проще интеграция (IoT event → sless function, один bus)

Контраргумент: single point of failure

Пользователь прав

Если Kafka — единственный брокер и он падает → оба сервиса мертвы (sless event-triggers + IoT). Blast radius = вся платформа.

А если два брокера:

  • Kafka упал → IoT не работает, но sless event-триггеры живы (RabbitMQ)
  • RabbitMQ упал → IoT работает, sless event-триггеры лежат
  • Полный outage = нужно чтобы упали ОБА одновременно (маловероятно)

Пересмотр решения

Это классический trade-off: операционная простота vs отказоустойчивость.

Для managed service платформы — отказоустойчивость важнее. Клиент платит за uptime.

Финальное решение: ДВА брокера

  • RabbitMQ → sless event-триггеры (уже написан, проще для task queue, независимый)
  • Kafka → IoT pipeline (replay, multi-consumer, масштаб)
  • Изоляция fault domains: падение одного не убивает другой сервис

Операционная сложность двух брокеров — приемлемая цена за изоляцию. Мониторинг/алерты — решаемо (Prometheus + Grafana для обоих).


Реальность: кубер сломан, выходные, нет облачных сервисов

Ситуация

  • Реалм пользователя не создаёт managed-сервисы (баг/инцидент)
  • До понедельника никого нет (шабат/выходные)
  • Kafka может оказаться в другом реалме, чем RabbitMQ
  • Нужно работать с тем что есть СЕЙЧАС

Мои мысли

Вариант A: Делаем MVP на RabbitMQ (который есть)

  • Плюс: RabbitMQ уже работает, ничего разворачивать не нужно
  • Плюс: event-dispatcher уже написан под AMQP
  • Плюс: можно прямо сейчас начать IoT-часть (CRD, контроллер, EMQX, bridge)
  • Плюс: demo будет работать к понедельнику
  • Минус: потом нужна миграция EMQX→Kafka вместо EMQX→RabbitMQ bridge
  • НО: мост MQTT→broker — это конфиг EMQX, а не наш код. Переключить EMQX bridge с RabbitMQ на Kafka = смена конфига, не переписывание

Вариант B: Поднять Kafka руками в кубере (Strimzi/Bitnami Helm)

  • Плюс: правильная архитектура с самого начала
  • Минус: Kafka в k8s = тяжело (3 ноды KRaft, storage, сетевые проблемы)
  • Минус: если реалм глючит — может и Kafka не развернуться
  • Минус: потом всё равно мигрировать на managed

Вариант C (мой выбор): MVP на RabbitMQ сейчас, архитектура ready for Kafka

Суть: делаем IoT bridge через абстракцию, не привязываясь к конкретному брокеру.

EMQX → [bridge config] → RabbitMQ (сейчас)
                        → Kafka (потом, смена конфига)

IoT event consumer → [interface] → POST → function
  сейчас: event-dispatcher (AMQP) уже есть
  потом:  iot-consumer (Kafka) — отдельный сервис

Ключевое: НАША кодовая база НЕ зависит от выбора брокера.

  • CRD IoTDevice — не зависит
  • IoT контроллер — не зависит
  • MQTT auth — не зависит
  • EMQX — bridge настраивается конфигом (RabbitMQ или Kafka)
  • Единственная точка замены: consumer, который читает из брокера и POST в функцию

Что менять при переходе RabbitMQ → Kafka

  1. EMQX bridge config: rabbitmqkafka (конфиг, не код)
  2. Consumer: отдельный iot-event-consumer вместо reuse event-dispatcher (~200 строк Go)
  3. Kafka deployment: managed или Strimzi

Всё. Наш IoT-оператор, CRD, device auth — не меняются вообще.