541 lines
18 KiB
Go
541 lines
18 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
|
||
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, "-", "_")
|
||
}
|