4 Commits
Author SHA1 Message Date
Naeel d57558c798 docs: IoT telemetry storage architecture decisions and thinking log 2026-04-04 18:41:29 +03:00
Naeel b23ae40975 security(iot): MQTT ACL isolation via EMQX HTTP authorization
Each IoT device can only pub/sub to its own topics: {namespace}/{deviceId}/#
Any attempt to access foreign topics → EMQX denies and disconnects.

Changes:
- internal/api/handler: add MQTTAcl handler (POST /internal/mqtt/acl)
- internal/api/router: register /internal/mqtt/acl route
- deployments/k8s/emqx.yaml: add HTTP authorization backend, no_match=deny
- Operator v0.1.52 deployed

Tested: own topic ALLOWED, foreign topic → authorization_permission_denied + disconnect
2026-04-04 17:51:51 +03:00
Naeel 6e3e473551 feat(iot): MQTT over WebSocket workaround via nginx-ingress
Port 1883 blocked by NSX-T Edge firewall (DevOps to open on Monday).
Temporary solution: EMQX WebSocket listener (8083) via nginx-ingress.

- Service emqx-ws: selector app=emqx, port 8083
- Ingress emqx-mqtt-websocket: iot.kube5s.ru/mqtt → emqx-ws:8083
- pathType: Exact (prevents trailing slash redirect breaking WS handshake)
- ssl-redirect: false (IoT devices cannot follow HTTP 301 redirects)
- DNS A record iot.kube5s.ru → 185.247.187.147 created by user

