diff --git a/doc/iot-mvp-plan.md b/doc/iot-mvp-plan.md index 5c04bea..c912aa5 100644 --- a/doc/iot-mvp-plan.md +++ b/doc/iot-mvp-plan.md @@ -737,3 +737,432 @@ IoT-хендлеры добавлять как методы того же Handle - 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` + +**Заменить** заглушку `
| Время | Устройство | Данные |
|---|---|---|
| Нет данных | ||