// 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() }