Indexierungs-Worker, der neu archivierte Mails asynchron in den Manticore-Index (SRC-01) einpflegt und Löschungen nachzieht. - indexworker/queue.go: Postgres-Jobqueue (mail_index_jobs), FOR UPDATE SKIP LOCKED, Stale-Lock-Wiedervorlage bei Worker-Absturz, arithmetischer Backoff bei Fail (kein String-Concat für Intervalle) — Konvention aus dms/internal/jobqueue (FDN-04), hier bewusst ohne DLQ (nicht Bestandteil der Akzeptanzkriterien dieser Kachel). - indexworker/worker.go: RunOnce verarbeitet index-/delete-Jobs über search.Client. - search: minimale Erweiterung um Client.Delete und deterministisches DocumentID(tenantSlug, messageID), damit Index/Delete für dieselbe Mail immer dasselbe Dokument treffen. Prüfungen (alle real durchgeführt, siehe mail/docs/SRC-02-PRUEFPROTOKOLL.md): 1. TestDequeue_WorkerCrashMidRunLosesNoJob: simulierter Worker-Absturz, Job wird nach Ablauf der Stale-Lock-Frist real erneut zugestellt. 2. TestDeleteJob_RemovesMailFromSearchResults: Lösch-Job entfernt Mail nachweislich aus Suchtreffern. 3. TestConsistency_DatabaseAndIndexMatchOnSample: DB-Job-Status und Index-Inhalt stichprobenartig real abgeglichen. Zusätzlich TestIndexJob_MakesMailSearchable für Akzeptanzkriterium 1. Kein Umbau: storage/crypto/encstorage/dedup unverändert, bestehendes SRC-01-Verhalten unverändert. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
178 lines
5.3 KiB
Go
178 lines
5.3 KiB
Go
// Integrationstest (SRC-02): echte Postgres- UND Manticore-Instanz,
|
|
// folgt derselben TEST_*-Env-Konvention wie mail/internal/search.
|
|
package indexworker
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"os"
|
|
"testing"
|
|
"time"
|
|
|
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/search"
|
|
)
|
|
|
|
func setupWorkerEnv(t *testing.T) (*Queue, *search.Client) {
|
|
t.Helper()
|
|
if os.Getenv("TEST_TENANT_DSN") == "" || os.Getenv("TEST_MANTICORE_URL") == "" {
|
|
t.Skip("TEST_TENANT_DSN/TEST_MANTICORE_URL nicht gesetzt, Integrationstest übersprungen")
|
|
}
|
|
queue := setupQueue(t, time.Minute)
|
|
searchClient := search.NewClient(os.Getenv("TEST_MANTICORE_URL"))
|
|
if err := searchClient.EnsureSchema(context.Background()); err != nil {
|
|
t.Fatalf("search-schema: %v", err)
|
|
}
|
|
return queue, searchClient
|
|
}
|
|
|
|
func drainQueue(t *testing.T, worker *Worker, maxJobs int) {
|
|
t.Helper()
|
|
ctx := context.Background()
|
|
for i := 0; i < maxJobs; i++ {
|
|
if err := worker.RunOnce(ctx); err != nil {
|
|
if err == ErrNoJobAvailable {
|
|
return
|
|
}
|
|
t.Fatalf("worker.RunOnce: %v", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestIndexJob_MakesMailSearchable ist die geforderte Funktionsprüfung zu
|
|
// Akzeptanzkriterium 1: eine eingereihte Indexierungsaufgabe macht die Mail
|
|
// nach Verarbeitung durch den Worker durchsuchbar.
|
|
func TestIndexJob_MakesMailSearchable(t *testing.T) {
|
|
queue, searchClient := setupWorkerEnv(t)
|
|
ctx := context.Background()
|
|
tenant := "mandant-src02-index"
|
|
|
|
payload, _ := json.Marshal(map[string]any{
|
|
"subject": "Jahresabschluss 2025",
|
|
"body": "Anbei der Jahresabschluss zur Prüfung.",
|
|
})
|
|
if _, err := queue.EnqueueIndex(ctx, tenant, "msg-idx-1", payload); err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
worker := NewWorker(queue, searchClient, "worker-test-index")
|
|
drainQueue(t, worker, 5)
|
|
|
|
results, err := searchClient.Search(ctx, tenant, "Jahresabschluss")
|
|
if err != nil {
|
|
t.Fatalf("search: %v", err)
|
|
}
|
|
found := false
|
|
for _, r := range results {
|
|
if r.MessageID == "msg-idx-1" {
|
|
found = true
|
|
}
|
|
}
|
|
if !found {
|
|
t.Fatalf("erwarteten treffer msg-idx-1 nach indexierung nicht gefunden, habe: %+v", results)
|
|
}
|
|
}
|
|
|
|
// TestDeleteJob_RemovesMailFromSearchResults ist die geforderte
|
|
// Pflichtprüfung 2: Löschung einer Mail entfernt sie zuverlässig aus
|
|
// Suchtreffern.
|
|
func TestDeleteJob_RemovesMailFromSearchResults(t *testing.T) {
|
|
queue, searchClient := setupWorkerEnv(t)
|
|
ctx := context.Background()
|
|
tenant := "mandant-src02-delete"
|
|
|
|
payload, _ := json.Marshal(map[string]any{
|
|
"subject": "Vertraulicher Vorgang Zeta",
|
|
"body": "Nur für internen Gebrauch.",
|
|
})
|
|
if _, err := queue.EnqueueIndex(ctx, tenant, "msg-del-1", payload); err != nil {
|
|
t.Fatalf("enqueue index: %v", err)
|
|
}
|
|
worker := NewWorker(queue, searchClient, "worker-test-delete")
|
|
drainQueue(t, worker, 5)
|
|
|
|
preResults, err := searchClient.Search(ctx, tenant, "Zeta")
|
|
if err != nil {
|
|
t.Fatalf("search vor löschung: %v", err)
|
|
}
|
|
preFound := false
|
|
for _, r := range preResults {
|
|
if r.MessageID == "msg-del-1" {
|
|
preFound = true
|
|
}
|
|
}
|
|
if !preFound {
|
|
t.Fatal("voraussetzung nicht erfüllt: mail vor löschung nicht auffindbar")
|
|
}
|
|
|
|
if _, err := queue.EnqueueDelete(ctx, tenant, "msg-del-1"); err != nil {
|
|
t.Fatalf("enqueue delete: %v", err)
|
|
}
|
|
drainQueue(t, worker, 5)
|
|
|
|
postResults, err := searchClient.Search(ctx, tenant, "Zeta")
|
|
if err != nil {
|
|
t.Fatalf("search nach löschung: %v", err)
|
|
}
|
|
for _, r := range postResults {
|
|
if r.MessageID == "msg-del-1" {
|
|
t.Fatal("gelöschte mail weiterhin in suchtreffern gefunden")
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestConsistency_DatabaseAndIndexMatchOnSample ist die geforderte
|
|
// Pflichtprüfung 3: Konsistenztest vergleicht Datenbankbestand (erfolgreich
|
|
// abgeschlossene Index-Jobs) mit Indexbestand stichprobenartig.
|
|
func TestConsistency_DatabaseAndIndexMatchOnSample(t *testing.T) {
|
|
queue, searchClient := setupWorkerEnv(t)
|
|
ctx := context.Background()
|
|
tenant := "mandant-src02-konsistenz"
|
|
|
|
messageIDs := []string{"msg-konsistenz-1", "msg-konsistenz-2", "msg-konsistenz-3"}
|
|
for _, mid := range messageIDs {
|
|
payload, _ := json.Marshal(map[string]any{
|
|
"subject": "Konsistenzprobe " + mid,
|
|
"body": "Inhalt zur Konsistenzprüfung.",
|
|
})
|
|
if _, err := queue.EnqueueIndex(ctx, tenant, mid, payload); err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
}
|
|
worker := NewWorker(queue, searchClient, "worker-test-konsistenz")
|
|
drainQueue(t, worker, 10)
|
|
|
|
for _, mid := range messageIDs {
|
|
job, err := findJobByMessageID(ctx, queue, tenant, mid)
|
|
if err != nil {
|
|
t.Fatalf("job für %s: %v", mid, err)
|
|
}
|
|
if job.Status != StatusSucceeded {
|
|
t.Fatalf("job für %s hat status %q, erwartet 'succeeded'", mid, job.Status)
|
|
}
|
|
|
|
results, err := searchClient.Search(ctx, tenant, "Konsistenzprobe")
|
|
if err != nil {
|
|
t.Fatalf("search: %v", err)
|
|
}
|
|
present := false
|
|
for _, r := range results {
|
|
if r.MessageID == mid {
|
|
present = true
|
|
}
|
|
}
|
|
if !present {
|
|
t.Fatalf("datenbank meldet job für %s als succeeded, aber index enthält kein passendes dokument", mid)
|
|
}
|
|
}
|
|
}
|
|
|
|
func findJobByMessageID(ctx context.Context, queue *Queue, tenantSlug, messageID string) (*Job, error) {
|
|
var j Job
|
|
err := queue.pool.QueryRow(ctx, `
|
|
SELECT id, job_type, tenant_slug, message_id, status, attempts
|
|
FROM mail_index_jobs WHERE tenant_slug = $1 AND message_id = $2 AND job_type = 'index'
|
|
ORDER BY id DESC LIMIT 1
|
|
`, tenantSlug, messageID).Scan(&j.ID, &j.JobType, &j.TenantSlug, &j.MessageID, &j.Status, &j.Attempts)
|
|
return &j, err
|
|
}
|