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