diff --git a/deployments/k8s/emqx.yaml b/deployments/k8s/emqx.yaml new file mode 100644 index 0000000..84f583c --- /dev/null +++ b/deployments/k8s/emqx.yaml @@ -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 diff --git a/deployments/k8s/iot-mqtt-bridge.yaml b/deployments/k8s/iot-mqtt-bridge.yaml new file mode 100644 index 0000000..e9ace7d --- /dev/null +++ b/deployments/k8s/iot-mqtt-bridge.yaml @@ -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 diff --git a/doc/thinking/2026-04-04.md b/doc/thinking/2026-04-04.md index 2fb6fc7..facfb43 100644 --- a/doc/thinking/2026-04-04.md +++ b/doc/thinking/2026-04-04.md @@ -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 константами diff --git a/examples/IOT/README.md b/examples/IOT/README.md new file mode 100644 index 0000000..ebd4ec3 --- /dev/null +++ b/examples/IOT/README.md @@ -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) diff --git a/examples/IOT/handler.py b/examples/IOT/handler.py new file mode 100644 index 0000000..f50b0a1 --- /dev/null +++ b/examples/IOT/handler.py @@ -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, + }) + } diff --git a/examples/IOT/main.tf b/examples/IOT/main.tf new file mode 100644 index 0000000..bb1386f --- /dev/null +++ b/examples/IOT/main.tf @@ -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)" +} diff --git a/examples/PG_TEST/.gitignore b/examples/PG_TEST/.gitignore new file mode 100644 index 0000000..d776a24 --- /dev/null +++ b/examples/PG_TEST/.gitignore @@ -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 diff --git a/examples/PG_TEST/postgres.tf.bak_db b/examples/PG_TEST/postgres.tf.bak_db deleted file mode 100644 index 3c9cb31..0000000 --- a/examples/PG_TEST/postgres.tf.bak_db +++ /dev/null @@ -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 -} diff --git a/examples/PG_TEST/postgres_extra.tf11 b/examples/PG_TEST/postgres_extra.tf11 deleted file mode 100644 index 6504fb4..0000000 --- a/examples/PG_TEST/postgres_extra.tf11 +++ /dev/null @@ -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] -} diff --git a/go.mod b/go.mod index 425ba09..1defc0b 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 5e3352e..47ab079 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/api/handler/iot_device_handler.go b/internal/api/handler/iot_device_handler.go new file mode 100644 index 0000000..8ff107f --- /dev/null +++ b/internal/api/handler/iot_device_handler.go @@ -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, "")) +} diff --git a/internal/api/router.go b/internal/api/router.go index d7af48d..6fe73aa 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -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 суброутеру. diff --git a/iot/cmd/mqtt-bridge/main.go b/iot/cmd/mqtt-bridge/main.go new file mode 100644 index 0000000..392e4bf --- /dev/null +++ b/iot/cmd/mqtt-bridge/main.go @@ -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) + } +}