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
This commit is contained in:
co-authored by
Claude Sonnet 5
parent
c3bf8100b1
commit
c5bdde0cc8
@@ -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()
|
||||
)
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user