netpulse-migrate замість PowerShell-скрипта: у контейнері немає ані psql, ані PowerShell, а тягнути клієнт Postgres в образ заради одного запуску — це половина дистрибутива на порожньому місці. Міграції вшиті через embed і переїхали в server/migrations: embed не бачить нічого за межами кореня свого модуля, а міграції поруч із бінарником, який їх накочує, не можуть розійтися версіями. Накочування під advisory-блокуванням: два інстанси при rolling update інакше застосували б ту саму міграцію двічі. Кожен файл в одній транзакції разом із записом у schema_migrations; виняток — continuous aggregates, які TimescaleDB забороняє в транзакції. Змінена вже застосована міграція зупиняє запуск: у різних інсталяціях інакше опиниться різна схема під одним номером. Перевірено на чистій базі: 23 міграції, 101 таблиця, повторний запуск каже «схема актуальна». Веб віддає сам API через embed: на self-hosted це прибирає з інструкції встановлення цілий компонент. Три політики кешування — назавжди для assets із хешем у імені, ніколи для index.html, коротко для решти. Знайдено живим прогоном: невідомий шлях під /api/ віддавав 200 з index.html, і клієнт падав на розборі HTML як JSON замість чесного 404. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
238 lines
7.8 KiB
Go
238 lines
7.8 KiB
Go
// Команда netpulse-migrate — накочування схеми.
|
||
//
|
||
// Бінарник, а не скрипт: у контейнері немає ані PowerShell, ані psql, а
|
||
// вимагати клієнт Postgres поруч із застосунком означає тягнути в образ
|
||
// половину дистрибутива заради одного запуску.
|
||
//
|
||
// Міграції вшиті в бінарник через embed: файл, який лежить поруч,
|
||
// рано чи пізно виявиться версією з іншого релізу.
|
||
package main
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"errors"
|
||
"flag"
|
||
"fmt"
|
||
"io/fs"
|
||
"os"
|
||
"os/signal"
|
||
"regexp"
|
||
"sort"
|
||
"strings"
|
||
"syscall"
|
||
"time"
|
||
|
||
"github.com/jackc/pgx/v5"
|
||
"github.com/jackc/pgx/v5/pgxpool"
|
||
schema "github.com/netpulse/netpulse/server/migrations"
|
||
)
|
||
|
||
// Шлях відносний до кореня модуля не працює: embed бачить лише те, що
|
||
// лежить у каталозі пакета або нижче. Тому міграції тягнуться через
|
||
// окремий пакет, який стоїть поруч із ними.
|
||
var migrations = schema.Files
|
||
|
||
// migrateLockKey — advisory-блокування на час накочування.
|
||
//
|
||
// Два інстанси, що стартують одночасно (rolling update, docker compose
|
||
// зі скейлом), інакше накотили б ту саму міграцію двічі: перевірка
|
||
// «чи застосовано» і сам запис — різні моменти часу.
|
||
const migrateLockKey = 0x6e70_6d67 // "npmg"
|
||
|
||
// continuousRe ловить те, що не можна виконати в транзакції.
|
||
//
|
||
// TimescaleDB відмовляє: CREATE MATERIALIZED VIEW WITH
|
||
// (timescaledb.continuous) поза транзакційним блоком. Такі файли
|
||
// накочуються без обгортки — ціна в тому, що збій посеред файлу лишає
|
||
// половину змін, тому їх свідомо тримають короткими.
|
||
var continuousRe = regexp.MustCompile(`timescaledb\.continuous`)
|
||
|
||
func main() {
|
||
if err := run(); err != nil {
|
||
fmt.Fprintln(os.Stderr, "netpulse-migrate:", err)
|
||
os.Exit(1)
|
||
}
|
||
}
|
||
|
||
func run() error {
|
||
dsn := flag.String("dsn", os.Getenv("NETPULSE_DSN"), "postgres://user:pass@host:5432/db")
|
||
dryRun := flag.Bool("dry-run", false, "лише показати, що буде застосовано")
|
||
timeout := flag.Duration("timeout", 10*time.Minute, "стеля на всі міграції")
|
||
flag.Parse()
|
||
|
||
if *dsn == "" {
|
||
return errors.New("не вказано -dsn (або NETPULSE_DSN)")
|
||
}
|
||
|
||
files, err := listMigrations()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||
defer stop()
|
||
ctx, cancel := context.WithTimeout(ctx, *timeout)
|
||
defer cancel()
|
||
|
||
pool, err := pgxpool.New(ctx, *dsn)
|
||
if err != nil {
|
||
return fmt.Errorf("підключення: %w", err)
|
||
}
|
||
defer pool.Close()
|
||
|
||
// Одне з'єднання на весь запуск: advisory-блокування живе рівно
|
||
// стільки, скільки з'єднання, яке його взяло.
|
||
conn, err := pool.Acquire(ctx)
|
||
if err != nil {
|
||
return fmt.Errorf("підключення: %w", err)
|
||
}
|
||
defer conn.Release()
|
||
|
||
if err := bootstrap(ctx, conn.Conn()); err != nil {
|
||
return err
|
||
}
|
||
|
||
if _, err := conn.Exec(ctx, `SELECT pg_advisory_lock($1)`, int64(migrateLockKey)); err != nil {
|
||
return fmt.Errorf("блокування: %w", err)
|
||
}
|
||
defer func() {
|
||
_, _ = conn.Exec(context.WithoutCancel(ctx),
|
||
`SELECT pg_advisory_unlock($1)`, int64(migrateLockKey))
|
||
}()
|
||
|
||
applied, err := appliedVersions(ctx, conn.Conn())
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
pending := 0
|
||
for _, f := range files {
|
||
version := strings.TrimSuffix(f, ".sql")
|
||
body, err := migrations.ReadFile(f)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
sum := sha256.Sum256(body)
|
||
checksum := hex.EncodeToString(sum[:])
|
||
|
||
if old, ok := applied[version]; ok {
|
||
// Змінена вже застосована міграція — це майже завжди
|
||
// помилка: у різних інсталяціях опиниться різна схема під
|
||
// одним номером. Кажемо про це й зупиняємось.
|
||
if old != checksum {
|
||
return fmt.Errorf(
|
||
"міграцію %s змінено після застосування (було %s…, стало %s…).\n"+
|
||
"Не правте застосовані міграції — додайте нову",
|
||
version, old[:8], checksum[:8])
|
||
}
|
||
continue
|
||
}
|
||
|
||
pending++
|
||
if *dryRun {
|
||
fmt.Printf("буде застосовано %s\n", version)
|
||
continue
|
||
}
|
||
|
||
fmt.Printf("застосовую %s\n", version)
|
||
if err := apply(ctx, conn.Conn(), version, checksum, string(body)); err != nil {
|
||
return fmt.Errorf("міграція %s: %w", version, err)
|
||
}
|
||
}
|
||
|
||
switch {
|
||
case pending == 0:
|
||
fmt.Println("схема актуальна")
|
||
case *dryRun:
|
||
fmt.Printf("непримінених міграцій: %d\n", pending)
|
||
default:
|
||
fmt.Printf("застосовано міграцій: %d\n", pending)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func listMigrations() ([]string, error) {
|
||
entries, err := fs.ReadDir(migrations, ".")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var out []string
|
||
for _, e := range entries {
|
||
if !e.IsDir() && strings.HasSuffix(e.Name(), ".sql") {
|
||
out = append(out, e.Name())
|
||
}
|
||
}
|
||
// Порядок — за іменем: номер на початку файлу і є версією.
|
||
sort.Strings(out)
|
||
if len(out) == 0 {
|
||
return nil, errors.New("у бінарнику немає жодної міграції")
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
func bootstrap(ctx context.Context, conn *pgx.Conn) error {
|
||
_, err := conn.Exec(ctx, `
|
||
CREATE TABLE IF NOT EXISTS public.schema_migrations (
|
||
version text PRIMARY KEY,
|
||
checksum text NOT NULL,
|
||
applied_at timestamptz NOT NULL DEFAULT now()
|
||
)
|
||
`)
|
||
return err
|
||
}
|
||
|
||
func appliedVersions(ctx context.Context, conn *pgx.Conn) (map[string]string, error) {
|
||
rows, err := conn.Query(ctx, `SELECT version, checksum FROM public.schema_migrations`)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
out := map[string]string{}
|
||
for rows.Next() {
|
||
var v, c string
|
||
if err := rows.Scan(&v, &c); err != nil {
|
||
return nil, err
|
||
}
|
||
out[v] = c
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
// apply накочує один файл.
|
||
//
|
||
// Зазвичай в одній транзакції разом із записом у schema_migrations: без
|
||
// цього збій посеред файлу лишає схему в стані, який ніхто не описував,
|
||
// а наступний запуск вважає міграцію незастосованою й повторює її.
|
||
func apply(ctx context.Context, conn *pgx.Conn, version, checksum, body string) error {
|
||
if continuousRe.MatchString(body) {
|
||
// TimescaleDB забороняє continuous aggregates у транзакції.
|
||
// Позначку ставимо після успіху: інакше збій лишив би міграцію
|
||
// «застосованою» без застосування.
|
||
if _, err := conn.Exec(ctx, body); err != nil {
|
||
return err
|
||
}
|
||
_, err := conn.Exec(ctx,
|
||
`INSERT INTO public.schema_migrations (version, checksum) VALUES ($1, $2)`,
|
||
version, checksum)
|
||
return err
|
||
}
|
||
|
||
tx, err := conn.Begin(ctx)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer func() { _ = tx.Rollback(ctx) }()
|
||
|
||
if _, err := tx.Exec(ctx, body); err != nil {
|
||
return err
|
||
}
|
||
if _, err := tx.Exec(ctx,
|
||
`INSERT INTO public.schema_migrations (version, checksum) VALUES ($1, $2)`,
|
||
version, checksum); err != nil {
|
||
return err
|
||
}
|
||
return tx.Commit(ctx)
|
||
}
|