**Sicherungsart.** Bisher entschied der Executor allein: Liegt ein Elternbackup
vor, wird inkrementell gesichert. Jetzt waehlbar je Auftrag —
- `incremental` (Standard, bisheriges Verhalten),
- `always_full`, oder
- inkrementell **mit einem festen Volltag** ("immer freitags").
Migration 000014 mit drei CHECKs. Der dritte lehnt "immer voll" zusammen mit
einem Wochentag ab: Dann ist ohnehin jeder Lauf voll, und die Regel gehoert in
die Datenbank, weil im Code jede Stelle sie einhalten muesste — eine vergisst
es. Real geprueft: der Widerspruch wird abgewiesen.
Der Wochentag wird in der **Zeitzone des Zeitplans** bestimmt. Rechnete der
Server in UTC, bekaeme ein Betreiber in Berlin seine Vollsicherung am
Donnerstagabend und wunderte sich, warum sie freitags fehlt. Vier Tests, der
entscheidende durch Mutation als fangend bestaetigt.
Zur Einordnung, weil es leicht verwechselt wird: Der Platzbedarf steigt bei
"immer voll" **nicht** nennenswert — unveraenderte Bloecke werden dedupliziert
und liegen weiterhin nur einmal im Repository. Was steigt, ist die Laufzeit.
Steht so in der Maske.
**Aufnahme-Token zeigte "undefined".** Das Feld heisst `token`, nicht
`enrollment_token` — Letzteres ist der Name im *Anfrage*koerper der
Registrierung. Der dritte Formfehler dieser Art; alle konsumierten Endpunkte
sind jetzt gegen den laufenden Dienst abgeglichen.
**Der Aufnahmedialog** hat jetzt eine vollstaendige Anleitung fuer Linux und
Windows mit fertig ausgefuellten Befehlen — Serveradresse und Token eingesetzt,
je Schritt einzeln kopierbar. Eine Anleitung mit Platzhaltern fuehrt
zuverlaessig dazu, dass jemand `<token>` woertlich einsetzt und dann eine
Fehlermeldung sucht, die nichts mit seinem Problem zu tun hat. Dazu die beiden
Stolperstellen: `--state` will eine Datei, und der Agent braucht Schreibzugriff
aufs Repository. Beim Windows-Weg steht dabei, dass der Dienst nie auf echter
Hardware lief.
**update.sh ruestet die Wiederherstellungsflaeche nach** — anlegen und in
ReadWritePaths eintragen. Ein Schritt, den man von Hand ausfuehren muss, wird
uebersehen und faellt erst im Ernstfall auf.
84 Tests im Frontend, alle Go-Tests gruen, shellcheck sauber.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
783 lines
26 KiB
Go
783 lines
26 KiB
Go
package jobs
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
"github.com/syncova/syncova/packages/scheduler"
|
|
)
|
|
|
|
// ErrJobNotFound meldet einen nicht vorhandenen Auftrag.
|
|
var ErrJobNotFound = errors.New("der sicherungsauftrag wurde nicht gefunden")
|
|
|
|
// ErrJobNameTaken meldet einen bereits vergebenen Namen.
|
|
var ErrJobNameTaken = errors.New("ein sicherungsauftrag dieses namens besteht bereits")
|
|
|
|
// ErrRunAlreadyActive meldet einen bereits laufenden Auftrag.
|
|
//
|
|
// Der Fehler entsteht am Teilindex der Datenbank, nicht an einer Prüfung im
|
|
// Code. Das ist der Unterschied zwischen „meistens richtig" und „richtig":
|
|
// Zwei Control-Server, die gleichzeitig denselben fälligen Auftrag sehen,
|
|
// prüfen beide erfolgreich und starten beide — die Datenbank lässt nur einen
|
|
// durch.
|
|
var ErrRunAlreadyActive = errors.New("für diesen auftrag läuft bereits ein sicherungslauf")
|
|
|
|
// uniqueViolationCode ist der PostgreSQL-Fehlercode für Eindeutigkeitsverstöße.
|
|
const uniqueViolationCode = "23505"
|
|
|
|
// PostgresStore legt Aufträge und Läufe in PostgreSQL ab.
|
|
type PostgresStore struct {
|
|
// connectionPool ist der Datenbankpool der Control Plane.
|
|
connectionPool *pgxpool.Pool
|
|
}
|
|
|
|
// NewPostgresStore erzeugt die Datenzugriffsschicht.
|
|
func NewPostgresStore(connectionPool *pgxpool.Pool) *PostgresStore {
|
|
return &PostgresStore{connectionPool: connectionPool}
|
|
}
|
|
|
|
// CreateJob legt einen Auftrag samt Quellen und Abhängigkeiten an.
|
|
//
|
|
// Alles geschieht in einer Transaktion: Ein Auftrag ohne seine Quellen wäre ein
|
|
// Auftrag, der erfolgreich durchliefe, ohne etwas zu sichern.
|
|
func (store *PostgresStore) CreateJob(createContext context.Context, newJob *Job) (uuid.UUID, error) {
|
|
if validationError := newJob.Validate(); validationError != nil {
|
|
return uuid.Nil, validationError
|
|
}
|
|
|
|
transaction, transactionError := store.connectionPool.Begin(createContext)
|
|
if transactionError != nil {
|
|
return uuid.Nil, fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError)
|
|
}
|
|
|
|
defer func() { _ = transaction.Rollback(createContext) }()
|
|
|
|
scheduleJSON, marshalError := json.Marshal(newJob.Schedule)
|
|
if marshalError != nil {
|
|
return uuid.Nil, fmt.Errorf("der zeitplan konnte nicht abgelegt werden: %w", marshalError)
|
|
}
|
|
|
|
retryJSON, retryMarshalError := json.Marshal(newJob.RetryPolicy)
|
|
if retryMarshalError != nil {
|
|
return uuid.Nil, fmt.Errorf("die wiederholungsstrategie konnte nicht abgelegt werden: %w", retryMarshalError)
|
|
}
|
|
|
|
const insertJobStatement = `
|
|
INSERT INTO backup_jobs (
|
|
name, description, status, priority, schedule_type, schedule_config,
|
|
repository_id, retention_policy_id, rpo_seconds, rto_seconds,
|
|
bandwidth_limit_bps, max_concurrency, retry_policy, next_run_at, created_by,
|
|
backup_mode, full_backup_weekday
|
|
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17)
|
|
RETURNING id`
|
|
|
|
var createdJobID uuid.UUID
|
|
|
|
scanError := transaction.QueryRow(createContext, insertJobStatement,
|
|
newJob.Name,
|
|
nullableText(newJob.Description),
|
|
string(newJob.Status),
|
|
string(newJob.Priority),
|
|
string(newJob.Schedule.ScheduleType),
|
|
scheduleJSON,
|
|
newJob.RepositoryID,
|
|
newJob.RetentionPolicyID,
|
|
nullableSeconds(newJob.RecoveryPointObjective),
|
|
nullableSeconds(newJob.RecoveryTimeObjective),
|
|
nullableBandwidth(newJob.BandwidthLimitBytesPerSecond),
|
|
newJob.MaximumConcurrency,
|
|
retryJSON,
|
|
newJob.NextRunAt,
|
|
newJob.CreatedBy,
|
|
string(normalizeBackupMode(newJob.BackupMode)),
|
|
nullableWeekday(newJob.FullBackupWeekday),
|
|
).Scan(&createdJobID)
|
|
|
|
if scanError != nil {
|
|
if isUniqueViolation(scanError) {
|
|
return uuid.Nil, fmt.Errorf("%w: %s", ErrJobNameTaken, newJob.Name)
|
|
}
|
|
|
|
return uuid.Nil, fmt.Errorf("der auftrag konnte nicht angelegt werden: %w", scanError)
|
|
}
|
|
|
|
if sourceError := insertJobSources(createContext, transaction, createdJobID, newJob.Sources); sourceError != nil {
|
|
return uuid.Nil, sourceError
|
|
}
|
|
|
|
if dependencyError := insertJobDependencies(createContext, transaction, createdJobID, newJob.DependsOnJobIDs); dependencyError != nil {
|
|
return uuid.Nil, dependencyError
|
|
}
|
|
|
|
if commitError := transaction.Commit(createContext); commitError != nil {
|
|
return uuid.Nil, fmt.Errorf("der auftrag konnte nicht festgeschrieben werden: %w", commitError)
|
|
}
|
|
|
|
return createdJobID, nil
|
|
}
|
|
|
|
// insertJobSources legt die Quellen eines Auftrags an.
|
|
func insertJobSources(insertContext context.Context, transaction pgx.Tx, jobIdentifier uuid.UUID, jobSources []JobSource) error {
|
|
const insertSourceStatement = `
|
|
INSERT INTO backup_job_sources (
|
|
job_id, source_type, source_id, source_name, agent_id, include_patterns,
|
|
exclude_patterns, cluster_id
|
|
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8)`
|
|
|
|
for _, jobSource := range jobSources {
|
|
includeJSON, includeError := json.Marshal(defaultToEmptySlice(jobSource.IncludePatterns))
|
|
if includeError != nil {
|
|
return fmt.Errorf("die einschlussregeln konnten nicht abgelegt werden: %w", includeError)
|
|
}
|
|
|
|
excludeJSON, excludeError := json.Marshal(defaultToEmptySlice(jobSource.ExcludePatterns))
|
|
if excludeError != nil {
|
|
return fmt.Errorf("die ausschlussregeln konnten nicht abgelegt werden: %w", excludeError)
|
|
}
|
|
|
|
if _, execError := transaction.Exec(insertContext, insertSourceStatement,
|
|
jobIdentifier,
|
|
string(jobSource.SourceType),
|
|
jobSource.SourceID,
|
|
nullableText(jobSource.SourceName),
|
|
jobSource.AgentID,
|
|
includeJSON,
|
|
excludeJSON,
|
|
jobSource.ClusterID,
|
|
); execError != nil {
|
|
return fmt.Errorf("die quelle %s konnte nicht angelegt werden: %w", jobSource.SourceID, execError)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// insertJobDependencies legt die Abhängigkeiten eines Auftrags an.
|
|
func insertJobDependencies(insertContext context.Context, transaction pgx.Tx, jobIdentifier uuid.UUID, dependencyIDs []uuid.UUID) error {
|
|
const insertDependencyStatement = `
|
|
INSERT INTO backup_job_dependencies (job_id, depends_on_job_id) VALUES ($1,$2)`
|
|
|
|
for _, dependencyID := range dependencyIDs {
|
|
if _, execError := transaction.Exec(insertContext, insertDependencyStatement,
|
|
jobIdentifier, dependencyID); execError != nil {
|
|
return fmt.Errorf("die abhängigkeit von %s konnte nicht angelegt werden: %w", dependencyID, execError)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// GetJob liest einen Auftrag samt Quellen und Abhängigkeiten.
|
|
func (store *PostgresStore) GetJob(readContext context.Context, jobIdentifier uuid.UUID) (*Job, error) {
|
|
const selectJobStatement = `
|
|
SELECT id, name, COALESCE(description,''), status, priority, schedule_config,
|
|
repository_id, retention_policy_id, rpo_seconds, rto_seconds,
|
|
bandwidth_limit_bps, max_concurrency, retry_policy, backup_mode, full_backup_weekday,
|
|
next_run_at, last_run_at, last_outcome, paused_at, created_by, created_at, updated_at
|
|
FROM backup_jobs
|
|
WHERE id = $1 AND deleted_at IS NULL`
|
|
|
|
loadedJob, scanError := scanJobRow(store.connectionPool.QueryRow(readContext, selectJobStatement, jobIdentifier))
|
|
if errors.Is(scanError, pgx.ErrNoRows) {
|
|
return nil, fmt.Errorf("%w: %s", ErrJobNotFound, jobIdentifier)
|
|
}
|
|
|
|
if scanError != nil {
|
|
return nil, fmt.Errorf("der auftrag konnte nicht gelesen werden: %w", scanError)
|
|
}
|
|
|
|
if sourceError := store.loadJobSources(readContext, loadedJob); sourceError != nil {
|
|
return nil, sourceError
|
|
}
|
|
|
|
if dependencyError := store.loadJobDependencies(readContext, loadedJob); dependencyError != nil {
|
|
return nil, dependencyError
|
|
}
|
|
|
|
return loadedJob, nil
|
|
}
|
|
|
|
// rowScanner deckt QueryRow und Rows gemeinsam ab.
|
|
type rowScanner interface {
|
|
// Scan liest die Spalten einer Zeile.
|
|
Scan(destinations ...any) error
|
|
}
|
|
|
|
// scanJobRow liest eine Auftragszeile.
|
|
func scanJobRow(scanner rowScanner) (*Job, error) {
|
|
var (
|
|
loadedJob Job
|
|
statusText string
|
|
priorityText string
|
|
scheduleJSON []byte
|
|
retryJSON []byte
|
|
rpoSeconds *int64
|
|
rtoSeconds *int64
|
|
bandwidthLimit *int64
|
|
lastOutcomeText *string
|
|
descriptionValue string
|
|
// Leerer Text und NULL bedeuten beide "incremental" — der Standard,
|
|
// der auch fuer Auftraege aus der Zeit vor dieser Spalte gilt.
|
|
backupModeText string
|
|
fullBackupWeekday *int16
|
|
)
|
|
|
|
scanError := scanner.Scan(
|
|
&loadedJob.ID,
|
|
&loadedJob.Name,
|
|
&descriptionValue,
|
|
&statusText,
|
|
&priorityText,
|
|
&scheduleJSON,
|
|
&loadedJob.RepositoryID,
|
|
&loadedJob.RetentionPolicyID,
|
|
&rpoSeconds,
|
|
&rtoSeconds,
|
|
&bandwidthLimit,
|
|
&loadedJob.MaximumConcurrency,
|
|
&retryJSON,
|
|
&backupModeText,
|
|
&fullBackupWeekday,
|
|
&loadedJob.NextRunAt,
|
|
&loadedJob.LastRunAt,
|
|
&lastOutcomeText,
|
|
&loadedJob.PausedAt,
|
|
&loadedJob.CreatedBy,
|
|
&loadedJob.CreatedAt,
|
|
&loadedJob.UpdatedAt,
|
|
)
|
|
if scanError != nil {
|
|
return nil, scanError
|
|
}
|
|
|
|
loadedJob.Description = descriptionValue
|
|
loadedJob.Status = JobStatus(statusText)
|
|
loadedJob.Priority = scheduler.Priority(priorityText)
|
|
loadedJob.BackupMode = normalizeBackupMode(BackupMode(backupModeText))
|
|
loadedJob.FullBackupWeekday = weekdayFromDatabase(fullBackupWeekday)
|
|
|
|
if unmarshalError := json.Unmarshal(scheduleJSON, &loadedJob.Schedule); unmarshalError != nil {
|
|
return nil, fmt.Errorf("der zeitplan des auftrags %s ist unlesbar: %w", loadedJob.ID, unmarshalError)
|
|
}
|
|
|
|
if len(retryJSON) > 0 {
|
|
if unmarshalError := json.Unmarshal(retryJSON, &loadedJob.RetryPolicy); unmarshalError != nil {
|
|
return nil, fmt.Errorf("die wiederholungsstrategie des auftrags %s ist unlesbar: %w", loadedJob.ID, unmarshalError)
|
|
}
|
|
}
|
|
|
|
// Eine leere Strategie in der Datenbank ergäbe null Versuche — also einen
|
|
// Auftrag, der bei der kleinsten Störung aufgibt.
|
|
if loadedJob.RetryPolicy.MaximumAttempts < 1 {
|
|
loadedJob.RetryPolicy = scheduler.DefaultRetryPolicy()
|
|
}
|
|
|
|
if rpoSeconds != nil {
|
|
loadedJob.RecoveryPointObjective = time.Duration(*rpoSeconds) * time.Second
|
|
}
|
|
|
|
if rtoSeconds != nil {
|
|
loadedJob.RecoveryTimeObjective = time.Duration(*rtoSeconds) * time.Second
|
|
}
|
|
|
|
if bandwidthLimit != nil {
|
|
loadedJob.BandwidthLimitBytesPerSecond = *bandwidthLimit
|
|
}
|
|
|
|
if lastOutcomeText != nil {
|
|
loadedJob.LastOutcome = scheduler.JobOutcome(*lastOutcomeText)
|
|
}
|
|
|
|
return &loadedJob, nil
|
|
}
|
|
|
|
// loadJobSources lädt die Quellen eines Auftrags.
|
|
func (store *PostgresStore) loadJobSources(readContext context.Context, targetJob *Job) error {
|
|
const selectSourcesStatement = `
|
|
SELECT id, source_type, source_id, COALESCE(source_name,''), agent_id,
|
|
include_patterns, exclude_patterns, cluster_id
|
|
FROM backup_job_sources
|
|
WHERE job_id = $1
|
|
ORDER BY created_at, source_id`
|
|
|
|
sourceRows, queryError := store.connectionPool.Query(readContext, selectSourcesStatement, targetJob.ID)
|
|
if queryError != nil {
|
|
return fmt.Errorf("die quellen des auftrags konnten nicht gelesen werden: %w", queryError)
|
|
}
|
|
|
|
defer sourceRows.Close()
|
|
|
|
targetJob.Sources = make([]JobSource, 0, 4)
|
|
|
|
for sourceRows.Next() {
|
|
var (
|
|
loadedSource JobSource
|
|
typeText string
|
|
includeJSON []byte
|
|
excludeJSON []byte
|
|
)
|
|
|
|
if scanError := sourceRows.Scan(&loadedSource.ID, &typeText, &loadedSource.SourceID,
|
|
&loadedSource.SourceName, &loadedSource.AgentID, &includeJSON, &excludeJSON,
|
|
&loadedSource.ClusterID); scanError != nil {
|
|
return fmt.Errorf("eine quelle konnte nicht gelesen werden: %w", scanError)
|
|
}
|
|
|
|
loadedSource.SourceType = SourceType(typeText)
|
|
|
|
if unmarshalError := json.Unmarshal(includeJSON, &loadedSource.IncludePatterns); unmarshalError != nil {
|
|
return fmt.Errorf("die einschlussregeln sind unlesbar: %w", unmarshalError)
|
|
}
|
|
|
|
if unmarshalError := json.Unmarshal(excludeJSON, &loadedSource.ExcludePatterns); unmarshalError != nil {
|
|
return fmt.Errorf("die ausschlussregeln sind unlesbar: %w", unmarshalError)
|
|
}
|
|
|
|
targetJob.Sources = append(targetJob.Sources, loadedSource)
|
|
}
|
|
|
|
return sourceRows.Err()
|
|
}
|
|
|
|
// loadJobDependencies lädt die Abhängigkeiten eines Auftrags.
|
|
func (store *PostgresStore) loadJobDependencies(readContext context.Context, targetJob *Job) error {
|
|
const selectDependenciesStatement = `
|
|
SELECT depends_on_job_id FROM backup_job_dependencies WHERE job_id = $1 ORDER BY created_at`
|
|
|
|
dependencyRows, queryError := store.connectionPool.Query(readContext, selectDependenciesStatement, targetJob.ID)
|
|
if queryError != nil {
|
|
return fmt.Errorf("die abhängigkeiten konnten nicht gelesen werden: %w", queryError)
|
|
}
|
|
|
|
defer dependencyRows.Close()
|
|
|
|
targetJob.DependsOnJobIDs = make([]uuid.UUID, 0, 2)
|
|
|
|
for dependencyRows.Next() {
|
|
var dependencyID uuid.UUID
|
|
|
|
if scanError := dependencyRows.Scan(&dependencyID); scanError != nil {
|
|
return fmt.Errorf("eine abhängigkeit konnte nicht gelesen werden: %w", scanError)
|
|
}
|
|
|
|
targetJob.DependsOnJobIDs = append(targetJob.DependsOnJobIDs, dependencyID)
|
|
}
|
|
|
|
return dependencyRows.Err()
|
|
}
|
|
|
|
// ListFilter schränkt eine Auftragsliste ein.
|
|
type ListFilter struct {
|
|
// Status beschränkt auf einen Zustand.
|
|
Status JobStatus
|
|
// RepositoryID beschränkt auf ein Repository.
|
|
RepositoryID *uuid.UUID
|
|
// SearchTerm sucht in Name und Beschreibung.
|
|
SearchTerm string
|
|
// Page ist die Seitennummer, beginnend bei 1.
|
|
Page int
|
|
// PageSize ist die Seitengröße.
|
|
PageSize int
|
|
}
|
|
|
|
// maximumPageSize begrenzt die Seitengröße.
|
|
//
|
|
// Ohne Grenze könnte ein Aufrufer mit page_size=1000000 die gesamte Tabelle in
|
|
// den Speicher des Dienstes ziehen.
|
|
const maximumPageSize = 200
|
|
|
|
// defaultPageSize ist die Seitengröße ohne Angabe.
|
|
const defaultPageSize = 50
|
|
|
|
// normalize bringt die Seitenangaben in einen brauchbaren Bereich.
|
|
func (listFilter *ListFilter) normalize() {
|
|
if listFilter.Page < 1 {
|
|
listFilter.Page = 1
|
|
}
|
|
|
|
if listFilter.PageSize < 1 {
|
|
listFilter.PageSize = defaultPageSize
|
|
}
|
|
|
|
if listFilter.PageSize > maximumPageSize {
|
|
listFilter.PageSize = maximumPageSize
|
|
}
|
|
}
|
|
|
|
// ListJobs liefert eine Seite von Aufträgen samt Gesamtzahl.
|
|
func (store *PostgresStore) ListJobs(listContext context.Context, listFilter ListFilter) ([]Job, int, error) {
|
|
listFilter.normalize()
|
|
|
|
const selectStatement = `
|
|
SELECT id, name, COALESCE(description,''), status, priority, schedule_config,
|
|
repository_id, retention_policy_id, rpo_seconds, rto_seconds,
|
|
bandwidth_limit_bps, max_concurrency, retry_policy, backup_mode, full_backup_weekday,
|
|
next_run_at, last_run_at, last_outcome, paused_at, created_by, created_at, updated_at,
|
|
COUNT(*) OVER () AS total_count
|
|
FROM backup_jobs
|
|
WHERE deleted_at IS NULL
|
|
AND ($1::text IS NULL OR status = $1)
|
|
AND ($2::uuid IS NULL OR repository_id = $2)
|
|
AND ($3::text IS NULL OR name ILIKE '%' || $3 || '%' OR description ILIKE '%' || $3 || '%')
|
|
ORDER BY name
|
|
LIMIT $4 OFFSET $5`
|
|
|
|
jobRows, queryError := store.connectionPool.Query(listContext, selectStatement,
|
|
nullableText(string(listFilter.Status)),
|
|
listFilter.RepositoryID,
|
|
nullableText(listFilter.SearchTerm),
|
|
listFilter.PageSize,
|
|
(listFilter.Page-1)*listFilter.PageSize,
|
|
)
|
|
if queryError != nil {
|
|
return nil, 0, fmt.Errorf("die auftragsliste konnte nicht gelesen werden: %w", queryError)
|
|
}
|
|
|
|
defer jobRows.Close()
|
|
|
|
loadedJobs := make([]Job, 0, listFilter.PageSize)
|
|
var totalCount int
|
|
|
|
for jobRows.Next() {
|
|
// Die Gesamtzahl kommt als zusätzliche Spalte aus derselben Abfrage.
|
|
// Eine zweite Zählabfrage könnte bei gleichzeitigen Änderungen eine
|
|
// andere Zahl liefern als die Seite selbst.
|
|
scannedJob, scanError := scanJobRowWithTotal(jobRows, &totalCount)
|
|
if scanError != nil {
|
|
return nil, 0, fmt.Errorf("ein auftrag konnte nicht gelesen werden: %w", scanError)
|
|
}
|
|
|
|
loadedJobs = append(loadedJobs, *scannedJob)
|
|
}
|
|
|
|
if rowsError := jobRows.Err(); rowsError != nil {
|
|
return nil, 0, rowsError
|
|
}
|
|
|
|
// Die Quellen werden nachgeladen, damit die Liste nicht durch einen Verbund
|
|
// mit einer Zeile je Quelle aufgebläht wird.
|
|
for jobIndex := range loadedJobs {
|
|
if sourceError := store.loadJobSources(listContext, &loadedJobs[jobIndex]); sourceError != nil {
|
|
return nil, 0, sourceError
|
|
}
|
|
}
|
|
|
|
return loadedJobs, totalCount, nil
|
|
}
|
|
|
|
// scanJobRowWithTotal liest eine Auftragszeile samt Gesamtzahl.
|
|
func scanJobRowWithTotal(jobRows pgx.Rows, totalCount *int) (*Job, error) {
|
|
var (
|
|
loadedJob Job
|
|
statusText string
|
|
priorityText string
|
|
scheduleJSON []byte
|
|
retryJSON []byte
|
|
rpoSeconds *int64
|
|
rtoSeconds *int64
|
|
bandwidthLimit *int64
|
|
lastOutcomeText *string
|
|
descriptionValue string
|
|
// Leerer Text und NULL bedeuten beide "incremental" — der Standard,
|
|
// der auch fuer Auftraege aus der Zeit vor dieser Spalte gilt.
|
|
backupModeText string
|
|
fullBackupWeekday *int16
|
|
)
|
|
|
|
scanError := jobRows.Scan(
|
|
&loadedJob.ID, &loadedJob.Name, &descriptionValue, &statusText, &priorityText, &scheduleJSON,
|
|
&loadedJob.RepositoryID, &loadedJob.RetentionPolicyID, &rpoSeconds, &rtoSeconds,
|
|
&bandwidthLimit, &loadedJob.MaximumConcurrency, &retryJSON,
|
|
&backupModeText, &fullBackupWeekday,
|
|
&loadedJob.NextRunAt, &loadedJob.LastRunAt, &lastOutcomeText, &loadedJob.PausedAt,
|
|
&loadedJob.CreatedBy, &loadedJob.CreatedAt, &loadedJob.UpdatedAt,
|
|
totalCount,
|
|
)
|
|
if scanError != nil {
|
|
return nil, scanError
|
|
}
|
|
|
|
loadedJob.Description = descriptionValue
|
|
loadedJob.Status = JobStatus(statusText)
|
|
loadedJob.Priority = scheduler.Priority(priorityText)
|
|
loadedJob.BackupMode = normalizeBackupMode(BackupMode(backupModeText))
|
|
loadedJob.FullBackupWeekday = weekdayFromDatabase(fullBackupWeekday)
|
|
|
|
if unmarshalError := json.Unmarshal(scheduleJSON, &loadedJob.Schedule); unmarshalError != nil {
|
|
return nil, unmarshalError
|
|
}
|
|
|
|
if len(retryJSON) > 0 {
|
|
_ = json.Unmarshal(retryJSON, &loadedJob.RetryPolicy)
|
|
}
|
|
|
|
if loadedJob.RetryPolicy.MaximumAttempts < 1 {
|
|
loadedJob.RetryPolicy = scheduler.DefaultRetryPolicy()
|
|
}
|
|
|
|
if rpoSeconds != nil {
|
|
loadedJob.RecoveryPointObjective = time.Duration(*rpoSeconds) * time.Second
|
|
}
|
|
|
|
if rtoSeconds != nil {
|
|
loadedJob.RecoveryTimeObjective = time.Duration(*rtoSeconds) * time.Second
|
|
}
|
|
|
|
if bandwidthLimit != nil {
|
|
loadedJob.BandwidthLimitBytesPerSecond = *bandwidthLimit
|
|
}
|
|
|
|
if lastOutcomeText != nil {
|
|
loadedJob.LastOutcome = scheduler.JobOutcome(*lastOutcomeText)
|
|
}
|
|
|
|
return &loadedJob, nil
|
|
}
|
|
|
|
// SetJobStatus ändert den Zustand eines Auftrags.
|
|
func (store *PostgresStore) SetJobStatus(updateContext context.Context, jobIdentifier uuid.UUID, newStatus JobStatus, actingUser *uuid.UUID) error {
|
|
const updateStatement = `
|
|
UPDATE backup_jobs
|
|
SET status = $2,
|
|
paused_at = CASE WHEN $2 = 'paused' THEN now() ELSE NULL END,
|
|
-- Die ausdrückliche Typangabe ist nötig: In einem CASE-Ausdruck kann
|
|
-- PostgreSQL den Typ eines NULL-Parameters nicht aus dem Zielfeld
|
|
-- ableiten und nimmt text an — was am UUID-Feld scheitert.
|
|
paused_by = CASE WHEN $2 = 'paused' THEN $3::uuid ELSE NULL END,
|
|
updated_at = now()
|
|
WHERE id = $1 AND deleted_at IS NULL`
|
|
|
|
commandTag, execError := store.connectionPool.Exec(updateContext, updateStatement,
|
|
jobIdentifier, string(newStatus), actingUser)
|
|
if execError != nil {
|
|
return fmt.Errorf("der zustand konnte nicht geändert werden: %w", execError)
|
|
}
|
|
|
|
if commandTag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%w: %s", ErrJobNotFound, jobIdentifier)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// SetNextRun hinterlegt den berechneten nächsten Zeitpunkt.
|
|
func (store *PostgresStore) SetNextRun(updateContext context.Context, jobIdentifier uuid.UUID, nextRunAt *time.Time) error {
|
|
const updateStatement = `UPDATE backup_jobs SET next_run_at = $2, updated_at = now() WHERE id = $1`
|
|
|
|
if _, execError := store.connectionPool.Exec(updateContext, updateStatement,
|
|
jobIdentifier, nextRunAt); execError != nil {
|
|
return fmt.Errorf("der nächste zeitpunkt konnte nicht gespeichert werden: %w", execError)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// SoftDeleteJob löscht einen Auftrag, ohne seine Historie mitzunehmen.
|
|
//
|
|
// Ein hart gelöschter Auftrag risse seine Läufe mit — und damit den Nachweis,
|
|
// dass gesichert wurde. Das Audit verlangt das Gegenteil.
|
|
func (store *PostgresStore) SoftDeleteJob(deleteContext context.Context, jobIdentifier uuid.UUID) error {
|
|
const updateStatement = `
|
|
UPDATE backup_jobs
|
|
SET deleted_at = now(), status = 'disabled', next_run_at = NULL, updated_at = now()
|
|
WHERE id = $1 AND deleted_at IS NULL`
|
|
|
|
commandTag, execError := store.connectionPool.Exec(deleteContext, updateStatement, jobIdentifier)
|
|
if execError != nil {
|
|
return fmt.Errorf("der auftrag konnte nicht gelöscht werden: %w", execError)
|
|
}
|
|
|
|
if commandTag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%w: %s", ErrJobNotFound, jobIdentifier)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// ClaimDueJobs übernimmt fällige Aufträge zur Ausführung.
|
|
//
|
|
// Das ist die kritischste Abfrage des Schedulers. Sie muss zwei Dinge zugleich
|
|
// leisten: fällige Aufträge finden **und** sie so übernehmen, dass kein zweiter
|
|
// Control-Server dieselben bekommt.
|
|
//
|
|
// Gelöst über FOR UPDATE SKIP LOCKED: Der erste Server sperrt die Zeilen, der
|
|
// zweite überspringt sie, statt zu warten. Ohne SKIP LOCKED bliebe der zweite
|
|
// Server bis zum Ende der ersten Transaktion stehen; ohne FOR UPDATE bekämen
|
|
// beide dieselben Aufträge.
|
|
func (store *PostgresStore) ClaimDueJobs(claimContext context.Context, currentTime time.Time, maximumJobs int, schedulerInstance string) ([]uuid.UUID, 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 selectDueStatement = `
|
|
SELECT id FROM backup_jobs
|
|
WHERE status = 'active'
|
|
AND deleted_at IS NULL
|
|
AND next_run_at IS NOT NULL
|
|
AND next_run_at <= $1
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM backup_job_runs
|
|
WHERE backup_job_runs.job_id = backup_jobs.id
|
|
AND backup_job_runs.status IN ('queued','running')
|
|
)
|
|
ORDER BY next_run_at
|
|
LIMIT $2
|
|
FOR UPDATE OF backup_jobs SKIP LOCKED`
|
|
|
|
dueRows, queryError := transaction.Query(claimContext, selectDueStatement, currentTime, maximumJobs)
|
|
if queryError != nil {
|
|
return nil, fmt.Errorf("die fälligen aufträge konnten nicht ermittelt werden: %w", queryError)
|
|
}
|
|
|
|
claimedJobIDs := make([]uuid.UUID, 0, maximumJobs)
|
|
|
|
for dueRows.Next() {
|
|
var jobIdentifier uuid.UUID
|
|
|
|
if scanError := dueRows.Scan(&jobIdentifier); scanError != nil {
|
|
dueRows.Close()
|
|
return nil, fmt.Errorf("ein fälliger auftrag konnte nicht gelesen werden: %w", scanError)
|
|
}
|
|
|
|
claimedJobIDs = append(claimedJobIDs, jobIdentifier)
|
|
}
|
|
|
|
dueRows.Close()
|
|
|
|
if rowsError := dueRows.Err(); rowsError != nil {
|
|
return nil, rowsError
|
|
}
|
|
|
|
// Der Lauf wird innerhalb derselben Transaktion angelegt. Damit ist die
|
|
// Übernahme abgeschlossen, sobald die Transaktion festgeschrieben ist —
|
|
// ein Absturz dazwischen hinterlässt keinen halb übernommenen Auftrag.
|
|
const insertRunStatement = `
|
|
INSERT INTO backup_job_runs (job_id, status, trigger, scheduled_for, correlation_id, scheduler_instance, heartbeat_at)
|
|
SELECT $1, 'queued', 'schedule', next_run_at, gen_random_uuid(), $2, now()
|
|
FROM backup_jobs WHERE id = $1`
|
|
|
|
for _, jobIdentifier := range claimedJobIDs {
|
|
if _, execError := transaction.Exec(claimContext, insertRunStatement,
|
|
jobIdentifier, schedulerInstance); execError != nil {
|
|
if isUniqueViolation(execError) {
|
|
// Ein anderer Server war schneller. Das ist kein Fehler,
|
|
// sondern genau der Fall, den der Teilindex verhindern soll.
|
|
continue
|
|
}
|
|
|
|
return nil, fmt.Errorf("der lauf für %s konnte nicht angelegt werden: %w", jobIdentifier, execError)
|
|
}
|
|
}
|
|
|
|
if commitError := transaction.Commit(claimContext); commitError != nil {
|
|
return nil, fmt.Errorf("die übernahme konnte nicht festgeschrieben werden: %w", commitError)
|
|
}
|
|
|
|
return claimedJobIDs, nil
|
|
}
|
|
|
|
// isUniqueViolation erkennt einen Eindeutigkeitsverstoß.
|
|
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
|
|
}
|
|
|
|
// nullableSeconds wandelt eine Dauer in Sekunden oder NULL.
|
|
func nullableSeconds(duration time.Duration) *int64 {
|
|
if duration <= 0 {
|
|
return nil
|
|
}
|
|
|
|
secondsValue := int64(duration.Seconds())
|
|
|
|
return &secondsValue
|
|
}
|
|
|
|
// nullableBandwidth wandelt eine Bandbreitengrenze in einen Wert oder NULL.
|
|
func nullableBandwidth(bytesPerSecond int64) *int64 {
|
|
if bytesPerSecond <= 0 {
|
|
return nil
|
|
}
|
|
|
|
return &bytesPerSecond
|
|
}
|
|
|
|
// defaultToEmptySlice ersetzt nil durch eine leere Liste.
|
|
//
|
|
// Ohne diesen Schritt landete "null" statt "[]" in der Spalte — und das Lesen
|
|
// ergäbe eine nil-Liste, die sich beim Vergleich anders verhält als eine leere.
|
|
func defaultToEmptySlice(patternList []string) []string {
|
|
if patternList == nil {
|
|
return []string{}
|
|
}
|
|
|
|
return patternList
|
|
}
|
|
|
|
// normalizeBackupMode fuellt eine fehlende Angabe mit dem Standard.
|
|
//
|
|
// Leer bedeutet `incremental`. Das gilt fuer neue Auftraege ohne Angabe **und**
|
|
// fuer alle, die vor der Migration 000014 entstanden sind — ihre Spalte traegt
|
|
// den Vorgabewert, aber ein leerer Wert aus einem Aufrufer darf nicht zu einem
|
|
// ungueltigen Zustand fuehren.
|
|
func normalizeBackupMode(requestedMode BackupMode) BackupMode {
|
|
if requestedMode == BackupModeAlwaysFull {
|
|
return BackupModeAlwaysFull
|
|
}
|
|
|
|
return BackupModeIncremental
|
|
}
|
|
|
|
// nullableWeekday uebersetzt einen Wochentag fuer die Datenbank.
|
|
//
|
|
// nil bedeutet: kein erzwungener Volltag. Ein Wochentag ausserhalb von 0..6
|
|
// wird verworfen statt gekappt — ein stillschweigend auf Sonntag gesetzter
|
|
// Montag waere ein Fehler, den niemand bemerkt.
|
|
func nullableWeekday(requestedWeekday *time.Weekday) *int16 {
|
|
if requestedWeekday == nil {
|
|
return nil
|
|
}
|
|
|
|
if *requestedWeekday < time.Sunday || *requestedWeekday > time.Saturday {
|
|
return nil
|
|
}
|
|
|
|
storedValue := int16(*requestedWeekday)
|
|
|
|
return &storedValue
|
|
}
|
|
|
|
// weekdayFromDatabase uebersetzt einen gespeicherten Wochentag zurueck.
|
|
func weekdayFromDatabase(storedValue *int16) *time.Weekday {
|
|
if storedValue == nil {
|
|
return nil
|
|
}
|
|
|
|
if *storedValue < 0 || *storedValue > 6 {
|
|
return nil
|
|
}
|
|
|
|
loadedWeekday := time.Weekday(*storedValue)
|
|
|
|
return &loadedWeekday
|
|
}
|