package disasterrecovery import ( "context" "encoding/json" "fmt" "os" "path/filepath" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) // ReadSnapshot liest den Sicherungssatz eines Repositorys. // // Ohne Datenbank: Nach einem Totalverlust will man zuerst sehen, was ueberhaupt // da ist, bevor man irgendetwas aufsetzt. func ReadSnapshot(repositoryRootPath string) (*Snapshot, error) { snapshotPath := filepath.Join(repositoryRootPath, filepath.FromSlash(SnapshotDirectory), SnapshotFileName) encodedSnapshot, readError := os.ReadFile(snapshotPath) if readError != nil { if os.IsNotExist(readError) { return nil, fmt.Errorf("das repository enthaelt keinen sicherungssatz der "+ "konfiguration (%s). die wiederherstellungspunkte lassen sich trotzdem "+ "nutzen — die anlage muss dann von hand eingerichtet werden", SnapshotDirectory) } return nil, fmt.Errorf("der sicherungssatz liess sich nicht lesen: %w", readError) } var snapshot Snapshot if decodeError := json.Unmarshal(encodedSnapshot, &snapshot); decodeError != nil { return nil, fmt.Errorf("der sicherungssatz ist beschaedigt: %w", decodeError) } if validationError := snapshot.Validate(); validationError != nil { return nil, validationError } return &snapshot, nil } // ImportResult beschreibt das Ergebnis einer Wiederherstellung. type ImportResult struct { // RepositoriesRestored ist die Zahl wiederhergestellter Ablagen. RepositoriesRestored int // RetentionPoliciesRestored ist die Zahl wiederhergestellter Regeln. RetentionPoliciesRestored int // JobsRestored ist die Zahl wiederhergestellter Auftraege. JobsRestored int // JobSourcesRestored ist die Zahl wiederhergestellter Quellen. JobSourcesRestored int // MaintenanceWindowsRestored ist die Zahl wiederhergestellter Fenster. MaintenanceWindowsRestored int // NotificationChannelsRestored ist die Zahl angelegter Kanaele. NotificationChannelsRestored int // UsersRestored ist die Zahl wiederhergestellter Konten. UsersRestored int // SettingsRestored ist die Zahl wiederhergestellter Einstellungen. SettingsRestored int // SkippedExisting benennt Objekte, die bereits vorhanden waren. // // Sie werden **nicht** ueberschrieben: Ein Wiederherstellen in eine // laufende Anlage darf deren Zustand nicht stillschweigend ersetzen. SkippedExisting []string // ManualStepsRequired benennt, was der Betreiber noch zu tun hat. ManualStepsRequired []string } // Importer spielt einen Sicherungssatz in die Datenbank ein. type Importer struct { // connectionPool ist der Datenbankpool der Control Plane. connectionPool *pgxpool.Pool } // NewImporter erzeugt den Wiederhersteller. func NewImporter(connectionPool *pgxpool.Pool) *Importer { return &Importer{connectionPool: connectionPool} } // ImportSnapshot spielt einen Sicherungssatz ein. // // Alles oder nichts: Der gesamte Vorgang laeuft in **einer** Transaktion. Ein // Abbruch auf halber Strecke hinterliesse Auftraege ohne Repositories und // Quellen ohne Auftraege — eine Anlage, die arbeitsfaehig aussieht und beim // ersten Lauf scheitert. // // Vorhandene Objekte werden uebersprungen, nicht ueberschrieben. Der Befehl // laesst sich damit gefahrlos zweimal ausfuehren, und er kann eine laufende // Anlage nicht beschaedigen, wenn ihn jemand versehentlich dort absetzt. func (importer *Importer) ImportSnapshot(importContext context.Context, snapshot *Snapshot) (*ImportResult, error) { if validationError := snapshot.Validate(); validationError != nil { return nil, validationError } currentSchemaVersion, versionError := importer.readSchemaVersion(importContext) if versionError != nil { return nil, versionError } // Ein Sicherungssatz aus einem hoeheren Schemastand verweist auf Spalten, // die es hier nicht gibt. Der Fehler faellt sonst erst mitten im Einspielen // auf — nach der Haelfte der Auftraege. if snapshot.SchemaVersion > currentSchemaVersion { return nil, fmt.Errorf("der sicherungssatz stammt aus schemastand %d, die datenbank "+ "steht auf %d. bringen sie das schema mit 'syncova-migrate up' auf den "+ "gleichen stand", snapshot.SchemaVersion, currentSchemaVersion) } importResult := &ImportResult{ SkippedExisting: make([]string, 0, 8), ManualStepsRequired: snapshot.OmittedForSecurity, } databaseTransaction, beginError := importer.connectionPool.Begin(importContext) if beginError != nil { return nil, fmt.Errorf("die transaktion liess sich nicht beginnen: %w", beginError) } defer func() { _ = databaseTransaction.Rollback(importContext) }() // Die Reihenfolge folgt den Fremdschluesseln. importSteps := []struct { description string apply func(context.Context, pgx.Tx, *Snapshot, *ImportResult) error }{ {"repositories", importRepositories}, {"aufbewahrungsregeln", importRetentionPolicies}, {"auftraege", importJobs}, {"wartungsfenster", importMaintenanceWindows}, {"benachrichtigungswege", importNotificationChannels}, {"konten", importUsers}, {"einstellungen", importSettings}, } for _, importStep := range importSteps { if applyError := importStep.apply(importContext, databaseTransaction, snapshot, importResult); applyError != nil { return nil, fmt.Errorf("die %s liessen sich nicht wiederherstellen: %w", importStep.description, applyError) } } if commitError := databaseTransaction.Commit(importContext); commitError != nil { return nil, fmt.Errorf("die wiederherstellung liess sich nicht abschliessen: %w", commitError) } return importResult, nil } // readSchemaVersion liest den Migrationsstand der Datenbank. func (importer *Importer) readSchemaVersion(readContext context.Context) (int, error) { var schemaVersion int scanError := importer.connectionPool.QueryRow(readContext, `SELECT version FROM schema_migrations LIMIT 1`).Scan(&schemaVersion) if scanError != nil { return 0, fmt.Errorf("der schemastand konnte nicht gelesen werden — ist das schema "+ "angelegt? ('syncova-migrate up'): %w", scanError) } return schemaVersion, nil } // importRepositories stellt die Ablagen wieder her. // // Der Zustand wird auf 'unavailable' gesetzt, nicht auf den gesicherten: Nach // einem Totalverlust steht nicht fest, ob die Ablage am alten Pfad ueberhaupt // erreichbar ist. Ein Repository, das als 'active' zurueckkommt und nicht da // ist, laesst den naechsten Sicherungslauf ins Leere greifen — und die // Uebersicht meldete alles in Ordnung. func importRepositories(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertStatement = ` INSERT INTO repositories (id, name, repository_type, location, repository_uuid, status, hardened, capacity_bytes, retention_seconds, minimum_retention_seconds) VALUES ($1::uuid, $2, $3, $4, NULLIF($5, ''), 'unavailable', $6, $7, $8, $9) ON CONFLICT (id) DO NOTHING` for _, record := range snapshot.Repositories { commandTag, executeError := databaseTransaction.Exec(importContext, insertStatement, record.ID, record.Name, record.RepositoryType, record.Location, record.RepositoryUUID, record.Hardened, record.CapacityBytes, record.RetentionSeconds, record.MinimumRetentionSeconds) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Repository "+record.Name) continue } importResult.RepositoriesRestored++ } return nil } // importRetentionPolicies stellt die Aufbewahrungsregeln wieder her. func importRetentionPolicies(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertStatement = ` INSERT INTO retention_policies (id, name, rules, keep_within_seconds, keep_last, keep_daily, keep_weekly, keep_monthly, keep_yearly, time_zone) VALUES ($1::uuid, $2, COALESCE($3::jsonb, '{}'::jsonb), $4, $5, $6, $7, $8, $9, NULLIF($10, '')) ON CONFLICT (id) DO NOTHING` for _, record := range snapshot.RetentionPolicies { commandTag, executeError := databaseTransaction.Exec(importContext, insertStatement, record.ID, record.Name, record.Rules, record.KeepWithinSeconds, record.KeepLast, record.KeepDaily, record.KeepWeekly, record.KeepMonthly, record.KeepYearly, record.TimeZone) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Aufbewahrungsregel "+record.Name) continue } importResult.RetentionPoliciesRestored++ } return nil } // importJobs stellt die Auftraege samt Quellen wieder her. // // Jeder Auftrag kommt **angehalten** zurueck. Der Grund ist der wichtigste // Sicherheitsgedanke dieses ganzen Pakets: Nach einem Totalverlust weiss niemand, // ob die Quellen noch existieren, ob die Repositories erreichbar sind und ob die // wiederhergestellten Daten die richtigen sind. Ein Zeitplan, der um zwei Uhr // nachts von selbst anlaeuft, koennte auf ein halb wiederhergestelltes System // schreiben — und der Betreiber erfuehre davon am naechsten Morgen. func importJobs(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertJobStatement = ` INSERT INTO backup_jobs (id, name, description, status, priority, schedule_type, schedule_config, repository_id, retention_policy_id, rpo_seconds, rto_seconds, bandwidth_limit_bps, max_concurrency, retry_policy, paused_at) VALUES ($1::uuid, $2, NULLIF($3, ''), 'paused', $4, $5, COALESCE($6::jsonb, '{}'::jsonb), $7::uuid, NULLIF($8, '')::uuid, $9, $10, $11, $12, COALESCE($13::jsonb, '{}'::jsonb), now()) ON CONFLICT (id) DO NOTHING` // COALESCE auf die Musterlisten: Eine Quelle ohne Ein- oder Ausschluesse // traegt in Go ein nil, und das schreibt pgx als NULL. Die Spalten sind // NOT NULL mit Vorgabewert '{}' — ohne COALESCE bricht die Wiederherstellung // genau bei der haeufigsten aller Quellen ab, naemlich der ohne Filter. // Im Nachweis aufgefallen. const insertSourceStatement = ` INSERT INTO backup_job_sources (id, job_id, source_type, source_id, source_name, include_patterns, exclude_patterns) VALUES ($1::uuid, $2::uuid, $3, $4, $5, COALESCE($6::jsonb, '[]'::jsonb), COALESCE($7::jsonb, '[]'::jsonb)) ON CONFLICT (id) DO NOTHING` for _, record := range snapshot.Jobs { commandTag, executeError := databaseTransaction.Exec(importContext, insertJobStatement, record.ID, record.Name, record.Description, record.Priority, record.ScheduleType, record.ScheduleConfig, record.RepositoryID, record.RetentionPolicyID, record.RPOSeconds, record.RTOSeconds, record.BandwidthLimitBPS, record.MaxConcurrency, record.RetryPolicy) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Auftrag "+record.Name) continue } importResult.JobsRestored++ for _, sourceRecord := range record.Sources { includePatternsJSON, includeError := encodePatternList(sourceRecord.IncludePatterns) if includeError != nil { return includeError } excludePatternsJSON, excludeError := encodePatternList(sourceRecord.ExcludePatterns) if excludeError != nil { return excludeError } sourceTag, sourceError := databaseTransaction.Exec(importContext, insertSourceStatement, sourceRecord.ID, record.ID, sourceRecord.SourceType, sourceRecord.SourceID, sourceRecord.SourceName, includePatternsJSON, excludePatternsJSON) if sourceError != nil { return sourceError } if sourceTag.RowsAffected() > 0 { importResult.JobSourcesRestored++ } } } return nil } // importMaintenanceWindows stellt die Wartungsfenster wieder her. func importMaintenanceWindows(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertWindowStatement = ` INSERT INTO maintenance_windows (id, name, window_kind, starts_at, ends_at, recurrence, enabled) VALUES ($1::uuid, $2, $3, $4, $5, $6::jsonb, $7) ON CONFLICT (id) DO NOTHING` const insertLinkStatement = ` INSERT INTO maintenance_window_jobs (window_id, job_id) VALUES ($1::uuid, $2::uuid) ON CONFLICT DO NOTHING` for _, record := range snapshot.MaintenanceWindows { commandTag, executeError := databaseTransaction.Exec(importContext, insertWindowStatement, record.ID, record.Name, record.WindowKind, record.StartsAt, record.EndsAt, record.Recurrence, record.Enabled) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Wartungsfenster "+record.Name) continue } importResult.MaintenanceWindowsRestored++ for _, jobIdentifier := range record.JobIDs { if _, linkError := databaseTransaction.Exec(importContext, insertLinkStatement, record.ID, jobIdentifier); linkError != nil { return linkError } } } return nil } // importNotificationChannels legt die Kanaele abgeschaltet an. // // Ohne Konfiguration kann ein Kanal nichts zustellen. Ihn eingeschaltet // anzulegen erzeugte den gefaehrlichsten aller Zustaende: eine Anlage, die // meldet — nur eben nirgendwohin. func importNotificationChannels(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertStatement = ` INSERT INTO notification_channels (id, name, channel_type, enabled, minimum_severity, configuration) VALUES ($1::uuid, $2, $3, false, $4, '{}'::jsonb) ON CONFLICT (id) DO NOTHING` for _, record := range snapshot.NotificationChannels { commandTag, executeError := databaseTransaction.Exec(importContext, insertStatement, record.ID, record.Name, record.ChannelType, record.MinimumSeverity) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Benachrichtigungsweg "+record.Name) continue } importResult.NotificationChannelsRestored++ } return nil } // importUsers stellt die Konten ohne Passwort wieder her. // // Sie kommen als 'disabled' zurueck. Ein Konto ohne Passwort, das als aktiv // gefuehrt wird, waere ein Konto, dessen Anmeldung nur deshalb scheitert, weil // kein Passwort passt — und wer das Passwort setzt, hat einen Zugang. Der Weg // ueber 'disabled' zwingt zu einer bewussten Freischaltung je Konto. func importUsers(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { // Ein Platzhalter, der als Argon2id-Hash nicht gueltig ist und deshalb auf // keine Eingabe passt. Die Spalte ist NOT NULL; ein leerer Wert waere // gefaehrlicher, weil eine kuenftige Pruefung ihn als „kein Passwort noetig" // auslegen koennte. const unusableHashPlaceholder = "$argon2id$disaster-recovery$no-password-set" const insertUserStatement = ` INSERT INTO users (id, username, email, password_hash, status, mfa_enabled) VALUES ($1::uuid, $2, NULLIF($3, ''), $4, 'disabled', false) ON CONFLICT (id) DO NOTHING` const insertRoleStatement = ` INSERT INTO user_roles (user_id, role_id) SELECT $1::uuid, r.id FROM roles r WHERE r.name = $2 ON CONFLICT DO NOTHING` for _, record := range snapshot.Users { commandTag, executeError := databaseTransaction.Exec(importContext, insertUserStatement, record.ID, record.Username, record.Email, unusableHashPlaceholder) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Konto "+record.Username) continue } importResult.UsersRestored++ for _, roleName := range record.Roles { if _, roleError := databaseTransaction.Exec(importContext, insertRoleStatement, record.ID, roleName); roleError != nil { return roleError } } } return nil } // importSettings stellt die Systemeinstellungen wieder her. func importSettings(importContext context.Context, databaseTransaction pgx.Tx, snapshot *Snapshot, importResult *ImportResult) error { const insertStatement = ` INSERT INTO system_settings (key, value_json) VALUES ($1, COALESCE($2::jsonb, '{}'::jsonb)) ON CONFLICT (key) DO NOTHING` for _, record := range snapshot.Settings { commandTag, executeError := databaseTransaction.Exec(importContext, insertStatement, record.Key, record.Value) if executeError != nil { return executeError } if commandTag.RowsAffected() == 0 { importResult.SkippedExisting = append(importResult.SkippedExisting, "Einstellung "+record.Key) continue } importResult.SettingsRestored++ } return nil } // encodePatternList kodiert eine Musterliste als JSON. // // Die Spalten include_patterns und exclude_patterns sind jsonb, nicht text[]. // Eine Go-Zeichenkettenliste unmittelbar zu uebergeben ergaebe ein // PostgreSQL-Array und damit einen Typfehler — im Nachweis aufgefallen, und // zwar erst beim zweiten Anlauf: Der erste Fehler war ein NULL, der zweite der // falsche Typ. func encodePatternList(patterns []string) ([]byte, error) { if len(patterns) == 0 { return []byte("[]"), nil } encodedPatterns, encodeError := json.Marshal(patterns) if encodeError != nil { return nil, fmt.Errorf("die musterliste liess sich nicht kodieren: %w", encodeError) } return encodedPatterns, nil }