Files
IoT/internal/storage/iotpg/iot_telemetry_store.go
T
Naeel 1a94241c62 feat: перенос IoT managed service из sless в отдельную репу
Компоненты:
- 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 адаптированы.
2026-04-12 14:29:43 +03:00

399 lines
14 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Создано: 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, "-", "_")
}