v0.1.59: IoT telemetry pipeline — Postgres storage + REST API + UI table
This commit is contained in:
@@ -0,0 +1,288 @@
|
||||
// Создано: 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"
|
||||
"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).
|
||||
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 {
|
||||
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()
|
||||
}
|
||||
|
||||
// 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, "-", "_")
|
||||
}
|
||||
Reference in New Issue
Block a user