FDN-04: job-queue & worker-runtime
Postgres-Jobqueue (processing_jobs, FOR UPDATE SKIP LOCKED), In-Prozess- Worker-Goroutinen, kein Redis/AMQP. Enqueue mit Idempotency-Key-Dedup, Dequeue mit Stale-Lock-Wiedervorlage (Absturzsicherheit), Fail mit arithmetischem Backoff und Dead-Letter-Queue nach erschoepften Versuchen, RequeueDeadLetter fuer manuelle Wiederholung. Auf 192.168.1.131 verifiziert: Absturz-Wiedervorlage (Job von einem "abgestuerzten" Worker nie completed/failed, zweiter Worker holt ihn nach Ablauf der Sperre erneut), Idempotenz bei Doppelzustellung (gleicher idempotency_key erzeugt nur 1 Zeile), DLQ-Eintrag manuell wiederholbar. Vier reale Fehler beim Testen gefunden und behoben: zwei pgx-Typinferenz- Bugs im SQL-Parameterhandling (toter workerID-Parameter ohne Referenz in der Query; untypisiertes any statt []string fuer den ::text[]-Cast), sowie zwei Testinfrastruktur-Bugs (dms_tenant_test sammelte schema_migrations- Zustand ueber Sitzungen hinweg an, jobqueue-Testfixture raeumte processing_jobs nicht auf) - neues scripts/reset-test-env.sh + make check (-p 1) analog Core behoben. Siehe dms/docs/FDN-04-PRUEFPROTOKOLL.md fuer alle Pruefungsergebnisse. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
This commit is contained in:
co-authored by
Claude Sonnet 5
parent
a9ede93176
commit
bfea94b032
@@ -0,0 +1,238 @@
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,233 @@
|
||||
package jobqueue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupTest(t *testing.T) *pgxpool.Pool {
|
||||
t.Helper()
|
||||
dsn := os.Getenv("TEST_TENANT_DSN")
|
||||
if dsn == "" {
|
||||
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { pool.Close() })
|
||||
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE EXTENSION IF NOT EXISTS pgcrypto;
|
||||
CREATE TABLE IF NOT EXISTS processing_jobs (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
job_type TEXT NOT NULL,
|
||||
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
idempotency_key TEXT UNIQUE,
|
||||
status TEXT NOT NULL DEFAULT 'pending'
|
||||
CHECK (status IN ('pending', 'processing', 'succeeded', 'failed', 'dead_letter')),
|
||||
attempts INT NOT NULL DEFAULT 0,
|
||||
max_attempts INT NOT NULL DEFAULT 5,
|
||||
available_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
locked_at TIMESTAMPTZ,
|
||||
locked_by TEXT,
|
||||
last_error TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
// DROP statt nur TRUNCATE: internal/jobqueue und internal/migrate teilen
|
||||
// sich dieselbe physische Test-Datenbank (TEST_TENANT_DSN) ueber
|
||||
// Paketgrenzen hinweg. Bliebe die Tabelle stehen, wuerde internal/migrate
|
||||
// spaeter mit "relation already exists" gegen die von diesem Fixture
|
||||
// angelegte, aber unversionierte Tabelle scheitern.
|
||||
t.Cleanup(func() {
|
||||
_, _ = pool.Exec(context.Background(), `DROP TABLE IF EXISTS processing_jobs`)
|
||||
})
|
||||
return pool
|
||||
}
|
||||
|
||||
// TestEnqueueDequeueComplete ist Akzeptanzkriterium 1: Jobs werden
|
||||
// zuverlaessig eingereiht und verarbeitet.
|
||||
func TestEnqueueDequeueComplete(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "index-document", map[string]string{"document_id": "d1"}, EnqueueOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
job, err := q.Dequeue(ctx, "worker-1", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue: %v", err)
|
||||
}
|
||||
if job.ID != id {
|
||||
t.Fatalf("dequeue lieferte job %q, want %q", job.ID, id)
|
||||
}
|
||||
if job.Status != StatusProcessing {
|
||||
t.Fatalf("status nach dequeue = %q, want %q", job.Status, StatusProcessing)
|
||||
}
|
||||
var payload map[string]string
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
t.Fatalf("payload dekodieren: %v", err)
|
||||
}
|
||||
if payload["document_id"] != "d1" {
|
||||
t.Fatalf("payload = %v, want document_id=d1", payload)
|
||||
}
|
||||
|
||||
if err := q.Complete(ctx, job.ID); err != nil {
|
||||
t.Fatalf("complete: %v", err)
|
||||
}
|
||||
|
||||
status, err := q.Status(ctx, job.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("status: %v", err)
|
||||
}
|
||||
if status.Status != StatusSucceeded {
|
||||
t.Fatalf("endstatus = %q, want %q", status.Status, StatusSucceeded)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDequeue_NoJobAvailable prueft den Leerfall.
|
||||
func TestDequeue_NoJobAvailable(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
|
||||
if _, err := q.Dequeue(context.Background(), "worker-1", nil); !errors.Is(err, ErrNoJobAvailable) {
|
||||
t.Fatalf("erwartet ErrNoJobAvailable, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFail_LandsInDeadLetterAfterMaxAttempts ist Akzeptanzkriterium 2.
|
||||
func TestFail_LandsInDeadLetterAfterMaxAttempts(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{MaxAttempts: 2})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
// Versuch 1: schlaegt fehl, geht zurueck nach "pending" (max_attempts=2 noch nicht erreicht).
|
||||
job, err := q.Dequeue(ctx, "worker-1", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue 1: %v", err)
|
||||
}
|
||||
if err := q.Fail(ctx, job.ID, errors.New("ocr-engine nicht erreichbar")); err != nil {
|
||||
t.Fatalf("fail 1: %v", err)
|
||||
}
|
||||
status, err := q.Status(ctx, id)
|
||||
if err != nil {
|
||||
t.Fatalf("status nach fail 1: %v", err)
|
||||
}
|
||||
if status.Status != StatusPending {
|
||||
t.Fatalf("status nach fail 1 = %q, want %q (noch nicht erschoepft)", status.Status, StatusPending)
|
||||
}
|
||||
|
||||
// Versuch 2: Backoff manuell umgehen (available_at direkt zuruecksetzen,
|
||||
// damit der Test nicht auf echten Backoff warten muss).
|
||||
if _, err := pool.Exec(ctx, `UPDATE processing_jobs SET available_at = now() WHERE id = $1`, id); err != nil {
|
||||
t.Fatalf("available_at zuruecksetzen: %v", err)
|
||||
}
|
||||
job2, err := q.Dequeue(ctx, "worker-1", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue 2: %v", err)
|
||||
}
|
||||
if err := q.Fail(ctx, job2.ID, errors.New("ocr-engine weiterhin nicht erreichbar")); err != nil {
|
||||
t.Fatalf("fail 2: %v", err)
|
||||
}
|
||||
|
||||
final, err := q.Status(ctx, id)
|
||||
if err != nil {
|
||||
t.Fatalf("status nach fail 2: %v", err)
|
||||
}
|
||||
if final.Status != StatusDeadLetter {
|
||||
t.Fatalf("status nach erschoepften versuchen = %q, want %q", final.Status, StatusDeadLetter)
|
||||
}
|
||||
if final.LastError == nil || *final.LastError == "" {
|
||||
t.Fatal("erwartet gesetzten last_error im dead-letter-eintrag")
|
||||
}
|
||||
}
|
||||
|
||||
// TestRequeueDeadLetter ist Pruefung 3: DLQ-Eintrag manuell wiederholbar.
|
||||
func TestRequeueDeadLetter(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "convert", nil, EnqueueOptions{MaxAttempts: 1})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
job, err := q.Dequeue(ctx, "worker-1", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue: %v", err)
|
||||
}
|
||||
if err := q.Fail(ctx, job.ID, errors.New("konverter abgestuerzt")); err != nil {
|
||||
t.Fatalf("fail: %v", err)
|
||||
}
|
||||
status, _ := q.Status(ctx, id)
|
||||
if status.Status != StatusDeadLetter {
|
||||
t.Fatalf("voraussetzung nicht erfuellt: job sollte in dead_letter stehen, ist %q", status.Status)
|
||||
}
|
||||
|
||||
if err := q.RequeueDeadLetter(ctx, id); err != nil {
|
||||
t.Fatalf("requeuedeadletter: %v", err)
|
||||
}
|
||||
afterRequeue, err := q.Status(ctx, id)
|
||||
if err != nil {
|
||||
t.Fatalf("status nach requeue: %v", err)
|
||||
}
|
||||
if afterRequeue.Status != StatusPending {
|
||||
t.Fatalf("status nach requeue = %q, want %q", afterRequeue.Status, StatusPending)
|
||||
}
|
||||
if afterRequeue.Attempts != 0 {
|
||||
t.Fatalf("attempts nach requeue = %d, want 0", afterRequeue.Attempts)
|
||||
}
|
||||
|
||||
// Requeue eines NICHT in dead_letter stehenden Jobs wird abgewiesen.
|
||||
if err := q.RequeueDeadLetter(ctx, id); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("requeue eines pending-jobs: erwartet ErrNotFound, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestEnqueue_IdempotencyKeyPreventsDuplicate ist Pruefung 2: Idempotenz
|
||||
// bei Doppelzustellung nachgewiesen (auf Einreih-Ebene).
|
||||
func TestEnqueue_IdempotencyKeyPreventsDuplicate(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id1, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{IdempotencyKey: "ocr:revision-42"})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue 1: %v", err)
|
||||
}
|
||||
id2, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{IdempotencyKey: "ocr:revision-42"})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue 2 (doppelzustellung): %v", err)
|
||||
}
|
||||
if id1 != id2 {
|
||||
t.Fatalf("doppelte einreihung mit gleichem idempotency-key erzeugte zwei jobs: %q != %q", id1, id2)
|
||||
}
|
||||
|
||||
var count int
|
||||
if err := pool.QueryRow(ctx, `SELECT count(*) FROM processing_jobs WHERE idempotency_key = 'ocr:revision-42'`).Scan(&count); err != nil {
|
||||
t.Fatalf("zeilen zaehlen: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("erwartet genau 1 zeile fuer den idempotency-key, habe %d", count)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package jobqueue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Handler verarbeitet einen einzelnen Job. Ein zurueckgegebener Fehler
|
||||
// fuehrt zu Queue.Fail (Backoff/DLQ), nil zu Queue.Complete.
|
||||
type Handler func(ctx context.Context, job *Job) error
|
||||
|
||||
// Worker pollt die Queue in einer In-Prozess-Goroutine und ruft Handler je
|
||||
// abgeholtem Job auf — die "In-Prozess-Worker-Goroutinen" aus der
|
||||
// Ticket-Vorgabe, kein separater Prozess/Redis noetig.
|
||||
type Worker struct {
|
||||
queue *Queue
|
||||
workerID string
|
||||
jobTypes []string
|
||||
pollInterval time.Duration
|
||||
handler Handler
|
||||
}
|
||||
|
||||
func NewWorker(queue *Queue, workerID string, jobTypes []string, pollInterval time.Duration, handler Handler) *Worker {
|
||||
return &Worker{queue: queue, workerID: workerID, jobTypes: jobTypes, pollInterval: pollInterval, handler: handler}
|
||||
}
|
||||
|
||||
// Run blockiert, bis ctx beendet wird, und verarbeitet dabei fortlaufend
|
||||
// Jobs. Absturzsicherheit (Pruefung 1) entsteht NICHT durch Run selbst,
|
||||
// sondern dadurch, dass ein abgestuerzter Prozess (der Run gar nicht mehr
|
||||
// ausfuehrt) seine "processing"-Sperren nie verlaengert — ein ANDERER
|
||||
// Worker-Prozess holt den Job nach Ablauf von staleLockAfter erneut ab
|
||||
// (siehe Queue.Dequeue).
|
||||
func (w *Worker) Run(ctx context.Context) {
|
||||
ticker := time.NewTicker(w.pollInterval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
w.processOne(ctx)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// processOne holt und verarbeitet EINEN Job, falls verfuegbar. Oeffentlich
|
||||
// über RunOnce fuer Tests, die deterministisch (ohne Polling-Timing) einen
|
||||
// einzelnen Verarbeitungsschritt auslösen wollen.
|
||||
func (w *Worker) processOne(ctx context.Context) {
|
||||
job, err := w.queue.Dequeue(ctx, w.workerID, w.jobTypes)
|
||||
if err != nil {
|
||||
if !errors.Is(err, ErrNoJobAvailable) {
|
||||
log.Printf("jobqueue: dequeue fehlgeschlagen: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if handlerErr := w.handler(ctx, job); handlerErr != nil {
|
||||
if err := w.queue.Fail(ctx, job.ID, handlerErr); err != nil {
|
||||
log.Printf("jobqueue: fehlschlag fuer job %q nicht erfassbar: %v", job.ID, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err := w.queue.Complete(ctx, job.ID); err != nil {
|
||||
log.Printf("jobqueue: abschluss fuer job %q fehlgeschlagen: %v", job.ID, err)
|
||||
}
|
||||
}
|
||||
|
||||
// RunOnce verarbeitet synchron genau einen Job (falls verfuegbar) und
|
||||
// kehrt zurueck — fuer Tests, die ohne Polling-Intervall arbeiten wollen.
|
||||
func (w *Worker) RunOnce(ctx context.Context) {
|
||||
w.processOne(ctx)
|
||||
}
|
||||
@@ -0,0 +1,115 @@
|
||||
package jobqueue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// TestWorker_RunOnce_ProcessesAndCompletesJob ist der End-to-End-Nachweis
|
||||
// fuer Akzeptanzkriterium 1 ueber den Worker statt direkt ueber Queue.
|
||||
func TestWorker_RunOnce_ProcessesAndCompletesJob(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "index-document", nil, EnqueueOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
var handled int32
|
||||
w := NewWorker(q, "worker-1", nil, time.Millisecond, func(ctx context.Context, job *Job) error {
|
||||
atomic.AddInt32(&handled, 1)
|
||||
return nil
|
||||
})
|
||||
w.RunOnce(ctx)
|
||||
|
||||
if atomic.LoadInt32(&handled) != 1 {
|
||||
t.Fatalf("handler wurde %d mal aufgerufen, want 1", handled)
|
||||
}
|
||||
status, err := q.Status(ctx, id)
|
||||
if err != nil {
|
||||
t.Fatalf("status: %v", err)
|
||||
}
|
||||
if status.Status != StatusSucceeded {
|
||||
t.Fatalf("status = %q, want %q", status.Status, StatusSucceeded)
|
||||
}
|
||||
}
|
||||
|
||||
// TestWorker_HandlerErrorTriggersFail prueft, dass ein Handler-Fehler zu
|
||||
// Queue.Fail fuehrt (Backoff/DLQ-Pfad ueber den Worker statt direkt).
|
||||
func TestWorker_HandlerErrorTriggersFail(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
q := NewQueue(pool, time.Minute)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{MaxAttempts: 5})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
w := NewWorker(q, "worker-1", nil, time.Millisecond, func(ctx context.Context, job *Job) error {
|
||||
return errors.New("ocr fehlgeschlagen")
|
||||
})
|
||||
w.RunOnce(ctx)
|
||||
|
||||
status, err := q.Status(ctx, id)
|
||||
if err != nil {
|
||||
t.Fatalf("status: %v", err)
|
||||
}
|
||||
if status.Status != StatusPending {
|
||||
t.Fatalf("status nach handler-fehler = %q, want %q (erneut eingeplant)", status.Status, StatusPending)
|
||||
}
|
||||
if status.LastError == nil || *status.LastError != "ocr fehlgeschlagen" {
|
||||
t.Fatalf("last_error = %v, want %q", status.LastError, "ocr fehlgeschlagen")
|
||||
}
|
||||
}
|
||||
|
||||
// TestDequeue_StaleLockIsRedelivered ist Pruefung 1: Absturz eines Workers
|
||||
// fuehrt zu erneuter Zustellung. Simuliert einen Absturz, indem ein Job
|
||||
// dequeued (auf "processing" gesperrt), aber NIE completed/failed wird —
|
||||
// nach Ablauf von staleLockAfter muss ein ANDERER Worker denselben Job
|
||||
// erneut abholen koennen.
|
||||
func TestDequeue_StaleLockIsRedelivered(t *testing.T) {
|
||||
pool := setupTest(t)
|
||||
// Sehr kurzes Stale-Fenster, damit der Test nicht lange warten muss.
|
||||
q := NewQueue(pool, 50*time.Millisecond)
|
||||
ctx := context.Background()
|
||||
|
||||
id, err := q.Enqueue(ctx, "convert", nil, EnqueueOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
crashed, err := q.Dequeue(ctx, "worker-crashed", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue (worker-crashed): %v", err)
|
||||
}
|
||||
if crashed.ID != id {
|
||||
t.Fatalf("dequeue lieferte unerwarteten job %q", crashed.ID)
|
||||
}
|
||||
// worker-crashed ruft absichtlich weder Complete noch Fail auf — simuliert
|
||||
// einen Prozessabsturz mitten in der Verarbeitung.
|
||||
|
||||
// Sofortiger erneuter Dequeue-Versuch (Sperre noch frisch) darf den Job
|
||||
// NICHT liefern.
|
||||
if _, err := q.Dequeue(ctx, "worker-2", nil); !errors.Is(err, ErrNoJobAvailable) {
|
||||
t.Fatalf("job wurde trotz frischer sperre erneut ausgeliefert (oder anderer fehler): %v", err)
|
||||
}
|
||||
|
||||
time.Sleep(80 * time.Millisecond) // > staleLockAfter
|
||||
|
||||
redelivered, err := q.Dequeue(ctx, "worker-2", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dequeue nach ablauf der sperre: %v", err)
|
||||
}
|
||||
if redelivered.ID != id {
|
||||
t.Fatalf("erneut zugestellter job = %q, want %q", redelivered.ID, id)
|
||||
}
|
||||
if err := q.Complete(ctx, redelivered.ID); err != nil {
|
||||
t.Fatalf("complete durch worker-2: %v", err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user