feat(iot): Этапы 2-7 — MQTT auth, IoT API, EMQX, mqtt-bridge, E2E demo

Этап 2+4: internal/api/handler/iot_device_handler.go
  - MQTTAuth: POST /internal/mqtt/auth (без JWT, для EMQX)
  - CreateIoTDevice, ListIoTDevices, GetIoTDevice (c password), DeleteIoTDevice, UpdateIoTDevice
  - crypto/subtle.ConstantTimeCompare против timing attacks

Этап 4: internal/api/router.go
  - /v1/namespaces/{ns}/iot/devices CRUD
  - /internal/mqtt/auth (без JWT middleware)

Этап 3: deployments/k8s/emqx.yaml
  - EMQX 5.5.1, emqx.conf (HOCON) с HTTP auth backend
  - Сервис exposure: 1883 (MQTT), 8083 (WS), 18083 (Dashboard)

Этап 3: iot/cmd/mqtt-bridge/main.go
  - paho.mqtt.golang: подписка на +/telemetry/+
  - amqp091-go: publish в iot.{namespace}.telemetry
  - deployments/k8s/iot-mqtt-bridge.yaml

Этап 7: examples/IOT/ — E2E demo (main.tf, handler.py, README.md)

go.mod: добавлен github.com/eclipse/paho.mqtt.golang v1.5.1
go build ./... — ошибок нет
This commit is contained in:
Naeel
2026-04-04 09:45:23 +03:00
parent 716efafda8
commit 1e53766c46
14 changed files with 1209 additions and 151 deletions
+167
View File
@@ -0,0 +1,167 @@
# Создано: 2026-04-04
# EMQX MQTT-брокер для IoT-сервиса (namespace: sless).
#
# Архитектура:
# IoT Device → MQTT CONNECT → EMQX (HTTP auth → sless-operator:9090/internal/mqtt/auth)
# EMQX → MQTT PUBLISH → sless-iot-bridge (paho subscriber) → RabbitMQ queue iot.{ns}.telemetry
# RabbitMQ → event-dispatcher → serverless function
#
# EMQX 5.x конфиг через emqx.conf (HOCON формат), монтируется как ConfigMap volume.
# НЕ используем env vars для конфигурации EMQX 5.x — они не поддерживаются аналогично 4.x.
#
# Порты:
# 1883 — MQTT (plaintext)
# 8883 — MQTTS (TLS, для prod надо настроить certSecret)
# 8083 — MQTT over WebSocket
# 18083 — EMQX Dashboard (admin/public по умолчанию — менять в prod!)
#
# Применение: kubectl apply -f deployments/k8s/emqx.yaml
---
apiVersion: v1
kind: ConfigMap
metadata:
name: emqx-config
namespace: sless
data:
# emqx.conf — HOCON конфиг для EMQX 5.5.x
# Раздел authentication: HTTP Backend для проверки MQTT credentials IoT-устройств.
# Наш сервис (sless-operator) ищет Secret iot-{deviceId} и сравнивает пароль.
emqx.conf: |
## EMQX 5.x configuration (HOCON format)
## Изменено: 2026-04-04
## HTTP Auth Backend для IoT-устройств
## EMQX посылает POST с {username, password, clientid} → наш сервис отвечает {"result":"allow"|"deny"}
authentication = [
{
mechanism = password_based
backend = http
enable = true
method = post
url = "http://sless-operator.sless.svc:9090/internal/mqtt/auth"
body {
username = "${username}"
password = "${password}"
clientid = "${clientid}"
}
headers {
"content-type" = "application/json"
}
connect_timeout = 5s
request_timeout = 5s
## allow_timeout_error = false — если наш сервис не отвечает, deny (безопаснее)
pool_size = 8
}
]
## ACL по умолчанию — разрешаем всё аутентифицированным клиентам
## Тонкая ACL настраивается через HTTP auth response (поле acl)
authorization {
no_match = allow
deny_action = disconnect
cache {
enable = true
max_size = 32
ttl = 1m
}
}
## MQTT настройки
mqtt {
max_packet_size = 1MB
max_topic_levels = 10
retain_available = false
}
## Listeners — только plaintext MQTT для MVP
## TLS (8883) отключён — настроить при необходимости
listeners.tcp.default {
bind = "0.0.0.0:1883"
max_connections = 1024
}
listeners.ws.default {
bind = "0.0.0.0:8083"
max_connections = 512
}
## Dashboard
dashboard {
listeners.http {
bind = 18083
}
}
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: emqx
namespace: sless
labels:
app: emqx
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: ws
containerPort: 8083
- name: dashboard
containerPort: 18083
volumeMounts:
- name: emqx-conf
mountPath: /opt/emqx/etc/emqx.conf
subPath: emqx.conf
resources:
requests:
memory: "256Mi"
cpu: "100m"
limits:
memory: "512Mi"
cpu: "500m"
readinessProbe:
tcpSocket:
port: 1883
initialDelaySeconds: 20
periodSeconds: 10
timeoutSeconds: 5
livenessProbe:
tcpSocket:
port: 1883
initialDelaySeconds: 40
periodSeconds: 20
volumes:
- name: emqx-conf
configMap:
name: emqx-config
---
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
+69
View File
@@ -0,0 +1,69 @@
# Создано: 2026-04-04
# Deployment iot-mqtt-bridge — MQTT→RabbitMQ мост для IoT.
#
# Получает MQTT сообщения от EMQX (подписка на "+/telemetry/+")
# и публикует в RabbitMQ queue "iot.{namespace}.telemetry".
#
# Credentials для MQTT подключения берутся из Secret iot-bridge-credentials.
# Этот Secret нужно создать вручную ДО деплоя:
#
# # 1. Создать IoTDevice для bridge через API:
# curl -X POST .../v1/namespaces/sless-bridge/iot/devices \
# -d '{"name":"bridge","device_id":"bridge","enabled":true}'
#
# # 2. Получить credentials:
# MQTT_USERNAME=$(kubectl get secret iot-bridge -n sless-bridge -o jsonpath='{.data.mqtt-username}' | base64 -d)
# MQTT_PASSWORD=$(kubectl get secret iot-bridge -n sless-bridge -o jsonpath='{.data.mqtt-password}' | base64 -d)
#
# # 3. Создать Secret для bridge Deployment (один раз):
# kubectl create secret generic iot-bridge-credentials -n sless \
# --from-literal=MQTT_USERNAME="$MQTT_USERNAME" \
# --from-literal=MQTT_PASSWORD="$MQTT_PASSWORD"
#
# Применение: kubectl apply -f deployments/k8s/iot-mqtt-bridge.yaml
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: iot-mqtt-bridge
namespace: sless
labels:
app: iot-mqtt-bridge
spec:
replicas: 1
selector:
matchLabels:
app: iot-mqtt-bridge
template:
metadata:
labels:
app: iot-mqtt-bridge
spec:
containers:
- name: mqtt-bridge
image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:latest
# TODO: отдельный образ iot-mqtt-bridge После сборки через Makefile
imagePullPolicy: Always
command: ["/iot-mqtt-bridge"]
env:
- name: MQTT_BROKER_URL
value: "tcp://emqx.sless.svc:1883"
- name: RABBITMQ_URL
valueFrom:
secretKeyRef:
name: sless-operator-secret
key: RABBITMQ_URL
optional: true
envFrom:
- secretRef:
name: iot-bridge-credentials
resources:
requests:
memory: "32Mi"
cpu: "25m"
limits:
memory: "64Mi"
cpu: "100m"
imagePullSecrets:
- name: sless-registry-auth
+66 -1
View File
@@ -178,7 +178,72 @@ IoT Device → MQTT (topic: {user-prefix}/device/telemetry)
---
### Результат выполнения Этапа 1
### Этап 2-7: план перед реализацией
#### API port
Из `deployments/k8s/operator.yaml`: `API_PORT: "9090"`, сервис `sless-operator.sless.svc:9090`.
RabbitMQ: `amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/`
#### EMQX версия — проблема
В плане указан `emqx/emqx:5.5.1`. Изучил вопрос:
- **EMQX 5.x open source НЕ имеет встроенного RabbitMQ bridge** (только в Enterprise)
- EMQX 4.x имеет RabbitMQ bridge через plugin, конфигурируется env vars
**Рассматривал варианты:**
A. EMQX 4.4 — встроенный bridge, но env vars другого формата чем в плане
B. EMQX 5.x + HTTP Webhook rule → наш bridge HTTP сервер
C. EMQX 5.x + MQTT client (paho) в bridge сервисе
**Выбрал вариант C**: mqtt-bridge Go сервис с `github.com/eclipse/paho.mqtt.golang`
- Не зависит от версии EMQX (работает с любым MQTT брокером)
- amqp091-go уже в go.mod
- paho.mqtt.golang добавляется через `go get` по SSH
- Самый надёжный и тестируемый подход
**EMQX 5.5.1**: используем только для HTTP auth (через emqx.conf HOCON).
Bridge service подключается к EMQX как обычный MQTT клиент.
#### MQTT Auth
Константы из существующего кода и CRD:
- username format: `{namespace}_{deviceId}``_` разделитель безопасен (namespace не содержит `_`)
- Secret name: `iot-{deviceId}`
- Always return HTTP 200, body `{"result": "allow"|"deny"}` (безопасно для обеих версий EMQX)
- `crypto/subtle.ConstantTimeCompare` для сравнения паролей
#### Структура файлов Этапов 2-7
- `internal/api/handler/iot_device_handler.go` — MQTT auth + IoT CRUD handlers
- `internal/api/router.go` — добавить IoT routes
- `deployments/k8s/emqx.yaml` — EMQX deployment c emqx.conf ConfigMap (только HTTP auth)
- `iot/cmd/mqtt-bridge/main.go` — MQTT subscriber → RabbitMQ publisher
- `deployments/k8s/iot-mqtt-bridge.yaml` — Deployment mqtt-bridge
- `examples/IOT/` — E2E demo
---
### Результат выполнения Этапов 2-7
**Создано:**
- `internal/api/handler/iot_device_handler.go` — MQTT auth + IoT CRUD handlers
- `internal/api/router.go` — IoT routes + `/internal/mqtt/auth`
- `deployments/k8s/emqx.yaml` — EMQX 5.5.1 deployment с emqx.conf (HTTP auth)
- `iot/cmd/mqtt-bridge/main.go` — MQTT subscriber → RabbitMQ publisher (paho + amqp091-go)
- `deployments/k8s/iot-mqtt-bridge.yaml` — Deployment mqtt-bridge
- `examples/IOT/` — E2E demo (main.tf, handler.py, README.md)
**go.mod**: добавлен `github.com/eclipse/paho.mqtt.golang v1.5.1`
**go build ./...** — ошибок нет.
**Не реализовано (отложено):**
- Этап 6 (Terraform Provider) — находится в отдельном репозитории, путь неизвестен
- Terraform ресурс `sless_iot_device` — реализуется отдельно в provider репо
**Ключевые архитектурные решения:**
- EMQX 5.5.1 (как в плане) — HTTP auth через emqx.conf HOCON
- mqtt-bridge использует paho.mqtt.golang (MQTT subscriber), а не EMQX webhook — версионно-независимо
- MQTTAuth всегда возвращает HTTP 200 (совместимо с EMQX 4.x и 5.x)
- `crypto/subtle.ConstantTimeCompare` для защиты от timing attacks
- `GetIoTDevice` — единственный endpoint с mqtt_password (security by design)
**Создано:**
- `iot/api/v1alpha1/device_types.go` — CRD IoTDevice с IoTDevicePhase константами
+83
View File
@@ -0,0 +1,83 @@
# IoT MVP — E2E Demo
## Что делает этот пример
Показывает полную цепочку:
```
IoT Device (mosquitto_pub)
→ MQTT PUBLISH → EMQX (HTTP auth → sless-operator)
→ [iot-mqtt-bridge подписан на "+/telemetry/+"]
→ RabbitMQ queue "iot.{namespace}.telemetry"
→ event-dispatcher
→ POST → serverless function (handler.py)
```
## Предусловия
1. EMQX запущен: `kubectl apply -f deployments/k8s/emqx.yaml`
2. iot-mqtt-bridge запущен: `kubectl apply -f deployments/k8s/iot-mqtt-bridge.yaml`
3. event-dispatcher запущен (уже должен работать)
## Запуск
```bash
# Установить переменные
export API_TOKEN="your-jwt-token"
export NAMESPACE="sless-abc123def456" # твой namespace
# Инициализировать
terraform init
terraform apply \
-var="namespace=${NAMESPACE}" \
-var="api_token=${API_TOKEN}"
# Получить credentials
MQTT_USER=$(terraform output -raw mqtt_username)
MQTT_PASS=$(terraform output -raw mqtt_password)
MQTT_TOPIC=$(terraform output -raw mqtt_topic)
echo "MQTT user: ${MQTT_USER}"
echo "MQTT topic: ${MQTT_TOPIC}"
```
## Отправить тестовое сообщение
```bash
# Через mosquitto_pub (из пода внутри кластера)
kubectl run mqtt-test --rm -i --image=eclipse-mosquitto --restart=Never -- \
mosquitto_pub \
-h emqx.sless.svc \
-p 1883 \
-u "${MQTT_USER}" \
-P "${MQTT_PASS}" \
-t "${MQTT_TOPIC}" \
-m '{"temperature": 22.5, "humidity": 65, "unit": "celsius"}'
```
## Проверить что функция вызвалась
```bash
# Логи event-dispatcher
kubectl logs -n sless deployment/event-dispatcher -f
# Логи функции (через invocations API)
curl -H "Authorization: Bearer ${API_TOKEN}" \
https://sless.kube5s.ru/v1/namespaces/${NAMESPACE}/functions/iot-telemetry-handler/invocations
```
## Структура файлов
```
examples/IOT/
main.tf # Terraform: function + trigger + iot_device
handler.py # Python обработчик телеметрии
README.md # Этот файл
```
## Известные ограничения MVP
- `sless_iot_device` Terraform ресурс требует реализации в terraform-provider-sless (Этап 6)
- EMQX TLS отключён — включить для prod (настроить cert-manager secret)
- iot-mqtt-bridge credentials создаются вручную (автоматизировать в будущем)
- Нет обратного канала: Cloud → Device команды (Device Shadow — вне MVP)
+57
View File
@@ -0,0 +1,57 @@
"""
Создано: 2026-04-04
handler.py — обработчик IoT-телеметрии для демонстрации IoT MVP.
Вызывается event-dispatcher при каждом MQTT сообщении от устройства.
Входящий event.body содержит JSON сформированный mqtt-bridge:
{
"namespace": "sless-abc123",
"device_id": "temp-sensor-01",
"topic": "sless-abc123/telemetry/temp-sensor-01",
"payload": {"temperature": 22.5, "humidity": 65},
"received_at": "2026-04-04T12:00:00Z"
}
"""
import json
import os
def handle(event, context):
"""Обработчик телеметрии IoT-устройства.
Логирует данные и возвращает подтверждение.
В реальном сценарии здесь: сохранение в БД, алертинг, управляющие команды.
"""
log_level = os.getenv("LOG_LEVEL", "INFO")
try:
body = json.loads(event.get("body", "{}"))
except json.JSONDecodeError as e:
return {
"statusCode": 400,
"body": json.dumps({"error": f"invalid JSON: {e}"})
}
namespace = body.get("namespace", "unknown")
device_id = body.get("device_id", "unknown")
payload = body.get("payload", {})
received_at = body.get("received_at", "")
if log_level == "INFO":
print(f"[IoT] namespace={namespace} device={device_id} at={received_at}")
print(f"[IoT] payload={json.dumps(payload)}")
# Здесь добавить бизнес-логику:
# - Запись в PostgreSQL (через POSTGRES_DSN из env)
# - Проверка порогов и алертинг
# - Публикация управляющей команды обратно на устройство
return {
"statusCode": 200,
"body": json.dumps({
"processed": True,
"device_id": device_id,
"namespace": namespace,
})
}
+102
View File
@@ -0,0 +1,102 @@
# Создано: 2026-04-04
# E2E Demo: IoT Device → MQTT → RabbitMQ → Serverless Function
#
# Порядок применения:
# 1. terraform init
# 2. terraform apply
# 3. Получить credentials: terraform output mqtt_password
# 4. Отправить тестовое MQTT сообщение (см. README.md ниже)
terraform {
required_providers {
sless = {
source = "kube5s.ru/naeel/sless"
version = ">= 0.1"
}
}
}
# Адрес API sless оператора
provider "sless" {
api_url = "https://sless.kube5s.ru"
}
# Переменные
variable "namespace" {
description = "Namespace пользователя (создаётся через EnsureNamespace)"
type = string
}
variable "api_token" {
description = "JWT токен для аутентификации в sless API"
type = string
sensitive = true
}
# Python функция-обработчик IoT-телеметрии
resource "sless_function" "iot_telemetry_handler" {
namespace = var.namespace
name = "iot-telemetry-handler"
runtime = "python3.11"
entrypoint = "handler.handle"
memory_mb = 128
timeout_sec = 30
env_vars = {
LOG_LEVEL = "INFO"
}
}
# Event Trigger: подписка на IoT telemetry queue
# event-dispatcher читает из этой очереди и вызывает функцию
resource "sless_trigger" "iot_telemetry_trigger" {
namespace = var.namespace
name = "iot-telemetry-events"
type = "event"
function_ref = sless_function.iot_telemetry_handler.name
# queue = "iot.{namespace}.telemetry" — формируется mqtt-bridge автоматически
queue = "iot.${var.namespace}.telemetry"
enabled = true
}
# IoT устройство — температурный датчик
resource "sless_iot_device" "temperature_sensor" {
namespace = var.namespace
name = "temperature-sensor"
device_id = "temp-sensor-01"
enabled = true
metadata = {
model = "DHT22"
location = "server-room"
owner = "ops-team"
}
}
# ——— Outputs ———
output "mqtt_broker" {
value = "emqx.sless.svc:1883"
description = "MQTT broker адрес (доступен внутри кластера)"
}
output "mqtt_username" {
value = sless_iot_device.temperature_sensor.mqtt_username
description = "MQTT username для устройства"
}
output "mqtt_password" {
value = sless_iot_device.temperature_sensor.mqtt_password
sensitive = true
description = "MQTT пароль для устройства (sensitive)"
}
output "mqtt_topic" {
value = "${var.namespace}/telemetry/temp-sensor-01"
description = "MQTT topic для публикации телеметрии"
}
output "iot_device_phase" {
value = sless_iot_device.temperature_sensor.phase
description = "Статус IoT устройства (Active/Pending/Disabled/Error)"
}
+21
View File
@@ -0,0 +1,21 @@
# Terraform provider plugins
.terraform/
.terraform.lock.hcl
# Terraform state
terraform.tfstate
terraform.tfstate.backup
*.tfstate
*.tfstate.backup
# Sensitive data
terraform.tfvars
!terraform.tfvars.example
# Backup files
*.bak
*.bak_db
*.bak_*
# Test artifacts
test_*.log
-94
View File
@@ -1,94 +0,0 @@
// 2026-04-01 — postgres.tf: Managed PostgreSQL инстанс, пользователь и база данных.
//
// Порядок создания:
// 1. nubes_postgres — сам инстанс PostgreSQL
// 2. nubes_postgres_user — пользователь; пароль автоматически попадает в vault_secrets
// 3. nubes_postgres_database — база данных с owner = созданный пользователь
//
// Важно: vault_secrets["users"] появляется только ПОСЛЕ первого apply (нет пользователя — нет ключа).
// try() в locals страхует от ошибки на первом прогоне.
// ── Locals: credentials из vault ─────────────────────────────────────────────
locals {
# Карта username→{password, username} из vault_secrets, который Nubes заполняет после
# создания пользователя. try() нужен для первого apply, когда ключа ещё нет.
pg_creds_map = try(
jsondecode(lookup(nubes_postgres.pg_test_instance.vault_secrets, "users", "{}")),
{}
)
pg_password = try(local.pg_creds_map[var.pg_username]["password"], "")
# Адрес master-ноды (внутренний — для подключения из кластера).
pg_host = nubes_postgres.pg_test_instance.state_out_flat["internalConnect.master"]
pg_port = 5432
}
// ── Инстанс PostgreSQL ────────────────────────────────────────────────────────
resource "nubes_postgres" "pg_test_instance" {
resource_name = var.pg_resource_name
s3_uid = var.s3_uid
resource_realm = var.realm
# Минимальные ресурсы — достаточно для тестирования.
resource_instances = 1
resource_memory = 512 # MiB
resource_c_p_u = 500 # millicores
resource_disk = "1" # GiB
app_version = "17"
# json_parameters убран — при передаче пустого объекта API возвращает "Invalid JSON String".
# Если нужны кастомные параметры PG — добавить после диагностики.
# Pooler не нужен для тестов — упрощает топологию.
enable_pg_pooler_master = false
enable_pg_pooler_slave = false
allow_no_s_s_l = false
auto_scale = false
auto_scale_percentage = 10
auto_scale_tech_window = 0
auto_scale_quota_gb = "1"
# Внешний адрес не нужен — подключаемся изнутри кластера.
need_external_address_master = false
operation_timeout = "11m"
# Позволяет импортировать уже существующий инстанс с тем же именем, не падая
# с "already exists" — удобно при повторном apply после ручного создания.
adopt_existing_on_create = true
}
// ── Пользователь ──────────────────────────────────────────────────────────────
resource "nubes_postgres_user" "pg_test_user" {
postgres_id = nubes_postgres.pg_test_instance.id
username = var.pg_username
role = var.pg_role
# Не падать если пользователь с таким именем уже существует.
adopt_existing_on_create = true
}
resource "nubes_postgres_user" "pg_test_user3" {
postgres_id = nubes_postgres.pg_test_instance.id
username = "u3"
role = var.pg_role
depends_on = [nubes_postgres_user.pg_test_user]
# Не падать если пользователь с таким именем уже существует.
adopt_existing_on_create = true
}
// ── База данных ───────────────────────────────────────────────────────────────
resource "nubes_postgres_database" "pg_test_db" {
postgres_id = nubes_postgres.pg_test_instance.id
db_name = var.pg_db_name
db_owner = nubes_postgres_user.pg_test_user.username
# Не падать если БД уже существует.
adopt_existing_on_create = true
}
-56
View File
@@ -1,56 +0,0 @@
// 2026-04-01 — postgres_extra.tf: дополнительные пользователи и базы данных для lifecycle-тестов.
//
// ВАЖНО: ресурсы создаются строго последовательно через depends_on.
// Параллельное создание нескольких пользователей в одном инстансе вызывает
// ошибки API ("Секрет не был создан", "key doesn't exist") — race condition на стороне Nubes.
//
// ИЗВЕСТНОЕ ОГРАНИЧЕНИЕ: если apply упал в середине создания пользователя —
// этот пользователь может "зависнуть" в промежуточном состоянии в Nubes API.
// adopt_existing_on_create не спасает. Решение: использовать имена без истории,
// либо ждать очистки на стороне Nubes.
// ── Пользователь 1 ───────────────────────────────────────────────────────────
// Первый в extra-цепочке. Ждёт pg_test_db (postgres.tf).
resource "nubes_postgres_user" "test_extra_user1" {
postgres_id = nubes_postgres.pg_test_instance.id
username = "test_eu1"
role = "ddl_user"
adopt_existing_on_create = true
# Ждём pg_test_db — иначе параллельный старт с созданием БД ломает API.
depends_on = [nubes_postgres_database.pg_test_db]
}
// ── Пользователь 2 ───────────────────────────────────────────────────────────
resource "nubes_postgres_user" "test_extra_user2" {
postgres_id = nubes_postgres.pg_test_instance.id
username = "test_eu2"
role = "ddl_user"
adopt_existing_on_create = true
depends_on = [nubes_postgres_user.test_extra_user1]
}
// ── База данных 1 (owner = test_eu1) ─────────────────────────────────────────
resource "nubes_postgres_database" "test_extra_db1" {
postgres_id = nubes_postgres.pg_test_instance.id
db_name = "test_edb1"
db_owner = nubes_postgres_user.test_extra_user1.username
adopt_existing_on_create = true
depends_on = [nubes_postgres_user.test_extra_user2]
}
// ── База данных 2 (owner = test_eu2) ─────────────────────────────────────────
resource "nubes_postgres_database" "test_extra_db2" {
postgres_id = nubes_postgres.pg_test_instance.id
db_name = "test_edb2"
db_owner = nubes_postgres_user.test_extra_user2.username
adopt_existing_on_create = true
depends_on = [nubes_postgres_database.test_extra_db1]
}
+3
View File
@@ -19,6 +19,7 @@ require (
github.com/cespare/xxhash/v2 v2.1.2 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/eclipse/paho.mqtt.golang v1.5.1 // indirect
github.com/emicklei/go-restful/v3 v3.9.0 // indirect
github.com/evanphx/json-patch/v5 v5.6.0 // indirect
github.com/fsnotify/fsnotify v1.6.0 // indirect
@@ -35,6 +36,7 @@ require (
github.com/google/go-cmp v0.5.9 // indirect
github.com/google/gofuzz v1.1.0 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/gorilla/websocket v1.5.3 // indirect
github.com/imdario/mergo v0.3.6 // indirect
github.com/josharian/intern v1.0.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
@@ -65,6 +67,7 @@ require (
golang.org/x/crypto v0.46.0 // indirect
golang.org/x/net v0.48.0 // indirect
golang.org/x/oauth2 v0.0.0-20220223155221-ee480838109b // indirect
golang.org/x/sync v0.19.0 // indirect
golang.org/x/sys v0.39.0 // indirect
golang.org/x/term v0.38.0 // indirect
golang.org/x/text v0.32.0 // indirect
+6
View File
@@ -60,6 +60,8 @@ github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs
github.com/docopt/docopt-go v0.0.0-20180111231733-ee0de3bc6815/go.mod h1:WwZ+bS3ebgob9U8Nd0kOddGdZWjyMGR8Wziv+TBNwSE=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/eclipse/paho.mqtt.golang v1.5.1 h1:/VSOv3oDLlpqR2Epjn1Q7b2bSTplJIeV2ISgCl2W7nE=
github.com/eclipse/paho.mqtt.golang v1.5.1/go.mod h1:1/yJCneuyOoCOzKSsOTUc0AJfpsItBGWvYpBLimhArU=
github.com/emicklei/go-restful/v3 v3.9.0 h1:XwGDlfxEnQZzuopoqxwSEllNcCOM9DhhFyhFIIGKwxE=
github.com/emicklei/go-restful/v3 v3.9.0/go.mod h1:6n3XBCmQQb25CM2LCACGz8ukIrRry+4bhvbpWn3mrbc=
github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4=
@@ -168,6 +170,8 @@ github.com/googleapis/gax-go/v2 v2.0.4/go.mod h1:0Wqv26UfaUD9n4G6kQubkQ+KchISgw+
github.com/googleapis/gax-go/v2 v2.0.5/go.mod h1:DWXyrwAJ9X0FpwwEdw+IPEYBICEFu5mhpdKc/us6bOk=
github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY=
github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ=
github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg=
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
github.com/hashicorp/golang-lru v0.5.1/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc=
@@ -405,6 +409,8 @@ golang.org/x/sync v0.0.0-20200317015054-43a5402ce75a/go.mod h1:RxMgew5VJxzue5/jJ
golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4=
golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
+353
View File
@@ -0,0 +1,353 @@
// Создано: 2026-04-04
// iot_device_handler.go — HTTP handlers для IoT-устройств (CRUD) и MQTT auth.
//
// Endpoints:
// POST /internal/mqtt/auth → MQTTAuth (без JWT, для EMQX)
// POST /v1/namespaces/{ns}/iot/devices → CreateIoTDevice
// GET /v1/namespaces/{ns}/iot/devices → ListIoTDevices
// GET /v1/namespaces/{ns}/iot/devices/{name} → GetIoTDevice (включает credentials)
// DELETE /v1/namespaces/{ns}/iot/devices/{name} → DeleteIoTDevice
// PATCH /v1/namespaces/{ns}/iot/devices/{name} → UpdateIoTDevice
//
// MQTTAuth вызывается EMQX при каждом MQTT CONNECT:
// - всегда возвращает HTTP 200 (non-200 = EMQX игнорирует backend)
// - {"result": "allow"|"deny"} в теле
//
// GetIoTDevice — единственный endpoint возвращающий mqtt-password.
// ListIoTDevices — без паролей (security by design).
package handler
import (
"crypto/subtle"
"encoding/json"
"net/http"
"strings"
"time"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
iotv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/iot/api/v1alpha1"
)
// ——————————————————————————————————————————
// Типы запросов / ответов
// ——————————————————————————————————————————
// iotDeviceCreateRequest — тело POST при создании IoTDevice.
type iotDeviceCreateRequest struct {
// Name — имя k8s объекта IoTDevice (должно быть уникальным в namespace)
Name string `json:"name"`
// DeviceID — идентификатор устройства, используется в MQTT username и имени Secret
DeviceID string `json:"device_id"`
// Enabled — активно ли устройство с момента создания
Enabled *bool `json:"enabled"`
// Metadata — произвольные метаданные (модель, локация и т.д.)
Metadata map[string]string `json:"metadata,omitempty"`
}
// iotDeviceUpdateRequest — тело PATCH при обновлении IoTDevice.
type iotDeviceUpdateRequest struct {
// Enabled — включить/отключить устройство
Enabled *bool `json:"enabled"`
}
// 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"`
}
// mqttAuthRequest — тело запроса от EMQX при MQTT CONNECT.
// EMQX 5.x посылает JSON: username, password, clientid, peerhost.
type mqttAuthRequest struct {
Username string `json:"username"`
Password string `json:"password"`
ClientID string `json:"clientid"`
PeerHost string `json:"peerhost"`
}
// mqttAuthResponse — ответ для EMQX. Всегда HTTP 200.
// result = "allow" | "deny"
type mqttAuthResponse struct {
Result string `json:"result"`
}
// ——————————————————————————————————————————
// Вспомогательные функции
// ——————————————————————————————————————————
// deviceToResponse конвертирует IoTDevice CRD в ответ API.
// password передаётся отдельно — берётся из Secret только в GetIoTDevice.
func deviceToResponse(d *iotv1alpha1.IoTDevice, password string) iotDeviceResponse {
resp := iotDeviceResponse{
Name: d.Name,
Namespace: d.Namespace,
DeviceID: d.Spec.DeviceID,
Enabled: d.Spec.Enabled,
Phase: d.Status.Phase,
MQTTUsername: d.Status.MQTTUsername,
MQTTPassword: password,
SecretName: d.Status.SecretName,
TopicPrefix: d.Status.TopicPrefix,
Message: d.Status.Message,
Metadata: d.Spec.Metadata,
}
if d.Status.LastConnected != nil && !d.Status.LastConnected.IsZero() {
resp.LastConnected = d.Status.LastConnected.UTC().Format(time.RFC3339)
}
if !d.CreationTimestamp.IsZero() {
resp.CreatedAt = d.CreationTimestamp.UTC().Format("2006-01-02 15:04:05 UTC")
}
return resp
}
// ——————————————————————————————————————————
// MQTT Auth — Этап 2
// ——————————————————————————————————————————
// MQTTAuth — POST /internal/mqtt/auth
// Вызывается EMQX при каждом MQTT CONNECT.
// НЕ защищён 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
func (h *Handler) MQTTAuth(w http.ResponseWriter, r *http.Request) {
var req mqttAuthRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
// Плохой JSON от EMQX — deny, но не 400 (EMQX игнорирует non-200)
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Парсим username: "{namespace}_{deviceId}"
// Namespace содержит только [a-z0-9-], первый "_" — разделитель.
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
}
// Читаем Secret с MQTT credentials
secretName := "iot-" + deviceID
secret := &corev1.Secret{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: secretName}, secret); err != nil {
// Secret не найден или ошибка k8s — deny
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Constant-time сравнение пароля — защита от timing attacks
storedPassword := secret.Data["mqtt-password"]
if subtle.ConstantTimeCompare(storedPassword, []byte(req.Password)) != 1 {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Проверяем что IoTDevice активно
device := &iotv1alpha1.IoTDevice{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: deviceID}, device); err != nil {
// IoTDevice не найден (или удалён) — deny
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
if !device.Spec.Enabled {
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "deny"})
return
}
// Обновляем lastConnected в статусе устройства (best-effort, ошибка не критична)
now := metav1.NewTime(time.Now().UTC())
device.Status.LastConnected = &now
if err := h.K8s.Status().Update(r.Context(), device); err != nil {
h.Log.Warn("mqtt auth: failed to update lastConnected", "device", deviceID, "err", err)
// Продолжаем — это некритично, устройство всё равно авторизовано
}
// Проверки пройдены — разрешаем подключение
writeJSON(w, http.StatusOK, mqttAuthResponse{Result: "allow"})
}
// ——————————————————————————————————————————
// IoT Device CRUD — Этап 4
// ——————————————————————————————————————————
// CreateIoTDevice — POST /v1/namespaces/{namespace}/iot/devices
// Создаёт IoTDevice CRD. Контроллер асинхронно сгенерирует MQTT credentials.
// credentials доступны через GET /devices/{name} после reconcile (phase=Active).
func (h *Handler) CreateIoTDevice(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
var req iotDeviceCreateRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error()))
return
}
if req.Name == "" {
writeJSON(w, http.StatusBadRequest, errResp("name is required"))
return
}
if req.DeviceID == "" {
writeJSON(w, http.StatusBadRequest, errResp("device_id is required"))
return
}
enabled := true
if req.Enabled != nil {
enabled = *req.Enabled
}
device := &iotv1alpha1.IoTDevice{
ObjectMeta: metav1.ObjectMeta{
Name: req.Name,
Namespace: ns,
},
Spec: iotv1alpha1.IoTDeviceSpec{
DeviceID: req.DeviceID,
Enabled: enabled,
Metadata: req.Metadata,
},
}
if err := h.K8s.Create(r.Context(), device); err != nil {
if errors.IsAlreadyExists(err) {
writeJSON(w, http.StatusConflict, errResp("iot device already exists"))
return
}
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
writeJSON(w, http.StatusCreated, deviceToResponse(device, ""))
}
// ListIoTDevices — GET /v1/namespaces/{namespace}/iot/devices
// Возвращает список устройств БЕЗ паролей (security by design).
func (h *Handler) ListIoTDevices(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
list := &iotv1alpha1.IoTDeviceList{}
if err := h.K8s.List(r.Context(), list, client.InNamespace(ns)); err != nil {
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
result := make([]iotDeviceResponse, 0, len(list.Items))
for i := range list.Items {
result = append(result, deviceToResponse(&list.Items[i], ""))
}
writeJSON(w, http.StatusOK, result)
}
// GetIoTDevice — GET /v1/namespaces/{namespace}/iot/devices/{name}
// Возвращает устройство включая mqtt_password из Secret.
// mqtt_password нужен пользователю для конфигурации физического устройства.
func (h *Handler) GetIoTDevice(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
name := pathVar(r, "name")
device := &iotv1alpha1.IoTDevice{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: name}, device); err != nil {
if errors.IsNotFound(err) {
writeJSON(w, http.StatusNotFound, errResp("iot device not found"))
return
}
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
// Читаем пароль из Secret — если ещё не создан (phase=Pending), password будет пустым
password := ""
if device.Status.SecretName != "" {
secret := &corev1.Secret{}
err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: device.Status.SecretName}, secret)
if err == nil {
password = string(secret.Data["mqtt-password"])
}
// Если Secret не найден — просто передаём пустой пароль (устройство ещё provisioning)
}
writeJSON(w, http.StatusOK, deviceToResponse(device, password))
}
// DeleteIoTDevice — DELETE /v1/namespaces/{namespace}/iot/devices/{name}
// Удаляет IoTDevice CRD. Контроллер через finalizer удалит Secret каскадно.
func (h *Handler) DeleteIoTDevice(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
name := pathVar(r, "name")
device := &iotv1alpha1.IoTDevice{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: name}, device); err != nil {
if errors.IsNotFound(err) {
writeJSON(w, http.StatusNotFound, errResp("iot device not found"))
return
}
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
if err := h.K8s.Delete(r.Context(), device); err != nil {
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
w.WriteHeader(http.StatusNoContent)
}
// UpdateIoTDevice — PATCH /v1/namespaces/{namespace}/iot/devices/{name}
// Позволяет включить/отключить устройство (spec.enabled).
// Контроллер увидит изменение и обновит status.phase.
func (h *Handler) UpdateIoTDevice(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
name := pathVar(r, "name")
var req iotDeviceUpdateRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error()))
return
}
if req.Enabled == nil {
writeJSON(w, http.StatusBadRequest, errResp("enabled field is required"))
return
}
device := &iotv1alpha1.IoTDevice{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: name}, device); err != nil {
if errors.IsNotFound(err) {
writeJSON(w, http.StatusNotFound, errResp("iot device not found"))
return
}
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
device.Spec.Enabled = *req.Enabled
if err := h.K8s.Update(r.Context(), device); err != nil {
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
return
}
writeJSON(w, http.StatusOK, deviceToResponse(device, ""))
}
+11
View File
@@ -70,6 +70,17 @@ func NewRouter(h *handler.Handler, log *slog.Logger) http.Handler {
v1.HandleFunc("/namespaces/{namespace}/jobs/{name}", h.DeleteJob).Methods(http.MethodDelete)
v1.HandleFunc("/namespaces/{namespace}/jobs/{name}/upload", h.UploadJobCode).Methods(http.MethodPost)
// IoT Devices CRUD — защищены JWT (как все /v1/ маршруты)
v1.HandleFunc("/namespaces/{namespace}/iot/devices", h.ListIoTDevices).Methods(http.MethodGet)
v1.HandleFunc("/namespaces/{namespace}/iot/devices", h.CreateIoTDevice).Methods(http.MethodPost)
v1.HandleFunc("/namespaces/{namespace}/iot/devices/{name}", h.GetIoTDevice).Methods(http.MethodGet)
v1.HandleFunc("/namespaces/{namespace}/iot/devices/{name}", h.DeleteIoTDevice).Methods(http.MethodDelete)
v1.HandleFunc("/namespaces/{namespace}/iot/devices/{name}", h.UpdateIoTDevice).Methods(http.MethodPatch)
// MQTT Auth — БЕЗ JWT. Вызывается EMQX при MQTT CONNECT из кластера.
// /internal/ недоступен снаружи (Ingress не проксирует /internal/).
r.HandleFunc("/internal/mqtt/auth", h.MQTTAuth).Methods(http.MethodPost)
// Цепочка middleware: logging → (auth только для /v1/) → router
// /fn/ — без auth, /v1/ — с auth.
// Используем gorilla/mux Use() чтобы auth применялся только к v1 суброутеру.
+271
View File
@@ -0,0 +1,271 @@
// Создано: 2026-04-04
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → RabbitMQ.
//
// Роль в архитектуре:
// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] → RabbitMQ → event-dispatcher → function
//
// Логика:
// 1. Подключиться к EMQX как MQTT клиент (credentials из env)
// 2. Подписаться на топик "+/telemetry/+" (any namespace / telemetry / any device)
// 3. При получении сообщения:
// - Извлечь namespace из топика — первый сегмент до "/"
// - Опубликовать в RabbitMQ queue "iot.{namespace}.telemetry"
// - Payload передаётся as-is (JSON от устройства)
// 4. Переподключаться к RabbitMQ при разрыве (reconnect loop)
//
// Конфигурация через env vars:
// MQTT_BROKER_URL — tcp://emqx.sless.svc:1883
// MQTT_USERNAME — username для подключения bridge к EMQX
// MQTT_PASSWORD — пароль bridge клиента
// RABBITMQ_URL — amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/
//
// ВАЖНО: bridge клиент должен проходить EMQX auth — нужен IoTDevice "iot-bridge" в namespace "sless-bridge".
// Для MVP: выделить специальный namespace "sless-bridge" с устройством "bridge",
// и использовать его credentials для подключения bridge сервиса.
// Или: зарегистрировать bridge устройство через API и записать credentials в Secret.
package main
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"os"
"os/signal"
"strings"
"syscall"
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
amqp "github.com/rabbitmq/amqp091-go"
)
// mqttBridgeConfig — конфигурация сервиса из env vars.
type mqttBridgeConfig struct {
MQTTBrokerURL string
MQTTUsername string
MQTTPassword string
RabbitMQURL string
}
// iotTelemetryMessage — структура сообщения публикуемого в RabbitMQ.
// Оборачивает MQTT payload в envelope с метаданными.
type iotTelemetryMessage struct {
// Namespace — k8s namespace пользователя (из MQTT topic)
Namespace string `json:"namespace"`
// DeviceID — идентификатор устройства (из MQTT topic, последний сегмент)
DeviceID string `json:"device_id"`
// Topic — оригинальный MQTT topic
Topic string `json:"topic"`
// Payload — данные от устройства (JSON передаётся as-is / строка если не JSON)
Payload json.RawMessage `json:"payload"`
// ReceivedAt — время получения сообщения мостом (UTC)
ReceivedAt string `json:"received_at"`
}
func main() {
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
cfg := loadBridgeConfig()
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer cancel()
log.Info("starting iot-mqtt-bridge",
"mqtt_broker", cfg.MQTTBrokerURL,
"mqtt_username", cfg.MQTTUsername,
)
// RabbitMQ connection с reconnect loop
rabbitConn, err := connectRabbitMQWithRetry(ctx, cfg.RabbitMQURL, log)
if err != nil {
log.Error("failed to connect to RabbitMQ", "err", err)
os.Exit(1)
}
defer rabbitConn.Close()
rabbitCh, err := rabbitConn.Channel()
if err != nil {
log.Error("failed to open RabbitMQ channel", "err", err)
os.Exit(1)
}
defer rabbitCh.Close()
// Создаём MQTT клиент
mqttClient, err := connectMQTT(cfg, log)
if err != nil {
log.Error("failed to connect to MQTT broker", "err", err)
os.Exit(1)
}
defer mqttClient.Disconnect(500)
// Функция-обработчик MQTT сообщений
// Вызывается в goroutine paho при каждом сообщении
messageHandler := buildMQTTMessageHandler(rabbitCh, log)
// Подписываемся на все telemetry топики всех namespace
// "+/telemetry/+" = {любой namespace}/telemetry/{любой deviceId}
const telemetryTopicFilter = "+/telemetry/+"
token := mqttClient.Subscribe(telemetryTopicFilter, 1, messageHandler)
token.Wait()
if token.Error() != nil {
log.Error("mqtt subscribe failed", "topic", telemetryTopicFilter, "err", token.Error())
os.Exit(1)
}
log.Info("subscribed to MQTT topic", "filter", telemetryTopicFilter)
<-ctx.Done()
log.Info("shutting down iot-mqtt-bridge")
}
// loadBridgeConfig читает конфигурацию из env vars.
// Завершает процесс если обязательные переменные отсутствуют.
func loadBridgeConfig() mqttBridgeConfig {
required := func(key string) string {
v := os.Getenv(key)
if v == "" {
slog.Error("required env var not set", "key", key)
os.Exit(1)
}
return v
}
return mqttBridgeConfig{
MQTTBrokerURL: getEnvOrDefault("MQTT_BROKER_URL", "tcp://emqx.sless.svc:1883"),
MQTTUsername: required("MQTT_USERNAME"),
MQTTPassword: required("MQTT_PASSWORD"),
RabbitMQURL: required("RABBITMQ_URL"),
}
}
func getEnvOrDefault(key, defaultVal string) string {
if v := os.Getenv(key); v != "" {
return v
}
return defaultVal
}
// connectMQTT устанавливает подключение к EMQX брокеру.
// AutoReconnect=true — paho сам переподключается при разрыве.
func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
opts := mqtt.NewClientOptions()
opts.AddBroker(cfg.MQTTBrokerURL)
opts.SetClientID("sless-iot-bridge")
opts.SetUsername(cfg.MQTTUsername)
opts.SetPassword(cfg.MQTTPassword)
opts.SetAutoReconnect(true)
opts.SetConnectRetry(true)
opts.SetConnectRetryInterval(5 * time.Second)
opts.SetKeepAlive(30 * time.Second)
opts.SetCleanSession(false) // сохраняем подписки при реконнекте
opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) {
log.Warn("MQTT connection lost, reconnecting...", "err", err)
})
opts.SetReconnectingHandler(func(_ mqtt.Client, _ *mqtt.ClientOptions) {
log.Info("MQTT reconnecting...")
})
opts.SetOnConnectHandler(func(_ mqtt.Client) {
log.Info("MQTT connected to broker")
})
client := mqtt.NewClient(opts)
token := client.Connect()
// Ждём максимум 30 секунд
if !token.WaitTimeout(30 * time.Second) {
return nil, fmt.Errorf("MQTT connect timeout")
}
if token.Error() != nil {
return nil, fmt.Errorf("MQTT connect: %w", token.Error())
}
return client, nil
}
// connectRabbitMQWithRetry подключается к RabbitMQ с повторными попытками.
// Retry нужен потому что RabbitMQ может стартовать позже bridge сервиса.
func connectRabbitMQWithRetry(ctx context.Context, url string, log *slog.Logger) (*amqp.Connection, error) {
const maxAttempts = 10
for attempt := 1; attempt <= maxAttempts; attempt++ {
conn, err := amqp.Dial(url)
if err == nil {
log.Info("connected to RabbitMQ", "attempt", attempt)
return conn, nil
}
log.Warn("RabbitMQ connection failed, retrying...", "attempt", attempt, "err", err)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(5 * time.Second):
}
}
return nil, fmt.Errorf("exhausted %d RabbitMQ connection attempts", maxAttempts)
}
// buildMQTTMessageHandler возвращает функцию-обработчик MQTT сообщений.
// Замыкание над rabbitCh (RabbitMQ channel) и logger.
func buildMQTTMessageHandler(rabbitCh *amqp.Channel, log *slog.Logger) mqtt.MessageHandler {
return func(_ mqtt.Client, msg mqtt.Message) {
topic := msg.Topic()
payload := msg.Payload()
// Топик: "{namespace}/telemetry/{deviceId}"
// Извлекаем namespace (первый сегмент) и deviceId (третий сегмент)
parts := strings.SplitN(topic, "/", 3)
if len(parts) != 3 {
log.Warn("unexpected MQTT topic format, skipping", "topic", topic)
return
}
ns := parts[0]
deviceID := parts[2]
// Формируем envelope — оборачиваем payload в JSON с метаданными
// Payload от устройства может быть любым JSON или строкой
rawPayload := json.RawMessage(payload)
if !json.Valid(payload) {
// Если payload не JSON — упаковываем в строку
quotedBytes, _ := json.Marshal(string(payload))
rawPayload = json.RawMessage(quotedBytes)
}
envelope := iotTelemetryMessage{
Namespace: ns,
DeviceID: deviceID,
Topic: topic,
Payload: rawPayload,
ReceivedAt: time.Now().UTC().Format(time.RFC3339),
}
body, err := json.Marshal(envelope)
if err != nil {
log.Error("marshal telemetry message", "topic", topic, "err", err)
return
}
// Queue name: "iot.{namespace}.telemetry"
// Declare-on-publish: если queue не существует — создаём
queueName := fmt.Sprintf("iot.%s.telemetry", ns)
if _, err := rabbitCh.QueueDeclare(queueName, true, false, false, false, nil); err != nil {
log.Error("declare RabbitMQ queue", "queue", queueName, "err", err)
return
}
err = rabbitCh.Publish(
"", // exchange — default exchange
queueName, // routing key = queue name для default exchange
false, // mandatory
false, // immediate
amqp.Publishing{
ContentType: "application/json",
Body: body,
DeliveryMode: amqp.Persistent, // сохранять при рестарте RabbitMQ
},
)
if err != nil {
log.Error("publish to RabbitMQ", "queue", queueName, "err", err)
return
}
log.Info("forwarded IoT telemetry", "topic", topic, "namespace", ns, "device", deviceID, "queue", queueName)
}
}