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
75 lines
2.1 KiB
Go
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)
|
|
}
|
|
}
|