// Package migrate implementiert FDN-02s Migrationsmechanik: versionierte, // rueckrollbare SQL-Migrationsdateien (kein ORM), mit einer // schema_migrations-Tabelle als Fortschrittsspeicher — dasselbe Prinzip wie // NEXARCH Core (internal/migrate), hier eigenstaendig implementiert, da DMS // ein eigenes Go-Modul ist und Cores internal/-Pakete nicht importieren // kann. package migrate import ( "context" "fmt" "os" "path/filepath" "sort" "strings" "github.com/jackc/pgx/v5/pgxpool" ) // Migration ist eine einzelne versionierte Migrationsdatei. type Migration struct { Version string // Dateiname ohne .up.sql/.down.sql, z.B. "0001_documents" UpSQL string DownSQL string } // Load liest alle *.up.sql/*.down.sql-Paare aus dir, sortiert nach // Dateiname (Akzeptanzkriterium: Migrationen sind versioniert). func Load(dir string) ([]Migration, error) { entries, err := os.ReadDir(dir) if err != nil { return nil, fmt.Errorf("migrationsverzeichnis %q lesen: %w", dir, err) } var versions []string for _, e := range entries { if e.IsDir() || !strings.HasSuffix(e.Name(), ".up.sql") { continue } versions = append(versions, strings.TrimSuffix(e.Name(), ".up.sql")) } sort.Strings(versions) migrations := make([]Migration, 0, len(versions)) for _, v := range versions { up, err := os.ReadFile(filepath.Join(dir, v+".up.sql")) if err != nil { return nil, fmt.Errorf("migration %q: up.sql lesen: %w", v, err) } down, err := os.ReadFile(filepath.Join(dir, v+".down.sql")) if err != nil { return nil, fmt.Errorf("migration %q: down.sql lesen (jede Migration braucht ein Rollback): %w", v, err) } migrations = append(migrations, Migration{Version: v, UpSQL: string(up), DownSQL: string(down)}) } return migrations, nil } func ensureTrackingTable(ctx context.Context, pool *pgxpool.Pool) error { _, err := pool.Exec(ctx, ` CREATE TABLE IF NOT EXISTS schema_migrations ( version TEXT PRIMARY KEY, applied_at TIMESTAMPTZ NOT NULL DEFAULT now() ) `) if err != nil { return fmt.Errorf("schema_migrations anlegen: %w", err) } return nil } func appliedVersions(ctx context.Context, pool *pgxpool.Pool) (map[string]bool, error) { rows, err := pool.Query(ctx, `SELECT version FROM schema_migrations`) if err != nil { return nil, fmt.Errorf("angewendete migrationen lesen: %w", err) } defer rows.Close() applied := map[string]bool{} for rows.Next() { var v string if err := rows.Scan(&v); err != nil { return nil, fmt.Errorf("migrationsversion lesen: %w", err) } applied[v] = true } return applied, rows.Err() } // Up wendet alle noch nicht angewendeten Migrationen in Reihenfolge an // (Akzeptanzkriterium 2: vorwaerts ausfuehrbar) — bereits angewendete // werden uebersprungen, damit Up auf einer leeren UND auf einer bestehenden // DB funktioniert (Pruefung 1). func Up(ctx context.Context, pool *pgxpool.Pool, migrations []Migration) (applied []string, err error) { if err := ensureTrackingTable(ctx, pool); err != nil { return nil, err } already, err := appliedVersions(ctx, pool) if err != nil { return nil, err } for _, m := range migrations { if already[m.Version] { continue } tx, err := pool.Begin(ctx) if err != nil { return applied, fmt.Errorf("transaktion fuer %q starten: %w", m.Version, err) } if _, err := tx.Exec(ctx, m.UpSQL); err != nil { _ = tx.Rollback(ctx) return applied, fmt.Errorf("migration %q anwenden: %w", m.Version, err) } if _, err := tx.Exec(ctx, `INSERT INTO schema_migrations (version) VALUES ($1)`, m.Version); err != nil { _ = tx.Rollback(ctx) return applied, fmt.Errorf("migration %q als angewendet markieren: %w", m.Version, err) } if err := tx.Commit(ctx); err != nil { return applied, fmt.Errorf("migration %q committen: %w", m.Version, err) } applied = append(applied, m.Version) } return applied, nil } // DownOne macht die zuletzt angewendete Migration rueckgaengig // (Akzeptanzkriterium 2: rueckwaerts ausfuehrbar) und liefert deren Version, // oder "" falls keine Migration angewendet war. func DownOne(ctx context.Context, pool *pgxpool.Pool, migrations []Migration) (version string, err error) { if err := ensureTrackingTable(ctx, pool); err != nil { return "", err } already, err := appliedVersions(ctx, pool) if err != nil { return "", err } var last *Migration for i := len(migrations) - 1; i >= 0; i-- { if already[migrations[i].Version] { last = &migrations[i] break } } if last == nil { return "", nil } tx, err := pool.Begin(ctx) if err != nil { return "", fmt.Errorf("transaktion fuer rollback von %q starten: %w", last.Version, err) } if _, err := tx.Exec(ctx, last.DownSQL); err != nil { _ = tx.Rollback(ctx) return "", fmt.Errorf("migration %q zurueckrollen: %w", last.Version, err) } if _, err := tx.Exec(ctx, `DELETE FROM schema_migrations WHERE version = $1`, last.Version); err != nil { _ = tx.Rollback(ctx) return "", fmt.Errorf("migration %q aus schema_migrations entfernen: %w", last.Version, err) } if err := tx.Commit(ctx); err != nil { return "", fmt.Errorf("rollback von %q committen: %w", last.Version, err) } return last.Version, nil }