doc: detailed telemetry pipeline plan for Sonnet (Postgres + REST API + UI)

This commit is contained in:
Naeel
2026-04-05 08:27:53 +03:00
parent d078d3156f
commit d51e33d876
2 changed files with 455 additions and 0 deletions
+429
View File
@@ -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`
**Заменить** заглушку `<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. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой!
+26
View File
@@ -0,0 +1,26 @@
# Thinking Log — 2026-04-05
## Агент: GitHub Copilot (Claude Opus 4.6)
---
## Задача: написать подробный план реализации Telemetry Pipeline для Sonnet
### Контекст
Пользователь: "я не понял как этим пользоваться" — вкладка Телеметрия заглушка.
Нужна полная цепочка: устройство -> данные в таблице.
### Анализ — что уже есть
IoT-компоненты реализованы на 85%: CRD, controller, EMQX, MQTT auth/ACL, mqtt-bridge, IoT Console.
mqtt-bridge: MQTT -> RabbitMQ (для function triggers). Данные НЕ сохраняются.
Postgres для invocations есть, для IoT телеметрии — нет.
### Ключевые решения при написании плана
1. Postgres сразу (пользователь настоял)
2. Отдельный Postgres инстанс (iot-postgres), но в namespace sless
3. DATABASE per tenant (изоляция)
4. mqtt-bridge: INSERT параллельно с RabbitMQ publish
5. REST API: GET /v1/.../iot/telemetry (только limit, без from/to)
6. Эмулятор: рандом temp 18-28, humidity 40-80
### Результат
Записан подробный план в doc/iot-mvp-plan.md — 7 шагов с файлами, кодом и YAML.