diff --git a/.gitignore b/.gitignore index 43ac389..ed92214 100644 --- a/.gitignore +++ b/.gitignore @@ -15,6 +15,14 @@ testbin/* hack/local.env Dockerfile.cross +# IoT compiled binaries — не коммитим, только в Docker образ +mqtt-bridge +kafka-consumer +iot-mqtt-bridge +iot-kafka-consumer +manager +sless + # Test binary, build with `go test -c` *.test diff --git a/Dockerfile b/Dockerfile index 682f4a9..0a1def4 100644 --- a/Dockerfile +++ b/Dockerfile @@ -26,6 +26,9 @@ RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} go build -a -o ma # Запускается в iot-mqtt-bridge Deployment через command: ["/iot-mqtt-bridge"]. # Один образ, два entrypoint — практично для MVP: один CI pipeline, один registry repo. RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} go build -a -o iot-mqtt-bridge ./iot/cmd/mqtt-bridge/ +# iot-kafka-consumer — читает из Kafka топика iot.telemetry и пишет в IoT Postgres. +# Запускается отдельным Deployment-ом через command: ["/iot-kafka-consumer"]. +RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} go build -a -o iot-kafka-consumer ./iot/cmd/kafka-consumer/ FROM alpine:3.19 # ca-certificates нужны для TLS (S3 HTTPS, DockerHub) @@ -34,6 +37,8 @@ WORKDIR / COPY --from=builder /workspace/manager . # iot-mqtt-bridge — второй бинарь, запускается отдельным Deployment-ом. COPY --from=builder /workspace/iot-mqtt-bridge . +# iot-kafka-consumer — третий бинарь, Kafka→Postgres pipeline. +COPY --from=builder /workspace/iot-kafka-consumer . # migrations нужны при старте — оператор читает SQL файлы для инициализации БД COPY migrations/ migrations/ # Запускаем от непривилегированного пользователя diff --git a/config/crd/bases/iot.kube5s.ru_iotdevices.yaml b/config/crd/bases/iot.kube5s.ru_iotdevices.yaml new file mode 100644 index 0000000..45ad7e5 --- /dev/null +++ b/config/crd/bases/iot.kube5s.ru_iotdevices.yaml @@ -0,0 +1,123 @@ +--- +apiVersion: apiextensions.k8s.io/v1 +kind: CustomResourceDefinition +metadata: + annotations: + controller-gen.kubebuilder.io/version: v0.14.0 + name: iotdevices.iot.kube5s.ru +spec: + group: iot.kube5s.ru + names: + kind: IoTDevice + listKind: IoTDeviceList + plural: iotdevices + singular: iotdevice + scope: Namespaced + versions: + - additionalPrinterColumns: + - jsonPath: .spec.deviceId + name: DeviceID + type: string + - jsonPath: .status.phase + name: Phase + type: string + - jsonPath: .spec.enabled + name: Enabled + type: boolean + - jsonPath: .status.mqttUsername + name: MQTTUser + type: string + - jsonPath: .metadata.creationTimestamp + name: Age + type: date + name: v1alpha1 + schema: + openAPIV3Schema: + description: |- + IoTDevice — ресурс для регистрации IoT-устройства в платформе. + Контроллер автоматически создаёт k8s Secret с MQTT-credentials. + properties: + apiVersion: + description: |- + APIVersion defines the versioned schema of this representation of an object. + Servers should convert recognized schemas to the latest internal value, and + may reject unrecognized values. + More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#resources + type: string + kind: + description: |- + Kind is a string value representing the REST resource this object represents. + Servers may infer this from the endpoint the client submits requests to. + Cannot be updated. + In CamelCase. + More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#types-kinds + type: string + metadata: + type: object + spec: + description: IoTDeviceSpec — желаемое состояние IoT-устройства. + properties: + deviceId: + description: |- + DeviceID — уникальный идентификатор устройства внутри namespace. + Используется как часть MQTT username и имени Secret. + Разрешены только строчные буквы, цифры и дефис — для совместимости с k8s именами. + maxLength: 48 + pattern: ^[a-z0-9][a-z0-9-]*[a-z0-9]$ + type: string + enabled: + default: true + description: |- + Enabled — активно ли устройство (может подключаться к MQTT). + Если false — контроллер устанавливает phase=Disabled, EMQX auth отклоняет подключение. + Secret с credentials НЕ удаляется — при re-enable пароль остаётся прежним. + type: boolean + metadata: + additionalProperties: + type: string + description: |- + Metadata — произвольные метаданные устройства (модель, локация и т.д.). + Хранятся только в CRD, не влияют на логику контроллера. + type: object + required: + - deviceId + - enabled + type: object + status: + description: IoTDeviceStatus — наблюдаемое состояние IoT-устройства (заполняет + контроллер). + properties: + lastConnected: + description: |- + LastConnected — время последнего MQTT-подключения устройства. + Заполняется MQTT auth-сервисом при каждом успешном CONNECT. + format: date-time + type: string + message: + description: Message — человекочитаемое сообщение о текущем статусе + или ошибке. + type: string + mqttUsername: + description: |- + MQTTUsername — имя пользователя для подключения к MQTT-брокеру. + Формат: {namespace}_{deviceId} — глобально уникален в рамках EMQX. + type: string + phase: + description: 'Phase — текущее состояние: Active, Disabled, Pending, + Error.' + type: string + secretName: + description: SecretName — имя k8s Secret в том же namespace, содержащего + mqtt-username и mqtt-password. + type: string + topicPrefix: + description: |- + TopicPrefix — MQTT topic prefix, на который разрешена публикация. + Формат: {namespace}/ — устройство не может публиковать в чужие namespace. + type: string + type: object + type: object + served: true + storage: true + subresources: + status: {} diff --git a/config/crd/bases/sless.kube5s.ru_services.yaml b/config/crd/bases/sless.kube5s.ru_services.yaml index 1969d01..394f915 100644 --- a/config/crd/bases/sless.kube5s.ru_services.yaml +++ b/config/crd/bases/sless.kube5s.ru_services.yaml @@ -85,11 +85,11 @@ spec: description: S3Key — ключ объекта в S3 (путь до zip архива) type: string timeoutSec: - default: 30 description: |- - TimeoutSec — таймаут HTTP-прокси в секундах (default: 30). - Ограничивает время ожидания ответа от пода в invoke.go. - Для длительных вызовов (batch, pgstorm) увеличить до нужного значения. + TimeoutSec — таймаут HTTP-прокси в секундах. + 0 (по умолчанию) = без ограничения времени выполнения. + Задай > 0 чтобы принудительно обрывать медленные вызовы. + Диапазон: 1–900. 0 = нет таймаута. format: int32 type: integer required: diff --git a/config/rbac/role.yaml b/config/rbac/role.yaml index f6f10f9..a7c19ca 100644 --- a/config/rbac/role.yaml +++ b/config/rbac/role.yaml @@ -26,7 +26,10 @@ rules: - secrets verbs: - create + - delete - get + - list + - watch - apiGroups: - "" resources: @@ -75,6 +78,32 @@ rules: - patch - update - watch +- apiGroups: + - iot.kube5s.ru + resources: + - iotdevices + verbs: + - create + - delete + - get + - list + - patch + - update + - watch +- apiGroups: + - iot.kube5s.ru + resources: + - iotdevices/finalizers + verbs: + - update +- apiGroups: + - iot.kube5s.ru + resources: + - iotdevices/status + verbs: + - get + - patch + - update - apiGroups: - networking.k8s.io resources: @@ -139,6 +168,32 @@ rules: - get - patch - update +- apiGroups: + - sless.kube5s.ru + resources: + - services + verbs: + - create + - delete + - get + - list + - patch + - update + - watch +- apiGroups: + - sless.kube5s.ru + resources: + - services/finalizers + verbs: + - update +- apiGroups: + - sless.kube5s.ru + resources: + - services/status + verbs: + - get + - patch + - update - apiGroups: - sless.kube5s.ru resources: diff --git a/deployments/k8s/iot-kafka-consumer.yaml b/deployments/k8s/iot-kafka-consumer.yaml new file mode 100644 index 0000000..7cbc1df --- /dev/null +++ b/deployments/k8s/iot-kafka-consumer.yaml @@ -0,0 +1,48 @@ +# Создано: 2026-04-06 +# Deployment iot-kafka-consumer — читает IoT телеметрию из Kafka → пишет в IoT Postgres. +# +# Consumer group "iot-pg-consumer" — можно масштабировать горизонтально без дублирования. +# Offset коммитится ТОЛЬКО после успешной записи в Postgres (at-least-once гарантия). +# +# Применение: kubectl apply -f deployments/k8s/iot-kafka-consumer.yaml + +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: iot-kafka-consumer + namespace: sless + labels: + app: iot-kafka-consumer +spec: + replicas: 1 + selector: + matchLabels: + app: iot-kafka-consumer + template: + metadata: + labels: + app: iot-kafka-consumer + spec: + containers: + - name: kafka-consumer + # Тот же образ что и оператор — все IoT бинари в одном образе. + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.67 + imagePullPolicy: Always + command: ["/iot-kafka-consumer"] + env: + - name: KAFKA_BROKERS + value: "kafka.sless.svc.cluster.local:9092" + envFrom: + # IOT_PG_DSN — master DSN для IoT Postgres (per-tenant DB) + - secretRef: + name: iot-postgres-secret + resources: + requests: + memory: "32Mi" + cpu: "25m" + limits: + memory: "64Mi" + cpu: "100m" + imagePullSecrets: + - name: sless-registry-auth diff --git a/deployments/k8s/iot-mqtt-bridge.yaml b/deployments/k8s/iot-mqtt-bridge.yaml index 23c2e5f..5134fdb 100644 --- a/deployments/k8s/iot-mqtt-bridge.yaml +++ b/deployments/k8s/iot-mqtt-bridge.yaml @@ -1,9 +1,9 @@ # Создано: 2026-04-04 -# Изменено: 2026-04-05 (добавлен IOT_PG_DSN, версия v0.1.59) -# Deployment iot-mqtt-bridge — MQTT→RabbitMQ мост для IoT. +# Изменено: 2026-04-06 (MQTT→Kafka: убран RABBITMQ_URL, добавлен KAFKA_BROKERS, v0.1.67) +# Deployment iot-mqtt-bridge — MQTT→Kafka мост для IoT. # # Получает MQTT сообщения от EMQX (подписка на "+/telemetry/+") -# и публикует в RabbitMQ queue "iot.{namespace}.telemetry". +# и публикует в Kafka топик "iot.telemetry" (ключ = namespace). # # Credentials для MQTT подключения берутся из Secret iot-bridge-credentials. # Этот Secret нужно создать вручную ДО деплоя: @@ -45,18 +45,14 @@ spec: - name: mqtt-bridge # Тот же образ что и оператор — оба бинаря в одном слое (manager + iot-mqtt-bridge). # При смене версии оператора — менять тег и здесь. - image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.59 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.67 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 + - name: KAFKA_BROKERS + value: "kafka.sless.svc.cluster.local:9092" envFrom: - secretRef: name: iot-bridge-credentials diff --git a/deployments/k8s/kafka.yaml b/deployments/k8s/kafka.yaml new file mode 100644 index 0000000..3090aa9 --- /dev/null +++ b/deployments/k8s/kafka.yaml @@ -0,0 +1,150 @@ +# Изменено: 2026-04-06 +# kafka.yaml — минимальный деплой Apache Kafka в KRaft mode (без Zookeeper). +# Образ: apache/kafka (официальный, бесплатный). +# Используется для IoT telemetry pipeline: mqtt-bridge → Kafka → iot-kafka-consumer → Postgres. +# Для prod: заменить на managed Kafka (Confluent/Aiven) — только изменить KAFKA_BROKERS в Secret. +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: kafka-config + namespace: sless +data: + # server.properties для KRaft mode (без Zookeeper). + # Нода совмещает роли controller + broker. + server.properties: | + process.roles=broker,controller + node.id=1 + controller.quorum.voters=1@localhost:9093 + listeners=PLAINTEXT://:9092,CONTROLLER://:9093 + inter.broker.listener.name=PLAINTEXT + controller.listener.names=CONTROLLER + listener.security.protocol.map=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT + advertised.listeners=PLAINTEXT://kafka.sless.svc.cluster.local:9092 + log.dirs=/var/kafka-data/logs + num.partitions=1 + default.replication.factor=1 + offsets.topic.replication.factor=1 + transaction.state.log.replication.factor=1 + transaction.state.log.min.isr=1 + log.retention.hours=168 + log.retention.check.interval.ms=300000 + auto.create.topics.enable=true +--- +apiVersion: apps/v1 +kind: StatefulSet +metadata: + name: kafka + namespace: sless + labels: + app: kafka +spec: + serviceName: kafka-headless + replicas: 1 + selector: + matchLabels: + app: kafka + template: + metadata: + labels: + app: kafka + spec: + # apache/kafka образ запускается как UID 1000 (kafka user). + # fsGroup=1000 — позволяет писать в PVC смонтированный как root. + securityContext: + fsGroup: 1000 + initContainers: + # Форматирует хранилище KRaft если ещё не отформатировано. + # KAFKA_CLUSTER_ID должен быть уникальным UUID — генерируется один раз. + - name: kafka-init + image: apache/kafka:3.7.0 + command: + - /bin/sh + - -c + - | + if [ ! -f /var/kafka-data/logs/meta.properties ]; then + echo "Formatting Kafka storage..." + /opt/kafka/bin/kafka-storage.sh format \ + -t "$(cat /var/kafka-data/cluster.id 2>/dev/null || \ + /opt/kafka/bin/kafka-storage.sh random-uuid | tee /var/kafka-data/cluster.id)" \ + -c /tmp/kafka-config/server.properties + fi + volumeMounts: + - name: kafka-data + mountPath: /var/kafka-data + - name: kafka-config + mountPath: /tmp/kafka-config + containers: + - name: kafka + image: apache/kafka:3.7.0 + command: + - /opt/kafka/bin/kafka-server-start.sh + - /tmp/kafka-config/server.properties + ports: + - containerPort: 9092 + name: client + - containerPort: 9093 + name: controller + resources: + requests: + cpu: 100m + memory: 256Mi + limits: + cpu: 500m + memory: 512Mi + volumeMounts: + - name: kafka-data + mountPath: /var/kafka-data + - name: kafka-config + mountPath: /tmp/kafka-config + readinessProbe: + tcpSocket: + port: 9092 + initialDelaySeconds: 30 + periodSeconds: 10 + failureThreshold: 6 + volumes: + - name: kafka-config + configMap: + name: kafka-config + volumeClaimTemplates: + - metadata: + name: kafka-data + spec: + accessModes: [ReadWriteOnce] + storageClassName: vcd-disk-ext4 + resources: + requests: + storage: 1Gi +--- +apiVersion: v1 +kind: Service +metadata: + name: kafka + namespace: sless + labels: + app: kafka +spec: + ports: + - name: client + port: 9092 + targetPort: 9092 + selector: + app: kafka +--- +apiVersion: v1 +kind: Service +metadata: + name: kafka-headless + namespace: sless + labels: + app: kafka +spec: + clusterIP: None + ports: + - name: client + port: 9092 + - name: controller + port: 9093 + selector: + app: kafka diff --git a/doc/iot-mvp-plan.md b/doc/iot-mvp-plan.md index c912aa5..3c12526 100644 --- a/doc/iot-mvp-plan.md +++ b/doc/iot-mvp-plan.md @@ -1166,3 +1166,235 @@ kubectl rollout restart -n sless deploy/sless-operator deploy/iot-mqtt-bridge 8. progress.md: обновлять до и после каждого шага 9. Коммит + пуш после каждого завершённого шага 10. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой! + + +--- + +# ПЛАН: Telemetry Pipeline — Postgres -> REST API -> UI + +> **Автор плана**: GitHub Copilot (Claude Opus 4.6) +> **Дата**: 2026-04-05 +> **Исполнитель**: Claude Sonnet +> **Ветка**: iot-pg-telemetry +> **Предусловия**: все компоненты до этого этапа РЕАЛИЗОВАНЫ и задеплоены (CRD, controller, EMQX, mqtt-bridge, IoT Console UI) + +--- + +## Цель + +Полная цепочка: IoT устройство (или эмулятор в UI) -> MQTT -> INSERT в Postgres -> REST API -> отображение в таблице на вкладке Телеметрия в IoT Console. + +**User story**: юзер входит токеном, регистрирует устройство, запускает эмулятор (рандомные temp/humidity), переходит на вкладку Телеметрия и видит таблицу с данными: время | устройство | payload. + +--- + +## Что уже готово (НЕ ТРОГАТЬ без крайней необходимости) + +| Компонент | Файл(ы) | Статус | +|-----------|---------|--------| +| IoTDevice CRD + types | iot/api/v1alpha1/device_types.go | DONE | +| IoTDevice controller | iot/controllers/iotdevice_controller.go | DONE | +| MQTT Auth + ACL | internal/api/handler/iot_device_handler.go | DONE | +| IoT API CRUD | internal/api/router.go + handler | DONE | +| EMQX deploy | deployments/k8s/emqx.yaml | DONE | +| mqtt-bridge MQTT->RabbitMQ | iot/cmd/mqtt-bridge/main.go | DONE | +| IoT Console UI | internal/api/ui/iot-console.html | DONE | +| TLS (HTTPS + WSS) | deployments/k8s/emqx-ws-ingress.yaml | DONE | +| Nubes branding | UI CSS | DONE | +| Existing Postgres (invocations) | deployments/k8s/postgres.yaml | DONE | + +--- + +## Архитектурное решение (принято 2026-04-04) + +Подробности: doc/decisions/iot-telemetry-storage-2026-04-04.md + +- **Один Postgres инстанс** для IoT (отдельный от sless Postgres для invocations) +- **Отдельная DATABASE per tenant** (не одна таблица с tenant_id!) +- Tenant DB: tenant_{namespace_hash}, User: tenant_{namespace_hash}, Password: UUID в k8s Secret +- Таблица: iot_telemetry(id BIGSERIAL, device_id TEXT, ts TIMESTAMPTZ, payload JSONB) +- Клиент читает ТОЛЬКО через REST API, не через прямой доступ к Postgres + +--- + +## Шаги реализации (порядок критичен!) + +### ШАГ 1: Postgres Deployment для IoT (namespace: sless) + +**Файл**: deployments/k8s/iot-postgres.yaml + +**Почему отдельный от sless postgres**: разные данные, разная нагрузка. +**Почему в namespace sless, а НЕ iot**: всё живёт в одном namespace, упрощение. + +**YAML манифест**: + + + +**Действие**: kubectl apply -f deployments/k8s/iot-postgres.yaml +**Проверка**: kubectl exec -n sless deploy/iot-postgres -- psql -U iot_admin -d iot_platform -c "SELECT 1" + +--- + +### ШАГ 2: Go-пакет IoT Postgres storage + +**Файл**: internal/storage/iotpg/iot_telemetry_store.go + +**Структура**: + +{ is a shell keyword + +**Методы (все обязательные)**: + +1. New(adminDSN string, log) (*IoTPostgresStore, error) -- подключение к iot_platform DB +2. EnsureTenantDB(ctx, namespace) error -- создать DATABASE + USER + таблицу если не существуют: + - SELECT 1 FROM pg_database WHERE datname = tenant_{ns} + - Если нет: CREATE USER, CREATE DATABASE, подключиться и CREATE TABLE + - Сохранить пароль в tenant_credentials таблице в iot_platform + - Таблица: iot_telemetry(id BIGSERIAL PK, device_id TEXT, ts TIMESTAMPTZ DEFAULT now(), payload JSONB) + - Индекс: idx_iot_telemetry_device_ts ON iot_telemetry(device_id, ts DESC) +3. InsertTelemetry(ctx, namespace, deviceID, payload json.RawMessage) error +4. QueryTelemetry(ctx, namespace, deviceID string, limit int) ([]TelemetryRow, error) +5. Close() error + +**Tenant DB provisioning**: таблица tenant_credentials в iot_platform: + + +**Кэширование**: sync.Map для *sql.DB per tenant. Lazy init при первом обращении. + +--- + +### ШАГ 3: Модифицировать mqtt-bridge -- добавить INSERT в Postgres + +**Файл**: iot/cmd/mqtt-bridge/main.go + +**Текущее поведение**: MQTT message -> envelope -> RabbitMQ. +**Новое поведение**: MQTT message -> INSERT в Postgres (tenant DB) + RabbitMQ (как было). + +**Изменения**: +1. Добавить env var IOT_PG_DSN +2. Подключиться к IoTPostgresStore при старте +3. В buildMQTTMessageHandler: + - store.EnsureTenantDB(ctx, namespace) -- идемпотентно + - store.InsertTelemetry(ctx, namespace, deviceID, payload) + - При ошибке INSERT -- логировать, НЕ блокировать RabbitMQ publish +4. RabbitMQ publish остаётся как было + +**YAML**: deployments/k8s/iot-mqtt-bridge.yaml -- добавить env IOT_PG_DSN из iot-postgres-secret + +--- + +### ШАГ 4: REST API endpoint для чтения телеметрии + +**Файл**: internal/api/handler/iot_telemetry_handler.go (НОВЫЙ) + +**Endpoint**: + + +**Параметры**: +- device -- фильтр по device_id (опционален) +- limit -- максимум записей (default: 50, max: 1000) + +**Response**: + + +**Сортировка**: ts DESC (новые сверху). + +--- + +### ШАГ 5: Инициализация IoTPostgresStore в main.go + +**Файл**: main.go + +1. Добавить поле IoTPG в handler.Handler struct (handler.go) +2. В main.go: if IOT_PG_DSN задан -> iotpg.New() -> передать в Handler +3. В router.go: зарегистрировать route /namespaces/{ns}/iot/telemetry + +**YAML**: deployments/k8s/operator.yaml -- добавить env IOT_PG_DSN + +--- + +### ШАГ 6: Обновить IoT Console UI -- вкладка Телеметрия + +**Файл**: internal/api/ui/iot-console.html + +**Заменить** заглушку coming-soon на реальную таблицу: + +| Время | Устройство | Данные | +|-------|-----------|--------| +| 2026-04-05 08:15 | sensor-01 | {"temperature": 22.5, "humidity": 65} | + +**JavaScript**: +- loadTelemetry() -- fetch GET /v1/.../iot/telemetry -> заполнить tbody +- Авто-обновление каждые 5с (чекбокс) +- Фильтр по устройству (select из списка devices) +- При переключении на вкладку -- автоматический loadTelemetry() + +**CSS**: таблица в стиле Nubes (navy фон, бордеры #0b2d50, текст #e2ecf6) + +--- + +### ШАГ 7: Улучшить эмулятор -- рандомные temp/humidity + +**Файл**: internal/api/ui/iot-console.html (секция эмулятора) + +**Новое поведение**: +- Чекбокс: "Генерировать случайные данные (temp/humidity)" (по умолчанию ON) +- Если ON: при каждой отправке payload = {temperature: random(18-28), humidity: random(40-80), ts: ISO} +- Если OFF: используется текстовое поле как сейчас + + + +--- + +## Деплой + +deployment.apps/iot-postgres condition met +deployment.apps/sless-operator restarted +deployment.apps/iot-mqtt-bridge restarted + +--- + +## Файлы СОЗДАТЬ + +| Файл | Описание | +|------|----------| +| deployments/k8s/iot-postgres.yaml | Deployment + Secret + Service | +| internal/storage/iotpg/iot_telemetry_store.go | Go: управление tenant DB + CRUD телеметрии | +| internal/api/handler/iot_telemetry_handler.go | REST handler GET /v1/.../iot/telemetry | + +## Файлы ИЗМЕНИТЬ + +| Файл | Что менять | +|------|-----------| +| internal/api/handler/handler.go | Добавить поле IoTPG *iotpg.IoTPostgresStore | +| internal/api/router.go | Route /namespaces/{ns}/iot/telemetry | +| main.go | Init IoTPostgresStore + передача в Handler | +| iot/cmd/mqtt-bridge/main.go | INSERT в Postgres при MQTT message | +| deployments/k8s/operator.yaml | env IOT_PG_DSN + версия v0.1.59 | +| deployments/k8s/iot-mqtt-bridge.yaml | env IOT_PG_DSN + версия v0.1.59 | +| internal/api/ui/iot-console.html | Telemetry tab + emulator random data | + +## Чего НЕ ДЕЛАТЬ + +- НЕ трогать CRD / controller / EMQX / RabbitMQ +- НЕ создавать namespace iot -- всё в sless +- НЕ делать processing данных -- RAW payload +- НЕ добавлять from/to фильтры -- хватит limit +- НЕ трогать Terraform provider +- НЕ рефакторить существующие файлы +- НЕ запускать команды локально -- только SSH + +--- + +## Правила для Sonnet + +1. Читай .github/copilot-instructions.md +2. Читай doc/decisions/iot-telemetry-storage-2026-04-04.md +3. Команды через SSH: ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213 +4. Файлы редактировать можно -- sshfs mount +5. ПЕРЕД go build -- проверить .gitignore +6. Комментарии: дата + назначение + почему +7. Thinking log: doc/thinking/2026-04-05.md +8. progress.md: обновлять до и после каждого шага +9. Коммит + пуш после каждого завершённого шага +10. ВЕРСИЮ ПОДНИМАТЬ перед каждой сборкой! diff --git a/doc/migration/cluster-migration-simple.md b/doc/migration/cluster-migration-simple.md new file mode 100644 index 0000000..4d1683c --- /dev/null +++ b/doc/migration/cluster-migration-simple.md @@ -0,0 +1,45 @@ +# Миграция на новый кластер (тестовые данные) + +> Дата: 2026-04-06 +> Сценарий: Все данные тестовые и неважны + +--- + +## Процесс + +1. **Clone + Build** + ```bash + git clone + make docker-build docker-push IMG=/sless:v1.0 + ``` + +2. **Deploy** + ```bash + kubectl create namespace sless + + # Создать Secrets (новые credentials) + kubectl create secret generic sless-operator-secret -n sless \ + --from-literal=POSTGRES_DSN="..." \ + --from-literal=S3_ACCESS_KEY="..." \ + --from-literal=S3_SECRET_KEY="..." \ + --from-literal=SLESS_API_TOKEN="..." \ + --from-literal=HARBOR_PASS="..." + + # Apply конфиги + kubectl apply -f deployments/k8s/rbac.yaml + kubectl apply -f deployments/k8s/ + ``` + +3. **Done** + - БД создадутся новые и пустые + - Registry пересоберётся + - Готово + +--- + +## Что не требуется +- ❌ pg_dump / восстановление БД +- ❌ Копирование PVC +- ❌ Миграция данных + +Всё пересоздаётся с нуля. diff --git a/go.mod b/go.mod index 1defc0b..6d50959 100644 --- a/go.mod +++ b/go.mod @@ -51,6 +51,7 @@ require ( github.com/modern-go/reflect2 v1.0.2 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/philhofer/fwd v1.2.0 // indirect + github.com/pierrec/lz4/v4 v4.1.15 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.14.0 // indirect github.com/prometheus/client_model v0.3.0 // indirect @@ -58,6 +59,7 @@ require ( github.com/prometheus/procfs v0.8.0 // indirect github.com/rabbitmq/amqp091-go v1.10.0 // indirect github.com/rs/xid v1.6.0 // indirect + github.com/segmentio/kafka-go v0.4.50 // indirect github.com/spf13/pflag v1.0.5 // indirect github.com/tinylib/msgp v1.6.1 // indirect go.uber.org/atomic v1.7.0 // indirect diff --git a/go.sum b/go.sum index 47ab079..b517b42 100644 --- a/go.sum +++ b/go.sum @@ -242,6 +242,8 @@ github.com/onsi/gomega v1.24.1 h1:KORJXNNTzJXzu4ScJWssJfJMnJ+2QJqhoQSRwNlze9E= github.com/onsi/gomega v1.24.1/go.mod h1:3AOiACssS3/MajrniINInwbfOOtfZvplPzuRSmvt1jM= github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM= +github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -279,6 +281,8 @@ github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMu github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= +github.com/segmentio/kafka-go v0.4.50 h1:mcyC3tT5WeyWzrFbd6O374t+hmcu1NKt2Pu1L3QaXmc= +github.com/segmentio/kafka-go v0.4.50/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E= github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88= diff --git a/iot/cmd/kafka-consumer/main.go b/iot/cmd/kafka-consumer/main.go new file mode 100644 index 0000000..00cd5e6 --- /dev/null +++ b/iot/cmd/kafka-consumer/main.go @@ -0,0 +1,150 @@ +// Создано: 2026-04-06 +// kafka-consumer/main.go — iot-kafka-consumer: читает IoT телеметрию из Kafka → пишет в Postgres. +// +// Роль в архитектуре: +// Kafka топик "iot.telemetry" → iot-kafka-consumer → IoT Postgres (per-tenant DB) +// +// Consumer group "iot-pg-consumer" — позволяет запускать несколько реплик без дублирования. +// При временной недоступности Postgres — Kafka хранит сообщения (retention 7 дней). +// +// Конфигурация через env vars: +// KAFKA_BROKERS — kafka.sless.svc.cluster.local:9092 (или managed Kafka в prod) +// IOT_PG_DSN — postgres://user:pass@host:5432/iotdb (master DSN для IoT Postgres) + +package main + +import ( + "context" + "encoding/json" + "log/slog" + "os" + "os/signal" + "strings" + "syscall" + + kafka "github.com/segmentio/kafka-go" + + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/iotpg" +) + +// kafkaConsumerConfig — конфигурация из env vars. +type kafkaConsumerConfig struct { + KafkaBrokers string +} + +// iotTelemetryMessage — envelope из Kafka (идентичен bridge). +type iotTelemetryMessage struct { + Namespace string `json:"namespace"` + DeviceID string `json:"device_id"` + Topic string `json:"topic"` + Payload json.RawMessage `json:"payload"` + ReceivedAt string `json:"received_at"` +} + +// iotTelemetryTopic — Kafka топик (должен совпадать с bridge). +const iotTelemetryTopic = "iot.telemetry" + +// iotConsumerGroup — идентификатор consumer group. +// При нескольких репликах Kafka распределяет партиции между ними. +const iotConsumerGroup = "iot-pg-consumer" + +func main() { + log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) + + cfg := loadConsumerConfig() + + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT) + defer cancel() + + log.Info("starting iot-kafka-consumer", + "kafka_brokers", cfg.KafkaBrokers, + "topic", iotTelemetryTopic, + "group", iotConsumerGroup, + ) + + // IoT Postgres — обязательный компонент для этого сервиса + iotStore, err := iotpg.NewFromEnv(log) + if err != nil || iotStore == nil { + log.Error("failed to connect to IoT Postgres — IOT_PG_DSN required", "err", err) + os.Exit(1) + } + defer iotStore.Close() + log.Info("connected to IoT Postgres") + + // Kafka reader с consumer group — автоматически коммитит offsets после обработки + reader := kafka.NewReader(kafka.ReaderConfig{ + Brokers: strings.Split(cfg.KafkaBrokers, ","), + Topic: iotTelemetryTopic, + GroupID: iotConsumerGroup, + MinBytes: 1, + MaxBytes: 1 << 20, // 1MB + }) + defer reader.Close() + + log.Info("kafka reader ready, waiting for messages...") + + for { + // FetchMessage — блокирует до следующего сообщения + kafkaMsg, err := reader.FetchMessage(ctx) + if err != nil { + if ctx.Err() != nil { + break // штатное завершение + } + log.Error("fetch from Kafka", "err", err) + continue + } + + if err := processKafkaTelemetry(ctx, kafkaMsg, iotStore, log); err != nil { + log.Error("process telemetry message", "err", err) + // НЕ коммитим offset — сообщение будет перечитано при следующем старте + continue + } + + // Коммитим offset только после успешной записи в Postgres + if err := reader.CommitMessages(ctx, kafkaMsg); err != nil { + log.Error("commit Kafka offset", "err", err) + } + } + + log.Info("shutting down iot-kafka-consumer") +} + +// processKafkaTelemetry десериализует сообщение из Kafka и записывает в Postgres. +func processKafkaTelemetry(ctx context.Context, msg kafka.Message, store *iotpg.IoTPostgresStore, log *slog.Logger) error { + var envelope iotTelemetryMessage + if err := json.Unmarshal(msg.Value, &envelope); err != nil { + // Битое сообщение — логируем и пропускаем (не блокируем очередь) + log.Warn("failed to unmarshal telemetry envelope, skipping", "err", err, "raw", string(msg.Value)) + return nil + } + + // EnsureTenantDB идемпотентен — кэшируется после первого вызова + if err := store.EnsureTenantDB(ctx, envelope.Namespace); err != nil { + return err + } + + if err := store.InsertTelemetry(ctx, envelope.Namespace, envelope.DeviceID, envelope.Payload); err != nil { + return err + } + + log.Info("telemetry saved to Postgres", + "namespace", envelope.Namespace, + "device", envelope.DeviceID, + "kafka_offset", msg.Offset, + ) + return nil +} + +// loadConsumerConfig читает конфигурацию из env vars. +func loadConsumerConfig() kafkaConsumerConfig { + return kafkaConsumerConfig{ + KafkaBrokers: getEnvOrDefault("KAFKA_BROKERS", "kafka.sless.svc.cluster.local:9092"), + } +} + +func getEnvOrDefault(key, defaultVal string) string { + if v := os.Getenv(key); v != "" { + return v + } + return defaultVal +} diff --git a/iot/cmd/mqtt-bridge/main.go b/iot/cmd/mqtt-bridge/main.go index 29420db..68f68d1 100644 --- a/iot/cmd/mqtt-bridge/main.go +++ b/iot/cmd/mqtt-bridge/main.go @@ -1,29 +1,26 @@ // Создано: 2026-04-04 -// Изменено: 2026-04-05 (добавлен INSERT в IoT Postgres) -// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → RabbitMQ. +// Изменено: 2026-04-06 (заменён RabbitMQ на Kafka) +// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → Kafka. // // Роль в архитектуре: -// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] → RabbitMQ → event-dispatcher → function +// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] +// → Kafka топик "iot.telemetry" +// → iot-kafka-consumer → Postgres (история телеметрии) +// → event-dispatcher → Serverless Functions (триггеры) // // Логика: // 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) +// 3. При получении сообщения — опубликовать в Kafka топик "iot.telemetry" +// 4. Payload оборачивается в envelope с метаданными (namespace, device_id, ts) // // Конфигурация через 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/ +// KAFKA_BROKERS — kafka.sless.svc.cluster.local:9092 (заменить на managed в prod) // -// ВАЖНО: bridge клиент должен проходить EMQX auth — нужен IoTDevice "iot-bridge" в namespace "sless-bridge". -// Для MVP: выделить специальный namespace "sless-bridge" с устройством "bridge", -// и использовать его credentials для подключения bridge сервиса. -// Или: зарегистрировать bridge устройство через API и записать credentials в Secret. +// Для возврата к Postgres напрямую: см. git история, коммиты до 2026-04-06. package main @@ -39,9 +36,7 @@ import ( "time" mqtt "github.com/eclipse/paho.mqtt.golang" - amqp "github.com/rabbitmq/amqp091-go" - - "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/iotpg" + kafka "github.com/segmentio/kafka-go" ) // mqttBridgeConfig — конфигурация сервиса из env vars. @@ -49,24 +44,28 @@ type mqttBridgeConfig struct { MQTTBrokerURL string MQTTUsername string MQTTPassword string - RabbitMQURL string + KafkaBrokers string } -// iotTelemetryMessage — структура сообщения публикуемого в RabbitMQ. -// Оборачивает MQTT payload в envelope с метаданными. +// iotTelemetryMessage — envelope сообщения публикуемого в Kafka. +// Потребители (iot-kafka-consumer, event-dispatcher) читают этот формат. type iotTelemetryMessage struct { - // Namespace — k8s namespace пользователя (из MQTT topic) + // Namespace — k8s namespace тенанта (из MQTT topic, первый сегмент) Namespace string `json:"namespace"` - // DeviceID — идентификатор устройства (из MQTT topic, последний сегмент) + // DeviceID — идентификатор устройства (из MQTT topic, третий сегмент) DeviceID string `json:"device_id"` // Topic — оригинальный MQTT topic Topic string `json:"topic"` - // Payload — данные от устройства (JSON передаётся as-is / строка если не JSON) + // Payload — данные от устройства (JSON as-is, или строка если не JSON) Payload json.RawMessage `json:"payload"` - // ReceivedAt — время получения сообщения мостом (UTC) + // ReceivedAt — время получения сообщения мостом (UTC, RFC3339) ReceivedAt string `json:"received_at"` } +// iotTelemetryTopic — Kafka топик для IoT телеметрии. +// Все устройства всех тенантов пишут в один топик, изоляция — по полю Namespace в payload. +const iotTelemetryTopic = "iot.telemetry" + func main() { log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) @@ -78,34 +77,19 @@ func main() { log.Info("starting iot-mqtt-bridge", "mqtt_broker", cfg.MQTTBrokerURL, "mqtt_username", cfg.MQTTUsername, + "kafka_brokers", cfg.KafkaBrokers, ) - // 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() - - // IoT Postgres — сохранение телеметрии (per-tenant DB). - // Опционально: если IOT_PG_DSN не задан — продолжаем работать без Postgres (только RabbitMQ) - iotPGStore, err := iotpg.NewFromEnv(log) - if err != nil { - log.Error("failed to connect to IoT Postgres", "err", err) - os.Exit(1) - } - if iotPGStore != nil { - defer iotPGStore.Close() - log.Info("connected to IoT Postgres for telemetry storage") + // Kafka writer — асинхронный, с автоматическим созданием топика. + kafkaWriter := &kafka.Writer{ + Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...), + Topic: iotTelemetryTopic, + Balancer: &kafka.LeastBytes{}, + // Позволяет продолжать работу при временной недоступности Kafka (буфер в памяти) + Async: false, + RequiredAcks: kafka.RequireOne, } + defer kafkaWriter.Close() // Создаём MQTT клиент mqttClient, err := connectMQTT(cfg, log) @@ -116,8 +100,7 @@ func main() { defer mqttClient.Disconnect(500) // Функция-обработчик MQTT сообщений - // Вызывается в goroutine paho при каждом сообщении - messageHandler := buildMQTTMessageHandler(ctx, rabbitCh, iotPGStore, log) + messageHandler := buildMQTTMessageHandler(ctx, kafkaWriter, log) // Подписываемся на все telemetry топики всех namespace // "+/telemetry/+" = {любой namespace}/telemetry/{любой deviceId} @@ -135,7 +118,6 @@ func main() { } // loadBridgeConfig читает конфигурацию из env vars. -// Завершает процесс если обязательные переменные отсутствуют. func loadBridgeConfig() mqttBridgeConfig { required := func(key string) string { v := os.Getenv(key) @@ -150,7 +132,7 @@ func loadBridgeConfig() mqttBridgeConfig { MQTTBrokerURL: getEnvOrDefault("MQTT_BROKER_URL", "tcp://emqx.sless.svc:1883"), MQTTUsername: required("MQTT_USERNAME"), MQTTPassword: required("MQTT_PASSWORD"), - RabbitMQURL: required("RABBITMQ_URL"), + KafkaBrokers: getEnvOrDefault("KAFKA_BROKERS", "kafka.sless.svc.cluster.local:9092"), } } @@ -162,7 +144,6 @@ func getEnvOrDefault(key, defaultVal string) string { } // connectMQTT устанавливает подключение к EMQX брокеру. -// AutoReconnect=true — paho сам переподключается при разрыве. func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) { opts := mqtt.NewClientOptions() opts.AddBroker(cfg.MQTTBrokerURL) @@ -173,7 +154,7 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) { opts.SetConnectRetry(true) opts.SetConnectRetryInterval(5 * time.Second) opts.SetKeepAlive(30 * time.Second) - opts.SetCleanSession(false) // сохраняем подписки при реконнекте + opts.SetCleanSession(false) opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) { log.Warn("MQTT connection lost, reconnecting...", "err", err) @@ -187,7 +168,6 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) { client := mqtt.NewClient(opts) token := client.Connect() - // Ждём максимум 30 секунд if !token.WaitTimeout(30 * time.Second) { return nil, fmt.Errorf("MQTT connect timeout") } @@ -197,40 +177,15 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, 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), iotStore (может быть nil) и logger. -// Порядок действий при получении сообщения: -// 1. INSERT в IoT Postgres (tenant DB) — если iotStore != nil -// 2. Publish в RabbitMQ — всегда (для event-dispatcher → function triggers) -// -// Ошибка INSERT не блокирует RabbitMQ publish — разные failure domain. -func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotStore *iotpg.IoTPostgresStore, log *slog.Logger) mqtt.MessageHandler { +// buildMQTTMessageHandler возвращает обработчик MQTT сообщений. +// При получении сообщения — публикует envelope в Kafka топик "iot.telemetry". +// Ключ сообщения Kafka = namespace, для партиционирования по тенанту. +func buildMQTTMessageHandler(ctx context.Context, w *kafka.Writer, 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) @@ -239,28 +194,13 @@ func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotSto ns := parts[0] deviceID := parts[2] - // Нормализуем payload: если это не JSON — оборачиваем в строку + // Нормализуем payload: если не JSON — оборачиваем в строку rawPayload := json.RawMessage(payload) if !json.Valid(payload) { quotedBytes, _ := json.Marshal(string(payload)) rawPayload = json.RawMessage(quotedBytes) } - // ШАГ 1: INSERT в IoT Postgres — сохраняем телеметрию в per-tenant DB - // EnsureTenantDB идемпотентен: кэшируется после первого вызова - if iotStore != nil { - if err := iotStore.EnsureTenantDB(ctx, ns); err != nil { - log.Error("ensure tenant DB", "namespace", ns, "err", err) - // НЕ возвращаемся — продолжаем RabbitMQ publish - } else if err := iotStore.InsertTelemetry(ctx, ns, deviceID, rawPayload); err != nil { - log.Error("insert telemetry", "topic", topic, "err", err) - // НЕ возвращаемся — RabbitMQ не должен зависеть от Postgres - } else { - log.Debug("telemetry saved to Postgres", "namespace", ns, "device", deviceID) - } - } - - // ШАГ 2: Publish в RabbitMQ (для event-dispatcher → function triggers) envelope := iotTelemetryMessage{ Namespace: ns, DeviceID: deviceID, @@ -275,30 +215,22 @@ func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotSto 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 - }, - ) + // Ключ = namespace — Kafka будет группировать сообщения одного тенанта + // на одну партицию (для упорядоченной обработки на consumer side) + err = w.WriteMessages(ctx, kafka.Message{ + Key: []byte(ns), + Value: body, + }) if err != nil { - log.Error("publish to RabbitMQ", "queue", queueName, "err", err) + log.Error("publish to Kafka", "topic", iotTelemetryTopic, "namespace", ns, "err", err) return } - log.Info("forwarded IoT telemetry", "topic", topic, "namespace", ns, "device", deviceID, "queue", queueName) + log.Info("forwarded IoT telemetry to Kafka", + "mqtt_topic", topic, + "namespace", ns, + "device", deviceID, + "kafka_topic", iotTelemetryTopic, + ) } }