package jobs import ( "context" "errors" "fmt" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5" "github.com/syncova/syncova/packages/scheduler" ) // RunStatus ist der Zustand eines Laufs. type RunStatus string const ( // RunStatusQueued ist übernommen, aber noch nicht begonnen. RunStatusQueued RunStatus = "queued" // RunStatusRunning ist in Ausführung. RunStatusRunning RunStatus = "running" // RunStatusSucceeded ist vollständig gelungen. RunStatusSucceeded RunStatus = "succeeded" // RunStatusPartialFailure hat Objekte übergangen. // // Ausdrücklich kein Erfolg (PROMPT.md §138). RunStatusPartialFailure RunStatus = "partial_failure" // RunStatusFailed ist gescheitert. RunStatusFailed RunStatus = "failed" // RunStatusCancelled wurde abgebrochen. RunStatusCancelled RunStatus = "cancelled" ) // RunTrigger benennt den Auslöser eines Laufs. type RunTrigger string const ( // TriggerSchedule ist der Zeitplan. TriggerSchedule RunTrigger = "schedule" // TriggerManual ist eine ausdrückliche Anforderung. TriggerManual RunTrigger = "manual" // TriggerRetry ist ein Wiederholungsversuch. TriggerRetry RunTrigger = "retry" // TriggerDependency ist die Folge eines vorausgesetzten Auftrags. TriggerDependency RunTrigger = "dependency" ) // Run ist ein einzelner Ausführungsvorgang eines Auftrags. type Run struct { // ID ist der öffentliche Bezeichner. ID uuid.UUID `json:"id"` // JobID ist der ausgeführte Auftrag. JobID uuid.UUID `json:"job_id"` // Status ist der Zustand. Status RunStatus `json:"status"` // Trigger benennt den Auslöser. Trigger RunTrigger `json:"trigger"` // AttemptNumber ist die Nummer des Versuchs, beginnend bei 1. AttemptNumber int `json:"attempt_number"` // ScheduledFor ist der geplante Zeitpunkt in UTC. ScheduledFor *time.Time `json:"scheduled_for,omitempty"` // StartedAt ist der tatsächliche Beginn in UTC. StartedAt *time.Time `json:"started_at,omitempty"` // CompletedAt ist das Ende in UTC. CompletedAt *time.Time `json:"completed_at,omitempty"` // BytesProcessed ist die gelesene Datenmenge. BytesProcessed int64 `json:"bytes_processed"` // BytesWritten ist die abgelegte Datenmenge. BytesWritten int64 `json:"bytes_written"` // FilesProcessed ist die Zahl bearbeiteter Objekte. FilesProcessed int64 `json:"files_processed"` // FilesSkipped ist die Zahl übergangener Objekte. FilesSkipped int64 `json:"files_skipped"` // ErrorCode ist die Fehlerkennung. ErrorCode string `json:"error_code,omitempty"` // ErrorMessage ist die verständliche Fehlermeldung. ErrorMessage string `json:"error_message,omitempty"` // FailureClass ordnet den Fehler ein. FailureClass scheduler.FailureClass `json:"failure_class,omitempty"` // CorrelationID verbindet den Lauf mit seinen Protokollzeilen. CorrelationID uuid.UUID `json:"correlation_id"` // SchedulerInstance benennt den ausführenden Control-Server. SchedulerInstance string `json:"scheduler_instance,omitempty"` // CreatedAt ist der Anlagezeitpunkt in UTC. CreatedAt time.Time `json:"created_at"` } // Duration liefert die Dauer eines abgeschlossenen Laufs. func (run *Run) Duration() time.Duration { if run.StartedAt == nil || run.CompletedAt == nil { return 0 } return run.CompletedAt.Sub(*run.StartedAt) } // IsFinished meldet einen abgeschlossenen Lauf. func (run *Run) IsFinished() bool { return run.Status != RunStatusQueued && run.Status != RunStatusRunning } // ErrRunNotFound meldet einen nicht vorhandenen Lauf. var ErrRunNotFound = errors.New("der sicherungslauf wurde nicht gefunden") // runColumnList sind die Spalten eines Laufs in fester Reihenfolge. const runColumnList = ` id, job_id, status, trigger, attempt_number, scheduled_for, started_at, completed_at, bytes_processed, bytes_written, files_processed, files_skipped, error_code, error_message, failure_class, correlation_id, scheduler_instance, created_at` // scanRun liest eine Laufzeile. func scanRun(scanner rowScanner) (*Run, error) { var ( loadedRun Run statusText string triggerText string errorCode *string errorMessage *string failureClass *string schedulerInstance *string ) scanError := scanner.Scan( &loadedRun.ID, &loadedRun.JobID, &statusText, &triggerText, &loadedRun.AttemptNumber, &loadedRun.ScheduledFor, &loadedRun.StartedAt, &loadedRun.CompletedAt, &loadedRun.BytesProcessed, &loadedRun.BytesWritten, &loadedRun.FilesProcessed, &loadedRun.FilesSkipped, &errorCode, &errorMessage, &failureClass, &loadedRun.CorrelationID, &schedulerInstance, &loadedRun.CreatedAt, ) if scanError != nil { return nil, scanError } loadedRun.Status = RunStatus(statusText) loadedRun.Trigger = RunTrigger(triggerText) if errorCode != nil { loadedRun.ErrorCode = *errorCode } if errorMessage != nil { loadedRun.ErrorMessage = *errorMessage } if failureClass != nil { loadedRun.FailureClass = scheduler.FailureClass(*failureClass) } if schedulerInstance != nil { loadedRun.SchedulerInstance = *schedulerInstance } return &loadedRun, nil } // CreateManualRun legt einen ausdrücklich angeforderten Lauf an. // // Der Teilindex der Datenbank lässt nur einen aktiven Lauf je Auftrag zu. Wer // einen bereits laufenden Auftrag erneut anstößt, bekommt deshalb einen // verständlichen Fehler statt eines zweiten Laufs. func (store *PostgresStore) CreateManualRun(createContext context.Context, jobIdentifier uuid.UUID, triggeredBy *uuid.UUID) (*Run, error) { const insertStatement = ` INSERT INTO backup_job_runs (job_id, status, trigger, scheduled_for, correlation_id, triggered_by) VALUES ($1, 'queued', 'manual', now(), gen_random_uuid(), $2) RETURNING ` + runColumnList createdRun, scanError := scanRun(store.connectionPool.QueryRow(createContext, insertStatement, jobIdentifier, triggeredBy)) if scanError != nil { if isUniqueViolation(scanError) { return nil, ErrRunAlreadyActive } return nil, fmt.Errorf("der lauf konnte nicht angelegt werden: %w", scanError) } return createdRun, nil } // StartRun setzt einen übernommenen Lauf auf „läuft". // // Die Bedingung auf den bisherigen Zustand ist kein Beiwerk: Sie verhindert, // dass ein zurückgewonnener Lauf von seinem alten, inzwischen wiederbelebten // Server erneut gestartet wird. func (store *PostgresStore) StartRun(startContext context.Context, runIdentifier uuid.UUID, schedulerInstance string) error { const updateStatement = ` UPDATE backup_job_runs SET status = 'running', started_at = now(), heartbeat_at = now(), scheduler_instance = $2 WHERE id = $1 AND status = 'queued'` commandTag, execError := store.connectionPool.Exec(startContext, updateStatement, runIdentifier, schedulerInstance) if execError != nil { return fmt.Errorf("der lauf konnte nicht gestartet werden: %w", execError) } if commandTag.RowsAffected() == 0 { return fmt.Errorf("%w: %s steht nicht mehr auf 'queued'", ErrRunNotFound, runIdentifier) } return nil } // RecordHeartbeat meldet einen laufenden Lauf als lebendig. // // Ohne diese Meldung liesse sich ein abgestürzter Control-Server nicht von // einem langsamen unterscheiden — und ein zurückgewonnener Lauf könnte einen // noch laufenden verdoppeln. func (store *PostgresStore) RecordHeartbeat(heartbeatContext context.Context, runIdentifier uuid.UUID) error { const updateStatement = ` UPDATE backup_job_runs SET heartbeat_at = now() WHERE id = $1 AND status = 'running'` if _, execError := store.connectionPool.Exec(heartbeatContext, updateStatement, runIdentifier); execError != nil { return fmt.Errorf("die lebendmeldung des laufs schlug fehl: %w", execError) } return nil } // RunOutcome beschreibt das Ergebnis eines Laufs. type RunOutcome struct { // Status ist der erreichte Zustand. Status RunStatus // BytesProcessed ist die gelesene Datenmenge. BytesProcessed int64 // BytesWritten ist die abgelegte Datenmenge. BytesWritten int64 // FilesProcessed ist die Zahl bearbeiteter Objekte. FilesProcessed int64 // FilesSkipped ist die Zahl übergangener Objekte. FilesSkipped int64 // ErrorCode ist die Fehlerkennung. ErrorCode string // ErrorMessage ist die verständliche Fehlermeldung. // // Sie enthält niemals Geheimnisse; die Redaktion geschieht vor dem Aufruf. ErrorMessage string // FailureClass ordnet den Fehler ein und entscheidet über Wiederholung. FailureClass scheduler.FailureClass } // FinishRun schreibt das Ergebnis eines Laufs fest. // // Der Ausgang wird zugleich am Auftrag vermerkt: Die Übersicht braucht ihn ohne // Verbund über eine wachsende Historientabelle, und die Abhängigkeitsprüfung // liest ihn dort. func (store *PostgresStore) FinishRun(finishContext context.Context, runIdentifier uuid.UUID, outcome RunOutcome) error { transaction, transactionError := store.connectionPool.Begin(finishContext) if transactionError != nil { return fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError) } defer func() { _ = transaction.Rollback(finishContext) }() const updateRunStatement = ` UPDATE backup_job_runs SET status = $2, completed_at = now(), bytes_processed = $3::bigint, bytes_written = $4, files_processed = $5, files_skipped = $6, error_code = $7, error_message = $8, failure_class = $9, -- Die ausdrückliche Typangabe ist nötig: $3 steht zugleich in einer -- Zuweisung an eine BIGINT-Spalte und in dieser Division, deren -- anderer Operand numeric ist. Ohne den Cast leitet PostgreSQL zwei -- verschiedene Typen für denselben Parameter ab und lehnt die -- gesamte Anweisung ab — jeder Lauf bliebe auf „running" stehen. throughput_bps = CASE WHEN started_at IS NOT NULL AND EXTRACT(EPOCH FROM (now() - started_at)) > 0 THEN ($3::bigint / EXTRACT(EPOCH FROM (now() - started_at)))::BIGINT ELSE NULL END WHERE id = $1 RETURNING job_id, started_at` var jobIdentifier uuid.UUID var startedAt *time.Time scanError := transaction.QueryRow(finishContext, updateRunStatement, runIdentifier, string(outcome.Status), outcome.BytesProcessed, outcome.BytesWritten, outcome.FilesProcessed, outcome.FilesSkipped, nullableText(outcome.ErrorCode), nullableText(outcome.ErrorMessage), nullableText(string(outcome.FailureClass)), ).Scan(&jobIdentifier, &startedAt) if errors.Is(scanError, pgx.ErrNoRows) { return fmt.Errorf("%w: %s", ErrRunNotFound, runIdentifier) } if scanError != nil { return fmt.Errorf("das ergebnis konnte nicht festgeschrieben werden: %w", scanError) } const updateJobStatement = ` UPDATE backup_jobs SET last_run_at = $2, last_outcome = $3, updated_at = now() WHERE id = $1` if _, execError := transaction.Exec(finishContext, updateJobStatement, jobIdentifier, startedAt, string(outcome.Status)); execError != nil { return fmt.Errorf("der ausgang konnte nicht am auftrag vermerkt werden: %w", execError) } if commitError := transaction.Commit(finishContext); commitError != nil { return fmt.Errorf("das ergebnis konnte nicht festgeschrieben werden: %w", commitError) } return nil } // ScheduleRetry legt einen Wiederholungsversuch an. // // Der neue Lauf entsteht erst, nachdem der gescheiterte abgeschlossen wurde — // sonst verhinderte der Teilindex ihn, der nur einen aktiven Lauf je Auftrag // zulässt. func (store *PostgresStore) ScheduleRetry(retryContext context.Context, jobIdentifier uuid.UUID, attemptNumber int, retryAt time.Time) (*Run, error) { const insertStatement = ` INSERT INTO backup_job_runs (job_id, status, trigger, attempt_number, scheduled_for, correlation_id) VALUES ($1, 'queued', 'retry', $2, $3, gen_random_uuid()) RETURNING ` + runColumnList createdRun, scanError := scanRun(store.connectionPool.QueryRow(retryContext, insertStatement, jobIdentifier, attemptNumber, retryAt)) if scanError != nil { if isUniqueViolation(scanError) { return nil, ErrRunAlreadyActive } return nil, fmt.Errorf("der wiederholungsversuch konnte nicht angelegt werden: %w", scanError) } return createdRun, nil } // ClaimQueuedRuns übernimmt anstehende Läufe zur Ausführung. // // Anders als ClaimDueJobs arbeitet diese Abfrage auf bereits angelegten Läufen: // Sie holt, was der Zeitplan, ein Anwender oder eine Wiederholung eingereiht // hat. Auch hier sperrt FOR UPDATE SKIP LOCKED, damit zwei Control-Server nicht // denselben Lauf ausführen. func (store *PostgresStore) ClaimQueuedRuns(claimContext context.Context, currentTime time.Time, maximumRuns int, schedulerInstance string) ([]Run, 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 ` + runColumnList + ` FROM backup_job_runs WHERE status = 'queued' AND (scheduled_for IS NULL OR scheduled_for <= $1) AND job_id IN (SELECT id FROM backup_jobs WHERE status = 'active' AND deleted_at IS NULL) ORDER BY scheduled_for NULLS FIRST LIMIT $2 FOR UPDATE SKIP LOCKED` runRows, queryError := transaction.Query(claimContext, selectStatement, currentTime, maximumRuns) if queryError != nil { return nil, fmt.Errorf("die anstehenden läufe konnten nicht ermittelt werden: %w", queryError) } claimedRuns := make([]Run, 0, maximumRuns) for runRows.Next() { claimedRun, scanError := scanRun(runRows) if scanError != nil { runRows.Close() return nil, fmt.Errorf("ein anstehender lauf konnte nicht gelesen werden: %w", scanError) } claimedRuns = append(claimedRuns, *claimedRun) } runRows.Close() if rowsError := runRows.Err(); rowsError != nil { return nil, rowsError } // Die Übernahme wird im selben Vorgang vermerkt. Damit kann kein zweiter // Server dieselben Läufe bekommen, sobald die Transaktion steht. const markStatement = ` UPDATE backup_job_runs SET scheduler_instance = $2, heartbeat_at = now() WHERE id = $1` for _, claimedRun := range claimedRuns { if _, execError := transaction.Exec(claimContext, markStatement, claimedRun.ID, schedulerInstance); execError != nil { return nil, fmt.Errorf("die übernahme des laufs %s schlug fehl: %w", claimedRun.ID, execError) } } if commitError := transaction.Commit(claimContext); commitError != nil { return nil, fmt.Errorf("die übernahme konnte nicht festgeschrieben werden: %w", commitError) } return claimedRuns, nil } // ReclaimStaleRuns gibt Läufe abgestürzter Control-Server frei. // // Das ist der wichtigste Selbstheilungsmechanismus der Ausführung. Stirbt ein // Server mitten im Lauf, bleibt dessen Zeile auf „running" stehen. Der // Teilindex lässt dann keinen weiteren Lauf dieses Auftrags zu — die Sicherung // fiele **dauerhaft** aus, ohne dass jemand einen Fehler sähe. Es gäbe nur // einen Auftrag, der nie wieder läuft. // // Ein verwaister Lauf wird deshalb als gescheitert markiert, nicht gelöscht: // Der Abbruch ist ein Ereignis, das ins Protokoll gehört. Die Fehlerklasse // „transient" sorgt dafür, dass die Wiederholungsstrategie greift. func (store *PostgresStore) ReclaimStaleRuns(reclaimContext context.Context, staleAfter time.Duration) (int, error) { const updateStatement = ` UPDATE backup_job_runs SET status = 'failed', completed_at = now(), error_code = 'SCHEDULER_LOST', error_message = $2, failure_class = 'transient' WHERE status IN ('queued', 'running') AND heartbeat_at IS NOT NULL AND heartbeat_at < now() - $1::interval RETURNING job_id, started_at` transaction, transactionError := store.connectionPool.Begin(reclaimContext) if transactionError != nil { return 0, fmt.Errorf("die transaktion konnte nicht begonnen werden: %w", transactionError) } defer func() { _ = transaction.Rollback(reclaimContext) }() const explanation = "Der ausführende Control-Server meldete sich nicht mehr. " + "Der Lauf wurde freigegeben, damit der Auftrag nicht dauerhaft blockiert bleibt." reclaimRows, queryError := transaction.Query(reclaimContext, updateStatement, staleAfter.String(), explanation) if queryError != nil { return 0, fmt.Errorf("verwaiste läufe konnten nicht freigegeben werden: %w", queryError) } type reclaimedRun struct { jobIdentifier uuid.UUID startedAt *time.Time } reclaimedRuns := make([]reclaimedRun, 0) for reclaimRows.Next() { var currentRun reclaimedRun if scanError := reclaimRows.Scan(¤tRun.jobIdentifier, ¤tRun.startedAt); scanError != nil { reclaimRows.Close() return 0, fmt.Errorf("ein verwaister lauf konnte nicht gelesen werden: %w", scanError) } reclaimedRuns = append(reclaimedRuns, currentRun) } reclaimRows.Close() if rowsError := reclaimRows.Err(); rowsError != nil { return 0, rowsError } const updateJobStatement = ` UPDATE backup_jobs SET last_outcome = 'failed', updated_at = now() WHERE id = $1` for _, currentRun := range reclaimedRuns { if _, execError := transaction.Exec(reclaimContext, updateJobStatement, currentRun.jobIdentifier); execError != nil { return 0, fmt.Errorf("der ausgang konnte nicht am auftrag vermerkt werden: %w", execError) } } if commitError := transaction.Commit(reclaimContext); commitError != nil { return 0, fmt.Errorf("die freigabe konnte nicht festgeschrieben werden: %w", commitError) } return len(reclaimedRuns), nil } // CancelRun bricht einen Lauf ab. func (store *PostgresStore) CancelRun(cancelContext context.Context, runIdentifier uuid.UUID, reason string) error { const updateStatement = ` UPDATE backup_job_runs 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, runIdentifier, reason) if execError != nil { return fmt.Errorf("der lauf konnte nicht abgebrochen werden: %w", execError) } if commandTag.RowsAffected() == 0 { return fmt.Errorf("%w: %s läuft nicht mehr", ErrRunNotFound, runIdentifier) } return nil } // GetRun liest einen einzelnen Lauf. func (store *PostgresStore) GetRun(readContext context.Context, runIdentifier uuid.UUID) (*Run, error) { const selectStatement = `SELECT ` + runColumnList + ` FROM backup_job_runs WHERE id = $1` loadedRun, scanError := scanRun(store.connectionPool.QueryRow(readContext, selectStatement, runIdentifier)) if errors.Is(scanError, pgx.ErrNoRows) { return nil, fmt.Errorf("%w: %s", ErrRunNotFound, runIdentifier) } if scanError != nil { return nil, fmt.Errorf("der lauf konnte nicht gelesen werden: %w", scanError) } return loadedRun, nil } // ListRuns liefert die Läufe eines Auftrags, neueste zuerst. func (store *PostgresStore) ListRuns(listContext context.Context, jobIdentifier uuid.UUID, page int, pageSize int) ([]Run, int, error) { pageFilter := ListFilter{Page: page, PageSize: pageSize} pageFilter.normalize() const selectStatement = ` SELECT ` + runColumnList + `, COUNT(*) OVER () AS total_count FROM backup_job_runs WHERE job_id = $1 ORDER BY created_at DESC LIMIT $2 OFFSET $3` runRows, queryError := store.connectionPool.Query(listContext, selectStatement, jobIdentifier, pageFilter.PageSize, (pageFilter.Page-1)*pageFilter.PageSize) if queryError != nil { return nil, 0, fmt.Errorf("die laufliste konnte nicht gelesen werden: %w", queryError) } defer runRows.Close() loadedRuns := make([]Run, 0, pageFilter.PageSize) var totalCount int for runRows.Next() { var ( loadedRun Run statusText string triggerText string errorCode *string errorMessage *string failureClass *string schedulerInstance *string ) if scanError := runRows.Scan( &loadedRun.ID, &loadedRun.JobID, &statusText, &triggerText, &loadedRun.AttemptNumber, &loadedRun.ScheduledFor, &loadedRun.StartedAt, &loadedRun.CompletedAt, &loadedRun.BytesProcessed, &loadedRun.BytesWritten, &loadedRun.FilesProcessed, &loadedRun.FilesSkipped, &errorCode, &errorMessage, &failureClass, &loadedRun.CorrelationID, &schedulerInstance, &loadedRun.CreatedAt, &totalCount, ); scanError != nil { return nil, 0, fmt.Errorf("ein lauf konnte nicht gelesen werden: %w", scanError) } loadedRun.Status = RunStatus(statusText) loadedRun.Trigger = RunTrigger(triggerText) if errorCode != nil { loadedRun.ErrorCode = *errorCode } if errorMessage != nil { loadedRun.ErrorMessage = *errorMessage } if failureClass != nil { loadedRun.FailureClass = scheduler.FailureClass(*failureClass) } if schedulerInstance != nil { loadedRun.SchedulerInstance = *schedulerInstance } loadedRuns = append(loadedRuns, loadedRun) } if rowsError := runRows.Err(); rowsError != nil { return nil, 0, rowsError } return loadedRuns, totalCount, nil }