// Package jobqueue implementiert FDN-04: eine Postgres-gestuetzte // Job-Queue mit Wiederholungslogik, Backoff und Dead-Letter-Queue — kein // Redis/AMQP (siehe Ticket-Vorgabe). FOR UPDATE SKIP LOCKED erlaubt // mehrere gleichzeitige Worker-Goroutinen (auch mehrinstanzfaehig, da der // Zustand ausschliesslich in Postgres liegt, keine In-Memory-Zaehler — // dieselbe Konvention wie Core internal/lockout, siehe "Bekannte Fehler // vermeiden" im Ticket). package jobqueue import ( "context" "encoding/json" "errors" "fmt" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) // Status-Werte spiegeln den CHECK-Constraint der Migration. const ( StatusPending = "pending" StatusProcessing = "processing" StatusSucceeded = "succeeded" StatusFailed = "failed" StatusDeadLetter = "dead_letter" ) // ErrNotFound wird geliefert, wenn ein angefragter Job nicht existiert. var ErrNotFound = errors.New("jobqueue: job nicht gefunden") // ErrNoJobAvailable wird von Dequeue geliefert, wenn aktuell kein // abholbarer Job vorhanden ist (kein Fehlerzustand, sondern der Normalfall // bei leerer Queue). var ErrNoJobAvailable = errors.New("jobqueue: kein job verfuegbar") // Job ist eine einzelne Aufgabe in der Queue. type Job struct { ID string JobType string Payload json.RawMessage Status string Attempts int MaxAttempts int LastError *string } // DefaultMaxAttempts/DefaultStaleLockAfter sind Standardwerte, ueberschreibbar // je Enqueue-Aufruf (MaxAttempts) bzw. am Queue selbst (StaleLockAfter). const ( DefaultMaxAttempts = 5 ) // Queue kapselt den Zugriff auf processing_jobs. type Queue struct { pool *pgxpool.Pool staleLockAfter time.Duration } // NewQueue erzeugt eine Queue. staleLockAfter legt fest, ab wann ein // Job, der als "processing" markiert ist, aber dessen Worker vermutlich // abgestuerzt ist, wieder als abholbar gilt (Pruefung 1: Absturz fuehrt zu // erneuter Zustellung) — kein Heartbeat-Mechanismus noetig, ein grosszuegiges // Zeitfenster genuegt fuer die "kleinste Loesung". func NewQueue(pool *pgxpool.Pool, staleLockAfter time.Duration) *Queue { return &Queue{pool: pool, staleLockAfter: staleLockAfter} } // EnqueueOptions steuert optionale Einreih-Parameter. type EnqueueOptions struct { // IdempotencyKey verhindert doppelte Einreihung derselben logischen // Aufgabe (Pruefung 2: Idempotenz bei Doppelzustellung) — leer bedeutet // kein Dedup-Anspruch. IdempotencyKey string MaxAttempts int } // Enqueue reiht einen neuen Job ein (Akzeptanzkriterium 1). Bei gesetztem // IdempotencyKey und bereits existierendem gleichen Key wird die ID des // BEREITS vorhandenen Jobs zurueckgegeben, kein Duplikat angelegt. func (q *Queue) Enqueue(ctx context.Context, jobType string, payload any, opts EnqueueOptions) (string, error) { payloadJSON, err := json.Marshal(payload) if err != nil { return "", fmt.Errorf("jobqueue: payload serialisieren: %w", err) } maxAttempts := opts.MaxAttempts if maxAttempts <= 0 { maxAttempts = DefaultMaxAttempts } var idempotencyKey any if opts.IdempotencyKey != "" { idempotencyKey = opts.IdempotencyKey } var id string err = q.pool.QueryRow(ctx, ` INSERT INTO processing_jobs (job_type, payload, idempotency_key, max_attempts) VALUES ($1, $2, $3, $4) ON CONFLICT (idempotency_key) DO UPDATE SET job_type = processing_jobs.job_type RETURNING id `, jobType, payloadJSON, idempotencyKey, maxAttempts).Scan(&id) if err != nil { return "", fmt.Errorf("jobqueue: job einreihen: %w", err) } return id, nil } // Dequeue holt GENAU EINEN abholbaren Job (faellig UND nicht bereits von // einem anderen Worker gesperrt, ODER dessen Sperre als abgestanden gilt) // und markiert ihn atomar als "processing" (Akzeptanzkriterium 1 / Pruefung // 1 — FOR UPDATE SKIP LOCKED erlaubt mehreren Worker-Goroutinen // gleichzeitigen Aufruf ohne sich gegenseitig zu blockieren oder denselben // Job doppelt zu holen). func (q *Queue) Dequeue(ctx context.Context, workerID string, jobTypes []string) (*Job, error) { tx, err := q.pool.Begin(ctx) if err != nil { return nil, fmt.Errorf("jobqueue: transaktion starten: %w", err) } defer func() { _ = tx.Rollback(ctx) }() var typeFilter []string if len(jobTypes) > 0 { typeFilter = jobTypes } row := tx.QueryRow(ctx, ` SELECT id, job_type, payload, status, attempts, max_attempts, last_error FROM processing_jobs WHERE ( (status = 'pending' AND available_at <= now()) OR (status = 'processing' AND locked_at <= now() - ($2 * interval '1 second')) ) AND ($1::text[] IS NULL OR job_type = ANY($1)) ORDER BY available_at FOR UPDATE SKIP LOCKED LIMIT 1 `, typeFilter, q.staleLockAfter.Seconds()) var j Job if err := row.Scan(&j.ID, &j.JobType, &j.Payload, &j.Status, &j.Attempts, &j.MaxAttempts, &j.LastError); err != nil { if errors.Is(err, pgx.ErrNoRows) { return nil, ErrNoJobAvailable } return nil, fmt.Errorf("jobqueue: naechsten job lesen: %w", err) } if _, err := tx.Exec(ctx, ` UPDATE processing_jobs SET status = 'processing', attempts = attempts + 1, locked_at = now(), locked_by = $2, updated_at = now() WHERE id = $1 `, j.ID, workerID); err != nil { return nil, fmt.Errorf("jobqueue: job sperren: %w", err) } if err := tx.Commit(ctx); err != nil { return nil, fmt.Errorf("jobqueue: dequeue committen: %w", err) } j.Status = StatusProcessing j.Attempts++ return &j, nil } // Complete markiert einen Job als erfolgreich abgeschlossen. func (q *Queue) Complete(ctx context.Context, jobID string) error { tag, err := q.pool.Exec(ctx, ` UPDATE processing_jobs SET status = 'succeeded', locked_at = NULL, locked_by = NULL, updated_at = now() WHERE id = $1 `, jobID) if err != nil { return fmt.Errorf("jobqueue: job abschliessen: %w", err) } if tag.RowsAffected() == 0 { return ErrNotFound } return nil } // Fail markiert einen Job als fehlgeschlagen (Akzeptanzkriterium 2): sind // die maximalen Versuche erreicht, wandert der Job in die Dead-Letter-Queue // (status='dead_letter'), sonst wird er mit exponentiellem Backoff erneut // eingeplant. Backoff-Berechnung nutzt arithmetischen Intervall-Cast // (attempts * interval), KEINE String-Konkatenation (siehe "Bekannte // Fehler vermeiden" im Ticket). func (q *Queue) Fail(ctx context.Context, jobID string, cause error) error { errMsg := cause.Error() tag, err := q.pool.Exec(ctx, ` UPDATE processing_jobs SET status = CASE WHEN attempts >= max_attempts THEN 'dead_letter' ELSE 'pending' END, available_at = now() + (LEAST(attempts, 10) * interval '30 seconds'), locked_at = NULL, locked_by = NULL, last_error = $2, updated_at = now() WHERE id = $1 `, jobID, errMsg) if err != nil { return fmt.Errorf("jobqueue: fehlschlag erfassen: %w", err) } if tag.RowsAffected() == 0 { return ErrNotFound } return nil } // Status liefert den aktuellen Zustand eines Jobs (Akzeptanzkriterium 3). func (q *Queue) Status(ctx context.Context, jobID string) (*Job, error) { var j Job err := q.pool.QueryRow(ctx, ` SELECT id, job_type, payload, status, attempts, max_attempts, last_error FROM processing_jobs WHERE id = $1 `, jobID).Scan(&j.ID, &j.JobType, &j.Payload, &j.Status, &j.Attempts, &j.MaxAttempts, &j.LastError) if err != nil { if errors.Is(err, pgx.ErrNoRows) { return nil, ErrNotFound } return nil, fmt.Errorf("jobqueue: job-status lesen: %w", err) } return &j, nil } // RequeueDeadLetter holt einen Job manuell aus der Dead-Letter-Queue zurueck // in "pending", mit zurueckgesetztem Versuchszaehler (Pruefung 3: DLQ-Eintrag // manuell wiederholbar). Nur fuer Jobs, die tatsaechlich in dead_letter // stehen — verhindert versehentliches Requeue eines noch laufenden Jobs. func (q *Queue) RequeueDeadLetter(ctx context.Context, jobID string) error { tag, err := q.pool.Exec(ctx, ` UPDATE processing_jobs SET status = 'pending', attempts = 0, available_at = now(), last_error = NULL, updated_at = now() WHERE id = $1 AND status = 'dead_letter' `, jobID) if err != nil { return fmt.Errorf("jobqueue: dead-letter-job erneut einreihen: %w", err) } if tag.RowsAffected() == 0 { return fmt.Errorf("jobqueue: job %q steht nicht in dead_letter (oder existiert nicht): %w", jobID, ErrNotFound) } return nil }