IMP-05: hot-folder-scanner-anbindung

Anbindung eines Hot-Folder/Scanner-Eingangs für E-Mail-Anhänge/
Dokumente außerhalb des IMAP-Postfachs, analog zum Ingestion-Pfad.

- store.go: Postgres-Store verzeichnet bereits importierte Dateien je
  Mandant/Postfach über SHA-256-Inhalts-Hash.
- watcher.go: ScanOnce verarbeitet den Eingangsordner, verschiebt
  Duplikate unauffällig und Verarbeitungsfehler gezielt in den
  Fehlerordner, ohne den Scan zu blockieren. Watch nutzt echtes fsnotify
  für Live-Ereignisse plus initialen ScanOnce beim Start.
- Neue minimale Abhängigkeit github.com/fsnotify/fsnotify ergänzt.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-05-PRUEFPROTOKOLL.md):
1. TestScanOnce_SameFileDroppedTwiceImportedOnce: identischer Inhalt
   unter zwei Dateinamen real nur einmal importiert.
2. TestScanOnce_CorruptFileMovedToErrorFolderTraceably: defekte Datei
   real im Fehlerordner, gute Nachbardatei real trotzdem verarbeitet.
3. TestScanOnce_ManyCyclesWithoutResourceLeak: 50 reale Zyklen ohne
   Goroutine-Leck.
Zusätzlich TestWatch_RealFsnotifyEventTriggersImport für die benannte
Technik.

