# IoT MVP — План реализации для Sonnet > **Автор плана**: GitHub Copilot (Claude Opus 4.6) > **Дата**: 2026-04-04 > **Исполнитель**: Claude Sonnet > **Ход рассуждений**: `doc/thinking/2026-04-04.md` --- ## Контекст Платформа **sless** — managed serverless functions. Нужно добавить **managed IoT service** как демо с возможностью усложнения. ### Согласованные решения | Вопрос | Решение | Обоснование | |--------|---------|-------------| | Репозиторий | Та же репа, код в `iot/` | Легко вынести потом, удобно для демо | | Message broker IoT | RabbitMQ (MVP), потом Kafka | RabbitMQ уже есть, архитектура broker-agnostic | | Message broker sless | RabbitMQ (не трогать) | Работает, отдельный fault domain | | MQTT-брокер | EMQX, деплой plain YAML | Не Helm, не Operator — достаточно для демо | | IoT-логика | CRD + controller (Go operator) | Консистентно с sless, сразу правильно | | Terraform | Расширяем текущий sless provider | Для демо ОК, переименование потом | | Device auth | EMQX HTTP Auth Backend → наш API | Динамическое добавление устройств | ### Что НЕ делаем (отложено) - Device Shadow / Digital Twin - Rules Engine (для демо — простая маршрутизация topic→queue) - Time-series storage (функция сама пишет в Postgres) - Cloud→Device commands - Client certificates (для демо: username/password) - Dashboard --- ## Существующая архитектура (НЕ ТРОГАТЬ) ``` Go module: gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless API Group: sless.kube5s.ru/v1alpha1 Namespace pattern: sless-{sha256(jwt.sub)[:16]} Контроллеры (main.go, строки 153-184): - FunctionReconciler - ServiceReconciler - TriggerReconciler - FunctionJobReconciler API-сервер: internal/api/router.go (gorilla/mux), порт cfg.APIPort - JWT auth middleware → namespace validation - Routes: /v1/namespaces/{namespace}/functions|services|triggers|jobs Trigger types: "http", "cron", "event" (api/v1alpha1/trigger_types.go) Event-dispatcher: services/event-dispatcher/ (AMQP consumer → POST в функцию) RabbitMQ: deployments/k8s/rabbitmq.yaml (namespace: sless) ``` --- ## Целевая архитектура MVP ``` IoT Device → MQTT connect (username=deviceId, password=deviceSecret) → EMQX (topic: {namespace}/telemetry/{deviceId}) → EMQX RabbitMQ Bridge → RabbitMQ (queue: iot.{namespace}) → event-dispatcher (существующий!) → POST → serverless function → function обрабатывает данные Аутентификация устройств: EMQX HTTP Auth Plugin → GET http://sless-iot-auth.sless.svc:8080/mqtt/auth → проверка credentials из k8s Secret → ACL (только свой namespace в topics) ``` --- ## Этапы реализации ### Этап 1: CRD IoTDevice и контроллер **Цель**: зарегистрировать IoT-устройство через CRD, автоматически создать credentials. #### 1.1. Создать CRD типы Файл: `iot/api/v1alpha1/device_types.go` ```go package v1alpha1 import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) // IoTDeviceSpec — спецификация IoT-устройства type IoTDeviceSpec struct { // DeviceID — уникальный идентификатор устройства внутри namespace DeviceID string `json:"deviceId"` // Metadata — произвольные метаданные устройства (модель, локация и т.д.) // +optional Metadata map[string]string `json:"metadata,omitempty"` // Enabled — активно ли устройство (может подключаться к MQTT) // +kubebuilder:default=true Enabled bool `json:"enabled"` } // IoTDeviceStatus — статус IoT-устройства type IoTDeviceStatus struct { // Phase — текущее состояние: Pending, Active, Disabled, Error Phase string `json:"phase,omitempty"` // MQTTUsername — имя пользователя для подключения к MQTT MQTTUsername string `json:"mqttUsername,omitempty"` // SecretName — имя k8s Secret с credentials SecretName string `json:"secretName,omitempty"` // TopicPrefix — разрешённый prefix для MQTT topics TopicPrefix string `json:"topicPrefix,omitempty"` // LastConnected — время последнего подключения (заполняется auth-сервисом) // +optional LastConnected *metav1.Time `json:"lastConnected,omitempty"` // Message — человекочитаемое сообщение о статусе Message string `json:"message,omitempty"` } // +kubebuilder:object:root=true // +kubebuilder:subresource:status // +kubebuilder:printcolumn:name="DeviceID",type=string,JSONPath=`.spec.deviceId` // +kubebuilder:printcolumn:name="Phase",type=string,JSONPath=`.status.phase` // +kubebuilder:printcolumn:name="Enabled",type=boolean,JSONPath=`.spec.enabled` type IoTDevice struct { metav1.TypeMeta `json:",inline"` metav1.ObjectMeta `json:"metadata,omitempty"` Spec IoTDeviceSpec `json:"spec,omitempty"` Status IoTDeviceStatus `json:"status,omitempty"` } // +kubebuilder:object:root=true type IoTDeviceList struct { metav1.TypeMeta `json:",inline"` metav1.ListMeta `json:"metadata,omitempty"` Items []IoTDevice `json:"items"` } ``` Файл: `iot/api/v1alpha1/groupversion_info.go` ```go package v1alpha1 import ( "k8s.io/apimachinery/pkg/runtime/schema" "sigs.k8s.io/controller-runtime/pkg/scheme" ) var ( // GroupVersion — API group для IoT ресурсов // ВАЖНО: отдельный group от sless.kube5s.ru — для будущего разделения GroupVersion = schema.GroupVersion{Group: "iot.kube5s.ru", Version: "v1alpha1"} SchemeBuilder = &scheme.Builder{GroupVersion: GroupVersion} AddToScheme = SchemeBuilder.AddToScheme ) func init() { SchemeBuilder.Register(&IoTDevice{}, &IoTDeviceList{}) } ``` **ВАЖНО**: API group `iot.kube5s.ru` — отдельная от `sless.kube5s.ru`. Причина: при разделении на отдельную репу CRD не будет конфликтовать. #### 1.2. Сгенерировать deepcopy и CRD манифесты ```bash # Из корня проекта: controller-gen object paths=./iot/api/v1alpha1/... controller-gen crd paths=./iot/api/v1alpha1/... output:crd:dir=iot/config/crd/bases ``` #### 1.3. Создать контроллер IoTDevice Файл: `iot/controllers/iotdevice_controller.go` Логика Reconcile: 1. Получить IoTDevice из пришедшего запроса 2. Если `DeletionTimestamp != nil` → удалить Secret, убрать finalizer 3. Добавить finalizer `iot.kube5s.ru/device-cleanup` если нет 4. Если `spec.enabled == false`: - Установить `status.phase = "Disabled"` - НЕ удалять Secret (устройство может быть включено обратно) 5. Если Secret не существует: - Сгенерировать пароль (32 байта crypto/rand → hex) - MQTTUsername = `{namespace}_{deviceId}` (namespace включён для уникальности MQTT username) - Создать Secret `iot-{deviceId}` в том же namespace с полями: - `mqtt-username`: `{namespace}_{deviceId}` - `mqtt-password`: сгенерированный пароль - Установить OwnerReference на IoTDevice (каскадное удаление) 6. Заполнить status: - `phase = "Active"` (или "Disabled" если !enabled) - `mqttUsername = {namespace}_{deviceId}` - `secretName = iot-{deviceId}` - `topicPrefix = {namespace}/` (устройство может публиковать только в topics с этим prefix) #### 1.4. Зарегистрировать контроллер в main.go Добавить в `main.go` после существующих SetupWithManager вызовов: ```go if err = (&iotcontrollers.IoTDeviceReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), Log: ctrl.Log.WithName("controllers").WithName("IoTDevice"), }).SetupWithManager(mgr); err != nil { setupLog.Error(err, "unable to create controller", "controller", "IoTDevice") os.Exit(1) } ``` Добавить в import: ```go iotv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/api/v1alpha1" iotcontrollers "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/controllers" ``` Добавить schema registration в init/scheme: ```go utilruntime.Must(iotv1alpha1.AddToScheme(scheme)) ``` --- ### Этап 2: MQTT Auth Service **Цель**: EMQX при каждом MQTT CONNECT проверяет credentials через наш HTTP-сервис. #### 2.1. Создать auth handler Файл: `iot/internal/mqttauth/mqtt_auth_handler.go` HTTP-сервис, отдельный порт (например 8081) или sub-router в основном API. **Эндпоинт**: `POST /mqtt/auth` (вызывается EMQX HTTP Auth Plugin) EMQX присылает JSON: ```json { "username": "sless-abc123def456_sensor-01", "password": "hex-encoded-secret", "clientid": "...", "peerhost": "10.0.0.5" } ``` Логика: 1. Распарсить username: `{namespace}_{deviceId}` 2. Найти Secret `iot-{deviceId}` в namespace `{namespace}` 3. Сравнить password с `mqtt-password` из Secret (constant-time comparison!) 4. Если совпало: - Проверить что IoTDevice существует и `enabled == true` - Вернуть 200 + JSON с ACL: ```json { "result": "allow", "is_superuser": false, "acl": [ {"permission": "allow", "action": "publish", "topic": "{namespace}/#"}, {"permission": "allow", "action": "subscribe", "topic": "{namespace}/#"}, {"permission": "deny", "action": "all", "topic": "#"} ] } ``` 5. Если не совпало → вернуть 200 + `{"result": "deny"}` **ВАЖНО**: НЕ возвращать 401/403 — EMQX интерпретирует HTTP-ошибки как "ignore this backend, try next". Всегда 200, result = "allow"/"deny". #### 2.2. Запустить auth-сервис Два варианта (решить при реализации): - **Вариант A**: отдельный binary `iot/cmd/mqtt-auth/main.go` + Deployment - **Вариант B**: добавить route в существующий API-сервер (проще для демо) Для демо — **вариант B**: добавить роут `/internal/mqtt/auth` в `internal/api/router.go`. Prefix `/internal/` = не защищён JWT (доступен только из кластера). --- ### Этап 3: Развёртывание EMQX **Цель**: MQTT-брокер, принимающий подключения от IoT-устройств. #### 3.1. Создать YAML Файл: `deployments/k8s/emqx.yaml` По аналогии с `deployments/k8s/rabbitmq.yaml`: ```yaml # Деплой EMQX MQTT-брокера для IoT-сервиса # 2026-04-04 apiVersion: v1 kind: ConfigMap metadata: name: emqx-config namespace: sless data: # HTTP Auth Backend — аутентификация устройств через наш API EMQX_AUTH__HTTP__AUTH_REQ__URL: "http://sless-api.sless.svc:8080/internal/mqtt/auth" EMQX_AUTH__HTTP__AUTH_REQ__METHOD: "post" EMQX_AUTH__HTTP__AUTH_REQ__CONTENT_TYPE: "json" # RabbitMQ Bridge (MVP) — маршрутизация MQTT → RabbitMQ # ПОТОМ заменится на Kafka bridge — только этот ConfigMap EMQX_BRIDGE__RABBIT__SERVER: "amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672" --- apiVersion: apps/v1 kind: Deployment metadata: name: emqx namespace: sless spec: replicas: 1 selector: matchLabels: app: emqx template: metadata: labels: app: emqx spec: containers: - name: emqx image: emqx/emqx:5.5.1 ports: - name: mqtt containerPort: 1883 - name: mqttssl containerPort: 8883 - name: ws containerPort: 8083 - name: dashboard containerPort: 18083 resources: requests: memory: "256Mi" cpu: "100m" limits: memory: "512Mi" cpu: "500m" envFrom: - configMapRef: name: emqx-config readinessProbe: tcpSocket: port: 1883 initialDelaySeconds: 15 periodSeconds: 10 --- apiVersion: v1 kind: Service metadata: name: emqx namespace: sless spec: selector: app: emqx ports: - name: mqtt port: 1883 targetPort: 1883 - name: ws port: 8083 targetPort: 8083 - name: dashboard port: 18083 targetPort: 18083 ``` **ВНИМАНИЕ**: конфиг EMQX 5.x сильно отличается от 4.x. При реализации: - Проверить актуальный формат env-переменных для EMQX 5.5 - HTTP Auth Plugin в EMQX 5 конфигурируется через Dashboard API или файл `etc/emqx.conf` - Возможно понадобится volume mount для `emqx.conf` вместо env #### 3.2. Настроить EMQX Rule + RabbitMQ Bridge EMQX Rule Engine (конфигурируется через EMQX HTTP API или Dashboard): ``` Rule SQL: SELECT * FROM '{namespace}/+/+' Action: Bridge to RabbitMQ Exchange: amq.topic Routing Key: iot.{namespace} Queue: iot.{namespace}.telemetry ``` Для MVP — можно сконфигурировать одно правило вручную или через EMQX REST API при старте (init-container или наш контроллер). **ВАЖНО для Sonnet**: EMQX 5.x использует `bridges` API — изучить документацию EMQX 5.5: - `POST /api/v5/bridges` — создание bridge - `POST /api/v5/rules` — создание правил - Или файл `etc/emqx.conf` (HOCON формат) --- ### Этап 4: API-эндпоинты для IoT **Цель**: REST API для управления устройствами (Terraform provider будет вызывать их). #### 4.1. Добавить IoT-роуты в router.go Файл: `internal/api/router.go` Новые routes (защищены JWT, как существующие): ``` POST /v1/namespaces/{namespace}/iot/devices → CreateIoTDevice GET /v1/namespaces/{namespace}/iot/devices → ListIoTDevices GET /v1/namespaces/{namespace}/iot/devices/{name} → GetIoTDevice DELETE /v1/namespaces/{namespace}/iot/devices/{name} → DeleteIoTDevice PATCH /v1/namespaces/{namespace}/iot/devices/{name} → UpdateIoTDevice (enable/disable) ``` Internal route (без JWT, только для EMQX из кластера): ``` POST /internal/mqtt/auth → MQTTAuth ``` #### 4.2. Создать IoT handler Файл: `iot/internal/api/iot_device_handler.go` (или добавить в `internal/api/handler/`) Хендлеры = тонкая обёртка над k8s API: - `CreateIoTDevice`: создаёт IoTDevice CRD объект → контроллер reconcile → Secret - `GetIoTDevice`: читает IoTDevice CRD + возвращает credentials из Secret - `DeleteIoTDevice`: удаляет IoTDevice CRD → контроллер cleanup через finalizer - `ListIoTDevices`: list IoTDevice в namespace - `UpdateIoTDevice`: patch spec.enabled **GET /devices/{name}** должен возвращать credentials (mqtt_username, mqtt_password) из Secret. Они нужны пользователю для конфигурации устройства. Credentials возвращаются **только при GET**, не хранятся в CRD status. --- ### Этап 5: Связь MQTT → Serverless Function **Цель**: IoT-устройство отправляет MQTT → вызывается serverless function. #### 5.1. Цепочка Пользователь создаёт через Terraform: 1. `sless_iot_device` → IoTDevice CRD → MQTT credentials 2. `sless_function` → Function → готовая serverless функция 3. `sless_trigger` type=event, queue="iot.{namespace}.telemetry" → event-dispatcher подписывается Event-dispatcher (уже работает!) читает из RabbitMQ queue → POST в функцию. #### 5.2. Автоматическое создание RabbitMQ queue IoT-контроллер при reconcile должен обеспечить (ensure) существование queue `iot.{namespace}.telemetry` в RabbitMQ. Варианты: - **A**: Контроллер создаёт queue через RMQ Management API (HTTP) — явно - **B**: Queue создаётся автоматически EMQX bridge + event-dispatcher consumer (declare on consume) Для MVP — **вариант B**: и EMQX bridge, и event-dispatcher делают QueueDeclare — кто первый, тот и создаст. Дурак-proof. --- ### Этап 6: Terraform Provider **Цель**: управление IoT-устройствами через Terraform. Terraform provider sless находится в **отдельной репе**. Нужно расширить его. Новый ресурс: `sless_iot_device` ```hcl resource "sless_iot_device" "sensor_01" { namespace = sless_namespace.my_ns.name name = "temperature-sensor" device_id = "sensor-01" enabled = true metadata = { model = "DHT22" location = "room-1" } } output "mqtt_username" { value = sless_iot_device.sensor_01.mqtt_username } output "mqtt_password" { value = sless_iot_device.sensor_01.mqtt_password sensitive = true } ``` CRUD маппинг: - Create → `POST /v1/namespaces/{ns}/iot/devices` - Read → `GET /v1/namespaces/{ns}/iot/devices/{name}` - Update → `PATCH /v1/namespaces/{ns}/iot/devices/{name}` - Delete → `DELETE /v1/namespaces/{ns}/iot/devices/{name}` --- ### Этап 7: E2E Demo **Цель**: показать полную цепочку device → function. #### Demo Terraform: ```hcl # 1. Функция-обработчик IoT данных resource "sless_function" "iot_handler" { namespace = var.namespace name = "iot-handler" runtime = "python3.11" source_dir = "./iot-handler" } # 2. Event trigger: подписка на IoT queue resource "sless_trigger" "iot_events" { namespace = var.namespace name = "iot-telemetry" type = "event" function_ref = sless_function.iot_handler.name queue = "iot.${var.namespace}.telemetry" enabled = true } # 3. IoT устройство resource "sless_iot_device" "sensor" { namespace = var.namespace name = "demo-sensor" device_id = "sensor-001" enabled = true } output "mqtt_host" { value = "emqx.sless.svc.cluster.local" } output "mqtt_username" { value = sless_iot_device.sensor.mqtt_username } output "mqtt_password" { value = sless_iot_device.sensor.mqtt_password sensitive = true } ``` #### Demo Python IoT handler (`iot-handler/handler.py`): ```python def handler(event, context): """Обработчик IoT-телеметрии. Вызывается event-dispatcher при новом MQTT сообщении.""" import json data = json.loads(event["body"]) print(f"Telemetry from device: {data}") return {"statusCode": 200, "body": json.dumps({"processed": True})} ``` #### Demo MQTT client (для тестирования): ```bash # Отправить MQTT сообщение (mosquitto_pub) mosquitto_pub \ -h emqx.sless.svc.cluster.local \ -p 1883 \ -u "sless-abc123_sensor-001" \ -P "generated-password" \ -t "sless-abc123/telemetry/sensor-001" \ -m '{"temperature": 22.5, "humidity": 65}' ``` --- ## Структура файлов (итого) ``` iot/ api/v1alpha1/ device_types.go # CRD IoTDevice groupversion_info.go # API group iot.kube5s.ru/v1alpha1 zz_generated.deepcopy.go # сгенерировано controller-gen config/ crd/bases/ # сгенерированные CRD YAML controllers/ iotdevice_controller.go # Reconcile: Secret, credentials internal/ mqttauth/ mqtt_auth_handler.go # HTTP Auth Backend для EMQX deployments/k8s/ emqx.yaml # EMQX deployment (новый файл) internal/api/ handler/ iot_devices.go # REST handlers для IoT devices (новый файл) router.go # + IoT routes (редактирование) main.go # + IoTDevice controller registration (редактирование) examples/ IOT/ # Demo пример main.tf handler.py ``` --- ## Порядок выполнения ``` 1. CRD types + deepcopy + manifests (iot/api/) 2. IoTDevice controller (iot/controllers/) 3. Register controller в main.go (main.go) 4. CRD apply в кластер (kubectl apply) 5. MQTT Auth handler (iot/internal/mqttauth/) 6. REST API endpoints для IoT (internal/api/) 7. EMQX deployment YAML (deployments/k8s/emqx.yaml) 8. EMQX конфигурация (auth backend + RMQ bridge) 9. Terraform provider resource (отдельная репа) 10. E2E demo (examples/IOT/) 11. Тестирование: device → MQTT → function ``` --- ## Что менять при переходе на Kafka (потом) | Компонент | Изменение | |-----------|-----------| | EMQX bridge config | `rabbitmq` → `kafka` (ConfigMap) | | Consumer | Новый `iot-event-consumer` (~200 строк Go) вместо event-dispatcher | | Kafka deploy | Managed сервис или Strimzi в кластере | | CRD / Controller | **БЕЗ ИЗМЕНЕНИЙ** | | MQTT Auth | **БЕЗ ИЗМЕНЕНИЙ** | | API endpoints | **БЕЗ ИЗМЕНЕНИЙ** | | Terraform | **БЕЗ ИЗМЕНЕНИЙ** | --- ## Правила для Sonnet 1. **Читай `doc/thinking/2026-04-04.md`** — там полный ход рассуждений и обоснования 2. **Читай `.github/copilot-instructions.md`** — правила проекта 3. **НЕ трогай** существующий код sless (controllers/, services/, internal/) без крайней необходимости 4. **Комментарии обязательны** — дата, назначение функций, "почему" для нетривиальной логики 5. **Именование** — уникальные осмысленные имена (iot_device_handler, NOT handler) 6. **Пиши в `doc/thinking/`** свои мысли при решении 7. **EMQX конфиг** — обязательно проверить актуальный формат для EMQX 5.5.x 8. **Security**: constant-time password comparison, не логировать credentials, ACL-изоляция по namespace --- ## Справка для нового агента (чтобы не искать) ### Terraform provider - Расположен **в этой же репе**: `terraform/provider/` - Go module: `terraform-provider-sless` (свой go.mod: `terraform/provider/go.mod`) - Структура: ``` terraform/provider/ main.go go.mod internal/ client/ # HTTP-клиент к sless API provider/ # provider.go — регистрация ресурсов resources/ # function_resource.go, trigger_resource.go, service_resource.go, job_resource.go hack/ build-and-publish.sh ``` - Новый ресурс `sless_iot_device` добавлять в `terraform/provider/internal/resources/iot_device_resource.go` - Зарегистрировать в `terraform/provider/internal/provider/provider.go` ### controller-gen - Бинарь: `bin/controller-gen` (в корне репы) - Генерация deepcopy: `./bin/controller-gen object paths=./iot/api/v1alpha1/...` - Генерация CRD: `./bin/controller-gen crd paths=./iot/api/v1alpha1/... output:crd:dir=iot/config/crd/bases` ### Как зарегистрированы существующие контроллеры (пример из main.go) ```go // main.go строки ~153-184 if err = (&controllers.FunctionReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), S3Client: s3Client, Config: cfg, }).SetupWithManager(mgr); err != nil { setupLog.Error(err, "unable to create controller", "controller", "Function") os.Exit(1) } ``` Новый IoT контроллер регистрируется аналогично, после этого блока. ### Scheme registration (main.go) ```go // Существующие: utilruntime.Must(slessv1alpha1.AddToScheme(scheme)) // Добавить: utilruntime.Must(iotv1alpha1.AddToScheme(scheme)) ``` ### Существующий handler паттерн (internal/api/handler/handler.go) ```go type Handler struct { K8s client.Client Scheme *runtime.Scheme S3 *minio.Client PG *sql.DB Log logr.Logger Config *config.Config } ``` IoT-хендлеры добавлять как методы того же Handler или создать отдельный IoTHandler. ### RabbitMQ connection string - Dev: `amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/` - Secret: `sless-operator-secret`, ключ `RABBITMQ_URL` ### Event-dispatcher — как он подписывается на queue - Файл: `services/event-dispatcher/dispatcher.go` - QueueDeclare (durable=true) при Subscribe - Consumer на queue → POST в `http://{functionRef}.{namespace}.svc.cluster.local:8080/` - Ack при 2xx, Nack+requeue при ошибке ### EMQX 5.x — ключевые отличия от 4.x - Конфиг: HOCON формат в `/opt/emqx/etc/emqx.conf`, NOT env variables для plugins - Auth: конфигурируется через `authentication` секцию в emqx.conf или REST API `POST /api/v5/authentication` - Bridges: REST API `POST /api/v5/bridges` или секция `bridges` в emqx.conf - Rules: REST API `POST /api/v5/rules` - Dashboard: порт 18083, default login admin/public - **Sonnet должен зайти на https://www.emqx.io/docs/en/v5.5/ и проверить формат конфигурации**