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) } }