package disasterrecovery import ( "context" "encoding/json" "fmt" "os" "path/filepath" "time" "github.com/jackc/pgx/v5/pgxpool" ) // Exporter erzeugt Sicherungssaetze der Control-Plane-Konfiguration. type Exporter struct { // connectionPool ist der Datenbankpool der Control Plane. connectionPool *pgxpool.Pool } // NewExporter erzeugt den Sicherungsersteller. func NewExporter(connectionPool *pgxpool.Pool) *Exporter { return &Exporter{connectionPool: connectionPool} } // BuildSnapshot liest die Konfiguration aus der Datenbank. func (exporter *Exporter) BuildSnapshot(buildContext context.Context, repositoryIdentifier string, createdBy string, productVersion string) (*Snapshot, error) { snapshot := &Snapshot{ FormatVersion: SnapshotFormatVersion, RepositoryID: repositoryIdentifier, CreatedAt: time.Now().UTC(), CreatedBy: createdBy, ProductVersion: productVersion, OmittedForSecurity: securityOmissions(), } schemaVersion, versionError := exporter.readSchemaVersion(buildContext) if versionError != nil { return nil, versionError } snapshot.SchemaVersion = schemaVersion // Die Reihenfolge folgt den Fremdschluesseln: Repositories und // Aufbewahrungsregeln zuerst, dann die Auftraege, die auf beide verweisen. // Beim Einspielen wird dieselbe Reihenfolge gebraucht. loadSteps := []struct { description string load func(context.Context, *Snapshot) error }{ {"repositories", exporter.loadRepositories}, {"aufbewahrungsregeln", exporter.loadRetentionPolicies}, {"auftraege", exporter.loadJobs}, {"wartungsfenster", exporter.loadMaintenanceWindows}, {"benachrichtigungswege", exporter.loadNotificationChannels}, {"konten", exporter.loadUsers}, {"einstellungen", exporter.loadSettings}, } for _, loadStep := range loadSteps { if loadError := loadStep.load(buildContext, snapshot); loadError != nil { return nil, fmt.Errorf("die %s konnten nicht gelesen werden: %w", loadStep.description, loadError) } } return snapshot, nil } // readSchemaVersion liest den Migrationsstand der Datenbank. func (exporter *Exporter) readSchemaVersion(readContext context.Context) (int, error) { var schemaVersion int scanError := exporter.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: %w", scanError) } return schemaVersion, nil } // loadRepositories liest die eingerichteten Ablagen. func (exporter *Exporter) loadRepositories(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT id::text, name, repository_type, location, COALESCE(repository_uuid, ''), status, hardened, capacity_bytes, retention_seconds, minimum_retention_seconds FROM repositories ORDER BY created_at` repositoryRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer repositoryRows.Close() for repositoryRows.Next() { var record RepositoryRecord if scanError := repositoryRows.Scan(&record.ID, &record.Name, &record.RepositoryType, &record.Location, &record.RepositoryUUID, &record.Status, &record.Hardened, &record.CapacityBytes, &record.RetentionSeconds, &record.MinimumRetentionSeconds); scanError != nil { return scanError } snapshot.Repositories = append(snapshot.Repositories, record) } return repositoryRows.Err() } // loadRetentionPolicies liest die Aufbewahrungsregeln. func (exporter *Exporter) loadRetentionPolicies(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT id::text, name, rules, keep_within_seconds, keep_last, keep_daily, keep_weekly, keep_monthly, keep_yearly, COALESCE(time_zone, '') FROM retention_policies ORDER BY created_at` policyRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer policyRows.Close() for policyRows.Next() { var record RetentionPolicyRecord if scanError := policyRows.Scan(&record.ID, &record.Name, &record.Rules, &record.KeepWithinSeconds, &record.KeepLast, &record.KeepDaily, &record.KeepWeekly, &record.KeepMonthly, &record.KeepYearly, &record.TimeZone); scanError != nil { return scanError } snapshot.RetentionPolicies = append(snapshot.RetentionPolicies, record) } return policyRows.Err() } // loadJobs liest die Sicherungsauftraege samt ihren Quellen. // // Geloeschte Auftraege bleiben draussen: Ein Soft-Delete bedeutet, dass jemand // den Auftrag beenden wollte. Ihn beim Wiederaufbau zurueckzubringen liesse eine // Sicherung wieder anlaufen, die abgeschaltet war. func (exporter *Exporter) loadJobs(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT id::text, name, COALESCE(description, ''), status, priority, schedule_type, schedule_config, repository_id::text, COALESCE(retention_policy_id::text, ''), rpo_seconds, rto_seconds, bandwidth_limit_bps, max_concurrency, retry_policy FROM backup_jobs WHERE deleted_at IS NULL ORDER BY created_at` jobRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer jobRows.Close() jobRecords := make([]JobRecord, 0, 16) for jobRows.Next() { var record JobRecord if scanError := jobRows.Scan(&record.ID, &record.Name, &record.Description, &record.Status, &record.Priority, &record.ScheduleType, &record.ScheduleConfig, &record.RepositoryID, &record.RetentionPolicyID, &record.RPOSeconds, &record.RTOSeconds, &record.BandwidthLimitBPS, &record.MaxConcurrency, &record.RetryPolicy); scanError != nil { return scanError } jobRecords = append(jobRecords, record) } if rowsError := jobRows.Err(); rowsError != nil { return rowsError } for jobIndex := range jobRecords { jobSources, sourcesError := exporter.loadJobSources(loadContext, jobRecords[jobIndex].ID) if sourcesError != nil { return sourcesError } jobRecords[jobIndex].Sources = jobSources } snapshot.Jobs = jobRecords return nil } // loadJobSources liest die Quellen eines Auftrags. func (exporter *Exporter) loadJobSources(loadContext context.Context, jobIdentifier string) ([]JobSourceRecord, error) { const selectStatement = ` SELECT id::text, source_type, source_id, source_name, include_patterns, exclude_patterns FROM backup_job_sources WHERE job_id = $1::uuid ORDER BY created_at` sourceRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement, jobIdentifier) if queryError != nil { return nil, queryError } defer sourceRows.Close() sourceRecords := make([]JobSourceRecord, 0, 4) for sourceRows.Next() { var record JobSourceRecord if scanError := sourceRows.Scan(&record.ID, &record.SourceType, &record.SourceID, &record.SourceName, &record.IncludePatterns, &record.ExcludePatterns); scanError != nil { return nil, scanError } sourceRecords = append(sourceRecords, record) } return sourceRecords, sourceRows.Err() } // loadMaintenanceWindows liest die Wartungsfenster samt Auftragszuordnung. func (exporter *Exporter) loadMaintenanceWindows(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT w.id::text, w.name, w.window_kind, w.starts_at, w.ends_at, w.recurrence, w.enabled, COALESCE(array_agg(j.job_id::text) FILTER (WHERE j.job_id IS NOT NULL), '{}') FROM maintenance_windows w LEFT JOIN maintenance_window_jobs j ON j.window_id = w.id GROUP BY w.id ORDER BY w.created_at` windowRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer windowRows.Close() for windowRows.Next() { var record MaintenanceWindowRecord if scanError := windowRows.Scan(&record.ID, &record.Name, &record.WindowKind, &record.StartsAt, &record.EndsAt, &record.Recurrence, &record.Enabled, &record.JobIDs); scanError != nil { return scanError } snapshot.MaintenanceWindows = append(snapshot.MaintenanceWindows, record) } return windowRows.Err() } // loadNotificationChannels liest die Benachrichtigungswege als Merkzettel. // // Die Spalten configuration und credentials_ciphertext werden **nicht einmal // abgefragt**. Sie versehentlich in einen Datenstrom zu bekommen, der in eine // Datei geschrieben wird, ist der Fehler, den man hinterher nicht mehr // zurueckholt: Die Datei liegt dann bereits im Repository. func (exporter *Exporter) loadNotificationChannels(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT id::text, name, channel_type, minimum_severity, enabled FROM notification_channels ORDER BY created_at` channelRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer channelRows.Close() for channelRows.Next() { var record NotificationChannelRecord if scanError := channelRows.Scan(&record.ID, &record.Name, &record.ChannelType, &record.MinimumSeverity, &record.WasEnabled); scanError != nil { return scanError } snapshot.NotificationChannels = append(snapshot.NotificationChannels, record) } return channelRows.Err() } // loadUsers liest die Konten mit ihren Rollennamen. // // password_hash wird nicht abgefragt — dieselbe Ueberlegung wie bei den // Kanaelen. Ein Argon2id-Hash ist kein Klartext, aber er laesst sich offline // angreifen, und ein Repository liegt naturgemaess ausserhalb der Anlage. func (exporter *Exporter) loadUsers(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = ` SELECT u.id::text, u.username, COALESCE(u.email, ''), u.status, COALESCE(array_agg(r.name) FILTER (WHERE r.name IS NOT NULL), '{}') FROM users u LEFT JOIN user_roles ur ON ur.user_id = u.id LEFT JOIN roles r ON r.id = ur.role_id WHERE u.deleted_at IS NULL GROUP BY u.id ORDER BY u.created_at` userRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer userRows.Close() for userRows.Next() { var record UserRecord if scanError := userRows.Scan(&record.ID, &record.Username, &record.Email, &record.Status, &record.Roles); scanError != nil { return scanError } snapshot.Users = append(snapshot.Users, record) } return userRows.Err() } // loadSettings liest die Systemeinstellungen. func (exporter *Exporter) loadSettings(loadContext context.Context, snapshot *Snapshot) error { const selectStatement = `SELECT key, value_json FROM system_settings ORDER BY key` settingRows, queryError := exporter.connectionPool.Query(loadContext, selectStatement) if queryError != nil { return queryError } defer settingRows.Close() for settingRows.Next() { var record SettingRecord if scanError := settingRows.Scan(&record.Key, &record.Value); scanError != nil { return scanError } snapshot.Settings = append(snapshot.Settings, record) } return settingRows.Err() } // WriteSnapshot legt einen Sicherungssatz im Repository ab. // // Geschrieben wird vierstufig wie jede andere Datei des Repositorys: Temp-Datei // im Zielverzeichnis, fsync der Datei, rename, fsync des Verzeichnisses. Ohne // die beiden fsync ueberlebt die Datei einen Stromausfall unvollstaendig — und // ein halber Sicherungssatz ist genau dann wertlos, wenn er gebraucht wird. func WriteSnapshot(repositoryRootPath string, snapshot *Snapshot) (string, error) { if validationError := snapshot.Validate(); validationError != nil { return "", validationError } snapshotDirectory := filepath.Join(repositoryRootPath, filepath.FromSlash(SnapshotDirectory)) if directoryError := os.MkdirAll(snapshotDirectory, 0o700); directoryError != nil { return "", fmt.Errorf("das verzeichnis des sicherungssatzes liess sich nicht anlegen: %w", directoryError) } encodedSnapshot, encodeError := json.MarshalIndent(snapshot, "", " ") if encodeError != nil { return "", fmt.Errorf("der sicherungssatz liess sich nicht kodieren: %w", encodeError) } snapshotPath := filepath.Join(snapshotDirectory, SnapshotFileName) if writeError := writeFileAtomically(snapshotPath, encodedSnapshot); writeError != nil { return "", writeError } // Zusaetzlich eine datierte Fassung. Der Sicherungssatz beschreibt einen // Stand; wer eine irrtuemliche Aenderung an der Konfiguration // zurueckdrehen will, braucht den vorherigen — die feste Datei allein waere // nach dem naechsten Lauf ueberschrieben. historyPath := filepath.Join(snapshotDirectory, fmt.Sprintf("control-plane-%s.json", snapshot.CreatedAt.Format("20060102-150405"))) if writeError := writeFileAtomically(historyPath, encodedSnapshot); writeError != nil { return "", writeError } return snapshotPath, nil } // writeFileAtomically schreibt eine Datei so, dass sie einen Absturz uebersteht. func writeFileAtomically(targetPath string, fileContent []byte) error { targetDirectory := filepath.Dir(targetPath) temporaryFile, createError := os.CreateTemp(targetDirectory, ".tmp-*") if createError != nil { return fmt.Errorf("die temporaere datei liess sich nicht anlegen: %w", createError) } temporaryPath := temporaryFile.Name() // Bei jedem Fehler nach diesem Punkt bleibt keine halbe Datei zurueck. defer func() { _ = temporaryFile.Close() _ = os.Remove(temporaryPath) }() if _, writeError := temporaryFile.Write(fileContent); writeError != nil { return fmt.Errorf("der sicherungssatz liess sich nicht schreiben: %w", writeError) } if syncError := temporaryFile.Sync(); syncError != nil { return fmt.Errorf("der sicherungssatz liess sich nicht auf den datentraeger schreiben: %w", syncError) } if closeError := temporaryFile.Close(); closeError != nil { return fmt.Errorf("die temporaere datei liess sich nicht schliessen: %w", closeError) } if chmodError := os.Chmod(temporaryPath, 0o600); chmodError != nil { return fmt.Errorf("die rechte liessen sich nicht setzen: %w", chmodError) } if renameError := os.Rename(temporaryPath, targetPath); renameError != nil { return fmt.Errorf("der sicherungssatz liess sich nicht an seinen platz bringen: %w", renameError) } return syncDirectory(targetDirectory) } // syncDirectory schreibt den Verzeichniseintrag auf den Datentraeger. func syncDirectory(directoryPath string) error { directoryHandle, openError := os.Open(directoryPath) if openError != nil { return fmt.Errorf("das verzeichnis liess sich nicht oeffnen: %w", openError) } defer func() { _ = directoryHandle.Close() }() if syncError := directoryHandle.Sync(); syncError != nil { return fmt.Errorf("das verzeichnis liess sich nicht schreiben: %w", syncError) } return nil }