Kein Umbau: kein bestehendes Paket angefasst.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
This commit is contained in:
sysops
2026-09-01 00:22:56 +02:00
co-authored by Claude Sonnet 5
parent dac7854440
commit 6e01cecca7
7 changed files with 542 additions and 0 deletions
@@ -0,0 +1,8 @@
CREATE TABLE IF NOT EXISTS mail_hotfolder_processed (
tenant_slug TEXT NOT NULL,
mailbox_name TEXT NOT NULL,
content_hash TEXT NOT NULL,
filename TEXT NOT NULL,
processed_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (tenant_slug, mailbox_name, content_hash)
)
+60
View File
@@ -0,0 +1,60 @@
// Package hotfolder implementiert IMP-05: Anbindung eines Hot-Folder/
// Scanner-Eingangs für E-Mail-Anhänge/Dokumente außerhalb des
// IMAP-Postfachs, analog zum Ingestion-Pfad. Kein Vorbild in archivmail
// für diesen Zuschnitt — Neubau.
package hotfolder
import (
"context"
_ "embed"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
//go:embed migrations/0001_mail_hotfolder_processed.sql
var schemaMigration string
// Store verzeichnet bereits verarbeitete Dateien je Mandant/Postfach
// über deren Inhalts-Hash — Grundlage für Akzeptanzkriterium 2 (kein
// Doppelimport bei identischem Inhalt, auch unter neuem Dateinamen).
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
func (s *Store) EnsureSchema(ctx context.Context) error {
if _, err := s.pool.Exec(ctx, schemaMigration); err != nil {
return fmt.Errorf("hotfolder: schema anlegen: %w", err)
}
return nil
}
// IsProcessed prüft, ob contentHash für tenantSlug/mailboxName bereits
// erfolgreich importiert wurde.
func (s *Store) IsProcessed(ctx context.Context, tenantSlug, mailboxName, contentHash string) (bool, error) {
var exists bool
err := s.pool.QueryRow(ctx, `
SELECT EXISTS(SELECT 1 FROM mail_hotfolder_processed WHERE tenant_slug = $1 AND mailbox_name = $2 AND content_hash = $3)
`, tenantSlug, mailboxName, contentHash).Scan(&exists)
if err != nil {
return false, fmt.Errorf("hotfolder: verarbeitungsstatus prüfen: %w", err)
}
return exists, nil
}
// MarkProcessed verzeichnet contentHash als erfolgreich importiert.
func (s *Store) MarkProcessed(ctx context.Context, tenantSlug, mailboxName, contentHash, filename string) error {
if _, err := s.pool.Exec(ctx, `
INSERT INTO mail_hotfolder_processed (tenant_slug, mailbox_name, content_hash, filename)
VALUES ($1, $2, $3, $4)
ON CONFLICT (tenant_slug, mailbox_name, content_hash) DO NOTHING
`, tenantSlug, mailboxName, contentHash, filename); err != nil {
return fmt.Errorf("hotfolder: als verarbeitet markieren: %w", err)
}
return nil
}
+197
View File
@@ -0,0 +1,197 @@
package hotfolder
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"os"
"path/filepath"
"github.com/fsnotify/fsnotify"
)
// Handler verarbeitet eine erkannte, noch nicht importierte Datei.
// Echte Ablage/Indexierung ist Sache späterer Kacheln — dieses Paket
// bereitet nur die Schnittstelle vor.
type Handler interface {
ProcessFile(ctx context.Context, tenantSlug, mailboxName, filename string, content []byte) error
}
// Watcher überwacht EIN Hot-Folder-Verzeichnis für EINEN Mandanten/EIN
// Postfach (Akzeptanzkriterium 1: Zuordnung ist strukturell — welches
// Verzeichnis zu welchem Mandanten/Postfach gehört, entscheidet der
// Aufrufer beim Konfigurieren des Watchers, nicht dieses Paket anhand
// von Dateiinhalten).
type Watcher struct {
tenantSlug string
mailboxName string
watchDir string
processedDir string
errorDir string
store *Store
handler Handler
}
// NewWatcher legt processedDir/errorDir an, falls sie noch nicht
// existieren.
func NewWatcher(tenantSlug, mailboxName, watchDir, processedDir, errorDir string, store *Store, handler Handler) (*Watcher, error) {
for _, dir := range []string{watchDir, processedDir, errorDir} {
if err := os.MkdirAll(dir, 0o755); err != nil {
return nil, fmt.Errorf("hotfolder: verzeichnis %s anlegen: %w", dir, err)
}
}
return &Watcher{
tenantSlug: tenantSlug,
mailboxName: mailboxName,
watchDir: watchDir,
processedDir: processedDir,
errorDir: errorDir,
store: store,
handler: handler,
}, nil
}
// ScanResult fasst einen abgeschlossenen Scan-Durchlauf zusammen.
type ScanResult struct {
Imported int
Duplicate int
Failed int
}
// ScanOnce verarbeitet alle regulären Dateien, die aktuell direkt in
// watchDir liegen (nicht rekursiv, processedDir/errorDir liegen
// außerhalb von watchDir und werden dadurch nie mit gescannt). Eine
// einzelne fehlerhafte Datei blockiert NICHT die übrigen
// (Akzeptanzkriterium 3) — sie landet im Fehlerordner, der Scan läuft
// mit der nächsten Datei weiter.
func (w *Watcher) ScanOnce(ctx context.Context) (ScanResult, error) {
entries, err := os.ReadDir(w.watchDir)
if err != nil {
return ScanResult{}, fmt.Errorf("hotfolder: verzeichnis lesen: %w", err)
}
var result ScanResult
for _, e := range entries {
if e.IsDir() {
continue
}
if err := ctx.Err(); err != nil {
return result, err
}
outcome := w.processOne(ctx, e.Name())
switch outcome {
case outcomeImported:
result.Imported++
case outcomeDuplicate:
result.Duplicate++
case outcomeFailed:
result.Failed++
}
}
return result, nil
}
type outcome int
const (
outcomeImported outcome = iota
outcomeDuplicate
outcomeFailed
)
// processOne verarbeitet GENAU EINE Datei — Fehler auf Dateiebene werden
// hier abgefangen (Fehlerordner statt Abbruch), niemals nach oben
// durchgereicht.
func (w *Watcher) processOne(ctx context.Context, filename string) outcome {
fullPath := filepath.Join(w.watchDir, filename)
content, err := os.ReadFile(fullPath)
if err != nil {
// Datei zwischen ReadDir und ReadFile verschwunden (z. B. vom
// Scanner noch nicht vollständig geschrieben) — kein Fehlerordner-
// Umzug möglich, einfach überspringen, nächster Scan versucht es
// erneut.
return outcomeFailed
}
hash := sha256.Sum256(content)
contentHash := hex.EncodeToString(hash[:])
alreadyDone, err := w.store.IsProcessed(ctx, w.tenantSlug, w.mailboxName, contentHash)
if err != nil {
w.moveTo(fullPath, w.errorDir, filename)
return outcomeFailed
}
if alreadyDone {
// Akzeptanzkriterium 2: identischer Inhalt wird nicht doppelt
// importiert — die redundante Kopie wandert unauffällig in den
// Verarbeitet-Ordner, ohne den Handler erneut aufzurufen.
w.moveTo(fullPath, w.processedDir, filename)
return outcomeDuplicate
}
if err := w.handler.ProcessFile(ctx, w.tenantSlug, w.mailboxName, filename, content); err != nil {
w.moveTo(fullPath, w.errorDir, filename)
return outcomeFailed
}
if err := w.store.MarkProcessed(ctx, w.tenantSlug, w.mailboxName, contentHash, filename); err != nil {
w.moveTo(fullPath, w.errorDir, filename)
return outcomeFailed
}
w.moveTo(fullPath, w.processedDir, filename)
return outcomeImported
}
// moveTo verschiebt eine Datei in ein Zielverzeichnis (Akzeptanzkriterium
// 3: Fehlerordner statt Blockade). Ein Fehlschlag beim Verschieben selbst
// wird bewusst nur best-effort behandelt — die Datei bleibt dann im
// Quellverzeichnis stehen und würde beim nächsten Scan erneut
// verarbeitet, was für bereits verarbeitete/fehlerhafte Dateien
// unschädlich ist (Store verhindert Doppelimport, ein wiederholter
// Fehlschlag landet wieder im Fehlerordner).
func (w *Watcher) moveTo(sourcePath, targetDir, filename string) {
_ = os.Rename(sourcePath, filepath.Join(targetDir, filename))
}
// Watch beobachtet watchDir live über fsnotify UND führt zu Beginn einen
// initialen ScanOnce aus (bereits vorhandene Dateien beim Start).
// Blockiert, bis ctx beendet wird.
func (w *Watcher) Watch(ctx context.Context) error {
if _, err := w.ScanOnce(ctx); err != nil {
return err
}
fsWatcher, err := fsnotify.NewWatcher()
if err != nil {
return fmt.Errorf("hotfolder: fsnotify-watcher erstellen: %w", err)
}
defer func() { _ = fsWatcher.Close() }()
if err := fsWatcher.Add(w.watchDir); err != nil {
return fmt.Errorf("hotfolder: verzeichnis beobachten: %w", err)
}
for {
select {
case <-ctx.Done():
return nil
case event, ok := <-fsWatcher.Events:
if !ok {
return nil
}
if event.Op&(fsnotify.Create|fsnotify.Write) == 0 {
continue
}
if _, err := w.ScanOnce(ctx); err != nil {
return err
}
case err, ok := <-fsWatcher.Errors:
if !ok {
return nil
}
return fmt.Errorf("hotfolder: fsnotify-fehler: %w", err)
}
}
}
+214
View File
@@ -0,0 +1,214 @@
// Integrationstest (IMP-05): echte Postgres-Instanz UND echtes
// Dateisystem, folgt derselben Testhost-Konvention wie
// mail/internal/dedup/folderstate — TEST_TENANT_DSN.
package hotfolder
import (
"context"
"os"
"path/filepath"
"runtime"
"sync"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
// recordingHandler zeichnet verarbeitete Dateien auf, kann gezielt für
// bestimmte Dateinamen fehlschlagen (simuliert eine defekte Datei).
type recordingHandler struct {
processed []string
failNames map[string]bool
}
func (h *recordingHandler) ProcessFile(_ context.Context, _, _, filename string, _ []byte) error {
if h.failNames[filename] {
return errFakeCorrupt
}
h.processed = append(h.processed, filename)
return nil
}
var errFakeCorrupt = &corruptFileError{}
type corruptFileError struct{}
func (*corruptFileError) Error() string { return "hotfolder: simuliert defekte datei" }
func setupWatcher(t *testing.T, handler Handler) (*Watcher, string) {
t.Helper()
dsn := os.Getenv("TEST_TENANT_DSN")
if dsn == "" {
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("pool: %v", err)
}
t.Cleanup(func() { pool.Close() })
store := NewStore(pool)
if err := store.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
tenant := "mandant-imp05-hotfolder"
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_hotfolder_processed WHERE tenant_slug LIKE 'mandant-%'`)
})
root := t.TempDir()
watchDir := filepath.Join(root, "eingang")
processedDir := filepath.Join(root, "verarbeitet")
errorDir := filepath.Join(root, "fehler")
watcher, err := NewWatcher(tenant, "INBOX", watchDir, processedDir, errorDir, store, handler)
if err != nil {
t.Fatalf("newwatcher: %v", err)
}
return watcher, watchDir
}
func writeFile(t *testing.T, dir, name, content string) {
t.Helper()
if err := os.WriteFile(filepath.Join(dir, name), []byte(content), 0o600); err != nil {
t.Fatalf("datei schreiben: %v", err)
}
}
// TestScanOnce_SameFileDroppedTwiceImportedOnce ist die geforderte
// Pflichtprüfung 1: gleiche Datei zweimal abgelegt wird nur einmal
// importiert.
func TestScanOnce_SameFileDroppedTwiceImportedOnce(t *testing.T) {
handler := &recordingHandler{failNames: map[string]bool{}}
watcher, watchDir := setupWatcher(t, handler)
ctx := context.Background()
writeFile(t, watchDir, "rechnung.pdf", "identischer inhalt")
result1, err := watcher.ScanOnce(ctx)
if err != nil {
t.Fatalf("erster scan: %v", err)
}
if result1.Imported != 1 {
t.Fatalf("erwartete 1 import im ersten scan, habe %d", result1.Imported)
}
// "Zweimal abgelegt": derselbe Inhalt landet unter NEUEM Dateinamen
// erneut im Eingang (z. B. Scanner mit Zeitstempel-Dateinamen).
writeFile(t, watchDir, "rechnung_kopie.pdf", "identischer inhalt")
result2, err := watcher.ScanOnce(ctx)
if err != nil {
t.Fatalf("zweiter scan: %v", err)
}
if result2.Imported != 0 {
t.Fatalf("erwartete 0 importe im zweiten scan (identischer inhalt bereits verarbeitet), habe %d", result2.Imported)
}
if result2.Duplicate != 1 {
t.Fatalf("erwartete 1 erkanntes duplikat, habe %d", result2.Duplicate)
}
if len(handler.processed) != 1 {
t.Fatalf("handler wurde erwartet genau 1x aufgerufen, habe %d: %v", len(handler.processed), handler.processed)
}
}
// TestScanOnce_CorruptFileMovedToErrorFolderTraceably ist die geforderte
// Pflichtprüfung 2: fehlerhafte Datei landet nachvollziehbar im
// Fehlerordner.
func TestScanOnce_CorruptFileMovedToErrorFolderTraceably(t *testing.T) {
handler := &recordingHandler{failNames: map[string]bool{"defekt.pdf": true}}
watcher, watchDir := setupWatcher(t, handler)
ctx := context.Background()
writeFile(t, watchDir, "defekt.pdf", "kaputter inhalt")
writeFile(t, watchDir, "gut.pdf", "guter inhalt")
result, err := watcher.ScanOnce(ctx)
if err != nil {
t.Fatalf("scan: %v", err)
}
if result.Failed != 1 || result.Imported != 1 {
t.Fatalf("erwartete 1 fehler + 1 import, habe: %+v", result)
}
if _, err := os.Stat(filepath.Join(watcher.errorDir, "defekt.pdf")); err != nil {
t.Fatalf("defekte datei liegt nicht nachvollziehbar im fehlerordner: %v", err)
}
if _, err := os.Stat(filepath.Join(watchDir, "defekt.pdf")); !os.IsNotExist(err) {
t.Fatal("defekte datei liegt noch im eingangsordner — hätte verschoben werden müssen")
}
if _, err := os.Stat(filepath.Join(watcher.processedDir, "gut.pdf")); err != nil {
t.Fatalf("die GUTE datei sollte trotz des defekten nachbarn real verarbeitet worden sein: %v", err)
}
}
// TestScanOnce_ManyCyclesWithoutResourceLeak ist die geforderte
// Pflichtprüfung 3: Dauertest über mehrere Scan-Zyklen ohne
// Ressourcenleck.
func TestScanOnce_ManyCyclesWithoutResourceLeak(t *testing.T) {
handler := &recordingHandler{failNames: map[string]bool{}}
watcher, watchDir := setupWatcher(t, handler)
ctx := context.Background()
before := runtime.NumGoroutine()
const cycles = 50
for i := 0; i < cycles; i++ {
writeFile(t, watchDir, "datei.txt", "inhalt-zyklus")
if _, err := watcher.ScanOnce(ctx); err != nil {
t.Fatalf("scan-zyklus %d: %v", i, err)
}
// Jeder Zyklus legt DIESELBE Datei erneut ab (identischer Inhalt,
// gleicher Dateiname) — nach dem ersten Mal muss jeder weitere
// Zyklus real als Duplikat erkannt werden, kein Ressourcenverbrauch
// pro Zyklus, der sich unbegrenzt aufbaut.
}
after := runtime.NumGoroutine()
// Großzügige Toleranz (Test-Runtime/GC-Hintergrundaktivität) — es
// geht um "kein unbegrenztes Wachstum", nicht um exakte Gleichheit.
if after > before+10 {
t.Fatalf("möglicher goroutine-leck über %d zyklen: vorher=%d nachher=%d", cycles, before, after)
}
entries, err := os.ReadDir(watcher.processedDir)
if err != nil {
t.Fatalf("verarbeitet-ordner lesen: %v", err)
}
if len(entries) != 1 {
t.Fatalf("erwartete genau 1 datei im verarbeitet-ordner nach %d zyklen (immer dieselbe verschoben/dedupliziert), habe %d", cycles, len(entries))
}
}
// TestWatch_RealFsnotifyEventTriggersImport belegt real die im Ticket
// benannte Technik (fsnotify): eine neu abgelegte Datei wird über ein
// echtes Dateisystem-Ereignis erkannt und importiert, ohne dass ein
// manueller ScanOnce-Aufruf nötig ist.
func TestWatch_RealFsnotifyEventTriggersImport(t *testing.T) {
handler := &recordingHandler{failNames: map[string]bool{}}
watcher, watchDir := setupWatcher(t, handler)
ctx, cancel := context.WithCancel(context.Background())
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
_ = watcher.Watch(ctx)
}()
t.Cleanup(func() {
cancel()
wg.Wait()
})
time.Sleep(100 * time.Millisecond) // Watcher real gestartet und lauscht
writeFile(t, watchDir, "live-ereignis.txt", "per fsnotify erkannt")
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
if _, err := os.Stat(filepath.Join(watcher.processedDir, "live-ereignis.txt")); err == nil {
return // real per fsnotify erkannt und verarbeitet
}
time.Sleep(20 * time.Millisecond)
}
t.Fatal("datei wurde nicht innerhalb der frist per echtem fsnotify-ereignis importiert")
}