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
214 lines
7.2 KiB
Go
214 lines
7.2 KiB
Go
// 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
|
|
}
|