Этап 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 ./... — ошибок нет
33 KiB
Лог мышления — 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
Что нового нужно создать
- MQTT-брокер — EMQX (есть multi-tenancy, WebSocket, rule engine, k8s operator)
- IoT Device Registry — CRD
IoTDevice(имя, credentials, metadata, user namespace) - IoT Bridge — сервис MQTT→RabbitMQ per user-namespace (или EMQX rule engine)
- Device Auth — MQTT username/password или client certificates, хранение в k8s secrets per namespace
- 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 constapi/v1alpha1/groupversion_info.go— паттерн регистрации API groupcontrollers/function_controller.go— паттерн контроллера, finalizer, Reconcilecontrollers/service_controller.go— второй контроллер для сравненияinternal/api/router.go— регистрация маршрутов, auth middlewaremain.go— как регистрируются контроллеры, как подключается APIgo.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.Read → hex.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)
iot/api/v1alpha1/device_types.go— CRD IoTDeviceiot/api/v1alpha1/groupversion_info.go— API group iot.kube5s.ru/v1alpha1iot/controllers/iotdevice_controller.go— контроллерiot/config/crd/bases/— директория для CRD YAML (создаётся controller-gen через SSH)- Обновление
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 handlersinternal/api/router.go— добавить IoT routesdeployments/k8s/emqx.yaml— EMQX deployment c emqx.conf ConfigMap (только HTTP auth)iot/cmd/mqtt-bridge/main.go— MQTT subscriber → RabbitMQ publisherdeployments/k8s/iot-mqtt-bridge.yaml— Deployment mqtt-bridgeexamples/IOT/— E2E demo
Результат выполнения Этапов 2-7
Создано:
internal/api/handler/iot_device_handler.go— MQTT auth + IoT CRUD handlersinternal/api/router.go— IoT routes +/internal/mqtt/authdeployments/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-bridgeexamples/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 attacksGetIoTDevice— единственный endpoint с mqtt_password (security by design)
Создано:
iot/api/v1alpha1/device_types.go— CRD IoTDevice с IoTDevicePhase константамиiot/api/v1alpha1/groupversion_info.go— API groupiot.kube5s.ru/v1alpha1iot/api/v1alpha1/zz_generated.deepcopy.go— сгенерирован controller-geniot/config/crd/bases/iot.kube5s.ru_iotdevices.yaml— CRD манифестiot/controllers/iotdevice_controller.go— контроллер с полным Reconcile loopmain.go— добавлена регистрация IoT схемы и контроллера
Проверка: go build ./... — ошибок нет.
Отклонения от плана (с обоснованием):
- Добавлен
+kubebuilder:validation:PatternнаDeviceID— необходимо чтобы имя Secretiot-{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-сервис (для справки):
- MQTT-брокер ✓
- Device Registry ✓
- Device Auth ✓
- Rules Engine (маршрутизация)
- Time-series storage (телеметрия)
- Device Shadow/Twin (состояние)
- Command Channel (cloud→device)
- Dashboard/мониторинг
MVP (демо с возможностью усложнения):
ДА, включаем:
- ✅ EMQX — деплой через YAML/Helm в кластер
- ✅ CRD
IoTDevice— имя, namespace, credentials (username/password), metadata - ✅ IoT-контроллер — reconcile IoTDevice → создаёт MQTT credentials в EMQX через HTTP API
- ✅ EMQX → RabbitMQ bridge — маршрутизация: MQTT topic → RabbitMQ queue
- ✅ Включение event-триггеров в sless API (снятие блокировки)
- ✅ Terraform:
sless_iot_device(CRUD) - ✅ 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 объективно лучше потому что:
- Replay — устройство отправило данные, функция упала → с Kafka можно перечитать. С RabbitMQ — данные потеряны (nack+requeue — не то же самое)
- Multi-consumer — одни и те же данные от датчика читают: функция-алерт, функция-аналитика, Trino (!) для запросов. С RabbitMQ нужны exchange fanout + N очередей (костыль)
- Time-series природа — IoT телеметрия = append-only лог. Kafka создан именно для этого
- Масштаб — 1000 устройств × 1 msg/sec = 1000 msg/sec (RabbitMQ справится). 100K устройств × 10 msg/sec = 1M msg/sec (RabbitMQ не справится)
- 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 для всего. Причины:
- Event-триггеры в sless ещё не запущены — переписать сейчас = 0 стоимости миграции
- Один брокер вместо двух — проще эксплуатация
- Retry через retry-topic — стандартный паттерн, ~50 строк кода
- Kafka для sless event-триггеров работает нормально (consumer group, offset commit = тот же ack)
- С прицелом: если 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
- EMQX bridge config:
rabbitmq→kafka(конфиг, не код) - Consumer: отдельный iot-event-consumer вместо reuse event-dispatcher (~200 строк Go)
- Kafka deployment: managed или Strimzi
Всё. Наш IoT-оператор, CRD, device auth — не меняются вообще.