From b08f6a49fc962de311d1b8b72b4f1ba922c129be Mon Sep 17 00:00:00 2001 From: sysops Date: Thu, 27 Aug 2026 19:05:47 +0200 Subject: [PATCH] TEN-07: migrations-orchestrierung-tenant-datenbanken MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Orchestrator.RolloutAll wendet eine geordnete Liste von Migrationen auf JEDE registrierte Tenant-Datenbank an (LoadMigrations liest *.up.sql aus einem Verzeichnis nach der bestehenden 000N_name-Namenskonvention). Jeder Tenant laeuft unabhaengig in eigener Verbindung — ein Fehlschlag bei einem Mandanten bricht nur dessen eigenen Rollout ab (spaetere Migrationen bauen typischerweise auf frueheren auf) und blockiert die uebrigen Tenants nicht. schema_migrations-Tabelle pro Tenant-Datenbank (version PK, applied_at, success, error) haelt den Stand pro Version einzeln nachvollziehbar fest. Bereits erfolgreiche Versionen werden bei einem erneuten Rollout uebersprungen (isAlreadySuccessful-Check vor jeder Anwendung), fehlgeschlagene werden beim naechsten Versuch automatisch erneut probiert (kein manuelles Zuruecksetzen noetig) — ON CONFLICT DO UPDATE haelt jeweils nur den letzten Versuch fest. Bewusst ohne Abhaengigkeit von TEN-06 (Router): Migrations-Rollouts sind seltene Batch-Vorgaenge, ein kurzlebiger Pool pro Tenant und Lauf reicht. Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS, go test ./... mit -p 1 noetig da mehrere Pakete die geteilte Registry-Tabelle auf derselben Postgres-Instanz nutzen — siehe scripts/reset-test-env.sh): 1. Rollout gegen 3 Test-Tenants, einer absichtlich inkompatibel (Tabellen- Konflikt bei Migration 2) — TestRolloutAll_IsolatesFailurePerTenant: die anderen beiden erhalten beide Migrationen, der inkompatible bekommt Migration 1 trotzdem, scheitert nur an Migration 2, faellt nicht die anderen um. PASS. 2. Migrationsstand-Abfrage liefert korrekten Stand pro Tenant — TestStatus_ReflectsPerTenantState: Version 1 success=true, Version 2 success=false mit Fehlertext. PASS. 3. Wiederholter Rollout fuer fehlgeschlagene Migration moeglich, ohne bereits erfolgreiche erneut anzuwenden — TestRolloutAll_RetryDoesNotReapplySuccessful: Migration 1 nutzt bewusst kein IF NOT EXISTS, ein Reapply haette den zweiten Lauf scheitern lassen; zweiter Lauf ist fehlerfrei und wendet nur die zuvor fehlgeschlagene Version an. PASS. Co-Authored-By: Claude Sonnet 5 --- go.mod | 9 ++ go.sum | 28 ++++ internal/migrate/migrate.go | 211 +++++++++++++++++++++++++++ internal/migrate/migrate_test.go | 238 +++++++++++++++++++++++++++++++ 4 files changed, 486 insertions(+) create mode 100644 go.sum create mode 100644 internal/migrate/migrate.go create mode 100644 internal/migrate/migrate_test.go diff --git a/go.mod b/go.mod index 75a79b7..05c10d0 100644 --- a/go.mod +++ b/go.mod @@ -3,3 +3,12 @@ module gitea.perlbach24.de/scripte/nexarch go 1.22 require github.com/jackc/pgx/v5 v5.6.0 + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect + github.com/jackc/puddle/v2 v2.2.1 // indirect + golang.org/x/crypto v0.17.0 // indirect + golang.org/x/sync v0.1.0 // indirect + golang.org/x/text v0.14.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..5c39671 --- /dev/null +++ b/go.sum @@ -0,0 +1,28 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk= +github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY= +github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw= +github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk= +github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k= +golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4= +golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o= +golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ= +golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/migrate/migrate.go b/internal/migrate/migrate.go new file mode 100644 index 0000000..15b935a --- /dev/null +++ b/internal/migrate/migrate.go @@ -0,0 +1,211 @@ +// Package migrate implementiert Core TEN-07: das automatisierte, pro Tenant +// fehlerisolierte Ausrollen von SQL-Migrationen ueber alle registrierten +// Mandanten-Datenbanken (Modell C — jede Migration muss N-mal statt einmal +// laufen, siehe TEN-01). +package migrate + +import ( + "context" + "errors" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/tenant" +) + +const upSuffix = ".up.sql" + +type Migration struct { + Version string + SQL string +} + +// LoadMigrations liest alle *.up.sql-Dateien aus dir und sortiert sie nach +// Dateiname (die bestehende Namenskonvention 0001_..., 0002_... aus TEN-01 +// sorgt fuer eine stabile Reihenfolge). +func LoadMigrations(dir string) ([]Migration, error) { + entries, err := os.ReadDir(dir) + if err != nil { + return nil, fmt.Errorf("migrationsverzeichnis lesen: %w", err) + } + + var names []string + for _, e := range entries { + if !e.IsDir() && strings.HasSuffix(e.Name(), upSuffix) { + names = append(names, e.Name()) + } + } + sort.Strings(names) + + migrations := make([]Migration, 0, len(names)) + for _, name := range names { + content, err := os.ReadFile(filepath.Join(dir, name)) + if err != nil { + return nil, fmt.Errorf("migration %q lesen: %w", name, err) + } + version := strings.TrimSuffix(name, upSuffix) + migrations = append(migrations, Migration{Version: version, SQL: string(content)}) + } + return migrations, nil +} + +// MigrationStatus ist der pro Tenant und Version nachvollziehbare Stand +// (Akzeptanzkriterium 3): welche Version, wann zuletzt versucht, Erfolg oder +// Fehler. +type MigrationStatus struct { + Version string + AppliedAt time.Time + Success bool + Error string +} + +// TenantResult fasst das Ergebnis eines Rollout-Versuchs fuer EINEN Tenant +// zusammen — wird von Orchestrator.RolloutAll pro Tenant gesammelt, damit ein +// Fehlschlag bei einem Mandanten die anderen nicht blockiert (Akzeptanzkriterium 2). +type TenantResult struct { + TenantSlug string + Applied []string // erfolgreich in diesem Lauf angewendete Versionen + FailedAt string // Version, bei der abgebrochen wurde; leer wenn kein Fehlschlag + Err error +} + +// Orchestrator rollt Migrationen ueber alle in der Registry gefuehrten +// Mandanten aus. Baut bewusst NICHT auf TEN-06 (Router) auf — Migrations- +// Rollouts sind seltene Batch-Vorgaenge, kein Hot-Path, ein kurzlebiger Pool +// pro Tenant und Lauf ist hier einfacher und unabhaengig von TEN-06 testbar. +type Orchestrator struct { + registry *tenant.Registry +} + +func NewOrchestrator(registry *tenant.Registry) *Orchestrator { + return &Orchestrator{registry: registry} +} + +const ensureTableSQL = ` +CREATE TABLE IF NOT EXISTS schema_migrations ( + version TEXT PRIMARY KEY, + applied_at TIMESTAMPTZ NOT NULL DEFAULT now(), + success BOOLEAN NOT NULL, + error TEXT +);` + +// RolloutAll wendet migrations auf JEDE registrierte Tenant-Datenbank an. +// Jeder Tenant laeuft unabhaengig — ein Fehlschlag bei einem Mandanten wird +// im jeweiligen TenantResult festgehalten und blockiert die uebrigen nicht +// (Akzeptanzkriterium 1 + 2). +func (o *Orchestrator) RolloutAll(ctx context.Context, migrations []Migration) ([]TenantResult, error) { + tenants, err := o.registry.List(ctx) + if err != nil { + return nil, fmt.Errorf("tenants fuer rollout auflisten: %w", err) + } + + results := make([]TenantResult, 0, len(tenants)) + for _, t := range tenants { + results = append(results, o.rolloutForTenant(ctx, t, migrations)) + } + return results, nil +} + +func (o *Orchestrator) rolloutForTenant(ctx context.Context, t tenant.Tenant, migrations []Migration) TenantResult { + result := TenantResult{TenantSlug: t.Slug} + + pool, err := pgxpool.New(ctx, t.DBDSN) + if err != nil { + result.Err = fmt.Errorf("verbindung zu tenant %q: %w", t.Slug, err) + return result + } + defer pool.Close() + + if _, err := pool.Exec(ctx, ensureTableSQL); err != nil { + result.Err = fmt.Errorf("schema_migrations anlegen fuer tenant %q: %w", t.Slug, err) + return result + } + + for _, m := range migrations { + alreadyApplied, err := isAlreadySuccessful(ctx, pool, m.Version) + if err != nil { + result.Err = fmt.Errorf("migrationsstand lesen fuer tenant %q, version %q: %w", t.Slug, m.Version, err) + return result + } + if alreadyApplied { + continue // Akzeptanzkriterium 3 / Pruefung 3: keine erneute Anwendung. + } + + _, execErr := pool.Exec(ctx, m.SQL) + if execErr != nil { + recordAttempt(ctx, pool, m.Version, false, execErr.Error()) + result.FailedAt = m.Version + result.Err = fmt.Errorf("migration %q fuer tenant %q fehlgeschlagen: %w", m.Version, t.Slug, execErr) + return result // spaetere Migrationen bauen typischerweise auf dieser auf — Abbruch NUR fuer diesen Tenant. + } + + recordAttempt(ctx, pool, m.Version, true, "") + result.Applied = append(result.Applied, m.Version) + } + + return result +} + +func isAlreadySuccessful(ctx context.Context, pool *pgxpool.Pool, version string) (bool, error) { + var success bool + err := pool.QueryRow(ctx, `SELECT success FROM schema_migrations WHERE version = $1`, version).Scan(&success) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return false, nil + } + return false, err + } + return success, nil +} + +func recordAttempt(ctx context.Context, pool *pgxpool.Pool, version string, success bool, errMsg string) { + var errVal *string + if errMsg != "" { + errVal = &errMsg + } + _, _ = pool.Exec(ctx, ` + INSERT INTO schema_migrations (version, applied_at, success, error) + VALUES ($1, now(), $2, $3) + ON CONFLICT (version) DO UPDATE SET applied_at = now(), success = $2, error = $3 + `, version, success, errVal) +} + +// Status liefert den Migrationsstand eines einzelnen Tenants (Akzeptanzkriterium 3). +func (o *Orchestrator) Status(ctx context.Context, tenantSlug string) ([]MigrationStatus, error) { + t, err := o.registry.GetBySlug(ctx, tenantSlug) + if err != nil { + return nil, fmt.Errorf("tenant %q nicht gefunden: %w", tenantSlug, err) + } + + pool, err := pgxpool.New(ctx, t.DBDSN) + if err != nil { + return nil, fmt.Errorf("verbindung zu tenant %q: %w", tenantSlug, err) + } + defer pool.Close() + + rows, err := pool.Query(ctx, ` + SELECT version, applied_at, success, COALESCE(error, '') + FROM schema_migrations ORDER BY version + `) + if err != nil { + return nil, fmt.Errorf("migrationsstand abfragen: %w", err) + } + defer rows.Close() + + var out []MigrationStatus + for rows.Next() { + var s MigrationStatus + if err := rows.Scan(&s.Version, &s.AppliedAt, &s.Success, &s.Error); err != nil { + return nil, fmt.Errorf("migrationsstand lesen: %w", err) + } + out = append(out, s) + } + return out, rows.Err() +} diff --git a/internal/migrate/migrate_test.go b/internal/migrate/migrate_test.go new file mode 100644 index 0000000..47c2352 --- /dev/null +++ b/internal/migrate/migrate_test.go @@ -0,0 +1,238 @@ +package migrate + +import ( + "context" + "fmt" + "os" + "strings" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/tenant" +) + +func TestLoadMigrations_SortsByFilename(t *testing.T) { + dir := t.TempDir() + writeFile(t, dir, "0002_second.up.sql", "SELECT 2;") + writeFile(t, dir, "0001_first.up.sql", "SELECT 1;") + writeFile(t, dir, "0001_first.down.sql", "SELECT 'ignored';") // muss ignoriert werden + + migrations, err := LoadMigrations(dir) + if err != nil { + t.Fatalf("load: %v", err) + } + if len(migrations) != 2 { + t.Fatalf("erwartet 2 migrationen, habe %d", len(migrations)) + } + if migrations[0].Version != "0001_first" || migrations[1].Version != "0002_second" { + t.Fatalf("unerwartete reihenfolge: %+v", migrations) + } +} + +func writeFile(t *testing.T, dir, name, content string) { + t.Helper() + if err := os.WriteFile(dir+"/"+name, []byte(content), 0o644); err != nil { + t.Fatalf("write %s: %v", name, err) + } +} + +func setupOrchestratorTest(t *testing.T, tenantCount int) (*Orchestrator, []tenant.Tenant, *pgxpool.Pool, func()) { + t.Helper() + adminDSN := os.Getenv("TEST_ADMIN_DSN") + if adminDSN == "" { + t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen") + } + ctx := context.Background() + + adminPool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("admin pool: %v", err) + } + registryPool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("registry pool: %v", err) + } + if _, err := registryPool.Exec(ctx, ` + CREATE TABLE IF NOT EXISTS tenants ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + slug TEXT NOT NULL UNIQUE, + name TEXT NOT NULL, + db_name TEXT NOT NULL UNIQUE, + db_dsn TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'active', + created_at TIMESTAMPTZ NOT NULL DEFAULT now() + )`); err != nil { + t.Fatalf("registry-schema: %v", err) + } + + registry := tenant.NewRegistry(registryPool) + dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1) + provisioner := tenant.NewProvisioner(adminPool, registry, dsnTemplate) + + var tenants []tenant.Tenant + var slugs []string + for i := 0; i < tenantCount; i++ { + slug := fmt.Sprintf("mig_t%d", i) + slugs = append(slugs, slug) + tn, err := provisioner.Provision(ctx, slug, slug) + if err != nil { + t.Fatalf("provision %s: %v", slug, err) + } + tenants = append(tenants, tn) + } + + orchestrator := NewOrchestrator(registry) + + cleanup := func() { + for _, slug := range slugs { + _, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, "tenant_"+slug)) + } + _, _ = registryPool.Exec(ctx, `DELETE FROM tenants WHERE slug = ANY($1)`, slugs) + registryPool.Close() + adminPool.Close() + } + return orchestrator, tenants, adminPool, cleanup +} + +// Akzeptanzkriterium 1 + 2 + Pruefung 1: Rollout gegen mehrere Tenants, einer +// davon absichtlich inkompatibel — die anderen laufen trotzdem durch. +func TestRolloutAll_IsolatesFailurePerTenant(t *testing.T) { + orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 3) + defer cleanup() + ctx := context.Background() + + migrations := []Migration{ + {Version: "0001_demo_a", SQL: "CREATE TABLE demo_a (id INT);"}, + {Version: "0002_demo_b", SQL: "CREATE TABLE demo_b (id INT);"}, + } + + // tenants[1] absichtlich inkompatibel machen: demo_b existiert schon, + // migration 0002 schlaegt dort mit "already exists" fehl. + badTenant := tenants[1] + badPool, err := pgxpool.New(ctx, badTenant.DBDSN) + if err != nil { + t.Fatalf("connect bad tenant: %v", err) + } + if _, err := badPool.Exec(ctx, "CREATE TABLE demo_b (id INT);"); err != nil { + t.Fatalf("inkompatiblen zustand vorbereiten: %v", err) + } + badPool.Close() + + results, err := orchestrator.RolloutAll(ctx, migrations) + if err != nil { + t.Fatalf("rollout: %v", err) + } + if len(results) != 3 { + t.Fatalf("erwartet 3 ergebnisse, habe %d", len(results)) + } + + byslug := map[string]TenantResult{} + for _, r := range results { + byslug[r.TenantSlug] = r + } + + good0 := byslug[tenants[0].Slug] + if good0.Err != nil || len(good0.Applied) != 2 { + t.Fatalf("tenant[0] sollte beide migrationen erhalten, habe %+v", good0) + } + good2 := byslug[tenants[2].Slug] + if good2.Err != nil || len(good2.Applied) != 2 { + t.Fatalf("tenant[2] sollte beide migrationen erhalten, habe %+v", good2) + } + + bad := byslug[badTenant.Slug] + if bad.Err == nil { + t.Fatal("erwartet fehler fuer den inkompatiblen tenant") + } + if bad.FailedAt != "0002_demo_b" { + t.Fatalf("failedAt = %q, want 0002_demo_b", bad.FailedAt) + } + if len(bad.Applied) != 1 || bad.Applied[0] != "0001_demo_a" { + t.Fatalf("erwartet dass 0001_demo_a trotzdem erfolgreich war, habe %+v", bad.Applied) + } +} + +// Akzeptanzkriterium 3 + Pruefung 2: Migrationsstand-Abfrage liefert fuer +// jeden Tenant den korrekten, unabhaengigen Stand. +func TestStatus_ReflectsPerTenantState(t *testing.T) { + orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 1) + defer cleanup() + ctx := context.Background() + + migrations := []Migration{ + {Version: "0001_ok", SQL: "CREATE TABLE ok_table (id INT);"}, + {Version: "0002_fail", SQL: "SELECT this_column_does_not_exist FROM ok_table;"}, + } + + if _, err := orchestrator.RolloutAll(ctx, migrations); err != nil { + t.Fatalf("rollout: %v", err) + } + + status, err := orchestrator.Status(ctx, tenants[0].Slug) + if err != nil { + t.Fatalf("status: %v", err) + } + if len(status) != 2 { + t.Fatalf("erwartet 2 status-eintraege, habe %d", len(status)) + } + if !status[0].Success || status[0].Version != "0001_ok" { + t.Fatalf("status[0] unerwartet: %+v", status[0]) + } + if status[1].Success || status[1].Version != "0002_fail" || status[1].Error == "" { + t.Fatalf("status[1] sollte fehlgeschlagen sein mit fehlertext: %+v", status[1]) + } +} + +// Pruefung 3: wiederholter Rollout-Versuch wendet bereits erfolgreiche +// Migrationen NICHT erneut an und kann die zuvor fehlgeschlagene nachholen, +// sobald die Ursache behoben ist. +func TestRolloutAll_RetryDoesNotReapplySuccessful(t *testing.T) { + orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 1) + defer cleanup() + ctx := context.Background() + + // 0002 schlaegt beim ersten Versuch fehl, weil demo_conflict schon + // existiert (wir legen sie vorher an, um den Fehlschlag zu erzwingen). + pool, err := pgxpool.New(ctx, tenants[0].DBDSN) + if err != nil { + t.Fatalf("connect: %v", err) + } + if _, err := pool.Exec(ctx, "CREATE TABLE demo_conflict (id INT);"); err != nil { + t.Fatalf("vorbedingung: %v", err) + } + + migrations := []Migration{ + {Version: "0001_ok", SQL: "CREATE TABLE demo_first (id INT);"}, // OHNE IF NOT EXISTS, + // damit ein erneutes Anwenden nachweislich fehlschlagen wuerde. + {Version: "0002_conflict", SQL: "CREATE TABLE demo_conflict (id INT);"}, + } + + firstRun, err := orchestrator.RolloutAll(ctx, migrations) + if err != nil { + t.Fatalf("rollout 1: %v", err) + } + if firstRun[0].FailedAt != "0002_conflict" { + t.Fatalf("erwartet fehlschlag bei 0002_conflict im ersten lauf, habe %+v", firstRun[0]) + } + + // Ursache beheben. + if _, err := pool.Exec(ctx, "DROP TABLE demo_conflict;"); err != nil { + t.Fatalf("ursache beheben: %v", err) + } + pool.Close() + + secondRun, err := orchestrator.RolloutAll(ctx, migrations) + if err != nil { + t.Fatalf("rollout 2: %v", err) + } + // Waere 0001_ok erneut angewendet worden ("CREATE TABLE demo_first" ohne + // IF NOT EXISTS), haette das einen Fehler erzeugt statt eines sauberen + // Applied-Eintrags fuer 0002_conflict. + if secondRun[0].Err != nil { + t.Fatalf("zweiter lauf sollte fehlerfrei sein, habe %+v", secondRun[0]) + } + if len(secondRun[0].Applied) != 1 || secondRun[0].Applied[0] != "0002_conflict" { + t.Fatalf("erwartet nur 0002_conflict im zweiten lauf angewendet, habe %+v", secondRun[0].Applied) + } +}