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