// Integrationstest (SRC-02): echte Postgres-Instanz, folgt derselben // Testhost-Konvention wie mail/internal/dedup — TEST_TENANT_DSN. package indexworker import ( "context" "encoding/json" "errors" "os" "testing" "time" "github.com/jackc/pgx/v5/pgxpool" ) func setupQueue(t *testing.T, staleLockAfter time.Duration) *Queue { t.Helper() dsn := os.Getenv("TEST_TENANT_DSN") if dsn == "" { t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen") } ctx := context.Background() pool, err := pgxpool.New(ctx, dsn) if err != nil { t.Fatalf("pool: %v", err) } t.Cleanup(func() { pool.Close() }) queue := NewQueue(pool, staleLockAfter) if err := queue.EnsureSchema(ctx); err != nil { t.Fatalf("schema: %v", err) } t.Cleanup(func() { _, _ = pool.Exec(context.Background(), `DELETE FROM mail_index_jobs WHERE tenant_slug LIKE 'mandant-src02-%'`) }) return queue } // TestDequeue_WorkerCrashMidRunLosesNoJob ist die geforderte Pflichtprüfung // 1: Absturz eines Workers (Job wird geholt, aber nie completed/failed) // führt nach Ablauf der Stale-Lock-Frist zu erneuter Zustellung an einen // zweiten Worker. func TestDequeue_WorkerCrashMidRunLosesNoJob(t *testing.T) { queue := setupQueue(t, 100*time.Millisecond) ctx := context.Background() jobID, err := queue.EnqueueDelete(ctx, "mandant-src02-crash", "msg-crash-1") if err != nil { t.Fatalf("enqueue: %v", err) } firstAttempt, err := queue.dequeue(ctx, "worker-1-abgestuerzt") if err != nil { t.Fatalf("erster dequeue: %v", err) } if firstAttempt.ID != jobID { t.Fatalf("erwartete job-id %d, habe %d", jobID, firstAttempt.ID) } // worker-1 "stürzt ab": kein complete(), kein fail() — Job bleibt // als "processing" mit veraltetem Lock stehen. if _, err := queue.dequeue(ctx, "worker-2-sofort"); !errors.Is(err, ErrNoJobAvailable) { t.Fatalf("erwartete kein verfügbarer job vor ablauf der stale-lock-frist, habe: %v", err) } time.Sleep(150 * time.Millisecond) secondAttempt, err := queue.dequeue(ctx, "worker-2-nach-timeout") if err != nil { t.Fatalf("zweiter dequeue nach stale-lock-ablauf: %v", err) } if secondAttempt.ID != jobID { t.Fatalf("erwartete erneute zustellung desselben jobs %d, habe %d", jobID, secondAttempt.ID) } if err := queue.complete(ctx, secondAttempt.ID); err != nil { t.Fatalf("complete: %v", err) } } // TestEnqueueAndStatus_RoundTrip deckt Akzeptanzkriterium 3 (Job-Status // abrufbar) auf Queue-Ebene ab. func TestEnqueueAndStatus_RoundTrip(t *testing.T) { queue := setupQueue(t, time.Minute) ctx := context.Background() payload, _ := json.Marshal(map[string]string{"subject": "Test"}) jobID, err := queue.EnqueueIndex(ctx, "mandant-src02-status", "msg-status-1", payload) if err != nil { t.Fatalf("enqueue: %v", err) } job, err := queue.Status(ctx, jobID) if err != nil { t.Fatalf("status: %v", err) } if job.Status != StatusPending { t.Fatalf("erwartete status 'pending', habe %q", job.Status) } }