diff --git a/mail/docs/SRC-02-PRUEFPROTOKOLL.md b/mail/docs/SRC-02-PRUEFPROTOKOLL.md new file mode 100644 index 0000000..520c846 --- /dev/null +++ b/mail/docs/SRC-02-PRUEFPROTOKOLL.md @@ -0,0 +1,58 @@ +# SRC-02 – Prüfprotokoll: Indexierungs-Worker & Synchronisierung + +Voraussetzung SRC-01, ARC-03 (beide Fertig). + +## Umsetzung + +- `mail/internal/indexworker/migrations/0001_mail_index_jobs.sql` — + statisches, versioniertes Schema (`go:embed`) für `mail_index_jobs` + (`job_type` index/delete, `status` pending/processing/succeeded/failed, + `attempts`/`max_attempts`, `available_at`, `locked_at`/`locked_by`). +- `mail/internal/indexworker/queue.go` — `Queue`: `EnqueueIndex`/ + `EnqueueDelete`, `dequeue` (Postgres `FOR UPDATE SKIP LOCKED` + + Stale-Lock-Wiedervorlage, gleiche Konvention wie + `dms/internal/jobqueue` aus FDN-04 — bewusst schlanker, keine DLQ, da + nicht Bestandteil der Akzeptanzkriterien dieser Kachel), `complete`/ + `fail` (arithmetischer Backoff, kein String-Concat für Intervalle), + `Status` (Akzeptanzkriterium 3 als Go-API). +- `mail/internal/indexworker/worker.go` — `Worker.RunOnce`: holt einen + Job, ruft je nach `job_type` `search.Client.Index`/`search.Client.Delete` + auf, markiert abschließend `complete`/`fail`. +- `mail/internal/search`: minimale Erweiterung um `Client.Delete` und + `DocumentID(tenantSlug, messageID)` (deterministische FNV-1a-ID, damit + Index und Delete für dieselbe Mail immer dasselbe Dokument referenzieren, + ohne zusätzlichen Zustand im Worker). +- Kein Umbau: `mail/internal/storage`/`mail/internal/crypto`/ + `mail/internal/encstorage`/`mail/internal/dedup` unverändert; + bestehende `search`-Tests/-Verhalten (SRC-01) unverändert. + +## Prüfungen + +| # | Prüfung | Ergebnis | +|---|---|---| +| 1 | Test: Worker-Neustart mitten im Lauf verliert keinen offenen Auftrag | **bestanden** – `TestDequeue_WorkerCrashMidRunLosesNoJob`: Job wird geholt und NICHT abgeschlossen (simulierter Absturz), vor Ablauf der Stale-Lock-Frist real kein zweiter Job verfügbar, nach Ablauf real erneut derselbe Job an einen zweiten Worker zugestellt | +| 2 | Test: Löschung einer Mail entfernt sie zuverlässig aus Suchtreffern | **bestanden** – `TestDeleteJob_RemovesMailFromSearchResults`: Mail indexiert und Auffindbarkeit real bestätigt, danach Lösch-Job verarbeitet, anschließende Suche liefert real keinen Treffer mehr | +| 3 | Konsistenztest vergleicht Datenbankbestand mit Indexbestand stichprobenartig | **bestanden** – `TestConsistency_DatabaseAndIndexMatchOnSample`: 3 Index-Jobs verarbeitet, je Stichprobe real geprüft, dass der DB-Job-Status `succeeded` UND das zugehörige Dokument tatsächlich im Manticore-Index auffindbar sind | + +Zusätzlich (Akzeptanzkriterium 1, Funktionsnachweis): `TestIndexJob_MakesMailSearchable` +— eingereihte Indexierungsaufgabe macht die Mail nach Worker-Verarbeitung +real durchsuchbar. + +## Build/Test-Ergebnis (192.168.1.131) + +``` +go build ./... -> clean +go vet ./... -> clean +golangci-lint run ./... -> 0 issues +TEST_TENANT_DSN=postgresql://nexarch_test:***@localhost:5432/tenant_acme?sslmode=disable \ +TEST_MANTICORE_URL=http://127.0.0.1:9308 \ + go test ./... -v -p 1 -> alle Pakete bestanden, inkl. internal/indexworker (5 Tests) +``` + +## Gesamtergebnis + +**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen +real erfüllt. SRC-02 ist der nächste Schritt in der Suche-Foundation-Kette +(Manticore-Schema → Schreib-/Suchzugriff → asynchrone Synchronisierung), +nicht nur eine nette Ergänzung — ohne ihn bliebe SRC-01 ein Index ohne +Befüllungspfad. Entsperrt QA-03. diff --git a/mail/internal/indexworker/migrations/0001_mail_index_jobs.sql b/mail/internal/indexworker/migrations/0001_mail_index_jobs.sql new file mode 100644 index 0000000..dd87808 --- /dev/null +++ b/mail/internal/indexworker/migrations/0001_mail_index_jobs.sql @@ -0,0 +1,16 @@ +CREATE TABLE IF NOT EXISTS mail_index_jobs ( + id BIGSERIAL PRIMARY KEY, + job_type TEXT NOT NULL CHECK (job_type IN ('index', 'delete')), + tenant_slug TEXT NOT NULL, + message_id TEXT NOT NULL, + payload JSONB NOT NULL DEFAULT '{}'::jsonb, + status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'processing', 'succeeded', 'failed')), + 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() +) diff --git a/mail/internal/indexworker/queue.go b/mail/internal/indexworker/queue.go new file mode 100644 index 0000000..4853fdd --- /dev/null +++ b/mail/internal/indexworker/queue.go @@ -0,0 +1,213 @@ +// Package indexworker implementiert SRC-02: einen Indexierungs-Worker, +// der neu archivierte Mails asynchron in den Manticore-Index (SRC-01) +// einpflegt und Löschungen/Metadatenänderungen nachzieht. Postgres- +// Jobqueue mit FOR UPDATE SKIP LOCKED, Stale-Lock-Wiedervorlage bei +// Worker-Absturz — dieselbe Konvention wie dms/internal/jobqueue (FDN-04), +// hier bewusst schlanker (kein Redis/AMQP, keine DLQ — nicht Bestandteil +// der Akzeptanzkriterien dieser Kachel). +package indexworker + +import ( + "context" + _ "embed" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +//go:embed migrations/0001_mail_index_jobs.sql +var schemaMigration string + +const ( + JobTypeIndex = "index" + JobTypeDelete = "delete" +) + +const ( + StatusPending = "pending" + StatusProcessing = "processing" + StatusSucceeded = "succeeded" + StatusFailed = "failed" +) + +// ErrNoJobAvailable wird von Dequeue geliefert, wenn aktuell kein +// abholbarer Job vorhanden ist (Normalfall bei leerer Queue). +var ErrNoJobAvailable = errors.New("indexworker: kein job verfügbar") + +// ErrNotFound wird geliefert, wenn ein angefragter Job nicht existiert. +var ErrNotFound = errors.New("indexworker: job nicht gefunden") + +const defaultMaxAttempts = 5 + +// Job ist eine einzelne Indexierungs-/Löschaufgabe. +type Job struct { + ID int64 + JobType string + TenantSlug string + MessageID string + Status string + Attempts int +} + +// Queue kapselt den Zugriff auf mail_index_jobs. +type Queue struct { + pool *pgxpool.Pool + staleLockAfter time.Duration +} + +// NewQueue erzeugt eine Queue. staleLockAfter legt fest, ab wann ein +// als "processing" markierter Job wieder abholbar gilt, weil sein Worker +// vermutlich abgestürzt ist (Akzeptanzkriterium 3: Worker-Ausfall verliert +// keine Indexierungsaufträge). +func NewQueue(pool *pgxpool.Pool, staleLockAfter time.Duration) *Queue { + return &Queue{pool: pool, staleLockAfter: staleLockAfter} +} + +// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert — +// gleiches Muster wie mail/internal/dedup (kein zentraler Migrationsläufer +// für Mandanten-Datenbanken im Mail-Modul vorhanden). +func (q *Queue) EnsureSchema(ctx context.Context) error { + if _, err := q.pool.Exec(ctx, schemaMigration); err != nil { + return fmt.Errorf("indexworker: schema anlegen: %w", err) + } + return nil +} + +// EnqueueIndex reiht eine Indexierungsaufgabe ein (Akzeptanzkriterium 1). +// payload enthält die für die Indexierung nötigen Felder (Betreff, Text, +// Anhangstext) als JSON — der Worker kennt keine Klartext-Beschaffung +// selbst, das ist Aufgabe des Aufrufers (analog dedup, das ebenfalls +// storage/crypto nicht kennt). +func (q *Queue) EnqueueIndex(ctx context.Context, tenantSlug, messageID string, payload []byte) (int64, error) { + return q.enqueue(ctx, JobTypeIndex, tenantSlug, messageID, payload) +} + +// EnqueueDelete reiht eine Löschaufgabe ein (Akzeptanzkriterium 2). +func (q *Queue) EnqueueDelete(ctx context.Context, tenantSlug, messageID string) (int64, error) { + return q.enqueue(ctx, JobTypeDelete, tenantSlug, messageID, []byte(`{}`)) +} + +func (q *Queue) enqueue(ctx context.Context, jobType, tenantSlug, messageID string, payload []byte) (int64, error) { + var id int64 + err := q.pool.QueryRow(ctx, ` + INSERT INTO mail_index_jobs (job_type, tenant_slug, message_id, payload, max_attempts) + VALUES ($1, $2, $3, $4, $5) + RETURNING id + `, jobType, tenantSlug, messageID, payload, defaultMaxAttempts).Scan(&id) + if err != nil { + return 0, fmt.Errorf("indexworker: job einreihen: %w", err) + } + return id, nil +} + +// dequeuedJob trägt zusätzlich den Payload, den nur das Paket selbst +// (worker.go) benötigt. +type dequeuedJob struct { + Job + Payload []byte +} + +// Dequeue holt GENAU EINEN abholbaren Job (fällig UND nicht gesperrt, ODER +// dessen Sperre abgestanden ist) und markiert ihn atomar als "processing" +// (Prüfung: Worker-Neustart mitten im Lauf verliert keinen offenen Auftrag +// — FOR UPDATE SKIP LOCKED erlaubt mehreren Worker-Goroutinen gleichzeitigen +// Aufruf ohne denselben Job doppelt zu holen). +func (q *Queue) dequeue(ctx context.Context, workerID string) (*dequeuedJob, error) { + tx, err := q.pool.Begin(ctx) + if err != nil { + return nil, fmt.Errorf("indexworker: transaktion starten: %w", err) + } + defer func() { _ = tx.Rollback(ctx) }() + + row := tx.QueryRow(ctx, ` + SELECT id, job_type, tenant_slug, message_id, payload, status, attempts + FROM mail_index_jobs + WHERE ( + (status = 'pending' AND available_at <= now()) + OR (status = 'processing' AND locked_at <= now() - ($1 * interval '1 second')) + ) + ORDER BY available_at + FOR UPDATE SKIP LOCKED + LIMIT 1 + `, q.staleLockAfter.Seconds()) + + var j dequeuedJob + if err := row.Scan(&j.ID, &j.JobType, &j.TenantSlug, &j.MessageID, &j.Payload, &j.Status, &j.Attempts); err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return nil, ErrNoJobAvailable + } + return nil, fmt.Errorf("indexworker: nächsten job lesen: %w", err) + } + + if _, err := tx.Exec(ctx, ` + UPDATE mail_index_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("indexworker: job sperren: %w", err) + } + + if err := tx.Commit(ctx); err != nil { + return nil, fmt.Errorf("indexworker: 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 int64) error { + tag, err := q.pool.Exec(ctx, ` + UPDATE mail_index_jobs SET status = 'succeeded', locked_at = NULL, locked_by = NULL, updated_at = now() + WHERE id = $1 + `, jobID) + if err != nil { + return fmt.Errorf("indexworker: job abschließen: %w", err) + } + if tag.RowsAffected() == 0 { + return ErrNotFound + } + return nil +} + +// fail markiert einen Job als fehlgeschlagen. Sind die maximalen Versuche +// erreicht, bleibt er dauerhaft 'failed' (keine DLQ, nicht Bestandteil +// dieser Kachel), sonst wird er mit arithmetischem Backoff (kein +// String-Concat für Intervalle) erneut eingeplant. +func (q *Queue) fail(ctx context.Context, jobID int64, cause error) error { + tag, err := q.pool.Exec(ctx, ` + UPDATE mail_index_jobs + SET status = CASE WHEN attempts >= max_attempts THEN 'failed' ELSE 'pending' END, + available_at = now() + (LEAST(attempts, 10) * interval '10 seconds'), + locked_at = NULL, locked_by = NULL, last_error = $2, updated_at = now() + WHERE id = $1 + `, jobID, cause.Error()) + if err != nil { + return fmt.Errorf("indexworker: fehlschlag erfassen: %w", err) + } + if tag.RowsAffected() == 0 { + return ErrNotFound + } + return nil +} + +// Status liefert den aktuellen Zustand eines Jobs (abrufbar über API, +// hier als Go-API — HTTP-Anbindung ist nicht Bestandteil dieser Kachel). +func (q *Queue) Status(ctx context.Context, jobID int64) (*Job, error) { + var j Job + err := q.pool.QueryRow(ctx, ` + SELECT id, job_type, tenant_slug, message_id, status, attempts + FROM mail_index_jobs WHERE id = $1 + `, jobID).Scan(&j.ID, &j.JobType, &j.TenantSlug, &j.MessageID, &j.Status, &j.Attempts) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("indexworker: job-status lesen: %w", err) + } + return &j, nil +} diff --git a/mail/internal/indexworker/queue_test.go b/mail/internal/indexworker/queue_test.go new file mode 100644 index 0000000..3d4fc66 --- /dev/null +++ b/mail/internal/indexworker/queue_test.go @@ -0,0 +1,100 @@ +// 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) + } +} diff --git a/mail/internal/indexworker/worker.go b/mail/internal/indexworker/worker.go new file mode 100644 index 0000000..8ef576f --- /dev/null +++ b/mail/internal/indexworker/worker.go @@ -0,0 +1,74 @@ +package indexworker + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "gitea.perlbach24.de/scripte/nexarch/mail/internal/search" +) + +// indexPayload sind die für die Indexierung nötigen Felder, wie sie beim +// EnqueueIndex als JSON übergeben werden. +type indexPayload struct { + Subject string `json:"subject"` + Body string `json:"body"` + AttachmentText string `json:"attachment_text"` + SentAtUnixEpoch int64 `json:"sent_at"` +} + +// Worker holt Jobs aus der Queue und pflegt sie in den Manticore-Index +// (SRC-01) ein bzw. entfernt sie daraus. +type Worker struct { + queue *Queue + searchClient *search.Client + id string +} + +func NewWorker(queue *Queue, searchClient *search.Client, workerID string) *Worker { + return &Worker{queue: queue, searchClient: searchClient, id: workerID} +} + +// RunOnce verarbeitet genau einen Job, falls vorhanden. Liefert +// ErrNoJobAvailable, wenn die Queue aktuell leer ist — kein Fehlerzustand. +func (w *Worker) RunOnce(ctx context.Context) error { + job, err := w.queue.dequeue(ctx, w.id) + if err != nil { + return err + } + + if procErr := w.process(ctx, job); procErr != nil { + if failErr := w.queue.fail(ctx, job.ID, procErr); failErr != nil { + return fmt.Errorf("indexworker: job %d fehlgeschlagen (%v) UND fehlschlag nicht erfassbar: %w", job.ID, procErr, failErr) + } + return nil + } + + return w.queue.complete(ctx, job.ID) +} + +func (w *Worker) process(ctx context.Context, job *dequeuedJob) error { + docID := search.DocumentID(job.TenantSlug, job.MessageID) + + switch job.JobType { + case JobTypeIndex: + var p indexPayload + if err := json.Unmarshal(job.Payload, &p); err != nil { + return fmt.Errorf("indexierungs-payload lesen: %w", err) + } + return w.searchClient.Index(ctx, search.Document{ + ID: docID, + TenantSlug: job.TenantSlug, + MessageID: job.MessageID, + Subject: p.Subject, + Body: p.Body, + AttachmentText: p.AttachmentText, + SentAtUnixEpoch: p.SentAtUnixEpoch, + }) + case JobTypeDelete: + return w.searchClient.Delete(ctx, docID) + default: + return errors.New("unbekannter job-typ: " + job.JobType) + } +} diff --git a/mail/internal/indexworker/worker_test.go b/mail/internal/indexworker/worker_test.go new file mode 100644 index 0000000..2c66d76 --- /dev/null +++ b/mail/internal/indexworker/worker_test.go @@ -0,0 +1,177 @@ +// 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 +} diff --git a/mail/internal/search/client.go b/mail/internal/search/client.go index 97b6a9c..e5419b7 100644 --- a/mail/internal/search/client.go +++ b/mail/internal/search/client.go @@ -106,6 +106,37 @@ func (c *Client) Index(ctx context.Context, doc Document) error { return nil } +// Delete entfernt ein Suchdokument anhand seiner ID (SRC-02 +// Akzeptanzkriterium 2: Löschungen werden im Index nachgezogen). Löschen +// eines nicht (mehr) vorhandenen Dokuments ist kein Fehler (idempotent). +func (c *Client) Delete(ctx context.Context, id uint64) error { + payload := map[string]any{ + "index": IndexName, + "id": id, + } + body, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("search: lösch-anfrage serialisieren: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/delete", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("search: lösch-anfrage bauen: %w", err) + } + req.Header.Set("Content-Type", "application/json") + + resp, err := c.http.Do(req) + if err != nil { + return fmt.Errorf("search: dokument löschen: %w", err) + } + defer func() { _ = resp.Body.Close() }() + respBody, _ := io.ReadAll(resp.Body) + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("search: dokument löschen, status %d: %s", resp.StatusCode, string(respBody)) + } + return nil +} + // Result ist ein Suchtreffer. type Result struct { MessageID string diff --git a/mail/internal/search/fields.go b/mail/internal/search/fields.go index 976eac2..e5dd4b4 100644 --- a/mail/internal/search/fields.go +++ b/mail/internal/search/fields.go @@ -9,6 +9,8 @@ // String-Zusammenbau von SQL), nicht über die SQL-Schnittstelle. package search +import "hash/fnv" + // IndexName ist der einzige Ort, an dem der Manticore-Indexname als // Literal steht. const IndexName = "mail_documents" @@ -27,3 +29,17 @@ const ( // searchableTextFields sind die Volltextfelder, über die eine Suchanfrage // läuft (Akzeptanzprüfung 3: Volltextsuche liefert erwartete Treffer). var searchableTextFields = []string{FieldSubject, FieldBody, FieldAttachmentText} + +// DocumentID berechnet deterministisch die Manticore-Dokument-ID aus +// Mandant und Message-ID (FNV-1a, 64 Bit). Deterministisch statt einer +// separat vergebenen ID, damit Re-Indexierung (Index) und Löschung +// (Delete) für dieselbe Mail immer dieselbe Dokument-ID referenzieren, +// ohne dass der Aufrufer sie zwischenspeichern muss (SRC-02: Löschungen +// müssen ohne zusätzlichen Zustand nachgezogen werden können). +func DocumentID(tenantSlug, messageID string) uint64 { + h := fnv.New64a() + _, _ = h.Write([]byte(tenantSlug)) + _, _ = h.Write([]byte{0}) + _, _ = h.Write([]byte(messageID)) + return h.Sum64() +}