Компоненты: - iot-operator: controller-manager (IoTDevice CRD) + REST API (порт 9090) - mqtt-bridge: MQTT (EMQX) → Kafka bridge - kafka-consumer: Kafka → Postgres pipeline Модуль: gitea.services.ngcloud.ru/Nail/IoT Все 3 бинарника собираются, import paths адаптированы.
399 lines
14 KiB
Go
399 lines
14 KiB
Go
// Создано: 2026-04-05
|
||
// iot_telemetry_store.go — управление per-tenant PostgreSQL databases для IoT телеметрии.
|
||
//
|
||
// Архитектура (принято 2026-04-04, см. doc/decisions/iot-telemetry-storage-2026-04-04.md):
|
||
// - Один Postgres инстанс (iot-postgres.sless.svc) — отдельный от sless postgres
|
||
// - Отдельная DATABASE per tenant: tenant_{namespace} (дефисы → подчёркивания)
|
||
// - Suперюзер iot_admin управляет всеми DBs; клиенты читают только через REST API
|
||
// - Пароли tenant хранятся в таблице tenant_credentials в management DB iot_platform
|
||
//
|
||
// Почему tenant_credentials в БД, а не в k8s Secret:
|
||
// mqtt-bridge вызывает InsertTelemetry в горячем пути MQTT.
|
||
// k8s API round-trip на каждое сообщение — неприемлемо.
|
||
//
|
||
// Подключение к tenant DB: суперюзер iot_admin, DSN строится заменой db name в adminDSN.
|
||
// Кэширование: sync.Map для *sql.DB per tenant (lazy init при первом обращении).
|
||
|
||
package iotpg
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"log/slog"
|
||
"os"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/google/uuid"
|
||
"github.com/lib/pq"
|
||
)
|
||
|
||
// IoTPostgresStore управляет per-tenant Postgres databases для IoT телеметрии.
|
||
type IoTPostgresStore struct {
|
||
adminDB *sql.DB
|
||
adminDSN string
|
||
tenants sync.Map
|
||
log *slog.Logger
|
||
}
|
||
|
||
// TelemetryRow — одна запись телеметрии из таблицы iot_telemetry.
|
||
type TelemetryRow struct {
|
||
ID int64 `json:"id"`
|
||
DeviceID string `json:"device_id"`
|
||
Ts time.Time `json:"ts"`
|
||
Payload json.RawMessage `json:"payload"`
|
||
}
|
||
|
||
// New подключается к management DB (iot_platform) и создаёт служебные таблицы.
|
||
func New(adminDSN string, log *slog.Logger) (*IoTPostgresStore, error) {
|
||
db, err := sql.Open("postgres", adminDSN)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("iotpg: open admin DB: %w", err)
|
||
}
|
||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||
defer cancel()
|
||
if err := db.PingContext(ctx); err != nil {
|
||
db.Close()
|
||
return nil, fmt.Errorf("iotpg: ping admin DB: %w", err)
|
||
}
|
||
db.SetMaxOpenConns(5)
|
||
db.SetMaxIdleConns(2)
|
||
db.SetConnMaxLifetime(5 * time.Minute)
|
||
|
||
store := &IoTPostgresStore{adminDB: db, adminDSN: adminDSN, log: log}
|
||
if err := store.initManagementSchema(ctx); err != nil {
|
||
db.Close()
|
||
return nil, fmt.Errorf("iotpg: init management schema: %w", err)
|
||
}
|
||
log.Info("iotpg: connected to IoT Postgres management DB")
|
||
return store, nil
|
||
}
|
||
|
||
// NewFromEnv создаёт store из env var IOT_PG_DSN.
|
||
// Возвращает (nil, nil) если переменная не задана — IoT Postgres опционален.
|
||
func NewFromEnv(log *slog.Logger) (*IoTPostgresStore, error) {
|
||
dsn := os.Getenv("IOT_PG_DSN")
|
||
if dsn == "" {
|
||
log.Info("iotpg: IOT_PG_DSN not set, IoT telemetry disabled")
|
||
return nil, nil
|
||
}
|
||
return New(dsn, log)
|
||
}
|
||
|
||
// initManagementSchema создаёт таблицу tenant_credentials в iot_platform.
|
||
func (s *IoTPostgresStore) initManagementSchema(ctx context.Context) error {
|
||
_, err := s.adminDB.ExecContext(ctx, `
|
||
CREATE TABLE IF NOT EXISTS tenant_credentials (
|
||
namespace TEXT PRIMARY KEY,
|
||
pg_password TEXT NOT NULL,
|
||
created_at TIMESTAMPTZ DEFAULT now()
|
||
)
|
||
`)
|
||
return err
|
||
}
|
||
|
||
// EnsureTenantDB создаёт DATABASE, USER и таблицу iot_telemetry для namespace.
|
||
// Идемпотентен — повторный вызов безопасен.
|
||
// Вызывается mqtt-bridge при первом сообщении от нового tenant.
|
||
func (s *IoTPostgresStore) EnsureTenantDB(ctx context.Context, namespace string) error {
|
||
dbName := tenantDBName(namespace)
|
||
userName := dbName
|
||
|
||
var exists bool
|
||
err := s.adminDB.QueryRowContext(ctx,
|
||
`SELECT EXISTS(SELECT 1 FROM pg_database WHERE datname = $1)`, dbName,
|
||
).Scan(&exists)
|
||
if err != nil {
|
||
return fmt.Errorf("iotpg: check tenant DB %s: %w", dbName, err)
|
||
}
|
||
|
||
if !exists {
|
||
password := uuid.New().String()
|
||
|
||
// CREATE USER через DO block — pg не поддерживает CREATE USER IF NOT EXISTS
|
||
_, err = s.adminDB.ExecContext(ctx, fmt.Sprintf(
|
||
`DO $$ BEGIN
|
||
IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = '%s') THEN
|
||
CREATE USER %s WITH PASSWORD '%s';
|
||
END IF;
|
||
END $$`, userName, userName, password,
|
||
))
|
||
if err != nil {
|
||
return fmt.Errorf("iotpg: create user %s: %w", userName, err)
|
||
}
|
||
|
||
// CREATE DATABASE нельзя в транзакции
|
||
if _, err = s.adminDB.ExecContext(ctx,
|
||
fmt.Sprintf(`CREATE DATABASE %s OWNER %s`, dbName, userName),
|
||
); err != nil {
|
||
return fmt.Errorf("iotpg: create database %s: %w", dbName, err)
|
||
}
|
||
|
||
if _, err = s.adminDB.ExecContext(ctx,
|
||
`INSERT INTO tenant_credentials (namespace, pg_password) VALUES ($1, $2)
|
||
ON CONFLICT (namespace) DO NOTHING`,
|
||
namespace, password,
|
||
); err != nil {
|
||
return fmt.Errorf("iotpg: save credentials %s: %w", namespace, err)
|
||
}
|
||
s.log.Info("iotpg: created tenant DB", "namespace", namespace, "db", dbName)
|
||
}
|
||
|
||
// Создаём таблицу в tenant DB (суперюзер имеет доступ)
|
||
tenantDB, err := s.getTenantDB(ctx, namespace)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
_, err = tenantDB.ExecContext(ctx, `
|
||
CREATE TABLE IF NOT EXISTS iot_telemetry (
|
||
id BIGSERIAL PRIMARY KEY,
|
||
device_id TEXT NOT NULL,
|
||
ts TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||
payload JSONB NOT NULL
|
||
);
|
||
CREATE INDEX IF NOT EXISTS idx_iot_telemetry_device_ts
|
||
ON iot_telemetry (device_id, ts DESC);
|
||
`)
|
||
return err
|
||
}
|
||
|
||
// InsertTelemetry записывает строку телеметрии в tenant DB.
|
||
func (s *IoTPostgresStore) InsertTelemetry(ctx context.Context, namespace, deviceID string, payload json.RawMessage) error {
|
||
tenantDB, err := s.getTenantDB(ctx, namespace)
|
||
if err != nil {
|
||
return fmt.Errorf("iotpg: get tenant DB for insert: %w", err)
|
||
}
|
||
_, err = tenantDB.ExecContext(ctx,
|
||
`INSERT INTO iot_telemetry (device_id, payload) VALUES ($1, $2)`,
|
||
deviceID, []byte(payload),
|
||
)
|
||
return err
|
||
}
|
||
|
||
// QueryTelemetry читает телеметрию из tenant DB (ts DESC).
|
||
// deviceID — фильтр (пустая строка = все устройства). limit — max записей (50..1000).
|
||
// Если tenant DB не существует (данных ещё нет) — возвращает пустой срез без ошибки.
|
||
func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, deviceID string, limit int) ([]TelemetryRow, error) {
|
||
if limit <= 0 {
|
||
limit = 50
|
||
}
|
||
if limit > 1000 {
|
||
limit = 1000
|
||
}
|
||
tenantDB, err := s.getTenantDB(ctx, namespace)
|
||
if err != nil {
|
||
// Если DB не существует — тенант ещё не отправлял данные, это нормально
|
||
if isDBNotExistErr(err) {
|
||
return []TelemetryRow{}, nil
|
||
}
|
||
return nil, fmt.Errorf("iotpg: get tenant DB for query: %w", err)
|
||
}
|
||
|
||
var rows *sql.Rows
|
||
if deviceID != "" {
|
||
rows, err = tenantDB.QueryContext(ctx,
|
||
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
||
WHERE device_id = $1 ORDER BY ts DESC LIMIT $2`,
|
||
deviceID, limit,
|
||
)
|
||
} else {
|
||
rows, err = tenantDB.QueryContext(ctx,
|
||
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
||
ORDER BY ts DESC LIMIT $1`,
|
||
limit,
|
||
)
|
||
}
|
||
if err != nil {
|
||
return nil, fmt.Errorf("iotpg: query telemetry for %s: %w", namespace, err)
|
||
}
|
||
defer rows.Close()
|
||
|
||
var result []TelemetryRow
|
||
for rows.Next() {
|
||
var r TelemetryRow
|
||
var rawPayload []byte
|
||
if err := rows.Scan(&r.ID, &r.DeviceID, &r.Ts, &rawPayload); err != nil {
|
||
return nil, fmt.Errorf("iotpg: scan row: %w", err)
|
||
}
|
||
r.Payload = json.RawMessage(rawPayload)
|
||
result = append(result, r)
|
||
}
|
||
return result, rows.Err()
|
||
}
|
||
|
||
// TenantPGStats — статистика телеметрии одного tenant за разные периоды.
|
||
type TenantPGStats struct {
|
||
Namespace string `json:"namespace"`
|
||
DBName string `json:"db_name"`
|
||
Total int64 `json:"total"`
|
||
Last1h int64 `json:"last_1h"`
|
||
Last24h int64 `json:"last_24h"`
|
||
Latest []TelemetryRow `json:"latest"`
|
||
Error string `json:"error,omitempty"`
|
||
}
|
||
|
||
// PostgresAdminStats — агрегированная статистика по всем tenant для страницы администратора.
|
||
type PostgresAdminStats struct {
|
||
Tenants []TenantPGStats `json:"tenants"`
|
||
TotalAll int64 `json:"total_all"`
|
||
Reachable bool `json:"reachable"`
|
||
}
|
||
|
||
// GetAdminStats собирает статистику по всем tenant из management DB.
|
||
// Используется только страницей администратора — не для tenant API.
|
||
func (s *IoTPostgresStore) GetAdminStats(ctx context.Context) (*PostgresAdminStats, error) {
|
||
// Список всех тенантов из management DB
|
||
nsRows, err := s.adminDB.QueryContext(ctx, `SELECT namespace FROM tenant_credentials ORDER BY namespace`)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("iotpg: list tenants: %w", err)
|
||
}
|
||
defer nsRows.Close()
|
||
|
||
var namespaces []string
|
||
for nsRows.Next() {
|
||
var ns string
|
||
if err := nsRows.Scan(&ns); err != nil {
|
||
return nil, err
|
||
}
|
||
namespaces = append(namespaces, ns)
|
||
}
|
||
if err := nsRows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
result := &PostgresAdminStats{
|
||
Reachable: true,
|
||
Tenants: make([]TenantPGStats, 0, len(namespaces)),
|
||
}
|
||
|
||
for _, ns := range namespaces {
|
||
stats := TenantPGStats{
|
||
Namespace: ns,
|
||
DBName: tenantDBName(ns),
|
||
}
|
||
|
||
tenantDB, err := s.getTenantDB(ctx, ns)
|
||
if err != nil {
|
||
stats.Error = err.Error()
|
||
result.Tenants = append(result.Tenants, stats)
|
||
continue
|
||
}
|
||
|
||
// Counts: total, last 1h, last 24h — одним запросом
|
||
err = tenantDB.QueryRowContext(ctx, `
|
||
SELECT
|
||
COUNT(*),
|
||
COUNT(*) FILTER (WHERE ts > NOW() - INTERVAL '1 hour'),
|
||
COUNT(*) FILTER (WHERE ts > NOW() - INTERVAL '24 hours')
|
||
FROM iot_telemetry`).Scan(&stats.Total, &stats.Last1h, &stats.Last24h)
|
||
if err != nil {
|
||
stats.Error = err.Error()
|
||
result.Tenants = append(result.Tenants, stats)
|
||
continue
|
||
}
|
||
result.TotalAll += stats.Total
|
||
|
||
// Последние 5 сообщений для предпросмотра
|
||
latestRows, err := tenantDB.QueryContext(ctx,
|
||
`SELECT id, device_id, ts, payload FROM iot_telemetry ORDER BY ts DESC LIMIT 5`)
|
||
if err == nil {
|
||
defer latestRows.Close()
|
||
for latestRows.Next() {
|
||
var r TelemetryRow
|
||
var rawPayload []byte
|
||
if err := latestRows.Scan(&r.ID, &r.DeviceID, &r.Ts, &rawPayload); err == nil {
|
||
r.Payload = json.RawMessage(rawPayload)
|
||
stats.Latest = append(stats.Latest, r)
|
||
}
|
||
}
|
||
}
|
||
|
||
result.Tenants = append(result.Tenants, stats)
|
||
}
|
||
|
||
return result, nil
|
||
}
|
||
|
||
// isDBNotExistErr проверяет что ошибка — «database does not exist» (PostgreSQL code 3D000).
|
||
// Используется в QueryTelemetry: если DB нет — просто нет данных, не ошибка системы.
|
||
func isDBNotExistErr(err error) bool {
|
||
var pqErr *pq.Error
|
||
if errors.As(err, &pqErr) {
|
||
// 3D000 = invalid_catalog_name (база данных не существует)
|
||
return pqErr.Code == "3D000"
|
||
}
|
||
return false
|
||
}
|
||
|
||
// Close закрывает все подключения (admin + tenant кэш).
|
||
func (s *IoTPostgresStore) Close() error {
|
||
s.tenants.Range(func(_, value any) bool {
|
||
if db, ok := value.(*sql.DB); ok {
|
||
db.Close()
|
||
}
|
||
return true
|
||
})
|
||
return s.adminDB.Close()
|
||
}
|
||
|
||
// getTenantDB возвращает *sql.DB для tenant DB из кэша или открывает новый.
|
||
func (s *IoTPostgresStore) getTenantDB(ctx context.Context, namespace string) (*sql.DB, error) {
|
||
if cached, ok := s.tenants.Load(namespace); ok {
|
||
return cached.(*sql.DB), nil
|
||
}
|
||
dsn := replaceDSNDatabase(s.adminDSN, tenantDBName(namespace))
|
||
db, err := sql.Open("postgres", dsn)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("iotpg: open tenant DB %s: %w", tenantDBName(namespace), err)
|
||
}
|
||
pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||
defer cancel()
|
||
if err := db.PingContext(pingCtx); err != nil {
|
||
db.Close()
|
||
return nil, fmt.Errorf("iotpg: ping tenant DB %s: %w", tenantDBName(namespace), err)
|
||
}
|
||
db.SetMaxOpenConns(5)
|
||
db.SetMaxIdleConns(2)
|
||
db.SetConnMaxLifetime(5 * time.Minute)
|
||
|
||
actual, loaded := s.tenants.LoadOrStore(namespace, db)
|
||
if loaded {
|
||
db.Close()
|
||
return actual.(*sql.DB), nil
|
||
}
|
||
return db, nil
|
||
}
|
||
|
||
// replaceDSNDatabase заменяет имя базы данных в DSN.
|
||
// Вход: "postgresql://user:pass@host:5432/iot_platform?sslmode=disable"
|
||
// Выход: "postgresql://user:pass@host:5432/tenant_abc?sslmode=disable"
|
||
func replaceDSNDatabase(dsn, newDBName string) string {
|
||
schemeEnd := strings.Index(dsn, "://")
|
||
if schemeEnd < 0 {
|
||
return dsn
|
||
}
|
||
hostPart := dsn[schemeEnd+3:]
|
||
slashIdx := strings.LastIndex(hostPart, "/")
|
||
if slashIdx < 0 {
|
||
return dsn
|
||
}
|
||
afterSlash := hostPart[slashIdx+1:]
|
||
suffix := ""
|
||
if qIdx := strings.Index(afterSlash, "?"); qIdx >= 0 {
|
||
suffix = afterSlash[qIdx:]
|
||
}
|
||
prefix := dsn[:schemeEnd+3+slashIdx+1]
|
||
return prefix + newDBName + suffix
|
||
}
|
||
|
||
// tenantDBName возвращает имя Postgres DATABASE для namespace.
|
||
// Дефисы заменяются на подчёркивания (pg не поддерживает дефисы в unquoted именах).
|
||
// Пример: "sless-abc123" → "tenant_sless_abc123"
|
||
func tenantDBName(namespace string) string {
|
||
return "tenant_" + strings.ReplaceAll(namespace, "-", "_")
|
||
}
|