Enterprise-Backup-, Recovery-, Verification-, Security- und Monitoring-Plattform fuer Proxmox VE, Windows, Linux und Dateisysteme. Der Leitsatz, der fast jede Entscheidung erklaert: Ein Backup gilt erst als vertrauenswuerdig, wenn Integritaet geprueft und Wiederherstellbarkeit nachgewiesen wurde. Deshalb steigt ein Wiederherstellungspunkt erst nach einem tatsaechlich durchgefuehrten Restore-Test auf "recoverable", und Unbekanntes geht in keine Bewertung als "gut" ein. Umfang (Phasen 0-23): - Repository Engine: inhaltsadressierte Bloecke, atomares Commit-Protokoll, Katalogaufbau allein aus den Manifesten — ohne Datenbank - Backup Engine: inhaltsabhaengiges Chunking, Deduplizierung trotz Verschluesselung, zstd, AES-256-GCM, Streaming mit Gegendruck - Agenten fuer Windows und Linux mit Auftragsabholung (Pull-Modell) - Proxmox-Provider mit beiden Zugriffswegen auf die Sicherungsarchive - Scheduler, Recovery Engine mit Pruefpunkt, Verification, Unveraenderlichkeit - Weboberflaeche, Kennzahlen, Meldungen, Berichte, Security Center, Ransomware-Heuristik (meldet, handelt nie) - Disaster Recovery, Haertung, Leistungsmessung, Chaos Testing - Eingefrorene Vertraege fuer API, Migrationen, Backup-Format und Repository - Auslieferungspaket fuer linux/amd64, linux/arm64 und windows/amd64 Nicht enthalten und als solches gekennzeichnet: Kapazitaetsprognose, Backup Copy, Changed Block Tracking bei Proxmox, erweiterte Attribute und ACLs. Gebaut, aber nie auf echter Hardware gefahren: der Windows-Dienst, die systemd-Einheit und der verpflichtende Proxmox-Meilenstein — ob eine wiederhergestellte VM startet, ist ungeprueft. Einzelheiten in CHANGELOG.md und docs/release-candidate.md. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
353 lines
12 KiB
Go
353 lines
12 KiB
Go
package agenttasks
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
// ErrNoTaskAvailable meldet, dass für einen Agenten nichts anliegt.
|
|
//
|
|
// Ein eigener Fehler und kein leeres Ergebnis: Der Agent fragt im Sekundentakt,
|
|
// und „nichts zu tun" ist der Normalfall — er darf sich nicht von einem echten
|
|
// Fehler unterscheiden lassen müssen.
|
|
var ErrNoTaskAvailable = errors.New("fuer diesen agenten liegt kein auftrag vor")
|
|
|
|
// ErrTaskNotOwned meldet einen Zugriff auf einen fremden Auftrag.
|
|
var ErrTaskNotOwned = errors.New("dieser auftrag gehoert einem anderen agenten")
|
|
|
|
// Store ist die Datenzugriffsschicht der Agentenaufträge.
|
|
type Store struct {
|
|
// connectionPool ist der Datenbankpool der Control Plane.
|
|
connectionPool *pgxpool.Pool
|
|
}
|
|
|
|
// NewStore erzeugt den Auftragsspeicher.
|
|
func NewStore(connectionPool *pgxpool.Pool) *Store {
|
|
return &Store{connectionPool: connectionPool}
|
|
}
|
|
|
|
// CreateBackupTask stellt einen Sicherungsauftrag ein.
|
|
func (store *Store) CreateBackupTask(createContext context.Context, agentIdentifier uuid.UUID,
|
|
jobRunIdentifier uuid.UUID, jobSourceIdentifier uuid.UUID, payload BackupPayload) (*Task, error) {
|
|
encodedPayload, encodeError := json.Marshal(payload)
|
|
if encodeError != nil {
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht kodieren: %w", encodeError)
|
|
}
|
|
|
|
const insertStatement = `
|
|
INSERT INTO agent_tasks (agent_id, job_run_id, job_source_id, task_type, status, task_payload)
|
|
VALUES ($1, $2, $3, 'backup', 'pending', $4)
|
|
RETURNING id, created_at`
|
|
|
|
createdTask := &Task{
|
|
AgentID: agentIdentifier,
|
|
JobRunID: jobRunIdentifier,
|
|
JobSourceID: jobSourceIdentifier,
|
|
TaskType: TypeBackup,
|
|
Status: StatusPending,
|
|
Backup: payload,
|
|
}
|
|
|
|
scanError := store.connectionPool.QueryRow(createContext, insertStatement,
|
|
agentIdentifier, jobRunIdentifier, jobSourceIdentifier, encodedPayload).Scan(
|
|
&createdTask.ID, &createdTask.CreatedAt)
|
|
if scanError != nil {
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht einstellen: %w", scanError)
|
|
}
|
|
|
|
return createdTask, nil
|
|
}
|
|
|
|
// CreateRestoreTask stellt einen Wiederherstellungsauftrag ein.
|
|
//
|
|
// Die drei Hürden gegen ein versehentliches Überschreiben liegen **vor** diesem
|
|
// Aufruf: Der Server prüft Recht und wörtliche Bestätigung, bevor er den
|
|
// Auftrag einstellt. Der Agent führt aus — er entscheidet nicht.
|
|
func (store *Store) CreateRestoreTask(createContext context.Context, agentIdentifier uuid.UUID,
|
|
jobRunIdentifier uuid.UUID, jobSourceIdentifier uuid.UUID, payload RestorePayload) (*Task, error) {
|
|
encodedPayload, encodeError := json.Marshal(struct {
|
|
Restore RestorePayload `json:"restore"`
|
|
}{Restore: payload})
|
|
if encodeError != nil {
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht kodieren: %w", encodeError)
|
|
}
|
|
|
|
const insertStatement = `
|
|
INSERT INTO agent_tasks (agent_id, job_run_id, job_source_id, task_type, status, task_payload)
|
|
VALUES ($1, $2, $3, 'restore', 'pending', $4)
|
|
RETURNING id, created_at`
|
|
|
|
createdTask := &Task{
|
|
AgentID: agentIdentifier,
|
|
JobRunID: jobRunIdentifier,
|
|
JobSourceID: jobSourceIdentifier,
|
|
TaskType: TypeRestore,
|
|
Status: StatusPending,
|
|
Restore: payload,
|
|
}
|
|
|
|
scanError := store.connectionPool.QueryRow(createContext, insertStatement,
|
|
agentIdentifier, jobRunIdentifier, jobSourceIdentifier, encodedPayload).Scan(
|
|
&createdTask.ID, &createdTask.CreatedAt)
|
|
if scanError != nil {
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht einstellen: %w", scanError)
|
|
}
|
|
|
|
return createdTask, nil
|
|
}
|
|
|
|
// ClaimNextTask holt den ältesten offenen Auftrag eines Agenten ab.
|
|
//
|
|
// Die Abholung nutzt FOR UPDATE SKIP LOCKED: Fragen zwei Vorgänge gleichzeitig,
|
|
// bekommt jeder einen anderen Auftrag statt beide denselben. Dieselbe Technik
|
|
// wie bei der Auftragsschleife (Phase 8).
|
|
//
|
|
// Der Teilindex auf (agent_id) WHERE status IN ('claimed','running') sorgt
|
|
// dafür, dass ein Agent nie zwei Aufträge gleichzeitig hält — zwei Sicherungen
|
|
// desselben Systems konkurrierten um dieselben Dateien und dieselbe
|
|
// Schreibsperre.
|
|
func (store *Store) ClaimNextTask(claimContext context.Context, agentIdentifier uuid.UUID) (*Task, error) {
|
|
databaseTransaction, beginError := store.connectionPool.Begin(claimContext)
|
|
if beginError != nil {
|
|
return nil, fmt.Errorf("die transaktion liess sich nicht beginnen: %w", beginError)
|
|
}
|
|
|
|
defer func() { _ = databaseTransaction.Rollback(claimContext) }()
|
|
|
|
const selectStatement = `
|
|
SELECT id, job_run_id, job_source_id, task_type, task_payload, created_at
|
|
FROM agent_tasks
|
|
WHERE agent_id = $1 AND status = 'pending'
|
|
ORDER BY created_at
|
|
LIMIT 1
|
|
FOR UPDATE SKIP LOCKED`
|
|
|
|
var (
|
|
claimedTask Task
|
|
rawTaskType string
|
|
encodedPayload []byte
|
|
)
|
|
|
|
scanError := databaseTransaction.QueryRow(claimContext, selectStatement, agentIdentifier).Scan(
|
|
&claimedTask.ID, &claimedTask.JobRunID, &claimedTask.JobSourceID,
|
|
&rawTaskType, &encodedPayload, &claimedTask.CreatedAt)
|
|
if scanError != nil {
|
|
if errors.Is(scanError, pgx.ErrNoRows) {
|
|
return nil, ErrNoTaskAvailable
|
|
}
|
|
|
|
return nil, fmt.Errorf("der naechste auftrag liess sich nicht lesen: %w", scanError)
|
|
}
|
|
|
|
const updateStatement = `
|
|
UPDATE agent_tasks
|
|
SET status = 'claimed', claimed_at = now(), heartbeat_at = now()
|
|
WHERE id = $1`
|
|
|
|
if _, executeError := databaseTransaction.Exec(claimContext, updateStatement,
|
|
claimedTask.ID); executeError != nil {
|
|
// Der Teilindex schlägt zu, wenn der Agent bereits einen Auftrag hält.
|
|
// Das ist kein Fehler, sondern die gewollte Begrenzung auf einen.
|
|
if isUniqueViolation(executeError) {
|
|
return nil, ErrNoTaskAvailable
|
|
}
|
|
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht uebernehmen: %w", executeError)
|
|
}
|
|
|
|
if commitError := databaseTransaction.Commit(claimContext); commitError != nil {
|
|
return nil, fmt.Errorf("die uebernahme liess sich nicht abschliessen: %w", commitError)
|
|
}
|
|
|
|
claimedTask.AgentID = agentIdentifier
|
|
claimedTask.TaskType = TaskType(rawTaskType)
|
|
claimedTask.Status = StatusClaimed
|
|
|
|
if decodeError := decodeTaskPayload(encodedPayload, &claimedTask); decodeError != nil {
|
|
return nil, decodeError
|
|
}
|
|
|
|
return &claimedTask, nil
|
|
}
|
|
|
|
// decodeTaskPayload liest den Auftragsinhalt je nach Art.
|
|
//
|
|
// Beide Arten teilen sich eine Spalte: Der Auftrag ist ein Vertrag, und ein
|
|
// zusätzliches Feld darf keine Migration auf jedem Agenten erzwingen.
|
|
func decodeTaskPayload(encodedPayload []byte, targetTask *Task) error {
|
|
if targetTask.TaskType == TypeRestore {
|
|
var restoreEnvelope struct {
|
|
Restore RestorePayload `json:"restore"`
|
|
}
|
|
|
|
if decodeError := json.Unmarshal(encodedPayload, &restoreEnvelope); decodeError != nil {
|
|
return fmt.Errorf("der auftragsinhalt ist unlesbar: %w", decodeError)
|
|
}
|
|
|
|
targetTask.Restore = restoreEnvelope.Restore
|
|
|
|
return nil
|
|
}
|
|
|
|
if decodeError := json.Unmarshal(encodedPayload, &targetTask.Backup); decodeError != nil {
|
|
return fmt.Errorf("der auftragsinhalt ist unlesbar: %w", decodeError)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// ReportProgress vermerkt Fortschritt und Lebendmeldung.
|
|
//
|
|
// Beides in einem Aufruf: Der Fortschritt ist ohnehin das Lebenszeichen. Ein
|
|
// getrennter Meldeweg brächte nur die Möglichkeit, dass eines von beiden
|
|
// vergessen wird.
|
|
func (store *Store) ReportProgress(progressContext context.Context, taskIdentifier uuid.UUID,
|
|
agentIdentifier uuid.UUID, bytesProcessed int64, filesProcessed int64) error {
|
|
const updateStatement = `
|
|
UPDATE agent_tasks
|
|
SET status = 'running',
|
|
started_at = COALESCE(started_at, now()),
|
|
heartbeat_at = now(),
|
|
bytes_processed = $3,
|
|
files_processed = $4
|
|
WHERE id = $1 AND agent_id = $2 AND status IN ('claimed', 'running')`
|
|
|
|
commandTag, executeError := store.connectionPool.Exec(progressContext, updateStatement,
|
|
taskIdentifier, agentIdentifier, bytesProcessed, filesProcessed)
|
|
if executeError != nil {
|
|
return fmt.Errorf("der fortschritt liess sich nicht vermerken: %w", executeError)
|
|
}
|
|
|
|
if commandTag.RowsAffected() == 0 {
|
|
return ErrTaskNotOwned
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// CompleteTask vermerkt das Ergebnis eines Auftrags.
|
|
func (store *Store) CompleteTask(completeContext context.Context, taskIdentifier uuid.UUID,
|
|
agentIdentifier uuid.UUID, result TaskResult) error {
|
|
result.Normalize()
|
|
|
|
const updateStatement = `
|
|
UPDATE agent_tasks
|
|
SET status = $3,
|
|
completed_at = now(),
|
|
heartbeat_at = now(),
|
|
bytes_processed = $4,
|
|
bytes_written = $5,
|
|
files_processed = $6,
|
|
files_skipped = $7,
|
|
backup_id_in_repository = NULLIF($8, ''),
|
|
error_code = NULLIF($9, ''),
|
|
error_message = NULLIF($10, ''),
|
|
failure_class = NULLIF($11, '')
|
|
WHERE id = $1 AND agent_id = $2 AND status IN ('claimed', 'running')`
|
|
|
|
commandTag, executeError := store.connectionPool.Exec(completeContext, updateStatement,
|
|
taskIdentifier, agentIdentifier, string(result.Status),
|
|
result.BytesProcessed, result.BytesWritten, result.FilesProcessed, result.FilesSkipped,
|
|
result.BackupIDInRepository, result.ErrorCode, result.ErrorMessage, result.FailureClass)
|
|
if executeError != nil {
|
|
return fmt.Errorf("das ergebnis liess sich nicht vermerken: %w", executeError)
|
|
}
|
|
|
|
if commandTag.RowsAffected() == 0 {
|
|
return ErrTaskNotOwned
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// LoadTask liest einen Auftrag.
|
|
func (store *Store) LoadTask(loadContext context.Context, taskIdentifier uuid.UUID) (*Task, error) {
|
|
const selectStatement = `
|
|
SELECT id, agent_id, job_run_id, job_source_id, task_type, status, task_payload,
|
|
bytes_processed, bytes_written, files_processed, files_skipped,
|
|
COALESCE(backup_id_in_repository, ''), COALESCE(error_code, ''),
|
|
COALESCE(error_message, ''), COALESCE(failure_class, ''),
|
|
created_at, started_at, completed_at
|
|
FROM agent_tasks
|
|
WHERE id = $1`
|
|
|
|
var (
|
|
loadedTask Task
|
|
rawTaskType string
|
|
rawStatus string
|
|
encodedPayload []byte
|
|
)
|
|
|
|
scanError := store.connectionPool.QueryRow(loadContext, selectStatement, taskIdentifier).Scan(
|
|
&loadedTask.ID, &loadedTask.AgentID, &loadedTask.JobRunID, &loadedTask.JobSourceID,
|
|
&rawTaskType, &rawStatus, &encodedPayload,
|
|
&loadedTask.BytesProcessed, &loadedTask.BytesWritten,
|
|
&loadedTask.FilesProcessed, &loadedTask.FilesSkipped,
|
|
&loadedTask.BackupIDInRepository, &loadedTask.ErrorCode,
|
|
&loadedTask.ErrorMessage, &loadedTask.FailureClass,
|
|
&loadedTask.CreatedAt, &loadedTask.StartedAt, &loadedTask.CompletedAt)
|
|
if scanError != nil {
|
|
return nil, fmt.Errorf("der auftrag liess sich nicht lesen: %w", scanError)
|
|
}
|
|
|
|
loadedTask.TaskType = TaskType(rawTaskType)
|
|
loadedTask.Status = TaskStatus(rawStatus)
|
|
|
|
if decodeError := decodeTaskPayload(encodedPayload, &loadedTask); decodeError != nil {
|
|
return nil, decodeError
|
|
}
|
|
|
|
return &loadedTask, nil
|
|
}
|
|
|
|
// ReclaimStaleTasks gibt Aufträge abgestürzter Agenten frei.
|
|
//
|
|
// Anders als bei einer Wiederherstellung wird der Auftrag **nicht** wieder
|
|
// eingereiht: Ein Agent, der mitten in einer Sicherung verschwindet, hat
|
|
// womöglich Blöcke geschrieben und eine Session offen. Der Auftrag gilt als
|
|
// gescheitert und wartet auf einen Blick — die Wiederholung entscheidet die
|
|
// Auftragsschleife anhand der Fehlerklasse.
|
|
//
|
|
// Die Frist muss deutlich über dem Meldeabstand des Agenten liegen, sonst gäbe
|
|
// sie einen Auftrag frei, der gerade läuft.
|
|
func (store *Store) ReclaimStaleTasks(reclaimContext context.Context, staleAfter time.Duration) (int, error) {
|
|
const updateStatement = `
|
|
UPDATE agent_tasks
|
|
SET status = 'failed',
|
|
completed_at = now(),
|
|
error_code = 'AGENT_LOST',
|
|
error_message = $2,
|
|
failure_class = 'transient'
|
|
WHERE status IN ('claimed', 'running')
|
|
AND heartbeat_at IS NOT NULL
|
|
AND heartbeat_at < now() - $1::interval`
|
|
|
|
const explanation = "Der ausfuehrende Agent meldete sich nicht mehr. Der Auftrag wurde nicht " +
|
|
"selbsttaetig wiederholt: Der Agent koennte bereits Bloecke geschrieben haben."
|
|
|
|
commandTag, executeError := store.connectionPool.Exec(reclaimContext, updateStatement,
|
|
staleAfter.String(), explanation)
|
|
if executeError != nil {
|
|
return 0, fmt.Errorf("verwaiste auftraege konnten nicht freigegeben werden: %w", executeError)
|
|
}
|
|
|
|
return int(commandTag.RowsAffected()), nil
|
|
}
|
|
|
|
// uniqueViolationCode ist der SQLSTATE eines Eindeutigkeitsverstosses.
|
|
const uniqueViolationCode = "23505"
|
|
|
|
// isUniqueViolation erkennt einen Eindeutigkeitsverstoss.
|
|
func isUniqueViolation(occurredError error) bool {
|
|
var pgError interface{ SQLState() string }
|
|
|
|
return errors.As(occurredError, &pgError) && pgError.SQLState() == uniqueViolationCode
|
|
}
|