package database import ( "database/sql" "errors" "fmt" "io/fs" "net/url" "strings" "github.com/golang-migrate/migrate/v4" migratepostgres "github.com/golang-migrate/migrate/v4/database/postgres" "github.com/golang-migrate/migrate/v4/source/iofs" // Der pgx-Treiber wird über database/sql angesprochen, weil golang-migrate // diese Schnittstelle erwartet. _ "github.com/jackc/pgx/v5/stdlib" ) // MigrationState beschreibt den Migrationsstand einer Datenbank. type MigrationState struct { // Version ist die zuletzt angewandte Migrationsversion. Version uint // IsDirty meldet, ob eine Migration abgebrochen ist und die Datenbank // in einem unklaren Zustand hinterlassen hat. IsDirty bool // HasAnyMigration meldet, ob überhaupt schon eine Migration lief. HasAnyMigration bool } // newMigrator baut eine golang-migrate-Instanz über die eingebetteten SQL-Dateien. // // Die Datenbankverbindung wird bewusst selbst geöffnet und als Instanz übergeben: // Überließe man golang-migrate das Öffnen, käme dessen eingebauter lib/pq-Treiber // zum Einsatz. Der versteht andere SSL-Modi als der pgx-Treiber des Verbindungspools, // womit dieselbe Konfiguration je nach Kommando funktionieren oder scheitern würde. func newMigrator(migrationFS fs.FS, connectionString string) (*migrate.Migrate, error) { migrationSource, sourceError := iofs.New(migrationFS, ".") if sourceError != nil { return nil, fmt.Errorf("die eingebetteten Migrationen konnten nicht gelesen werden: %w", sourceError) } // "pgx" ist der von pgx/v5/stdlib registrierte database/sql-Treiber. migrationDatabase, openError := sql.Open("pgx", connectionString) if openError != nil { return nil, fmt.Errorf("die Datenbankverbindung für den Migrationslauf konnte nicht geöffnet werden: %s", redactPasswordInText(connectionString, openError.Error())) } databaseDriver, driverError := migratepostgres.WithInstance(migrationDatabase, &migratepostgres.Config{}) if driverError != nil { _ = migrationDatabase.Close() return nil, fmt.Errorf("der Migrationslauf konnte nicht vorbereitet werden: %s", redactPasswordInText(connectionString, driverError.Error())) } migrator, migratorError := migrate.NewWithInstance("iofs", migrationSource, "postgres", databaseDriver) if migratorError != nil { _ = migrationDatabase.Close() // Die Meldung von golang-migrate kann die DSN samt Passwort enthalten. // Sie wird deshalb redigiert weitergereicht statt verworfen: ohne Ursache // wäre ein Verbindungsproblem praktisch nicht diagnostizierbar. return nil, fmt.Errorf("der Migrationslauf konnte nicht vorbereitet werden: %s", redactPasswordInText(connectionString, migratorError.Error())) } return migrator, nil } // redactPasswordInText ersetzt das Passwort der DSN in einem Fehlertext. // // Fremdbibliotheken nehmen die vollständige Verbindungszeichenkette gerne in // ihre Meldungen auf. Diese Funktion erlaubt es, solche Meldungen dennoch // weiterzureichen, ohne ein Geheimnis preiszugeben (PROMPT.md §12). func redactPasswordInText(connectionString string, originalText string) string { parsedConnectionURL, parseError := url.Parse(connectionString) if parseError != nil || parsedConnectionURL.User == nil { return originalText } databasePassword, hasPassword := parsedConnectionURL.User.Password() if !hasPassword || databasePassword == "" { return originalText } // Sowohl die Klartext- als auch die URL-kodierte Form ersetzen, da // Fehlermeldungen die DSN in beiden Varianten enthalten können. redactedText := strings.ReplaceAll(originalText, databasePassword, "***") return strings.ReplaceAll(redactedText, url.QueryEscape(databasePassword), "***") } // MigrateUp wendet alle ausstehenden Migrationen an. // // Die Funktion wird ausschließlich aus dem eigenen Migrationskommando // aufgerufen. Der Produktionsstart darf das Schema niemals stillschweigend // verändern (SYNCOVA_DATABASE.md §18). func MigrateUp(migrationFS fs.FS, connectionString string) (appliedMigrations bool, migrationError error) { migrator, migratorError := newMigrator(migrationFS, connectionString) if migratorError != nil { return false, migratorError } defer closeMigrator(migrator, &migrationError) if upError := migrator.Up(); upError != nil { // ErrNoChange bedeutet: das Schema ist bereits aktuell — kein Fehlerfall. if errors.Is(upError, migrate.ErrNoChange) { return false, nil } return false, fmt.Errorf("die Migration konnte nicht angewandt werden: %w", upError) } return true, nil } // MigrateDownOneStep nimmt genau eine Migration zurück. // // Der Schritt ist bewusst einzeln: ein versehentlicher Rücklauf auf Version 0 // würde die gesamte Control Plane löschen (PROMPT.md §141). func MigrateDownOneStep(migrationFS fs.FS, connectionString string) (migrationError error) { migrator, migratorError := newMigrator(migrationFS, connectionString) if migratorError != nil { return migratorError } defer closeMigrator(migrator, &migrationError) if stepError := migrator.Steps(-1); stepError != nil { if errors.Is(stepError, migrate.ErrNoChange) { return nil } return fmt.Errorf("die Migration konnte nicht zurückgenommen werden: %w", stepError) } return nil } // ForceVersion setzt den Migrationsstand ohne SQL auszuführen. // // Der einzige Zweck ist die Rettung nach einer abgebrochenen Migration. Bricht // ein Lauf mitten in der Datei ab — etwa an einem Constraint-Verstoß —, merkt // golang-migrate die Datenbank als „dirty" vor und verweigert jeden weiteren // Schritt. Ohne dieses Kommando gäbe es keinen Weg zurück außer von Hand in // der Datenbank; genau das soll ein Migrationswerkzeug ersparen. // // Die Funktion führt **kein** SQL aus. Sie behauptet lediglich einen Stand. // Wer sie aufruft, muss vorher sichergestellt haben, dass das Schema diesem // Stand auch entspricht — sonst arbeitet die Anwendung gegen ein Schema, das // sie für etwas anderes hält. func ForceVersion(migrationFS fs.FS, connectionString string, targetVersion int) (migrationError error) { migrator, migratorError := newMigrator(migrationFS, connectionString) if migratorError != nil { return migratorError } defer closeMigrator(migrator, &migrationError) if forceError := migrator.Force(targetVersion); forceError != nil { return fmt.Errorf("der Migrationsstand konnte nicht auf %d gesetzt werden: %w", targetVersion, forceError) } return nil } // CurrentState liest den Migrationsstand der Datenbank. func CurrentState(migrationFS fs.FS, connectionString string) (state MigrationState, migrationError error) { migrator, migratorError := newMigrator(migrationFS, connectionString) if migratorError != nil { return MigrationState{}, migratorError } defer closeMigrator(migrator, &migrationError) currentVersion, isDirty, versionError := migrator.Version() if versionError != nil { // ErrNilVersion bedeutet: die Datenbank ist noch unmigriert. if errors.Is(versionError, migrate.ErrNilVersion) { return MigrationState{HasAnyMigration: false}, nil } return MigrationState{}, fmt.Errorf("der Migrationsstand konnte nicht gelesen werden: %w", versionError) } return MigrationState{Version: currentVersion, IsDirty: isDirty, HasAnyMigration: true}, nil } // VerifySchemaIsUpToDate prüft, ob das Schema zur Programmversion passt. // // Der API-Dienst ruft die Funktion beim Start auf und verweigert den Dienst bei // Abweichung, statt gegen ein unpassendes Schema zu arbeiten. func VerifySchemaIsUpToDate(migrationFS fs.FS, connectionString string) error { migrationState, stateError := CurrentState(migrationFS, connectionString) if stateError != nil { return stateError } // Ein "dirty" Zustand heißt: eine Migration ist mittendrin abgebrochen. // Weiterarbeiten hieße, auf einem unbekannten Schema zu operieren. if migrationState.IsDirty { return fmt.Errorf("das Datenbankschema ist in einem unklaren Zustand (Version %d, abgebrochene Migration). "+ "Bitte den Migrationslauf prüfen und mit 'syncova-migrate status' den Zustand klären", migrationState.Version) } if !migrationState.HasAnyMigration { return errors.New("die Datenbank enthält noch kein Syncova-Schema. Bitte zuerst 'syncova-migrate up' ausführen") } expectedVersion, versionError := latestAvailableVersion(migrationFS) if versionError != nil { return versionError } if migrationState.Version != expectedVersion { return fmt.Errorf("das Datenbankschema hat Version %d, diese Programmversion erwartet %d. "+ "Bitte 'syncova-migrate up' ausführen", migrationState.Version, expectedVersion) } return nil } // latestAvailableVersion ermittelt die höchste mitgelieferte Migrationsversion. func latestAvailableVersion(migrationFS fs.FS) (uint, error) { migrationSource, sourceError := iofs.New(migrationFS, ".") if sourceError != nil { return 0, fmt.Errorf("die eingebetteten Migrationen konnten nicht gelesen werden: %w", sourceError) } defer func() { _ = migrationSource.Close() }() firstVersion, firstError := migrationSource.First() if firstError != nil { return 0, fmt.Errorf("es sind keine Migrationen eingebettet: %w", firstError) } // Von der ersten Version aus vorwärts laufen, bis keine weitere folgt. highestVersion := firstVersion for { nextVersion, nextError := migrationSource.Next(highestVersion) if nextError != nil { return highestVersion, nil } highestVersion = nextVersion } } // closeMigrator schließt Quelle und Datenbankverbindung des Migrators. // // Auftretende Fehler werden nur dann gemeldet, wenn die eigentliche Operation // erfolgreich war: sonst würde ein Aufräumfehler die echte Ursache verdecken. func closeMigrator(migrator *migrate.Migrate, operationError *error) { sourceCloseError, databaseCloseError := migrator.Close() if *operationError != nil { return } if sourceCloseError != nil { *operationError = fmt.Errorf("die Migrationsquelle konnte nicht geschlossen werden: %w", sourceCloseError) return } if databaseCloseError != nil { *operationError = fmt.Errorf("die Datenbankverbindung des Migrators konnte nicht geschlossen werden: %w", databaseCloseError) } }