Компоненты: - iot-operator: controller-manager (IoTDevice CRD) + REST API (порт 9090) - mqtt-bridge: MQTT (EMQX) → Kafka bridge - kafka-consumer: Kafka → Postgres pipeline Модуль: gitea.services.ngcloud.ru/Nail/IoT Все 3 бинарника собираются, import paths адаптированы.
1401 lines
53 KiB
Markdown
1401 lines
53 KiB
Markdown
# 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`
|
||
|
||
**Заменить** заглушку `<div class="coming-soon">` на реальную таблицу.
|
||
|
||
**HTML**:
|
||
```html
|
||
<div id="tab-telemetry">
|
||
<div style="display:flex; justify-content:space-between; align-items:center; margin-bottom:16px;">
|
||
<h3>Телеметрия</h3>
|
||
<div>
|
||
<select id="telemetry-device-filter">
|
||
<option value="">Все устройства</option>
|
||
</select>
|
||
<button onclick="loadTelemetry()">Обновить</button>
|
||
<label><input type="checkbox" id="telemetry-auto-refresh"> Авто (5с)</label>
|
||
</div>
|
||
</div>
|
||
<table class="telemetry-table">
|
||
<thead><tr><th>Время</th><th>Устройство</th><th>Данные</th></tr></thead>
|
||
<tbody id="telemetry-body">
|
||
<tr><td colspan="3" style="text-align:center;">Нет данных</td></tr>
|
||
</tbody>
|
||
</table>
|
||
</div>
|
||
```
|
||
|
||
**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. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой!
|