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
116 lines
3.5 KiB
Go
116 lines
3.5 KiB
Go
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)
|
|
}
|
|
}
|