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