package recovery import ( "context" "encoding/json" "errors" "fmt" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) // RestoreStatus ist der Zustand eines Wiederherstellungsauftrags. type RestoreStatus string const ( // RestoreStatusQueued wartet auf Ausfuehrung. RestoreStatusQueued RestoreStatus = "queued" // RestoreStatusRunning laeuft. RestoreStatusRunning RestoreStatus = "running" // RestoreStatusSucceeded ist vollstaendig gelungen. RestoreStatusSucceeded RestoreStatus = "succeeded" // RestoreStatusPartialFailure hat Objekte uebergangen. RestoreStatusPartialFailure RestoreStatus = "partial_failure" // RestoreStatusFailed ist gescheitert. RestoreStatusFailed RestoreStatus = "failed" // RestoreStatusCancelled wurde abgebrochen. RestoreStatusCancelled RestoreStatus = "cancelled" ) // IsFinished meldet einen abgeschlossenen Auftrag. func (restoreStatus RestoreStatus) IsFinished() bool { return restoreStatus != RestoreStatusQueued && restoreStatus != RestoreStatusRunning } // TargetType benennt die Art eines Wiederherstellungsziels. type TargetType string const ( // TargetFilesystem schreibt in ein beliebiges Verzeichnis. TargetFilesystem TargetType = "filesystem" // TargetOriginalLocation schreibt an den Ursprungsort zurueck. // // Der gefaehrlichste Fall: Dort liegen die Daten, die man eigentlich retten // will. Er verlangt deshalb eine eigene Berechtigung und eine ausdrueckliche // Bestaetigung. TargetOriginalLocation TargetType = "original_location" ) // RestoreJob ist ein Wiederherstellungsauftrag. type RestoreJob struct { // ID ist der oeffentliche Bezeichner. ID uuid.UUID `json:"id"` // BackupID ist das wiederherzustellende Backup in der Control Plane. BackupID uuid.UUID `json:"backup_id"` // SourceType benennt die Art der urspruenglichen Quelle. SourceType string `json:"source_type"` // TargetType benennt die Art des Ziels. TargetType TargetType `json:"target_type"` // TargetRef ist der Zielpfad. TargetRef string `json:"target_ref"` // PathPrefix beschraenkt auf einen Teilbaum. PathPrefix string `json:"path_prefix,omitempty"` // Status ist der Zustand. Status RestoreStatus `json:"status"` // OverwriteExisting erlaubt das Ueberschreiben vorhandener Daten. OverwriteExisting bool `json:"overwrite_existing"` // RestorePermissions setzt die urspruenglichen Rechte. RestorePermissions bool `json:"restore_permissions"` // VerifyContent prueft jede Datei gegen ihre Pruefsumme. VerifyContent bool `json:"verify_content"` // ValidationReport ist das Ergebnis der Vorabpruefung. ValidationReport *ValidationReport `json:"validation_report,omitempty"` // StartedAt ist der Beginn in UTC. StartedAt *time.Time `json:"started_at,omitempty"` // CompletedAt ist das Ende in UTC. CompletedAt *time.Time `json:"completed_at,omitempty"` // BytesRestored ist die zurueckgeschriebene Datenmenge. BytesRestored int64 `json:"bytes_restored"` // FilesRestored ist die Zahl zurueckgeschriebener Objekte. FilesRestored int64 `json:"files_restored"` // FilesSkipped ist die Zahl uebergangener Objekte. FilesSkipped int64 `json:"files_skipped"` // ErrorCode ist die Fehlerkennung. ErrorCode string `json:"error_code,omitempty"` // ErrorMessage ist die verstaendliche Fehlermeldung. ErrorMessage string `json:"error_message,omitempty"` // CorrelationID verbindet den Auftrag mit seinen Protokollzeilen. CorrelationID uuid.UUID `json:"correlation_id"` // CreatedBy benennt den Anfordernden. CreatedBy *uuid.UUID `json:"created_by,omitempty"` // CreatedAt ist der Anlagezeitpunkt in UTC. CreatedAt time.Time `json:"created_at"` } // Checkpoint haelt den Fortschritt einer Wiederherstellung. type Checkpoint struct { // LastCompletedPath ist der zuletzt vollstaendig zurueckgeschriebene Pfad. // // Ein einzelner Pfad genuegt, weil die Objekte im Manifest fest sortiert // sind: Alles lexikographisch davor ist erledigt. LastCompletedPath string `json:"last_completed_path"` // FilesCompleted ist die Zahl fertiger Objekte. FilesCompleted int64 `json:"files_completed"` // BytesCompleted ist die Menge fertiger Daten. BytesCompleted int64 `json:"bytes_completed"` // UpdatedAt ist der Zeitpunkt der letzten Fortschreibung in UTC. UpdatedAt time.Time `json:"updated_at"` } // RestoreSession ist eine Sitzung mit Pruefpunkt. type RestoreSession struct { // ID ist der oeffentliche Bezeichner. ID uuid.UUID `json:"id"` // RestoreJobID ist der zugehoerige Auftrag. RestoreJobID uuid.UUID `json:"restore_job_id"` // State ist der Zustand der Sitzung. State string `json:"state"` // Checkpoint haelt den Fortschritt. Checkpoint Checkpoint `json:"checkpoint"` // AttemptNumber zaehlt die Fortsetzungen, beginnend bei 1. AttemptNumber int `json:"attempt_number"` } // ErrRestoreNotFound meldet einen nicht vorhandenen Auftrag. var ErrRestoreNotFound = errors.New("der wiederherstellungsauftrag wurde nicht gefunden") // ErrTargetBusy meldet ein bereits belegtes Ziel. var ErrTargetBusy = errors.New("in dieses ziel wird bereits wiederhergestellt") // uniqueViolationCode ist der PostgreSQL-Fehlercode fuer Eindeutigkeitsverstoesse. const uniqueViolationCode = "23505" // Store legt Wiederherstellungsauftraege in PostgreSQL ab. type Store struct { // connectionPool ist der Datenbankpool der Control Plane. connectionPool *pgxpool.Pool } // NewStore erzeugt die Datenzugriffsschicht. func NewStore(connectionPool *pgxpool.Pool) *Store { return &Store{connectionPool: connectionPool} } // restoreColumnList sind die Spalten eines Auftrags in fester Reihenfolge. const restoreColumnList = ` id, backup_id, source_type, target_type, target_ref, path_prefix, status, overwrite_existing, restore_permissions, verify_content, validation_result, started_at, completed_at, bytes_restored, files_restored, files_skipped, error_code, error_message, correlation_id, created_by, created_at` // rowScanner deckt QueryRow und Rows gemeinsam ab. type rowScanner interface { // Scan liest die Spalten einer Zeile. Scan(destinations ...any) error } // scanRestoreJob liest eine Auftragszeile. func scanRestoreJob(scanner rowScanner) (*RestoreJob, error) { var ( loadedJob RestoreJob targetTypeText string statusText string pathPrefix *string validationJSON []byte errorCode *string errorMessage *string ) scanError := scanner.Scan( &loadedJob.ID, &loadedJob.BackupID, &loadedJob.SourceType, &targetTypeText, &loadedJob.TargetRef, &pathPrefix, &statusText, &loadedJob.OverwriteExisting, &loadedJob.RestorePermissions, &loadedJob.VerifyContent, &validationJSON, &loadedJob.StartedAt, &loadedJob.CompletedAt, &loadedJob.BytesRestored, &loadedJob.FilesRestored, &loadedJob.FilesSkipped, &errorCode, &errorMessage, &loadedJob.CorrelationID, &loadedJob.CreatedBy, &loadedJob.CreatedAt, ) if scanError != nil { return nil, scanError } loadedJob.TargetType = TargetType(targetTypeText) loadedJob.Status = RestoreStatus(statusText) if pathPrefix != nil { loadedJob.PathPrefix = *pathPrefix } if errorCode != nil { loadedJob.ErrorCode = *errorCode } if errorMessage != nil { loadedJob.ErrorMessage = *errorMessage } if len(validationJSON) > 0 { var parsedReport ValidationReport if unmarshalError := json.Unmarshal(validationJSON, &parsedReport); unmarshalError == nil { loadedJob.ValidationReport = &parsedReport } } return &loadedJob, nil } // CreateRequest beschreibt einen anzulegenden Wiederherstellungsauftrag. type CreateRequest struct { // BackupID ist das wiederherzustellende Backup. BackupID uuid.UUID // SourceType benennt die Art der urspruenglichen Quelle. SourceType string // TargetType benennt die Art des Ziels. TargetType TargetType // TargetRef ist der Zielpfad. TargetRef string // PathPrefix beschraenkt auf einen Teilbaum. PathPrefix string // OverwriteExisting erlaubt das Ueberschreiben vorhandener Daten. OverwriteExisting bool // RestorePermissions setzt die urspruenglichen Rechte. RestorePermissions bool // VerifyContent prueft jede Datei gegen ihre Pruefsumme. VerifyContent bool // ValidationReport ist das Ergebnis der Vorabpruefung. // // Es wird mitgespeichert, damit spaeter belegbar ist, was vor dem Start // bekannt war — etwa dass ein Ziel ueberschrieben wurde. ValidationReport *ValidationReport // CreatedBy benennt den Anfordernden. CreatedBy *uuid.UUID } // CreateRestore legt einen Wiederherstellungsauftrag samt Sitzung an. func (store *Store) CreateRestore(createContext context.Context, createRequest CreateRequest) (*RestoreJob, error) { transaction, transactionError := store.connectionPool.Begin(createContext) if transactionError != nil { return nil, fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError) } defer func() { _ = transaction.Rollback(createContext) }() var validationJSON []byte if createRequest.ValidationReport != nil { marshalled, marshalError := json.Marshal(createRequest.ValidationReport) if marshalError != nil { return nil, fmt.Errorf("der pruefbericht konnte nicht abgelegt werden: %w", marshalError) } validationJSON = marshalled } const insertStatement = ` INSERT INTO restore_jobs ( backup_id, source_type, target_type, target_ref, path_prefix, overwrite_existing, restore_permissions, verify_content, validation_result, correlation_id, created_by ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,gen_random_uuid(),$10) RETURNING ` + restoreColumnList createdJob, scanError := scanRestoreJob(transaction.QueryRow(createContext, insertStatement, createRequest.BackupID, createRequest.SourceType, string(createRequest.TargetType), createRequest.TargetRef, nullableText(createRequest.PathPrefix), createRequest.OverwriteExisting, createRequest.RestorePermissions, createRequest.VerifyContent, validationJSON, createRequest.CreatedBy, )) if scanError != nil { if isUniqueViolation(scanError) { // Zwei gleichzeitige Wiederherstellungen in dasselbe Ziel schrieben // sich gegenseitig zu — und zwar ohne dass es auffiele. return nil, fmt.Errorf("%w: %s", ErrTargetBusy, createRequest.TargetRef) } return nil, fmt.Errorf("der wiederherstellungsauftrag konnte nicht angelegt werden: %w", scanError) } // Die Sitzung entsteht mit dem Auftrag. Ohne sie gaebe es beim ersten // Abbruch keinen Ort, an dem der Fortschritt stuende. const insertSessionStatement = ` INSERT INTO restore_sessions (restore_job_id, state) VALUES ($1, 'active')` if _, execError := transaction.Exec(createContext, insertSessionStatement, createdJob.ID); execError != nil { return nil, fmt.Errorf("die sitzung konnte nicht angelegt werden: %w", execError) } if commitError := transaction.Commit(createContext); commitError != nil { return nil, fmt.Errorf("der auftrag konnte nicht festgeschrieben werden: %w", commitError) } return createdJob, nil } // GetRestore liest einen Wiederherstellungsauftrag. func (store *Store) GetRestore(readContext context.Context, restoreIdentifier uuid.UUID) (*RestoreJob, error) { const selectStatement = `SELECT ` + restoreColumnList + ` FROM restore_jobs WHERE id = $1` loadedJob, scanError := scanRestoreJob(store.connectionPool.QueryRow(readContext, selectStatement, restoreIdentifier)) if errors.Is(scanError, pgx.ErrNoRows) { return nil, fmt.Errorf("%w: %s", ErrRestoreNotFound, restoreIdentifier) } if scanError != nil { return nil, fmt.Errorf("der auftrag konnte nicht gelesen werden: %w", scanError) } return loadedJob, nil } // ListRestores liefert eine Seite von Wiederherstellungsauftraegen. func (store *Store) ListRestores(listContext context.Context, page int, pageSize int) ([]RestoreJob, int, error) { if page < 1 { page = 1 } if pageSize < 1 || pageSize > 200 { pageSize = 50 } const selectStatement = ` SELECT ` + restoreColumnList + `, COUNT(*) OVER () AS total_count FROM restore_jobs ORDER BY created_at DESC LIMIT $1 OFFSET $2` restoreRows, queryError := store.connectionPool.Query(listContext, selectStatement, pageSize, (page-1)*pageSize) if queryError != nil { return nil, 0, fmt.Errorf("die auftragsliste konnte nicht gelesen werden: %w", queryError) } defer restoreRows.Close() loadedJobs := make([]RestoreJob, 0, pageSize) var totalCount int for restoreRows.Next() { var ( loadedJob RestoreJob targetTypeText string statusText string pathPrefix *string validationJSON []byte errorCode *string errorMessage *string ) if scanError := restoreRows.Scan( &loadedJob.ID, &loadedJob.BackupID, &loadedJob.SourceType, &targetTypeText, &loadedJob.TargetRef, &pathPrefix, &statusText, &loadedJob.OverwriteExisting, &loadedJob.RestorePermissions, &loadedJob.VerifyContent, &validationJSON, &loadedJob.StartedAt, &loadedJob.CompletedAt, &loadedJob.BytesRestored, &loadedJob.FilesRestored, &loadedJob.FilesSkipped, &errorCode, &errorMessage, &loadedJob.CorrelationID, &loadedJob.CreatedBy, &loadedJob.CreatedAt, &totalCount, ); scanError != nil { return nil, 0, fmt.Errorf("ein auftrag konnte nicht gelesen werden: %w", scanError) } loadedJob.TargetType = TargetType(targetTypeText) loadedJob.Status = RestoreStatus(statusText) if pathPrefix != nil { loadedJob.PathPrefix = *pathPrefix } if errorCode != nil { loadedJob.ErrorCode = *errorCode } if errorMessage != nil { loadedJob.ErrorMessage = *errorMessage } loadedJobs = append(loadedJobs, loadedJob) } return loadedJobs, totalCount, restoreRows.Err() } // ClaimQueuedRestores uebernimmt anstehende Auftraege zur Ausfuehrung. // // Wie bei den Sicherungen sperrt FOR UPDATE SKIP LOCKED: Der erste Server // nimmt die Zeilen, der zweite ueberspringt sie, statt zu warten. func (store *Store) ClaimQueuedRestores(claimContext context.Context, maximumJobs int, schedulerInstance string) ([]RestoreJob, error) { transaction, transactionError := store.connectionPool.Begin(claimContext) if transactionError != nil { return nil, fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError) } defer func() { _ = transaction.Rollback(claimContext) }() const selectStatement = ` SELECT ` + restoreColumnList + ` FROM restore_jobs WHERE status = 'queued' ORDER BY created_at LIMIT $1 FOR UPDATE SKIP LOCKED` restoreRows, queryError := transaction.Query(claimContext, selectStatement, maximumJobs) if queryError != nil { return nil, fmt.Errorf("die anstehenden auftraege konnten nicht ermittelt werden: %w", queryError) } claimedJobs := make([]RestoreJob, 0, maximumJobs) for restoreRows.Next() { claimedJob, scanError := scanRestoreJob(restoreRows) if scanError != nil { restoreRows.Close() return nil, fmt.Errorf("ein anstehender auftrag konnte nicht gelesen werden: %w", scanError) } claimedJobs = append(claimedJobs, *claimedJob) } restoreRows.Close() if rowsError := restoreRows.Err(); rowsError != nil { return nil, rowsError } const markStatement = ` UPDATE restore_jobs SET scheduler_instance = $2, heartbeat_at = now() WHERE id = $1` for _, claimedJob := range claimedJobs { if _, execError := transaction.Exec(claimContext, markStatement, claimedJob.ID, schedulerInstance); execError != nil { return nil, fmt.Errorf("die uebernahme des auftrags %s schlug fehl: %w", claimedJob.ID, execError) } } if commitError := transaction.Commit(claimContext); commitError != nil { return nil, fmt.Errorf("die uebernahme konnte nicht festgeschrieben werden: %w", commitError) } return claimedJobs, nil } // StartRestore setzt einen uebernommenen Auftrag auf „laeuft". func (store *Store) StartRestore(startContext context.Context, restoreIdentifier uuid.UUID, schedulerInstance string) error { const updateStatement = ` UPDATE restore_jobs SET status = 'running', started_at = COALESCE(started_at, now()), heartbeat_at = now(), scheduler_instance = $2 WHERE id = $1 AND status = 'queued'` commandTag, execError := store.connectionPool.Exec(startContext, updateStatement, restoreIdentifier, schedulerInstance) if execError != nil { return fmt.Errorf("der auftrag konnte nicht gestartet werden: %w", execError) } if commandTag.RowsAffected() == 0 { return fmt.Errorf("%w: %s steht nicht mehr auf 'queued'", ErrRestoreNotFound, restoreIdentifier) } return nil } // RecordHeartbeat meldet einen laufenden Auftrag als lebendig. func (store *Store) RecordHeartbeat(heartbeatContext context.Context, restoreIdentifier uuid.UUID) error { const updateStatement = ` UPDATE restore_jobs SET heartbeat_at = now() WHERE id = $1 AND status = 'running'` if _, execError := store.connectionPool.Exec(heartbeatContext, updateStatement, restoreIdentifier); execError != nil { return fmt.Errorf("die lebendmeldung schlug fehl: %w", execError) } return nil } // GetActiveSession liest die aktive Sitzung eines Auftrags. func (store *Store) GetActiveSession(readContext context.Context, restoreIdentifier uuid.UUID) (*RestoreSession, error) { const selectStatement = ` SELECT id, restore_job_id, state, checkpoint, attempt_number FROM restore_sessions WHERE restore_job_id = $1 AND state = 'active'` var ( loadedSession RestoreSession checkpointJSON []byte ) scanError := store.connectionPool.QueryRow(readContext, selectStatement, restoreIdentifier).Scan( &loadedSession.ID, &loadedSession.RestoreJobID, &loadedSession.State, &checkpointJSON, &loadedSession.AttemptNumber) if errors.Is(scanError, pgx.ErrNoRows) { return nil, nil } if scanError != nil { return nil, fmt.Errorf("die sitzung konnte nicht gelesen werden: %w", scanError) } if len(checkpointJSON) > 0 { _ = json.Unmarshal(checkpointJSON, &loadedSession.Checkpoint) } return &loadedSession, nil } // SaveCheckpoint schreibt den Fortschritt einer Sitzung fort. // // Der Pruefpunkt entsteht **nachdem** eine Datei vollstaendig am Platz liegt. // Ein Pruefpunkt auf eine halb geschriebene Datei waere schlimmer als keiner: // Die Fortsetzung uebersprange sie als erledigt. func (store *Store) SaveCheckpoint(saveContext context.Context, sessionIdentifier uuid.UUID, checkpoint Checkpoint) error { checkpoint.UpdatedAt = time.Now().UTC() checkpointJSON, marshalError := json.Marshal(checkpoint) if marshalError != nil { return fmt.Errorf("der pruefpunkt konnte nicht abgelegt werden: %w", marshalError) } const updateStatement = ` UPDATE restore_sessions SET checkpoint = $2, updated_at = now() WHERE id = $1` if _, execError := store.connectionPool.Exec(saveContext, updateStatement, sessionIdentifier, checkpointJSON); execError != nil { return fmt.Errorf("der pruefpunkt konnte nicht gespeichert werden: %w", execError) } return nil } // CloseSession beendet eine Sitzung. func (store *Store) CloseSession(closeContext context.Context, sessionIdentifier uuid.UUID, finalState string) error { const updateStatement = ` UPDATE restore_sessions SET state = $2, updated_at = now() WHERE id = $1` if _, execError := store.connectionPool.Exec(closeContext, updateStatement, sessionIdentifier, finalState); execError != nil { return fmt.Errorf("die sitzung konnte nicht beendet werden: %w", execError) } return nil } // RestoreOutcome beschreibt das Ergebnis einer Wiederherstellung. type RestoreOutcome struct { // Status ist der erreichte Zustand. Status RestoreStatus // BytesRestored ist die zurueckgeschriebene Datenmenge. BytesRestored int64 // FilesRestored ist die Zahl zurueckgeschriebener Objekte. FilesRestored int64 // FilesSkipped ist die Zahl uebergangener Objekte. FilesSkipped int64 // ErrorCode ist die Fehlerkennung. ErrorCode string // ErrorMessage ist die verstaendliche Fehlermeldung. ErrorMessage string } // FinishRestore schreibt das Ergebnis fest. func (store *Store) FinishRestore(finishContext context.Context, restoreIdentifier uuid.UUID, outcome RestoreOutcome) error { const updateStatement = ` UPDATE restore_jobs SET status = $2, completed_at = now(), bytes_restored = $3, files_restored = $4, files_skipped = $5, error_code = $6, error_message = $7 WHERE id = $1` commandTag, execError := store.connectionPool.Exec(finishContext, updateStatement, restoreIdentifier, string(outcome.Status), outcome.BytesRestored, outcome.FilesRestored, outcome.FilesSkipped, nullableText(outcome.ErrorCode), nullableText(outcome.ErrorMessage), ) if execError != nil { return fmt.Errorf("das ergebnis konnte nicht festgeschrieben werden: %w", execError) } if commandTag.RowsAffected() == 0 { return fmt.Errorf("%w: %s", ErrRestoreNotFound, restoreIdentifier) } return nil } // CancelRestore bricht einen Auftrag ab. func (store *Store) CancelRestore(cancelContext context.Context, restoreIdentifier uuid.UUID, reason string) error { const updateStatement = ` UPDATE restore_jobs SET status = 'cancelled', completed_at = now(), error_code = 'CANCELLED', error_message = $2 WHERE id = $1 AND status IN ('queued', 'running')` commandTag, execError := store.connectionPool.Exec(cancelContext, updateStatement, restoreIdentifier, reason) if execError != nil { return fmt.Errorf("der auftrag konnte nicht abgebrochen werden: %w", execError) } if commandTag.RowsAffected() == 0 { return fmt.Errorf("%w: %s laeuft nicht mehr", ErrRestoreNotFound, restoreIdentifier) } return nil } // ReclaimStaleRestores gibt Auftraege abgestuerzter Control-Server frei. // // Anders als bei Sicherungen wird der Auftrag **nicht** erneut eingereiht: Eine // halb geschriebene Wiederherstellung erneut zu starten kann Daten // beschaedigen, die der erste Versuch bereits am Platz hatte. Sie wird als // gescheitert vermerkt und wartet auf einen Blick — die Sitzung mit ihrem // Pruefpunkt bleibt bestehen, damit eine Fortsetzung moeglich ist. func (store *Store) ReclaimStaleRestores(reclaimContext context.Context, staleAfter time.Duration) (int, error) { const updateStatement = ` UPDATE restore_jobs SET status = 'failed', completed_at = now(), error_code = 'SCHEDULER_LOST', error_message = $2 WHERE status IN ('queued', 'running') AND heartbeat_at IS NOT NULL AND heartbeat_at < now() - $1::interval` const explanation = "Der ausfuehrende Control-Server meldete sich nicht mehr. Die Wiederherstellung " + "wurde nicht automatisch fortgesetzt: Ein zweiter Lauf koennte bereits zurueckgeschriebene " + "Daten beschaedigen. Der Pruefpunkt der Sitzung bleibt erhalten." commandTag, execError := store.connectionPool.Exec(reclaimContext, updateStatement, staleAfter.String(), explanation) if execError != nil { return 0, fmt.Errorf("verwaiste auftraege konnten nicht freigegeben werden: %w", execError) } return int(commandTag.RowsAffected()), nil } // isUniqueViolation erkennt einen Eindeutigkeitsverstoss. func isUniqueViolation(occurredError error) bool { var pgError interface{ SQLState() string } return errors.As(occurredError, &pgError) && pgError.SQLState() == uniqueViolationCode } // nullableText wandelt eine leere Zeichenkette in NULL. func nullableText(textValue string) *string { if textValue == "" { return nil } return &textValue } // ErrRestoreNotResumable meldet einen nicht fortsetzbaren Auftrag. var ErrRestoreNotResumable = errors.New("dieser wiederherstellungsauftrag laesst sich nicht fortsetzen") // ResumeRestore reiht einen unterbrochenen Auftrag erneut ein. // // Fortsetzbar ist nur, was abgeschlossen und nicht gelungen ist: ein // gescheiterter oder abgebrochener Lauf. Einen erfolgreichen erneut // einzureihen schriebe dieselben Daten ein zweites Mal; einen laufenden // erzeugte zwei Schreiber auf demselben Ziel. // // Die Sitzung wird **nicht** angefasst — ihr Pruefpunkt ist der Grund, warum // die Fortsetzung ueberhaupt etwas spart. func (store *Store) ResumeRestore(resumeContext context.Context, restoreIdentifier uuid.UUID) (*RestoreJob, error) { transaction, transactionError := store.connectionPool.Begin(resumeContext) if transactionError != nil { return nil, fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError) } defer func() { _ = transaction.Rollback(resumeContext) }() const updateStatement = ` UPDATE restore_jobs SET status = 'queued', completed_at = NULL, error_code = NULL, error_message = NULL, heartbeat_at = NULL, scheduler_instance = NULL WHERE id = $1 AND status IN ('failed', 'cancelled', 'partial_failure') RETURNING ` + restoreColumnList resumedJob, scanError := scanRestoreJob(transaction.QueryRow(resumeContext, updateStatement, restoreIdentifier)) if errors.Is(scanError, pgx.ErrNoRows) { return nil, ErrRestoreNotResumable } if scanError != nil { if isUniqueViolation(scanError) { // In dasselbe Ziel wird bereits geschrieben. return nil, ErrTargetBusy } return nil, fmt.Errorf("der auftrag konnte nicht fortgesetzt werden: %w", scanError) } // Ist die Sitzung geschlossen — etwa nach einem frueher gelungenen Lauf —, // wird eine neue eroeffnet. Ohne aktive Sitzung gaebe es keinen Ort fuer // den Pruefpunkt des neuen Versuchs. const ensureSessionStatement = ` INSERT INTO restore_sessions (restore_job_id, state, attempt_number) SELECT $1, 'active', COALESCE(MAX(attempt_number), 0) + 1 FROM restore_sessions WHERE restore_job_id = $1 ON CONFLICT DO NOTHING` if _, execError := transaction.Exec(resumeContext, ensureSessionStatement, restoreIdentifier); execError != nil { return nil, fmt.Errorf("die sitzung konnte nicht sichergestellt werden: %w", execError) } if commitError := transaction.Commit(resumeContext); commitError != nil { return nil, fmt.Errorf("die fortsetzung konnte nicht festgeschrieben werden: %w", commitError) } return resumedJob, nil }