Tested: Connected rc=0 via paho-mqtt WebSocket from inside cluster.
2026-04-04 14:06:17 +03:00
Naeel 857d057af9 feat(iot): деплой IoT MVP — Dockerfile, RBAC, EMQX fix, operator v0.1.50, mqtt-bridge, doc/iot 2026-04-04 10:29:47 +03:00
15 changed files with 1116 additions and 45 deletions
+10 -1
View File
@@ -1,4 +1,4 @@
# Изменено: 2026-03-07
# Изменено: 2026-04-04 — добавлена IoT поддержка: COPY iot/ + сборка iot-mqtt-bridge бинаря
# Multi-stage build для sless оператора.
# Stage 1: сборка бинаря (golang:1.23-alpine)
# Stage 2: минимальный образ (alpine:3.19, не distroless — нужен ca-certificates для S3/HTTPS)
@@ -17,14 +17,23 @@ COPY api/ api/
COPY controllers/ controllers/
COPY internal/ internal/
COPY migrations/ migrations/
# iot/ — IoT CRD types, controller, mqtt-bridge cmd.
# Обязательно: main.go импортирует iot/api/v1alpha1 и iot/controllers — без этого go build упадёт.
COPY iot/ iot/
RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} go build -a -o manager main.go
# iot-mqtt-bridge — отдельный бинарь в том же образе.
# Запускается в iot-mqtt-bridge Deployment через command: ["/iot-mqtt-bridge"].
# Один образ, два entrypoint — практично для MVP: один CI pipeline, один registry repo.
RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} go build -a -o iot-mqtt-bridge ./iot/cmd/mqtt-bridge/
FROM alpine:3.19
# ca-certificates нужны для TLS (S3 HTTPS, DockerHub)
RUN apk add --no-cache ca-certificates
WORKDIR /
COPY --from=builder /workspace/manager .
# iot-mqtt-bridge — второй бинарь, запускается отдельным Deployment-ом.
COPY --from=builder /workspace/iot-mqtt-bridge .
# migrations нужны при старте — оператор читает SQL файлы для инициализации БД
COPY migrations/ migrations/
# Запускаем от непривилегированного пользователя
+51
View File
@@ -0,0 +1,51 @@
# Создано: 2026-04-04
# ВРЕМЕННЫЙ WORKAROUND: MQTT over WebSocket через порт 443.
# Причина: порт 1883 заблокирован NSX-T Edge firewall на уровне облака.
# Решение: EMQX WebSocket listener (8083) проксируется через nginx-ingress.
#
# IoT устройство подключается: ws://iot.kube5s.ru/mqtt
# DNS A-запись: iot.kube5s.ru → 185.247.187.147 (создана через Nubes API, zoneUid=498096ee)
#
# Когда DevOps откроет порт 1883 — этот файл можно удалить.
---
apiVersion: v1
kind: Service
metadata:
name: emqx-ws
namespace: sless
# Отдельный Service чтобы не путать — WebSocket порт для Ingress
spec:
selector:
app: emqx
ports:
- name: mqtt-ws
port: 8083
targetPort: 8083
protocol: TCP
---
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: emqx-mqtt-websocket
namespace: sless
annotations:
kubernetes.io/ingress.class: nginx
# ssl-redirect=false: IoT устройства не умеют следовать 301 редиректам при WebSocket upgrade
nginx.ingress.kubernetes.io/ssl-redirect: "false"
# WebSocket: nginx-ingress автоматически добавляет Upgrade/Connection при proxy-http-version=1.1
nginx.ingress.kubernetes.io/proxy-http-version: "1.1"
nginx.ingress.kubernetes.io/proxy-read-timeout: "3600"
nginx.ingress.kubernetes.io/proxy-send-timeout: "3600"
spec:
ingressClassName: nginx
rules:
- host: iot.kube5s.ru
http:
paths:
- path: /mqtt
pathType: Exact
backend:
service:
name: emqx-ws
port:
number: 8083
+33 -3
View File
@@ -31,6 +31,15 @@ data:
## EMQX 5.x configuration (HOCON format)
## Изменено: 2026-04-04
## Обязательные поля node — без них EMQX 5.x падает при старте
## node.cookie — секрет кластерного Erlang-соединения, для single-node любая строка
## node.data_dir — директория данных (mnesia, конфиги), должна существовать в контейнере
node {
name = "emqx@127.0.0.1"
cookie = "sless-emqx-cookie-mvp"
data_dir = "/opt/emqx/data"
}
## HTTP Auth Backend для IoT-устройств
## EMQX посылает POST с {username, password, clientid} → наш сервис отвечает {"result":"allow"|"deny"}
authentication = [
@@ -55,16 +64,37 @@ data:
}
]
## ACL по умолчанию — разрешаем всё аутентифицированным клиентам
## Тонкая ACL настраивается через HTTP auth response (поле acl)
## Authorization (ACL) — HTTP backend для изоляции топиков по устройству.
## no_match = deny: если HTTP backend недоступен или не ответил — запрещаем.
## Endpoint /internal/mqtt/acl возвращает allow только для топиков {ns}/{deviceId}/#
authorization {
no_match = allow
no_match = deny
deny_action = disconnect
cache {
enable = true
max_size = 32
ttl = 1m
}
sources = [
{
type = http
enable = true
method = post
url = "http://sless-operator.sless.svc:9090/internal/mqtt/acl"
body {
username = "${username}"
clientid = "${clientid}"
action = "${action}"
topic = "${topic}"
}
headers {
"content-type" = "application/json"
}
connect_timeout = 5s
request_timeout = 5s
pool_size = 8
}
]
}
## MQTT настройки
+3 -2
View File
@@ -42,8 +42,9 @@ spec:
spec:
containers:
- name: mqtt-bridge
image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:latest
# TODO: отдельный образ iot-mqtt-bridge После сборки через Makefile
# Тот же образ что и оператор — оба бинаря в одном слое (manager + iot-mqtt-bridge).
# При смене версии оператора — менять тег и здесь.
image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.50
imagePullPolicy: Always
command: ["/iot-mqtt-bridge"]
env:
+2 -1
View File
@@ -74,7 +74,8 @@ spec:
containers:
- name: operator
# При обновлении версии оператора — менять тег здесь (не latest!)
image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.49
# v0.1.50 — добавлены IoT controller, IoT REST API, MQTT auth endpoint
image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.50
# Always — чтобы всегда тянуть по точному тегу (не кешировать старый)
imagePullPolicy: Always
ports:
+13 -2
View File
@@ -1,4 +1,4 @@
# Изменено: 2026-03-20 (добавлен Service CRD sless.kube5s.ru — services + status + finalizers)
# Изменено: 2026-04-04 — добавлены IoT CRD права (iot.kube5s.ru)
# RBAC для sless оператора.
# ServiceAccount + ClusterRole + ClusterRoleBinding.
# ClusterRole нужен (не namespaced Role) потому что оператор создаёт
@@ -15,7 +15,7 @@ kind: ClusterRole
metadata:
name: sless-operator
rules:
# Наши CRD
# Наши CRD (sless.kube5s.ru)
- apiGroups: ["sless.kube5s.ru"]
resources: ["functions", "triggers", "functionjobs", "services"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
@@ -25,6 +25,17 @@ rules:
- apiGroups: ["sless.kube5s.ru"]
resources: ["functions/finalizers", "triggers/finalizers", "functionjobs/finalizers", "services/finalizers"]
verbs: ["update"]
# IoT CRD (iot.kube5s.ru) — IoTDevice lifecycle + Secret генерация в контроллере
# Права нужны во всех namespace где пользователи создают IoT-устройства
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices/status"]
verbs: ["get", "update", "patch"]
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices/finalizers"]
verbs: ["update"]
# Deployments для функций
- apiGroups: ["apps"]
resources: ["deployments"]
+54 -12
View File
@@ -1,32 +1,74 @@
# Архитектура системы
Последнее обновление: 2026-03-18 (v0.1.34 + funcs-service v0.2.0)
Последнее обновление: 2026-04-04 (IoT telemetry storage architecture decision)
## Общее описание
Managed Serverless Functions Service для облачного провайдера nubes.ru.
Managed Serverless Functions Service + IoT Platform для облачного провайдера nubes.ru.
Два независимых компонента: sless (serverless) и iot (IoT), каждый со своим оператором.
Пользователь загружает код через Terraform, сервис его собирает (kaniko) и запускает
по HTTP-триггеру, расписанию (cron) или вручную через one-shot Job.
по HTTP-триггеру, расписанию (cron), вручную или по событию от IoT устройства.
## Namespace Layout
```
namespace: sless — платформа serverless
(sless-operator, event-dispatcher, RabbitMQ, Postgres invocations)
namespace: sless-{hash} — tenant функции (function pods каждого клиента)
namespace: iot — платформа IoT
(iot-operator, EMQX, Postgres telemetry)
namespace: iot-{hash} — tenant IoT ресурсы (IoTDevice CRDs)
```
## Стек
| Компонент | Технология | Где запущен |
|-----------|-----------|-------------|
| Operator (API + Controllers) | Go (controller-runtime) | Kubernetes, namespace `sless` |
| funcs-service (глобальная консоль) | Go (net/http) | Kubernetes, namespace `sless` |
| PostgreSQL | PostgreSQL 16 | Kubernetes, namespace `sless` |
| sless-operator (API + Controllers) | Go (controller-runtime) | namespace `sless` |
| iot-operator (API + Controllers) | Go (controller-runtime) | namespace `iot` |
| PostgreSQL (invocations) | PostgreSQL 16 | namespace `sless` |
| PostgreSQL (telemetry) | PostgreSQL 16 | namespace `iot` |
| EMQX | EMQX 5.5.1 | namespace `iot` |
| RabbitMQ | RabbitMQ 3 | namespace `sless` |
| event-dispatcher | Go | namespace `sless` |
| iot-mqtt-bridge | Go | namespace `iot` |
| S3 | Ceph (облачный) | `s3.msk-1.ngcloud.ru` |
| Container Registry | DockerHub (`naeel/`) | внешний |
| Container Registry | PearlHarbor (Nubes) | внешний |
| Builder | kaniko (k8s Job) | namespace пользователя |
| Функции (HTTP) | k8s Deployment + Service | namespace пользователя |
| Функции (one-shot) | k8s Job | namespace пользователя |
| Функции (cron) | k8s CronJob | namespace пользователя |
| Функции (HTTP) | k8s Deployment + Service | namespace sless-{hash} |
| Функции (one-shot) | k8s Job | namespace sless-{hash} |
| Функции (cron) | k8s CronJob | namespace sless-{hash} |
| Terraform Provider | Go (plugin framework v6) | localhost/CI |
| nubes API | REST (облако) | `deck-api.ngcloud.ru` |
> Redis и RabbitMQ — отложены до v2.
## IoT Data Flow
## Компонент: funcs-service
```
IoT устройство
↓ ws://iot.kube5s.ru:80/mqtt (WebSocket, пока 1883 закрыт)
EMQX (namespace iot)
↓ ACL: каждое устройство видит только свои топики {ns}/{deviceId}/#
iot-mqtt-bridge
RabbitMQ (namespace sless)
event-dispatcher
↓ параллельно:
1. INSERT INTO tenant_{ns}.iot_telemetry ← автоматически
2. Вызов serverless function (если настроена)
```
## Изоляция данных
- MQTT: ACL по username → топики только своего устройства
- Postgres: отдельная DATABASE per tenant, разные credentials
- k8s: отдельный namespace per tenant
## Связь sless ↔ iot
- Общий идентификатор tenant: `{hash}` в именах namespace
- Коммуникация через RabbitMQ endpoint (не через Go пакеты)
- Loose coupling — могут быть в разных кластерах
Глобальный HTTP сервис — **одна копия** на весь кластер, для всех пользователей.
@@ -0,0 +1,174 @@
# Решение: IoT Telemetry Storage Architecture
# Дата: 2026-04-04
# Агент: GitHub Copilot (Claude Sonnet 4.6)
# Статус: ПРИНЯТО
---
## Контекст
IoT платформа принимает данные с датчиков через MQTT. Данные проходят:
EMQX → iot-mqtt-bridge → RabbitMQ → event-dispatcher → function pod.
Проблема: данные не сохраняются. Функция получает событие и забывает его.
Для клиентов (мониторинг объектов, счётчики, производство) нужно:
- Автоматическое хранение всей телеметрии
- Доступ к историческим данным
- Низкий порог входа — не требовать от клиента настройки БД
---
## Решения
### 1. Хранилище — Postgres, отдельная DATABASE per tenant
**Выбрано**: один Postgres инстанс, отдельная DATABASE на каждого клиента.
**Отклонено**: одна таблица с tenant_id колонкой.
- Причина: изоляция только программная. Ошибка в WHERE → утечка чужих данных.
**Структура**:
```
Postgres (StatefulSet в namespace iot)
├── sless_platform — системные данные платформы (tenants, etc)
├── tenant_{hash} — данные клиента A (полная изоляция)
└── tenant_{hash} — данные клиента B (полная изоляция)
```
**Безопасность**:
- Каждый tenant имеет свой Postgres USER с уникальным паролем (UUID)
- Пароль генерируется при создании tenant, хранится в k8s Secret
- Клиент B физически не может подключиться к DATABASE клиента A
---
### 2. Доступ клиента — только через REST API
**Выбрано**: клиент читает телеметрию через REST API платформы.
**Отклонено**: прямой доступ к Postgres через connection string.
- Причина: Postgres внутри кластера, не должен торчать наружу. Security.
**API**:
```
GET /v1/namespaces/{ns}/iot/telemetry
?device={device_id}
&from={RFC3339}
&to={RFC3339}
&limit={int}
GET /v1/namespaces/{ns}/iot/devices/{id}/last
```
Авторизация — Bearer токен (тот же механизм что и для functions).
---
### 3. Схема таблицы telemetry
```sql
CREATE TABLE iot_telemetry (
id BIGSERIAL PRIMARY KEY,
device_id TEXT NOT NULL,
ts TIMESTAMPTZ NOT NULL DEFAULT now(),
payload JSONB NOT NULL
);
CREATE INDEX idx_iot_telemetry_device_ts
ON iot_telemetry (device_id, ts DESC);
```
**Почему JSONB**: у каждого клиента разные наборы данных:
- датчик температуры: `{"temp": 22.5, "humidity": 60}`
- GPS трекер: `{"lat": 55.75, "lon": 37.61, "speed": 60}`
- счётчик воды: `{"liters": 1234.5, "flow": 0.3}`
Фиксированная схема невозможна. JSONB + индекс по (device_id, ts) даёт
достаточную производительность для малого и среднего бизнеса.
---
### 4. schema.sql при деплое функции
Клиент может положить `schema.sql` рядом с функцией:
```
my-function/
├── handler.py
├── schema.sql ← CREATE TABLE IF NOT EXISTS my_alerts (...)
└── requirements.txt
```
При деплое оператор выполняет `schema.sql` в БД tenant'а.
Это позволяет клиентам без знания Python настраивать дополнительные таблицы.
---
### 5. DB_DSN в функцию
При запуске function pod оператор прокидывает `DB_DSN` из Secret в env var:
```
DB_DSN=postgresql://tenant_abc:password@iot-postgres.iot.svc:5432/tenant_abc
```
Функция использует стандартный драйвер, не знает о деталях платформы.
---
### 6. Разделение sless и iot операторов
**Решение**: sless-operator и iot-operator — ОТДЕЛЬНЫЕ компоненты с раздельными namespace.
**Мотивация**:
- В будущем могут быть в разных кластерах
- Независимый деплой и масштабирование
- Разные команды могут владеть компонентами
- Нет cross-dependency в коде (loose coupling)
**Namespace layout**:
```
namespace: sless — платформа serverless
(sless-operator, event-dispatcher, RabbitMQ, Postgres invocations)
namespace: sless-{hash} — tenant функции (function pods каждого клиента)
namespace: iot — платформа IoT
(iot-operator, EMQX, Postgres telemetry)
namespace: iot-{hash} — tenant IoT ресурсы (IoTDevice CRDs)
```
**Связь между sless и iot**:
- Общий идентификатор tenant: `{hash}` одинаковый в sless-{hash} и iot-{hash}
- MQTT событие → RabbitMQ в namespace sless → function pod в sless-{hash}
- iot-operator НЕ импортирует Go пакеты sless-operator
- Общение только через k8s API и RabbitMQ endpoints
**Postgres**:
- sless: отдельный Postgres для invocations логов
- iot: отдельный Postgres для telemetry per-tenant
- Разные StatefulSet, разные PVC, разные credentials
---
### 7. Postgres инстанс для IoT
**Выбрано**: `postgres:16-alpine` StatefulSet в namespace `iot`.
**Причина**: простота для разработки. При передаче в production девопсы
заменят на Managed Postgres от Nubes — connection string поменяется, код не меняется.
**Ресурсы**:
- PVC: 10Gi (начальный размер, увеличивается по мере роста)
- Memory limit: 512Mi
- CPU: 0.5 cores
---
## План реализации
1. StatefulSet Postgres в namespace `iot`
2. iot-operator: provisioning при создании IoTDevice namespace
- CREATE USER tenant_{ns} PASSWORD '{uuid}'
- CREATE DATABASE tenant_{ns} OWNER tenant_{ns}
- CREATE TABLE iot_telemetry + индекс
3. iot-mqtt-bridge: INSERT telemetry при получении MQTT сообщения
4. iot-operator: REST API `/v1/namespaces/{ns}/iot/telemetry`
5. sless-operator: при деплое function → прокинуть DB_DSN + выполнить schema.sql
+306
View File
@@ -0,0 +1,306 @@
# IoT MVP — Инженерная документация деплоя
> Создано: 2026-04-04
> Ветка: Ioter
> Автор: GitHub Copilot (Claude Sonnet 4.6)
---
## Архитектура IoT стека
```
IoT Device (физическое)
│ MQTT CONNECT (username="{ns}_{deviceId}", password=hex)
EMQX 5.5.1 (sless/emqx)
│ HTTP POST /internal/mqtt/auth → sless-operator:9090
│ (auth backend: проверяет Secret iot-{deviceId} в k8s)
│ MQTT PUBLISH → topic: "{ns}/telemetry/{deviceId}"
iot-mqtt-bridge (sless/iot-mqtt-bridge)
│ paho.mqtt.golang, подписка на "+/telemetry/+"
│ parse topic → namespace из первого сегмента
RabbitMQ (sless/rabbitmq)
│ queue: "iot.{namespace}.telemetry"
event-dispatcher (sless/event-dispatcher)
│ Trigger type=event, queue=iot.{namespace}.telemetry
Serverless Function (пользовательский handler)
```
---
## Компоненты
### 1. CRD IoTDevice
**Расположение:** `iot/api/v1alpha1/device_types.go`
**API group:** `iot.kube5s.ru/v1alpha1`
**Манифест:** `iot/config/crd/bases/iot.kube5s.ru_iotdevices.yaml`
Поля Spec:
| Поле | Тип | Обязательное | Описание |
|------|-----|--------------|----------|
| `deviceId` | string | да | Идентификатор устройства. Pattern: `^[a-z0-9][a-z0-9-]*[a-z0-9]$` |
| `enabled` | bool | нет | Активно ли устройство (default: true) |
| `metadata` | map[string]string | нет | Произвольные метаданные (модель, локация) |
Поля Status:
| Поле | Описание |
|------|----------|
| `phase` | `Active` / `Disabled` / `Pending` / `Error` |
| `mqttUsername` | `{namespace}_{deviceId}` |
| `secretName` | Имя k8s Secret с credentials |
| `topicPrefix` | `{namespace}/` |
| `message` | Сообщение об ошибке если phase=Error |
### 2. IoT Controller
**Файл:** `iot/controllers/iotdevice_controller.go`
**Логика Reconcile:**
```
IoTDevice CREATE/UPDATE
1. Добавить finalizer "iot.kube5s.ru/device-cleanup"
2. Если Secret iot-{deviceId} не существует:
- Сгенерировать пароль: crypto/rand 32 bytes → hex (64 символа)
- OwnerReference → Secret удаляется каскадно при удалении IoTDevice
- Secret keys: mqtt-username, mqtt-password, device-id
3. Обновить Status: phase=Active, mqttUsername, secretName, topicPrefix
4. Если enabled=false → phase=Disabled
IoTDevice DELETE
1. Проверить finalizer
2. Secret удаляется каскадно (OwnerReference)
3. Убрать finalizer → k8s завершает удаление
```
### 3. IoT REST API
**Файл:** `internal/api/handler/iot_device_handler.go`
| Endpoint | Auth | Описание |
|----------|------|----------|
| `POST /internal/mqtt/auth` | Нет (internal) | MQTT auth backend для EMQX |
| `POST /v1/namespaces/{ns}/iot/devices` | JWT | Создать IoTDevice |
| `GET /v1/namespaces/{ns}/iot/devices` | JWT | Список (без паролей) |
| `GET /v1/namespaces/{ns}/iot/devices/{name}` | JWT | Получить (включая mqtt_password из Secret) |
| `DELETE /v1/namespaces/{ns}/iot/devices/{name}` | JWT | Удалить |
| `PATCH /v1/namespaces/{ns}/iot/devices/{name}` | JWT | Обновить enabled |
**MQTT Auth endpoint:**
- Всегда HTTP 200 (EMQX игнорирует non-200)
- Парсит `username``{namespace}_{deviceId}` (разделитель первый `_`)
- Ищет k8s Secret `iot-{deviceId}` в namespace
- `crypto/subtle.ConstantTimeCompare` для защиты от timing attack
### 4. EMQX 5.5.1
**Манифест:** `deployments/k8s/emqx.yaml`
**Конфиг:** HOCON `emqx.conf`, монтируется как ConfigMap volume
**Критически важные поля (без них EMQX 5.x не стартует):**
```hocon
node {
name = "emqx@127.0.0.1" # Обязательно для single-node
cookie = "..." # Erlang cluster cookie (любая строка для single-node)
data_dir = "/opt/emqx/data" # Директория данных Mnesia
}
```
> ⚠️ EMQX 5.x: поля `node.cookie` и `node.data_dir` — **обязательные** (mandatory),
> в отличие от 4.x где были значения по умолчанию.
> При обновлении ConfigMap нужен `kubectl rollout restart` — Deployment не перезапускается автоматически.
**Auth backend:**
```hocon
authentication = [{
mechanism = password_based
backend = http
method = post
url = "http://sless-operator.sless.svc:9090/internal/mqtt/auth"
}]
```
### 5. iot-mqtt-bridge
**Код:** `iot/cmd/mqtt-bridge/main.go`
**Манифест:** `deployments/k8s/iot-mqtt-bridge.yaml`
**Образ:** тот же что и оператор (`sless-operator:v0.1.50`), бинарь `/iot-mqtt-bridge`
**Логика:**
1. Подключиться к EMQX как MQTT клиент (credentials из Secret `iot-bridge-credentials`)
2. Подписаться на `+/telemetry/+` (все namespace, все устройства)
3. При получении: извлечь namespace из topic[0], publish в RabbitMQ `iot.{namespace}.telemetry`
4. Reconnect loop при обрыве соединения
**Envelope в RabbitMQ:**
```json
{
"namespace": "sless-user123",
"device_id": "sensor-01",
"topic": "sless-user123/telemetry/sensor-01",
"payload": "<base64 of raw MQTT payload>",
"received_at": "2026-04-04T07:19:30Z"
}
```
### 6. Terraform Provider
**Файл:** `terraform/provider/internal/resources/iot_device_resource.go`
**Ресурс:** `sless_iot_device`
**Версия провайдера:** `0.1.2`
```hcl
resource "sless_iot_device" "temperature_sensor" {
name = "temp-sensor-01"
device_id = "temp-sensor-01"
enabled = true
metadata = {
model = "DHT22"
location = "Warehouse A"
}
}
output "mqtt_password" {
value = sless_iot_device.temperature_sensor.mqtt_password
sensitive = true
}
```
---
## Процедура первого деплоя
### Предварительные условия
- Кластер с namespace `sless`
- sless-operator запущен (или будет запущен в шаге 3)
- RabbitMQ доступен в кластере
### Шаги
**1. Применить CRD (один раз, cluster-wide)**
```bash
kubectl apply -f iot/config/crd/bases/iot.kube5s.ru_iotdevices.yaml
```
**2. Обновить RBAC (добавить права на iot.kube5s.ru)**
```bash
kubectl apply -f deployments/k8s/rbac.yaml
```
**3. Применить EMQX**
```bash
kubectl apply -f deployments/k8s/emqx.yaml
kubectl rollout status deployment/emqx -n sless
```
**4. Применить оператор (с IoT поддержкой)**
```bash
kubectl apply -f deployments/k8s/operator.yaml
kubectl rollout status deployment/sless-operator -n sless
```
**5. Bootstrap credentials для mqtt-bridge**
Создать системное IoTDevice устройство для bridge:
```bash
TOKEN=$(kubectl get secret sless-operator-secret -n sless \
-o jsonpath="{.data.SLESS_API_TOKEN}" | base64 -d)
# Создать IoTDevice
curl -X POST https://sless.kube5s.ru/v1/namespaces/sless/iot/devices \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d '{"name":"iot-bridge","device_id":"iot-bridge","enabled":true}'
# Подождать 5с пока контроллер создаст Secret
sleep 5
# Получить credentials
CREDS=$(curl -s https://sless.kube5s.ru/v1/namespaces/sless/iot/devices/iot-bridge \
-H "Authorization: Bearer $TOKEN")
MQTT_USER=$(echo $CREDS | jq -r .mqtt_username)
MQTT_PASS=$(echo $CREDS | jq -r .mqtt_password)
# Создать Secret для bridge Deployment
kubectl create secret generic iot-bridge-credentials -n sless \
--from-literal=MQTT_USERNAME="$MQTT_USER" \
--from-literal=MQTT_PASSWORD="$MQTT_PASS"
```
**6. Применить mqtt-bridge**
```bash
kubectl apply -f deployments/k8s/iot-mqtt-bridge.yaml
kubectl rollout status deployment/iot-mqtt-bridge -n sless
```
### Ожидаемый результат
```
emqx-xxx 1/1 Running
iot-mqtt-bridge-xxx 1/1 Running
sless-operator-xxx 1/1 Running
```
---
## Известные ошибки и решения
### EMQX CrashLoopBackOff: required_field node.cookie/node.data_dir
**Симптом:** `escript: exception throw: {emqx_conf_schema, [{kind=>validation_error, path=>"node.cookie", reason=>required_field}]}`
**Причина:** EMQX 5.x требует явного задания `node { cookie, data_dir }` в конфиге.
**Решение:** Добавить в `emqx.conf`:
```hocon
node {
name = "emqx@127.0.0.1"
cookie = "your-cookie-string"
data_dir = "/opt/emqx/data"
}
```
После `kubectl apply` — сделать `kubectl rollout restart deployment/emqx -n sless`.
---
### RBAC forbidden: iotdevices.iot.kube5s.ru
**Симптом:** `{"error":"iotdevices.iot.kube5s.ru is forbidden: User \"system:serviceaccount:sless:sless-operator\" cannot create resource"}`
**Причина:** ClusterRole `sless-operator` не включает API group `iot.kube5s.ru`.
**Решение:** Добавить в `deployments/k8s/rbac.yaml` и применить:
```yaml
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices/status"]
verbs: ["get", "update", "patch"]
- apiGroups: ["iot.kube5s.ru"]
resources: ["iotdevices/finalizers"]
verbs: ["update"]
```
---
### mqtt-bridge: multiple restarts при старте
**Симптом:** `iot-mqtt-bridge RESTARTS=3`
**Причина:** bridge пытается подключиться к EMQX который ещё не готов. Нормальное поведение.
**Решение:** Bridge имеет reconnect loop — после старта EMQX подключение восстанавливается автоматически. Ничего делать не нужно.
---
## Версии образов
| Версия | Дата | Изменения |
|--------|------|-----------|
| v0.1.50 | 2026-04-04 | IoT controller + IoT API + iot-mqtt-bridge бинарь |
| v0.1.49 | ранее | До IoT |
+47 -1
View File
@@ -1597,5 +1597,51 @@ G15 перезапущен → **21/21 PASS ✅**
| 4 | nginx `client_max_body_size` ограничивает upload → 413 (не настроено явно) | G13F-4 NOTE |
### Версия оператора
`v0.1.51` — задеплоен, работает
`v0.1.52` — задеплоен, работает
---
## 2026-04-04 — IoT WebSocket workaround + MQTT ACL + Архитектура телеметрии
### Выполнено
#### MQTT WebSocket workaround (порт 1883 закрыт NSX-T)
- Создан Ingress `emqx-mqtt-websocket`: `iot.kube5s.ru/mqtt → EMQX:8083`
- DNS A-запись `iot.kube5s.ru → 185.247.187.147` создана пользователем
- Отлажена цепочка: убран `configuration-snippet` (заблокирован в nginx v1.12.6),
добавлен `ssl-redirect: false`, `pathType: Exact`
- Тест: `Connected rc=0` через paho-mqtt WebSocket ✅
- Коммит: `6e3e473` (ветка Ioter)
#### MQTT ACL изоляция топиков
- Найдена уязвимость: `authorization { no_match = allow }` — любой клиент мог
читать топики других клиентов после успешного CONNECT
- EMQX 5.x: ACL в ответе auth игнорируется (это EMQX 4.x фича)
- Добавлен endpoint `POST /internal/mqtt/acl` в sless-operator
- Обновлён `emqx.conf`: HTTP authorization backend, `no_match = deny`
- Тест: `sless/iot-bridge/#` → ALLOWED, `sless/other-device/telemetry` → DENIED + disconnect
- Лог EMQX: `authorization_permission_denied`
- Оператор v0.1.52 задеплоен
- Коммит: `b23ae40` (ветка Ioter)
### Архитектурные решения (обсуждение, не реализовано)
Принято решение о хранении IoT телеметрии:
- Отдельная DATABASE per tenant в одном Postgres инстансе
- REST API для доступа (не прямой доступ к Postgres)
- JSONB payload (разные данные у разных клиентов)
- Отдельный iot-operator независимо от sless-operator
- schema.sql при деплое функции
- DB_DSN в env var функции
Подробно: `doc/decisions/iot-telemetry-storage-2026-04-04.md`
### Следующий этап (ветка iot-pg-telemetry)
- [ ] Postgres StatefulSet в namespace `iot`
- [ ] Provisioning БД при создании tenant
- [ ] INSERT telemetry из iot-mqtt-bridge
- [ ] REST API чтения телеметрии
- [ ] DB_DSN в function pod env
- [ ] schema.sql при деплое функции
+215
View File
@@ -0,0 +1,215 @@
# Лог мышления — 2026-04-04-02
## Агент: GitHub Copilot (Claude Sonnet 4.6)
Продолжение сессии от 2026-04-04 (файл 2026-04-04.md достиг 681 строки).
---
## Задача: MQTT WebSocket через nginx-ingress — отладка 404 → 400
### Контекст
Порт 1883 заблокирован NSX-T Edge firewall на уровне облака.
Временный workaround: EMQX WebSocket listener (8083) → nginx-ingress → `iot.kube5s.ru/mqtt`.
### Что нашёл
**Проблема 1 — 404:** Ingress не получил ADDRESS (IP не присвоен).
- Причина: annotation `nginx.ingress.kubernetes.io/configuration-snippet` заблокирован в nginx-controller v1.12.6.
По умолчанию `allow-snippet-annotations=false`.
- Решение: убрал `configuration-snippet` из Ingress.
- Результат: Ingress получил sync (Scheduled for sync), но ADDRESS всё ещё пустой — это нормально для данного контроллера.
**Проблема 2 — 400:** После удаления snippet стало 400 вместо 404.
- nginx конфиг для `iot.kube5s.ru` правильный:
- `proxy_set_header Upgrade $http_upgrade` — ЕСТЬ (nginx-ingress добавляет автоматически)
- `proxy_http_version 1.1` — ЕСТЬ
- `proxy_send_timeout/read_timeout 3600s` — ЕСТЬ
- `proxy_pass → sless-emqx-ws-8083` — ЕСТЬ
- Значит 400 приходит от EMQX, а не от nginx.
### Гипотезы почему EMQX возвращает 400
EMQX WebSocket listener конфиг:
```
listeners.ws.default {
bind = "0.0.0.0:8083"
max_connections = 512
}
```
Нет явного `websocket.path`. По умолчанию EMQX ожидает путь `/mqtt`.
Но nginx из-за `pathType: Prefix` на `/mqtt` может передавать путь без trailing slash.
**Вариант A**: nginx передаёт запрос на `/mqtt` (без `/`) — EMQX ожидает `/mqtt` → должно работать.
**Вариант B**: EMQX в 5.5.1 ожидает `ws://host/mqtt` без всякого proxy path rewrite и есть баг с `X-Forwarded-Proto` — nginx шлёт `http`, а EMQX может редиректить.
**Вариант C**: `ssl_redirect: true` в nginx конфиге означает что HTTP запрос перенаправляется на HTTPS, и там 400.
### Следующий шаг
Проверить `$ssl_redirect = true` в nginx server block для нашего хоста:
- Если да — nginx делает redirect HTTP→HTTPS
- Добавить аннотацию `nginx.ingress.kubernetes.io/ssl-redirect: "false"` или тестировать через HTTPS
### Правило (новое)
Если файл `doc/thinking/ГГГГ-ММ-ДД.md` стал большим — создавать `ГГГГ-ММ-ДД-02.md`, `ГГГГ-ММ-ДД-03.md` и т.д.
---
## Задача: Анализ безопасности MQTT multitenancy
### Что исследовал
Пользователь спросил об угрозах межтенантного проникновения.
Изучил:
1. `emqx.conf` в поде — `authorization { no_match = allow }` — после CONNECT всё разрешено
2. `internal/api/handler/iot_device_handler.go``MQTTAuth` возвращает только `{"result":"allow"}` без ACL rules
### Вывод
**Auth (CONNECT) защищён:**
- HTTP auth endpoint проверяет namespace+deviceId+password (constant-time compare)
- enabled=true проверяется
- Secret изолирован по namespace
**ACL на pub/sub НЕТ:**
- `no_match = allow` — аутентифицированный клиент может SUBSCRIBE на любой топик
- EMQX HTTP auth plugin поддерживает возврат ACL rules в ответе на auth
- Формат ответа: `{"result":"allow","acl":[{"permission":"allow","action":"all","topic":"sless/ns/+"}]}`
- Текущий `mqttAuthResponse` содержит только `Result string` — ACL поле отсутствует
### Риски по приоритету
1. **Критично**: User A может SUBSCRIBE `sless/#` и читать все IoT данные всех пользователей
2. **Средне**: Нет rate limit на MQTT — один клиент может flood брокер
3. **Низко**: Нет TLS на 8083 (WebSocket без шифрования) — данные видны в сети
### План фикса
Добавить в `mqttAuthResponse` поле `ACL []aclRule` и возвращать из `MQTTAuth`:
```json
{
"result": "allow",
"acl": [
{"permission": "allow", "action": "publish", "topic": "sless/{ns}/{deviceId}"},
{"permission": "allow", "action": "subscribe", "topic": "sless/{ns}/{deviceId}"},
{"permission": "deny", "action": "all", "topic": "#"}
]
}
```
Ждём подтверждения от пользователя перед реализацией.
---
## Архитектурная дискуссия — IoT телеметрия и хранение данных
### Контекст разговора
Пользователь задал вопрос: "куда пишутся данные с IoT датчиков?"
Выяснилось что сейчас данные теряются — function pod получает событие но никуда не сохраняет. Это нормально для serverless (пользователь сам решает), но для IoT платформы нужно автоматическое хранение.
### Анализ сценариев использования
Реалистичные клиенты для Nubes (облачный провайдер СНГ, малый/средний бизнес):
1. Мониторинг объектов (склады, серверные, торговые точки) — температура, влажность, протечка
2. Умные счётчики / ЖКХ — снятие показаний без выезда
3. Небольшое производство / агро — теплицы, мини-заводы
Общий паттерн для всех: датчик → данные в БД → алерт если порог → график
### Решение по хранению данных
**Вопрос**: один большой Postgres или отдельный на каждого?
**Ответ**: один Postgres инстанс, но отдельная DATABASE на каждого tenant.
Причины:
- Вариант со одной таблицей + tenant_id — изоляция программная, ошибка в коде = утечка
- Отдельная DATABASE — физическая изоляция, разные connection string, разные пароли
- Клиент B не может подключиться к DATABASE клиента A даже при баге в коде платформы
Структура:
```
Postgres инстанс
├── sless_platform DB — системные таблицы (tenants, invocations)
├── tenant_abc DB — только данные клиента A
└── tenant_def DB — только данные клиента B
```
### Решение по доступу клиента
**Вопрос**: давать клиенту прямой доступ к Postgres?
**Ответ**: нет. Только через REST API платформы.
Причины:
- Postgres внутри кластера, снаружи не торчит (security)
- Единый endpoint `iot.kube5s.ru`
- Легко добавить rate limit, биллинг, кеш
- Клиент не зависит от деталей реализации хранилища
API:
```
GET /v1/namespaces/{ns}/iot/telemetry?device=X&from=T&to=T
GET /v1/namespaces/{ns}/iot/devices/{id}/last
```
### Решение по schema.sql
Клиент может положить `schema.sql` рядом с функцией. При деплое платформа выполняет его в БД tenant'а.
Это даёт низкий порог входа — клиент не шарит в Python, но может написать SQL по шаблону.
### Ключевое архитектурное решение — разделение операторов
**Решение**: sless-operator и iot-operator — ОТДЕЛЬНЫЕ компоненты.
Пока в одном кластере, но сделать так чтобы могли быть в разных.
**Namespace layout:**
```
namespace: sless — платформа sless (operator, event-dispatcher, RabbitMQ, Postgres invocations)
namespace: sless-{hash} — tenant функции (function pods)
namespace: iot — платформа IoT (iot-operator, EMQX, Postgres telemetry)
namespace: iot-{hash} — tenant IoT (IoTDevice CRDs)
```
**Связь**:
- Общий идентификатор tenant: `{hash}` одинаковый в обоих namespace
- MQTT событие → RabbitMQ в sless → function pod в sless-{hash}
- IoT operator НЕ импортирует пакеты sless-operator (loose coupling)
- Общение только через k8s API и RabbitMQ
**Postgres**:
- sless имеет свой Postgres (invocations)
- iot имеет свой Postgres (telemetry per tenant)
- Разные StatefulSet, разные PVC
### Что делает пользователь
Клиент:
1. Подключает устройство → данные автоматически пишутся в его `iot_telemetry`
2. Пишет функцию которая реагирует на события
3. Функция получает `DB_DSN` в env var (автоматически из Secret)
4. Может делать SELECT/INSERT в свою БД через обычный SQL в коде функции
5. Может читать телеметрию через REST API
### Plan — следующие шаги (этап IoT Postgres)
1. Поднять Postgres StatefulSet в namespace `iot`
2. В iot-operator при создании IoTDevice namespace → `CREATE USER`, `CREATE DATABASE`, `CREATE TABLE iot_telemetry`, `CREATE TABLE iot_devices`
3. Credentials → k8s Secret `iot-tenant-{ns}-pg`
4. В iot-mqtt-bridge при получении MQTT сообщения → INSERT в tenant БД
5. REST API endpoint для чтения телеметрии
6. При деплое function → прокинуть `DB_DSN` в env var из Secret
7. При деплое function → если есть `schema.sql` → выполнить в tenant БД
### Технические решения
- Postgres: `postgres:16-alpine` StatefulSet с PVC 10Gi в namespace `iot`
- Connection pool: pgxpool (pgx v5) per-tenant, lazy init, max 5 conn per tenant
- Таблица telemetry: `(id bigserial, device_id text, ts timestamptz default now(), payload jsonb)`
- Индекс: `(device_id, ts DESC)` для быстрых запросов по устройству за период
- Retention: пока без TTL, добавить позже через pg_partman или cron job
+108
View File
@@ -571,3 +571,111 @@ IoT event consumer → [interface] → POST → function
3. Kafka deployment: managed или Strimzi
Всё. Наш IoT-оператор, CRD, device auth — не меняются вообще.
---
## Агент: GitHub Copilot (Claude Sonnet 4.6) — ДЕПЛОЙ IoT MVP
### Задача от пользователя
Запустить всё. Документировать мысли и инженерную часть подробно.
### Анализ ситуации перед деплоем
**Что есть в git (ветка Ioter, коммиты 716efaf + 1e53766 + b920dc5):**
- CRD IoTDevice (iot/config/crd/bases/iot.kube5s.ru_iotdevices.yaml) — готово, не задеплоено
- IoT controller + IoT API handlers — в бинаре `manager`, НО образ не пересобран (v0.1.49 = без IoT)
- EMQX YAML (deployments/k8s/emqx.yaml) — готово, не задеплоено
- iot-mqtt-bridge YAML (deployments/k8s/iot-mqtt-bridge.yaml) — готово, не задеплоено
- mqtt-bridge бинарь (iot/cmd/mqtt-bridge/main.go) — есть код, но НЕ в образе
**Критические блокеры:**
1. **Dockerfile не включает `iot/`** — оператор собирается без IoT кода. Нет, стоп.
После проверки: iot/api включается через импорты в main.go — Go компилятор тянет весь нужный код.
НО: Dockerfile копирует только `main.go api/ controllers/ internal/ migrations/`.
IoT код живёт в `iot/` — и этот каталог НЕ копируется в Dockerfile!
Значит: `go build -o manager main.go` внутри Docker упадёт — не найдёт пакеты `iot/api/v1alpha1`, `iot/controllers/`.
2. **mqtt-bridge бинарь** — отдельная точка входа `iot/cmd/mqtt-bridge/main.go`.
В текущем Dockerfile собирается только `manager`. Нужно добавить сборку `iot-mqtt-bridge`.
3. **iot-mqtt-bridge.yaml** использует образ оператора (`sless-operator:latest`) с командой `/iot-mqtt-bridge`.
Значит оба бинаря должны быть в одном образе. Это правильный подход — один образ, разные entrypoint.
**Порядок действий (решение):**
```
1. Обновить Dockerfile:
- Добавить COPY iot/ iot/
- Добавить RUN go build -o iot-mqtt-bridge iot/cmd/mqtt-bridge/main.go
- Добавить COPY --from=builder /workspace/iot-mqtt-bridge .
2. Обновить operator.yaml: тег v0.1.49 → v0.1.50
3. SSH: docker build + push v0.1.50
4. kubectl apply CRD IoTDevice (один раз, cluster-wide)
5. kubectl apply EMQX (EMQX deployment + svc + configmap)
6. kubectl apply operator v0.1.50 (подхватит IoT controller + IoT API)
7. Bootstrap mqtt-bridge:
- Оператор должен быть живым (шаг 6)
- Создать IoTDevice "iot-bridge" через API → контроллер сгенерирует Secret в namespace sless-bridge
- Из Secret взять mqtt_username + mqtt_password
- kubectl create secret generic iot-bridge-credentials -n sless
- kubectl apply iot-mqtt-bridge.yaml
8. Проверка end-to-end
```
**Риски и как их обходить:**
- `sless-bridge` namespace может не существовать → создать заранее через kubectl
- EMQX может быть не готов к моменту запуска bridge → bridge сам делает retry (в коде есть reconnect loop)
- IoT API требует JWT-токен → при bootstrap curl с токеном из sless-operator-secret
**Почему один образ для operator + bridge:**
Это не идеально с т.з. SRP, но практично:
- Не нужен отдельный CI pipeline
- Не нужен отдельный registry repo
- Bridge — простой процесс (~100 строк Go), не нагружает образ
- В будущем можно разделить, порог изменений низкий
**Итог по мышлению:** Plan is solid. Начинаю выполнение.
### Проблемы, найденные при выполнении (до → решение)
**Проблема 1 — RBAC не настроен для iot.kube5s.ru:**
- Попытка создать IoTDevice через API → 403 Forbidden
- `sless-operator` ServiceAccount не имел прав на `iotdevices.iot.kube5s.ru`
- Причина: CRD для IoT — новая API-группа, в rbac.yaml её не было
- Решение: добавил в ClusterRole правила на `iot.kube5s.ru` (get/list/watch/create/update/patch/delete + status + finalizers)
- `kubectl apply -f rbac.yaml` → configured
- Вывод: при добавлении нового CRD API group ВСЕГДА нужно обновлять ClusterRole
**Проблема 2 — EMQX 5.x требует обязательные поля node.cookie и node.data_dir:**
- EMQX CrashLoopBackOff с ошибкой: `required_field: node.cookie, node.data_dir`
- В нашем emqx.conf (HOCON) эти поля отсутствовали — думал что для single-node они необязательны
- На самом деле в EMQX 5.x они mandatory (в отличие от 4.x где были defaults)
- Решение: добавил `node {}` секцию: name=emqx@127.0.0.1, cookie=sless-emqx-cookie-mvp, data_dir=/opt/emqx/data
- kubectl apply обновил ConfigMap, rollout restart → EMQX поднялся
- Вывод: при обновлении ConfigMap Deployment не перезапускается автоматически — нужен `kubectl rollout restart`
**Проблема 3 — kubectl logs берёт старый (crashing) pod:**
- deployment/emqx — логи шли со старого пода в CrashLoopBackOff
- Нужно указывать pod name явно для нового пода
- Это нормальное поведение kubectl — нет флага "новый pod"
### Итоговый статус деплоя
```
emqx-6f9689fc99-4mbhr 1/1 Running ✅
iot-mqtt-bridge-7d784d7d6b-n45fp 1/1 Running ✅ (3 restarts — reconnect loop до старта EMQX)
sless-operator-579dd6dcd5-fk2n8 1/1 Running ✅
```
CRD применён: `iotdevices.iot.kube5s.ru created`
IoTDevice iot-bridge создан: phase=Active, credentials в secret iot-iot-bridge
Secret iot-bridge-credentials создан в namespace sless
+95 -21
View File
@@ -58,19 +58,19 @@ type iotDeviceUpdateRequest struct {
// iotDeviceResponse — ответ при чтении одного IoTDevice.
// MQTTPassword заполняется только из GetIoTDevice (чтение из Secret).
type iotDeviceResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
DeviceID string `json:"device_id"`
Enabled bool `json:"enabled"`
Phase iotv1alpha1.IoTDevicePhase `json:"phase"`
MQTTUsername string `json:"mqtt_username,omitempty"`
MQTTPassword string `json:"mqtt_password,omitempty"` // только в GET /devices/{name}
SecretName string `json:"secret_name,omitempty"`
TopicPrefix string `json:"topic_prefix,omitempty"`
LastConnected string `json:"last_connected,omitempty"`
Message string `json:"message,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
CreatedAt string `json:"created_at,omitempty"`
Name string `json:"name"`
Namespace string `json:"namespace"`
DeviceID string `json:"device_id"`
Enabled bool `json:"enabled"`
Phase iotv1alpha1.IoTDevicePhase `json:"phase"`
MQTTUsername string `json:"mqtt_username,omitempty"`
MQTTPassword string `json:"mqtt_password,omitempty"` // только в GET /devices/{name}
SecretName string `json:"secret_name,omitempty"`
TopicPrefix string `json:"topic_prefix,omitempty"`
LastConnected string `json:"last_connected,omitempty"`
Message string `json:"message,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
CreatedAt string `json:"created_at,omitempty"`
}
// mqttAuthRequest — тело запроса от EMQX при MQTT CONNECT.
@@ -84,8 +84,20 @@ type mqttAuthRequest struct {
// mqttAuthResponse — ответ для EMQX. Всегда HTTP 200.
// result = "allow" | "deny"
// ACL — список правил pub/sub, изолирует топики по устройству.
type mqttAuthResponse struct {
Result string `json:"result"`
Result string `json:"result"`
ACL []aclRule `json:"acl,omitempty"`
}
// aclRule — одно правило ACL для EMQX HTTP auth plugin.
// permission: "allow" | "deny"
// action: "publish" | "subscribe" | "all"
// topic: точный топик или wildcard (#, +)
type aclRule struct {
Permission string `json:"permission"`
Action string `json:"action"`
Topic string `json:"topic"`
}
// ——————————————————————————————————————————
@@ -126,11 +138,11 @@ func deviceToResponse(d *iotv1alpha1.IoTDevice, password string) iotDeviceRespon
// НЕ защищён JWT middleware — доступен только из кластера (путь /internal/).
//
// Логика аутентификации:
// 1. Распарсить username → namespace + deviceId
// 2. Получить Secret iot-{deviceId} в namespace
// 3. Constant-time сравнение пароля (защита от timing attacks)
// 4. Проверить что IoTDevice существует и enabled=true
// 5. Обновить status.lastConnected в IoTDevice
// 1. Распарсить username → namespace + deviceId
// 2. Получить Secret iot-{deviceId} в namespace
// 3. Constant-time сравнение пароля (защита от timing attacks)
// 4. Проверить что IoTDevice существует и enabled=true
// 5. Обновить status.lastConnected в IoTDevice
func (h *Handler) MQTTAuth(w http.ResponseWriter, r *http.Request) {
var req mqttAuthRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
@@ -189,8 +201,20 @@ func (h *Handler) MQTTAuth(w http.ResponseWriter, r *http.Request) {
// Продолжаем — это некритично, устройство всё равно авторизовано
}
// Проверки пройдены — разрешаем подключение
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "allow"})
// Проверки пройдены — разрешаем подключение.
// ACL ограничивает устройство только его собственным топиком:
// publish: {namespace}/{deviceId} (данные устройства)
// subscribe: {namespace}/{deviceId} (команды устройству, если нужны)
// deny all: всё остальное запрещено — нельзя читать чужие данные
ownerTopic := ns + "/" + deviceID + "/#"
writeJSON(w, http.StatusOK, mqttAuthResponse{
Result: "allow",
ACL: []aclRule{
{Permission: "allow", Action: "publish", Topic: ownerTopic},
{Permission: "allow", Action: "subscribe", Topic: ownerTopic},
{Permission: "deny", Action: "all", Topic: "#"},
},
})
}
// ——————————————————————————————————————————
@@ -351,3 +375,53 @@ func (h *Handler) UpdateIoTDevice(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, deviceToResponse(device, ""))
}
// ——————————————————————————————————————————
// MQTT ACL — авторизация pub/sub
// ——————————————————————————————————————————
// mqttAclRequest — тело запроса от EMQX при каждом publish/subscribe.
type mqttAclRequest struct {
Username string `json:"username"`
ClientID string `json:"clientid"`
Action string `json:"action"` // "publish" | "subscribe"
Topic string `json:"topic"`
}
// MQTTAcl — POST /internal/mqtt/acl
// Вызывается EMQX для каждого pub/sub действия.
// НЕ защищён JWT — доступен только из кластера.
//
// Логика: клиент видит только топики вида {namespace}/{deviceId}/#
// Любой другой топик — deny и disconnect.
func (h *Handler) MQTTAcl(w http.ResponseWriter, r *http.Request) {
var req mqttAclRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Парсим username → namespace + deviceId (формат: "{ns}_{deviceId}")
idx := strings.Index(req.Username, "_")
if idx < 0 {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
ns := req.Username[:idx]
deviceID := req.Username[idx+1:]
if ns == "" || deviceID == "" {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Разрешаем только топики этого устройства: {ns}/{deviceId}/...
// Используем strings.HasPrefix — wildcard не нужен, проверяем prefix реального топика.
allowedPrefix := ns + "/" + deviceID + "/"
exactMatch := ns + "/" + deviceID
if strings.HasPrefix(req.Topic, allowedPrefix) || req.Topic == exactMatch {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "allow"})
return
}
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
}
+3
View File
@@ -80,6 +80,9 @@ func NewRouter(h *handler.Handler, log *slog.Logger) http.Handler {
// MQTT Auth — БЕЗ JWT. Вызывается EMQX при MQTT CONNECT из кластера.
// /internal/ недоступен снаружи (Ingress не проксирует /internal/).
r.HandleFunc("/internal/mqtt/auth", h.MQTTAuth).Methods(http.MethodPost)
// MQTT ACL — БЕЗ JWT. Вызывается EMQX при каждом pub/sub для проверки прав.
// Изолирует клиента в пределах его топиков: {namespace}/{deviceId}/#
r.HandleFunc("/internal/mqtt/acl", h.MQTTAcl).Methods(http.MethodPost)
// Цепочка middleware: logging → (auth только для /v1/) → router
// /fn/ — без auth, /v1/ — с auth.
+2 -2
View File
@@ -27,8 +27,6 @@ import (
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/controllers"
iotv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/api/v1alpha1"
iotcontrollers "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/controllers"
slessapi "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api/handler"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder"
@@ -36,6 +34,8 @@ import (
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/harbor"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/postgres"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/s3"
iotv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/api/v1alpha1"
iotcontrollers "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/controllers"
//+kubebuilder:scaffold:imports
)