Files
IoT/internal/storage/iotpg/iot_telemetry_store.go
T

541 lines
18 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
tenantMu sync.Map // namespace → *sync.Mutex (дедупликация открытия)
ensured sync.Map // namespace → struct{} (EnsureTenantDB уже выполнен)
log *slog.Logger
// батчинг вставок (фикс MEDIUM ревью 2026-08-16: 1000 msg/s = 1000 INSERT)
mu sync.Mutex
batches map[string][]telemetryInsert
stopFlush chan struct{}
flushWG sync.WaitGroup
}
// telemetryInsert — строка для батч-вставки.
type telemetryInsert struct {
deviceID string
payload []byte
}
// telemetryBatchSize — размер батча перед синхронным flush.
const telemetryBatchSize = 100
// telemetryFlushInterval — период фонового flush неполных батчей.
const telemetryFlushInterval = 200 * time.Millisecond
// 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(10)
db.SetMaxIdleConns(2)
db.SetConnMaxLifetime(5 * time.Minute)
store := &IoTPostgresStore{
adminDB: db,
adminDSN: adminDSN,
batches: make(map[string][]telemetryInsert),
stopFlush: make(chan struct{}),
log: log,
}
if err := store.initManagementSchema(ctx); err != nil {
db.Close()
return nil, fmt.Errorf("iotpg: init management schema: %w", err)
}
store.startFlusher()
log.Info("iotpg: connected to IoT Postgres management DB")
return store, nil
}
// startFlusher — фоновая горутина периодического flush неполных батчей.
func (s *IoTPostgresStore) startFlusher() {
s.flushWG.Add(1)
go func() {
defer s.flushWG.Done()
t := time.NewTicker(telemetryFlushInterval)
defer t.Stop()
for {
select {
case <-s.stopFlush:
return
case <-t.C:
s.flushAll(false)
}
}
}()
}
// flushAll — flush всех накопленных батчей. reportErr=false: только лог.
func (s *IoTPostgresStore) flushAll(reportErr bool) error {
s.mu.Lock()
if len(s.batches) == 0 {
s.mu.Unlock()
return nil
}
batches := s.batches
s.batches = make(map[string][]telemetryInsert)
s.mu.Unlock()
var firstErr error
for ns, rows := range batches {
if err := s.flushTenant(ns, rows); err != nil {
s.log.Error("iotpg: flush batch", "namespace", ns, "rows", len(rows), "err", err)
if firstErr == nil {
firstErr = err
}
}
}
return firstErr
}
// flushTenant — INSERT батча строк одного тенанта (multi-VALUES).
func (s *IoTPostgresStore) flushTenant(ns string, rows []telemetryInsert) error {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
tenantDB, err := s.getTenantDB(ctx, ns)
if err != nil {
return err
}
// INSERT INTO iot_telemetry (device_id, payload) VALUES ($1,$2),($3,$4),...
var sb strings.Builder
sb.WriteString("INSERT INTO iot_telemetry (device_id, payload) VALUES ")
args := make([]any, 0, len(rows)*2)
for i, r := range rows {
if i > 0 {
sb.WriteString(",")
}
sb.WriteString(fmt.Sprintf("($%d,$%d)", i*2+1, i*2+2))
args = append(args, r.deviceID, r.payload)
}
_, err = tenantDB.ExecContext(ctx, sb.String(), args...)
return err
}
// 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.
// Идемпотентен — повторный вызов безопасен (кэш ensured пропускает проверки).
func (s *IoTPostgresStore) EnsureTenantDB(ctx context.Context, namespace string) error {
if _, ok := s.ensured.Load(namespace); ok {
return nil
}
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()
// Роль создаём только если её нет (DO-блоки НЕ принимают параметры —
// прецедент 2026-08-16: "got 2 parameters but the statement requires 0").
var roleExists bool
err = s.adminDB.QueryRowContext(ctx,
`SELECT EXISTS(SELECT 1 FROM pg_roles WHERE rolname = $1)`, userName,
).Scan(&roleExists)
if err != nil {
return fmt.Errorf("iotpg: check role %s: %w", userName, err)
}
if !roleExists {
_, err = s.adminDB.ExecContext(ctx,
`CREATE USER `+pq.QuoteIdentifier(userName)+` WITH PASSWORD `+pq.QuoteLiteral(password),
)
if err != nil {
var pqErr *pq.Error
if errors.As(err, &pqErr) && pqErr.Code == "42710" { // duplicate_object
// гонка: роль создал параллельный вызов — ок
} else {
return fmt.Errorf("iotpg: create user %s: %w", userName, err)
}
}
}
// PG15+: GRANT role TO current_user перед CREATE DATABASE ... OWNER
if _, err = s.adminDB.ExecContext(ctx,
`GRANT `+pq.QuoteIdentifier(userName)+` TO CURRENT_USER`,
); err != nil {
return fmt.Errorf("iotpg: grant role %s: %w", userName, err)
}
// CREATE DATABASE нельзя в транзакции
if _, err = s.adminDB.ExecContext(ctx,
`CREATE DATABASE `+pq.QuoteIdentifier(dbName)+` OWNER `+pq.QuoteIdentifier(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);
`)
if err != nil {
return err
}
s.ensured.Store(namespace, struct{}{})
return nil
}
// InsertTelemetry ставит строку в батч-буфер тенанта.
// При накоплении telemetryBatchSize строк батч пишется синхронно (ошибка
// возвращается вызывающему); неполные батчи дописывает фоновый flusher.
func (s *IoTPostgresStore) InsertTelemetry(ctx context.Context, namespace, deviceID string, payload json.RawMessage) error {
s.mu.Lock()
s.batches[namespace] = append(s.batches[namespace], telemetryInsert{
deviceID: deviceID,
payload: append([]byte(nil), payload...),
})
full := len(s.batches[namespace]) >= telemetryBatchSize
var rows []telemetryInsert
if full {
rows = s.batches[namespace]
delete(s.batches, namespace)
}
s.mu.Unlock()
if full {
return s.flushTenant(namespace, rows)
}
return nil
}
// QueryTelemetry читает телеметрию из tenant DB (ts DESC).
// deviceID — фильтр (пустая строка = все устройства). limit — max записей
// (50..1000), offset — смещение для пагинации.
// Если tenant DB не существует (данных ещё нет) — пустой срез без ошибки.
func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, deviceID string, limit, offset int) ([]TelemetryRow, error) {
if limit <= 0 {
limit = 50
}
if limit > 1000 {
limit = 1000
}
if offset < 0 {
offset = 0
}
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 OFFSET $3`,
deviceID, limit, offset,
)
} else {
rows, err = tenantDB.QueryContext(ctx,
`SELECT id, device_id, ts, payload FROM iot_telemetry
ORDER BY ts DESC LIMIT $1 OFFSET $2`,
limit, offset,
)
}
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 {
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)
}
}
latestRows.Close() // закрываем сразу, не defer в цикле
}
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 останавливает flusher, дописывает остаток и закрывает все подключения.
func (s *IoTPostgresStore) Close() error {
close(s.stopFlush)
s.flushWG.Wait()
_ = s.flushAll(true)
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 из кэша или открывает новый.
// Открытие дедуплицируется per-namespace мьютексом (гонка Load→LoadOrStore
// из ревью 2026-08-16).
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
}
m, _ := s.tenantMu.LoadOrStore(namespace, &sync.Mutex{})
mu := m.(*sync.Mutex)
mu.Lock()
defer mu.Unlock()
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, "-", "_")
}