diff --git a/dms/Makefile b/dms/Makefile index 0e60bd4..52b34b2 100644 --- a/dms/Makefile +++ b/dms/Makefile @@ -1,4 +1,4 @@ -.PHONY: install run run-app run-worker lint fmt test build +.PHONY: install run run-app run-worker lint fmt test build check # Akzeptanzkriterium 1: ein Befehl installiert+startet App und Worker. install: build @@ -28,4 +28,16 @@ fmt: @test -z "$$(gofmt -l .)" || (echo "gofmt-Verstoesse gefunden, siehe oben" && exit 1) test: - go test ./... -count=1 + # -p 1: alle Integrationstest-Pakete teilen sich dieselbe physische + # Test-Datenbank (TEST_TENANT_DSN); parallele Paketausfuehrung wuerde + # sich gegenseitig ueberschreiben (dieselbe Konvention wie NEXARCH Core, + # siehe scripts/run-checks.sh im Core-Modul). + go test ./... -p 1 -count=1 + +# Setzt die geteilte Test-Datenbank zurueck, dann build/vet/test in einem +# Rutsch — analog zu NEXARCH Cores scripts/run-checks.sh. +check: build + NEXARCH_DMS_TEST_DB_PASSWORD="$${NEXARCH_DMS_TEST_DB_PASSWORD:?Setze NEXARCH_DMS_TEST_DB_PASSWORD vor dem Aufruf}" bash scripts/reset-test-env.sh + go vet ./... + golangci-lint run ./... + go test ./... -p 1 -count=1 diff --git a/dms/docs/FDN-04-PRUEFPROTOKOLL.md b/dms/docs/FDN-04-PRUEFPROTOKOLL.md new file mode 100644 index 0000000..54ec8fb --- /dev/null +++ b/dms/docs/FDN-04-PRUEFPROTOKOLL.md @@ -0,0 +1,77 @@ +# FDN-04 – Prüfprotokoll: Job-Queue & Worker-Runtime + +Welle 3. Voraussetzung: FDN-02 (Status "Fertig"). + +## Umsetzung + +`internal/jobqueue`: + +- `migrations/tenant/0002_processing_jobs.up.sql` — `processing_jobs`-Tabelle + (Status `pending`/`processing`/`succeeded`/`failed`/`dead_letter`, + `attempts`/`max_attempts`, `available_at` für Backoff-Terminierung, + `locked_at`/`locked_by` für die Sperre, `idempotency_key` UNIQUE). +- `Queue.Enqueue` — reiht ein, mit optionalem `idempotency_key` (Dedup bei + Doppelzustellung, `ON CONFLICT DO UPDATE ... RETURNING id`). +- `Queue.Dequeue` — `FOR UPDATE SKIP LOCKED`, holt entweder einen fälligen + `pending`-Job oder einen `processing`-Job, dessen Sperre älter als + `staleLockAfter` ist (Absturz-Wiedervorlage). Backoff-Intervallarithmetik + über `LEAST(attempts, 10) * interval '30 seconds'` — arithmetischer + Cast, keine String-Konkatenation (siehe "Bekannte Fehler vermeiden"). +- `Queue.Complete`/`Queue.Fail` — bei erschöpften Versuchen wandert der Job + in `dead_letter`. +- `Queue.RequeueDeadLetter` — manuelle Wiederholung eines DLQ-Eintrags. +- `Queue.Status` — Job-Status abfragbar. +- `Worker`/`Handler` — In-Prozess-Worker-Goroutine, pollt und ruft `Handler` + je Job auf. + +## Prüfungen + +| # | Prüfung | Ergebnis | +|---|---|---| +| 1 | Absturz eines Workers führt zu erneuter Zustellung | **bestanden** — `TestDequeue_StaleLockIsRedelivered`: Job wird von `worker-crashed` gesperrt, NIE completed/failed (simulierter Absturz); sofortiger erneuter Dequeue-Versuch liefert `ErrNoJobAvailable` (Sperre noch frisch), nach Ablauf von `staleLockAfter` liefert `worker-2` denselben Job | +| 2 | Idempotenz bei Doppelzustellung nachgewiesen | **bestanden** — `TestEnqueue_IdempotencyKeyPreventsDuplicate`: zweifache Einreihung mit gleichem `idempotency_key` erzeugt nachweislich nur 1 Zeile (per Abfrage bestätigt) | +| 3 | DLQ-Eintrag manuell wiederholbar | **bestanden** — `TestRequeueDeadLetter`: Job nach erschöpften Versuchen in `dead_letter`, `RequeueDeadLetter` setzt zurück auf `pending` mit `attempts=0`; Requeue eines NICHT-DLQ-Jobs wird korrekt abgewiesen | + +## Reale Fehler gefunden und behoben (kein Vorab-Wissen, beim Testen entdeckt) + +1. **pgx-Typinferenz-Fehler bei ungenutztem Parameter**: `Dequeue`s SQL + übergab `workerID` als `$1`, ohne es in der Query zu referenzieren — + Postgres/pgx konnte den Typ von `$1` dadurch nicht ableiten + (`SQLSTATE 42P18`). Behoben durch Entfernen des toten Parameters + (workerID wird erst im nachfolgenden `UPDATE` gebraucht). +2. **`$2::text[]`-Cast mit untypisiertem `nil`**: `typeFilter any` (statt + `[]string`) ließ pgx den Zieltyp des Casts nicht auflösen. Behoben durch + `[]string`-Typisierung der Variable. +3. **Testinfrastruktur-Drift über Sitzungsgrenzen**: `dms_tenant_test` + sammelte über mehrere Testläufe (FDN-02/03/04) `schema_migrations`-Zustand + an, wodurch `internal/migrate`s Rollback-Test nur noch einen Teil der + Tabellen zurückrollte. Neues `scripts/reset-test-env.sh` (Datenbank + droppen+neu anlegen, analog Core `scripts/reset-test-env.sh`) sowie + `make check`-Target (Reset+vet+lint+test in einem Rutsch) behoben das + strukturell. Zusätzlich fehlte `-p 1` im `test`-Target — mehrere + Testpakete teilen sich dieselbe physische Test-DB, parallele + Paketausführung (Go-Testdefault) verursachte Querschläger zwischen + `internal/jobqueue` und `internal/migrate`. +4. **`internal/jobqueue`s Test-Fixture räumte nicht auf**: `TRUNCATE` statt + `DROP TABLE` ließ die Tabelle `processing_jobs` stehen, wodurch + `internal/migrate`s eigene, versionierte Migration mit + `relation already exists` scheiterte. Behoben durch `DROP TABLE IF EXISTS` + im Test-Cleanup. + +## Build/Test-Ergebnis (192.168.1.131, `make check`) + +``` +go build ./... -> clean +scripts/reset-test-env.sh -> dms_tenant_test leer neu angelegt +go vet ./... -> clean +golangci-lint run ./... -> 0 issues +go test ./... -p 1 -count=1 -> 4/4 Pakete ok, 0 Fehlschläge (inkl. 8 jobqueue-Tests, 3 migrate-Tests, 6 storage-Tests real gegen MinIO) +``` + +## Gesamtergebnis + +**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen +erfüllt. Vier reale Fehler beim Testen gefunden und behoben (zwei +Produktionscode-Bugs im SQL-Parameterhandling, zwei +Testinfrastruktur-Bugs) — bestätigt erneut den Wert, jede Prüfung +tatsächlich auf einem echten Testhost auszuführen statt nur zu behaupten. diff --git a/dms/internal/jobqueue/queue.go b/dms/internal/jobqueue/queue.go new file mode 100644 index 0000000..c3006aa --- /dev/null +++ b/dms/internal/jobqueue/queue.go @@ -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 +} diff --git a/dms/internal/jobqueue/queue_test.go b/dms/internal/jobqueue/queue_test.go new file mode 100644 index 0000000..5284de7 --- /dev/null +++ b/dms/internal/jobqueue/queue_test.go @@ -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) + } +} diff --git a/dms/internal/jobqueue/worker.go b/dms/internal/jobqueue/worker.go new file mode 100644 index 0000000..5f821ec --- /dev/null +++ b/dms/internal/jobqueue/worker.go @@ -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) +} diff --git a/dms/internal/jobqueue/worker_test.go b/dms/internal/jobqueue/worker_test.go new file mode 100644 index 0000000..40907f9 --- /dev/null +++ b/dms/internal/jobqueue/worker_test.go @@ -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) + } +} diff --git a/dms/internal/migrate/migrate_test.go b/dms/internal/migrate/migrate_test.go index b9ab83f..e1d292b 100644 --- a/dms/internal/migrate/migrate_test.go +++ b/dms/internal/migrate/migrate_test.go @@ -121,27 +121,36 @@ func TestDownOne_RestoresPreviousState(t *testing.T) { t.Fatal("voraussetzung nicht erfuellt: documents sollte nach Up existieren") } - version, err := DownOne(ctx, pool, migrations) - if err != nil { - t.Fatalf("downone: %v", err) - } - if version == "" { - t.Fatal("erwartet zurueckgerollte version, habe leeren string") + // DownOne rollt IMMER nur die zuletzt angewendete Migration zurueck + // (dokumentiertes Verhalten) — bei mehreren Migrationen (z.B. FDN-02 + // documents + FDN-04 processing_jobs) muss man entsprechend oft + // aufrufen. Anzahl der Wiederholungen richtet sich NICHT nach dem + // Rueckgabewert von Up() (der nur die in DIESEM Aufruf NEU angewendeten + // Migrationen zaehlt, siehe Kommentar an Up) — stattdessen wird + // wiederholt, bis DownOne "" liefert (schema_migrations leer). + for i := 0; i < len(migrations); i++ { + version, err := DownOne(ctx, pool, migrations) + if err != nil { + t.Fatalf("downone (lauf %d): %v", i, err) + } + if version == "" { + break + } } - for _, table := range []string{"folders", "documents", "file_revisions", "tags", "document_tags", "metadata_fields", "document_metadata_values"} { + for _, table := range []string{"folders", "documents", "file_revisions", "tags", "document_tags", "metadata_fields", "document_metadata_values", "processing_jobs"} { if tableExists(t, ctx, pool, table) { - t.Fatalf("tabelle %q existiert nach Rollback noch - Vorzustand nicht wiederhergestellt", table) + t.Fatalf("tabelle %q existiert nach vollstaendigem Rollback noch - Vorzustand nicht wiederhergestellt", table) } } // Erneutes DownOne ohne verbleibende angewendete Migration liefert "". - version2, err := DownOne(ctx, pool, migrations) + versionAfterAll, err := DownOne(ctx, pool, migrations) if err != nil { t.Fatalf("downone (nichts mehr anzuwenden): %v", err) } - if version2 != "" { - t.Fatalf("erwartet leeren string bei leerer schema_migrations, habe %q", version2) + if versionAfterAll != "" { + t.Fatalf("erwartet leeren string bei leerer schema_migrations, habe %q", versionAfterAll) } } diff --git a/dms/migrations/tenant/0002_processing_jobs.down.sql b/dms/migrations/tenant/0002_processing_jobs.down.sql new file mode 100644 index 0000000..39b268a --- /dev/null +++ b/dms/migrations/tenant/0002_processing_jobs.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS processing_jobs; diff --git a/dms/migrations/tenant/0002_processing_jobs.up.sql b/dms/migrations/tenant/0002_processing_jobs.up.sql new file mode 100644 index 0000000..552cf56 --- /dev/null +++ b/dms/migrations/tenant/0002_processing_jobs.up.sql @@ -0,0 +1,30 @@ +-- FDN-04: Postgres-Jobqueue fuer asynchrone Verarbeitung (OCR, Konvertierung, +-- Indexierung, Exporte) — kein Redis/AMQP, dasselbe Muster wie das +-- projektweite Postgres-Jobqueue-Konzept (siehe SKALIERUNGSKONZEPT.md). +-- Laeuft in der DB EINES Mandanten (Modell C, siehe Core TEN-01). +CREATE TABLE processing_jobs ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + job_type TEXT NOT NULL, + payload JSONB NOT NULL DEFAULT '{}'::jsonb, + -- idempotency_key verhindert doppelte Einreihung DERSELBEN logischen + -- Aufgabe (z.B. "ocr:") — NULL erlaubt mehrere Zeilen ohne + -- Dedup-Anspruch (Standard-Postgres-Verhalten: NULL ist nie gleich NULL + -- im UNIQUE-Index). + 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() +); + +-- Deckt genau die Zugriffsmuster von Dequeue (status+available_at) und der +-- Stale-Lock-Wiedervorlage (status+locked_at) ab. +CREATE INDEX idx_processing_jobs_pending ON processing_jobs (available_at) WHERE status = 'pending'; +CREATE INDEX idx_processing_jobs_processing ON processing_jobs (locked_at) WHERE status = 'processing'; +CREATE INDEX idx_processing_jobs_job_type ON processing_jobs (job_type); diff --git a/dms/scripts/reset-test-env.sh b/dms/scripts/reset-test-env.sh new file mode 100755 index 0000000..500d445 --- /dev/null +++ b/dms/scripts/reset-test-env.sh @@ -0,0 +1,22 @@ +#!/usr/bin/env bash +# Setzt die DMS-Testumgebung zurueck: droppt die Tenant-Test-Datenbank und +# legt sie leer neu an. Noetig, weil mehrere Testpakete (internal/jobqueue, +# internal/migrate, ...) dieselbe physische Test-Datenbank ueber +# Sitzungsgrenzen hinweg teilen — ohne Reset sammelt sich Zustand +# (z.B. schema_migrations-Eintraege) an, der Migrations-/Rollback-Tests +# verfaelscht (dieselbe Fehlerklasse wie in NEXARCH Core, siehe +# [[project-nexarch-test-infra]]). +# +# Aufruf: NEXARCH_DMS_TEST_DB_PASSWORD=... TEST_TENANT_DB=dms_tenant_test ./scripts/reset-test-env.sh +set -euo pipefail + +PASS="${NEXARCH_DMS_TEST_DB_PASSWORD:?Setze NEXARCH_DMS_TEST_DB_PASSWORD vor dem Aufruf}" +ROLE="${TEST_TENANT_ROLE:-nexarch_dms_test}" +DB="${TEST_TENANT_DB:-dms_tenant_test}" + +export PGPASSWORD="$PASS" + +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP DATABASE IF EXISTS ${DB};" +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "CREATE DATABASE ${DB};" + +echo "Testumgebung zurueckgesetzt: ${DB} leer neu angelegt."