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