Files
archivdms/internal/jobqueue/jobqueue.go
T
patrick 9a24ea29e1 FDN-01: repository & projektgerüst
Git-Repository für bestehenden archivdms-Code initialisiert, Branch-/Commit-Konvention (feature/<ticket>-<slug>-Branches, Ticket-Prefix in Commit-Nachricht) etabliert.
2026-08-11 21:27:53 +02:00

259 lines
8.6 KiB
Go

// Package jobqueue ist die Mandanten-faire Arbeitswarteschlange für die
// asynchrone Dokument-Nachverarbeitung (OCR, Taxonomie-Autozuordnung,
// on_upload-Workflows).
//
// Architektur (bewusst schlank):
//
// - Backend: Postgres-Tabelle processing_jobs (internal/storage/
// processing_jobs.go), KEIN Redis/AMQP. Job-Insert und documents-Insert
// laufen in derselben Transaktion, damit nie ein Dokument ohne Job
// entsteht.
// - Worker: Goroutinen IM SELBEN Backend-Prozess, kein separater Dienst und
// kein Container. Anzahl aus der Config (jobqueue.workers).
// - Fairness: der Dispatcher arbeitet Round-Robin über die Mandanten. Pro
// Runde wird je Mandant mit fälligen Jobs GENAU EINER gezogen
// (ClaimNextJobForTenant), erst danach beginnt die nächste Runde. Ein
// Mandant mit 5.000 Batch-Scans kann damit die Verarbeitung der anderen
// Mandanten verzögern, aber nicht aushungern (globales FIFO würde genau
// das tun).
// - Locking: FOR UPDATE SKIP LOCKED, dadurch können beliebig viele Worker
// (und theoretisch mehrere Prozesse) parallel ziehen, ohne dass ein Job
// doppelt läuft.
// - Reaper: hängengebliebene 'processing'-Jobs (Prozess-Neustart, toter
// OCR-Subprozess) werden nach jobqueue.job_timeout_seconds zurückgesetzt
// und mit exponentiellem Backoff neu eingeplant; ab max_retries bleiben
// sie dauerhaft 'failed' (kein Automatik-Retry mehr, manueller Retry ist
// Phase 3/Frontend).
//
// WORM bleibt außen vor: der Worker liest die archivierte Datei nur und
// schreibt ausschließlich abgeleitete Metadaten.
package jobqueue
import (
"context"
"errors"
"fmt"
"log/slog"
"sync"
"time"
"archivdms/config"
"archivdms/internal/storage"
)
// ProcessFunc verarbeitet einen einzelnen Job. Implementiert von
// internal/api.(*Server).ProcessDocumentJob und als Funktionswert
// hereingereicht (statt internal/api zu importieren) — dasselbe Muster wie
// sftpserver.UploadFunc, um einen Import-Zyklus zu vermeiden. Die
// Verdrahtung passiert in cmd/archivdms/main.go.
type ProcessFunc func(ctx context.Context, tenantID, documentID int64, deriveTitle bool) error
// Dispatcher zieht Jobs Round-Robin über die Mandanten und verteilt sie an
// einen Pool von Worker-Goroutinen.
type Dispatcher struct {
cfg config.JobQueueConfig
store *storage.Store
process ProcessFunc
logger *slog.Logger
jobs chan *storage.ProcessingJob
stopOnce sync.Once
stopCh chan struct{}
wg sync.WaitGroup
}
// New erzeugt einen Dispatcher. Start startet Worker, Dispatch-Loop und
// Reaper; Stop fährt alles sauber herunter.
func New(cfg config.JobQueueConfig, store *storage.Store, process ProcessFunc, logger *slog.Logger) *Dispatcher {
return &Dispatcher{
cfg: cfg,
store: store,
process: process,
logger: logger,
// Ungepuffert: der Dispatch-Loop blockiert, solange alle Worker
// beschäftigt sind. Genau erwünscht — so werden nur so viele Jobs auf
// 'processing' gesetzt, wie auch tatsächlich gerade laufen können, und
// ein Prozess-Neustart lässt keine unnötig große Menge Jobs im
// Reaper-Timeout hängen.
jobs: make(chan *storage.ProcessingJob),
stopCh: make(chan struct{}),
}
}
// Start startet den Worker-Pool, den Round-Robin-Dispatch-Loop und den
// Reaper. Nicht blockierend.
func (d *Dispatcher) Start(ctx context.Context) {
workers := d.cfg.ResolvedWorkers()
for i := 0; i < workers; i++ {
d.wg.Add(1)
go d.worker(ctx, i+1)
}
d.wg.Add(2)
go d.dispatchLoop(ctx)
go d.reapLoop(ctx)
d.logger.Info("job queue started",
"workers", workers,
"poll_interval", d.cfg.ResolvedPollInterval(),
"job_timeout", d.cfg.ResolvedJobTimeout(),
"max_retries", d.cfg.ResolvedMaxRetries())
}
// Stop signalisiert allen Goroutinen das Ende und wartet auf sie.
func (d *Dispatcher) Stop() {
d.stopOnce.Do(func() { close(d.stopCh) })
d.wg.Wait()
}
// dispatchLoop pollt die Queue und verteilt Jobs Round-Robin über Mandanten.
func (d *Dispatcher) dispatchLoop(ctx context.Context) {
defer d.wg.Done()
ticker := time.NewTicker(d.cfg.ResolvedPollInterval())
defer ticker.Stop()
for {
select {
case <-d.stopCh:
close(d.jobs)
return
case <-ctx.Done():
close(d.jobs)
return
case <-ticker.C:
d.dispatchRound(ctx)
}
}
}
// dispatchRound führt so lange Round-Robin-Runden aus, wie noch Mandanten
// mit fälligen Jobs übrig sind. Pro Runde wird je Mandant genau ein Job
// gezogen und an den Worker-Pool übergeben — dadurch wechseln sich die
// Mandanten ab, statt dass Mandant A komplett leergeräumt wird, bevor
// Mandant B drankommt.
func (d *Dispatcher) dispatchRound(ctx context.Context) {
for {
tenants, err := d.store.TenantsWithDueJobs(ctx)
if err != nil {
d.logger.Warn("job queue: listing tenants with due jobs failed", "err", err)
return
}
if len(tenants) == 0 {
return
}
dispatched := 0
for _, tenantID := range tenants {
select {
case <-d.stopCh:
return
case <-ctx.Done():
return
default:
}
job, err := d.store.ClaimNextJobForTenant(ctx, tenantID)
if err != nil {
if errors.Is(err, storage.ErrNoJob) {
continue // Runde hat sich zwischenzeitlich erledigt
}
d.logger.Warn("job queue: claim failed", "tenant_id", tenantID, "err", err)
continue
}
select {
case d.jobs <- job:
dispatched++
case <-d.stopCh:
// Beim Herunterfahren den bereits geclaimten Job nicht
// verlieren: sofort wieder einreihen (retry_count bleibt
// unangetastet), sonst müsste erst der Reaper-Timeout
// ablaufen.
if rerr := d.store.RequeueJob(context.Background(), job.ID, job.TenantID); rerr != nil {
d.logger.Warn("job queue: requeue on shutdown failed", "job_id", job.ID, "err", rerr)
}
return
case <-ctx.Done():
if rerr := d.store.RequeueJob(context.Background(), job.ID, job.TenantID); rerr != nil {
d.logger.Warn("job queue: requeue on shutdown failed", "job_id", job.ID, "err", rerr)
}
return
}
}
if dispatched == 0 {
return
}
}
}
// worker verarbeitet Jobs aus dem Kanal, einer nach dem anderen.
func (d *Dispatcher) worker(ctx context.Context, num int) {
defer d.wg.Done()
for job := range d.jobs {
d.runJob(ctx, num, job)
}
}
func (d *Dispatcher) runJob(ctx context.Context, worker int, job *storage.ProcessingJob) {
started := time.Now()
// Eigener Timeout je Job, damit ein hängender OCR-Subprozess einen Worker
// nicht dauerhaft belegt. Bewusst NICHT vom Request-Kontext abgeleitet —
// die Verarbeitung ist vom Upload-Request entkoppelt.
jobCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), d.cfg.ResolvedJobTimeout())
defer cancel()
err := d.process(jobCtx, job.TenantID, job.DocumentID, job.DeriveTitle)
if err != nil {
requeued, mErr := d.store.MarkJobFailed(context.WithoutCancel(ctx), job.ID, job.TenantID, job.DocumentID, err.Error(), d.cfg.ResolvedMaxRetries())
if mErr != nil {
d.logger.Error("job queue: recording job failure failed", "job_id", job.ID, "err", mErr)
}
d.logger.Warn("job queue: job failed",
"worker", worker, "job_id", job.ID, "tenant_id", job.TenantID, "document_id", job.DocumentID,
"retry_count", job.RetryCount, "will_retry", requeued, "duration", time.Since(started), "err", err)
return
}
if err := d.store.MarkJobDone(context.WithoutCancel(ctx), job.ID, job.TenantID, job.DocumentID); err != nil {
d.logger.Error("job queue: marking job done failed", "job_id", job.ID, "err", err)
return
}
d.logger.Info("job queue: job done",
"worker", worker, "job_id", job.ID, "tenant_id", job.TenantID, "document_id", job.DocumentID,
"duration", time.Since(started))
}
// reapLoop setzt regelmäßig hängengebliebene 'processing'-Jobs zurück.
func (d *Dispatcher) reapLoop(ctx context.Context) {
defer d.wg.Done()
interval := d.cfg.ResolvedJobTimeout() / 2
if interval < 5*time.Second {
interval = 5 * time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-d.stopCh:
return
case <-ctx.Done():
return
case <-ticker.C:
n, err := d.store.ReapStaleJobs(ctx, d.cfg.ResolvedJobTimeout(), d.cfg.ResolvedMaxRetries())
if err != nil {
d.logger.Warn("job queue: reaper failed", "err", err)
continue
}
if n > 0 {
d.logger.Warn("job queue: reset stale processing jobs", "count", n, "timeout", d.cfg.ResolvedJobTimeout())
}
}
}
}
// String beschreibt die aktive Konfiguration (Diagnose-/Log-Hilfe).
func (d *Dispatcher) String() string {
return fmt.Sprintf("jobqueue(workers=%d poll=%s timeout=%s max_retries=%d)",
d.cfg.ResolvedWorkers(), d.cfg.ResolvedPollInterval(), d.cfg.ResolvedJobTimeout(), d.cfg.ResolvedMaxRetries())
}