Compare commits

...
Author SHA1 Message Date
sysopsandClaude Sonnet 5 c5bdde0cc8 SRC-02: indexierungs-worker-synchronisierung
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
2026-08-31 09:47:50 +02:00
8 changed files with 685 additions and 0 deletions
+58
View File
@@ -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.
@@ -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()
)
+213
View File
@@ -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
}
+100
View File
@@ -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)
}
}
+74
View File
@@ -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)
}
}
+177
View File
@@ -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
}
+31
View File
@@ -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
+16
View File
@@ -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()
}