feat(core): add naming, RCON, store, config, and image-build libraries
Foundational libraries: deterministic resource naming, the RCON client, the Postgres store with embedded SQL migrations, configuration loading, and container image-build helpers.
This commit is contained in:
18 files changed
+3475
No files matched your search
@@ -0,0 +1,136 @@
|
||||
// Package store owns the Felis business-layer database (spec §6): the embedded
|
||||
// schema migrations and the typed access layer. Migrations are applied by
|
||||
// `felis migrate up` under a Postgres advisory lock so concurrent api/operator
|
||||
// replicas cannot race each other.
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"embed"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// AdvisoryLockKey is the fixed pg_advisory_lock key guarding migrations. The
|
||||
// value is the ASCII bytes of "felis"; any replica running migrations contends
|
||||
// on the same key.
|
||||
const AdvisoryLockKey int64 = 0x66656c6973 // "felis"
|
||||
|
||||
//go:embed migrations/*.sql
|
||||
var migrationsFS embed.FS
|
||||
|
||||
// Migration is a single ordered schema step loaded from the embedded FS.
|
||||
type Migration struct {
|
||||
Version int
|
||||
Name string
|
||||
SQL string
|
||||
}
|
||||
|
||||
// Driver is the database-facing seam the migration engine drives. Splitting it
|
||||
// out lets the ordering/idempotency/lock logic be tested without a live
|
||||
// Postgres; PostgresDriver is the production implementation.
|
||||
type Driver interface {
|
||||
// Lock acquires the migration advisory lock, blocking until held.
|
||||
Lock(ctx context.Context) error
|
||||
// Unlock releases the advisory lock.
|
||||
Unlock(ctx context.Context) error
|
||||
// EnsureVersionTable creates the schema_migrations bookkeeping table.
|
||||
EnsureVersionTable(ctx context.Context) error
|
||||
// AppliedVersions returns the set of versions already applied.
|
||||
AppliedVersions(ctx context.Context) (map[int]struct{}, error)
|
||||
// Apply runs one migration and records it, atomically.
|
||||
Apply(ctx context.Context, m Migration) error
|
||||
}
|
||||
|
||||
// LoadMigrations parses the embedded migrations into an ascending, gap-tolerant
|
||||
// but duplicate-free list.
|
||||
func LoadMigrations() ([]Migration, error) {
|
||||
entries, err := fs.ReadDir(migrationsFS, "migrations")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read embedded migrations: %w", err)
|
||||
}
|
||||
var ms []Migration
|
||||
for _, e := range entries {
|
||||
if e.IsDir() || !strings.HasSuffix(e.Name(), ".sql") {
|
||||
continue
|
||||
}
|
||||
version, name, err := parseMigrationName(e.Name())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
data, err := migrationsFS.ReadFile("migrations/" + e.Name())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read migration %q: %w", e.Name(), err)
|
||||
}
|
||||
if strings.TrimSpace(string(data)) == "" {
|
||||
return nil, fmt.Errorf("migration %q is empty", e.Name())
|
||||
}
|
||||
ms = append(ms, Migration{Version: version, Name: name, SQL: string(data)})
|
||||
}
|
||||
sort.Slice(ms, func(i, j int) bool { return ms[i].Version < ms[j].Version })
|
||||
for i := 1; i < len(ms); i++ {
|
||||
if ms[i].Version == ms[i-1].Version {
|
||||
return nil, fmt.Errorf("duplicate migration version %d (%s, %s)", ms[i].Version, ms[i-1].Name, ms[i].Name)
|
||||
}
|
||||
}
|
||||
if len(ms) == 0 {
|
||||
return nil, fmt.Errorf("no migrations found")
|
||||
}
|
||||
return ms, nil
|
||||
}
|
||||
|
||||
// Up applies every pending migration in ascending order, exactly once, under
|
||||
// the advisory lock. It is safe to run concurrently from multiple replicas: the
|
||||
// lock serializes them and AppliedVersions makes the work idempotent.
|
||||
func Up(ctx context.Context, d Driver, migrations []Migration) (applied []int, err error) {
|
||||
if err := d.Lock(ctx); err != nil {
|
||||
return nil, fmt.Errorf("acquire migration lock: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
if uerr := d.Unlock(ctx); uerr != nil && err == nil {
|
||||
err = fmt.Errorf("release migration lock: %w", uerr)
|
||||
}
|
||||
}()
|
||||
|
||||
if err := d.EnsureVersionTable(ctx); err != nil {
|
||||
return nil, fmt.Errorf("ensure version table: %w", err)
|
||||
}
|
||||
done, err := d.AppliedVersions(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read applied versions: %w", err)
|
||||
}
|
||||
|
||||
ordered := append([]Migration(nil), migrations...)
|
||||
sort.Slice(ordered, func(i, j int) bool { return ordered[i].Version < ordered[j].Version })
|
||||
for _, m := range ordered {
|
||||
if _, ok := done[m.Version]; ok {
|
||||
continue
|
||||
}
|
||||
if err := d.Apply(ctx, m); err != nil {
|
||||
return applied, fmt.Errorf("apply migration %04d_%s: %w", m.Version, m.Name, err)
|
||||
}
|
||||
applied = append(applied, m.Version)
|
||||
}
|
||||
return applied, nil
|
||||
}
|
||||
|
||||
// parseMigrationName turns "0001_init.sql" into (1, "init").
|
||||
func parseMigrationName(filename string) (int, string, error) {
|
||||
base := strings.TrimSuffix(filename, ".sql")
|
||||
idx := strings.IndexByte(base, '_')
|
||||
if idx <= 0 {
|
||||
return 0, "", fmt.Errorf("migration %q must be named NNNN_name.sql", filename)
|
||||
}
|
||||
version, err := strconv.Atoi(base[:idx])
|
||||
if err != nil {
|
||||
return 0, "", fmt.Errorf("migration %q has a non-numeric version: %w", filename, err)
|
||||
}
|
||||
name := base[idx+1:]
|
||||
if name == "" {
|
||||
return 0, "", fmt.Errorf("migration %q is missing a name", filename)
|
||||
}
|
||||
return version, name, nil
|
||||
}
|
||||
@@ -0,0 +1,157 @@
|
||||
package store_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"felis.lolicon.best/internal/store"
|
||||
)
|
||||
|
||||
func TestLoadMigrationsOrderedAndWellFormed(t *testing.T) {
|
||||
ms, err := store.LoadMigrations()
|
||||
if err != nil {
|
||||
t.Fatalf("LoadMigrations: %v", err)
|
||||
}
|
||||
if len(ms) == 0 {
|
||||
t.Fatal("expected at least one migration")
|
||||
}
|
||||
if ms[0].Version != 1 || ms[0].Name != "init" {
|
||||
t.Errorf("first migration = %d_%s, want 0001_init", ms[0].Version, ms[0].Name)
|
||||
}
|
||||
for i := 1; i < len(ms); i++ {
|
||||
if ms[i].Version <= ms[i-1].Version {
|
||||
t.Errorf("migrations not strictly ascending at %d: %d then %d", i, ms[i-1].Version, ms[i].Version)
|
||||
}
|
||||
}
|
||||
// The init migration must define the core business tables (spec §6).
|
||||
for _, want := range []string{"CREATE TABLE users", "CREATE TABLE servers", "CREATE TABLE world_backups"} {
|
||||
if !strings.Contains(ms[0].SQL, want) {
|
||||
t.Errorf("init migration missing %q", want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// recordingDriver captures the migration engine's calls without a database.
|
||||
type recordingDriver struct {
|
||||
already map[int]struct{}
|
||||
applied []int
|
||||
locked bool
|
||||
unlocked bool
|
||||
ensured bool
|
||||
appliedWhileUnsafe bool // true if Apply ran while not locked or already unlocked
|
||||
failOn int // version whose Apply should fail (0 = never)
|
||||
}
|
||||
|
||||
func (d *recordingDriver) Lock(context.Context) error { d.locked = true; return nil }
|
||||
func (d *recordingDriver) Unlock(context.Context) error { d.unlocked = true; return nil }
|
||||
func (d *recordingDriver) EnsureVersionTable(context.Context) error {
|
||||
if !d.locked || d.unlocked {
|
||||
d.appliedWhileUnsafe = true
|
||||
}
|
||||
d.ensured = true
|
||||
return nil
|
||||
}
|
||||
func (d *recordingDriver) AppliedVersions(context.Context) (map[int]struct{}, error) {
|
||||
if d.already == nil {
|
||||
return map[int]struct{}{}, nil
|
||||
}
|
||||
return d.already, nil
|
||||
}
|
||||
func (d *recordingDriver) Apply(_ context.Context, m store.Migration) error {
|
||||
if !d.locked || d.unlocked {
|
||||
d.appliedWhileUnsafe = true
|
||||
}
|
||||
if d.failOn != 0 && m.Version == d.failOn {
|
||||
return errors.New("boom")
|
||||
}
|
||||
d.applied = append(d.applied, m.Version)
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestUpAppliesAllPendingInOrder(t *testing.T) {
|
||||
d := &recordingDriver{}
|
||||
ms := []store.Migration{
|
||||
{Version: 3, Name: "c", SQL: "x"},
|
||||
{Version: 1, Name: "a", SQL: "y"},
|
||||
{Version: 2, Name: "b", SQL: "z"},
|
||||
}
|
||||
applied, err := store.Up(context.Background(), d, ms)
|
||||
if err != nil {
|
||||
t.Fatalf("Up: %v", err)
|
||||
}
|
||||
if got := strings.Trim(strings.Join(intsToStrings(applied), ","), ""); got != "1,2,3" {
|
||||
t.Errorf("applied = %v, want [1 2 3]", applied)
|
||||
}
|
||||
if !d.locked || !d.unlocked || !d.ensured {
|
||||
t.Errorf("lifecycle flags: locked=%v unlocked=%v ensured=%v", d.locked, d.unlocked, d.ensured)
|
||||
}
|
||||
if d.appliedWhileUnsafe {
|
||||
t.Error("work ran outside the advisory lock")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpSkipsAlreadyApplied(t *testing.T) {
|
||||
d := &recordingDriver{already: map[int]struct{}{1: {}}}
|
||||
ms := []store.Migration{
|
||||
{Version: 1, Name: "a", SQL: "y"},
|
||||
{Version: 2, Name: "b", SQL: "z"},
|
||||
}
|
||||
applied, err := store.Up(context.Background(), d, ms)
|
||||
if err != nil {
|
||||
t.Fatalf("Up: %v", err)
|
||||
}
|
||||
if len(applied) != 1 || applied[0] != 2 {
|
||||
t.Errorf("applied = %v, want [2]", applied)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpStopsOnErrorButStillUnlocks(t *testing.T) {
|
||||
d := &recordingDriver{failOn: 2}
|
||||
ms := []store.Migration{
|
||||
{Version: 1, Name: "a", SQL: "y"},
|
||||
{Version: 2, Name: "b", SQL: "z"},
|
||||
{Version: 3, Name: "c", SQL: "x"},
|
||||
}
|
||||
applied, err := store.Up(context.Background(), d, ms)
|
||||
if err == nil {
|
||||
t.Fatal("expected an error when a migration fails")
|
||||
}
|
||||
if len(applied) != 1 || applied[0] != 1 {
|
||||
t.Errorf("applied = %v, want only [1] before the failure", applied)
|
||||
}
|
||||
if !d.unlocked {
|
||||
t.Error("advisory lock must be released even when a migration fails")
|
||||
}
|
||||
}
|
||||
|
||||
func intsToStrings(in []int) []string {
|
||||
out := make([]string, len(in))
|
||||
for i, v := range in {
|
||||
out[i] = itoa(v)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func itoa(v int) string {
|
||||
if v == 0 {
|
||||
return "0"
|
||||
}
|
||||
neg := v < 0
|
||||
if neg {
|
||||
v = -v
|
||||
}
|
||||
var buf [20]byte
|
||||
i := len(buf)
|
||||
for v > 0 {
|
||||
i--
|
||||
buf[i] = byte('0' + v%10)
|
||||
v /= 10
|
||||
}
|
||||
if neg {
|
||||
i--
|
||||
buf[i] = '-'
|
||||
}
|
||||
return string(buf[i:])
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
-- Felis business-layer schema (spec §6). The CRD is the lifecycle
|
||||
-- source-of-truth; this database owns only what the CRD cannot express:
|
||||
-- ownership/claim, account links, quotas, image admission, builds, backups,
|
||||
-- audit. Authoritative CRD fields are never duplicated here.
|
||||
|
||||
CREATE TYPE user_role AS ENUM ('admin','user');
|
||||
CREATE TYPE build_status AS ENUM ('pending','building','succeeded','failed','cancelled');
|
||||
CREATE TYPE backup_status AS ENUM ('present','expired','deleted');
|
||||
|
||||
CREATE TABLE users (
|
||||
id text PRIMARY KEY, username text UNIQUE NOT NULL, email text,
|
||||
role user_role NOT NULL DEFAULT 'user', created_at timestamptz NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Identity bridge: web identity <-> MC UUID (claim/owner-only/allowlist rely on it).
|
||||
CREATE TABLE account_links (
|
||||
user_id text NOT NULL REFERENCES users(id), mc_uuid uuid NOT NULL,
|
||||
verified_at timestamptz NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (user_id, mc_uuid), UNIQUE (mc_uuid)
|
||||
);
|
||||
CREATE TABLE account_link_codes ( code text PRIMARY KEY, mc_uuid uuid NOT NULL, expires_at timestamptz NOT NULL );
|
||||
|
||||
CREATE TABLE quotas (
|
||||
user_id text PRIMARY KEY REFERENCES users(id),
|
||||
max_servers int, max_cpu_milli int, max_memory_mb int, max_storage_gb int
|
||||
);
|
||||
|
||||
-- Business projection: the CRD lives in K8s; this stores only the
|
||||
-- ownership/activity/warning that the CRD cannot express, plus a fast-query cache.
|
||||
CREATE TABLE servers (
|
||||
name text PRIMARY KEY, -- matches CRD metadata.name
|
||||
owner_id text REFERENCES users(id), -- NULL until claimed; reaper resets to NULL
|
||||
claimed_at timestamptz,
|
||||
last_active_at timestamptz NOT NULL DEFAULT now(), -- max(last human join, created_at)
|
||||
warned_3d_at timestamptz, warned_1d_at timestamptz, -- reaper warning dedup; cleared on renewal
|
||||
cached_phase text, -- CRD status projection, non-authoritative
|
||||
created_at timestamptz NOT NULL DEFAULT now(), deleted_at timestamptz
|
||||
);
|
||||
CREATE TABLE server_aliases ( subdomain text PRIMARY KEY, server_name text NOT NULL REFERENCES servers(name) );
|
||||
CREATE TABLE server_allowlist ( -- autostartPolicy=allowlist; first join auto-appends
|
||||
server_name text NOT NULL REFERENCES servers(name), mc_uuid uuid NOT NULL,
|
||||
added_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (server_name, mc_uuid)
|
||||
);
|
||||
|
||||
-- Image admission (dynamic, auditable -> DB, not toml).
|
||||
CREATE TABLE image_whitelist (
|
||||
image_ref text PRIMARY KEY, -- registry/foo:1.0 or registry/foo:*
|
||||
source text NOT NULL DEFAULT 'built', -- built (cluster build) | external (pushed)
|
||||
build_id text, added_by text NOT NULL, enabled boolean NOT NULL DEFAULT true,
|
||||
added_at timestamptz NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE image_builds (
|
||||
id text PRIMARY KEY, image_ref text NOT NULL, status build_status NOT NULL DEFAULT 'pending',
|
||||
dockerfile text NOT NULL, -- archived for audit
|
||||
context_ref text, base_image text, -- resolved FROM, audit
|
||||
requested_by text NOT NULL, job_name text, log_ref text, error text,
|
||||
created_at timestamptz NOT NULL DEFAULT now(), finished_at timestamptz
|
||||
);
|
||||
|
||||
-- World backups (reaper output; not FK'd to servers, which may be reset/deleted).
|
||||
CREATE TABLE world_backups (
|
||||
id text PRIMARY KEY, server_name text NOT NULL, former_owner text,
|
||||
backup_ref text NOT NULL, -- WorldArchiver location (ArchiveRef)
|
||||
size_bytes bigint, reason text NOT NULL, -- inactive_15d | manual
|
||||
status backup_status NOT NULL DEFAULT 'present',
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
expires_at timestamptz NOT NULL, -- created_at + 3mo
|
||||
deleted_at timestamptz
|
||||
);
|
||||
|
||||
CREATE TABLE audit_logs (
|
||||
id bigserial PRIMARY KEY, actor text NOT NULL, source text NOT NULL, action text NOT NULL,
|
||||
server_name text, request_id text, payload jsonb, created_at timestamptz NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE tokens ( id text PRIMARY KEY, name text NOT NULL, token_hash text NOT NULL, scope jsonb NOT NULL, expires_at timestamptz );
|
||||
@@ -0,0 +1,42 @@
|
||||
-- User-submitted modpack approval lane (a user-directed extension over the §16
|
||||
-- build subsystem; see internal/submit for provenance). This is the UNTRUSTED-
|
||||
-- origin counterpart to the admin build path (POST /images/build): an ordinary
|
||||
-- user may upload a modpack but cannot start a build directly. Each upload lands
|
||||
-- here as pending_review; an admin must approve it before anything is built, and
|
||||
-- the approved submission then routes through the SAME Trivy-gated Kaniko build
|
||||
-- as an admin build (build subsystem §16). Approval is a human gate layered in
|
||||
-- FRONT of the automatic scan, never instead of it — a CRITICAL CVE still fails
|
||||
-- the build and nothing is admitted even after a human approved.
|
||||
--
|
||||
-- Source of truth (spec §1): this row is the Postgres BUSINESS authority for the
|
||||
-- approval (verdict + reviewer); the build EXECUTION lives in image_builds,
|
||||
-- linked by build_id once Builder.Submit succeeds. The approval never copies the
|
||||
-- build's authoritative fields.
|
||||
--
|
||||
-- Trust note: the platform derives BOTH the push target (image_ref) and the
|
||||
-- build context (context_ref) from the submission id — neither is free-form user
|
||||
-- input — so an untrusted submitter can never point the build at an arbitrary
|
||||
-- source or collide with the platform image namespace. There is deliberately no
|
||||
-- `origin` column: image_submissions is ONLY the user-upload lane (the platform
|
||||
-- uses the direct build path), and submitted_by already records the origin.
|
||||
|
||||
CREATE TYPE submission_status AS ENUM ('pending_review','approved','rejected');
|
||||
|
||||
CREATE TABLE image_submissions (
|
||||
id text PRIMARY KEY, -- lowercase, namespaces the derived image/context refs
|
||||
submitted_by text NOT NULL, -- uploading user's id (untrusted origin)
|
||||
display_name text NOT NULL, -- human-friendly label for the modpack
|
||||
context_ref text NOT NULL, -- DERIVED pinned build context (not user-supplied)
|
||||
status submission_status NOT NULL DEFAULT 'pending_review',
|
||||
image_ref text, -- DERIVED {registry}/user-uploads/{id}:latest, set at approve
|
||||
build_id text, -- image_builds.id, set only after Builder.Submit succeeds
|
||||
reviewed_by text, -- admin who approved/rejected
|
||||
reject_reason text, -- set on rejection
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
reviewed_at timestamptz
|
||||
);
|
||||
|
||||
-- The admin review queue scans by status (pending first); the per-user index
|
||||
-- serves the "my submissions" list.
|
||||
CREATE INDEX image_submissions_status_idx ON image_submissions (status, created_at);
|
||||
CREATE INDEX image_submissions_submitted_by_idx ON image_submissions (submitted_by, created_at DESC);
|
||||
@@ -0,0 +1,94 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib" // register the "pgx" database/sql driver
|
||||
)
|
||||
|
||||
// PostgresDriver is the production Driver, backed by a database/sql pool using
|
||||
// the pgx stdlib driver.
|
||||
type PostgresDriver struct {
|
||||
db *sql.DB
|
||||
}
|
||||
|
||||
// Open dials dsn and returns a PostgresDriver. The caller owns Close.
|
||||
func Open(ctx context.Context, dsn string) (*PostgresDriver, error) {
|
||||
db, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open postgres: %w", err)
|
||||
}
|
||||
if err := db.PingContext(ctx); err != nil {
|
||||
db.Close()
|
||||
return nil, fmt.Errorf("ping postgres: %w", err)
|
||||
}
|
||||
return &PostgresDriver{db: db}, nil
|
||||
}
|
||||
|
||||
// DB exposes the underlying pool for the access layer.
|
||||
func (d *PostgresDriver) DB() *sql.DB { return d.db }
|
||||
|
||||
// Close releases the pool.
|
||||
func (d *PostgresDriver) Close() error { return d.db.Close() }
|
||||
|
||||
// Lock takes the session-level advisory lock that serializes migrations.
|
||||
func (d *PostgresDriver) Lock(ctx context.Context) error {
|
||||
_, err := d.db.ExecContext(ctx, "SELECT pg_advisory_lock($1)", AdvisoryLockKey)
|
||||
return err
|
||||
}
|
||||
|
||||
// Unlock releases the advisory lock.
|
||||
func (d *PostgresDriver) Unlock(ctx context.Context) error {
|
||||
_, err := d.db.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", AdvisoryLockKey)
|
||||
return err
|
||||
}
|
||||
|
||||
// EnsureVersionTable creates the bookkeeping table if absent.
|
||||
func (d *PostgresDriver) EnsureVersionTable(ctx context.Context) error {
|
||||
const ddl = `CREATE TABLE IF NOT EXISTS schema_migrations (
|
||||
version int PRIMARY KEY,
|
||||
name text NOT NULL,
|
||||
applied_at timestamptz NOT NULL DEFAULT now()
|
||||
)`
|
||||
_, err := d.db.ExecContext(ctx, ddl)
|
||||
return err
|
||||
}
|
||||
|
||||
// AppliedVersions reads the set of recorded versions.
|
||||
func (d *PostgresDriver) AppliedVersions(ctx context.Context) (map[int]struct{}, error) {
|
||||
rows, err := d.db.QueryContext(ctx, "SELECT version FROM schema_migrations")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
out := map[int]struct{}{}
|
||||
for rows.Next() {
|
||||
var v int
|
||||
if err := rows.Scan(&v); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[v] = struct{}{}
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// Apply runs the migration body and records it in one transaction, so a failure
|
||||
// never leaves a half-applied version marked as done.
|
||||
func (d *PostgresDriver) Apply(ctx context.Context, m Migration) error {
|
||||
tx, err := d.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback() //nolint:errcheck // rollback after a successful commit is a no-op
|
||||
|
||||
if _, err := tx.ExecContext(ctx, m.SQL); err != nil {
|
||||
return fmt.Errorf("exec body: %w", err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx,
|
||||
"INSERT INTO schema_migrations (version, name) VALUES ($1, $2)", m.Version, m.Name); err != nil {
|
||||
return fmt.Errorf("record version: %w", err)
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
Reference in new issue
Block a user