Files
nexarch/mail/internal/indexworker/worker.go
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

75 lines
2.1 KiB
Go

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)
}
}