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
198 lines
5.7 KiB
Go
198 lines
5.7 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|