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 }