Files
archivmail/internal/index/worker.go
T
sysopsandClaude Sonnet 5 798cb2817c fix: Crash-Robustheit + Performance in Backend und Frontend härten
Backend: recover() in allen langlebigen Goroutinen (neues internal/safego-
Paket), MIME-Multipart-Tiefenlimit gegen Stack-Overflow, IMAP-Zeilenlängen-
und FETCH-Result-Limits gegen OOM, MBOX-Buffer-Aliasing-Bug (Datenkorruption
beim Import), ungeprüfte Type Assertions abgesichert, Data Race im
API-Key-Rate-Limiter behoben, SMTP-Session-Panic führt jetzt zu 451-Retry
statt Prozessabsturz. Performance: O(n²)-String-Concat in Mailparser und
IMAP-Parser durch strings.Builder ersetzt.

Frontend: Error Boundaries für Root und Mail-Detailansicht ergänzt (gab es
vorher nicht), zahlreiche Guards gegen nil-Slices aus dem Backend-JSON die
sonst .map()/.length-Crashes/White-Screens auslösten, defekte JSON-Antworten
in api/core.ts abgefangen, zwei React-Key-Bugs bei löschbaren Listen
korrigiert.

Verifiziert auf 192.168.1.132: Build und Tests der geänderten Pakete
fehlerfrei, keine Regressionen gegenüber vorbestehendem Stand.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019j28kGcaJAhBnrYX34hGdt
2026-08-05 13:39:00 +02:00

103 lines
2.5 KiB
Go

package index
import (
"fmt"
"log/slog"
"sync"
)
// IndexWorker processes MailDocument indexing requests asynchronously via a
// buffered channel. It serialises writes to the underlying Indexer so that
// backends requiring a single concurrent writer remain safe.
type IndexWorker struct {
idx Indexer
queue chan MailDocument
done chan struct{}
wg sync.WaitGroup
logger *slog.Logger
}
// NewWorker creates a new IndexWorker with the given queue capacity.
func NewWorker(idx Indexer, queueSize int, logger *slog.Logger) *IndexWorker {
if queueSize <= 0 {
queueSize = 1000
}
return &IndexWorker{
idx: idx,
queue: make(chan MailDocument, queueSize),
done: make(chan struct{}),
logger: logger,
}
}
// Submit enqueues a document for background indexing. If the queue is full,
// the document is dropped and a warning is logged.
func (w *IndexWorker) Submit(doc MailDocument) {
select {
case w.queue <- doc:
// queued
default:
w.logger.Warn("index worker: queue full, dropping document", "id", doc.ID)
}
}
// Start launches the background goroutine that processes the queue.
func (w *IndexWorker) Start() {
w.wg.Add(1)
go func() {
defer w.wg.Done()
w.logger.Info("index worker: started", "queue_size", cap(w.queue))
for {
select {
case doc, ok := <-w.queue:
if !ok {
// Channel closed, drain complete
return
}
w.indexDoc(doc, "")
case <-w.done:
// Drain remaining items in the queue before exiting
for {
select {
case doc, ok := <-w.queue:
if !ok {
return
}
w.indexDoc(doc, " (drain)")
default:
return
}
}
}
}
}()
}
// indexDoc indexes a single document. A panic in the backend must not kill the
// long-running worker goroutine (and with it the whole process).
func (w *IndexWorker) indexDoc(doc MailDocument, phase string) {
defer func() {
if r := recover(); r != nil {
w.logger.Error("index worker: recovered from panic"+phase,
"id", doc.ID, "panic", fmt.Sprintf("%v", r))
}
}()
if err := w.idx.IndexSync(doc); err != nil {
w.logger.Error("index worker: index failed"+phase, "id", doc.ID, "err", err)
}
}
// Stop signals the worker to drain remaining items and stop. It blocks until
// the worker goroutine has exited.
func (w *IndexWorker) Stop() {
close(w.done)
w.wg.Wait()
w.logger.Info("index worker: stopped")
}
// QueueLen returns the current number of items waiting in the queue.
func (w *IndexWorker) QueueLen() int {
return len(w.queue)
}