# 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/ и проверить формат конфигурации** --- # ПЛАН: Telemetry Pipeline — Postgres → REST API → UI > **Автор плана**: GitHub Copilot (Claude Opus 4.6) > **Дата**: 2026-04-05 > **Исполнитель**: Claude Sonnet > **Ветка**: `iot-pg-telemetry` > **Предусловия**: все компоненты до этого этапа РЕАЛИЗОВАНЫ и задеплоены (CRD, controller, EMQX, mqtt-bridge, IoT Console UI) --- ## Цель Полная цепочка: IoT устройство (или эмулятор в UI) → MQTT → INSERT в Postgres → REST API → отображение в таблице на вкладке «Телеметрия» в IoT Console. **User story**: юзер входит токеном, регистрирует устройство, запускает эмулятор (рандомные temp/humidity), переходит на вкладку Телеметрия и видит таблицу с данными: время | устройство | payload. --- ## Что уже готово (НЕ ТРОГАТЬ без крайней необходимости) | Компонент | Файл(ы) | Статус | |-----------|---------|--------| | IoTDevice CRD + types | `iot/api/v1alpha1/device_types.go` | ✅ | | IoTDevice controller | `iot/controllers/iotdevice_controller.go` | ✅ | | MQTT Auth + ACL | `internal/api/handler/iot_device_handler.go` | ✅ | | IoT API CRUD | `internal/api/router.go` + handler | ✅ | | EMQX deploy | `deployments/k8s/emqx.yaml` | ✅ | | mqtt-bridge MQTT→RabbitMQ | `iot/cmd/mqtt-bridge/main.go` | ✅ | | IoT Console UI | `internal/api/ui/iot-console.html` | ✅ | | TLS (HTTPS + WSS) | `deployments/k8s/emqx-ws-ingress.yaml` | ✅ | | Nubes branding | UI CSS | ✅ | | Existing Postgres (invocations) | `deployments/k8s/postgres.yaml` | ✅ | --- ## Архитектурное решение (принято 2026-04-04) Подробности: `doc/decisions/iot-telemetry-storage-2026-04-04.md` - **Один Postgres инстанс** для IoT (отдельный от sless Postgres для invocations) - **Отдельная DATABASE per tenant** (не одна таблица с tenant_id!) - Tenant DB: `tenant_{namespace_hash}`, User: `tenant_{namespace_hash}`, Password: UUID в k8s Secret - Таблица: `iot_telemetry(id BIGSERIAL, device_id TEXT, ts TIMESTAMPTZ, payload JSONB)` - Клиент читает ТОЛЬКО через REST API, не через прямой доступ к Postgres --- ## Шаги реализации (порядок критичен!) ### ШАГ 1: Postgres Deployment для IoT (namespace: sless) **Файл**: `deployments/k8s/iot-postgres.yaml` **Почему отдельный от sless postgres**: разные данные, разная нагрузка. Sless postgres хранит invocations логи. IoT postgres хранит телеметрию — может быть значительно больше по объёму. **Почему в namespace `sless`, а НЕ `iot`**: пока всё живёт в одном namespace `sless`. Отдельный namespace `iot` только усложнит без выигрыша. Отдельный Deployment с другим именем (`iot-postgres`) достаточно для изоляции. **YAML манифест**: ```yaml apiVersion: v1 kind: Secret metadata: name: iot-postgres-secret namespace: sless stringData: POSTGRES_PASSWORD: "iot-pg-super-2026" POSTGRES_USER: "iot_admin" POSTGRES_DB: "iot_platform" IOT_PG_DSN: "postgresql://iot_admin:iot-pg-super-2026@iot-postgres.sless.svc:5432/iot_platform?sslmode=disable" --- apiVersion: apps/v1 kind: Deployment metadata: name: iot-postgres namespace: sless spec: replicas: 1 selector: matchLabels: app: iot-postgres template: metadata: labels: app: iot-postgres spec: containers: - name: postgres image: postgres:16-alpine ports: - containerPort: 5432 envFrom: - secretRef: name: iot-postgres-secret resources: requests: memory: "128Mi" cpu: "100m" limits: memory: "512Mi" cpu: "500m" --- apiVersion: v1 kind: Service metadata: name: iot-postgres namespace: sless spec: selector: app: iot-postgres ports: - port: 5432 targetPort: 5432 ``` **Действие**: `kubectl apply -f deployments/k8s/iot-postgres.yaml` **Проверка**: `kubectl exec -n sless deploy/iot-postgres -- psql -U iot_admin -d iot_platform -c "SELECT 1"` --- ### ШАГ 2: Go-пакет IoT Postgres storage **Файл**: `internal/storage/iotpg/iot_telemetry_store.go` **Назначение**: управление tenant databases + CRUD телеметрии. Один пакет, одна структура. **Структура**: ```go package iotpg type IoTPostgresStore struct { adminDB *sql.DB // подключение к iot_platform (суперюзер) tenants sync.Map // кэш *sql.DB per tenant namespace log *slog.Logger } type TelemetryRow struct { ID int64 `json:"id"` DeviceID string `json:"device_id"` Ts time.Time `json:"ts"` Payload json.RawMessage `json:"payload"` } ``` **Методы (все обязательные)**: 1. `New(adminDSN string, log *slog.Logger) (*IoTPostgresStore, error)` — подключение к iot_platform DB. При старте создать таблицу tenant_credentials если не существует. 2. `EnsureTenantDB(ctx, namespace string) error` — создать DATABASE + USER + таблицу если не существуют: - `SELECT 1 FROM pg_database WHERE datname = 'tenant_{ns}'` - Если нет: `CREATE USER tenant_{ns} WITH PASSWORD '{uuid}'` - `CREATE DATABASE tenant_{ns} OWNER tenant_{ns}` - Подключиться к tenant_{ns} и: `CREATE TABLE IF NOT EXISTS iot_telemetry (...)` - Индекс: `CREATE INDEX IF NOT EXISTS idx_iot_telemetry_device_ts ON iot_telemetry(device_id, ts DESC)` - Записать пароль в tenant_credentials 3. `InsertTelemetry(ctx, namespace, deviceID string, payload json.RawMessage) error` — INSERT одной записи в tenant DB 4. `QueryTelemetry(ctx, namespace string, deviceID string, limit int) ([]TelemetryRow, error)` — SELECT из tenant DB 5. `Close() error` **Таблица tenant_credentials** (в iot_platform DB): ```sql CREATE TABLE IF NOT EXISTS tenant_credentials ( namespace TEXT PRIMARY KEY, pg_password TEXT NOT NULL, created_at TIMESTAMPTZ DEFAULT now() ); ``` **Кэширование**: `sync.Map` для `*sql.DB` per tenant. Lazy init при первом обращении. DSN для tenant: `postgresql://tenant_{ns}:{password}@iot-postgres.sless.svc:5432/tenant_{ns}?sslmode=disable` --- ### ШАГ 3: Модифицировать mqtt-bridge — добавить INSERT в Postgres **Файл**: `iot/cmd/mqtt-bridge/main.go` **Текущее поведение**: MQTT message → envelope → RabbitMQ. **Новое поведение**: MQTT message → INSERT в Postgres (tenant DB) + publish в RabbitMQ (как было). **Изменения**: 1. Добавить env var: `IOT_PG_DSN` (admin DSN для iot-postgres) 2. При старте: если IOT_PG_DSN задан → подключиться к IoTPostgresStore 3. В `buildMQTTMessageHandler`: - После парсинга namespace и deviceID из topic - Если store != nil: - `store.EnsureTenantDB(ctx, namespace)` — идемпотентно, кэшируется - `store.InsertTelemetry(ctx, namespace, deviceID, payload)` - При ошибке INSERT → логировать, НЕ останавливать publish в RabbitMQ - RabbitMQ publish остаётся как было **YAML обновление**: `deployments/k8s/iot-mqtt-bridge.yaml` — добавить env: ```yaml - name: IOT_PG_DSN valueFrom: secretKeyRef: name: iot-postgres-secret key: IOT_PG_DSN ``` --- ### ШАГ 4: REST API endpoint для чтения телеметрии **Файл**: `internal/api/handler/iot_telemetry_handler.go` (НОВЫЙ) **Endpoint**: ``` GET /v1/namespaces/{namespace}/iot/telemetry?device={deviceId}&limit={N} ``` **Параметры**: - `device` — фильтр по device_id (опционален: если нет — все устройства namespace) - `limit` — максимум записей (default: 50, max: 1000) - Авторизация: Bearer JWT → namespace validation (как все /v1/ маршруты) **Response**: ```json { "items": [ { "id": 1, "device_id": "sensor-01", "ts": "2026-04-05T08:15:30Z", "payload": {"temperature": 22.5, "humidity": 65} } ], "count": 1 } ``` **Порядок сортировки**: `ts DESC` (новые сверху). **Реализация в handler**: ```go func (h *Handler) ListIoTTelemetry(w http.ResponseWriter, r *http.Request) { ns := mux.Vars(r)["namespace"] deviceID := r.URL.Query().Get("device") limitStr := r.URL.Query().Get("limit") // парсинг limit, default=50, max=1000 rows, err := h.IoTPG.QueryTelemetry(r.Context(), ns, deviceID, limit) // writeJSON } ``` --- ### ШАГ 5: Инициализация IoTPostgresStore в main.go и handler **Файл handler.go** — добавить поле: ```go type Handler struct { K8s client.Client Scheme *runtime.Scheme S3 *s3.Client PG *postgres.Store IoTPG *iotpg.IoTPostgresStore // ← НОВОЕ Log *slog.Logger } ``` **Файл main.go** — после создания Handler: ```go var iotPGStore *iotpg.IoTPostgresStore if iotDSN := os.Getenv("IOT_PG_DSN"); iotDSN != "" { iotPGStore, err = iotpg.New(iotDSN, log) if err != nil { log.Error("failed to connect to IoT Postgres", "err", err) os.Exit(1) } defer iotPGStore.Close() } ``` **Файл router.go** — добавить route: ```go v1.HandleFunc("/namespaces/{namespace}/iot/telemetry", h.ListIoTTelemetry).Methods(http.MethodGet) ``` **YAML**: `deployments/k8s/operator.yaml` — добавить env IOT_PG_DSN из iot-postgres-secret --- ### ШАГ 6: Обновить IoT Console UI — вкладка «Телеметрия» **Файл**: `internal/api/ui/iot-console.html` **Заменить** заглушку `
` на реальную таблицу. **HTML**: ```html

