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.
This commit is contained in:
@@ -0,0 +1,258 @@
|
||||
// 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())
|
||||
}
|
||||
Reference in New Issue
Block a user