Телеметрия

ВремяУстройствоДанные
Нет данных
``` **CSS**: таблица в стиле Nubes — navy фон, бордеры #0b2d50, текст #e2ecf6. **JavaScript**: - `loadTelemetry()` — fetch GET `/v1/namespaces/{ns}/iot/telemetry?limit=100` → заполнить tbody - Авто-обновление каждые 5с (чекбокс) - Фильтр по устройству (select из списка devices) - При переключении на вкладку → автоматический loadTelemetry() --- ### ШАГ 7: Улучшить эмулятор — рандомные temp/humidity **Файл**: `internal/api/ui/iot-console.html` (секция эмулятора) **Новое поведение**: - Чекбокс: «Генерировать случайные данные (temp/humidity)» — по умолчанию ON - Если ON: при каждой отправке payload = `{temperature: random(18-28), humidity: random(40-80), ts: ISO}` - Если OFF: используется текстовое поле как сейчас ```javascript function generateSensorPayload() { return JSON.stringify({ temperature: +(18 + Math.random() * 10).toFixed(1), humidity: +(40 + Math.random() * 40).toFixed(1), ts: new Date().toISOString() }); } ``` --- ## Деплой и тестирование ### Сборка ```bash cd ~/terra/sless CGO_ENABLED=0 go build -o sless ./main.go CGO_ENABLED=0 go build -o iot-mqtt-bridge ./iot/cmd/mqtt-bridge/ # УВЕЛИЧИТЬ ВЕРСИЮ! docker build --no-cache -t pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.59 . docker push pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.59 ``` ### Обновить версии в YAML ```bash sed -i 's/v0\.1\.58/v0.1.59/g' deployments/k8s/operator.yaml sed -i 's/v0\.1\.53/v0.1.59/g' deployments/k8s/iot-mqtt-bridge.yaml ``` ### Деплой ```bash kubectl apply -f deployments/k8s/iot-postgres.yaml kubectl wait -n sless deploy/iot-postgres --for=condition=available --timeout=60s kubectl apply -f deployments/k8s/operator.yaml kubectl apply -f deployments/k8s/iot-mqtt-bridge.yaml kubectl rollout restart -n sless deploy/sless-operator deploy/iot-mqtt-bridge ``` ### Проверка (E2E) 1. `kubectl exec -n sless deploy/iot-postgres -- psql -U iot_admin -d iot_platform -c "SELECT 1"` 2. `kubectl logs -n sless deploy/sless-operator | grep "IoT Postgres"` 3. Открыть `https://iot.kube5s.ru/console` 4. Ввести токен → вкладка Credentials → должно быть устройство 5. Вкладка Emulator → подключиться, включить «случайные данные», запустить авто-отправку 6. Вкладка Telemetry → должны появляться строки в таблице (avto-refresh 5с) --- ## Файлы СОЗДАТЬ | # | Файл | Описание | |---|------|----------| | 1 | `deployments/k8s/iot-postgres.yaml` | Deployment + Secret + Service | | 2 | `internal/storage/iotpg/iot_telemetry_store.go` | Go: tenant DB management + telemetry CRUD | | 3 | `internal/api/handler/iot_telemetry_handler.go` | REST handler GET telemetry | ## Файлы ИЗМЕНИТЬ | # | Файл | Что менять | |---|------|-----------| | 1 | `internal/api/handler/handler.go` | Добавить поле `IoTPG *iotpg.IoTPostgresStore` | | 2 | `internal/api/router.go` | Route `/namespaces/{ns}/iot/telemetry` | | 3 | `main.go` | Init IoTPostgresStore + передача в Handler | | 4 | `iot/cmd/mqtt-bridge/main.go` | INSERT в Postgres при MQTT message | | 5 | `deployments/k8s/operator.yaml` | env IOT_PG_DSN + версия | | 6 | `deployments/k8s/iot-mqtt-bridge.yaml` | env IOT_PG_DSN + версия | | 7 | `internal/api/ui/iot-console.html` | Telemetry tab + emulator random data | ## Чего НЕ ДЕЛАТЬ - НЕ трогать CRD / controller / EMQX / RabbitMQ — всё работает - НЕ создавать namespace `iot` — пока всё в `sless` - НЕ делать processing данных — пока RAW payload - НЕ добавлять from/to фильтры по времени — хватит limit - НЕ трогать Terraform provider — это следующий этап - НЕ рефакторить существующие файлы - НЕ запускать команды локально — только SSH --- ## Правила для Sonnet (напоминание) 1. Читай `.github/copilot-instructions.md` — правила проекта 2. Читай `doc/decisions/iot-telemetry-storage-2026-04-04.md` — принятое решение 3. Все команды — через SSH: `ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213` 4. Файлы редактировать можно — sshfs mount видна на VM 5. ПЕРЕД go build — проверить .gitignore 6. Комментарии: дата + назначение + «почему» 7. Thinking log: `doc/thinking/2026-04-05.md` 8. progress.md: обновлять до и после каждого шага 9. Коммит + пуш после каждого завершённого шага 10. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой! --- # ПЛАН: Telemetry Pipeline — Postgres -> REST API -> UI > **Автор плана**: GitHub Copilot (Claude Opus 4.6) > **Дата**: 2026-04-05 > **Исполнитель**: Claude Sonnet > **Ветка**: iot-pg-telemetry > **Предусловия**: все компоненты до этого этапа РЕАЛИЗОВАНЫ и задеплоены (CRD, controller, EMQX, mqtt-bridge, IoT Console UI) --- ## Цель Полная цепочка: IoT устройство (или эмулятор в UI) -> MQTT -> INSERT в Postgres -> REST API -> отображение в таблице на вкладке Телеметрия в IoT Console. **User story**: юзер входит токеном, регистрирует устройство, запускает эмулятор (рандомные temp/humidity), переходит на вкладку Телеметрия и видит таблицу с данными: время | устройство | payload. --- ## Что уже готово (НЕ ТРОГАТЬ без крайней необходимости) | Компонент | Файл(ы) | Статус | |-----------|---------|--------| | IoTDevice CRD + types | iot/api/v1alpha1/device_types.go | DONE | | IoTDevice controller | iot/controllers/iotdevice_controller.go | DONE | | MQTT Auth + ACL | internal/api/handler/iot_device_handler.go | DONE | | IoT API CRUD | internal/api/router.go + handler | DONE | | EMQX deploy | deployments/k8s/emqx.yaml | DONE | | mqtt-bridge MQTT->RabbitMQ | iot/cmd/mqtt-bridge/main.go | DONE | | IoT Console UI | internal/api/ui/iot-console.html | DONE | | TLS (HTTPS + WSS) | deployments/k8s/emqx-ws-ingress.yaml | DONE | | Nubes branding | UI CSS | DONE | | Existing Postgres (invocations) | deployments/k8s/postgres.yaml | DONE | --- ## Архитектурное решение (принято 2026-04-04) Подробности: doc/decisions/iot-telemetry-storage-2026-04-04.md - **Один Postgres инстанс** для IoT (отдельный от sless Postgres для invocations) - **Отдельная DATABASE per tenant** (не одна таблица с tenant_id!) - Tenant DB: tenant_{namespace_hash}, User: tenant_{namespace_hash}, Password: UUID в k8s Secret - Таблица: iot_telemetry(id BIGSERIAL, device_id TEXT, ts TIMESTAMPTZ, payload JSONB) - Клиент читает ТОЛЬКО через REST API, не через прямой доступ к Postgres --- ## Шаги реализации (порядок критичен!) ### ШАГ 1: Postgres Deployment для IoT (namespace: sless) **Файл**: deployments/k8s/iot-postgres.yaml **Почему отдельный от sless postgres**: разные данные, разная нагрузка. **Почему в namespace sless, а НЕ iot**: всё живёт в одном namespace, упрощение. **YAML манифест**: **Действие**: kubectl apply -f deployments/k8s/iot-postgres.yaml **Проверка**: kubectl exec -n sless deploy/iot-postgres -- psql -U iot_admin -d iot_platform -c "SELECT 1" --- ### ШАГ 2: Go-пакет IoT Postgres storage **Файл**: internal/storage/iotpg/iot_telemetry_store.go **Структура**: { is a shell keyword **Методы (все обязательные)**: 1. New(adminDSN string, log) (*IoTPostgresStore, error) -- подключение к iot_platform DB 2. EnsureTenantDB(ctx, namespace) error -- создать DATABASE + USER + таблицу если не существуют: - SELECT 1 FROM pg_database WHERE datname = tenant_{ns} - Если нет: CREATE USER, CREATE DATABASE, подключиться и CREATE TABLE - Сохранить пароль в tenant_credentials таблице в iot_platform - Таблица: iot_telemetry(id BIGSERIAL PK, device_id TEXT, ts TIMESTAMPTZ DEFAULT now(), payload JSONB) - Индекс: idx_iot_telemetry_device_ts ON iot_telemetry(device_id, ts DESC) 3. InsertTelemetry(ctx, namespace, deviceID, payload json.RawMessage) error 4. QueryTelemetry(ctx, namespace, deviceID string, limit int) ([]TelemetryRow, error) 5. Close() error **Tenant DB provisioning**: таблица tenant_credentials в iot_platform: **Кэширование**: sync.Map для *sql.DB per tenant. Lazy init при первом обращении. --- ### ШАГ 3: Модифицировать mqtt-bridge -- добавить INSERT в Postgres **Файл**: iot/cmd/mqtt-bridge/main.go **Текущее поведение**: MQTT message -> envelope -> RabbitMQ. **Новое поведение**: MQTT message -> INSERT в Postgres (tenant DB) + RabbitMQ (как было). **Изменения**: 1. Добавить env var IOT_PG_DSN 2. Подключиться к IoTPostgresStore при старте 3. В buildMQTTMessageHandler: - store.EnsureTenantDB(ctx, namespace) -- идемпотентно - store.InsertTelemetry(ctx, namespace, deviceID, payload) - При ошибке INSERT -- логировать, НЕ блокировать RabbitMQ publish 4. RabbitMQ publish остаётся как было **YAML**: deployments/k8s/iot-mqtt-bridge.yaml -- добавить env IOT_PG_DSN из iot-postgres-secret --- ### ШАГ 4: REST API endpoint для чтения телеметрии **Файл**: internal/api/handler/iot_telemetry_handler.go (НОВЫЙ) **Endpoint**: **Параметры**: - device -- фильтр по device_id (опционален) - limit -- максимум записей (default: 50, max: 1000) **Response**: **Сортировка**: ts DESC (новые сверху). --- ### ШАГ 5: Инициализация IoTPostgresStore в main.go **Файл**: main.go 1. Добавить поле IoTPG в handler.Handler struct (handler.go) 2. В main.go: if IOT_PG_DSN задан -> iotpg.New() -> передать в Handler 3. В router.go: зарегистрировать route /namespaces/{ns}/iot/telemetry **YAML**: deployments/k8s/operator.yaml -- добавить env IOT_PG_DSN --- ### ШАГ 6: Обновить IoT Console UI -- вкладка Телеметрия **Файл**: internal/api/ui/iot-console.html **Заменить** заглушку coming-soon на реальную таблицу: | Время | Устройство | Данные | |-------|-----------|--------| | 2026-04-05 08:15 | sensor-01 | {"temperature": 22.5, "humidity": 65} | **JavaScript**: - loadTelemetry() -- fetch GET /v1/.../iot/telemetry -> заполнить tbody - Авто-обновление каждые 5с (чекбокс) - Фильтр по устройству (select из списка devices) - При переключении на вкладку -- автоматический loadTelemetry() **CSS**: таблица в стиле Nubes (navy фон, бордеры #0b2d50, текст #e2ecf6) --- ### ШАГ 7: Улучшить эмулятор -- рандомные temp/humidity **Файл**: internal/api/ui/iot-console.html (секция эмулятора) **Новое поведение**: - Чекбокс: "Генерировать случайные данные (temp/humidity)" (по умолчанию ON) - Если ON: при каждой отправке payload = {temperature: random(18-28), humidity: random(40-80), ts: ISO} - Если OFF: используется текстовое поле как сейчас --- ## Деплой deployment.apps/iot-postgres condition met deployment.apps/sless-operator restarted deployment.apps/iot-mqtt-bridge restarted --- ## Файлы СОЗДАТЬ | Файл | Описание | |------|----------| | deployments/k8s/iot-postgres.yaml | Deployment + Secret + Service | | internal/storage/iotpg/iot_telemetry_store.go | Go: управление tenant DB + CRUD телеметрии | | internal/api/handler/iot_telemetry_handler.go | REST handler GET /v1/.../iot/telemetry | ## Файлы ИЗМЕНИТЬ | Файл | Что менять | |------|-----------| | internal/api/handler/handler.go | Добавить поле IoTPG *iotpg.IoTPostgresStore | | internal/api/router.go | Route /namespaces/{ns}/iot/telemetry | | main.go | Init IoTPostgresStore + передача в Handler | | iot/cmd/mqtt-bridge/main.go | INSERT в Postgres при MQTT message | | deployments/k8s/operator.yaml | env IOT_PG_DSN + версия v0.1.59 | | deployments/k8s/iot-mqtt-bridge.yaml | env IOT_PG_DSN + версия v0.1.59 | | internal/api/ui/iot-console.html | Telemetry tab + emulator random data | ## Чего НЕ ДЕЛАТЬ - НЕ трогать CRD / controller / EMQX / RabbitMQ - НЕ создавать namespace iot -- всё в sless - НЕ делать processing данных -- RAW payload - НЕ добавлять from/to фильтры -- хватит limit - НЕ трогать Terraform provider - НЕ рефакторить существующие файлы - НЕ запускать команды локально -- только SSH --- ## Правила для Sonnet 1. Читай .github/copilot-instructions.md 2. Читай doc/decisions/iot-telemetry-storage-2026-04-04.md 3. Команды через SSH: ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213 4. Файлы редактировать можно -- sshfs mount 5. ПЕРЕД go build -- проверить .gitignore 6. Комментарии: дата + назначение + почему 7. Thinking log: doc/thinking/2026-04-05.md 8. progress.md: обновлять до и после каждого шага 9. Коммит + пуш после каждого завершённого шага 10. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой!