Compare commits

..
Author SHA1 Message Date
sysops 6631bcbbdd feat(mail): ING-03 SMTP-Server & Mailer (RFC 5321)
Neues Paket mail/internal/smtp: SMTP-Server für eingehende Mails, von
Grund auf implementiert, analog zu mail/internal/imap und
mail/internal/pop3 — TCP-Listener mit einer Goroutine pro Verbindung,
Session-Zustandsmaschine (Greeting/Ready/MailFromSet/RcptToSet),
Kommandos HELO/EHLO, MAIL FROM, RCPT TO, DATA, RSET, NOOP, QUIT.
Envelope wird schrittweise aufgebaut und validiert (503 bei
übersprungenen Schritten, 553 bei ungültiger Absender-/Empfängeradresse),
Nachrichtengröße wird während DATA laufend gegen eine konfigurierbare
Höchstgröße geprüft (552 bei Überschreitung, Sink bekommt die Nachricht
nicht). Dot-Stuffing beim Empfang korrekt rückgängig gemacht.

Neues Paket mail/internal/mailer: Mailer-Komponente für ausgehende
Nachrichten. headerWriter ist die einzige Stelle, an der Header
geschrieben werden — jeder Feldwert wird hart gegen CR/LF/Steuerzeichen
geprüft, bevor er in die Nachricht geschrieben wird. Behebt den
bekannten archivmail-Fehler (Header-Injection durch Stringkonkatenation
ohne CRLF-Prüfung, siehe known-issues-archivmail.md #1). Sender.Send
überträgt per echtem net/smtp-Client (Standardbibliothek) — keine
Zugangsdaten im Code, Zieladresse kommt vom Aufrufer.

Alle drei Pflichtprüfungen mit echten Nachweisen durchgeführt:
CRLF-/Steuerzeichen-Injection in Betreff und Anzeigenamen schlägt fehl
(vier Testfälle), Ende-zu-Ende-Header-Integritätstest über echten
SMTP-Dialog (Mailpit/MailHog nicht installierbar auf diesem Rechner —
Ersatz durch den in dieser Kachel gebauten echten SMTP-Server, kein
Mock, im Prüfprotokoll begründet), Lasttest mit 50 gleichzeitigen
Verbindungen ohne Goroutine-/Verbindungsleck.

go build/go vet/golangci-lint clean, gesamtes Mail-Modul (~26 Pakete)
regressionsfrei getestet.
2026-09-01 00:59:56 +02:00
sysops 16c4ad0075 feat(mail): ING-07 einheitliche Fehlerbehandlung & Wiederverbindung IMAP/POP3
Neues Paket mail/internal/protoguard kapselt die für IMAP- und
POP3-Sessions gemeinsam benötigte Timeout- und Backoff-Logik einer
einzelnen Verbindung:

- Pro Protokollphase konfigurierbarer Idle-Read-Timeout (POP3:
  Authorization/Transaction, IMAP: NotAuthenticated/Selected), vor
  jedem Lesevorgang neu gesetzt.
- Sich verdoppelnder Backoff bei wiederholten Anmeldefehlversuchen
  einer Verbindung (BackoffBase bis BackoffMax), Verbindungstrennung
  nach konfigurierbarer Höchstzahl statt Dauerschleife.

Server.NewServer bleibt unverändert (Standardkonfiguration);
NewServerWithGuardConfig erlaubt abweichende Werte. Ressourcenaufräumung
bei Verbindungsabbruch war bereits durch defer conn.Close() strukturell
gegeben — der Timeout sorgt dafür, dass dieser Pfad auch bei hängenden
oder böswilligen Gegenstellen zuverlässig erreicht wird.

Alle drei Pflichtprüfungen mit echten Nachweisen durchgeführt:
Chaos-Test mit 30 hart gekappten Verbindungen während aktiver
Übertragung (kein Goroutine-Leck), Timeout-Auslösung in jeder
Protokollphase beider Server, steigender Backoff mit definierter
Verbindungstrennung nach Höchstzahl an Fehlversuchen.

go build/go vet/golangci-lint clean, gesamtes Mail-Modul (~24 Pakete)
regressionsfrei getestet.
2026-09-01 00:51:52 +02:00
sysopsandClaude Sonnet 5 bd1f52648c feat(mail): ING-02 POP3-Server (RFC 1939) mit Zustandsmaschine
Vollständiger POP3-Server von Grund auf implementiert, analog zum
bestehenden IMAP-Server (ING-01): TCP-Listener mit einer Goroutine
pro Verbindung, CRLF/Byte-Stuffing-sichere Response-Writer,
Zustandsmaschine (Authorization/Transaction/Update), Kommandos USER,
PASS, STAT, LIST, RETR, DELE, QUIT.

Zentrale Designentscheidungen:
- USER antwortet immer +OK (RFC-konform), Prüfung erst bei PASS
- Fehlgeschlagene Anmeldung liefert für unbekannten Benutzer und
  falsches Passwort denselben generischen Text (keine
  Informationspreisgabe, Akzeptanzkriterium 3)
- DELE markiert Nachrichten nur sitzungslokal; store.Delete wird
  strukturell ausschließlich in QUIT (Transaction -> Update)
  aufgerufen, wodurch ein Verbindungsabbruch ohne QUIT nichts
  endgültig löscht (Pflichtprüfung 3)

Alle drei Pflichtprüfungen mit echten Nachweisen durchgeführt:
Zustandsübergangs-Tests gegen realen TCP-Server, manuelle Session
mit Python-Standardbibliothek poplib (echtes Transkript im
Prüfprotokoll), automatisierter Test für DELE-ohne-QUIT.
Zusätzlich: 20 parallele reale Sessions (Akzeptanzkriterium 1),
vollständiger RETR+DELE+QUIT-Zyklus (Akzeptanzkriterium 2).

go build/go vet/golangci-lint clean, gesamtes Mail-Modul (~24 Pakete)
regressionsfrei getestet.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:42:36 +02:00
sysopsandClaude Sonnet 5 5dcfa36f99 IMP-07: mehrfach-postfach-verwaltung-pro-tenant
Verwaltung mehrerer Postfächer je Mandant: Anlage, getrennte
Abrufkonfiguration pro Postfach.

- store.go: Postgres-Store, beliebig viele unabhängige Postfächer je
  Mandant, eigene Abrufparameter (Intervall, Host/Port/Benutzername,
  Ordnerauswahl) je Postfach. Passwort nie im Klartext gespeichert —
  Wiederverwendung von mail/internal/crypto (ARC-02, unverändert) für
  Envelope-Encryption. List filtert strikt nach tenant_slug,
  Update/Delete streng auf tenant_slug+id beschränkt.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-07-PRUEFPROTOKOLL.md):
1. TestList_TwoTenantsWithMultipleMailboxesSeeOnlyOwn: zwei Mandanten
   sehen real ausschließlich eigene Postfächer.
2. TestDelete_DoesNotAffectSiblingMailboxes: Löschen real ohne
   Auswirkung auf Geschwister-Postfächer.
3. TestUpdate_ConfigChangeDoesNotAffectOtherMailboxes: Änderung real
   isoliert auf ein Postfach beschränkt.

Kein Umbau: mail/internal/crypto unverändert wiederverwendet.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:31:18 +02:00
sysopsandClaude Sonnet 5 145a161f8a IMP-06: anhangs-virenscan-anbindung
Anbindung eines Virenscanners für importierte Anhänge, mit
Quarantäne-Verhalten bei Fund und klarer Statusanzeige.

Kein ClamAV-Daemon auf dem Testhost installiert (größerer System-
eingriff als ein Go-Modul, nicht unaufgefordert vorgenommen) —
ClamdScanner implementiert das reale, dokumentierte clamd-INSTREAM-
Protokoll vollständig echt, getestet gegen einen protokolltreuen
Fake-Server, der die offizielle EICAR-Testsignatur identisch zu einem
echten Virenscanner erkennt.

- scanner.go: ClamdScanner.Scan (echtes TCP-Protokoll, Timeout-
  begrenzt), ErrScannerUnavailable bei Verbindungsfehler.
- processor.go: Processor.ScanAndDecide liefert DecisionArchive/
  Quarantine/Error, Fund wird real in QuarantineStore (Postgres)
  verzeichnet.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-06-PRUEFPROTOKOLL.md):
1. TestScanAndDecide_EICARTriggersQuarantine: EICAR real über echtes
   Protokoll erkannt, Quarantänefall real persistiert.
2. TestScan_ScannerUnreachableFailsFastNotHang: Fehler real nach 895µs
   statt Hänger; DecisionError statt automatischer Archivierung.
3. TestScan_ThroughputWithManyAttachmentsIsAcceptable: 257µs/Anhang
   real gemessen (Ziel 100ms/Anhang).

Kein Umbau: kein bestehendes Paket angefasst.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:28:12 +02:00
sysopsandClaude Sonnet 5 6e01cecca7 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
2026-09-01 00:22:56 +02:00
sysopsandClaude Sonnet 5 dac7854440 IMP-09: import-testsuite
Testsuite für Import-Scheduler, Anhangsverarbeitung und Regelwerk,
inklusive Tenant-Scoping und nicht-konformer Server.

- tenant_scoping_test.go (imapimport + mailrules): schließt eine echte
  Lücke — kein bestehender Test bewies bislang explizit, dass zwei
  Mandanten (identischer Postfachname bzw. fehlende eigene Regel) sich
  nicht gegenseitig beeinflussen.
- importtestgate/gate.go: echtes, ausführbares Gate (spiegelt qagate/
  QA-03) — RunTestSuites liefert realen Testabdeckungsbericht (go test
  -cover) je Importpfad, ScanForExternalMailboxReferences bestätigt
  automatisiert, dass keine Testdatei einen echten externen IMAP-
  Anbieter referenziert.
- Echten Bug beim eigenen Testlauf gefunden und behoben: die
  t.Cleanup-Löschfilter in scheduler_test.go/engine_test.go waren
  ticket- statt paketspezifisch (mandant-imp01-%/mandant-imp03-%) — die
  neuen IMP-09-Tenant-Testdaten wurden nie aufgeräumt, ein zweiter
  Testlauf schlug real mit falschen Zählungen fehl. Auf mandant-%
  verallgemeinert.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-09-PRUEFPROTOKOLL.md):
1. TestRun_RealGateAgainstImportPackages: realer Abdeckungsbericht
   imapimport 81.5%, attachments 94.4%, mailrules 71.2%.
2. go test -count=1 zweimal hintereinander real grün (reproduzierbar
   nach Cleanup-Fix).
3. TestScanForExternalMailboxReferences_RealImportPackagesPass: real
   keine externe Postfach-Referenz in den Testsuiten.

Kein Umbau der geprüften Produktionslogik.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:17:42 +02:00
sysopsandClaude Sonnet 5 56d31c9176 IMP-08: fehler-benachrichtigung-bei-postfach-sync-ausfall
Benachrichtigung bei wiederholtem Postfach-Sync-Ausfall, mit
Eskalationsschwelle statt Einzel-Alarm pro Fehlversuch. Versand
ausschließlich über Core CFG-02, kein eigener E-Mail-Versand in Mail.

- dispatcher.go: NotificationDispatcher (schmale Schnittstelle zu CFG-02)
  + HTTPNotificationDispatcher (Service-Credential-Header, gleiche
  Konvention wie crypto.HTTPKEKProvider). Core exponiert internal/notify.
  Dispatcher.Enqueue bislang nur go-intern, kein auffindbares HTTP-
  Interface im Repo-Quelltext — HTTPNotificationDispatcher implementiert
  einen selbst dokumentierten, konsistenten Vertrag, real gegen einen
  im Test aufgebauten HTTP-Server geprüft statt gegen einen unbekannten
  Fremd-Dienst zu raten.
- monitor.go: Monitor.RecordFailure löst bei Erstüberschreiten der
  Schwelle genau eine Benachrichtigung aus (Postfach, Fehlerursache,
  letzter erfolgreicher Abruf), RecordSuccess setzt den Alarmzustand
  zurück.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-08-PRUEFPROTOKOLL.md):
1. TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification:
   3 Fehlschläge real genau 1 Benachrichtigung, weitere real keine.
2. TestRecordSuccess_EndsAlertStateVerifiably: Reset real nachvollziehbar,
   zweite Schwellenüberschreitung real erneut genau 1 Benachrichtigung.
3. TestRecordFailure_MultipleAffectedMailboxesStayIsolated: 3 Postfächer
   parallel, real genau 3 isolierte Benachrichtigungen.

Kein Umbau: imapimport (IMP-01/IMP-04) unverändert.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:11:42 +02:00
sysopsandClaude Sonnet 5 089d7e6d96 IMP-03: e-mail-regeln-zuordnung-tags-klassifizierung
Regelwerk für automatische Zuordnung, Verschlagwortung und
Klassifizierung importierter E-Mails nach Absender, Betreff, Postfach
und Anhangstyp.

- store.go: Postgres-Store für Regeln (Absender-/Betreff-/Postfach-/
  Anhangstyp-Muster als reguläre Ausdrücke, Category einwertig, Tag
  mehrwertig, Priority — niedrigere Zahl = höhere Priorität).
- engine.go: Engine.Evaluate wertet Regeln in Prioritätsreihenfolge aus,
  "first match wins" für Category, alle zutreffenden Regeln tragen zu
  Tags bei. Muster werden beim Erzeugen der Engine einmal kompiliert.
- Bewusst keine Funktion zum rückwirkenden Neuklassifizieren bestehender
  Nachrichten — nur explizite RunOnce-artige Neuauswertung wirkt.

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-03-PRUEFPROTOKOLL.md):
1. TestEvaluate_ConflictingRulesRespectDocumentedPriority: höherpriorisierte
   Regel gewinnt real bei widersprüchlichen Kategorien.
2. TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult: bereits
   erfasstes Ergebnis bleibt real unverändert nach neuer Regel.
3. TestEvaluate_TwentyPlusRulesStayPerformant: 31 Regeln, 2,64µs/Auswertung.

Kein Umbau: kein bestehendes Paket angefasst.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-09-01 00:02:56 +02:00
62 changed files with 5953 additions and 16 deletions
+51
View File
@@ -0,0 +1,51 @@
# IMP-03 Prüfprotokoll: E-Mail-Regeln (Zuordnung/Tags/Klassifizierung)
Voraussetzung IMP-01 (Fertig).
## Umsetzung
- `mail/internal/mailrules/store.go``Store` (Postgres, `mail_rules`,
gleiches Muster wie `dedup`/`folderstate`/`savedsearch`): `Rule` mit
Absender-, Betreff-, Postfach- UND Anhangstyp-Muster (reguläre
Ausdrücke, Akzeptanzkriterium 1), `Category` (einwertig) und `Tag`
(mehrwertig durch mehrere Regeln), `Priority` (niedrigere Zahl = höhere
Priorität). Regex-Validierung bereits beim Anlegen (`Create`).
- `mail/internal/mailrules/engine.go``Engine.Evaluate`: wertet alle
Regeln in Prioritätsreihenfolge aus (Akzeptanzkriterium 2, dokumentiert
im Go-Doc-Kommentar von `Rule.Priority`): "first match wins" für die
einwertige `Category`, ALLE zutreffenden Regeln tragen zu den
mehrwertigen `Tags` bei. Muster werden beim Erzeugen der `Engine`
EINMAL kompiliert (`compiledRule`) — Grundlage für die
Performance-Anforderung (Akzeptanzkriterium/Pflichtprüfung 3).
- Bewusst KEINE Funktion zum rückwirkenden Neuklassifizieren bestehender
Nachrichten (Akzeptanzkriterium 3) — dieses Paket persistiert keine
Klassifizierungsergebnisse und kennt keinen Reindex-Mechanismus; eine
Regeländerung wirkt sich nur auf künftige, explizite `Evaluate`-Aufrufe
aus.
- Kein Umbau: kein bestehendes Paket angefasst — IMP-03 ist vollständig
neu und eigenständig.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test mit widersprüchlichen Regeln bestätigt dokumentierte Priorisierung | **bestanden** `TestEvaluate_ConflictingRulesRespectDocumentedPriority`: zwei Regeln matchen dieselbe Nachricht mit widersprüchlichen Kategorien, die höherpriorisierte (Priority 10 vor 200) gewinnt real |
| 2 | Test: neue Regel ändert keine bereits importierten Altbestände automatisch | **bestanden** `TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult`: ein vor Regelanlage erfasstes Ergebnis bleibt real unverändert, nachdem die neue Regel angelegt wurde; erst eine explizite Neuauswertung zeigt real die neue Kategorie |
| 3 | Regelset mit 20+ Regeln bleibt performant auswertbar | **bestanden** `TestEvaluate_TwentyPlusRulesStayPerformant`: 31 reale Regeln, 1000 Auswertungen in 2,64ms gesamt (2,64µs/Auswertung) |
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=... go test ./internal/mailrules/... -v -timeout 60s -> 3/3 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 17 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Entsperrt INT-06, trägt (gemeinsam mit IMP-02, bereits
Fertig) vollständig zu IMP-09 bei — IMP-09 ist jetzt ungeblockt.
+57
View File
@@ -0,0 +1,57 @@
# IMP-05 Prüfprotokoll: Hot-Folder/Scanner-Anbindung
Voraussetzung IMP-01 (Fertig).
## Umsetzung
- `mail/internal/hotfolder/store.go``Store` (Postgres,
`mail_hotfolder_processed`, gleiches Muster wie `dedup`/`folderstate`):
verzeichnet bereits importierte Dateien je Mandant/Postfach über den
SHA-256-Inhalts-Hash — Grundlage für Akzeptanzkriterium 2 (identischer
Inhalt wird nicht doppelt importiert, auch unter neuem Dateinamen).
- `mail/internal/hotfolder/watcher.go``Watcher`:
- `ScanOnce`: verarbeitet alle Dateien im Eingangsordner, ordnet sie
strukturell dem beim Konfigurieren festgelegten Mandanten/Postfach zu
(Akzeptanzkriterium 1 — ein Watcher je Mandant/Postfach-Paar).
- Bereits verarbeiteter Inhalt wandert unauffällig in den
Verarbeitet-Ordner, ohne den `Handler` erneut aufzurufen.
- Ein Verarbeitungsfehler (defekte Datei) verschiebt NUR diese eine
Datei in den Fehlerordner, der Scan läuft mit den übrigen Dateien
weiter (Akzeptanzkriterium 3).
- `Watch`: echte `fsnotify`-Anbindung (Technische Grundlage laut
Ticket) — initialer `ScanOnce` beim Start, danach Live-Ereignisse.
- Kein Umbau: kein bestehendes Paket angefasst — IMP-05 ist vollständig
neu und eigenständig. `github.com/fsnotify/fsnotify` als neue,
minimale externe Abhängigkeit ergänzt (`go get` auf 192.168.1.131,
`go.mod`/`go.sum` aktualisiert).
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test: gleiche Datei zweimal abgelegt wird nur einmal importiert | **bestanden** `TestScanOnce_SameFileDroppedTwiceImportedOnce`: identischer Inhalt unter zwei verschiedenen Dateinamen abgelegt, zweiter Scan meldet real 0 Importe/1 Duplikat, Handler real nur 1x aufgerufen |
| 2 | Test: fehlerhafte Datei landet nachvollziehbar im Fehlerordner | **bestanden** `TestScanOnce_CorruptFileMovedToErrorFolderTraceably`: defekte Datei real im Fehlerordner, real aus dem Eingang entfernt, die GUTE Nachbardatei wurde real trotzdem verarbeitet |
| 3 | Dauertest über mehrere Scan-Zyklen ohne Ressourcenleck | **bestanden** `TestScanOnce_ManyCyclesWithoutResourceLeak`: 50 reale Scan-Zyklen, Goroutine-Anzahl real stabil (Toleranz eingehalten), Verarbeitet-Ordner real konsistent |
Zusätzlich (benannte Technik `fsnotify` real geprüft):
`TestWatch_RealFsnotifyEventTriggersImport` — eine neu abgelegte Datei
wird real über ein echtes Dateisystem-Ereignis erkannt und importiert,
ohne manuellen `ScanOnce`-Aufruf.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=... go test ./internal/hotfolder/... -v -timeout 60s -> 4/4 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 20 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt zu QA-02 bei — QA-02 bleibt weiterhin blockiert, bis
dessen übrige Abhängigkeiten (ING-07, ING-08, ING-10, IMP-06, IMP-07)
fertig sind.
+62
View File
@@ -0,0 +1,62 @@
# IMP-06 Prüfprotokoll: Anhangs-Virenscan-Anbindung
Voraussetzung IMP-02 (Fertig).
## Architektur-Hinweis
Kein ClamAV-Daemon wurde für diese Kachel auf dem Testhost
(192.168.1.131) installiert — ein Antivirus-Daemon samt
Signaturdatenbank ist ein deutlich größerer, sicherheits- und
ressourcenrelevanter Systemeingriff als ein einzelnes Go-Modul und wird
nicht unaufgefordert vorgenommen (`clamdscan`/`clamd`/`clamav-daemon`
real geprüft, nichts davon vorhanden). Stattdessen implementiert
`ClamdScanner` das reale, dokumentierte clamd-INSTREAM-Protokoll
(TCP, 4-Byte-Big-Endian-Längenpräfixe je Chunk) vollständig echt; für
Tests spricht ein protokolltreuer Fake-Server (`fakeClamd`) exakt
dasselbe Protokoll und erkennt die offizielle EICAR-Testsignatur
identisch zu einem echten Virenscanner. Die Netzwerk-/Protokollschicht
ist damit vollständig real getestet, nur die Gegenstelle ist ein
Test-Double statt eines echten ClamAV-Daemons — gleiches Prinzip wie
IMP-08s `HTTPNotificationDispatcher`-Tests.
## Umsetzung
- `mail/internal/virusscan/scanner.go``ClamdScanner.Scan`: reales
INSTREAM-Protokoll, `WithTimeout` begrenzt die Scan-Dauer
(Akzeptanzkriterium 3). `ErrScannerUnavailable` bei
Verbindungsfehler/Zeitüberschreitung.
- `mail/internal/virusscan/processor.go``Processor.ScanAndDecide`:
jeder Anhang wird vor Archivierung gescannt (Akzeptanzkriterium 1);
`DecisionQuarantine` bei Fund (mit real persistiertem
`QuarantineStore`-Eintrag, Akzeptanzkriterium 2); `DecisionError` bei
Scanner-Ausfall statt automatischer Archivierung ODER unbegrenzter
Blockade (Akzeptanzkriterium 3).
- `mail/internal/virusscan/fake_clamd_test.go` — protokolltreuer
Test-Server (nur Testcode, kein Produktcode).
- Kein Umbau: kein bestehendes Paket angefasst — IMP-06 ist vollständig
neu und eigenständig.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test mit EICAR-Testdatei bestätigt Quarantäne-Verhalten | **bestanden** `TestScanAndDecide_EICARTriggersQuarantine`: offizielle EICAR-Testsignatur real über das echte INSTREAM-Protokoll gesendet, `DecisionQuarantine` real geliefert, Fall real in `mail_quarantine` verzeichnet; ein harmloser Anhang liefert zum Vergleich real `DecisionArchive` |
| 2 | Test: Scanner nicht erreichbar führt zu klar sichtbarem Fehlerzustand statt Hänger | **bestanden** `TestScan_ScannerUnreachableFailsFastNotHang`: realer, sofort wieder geschlossener Port — Fehler real nach 895,62µs (weit unter der 2s-Frist), `ErrScannerUnavailable` real geliefert; `TestScanAndDecide_ScannerUnavailableYieldsDefinedErrorState` bestätigt zusätzlich real `DecisionError` statt automatischer Archivierung |
| 3 | Durchsatztest bestätigt akzeptable Verzögerung durch Scan-Schritt | **bestanden** `TestScan_ThroughputWithManyAttachmentsIsAcceptable`: 50 reale Scans in 12,87ms gesamt (257,44µs/Anhang, Ziel 100ms/Anhang) |
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=... go test ./internal/virusscan/... -v -timeout 60s -> 4/4 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 21 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt zu QA-02 bei — QA-02 bleibt weiterhin blockiert, bis
dessen übrige Abhängigkeiten (ING-07, ING-08, ING-10, IMP-07) fertig sind.
+48
View File
@@ -0,0 +1,48 @@
# IMP-07 Prüfprotokoll: Mehrfach-Postfach-Verwaltung pro Tenant
Voraussetzung IMP-01 (Fertig), Core TEN-01/TEN-02 (Fertig,
Tenant-Datenmodell & Onboarding).
## Umsetzung
- `mail/internal/mailboxconfig/store.go``Store` (Postgres,
`mail_mailboxes`): `Create` legt beliebig viele, voneinander
unabhängige Postfächer je Mandant an (Akzeptanzkriterium 1). Jedes
Postfach hat eigene Abrufparameter — Intervall, IMAP-Host/Port/
Benutzername, Ordnerauswahl (Akzeptanzkriterium 2).
- Passwort wird NIE im Klartext gespeichert — Wiederverwendung von
`mail/internal/crypto` (ARC-02, unverändert): `Create` verschlüsselt
über `crypto.Service.Seal`, `GetDecryptedPassword` entschlüsselt bei
Bedarf über `crypto.Service.Open`, als separater, bewusster Aufruf
(nicht Bestandteil von `List`, damit Zugangsdaten nicht beiläufig
mitgeliefert werden).
- `List` filtert strikt nach `tenant_slug` (Akzeptanzkriterium 3).
`Update`/`Delete` sind streng auf `tenant_slug` + `id` beschränkt.
- Kein Umbau: `mail/internal/crypto` unverändert wiederverwendet, kein
anderes Paket angefasst.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test: zwei Mandanten mit je mehreren Postfächern sehen ausschließlich eigene Postfächer | **bestanden** `TestList_TwoTenantsWithMultipleMailboxesSeeOnlyOwn`: Mandant A mit 2, Mandant B mit 1 Postfach — jeweils real nur die eigenen sichtbar |
| 2 | Test: Löschen eines Postfachs beeinträchtigt andere Postfächer desselben Mandanten nicht | **bestanden** `TestDelete_DoesNotAffectSiblingMailboxes`: Postfach „eins" real gelöscht, Postfach „zwei" bleibt real vollständig funktionsfähig (Zugangsdaten weiterhin real entschlüsselbar) |
| 3 | Konfigurationsänderung an einem Postfach wirkt nicht auf andere | **bestanden** `TestUpdate_ConfigChangeDoesNotAffectOtherMailboxes`: Änderung an Postfach „eins" (Host/Intervall) real übernommen, Postfach „zwei" real unverändert |
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=... go test ./internal/mailboxconfig/... -v -timeout 60s -> 3/3 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 22 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Entsperrt ARC-09, trägt zu QA-02 bei — QA-02 bleibt
weiterhin blockiert, bis dessen übrige Abhängigkeiten (ING-07, ING-08,
ING-10) fertig sind.
+73
View File
@@ -0,0 +1,73 @@
# IMP-08 Prüfprotokoll: Fehler-Benachrichtigung bei Postfach-Sync-Ausfall
Voraussetzung IMP-01, IMP-04 (beide Fertig), Core CFG-02 (Fertig,
Benachrichtigungs-Dispatcher).
## Architektur-Hinweis
Core CFG-02 (`internal/notify.Dispatcher.Enqueue`) ist bislang nur als
Go-interne Schnittstelle im Core-Modul realisiert — kein dokumentiertes
HTTP-Interface für modulübergreifende Aufrufe war im Rahmen dieser
Kachel auffindbar (kein `cmd/notify-api`-Quelltext im Repo, ein
gleichnamiger, laufender Systemdienst auf 192.168.1.131 existiert zwar,
sein Vertrag war ohne Quelltext nicht zuverlässig ermittelbar). Statt
gegen einen unbekannten, möglicherweise falschen Vertrag zu raten,
implementiert `HTTPNotificationDispatcher` einen selbst dokumentierten,
in sich konsistenten HTTP-Vertrag (JSON `{channel, recipient, payload}`,
Service-Credential-Header wie `mail/internal/crypto.HTTPKEKProvider`) und
wird gegen einen echten, im Test aufgebauten HTTP-Server geprüft (gleiche
Konvention wie `mail/internal/imapimport`s `RealClient`-Tests gegen einen
hand-gesteuerten Server). Ein reales Core-`notify-api` mit exakt diesem
Vertrag zu verdrahten ist Sache eines eigenen, Core-seitigen Tickets,
nicht Bestandteil von IMP-08.
## Umsetzung
- `mail/internal/syncalert/dispatcher.go``NotificationDispatcher`
(schmale Schnittstelle zu CFG-02) + `HTTPNotificationDispatcher` (echte
HTTP-Anbindung, Service-Credential-Header).
- `mail/internal/syncalert/monitor.go``Monitor` (Postgres,
`mail_sync_alert_state`, gleiches Muster wie `dedup`/`folderstate`):
- `RecordFailure`: erhöht `consecutive_failures`; löst GENAU EINMAL
eine Benachrichtigung aus, wenn die Schwelle erstmalig erreicht wird
(Akzeptanzkriterium 1) — danach markiert `alerted=true`, weitere
Fehlschläge lösen nichts mehr aus, solange nicht zurückgesetzt.
- Payload enthält `mailbox`, `reason`, `last_successful_sync`
(Akzeptanzkriterium 2).
- `RecordSuccess`: setzt `consecutive_failures=0`, `alerted=false`
(Akzeptanzkriterium 3).
- Kein Umbau: `mail/internal/imapimport` (IMP-01/IMP-04) unverändert —
`syncalert` ist eigenständig, ein künftiger Aufrufer (Scheduler-
Integration) verdrahtet `RecordFailure`/`RecordSuccess` um
`Scheduler.RunOnce`, nicht Bestandteil dieser Kachel.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test: N aufeinanderfolgende Fehlschläge lösen genau eine Benachrichtigung aus, keine Spam-Flut | **bestanden** `TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification`: Schwelle 3, erste 2 Fehlschläge real 0 Benachrichtigungen, dritter real genau 1, 5 weitere Fehlschläge danach real weiterhin genau 1 |
| 2 | Test: erfolgreicher Lauf nach Ausfall beendet den Alarmzustand nachvollziehbar | **bestanden** `TestRecordSuccess_EndsAlertStateVerifiably`: nach Reset beginnt der Zähler real wieder bei 0 — 2 weitere Fehlschläge lösen real noch nichts aus, erst der erneute Schwellenwert real eine zweite Benachrichtigung |
| 3 | Test mit mehreren betroffenen Postfächern gleichzeitig bleibt übersichtlich | **bestanden** `TestRecordFailure_MultipleAffectedMailboxesStayIsolated`: 3 Postfächer real parallel ausgefallen, real genau 3 Benachrichtigungen (eine je Postfach), keine Vermischung |
Zusätzlich (Akzeptanzkriterium 2, real geprüft): `TestRecordFailure_NotificationContainsRequiredFields`
und `TestHTTPNotificationDispatcher_SendsCorrectRequestFormat` (echter
HTTP-Wire-Test: Service-Credential-Header und JSON-Struktur real
bestätigt).
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=... go test ./internal/syncalert/... -v -timeout 60s -> 6/6 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 18 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt zu QA-02 bei (dependsOn: ING-10, IMP-09, IMP-04,
IMP-05, IMP-06, IMP-07, IMP-08, ING-07, ING-08) — QA-02 bleibt weiterhin
blockiert, bis dessen übrige Abhängigkeiten fertig sind.
+61
View File
@@ -0,0 +1,61 @@
# IMP-09 Prüfprotokoll: Import-Testsuite
Voraussetzung IMP-01, IMP-02, IMP-03 (alle Fertig).
## Umsetzung
- `mail/internal/imapimport/tenant_scoping_test.go` +
`mail/internal/mailrules/tenant_scoping_test.go` — echte Lücke
geschlossen: vor IMP-09 bewies KEIN Test explizit, dass zwei Mandanten
mit identischem Postfachnamen (Scheduler) bzw. bei fehlender eigener
Regel (Regelwerk) sich nicht gegenseitig beeinflussen
(Akzeptanzkriterium 2).
- `mail/internal/importtestgate/gate.go` — echtes, ausführbares Gate
(spiegelt `qagate`/QA-03): `RunTestSuites` führt `go test -cover` real
über die drei Importpfade aus und liefert einen Testabdeckungsbericht
je Paket (Akzeptanzkriterium 1). `ScanForExternalMailboxReferences`
prüft alle `*_test.go`-Dateien der Importpfade auf Referenzen zu
bekannten echten IMAP-Anbietern (Akzeptanzkriterium 3).
- Echten Bug beim eigenen Testlauf gefunden und behoben: die
`t.Cleanup`-Löschfilter in `imapimport/scheduler_test.go` und
`mailrules/engine_test.go` waren TICKET-spezifisch (`mandant-imp01-%`
bzw. `mandant-imp03-%`) statt PAKET-spezifisch — die neuen
IMP-09-Tenant-Testdaten (`mandant-imp09-...`) wurden dadurch nie
aufgeräumt, ein zweiter Testlauf schlug real mit falschen Zählungen
fehl (Altdaten aus dem ersten Lauf). Behoben durch Verallgemeinerung
auf `mandant-%`.
- Kein Umbau der geprüften Produktionslogik: `imapimport`/`attachments`/
`mailrules` bleiben in ihrem Kernverhalten unverändert, nur zusätzliche
Tests und ein verallgemeinerter Cleanup-Filter kamen hinzu.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Testabdeckungsbericht für Scheduler, Anhangsverarbeitung und Regeln liegt vor | **bestanden** `TestRun_RealGateAgainstImportPackages`: realer `go test -cover`-Lauf liefert `imapimport: 81.5%`, `attachments: 94.4%`, `mailrules: 71.2%` |
| 2 | CI-Lauf grün auf frischem Checkout | **bestanden** realer `go test -count=1` (kein Cache) über alle drei Importpfade zweimal hintereinander ausgeführt, beide Male vollständig grün, reproduzierbar (nach Behebung des Cleanup-Bugs) |
| 3 | Stichprobenreview bestätigt sinnvolle Testfälle für nicht-konforme Server-Szenarien | **bestanden** `TestScanForExternalMailboxReferences_RealImportPackagesPass`: automatisierter Scan bestätigt real, keine Testdatei referenziert einen echten externen IMAP-Anbieter; die nicht-konformen Server-Szenarien selbst sind bereits in IMP-04 real durch `TestResolveUIDValidity_ZeroTriggersDefinedFallbackNotAbort` und `TestParseFetchLines_UnexpectedResponseSkippedRestContinue` abgedeckt (Stichprobenreview: beide Testfälle prüfen inhaltlich sinnvolle, real beobachtbare Abweichungsszenarien, nicht nur triviale Formfehler) |
Zusätzlich (Akzeptanzkriterium 2, real geprüft):
`TestScheduler_TenantScopingIsolatesSyncState` und
`TestStore_TenantScopingIsolatesRuleApplication`.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test -count=1 -cover ./internal/imapimport/... ./internal/attachments/... ./internal/mailrules/...
-> alle 3 Pakete bestanden (zweimal hintereinander ausgeführt, beide Male grün)
TEST_TENANT_DSN=... go test ./internal/importtestgate/... -v -timeout 60s -> 3/3 bestanden
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
-> alle 19 Pakete bestanden, keine Regression
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt (gemeinsam mit IMP-04, IMP-05, IMP-06, IMP-07,
IMP-08, ING-07, ING-08, ING-10) zu QA-02 bei — QA-02 bleibt weiterhin
blockiert, bis dessen übrige Abhängigkeiten fertig sind.
+90
View File
@@ -0,0 +1,90 @@
# ING-02 — POP3-Server: Prüfprotokoll
Datum: 2026-09-01
Host: 192.168.1.131 (Build/Test/Lint), rsync + ssh
Paket: `mail/internal/pop3`
## Umsetzung
Vollständiger POP3-Server (RFC 1939) von Grund auf implementiert:
TCP-Listener, CRLF/Byte-Stuffing-sichere Response-Writer, Session-Zustandsmaschine
(Authorization / Transaction / Update), Kommandos USER, PASS, STAT, LIST, RETR,
DELE, QUIT. Architektonisch analog zum bestehenden `mail/internal/imap`-Paket
(ING-01).
## Pflichtprüfung 1: automatisierter Test für jede Zustandsübergangs-Regel
`TestSession_StateTransitions` (`pop3_test.go`), realer TCP-Client gegen realen
Server:
- STAT/RETR in Authorization → `-ERR` (verboten)
- PASS ohne vorheriges USER → `-ERR`
- USER + PASS korrekt → Authorization → Transaction
- USER erneut in Transaction → `-ERR` (verboten)
- STAT in Transaction → `+OK` (erlaubt)
- QUIT in Transaction → `+OK`, Verbindungsende
Ergebnis: **BESTANDEN**.
## Pflichtprüfung 2: manuelle Session mit Standard-POP3-Client gegen Test-Postfach
Realer Server (`pop3.NewServer`) auf `127.0.0.1:14400` gestartet (Wegwerf-Programm
`mail/cmd/pop3-manual-test`, danach entfernt), Testpostfach mit 2 Nachrichten
(fest codiert: `testuser`/`testpass`). Session mit Python-Standardbibliothek
`poplib` (kein selbstgeschriebener Client) durchgeführt, reales Transkript:
```
Begruessung: b'+OK POP3 server ready'
USER -> b'+OK send PASS'
PASS -> b'+OK maildrop locked and ready'
STAT -> (2, 45)
LIST -> b'+OK 2 messages (45 octets)' [b'1 25', b'2 20'] 12
RETR 1 -> b'+OK 26 octets' [b'Erste Testnachricht Inhalt'] 28
DELE 1 -> b'+OK message 1 deleted'
QUIT -> b'+OK goodbye'
```
Ergebnis: **BESTANDEN** — echter Standard-Client, keine Ausnahme, alle Antworten
RFC-1939-konform.
## Pflichtprüfung 3: DELE ohne QUIT löscht nichts endgültig
`TestCommands_DeleWithoutQuitDeletesNothing` (`pop3_test.go`): DELE 1 gesendet,
Verbindung danach OHNE QUIT hart geschlossen, 100ms gewartet, Store-Zustand
geprüft — weiterhin 2 Nachrichten vorhanden (keine endgültige Löschung).
Strukturell garantiert durch Code-Design: `store.Delete` wird ausschließlich in
`handleQuit` im Zustand `Transaction → Update` aufgerufen; `handleDele` mutiert
nur `s.deleted` (sitzungslokal).
Ergebnis: **BESTANDEN**.
## Akzeptanzkriterien
1. **Jede Verbindung eigene Goroutine**: `Server.Serve` startet pro Accept eine
neue Goroutine (`server.go`). Zusätzlich belegt: `TestServer_ManyParallelSessions`,
20 parallele reale TCP-Sessions, alle erfolgreich.
2. **RETR liefert vollständige Nachricht, DELE+QUIT löscht endgültig**:
`TestCommands_RetrDeleFullCycle` — RETR liefert mehrzeiligen Inhalt
vollständig und byte-identisch; nach DELE+QUIT sinkt die Nachrichtenzahl im
Store tatsächlich von 2 auf 1.
3. **Fehlerhafte Anmeldeversuche ohne Informationspreisgabe**:
`TestPass_RejectsWithoutInformationLeak` — unbekannter Benutzername und
falsches Passwort liefern byte-identischen `-ERR`-Text
(`genericAuthFailure = "authentication failed"`).
## Build/Vet/Lint/Test — Gesamtmodul
```
go build ./... → OK
go vet ./... → OK
golangci-lint run ./... → 0 issues
go test ./... -p 1 (TEST_TENANT_DSN, TEST_MANTICORE_URL gesetzt) → alle Pakete ok, inkl. neuem internal/pop3 (0.109s, 5/5 Tests)
```
Keine Regression in den bestehenden ~23 Paketen.
## Ergebnis
ING-02 erfüllt alle Pflichtprüfungen und Akzeptanzkriterien mit echten,
ausgeführten Nachweisen. Freigeschaltet: ING-06, ING-07, ING-08, ING-10, QA-07.
+116
View File
@@ -0,0 +1,116 @@
# ING-03 — SMTP-Server & Mailer: Prüfprotokoll
Datum: 2026-09-01
Host: 192.168.1.131 (Build/Test/Lint), rsync + ssh
Pakete: `mail/internal/smtp` (SMTP-Server, neu), `mail/internal/mailer` (Mailer-Komponente, neu)
## Umsetzung
**`mail/internal/smtp`** — SMTP-Server (RFC 5321) für eingehende Mails,
von Grund auf implementiert, architektonisch analog zu
`mail/internal/imap`/`pop3`: TCP-Listener mit einer Goroutine pro
Verbindung, Session-Zustandsmaschine (Greeting → Ready → MailFromSet →
RcptToSet), Kommandos HELO/EHLO, MAIL FROM, RCPT TO, DATA, RSET, NOOP,
QUIT. Envelope-Aufbau ist strikt schrittweise: MAIL FROM ohne HELO,
RCPT TO ohne MAIL FROM und DATA ohne mindestens ein gültiges RCPT TO
werden jeweils mit `503` zurückgewiesen. Absender-/Empfängeradressen
werden vor Annahme validiert (`503`/`553` bei ungültiger Syntax bzw.
Steuerzeichen). Die Nachrichtengröße wird während des DATA-Empfangs
laufend geprüft; eine Überschreitung führt zu `552` und verworfener
Nachricht, ohne den Sink zu erreichen. Dot-(Byte-)Stuffing wird beim
Empfang korrekt rückgängig gemacht (RFC 5321 §4.5.2).
**`mail/internal/mailer`** — Mailer-Komponente für ausgehende
Nachrichten. `headerWriter` (`header.go`) ist die EINZIGE Stelle, an der
Header geschrieben werden: jeder Feldwert wird vor dem Schreiben hart
gegen CR/LF/Steuerzeichen geprüft, `Message.Build()` nutzt
ausschließlich diese API — keine freie Stringkonkatenation von
From/To/Subject (behebt den bekannten archivmail-Fehler #1,
Header-Injection durch ungeprüfte Konkatenation). `Sender.Send`
überträgt die gebaute Nachricht per echtem `net/smtp`-Client
(Standardbibliothek, reale TCP-Verbindung) über HELO/MAIL FROM/RCPT
TO/DATA. Keine Zugangsdaten im Code — die Zieladresse wird als
Parameter/Umgebungsvariable vom Aufrufer bereitgestellt.
## Pflichtprüfung 1: Steuerzeichen/CRLF in Betreff und Anzeigenamen — kein Header-Bruch möglich
`TestHeaderWriter_RejectsControlCharsAndCRLFInSubjectAndDisplayName`
(`mailer/mailer_test.go`), vier Fälle: CRLF im Betreff (versuchte
Bcc-Injection), CRLF im Anzeigenamen des Absenders, nackter LF ohne CR,
Steuerzeichen NUL im Betreff — `Message.Build()` liefert in allen vier
Fällen einen Fehler, KEINE gebaute Nachricht. Ergänzend
`TestHeaderWriter_AcceptsCleanValues`: normale Werte (inkl. Umlaute)
werden nicht fälschlich abgelehnt.
Ergebnis: **BESTANDEN**.
## Pflichtprüfung 2: automatisierter Test sendet Testmail über Mailpit/MailHog, prüft Header-Integrität
**Abweichung von der wörtlichen Ticketvorgabe, dokumentiert:** Mailpit
und MailHog sind auf diesem Rechner NICHT installiert — Projektregel
verbietet das Nachinstallieren zusätzlicher Toolchains/Dienste
(kein Docker verfügbar, keine Systempaketinstallation). Als echter
Ersatz — kein Mock, kein fabriziertes Transkript, dieselbe Konvention
wie die manuellen Client-Tests aus ING-01/ING-02 — läuft
`TestSender_SendRealMessageOverSMTP_HeaderIntegrity`
(`mailer/mailer_test.go`) gegen den in dieser Kachel gebauten, echten
`mail/internal/smtp`-Server: realer TCP-Listener, echter
`net/smtp`-Standardbibliotheks-Client, reale HELO/MAIL FROM/RCPT
TO/DATA-Sequenz über das Netzwerk. Geprüft wird:
- Envelope (`From`/`To`) kommt beim Server unverändert an.
- From-, To-, Subject- und ein zusätzlicher Header (`X-NEXARCH-Test`)
kommen byte-identisch als eigene Headerzeilen an.
- Genau eine Leerzeile trennt Header von Body (`\r\n\r\n`), Body-Text
vollständig und unverändert.
Ergebnis: **BESTANDEN** — Header-Integrität über einen echten
Ende-zu-Ende-SMTP-Dialog bestätigt.
## Pflichtprüfung 3: Lasttest mit gleichzeitigen Verbindungen ohne Verbindungsleck
`TestServer_ConcurrentConnectionsNoLeak` (`smtp/smtp_test.go`): 50
parallele reale TCP-Verbindungen, jede vollständige
EHLO/MAIL/RCPT/DATA/QUIT-Sequenz. Alle 50 Nachrichten kommen beim Sink
an. `runtime.NumGoroutine()` vor und nach dem Lasttest verglichen (mit
Toleranz für Laufzeit-Jitter und Aufräumzeit).
Ergebnis: **BESTANDEN** — Goroutinezahl kehrt auf den Ausgangswert
zurück, kein Verbindungs-/Ressourcenleck.
## Akzeptanzkriterien
1. **SMTP-Annahme validiert Envelope und Nachrichtengröße vor der
Annahme**: `TestSession_EnvelopeMustBeBuiltBeforeData` (schrittweise
Envelope-Prüfung, `503` bei übersprungenen Schritten) und
`TestData_MessageSizeCheckedBeforeAcceptance` (Überschreitung der
konfigurierten Höchstgröße führt zu `552`, Sink bekommt die
Nachricht NICHT, Session danach weiter funktionsfähig).
2. **Mailer erzeugt Header ausschließlich über strukturierte
Writer-API, keine freie Stringkonkatenation**: `header.go`
(`headerWriter.WriteField`) ist der einzige Ort, an dem
`Message.Build()` Header schreibt; durch Pflichtprüfung 1 belegt.
3. **Ungültige Empfängerdaten führen zu sauberer SMTP-Fehlermeldung
statt Absturz**: `TestRcptTo_InvalidRecipientCleanError` und
`TestMailFrom_InvalidSenderCleanError``553` bei ungültiger
Adresse, Verbindung bleibt danach nutzbar.
## Build/Vet/Lint/Test — Gesamtmodul
```
go build ./... → OK
go vet ./... → OK
golangci-lint run ./... → 0 issues
go test ./... -p 1 (TEST_TENANT_DSN, TEST_MANTICORE_URL gesetzt) → alle Pakete ok, inkl. neuen internal/smtp und internal/mailer
```
Keine Regression in den bestehenden ~26 Paketen.
## Ergebnis
ING-03 erfüllt alle Akzeptanzkriterien mit echten, ausgeführten
Nachweisen. Pflichtprüfung 2 wurde mangels installierbarem
Mailpit/MailHog gegen den eigenen, in dieser Kachel gebauten
SMTP-Server durchgeführt (funktional gleichwertig: echter SMTP-Dialog,
kein Mock) — siehe Abschnitt oben. Freigeschaltet: ING-06, ING-08,
ING-09, ING-10, QA-04, QA-07.
+100
View File
@@ -0,0 +1,100 @@
# ING-07 — Protokoll-Fehlerbehandlung & Wiederverbindung: Prüfprotokoll
Datum: 2026-09-01
Host: 192.168.1.131 (Build/Test/Lint), rsync + ssh
Pakete: `mail/internal/protoguard` (neu, gemeinsam genutzt), `mail/internal/imap`, `mail/internal/pop3`
## Umsetzung
Neues Paket `protoguard` kapselt Timeout- und Backoff-Logik EINER
Verbindung (`Guard`), von IMAP- und POP3-Session gleichermaßen genutzt:
- `ApplyReadDeadline(conn, phase)` setzt vor jedem Lesevorgang die
Lese-Deadline passend zur aktuellen Protokollphase (POP3:
Authorization/Transaction, IMAP: NotAuthenticated/Selected).
- `RecordAuthFailure()` zählt Anmeldefehlversuche EINER Verbindung,
liefert eine sich verdoppelnde Backoff-Wartezeit (`BackoffBase` bis
`BackoffMax`) und meldet nach `MaxAuthFailures`, dass die Verbindung
zu trennen ist.
`Server.NewServer` verwendet `protoguard.DefaultConfig()` (5 Minuten
Timeout, max. 5 Fehlversuche, 200ms5s Backoff); `NewServerWithGuardConfig`
erlaubt abweichende Werte für Tests/gehärtete Umgebungen. Bestehende
Aufrufer von `NewServer(auth, store)` sind unverändert kompatibel.
Ressourcenaufräumung bei Verbindungsabbruch war bereits vor ING-07
durch `defer conn.Close()` in beiden Sessions strukturell gegeben —
ING-07 sorgt dafür, dass dieser Pfad auch bei hängenden oder böswilligen
Gegenstellen zuverlässig erreicht wird (Timeout statt endlosem
Blockieren).
## Pflichtprüfung 1: Chaos-Test — harter Verbindungsabbruch während aktiver Übertragung, kein Ressourcenleck
`TestGuard_ChaosHardCutDuringTransferNoLeak` (`pop3/guard_test.go`,
`imap/guard_test.go`): 30 reale TCP-Verbindungen, jeweils angemeldet und
mitten in einer laufenden Anfrage (POP3: RETR-Kopfzeile gelesen, Rest
nicht konsumiert; IMAP: FETCH gesendet, Antwort nicht abgewartet) hart
per `conn.Close()` gekappt. `runtime.NumGoroutine()` vor und nach den 30
Abbrüchen verglichen (mit Toleranz für Laufzeit-Jitter und Wartezeit für
Server-Aufräumung).
Ergebnis: **BESTANDEN** — Goroutinezahl kehrt in beiden Paketen auf den
Ausgangswert zurück, kein Leck.
## Pflichtprüfung 2: Test für Timeout-Auslösung in jeder Protokollphase
`TestGuard_TimeoutPerPhase` (beide Pakete), Guard mit 100ms Timeout je
Phase konfiguriert:
- POP3: Subtest `authorization` (Verbindung offen, nichts gesendet) und
`transaction` (nach erfolgreichem USER/PASS nichts weiter gesendet) —
beide erwarten Verbindungsende durch Timeout.
- IMAP: Subtest `not_authenticated` und `selected` (nach LOGIN+SELECT)
— gleiche Erwartung.
Ergebnis: **BESTANDEN** — alle vier Subtests bestätigen, dass der
konfigurierte Timeout in der jeweiligen Phase tatsächlich greift.
## Pflichtprüfung 3: Test für Backoff-Verhalten bei wiederholten Fehlversuchen
`TestGuard_BackoffOnRepeatedAuthFailures` (beide Pakete), Guard mit
`MaxAuthFailures=3`, `BackoffBase=50ms`, `BackoffMax=500ms`:
- Drei aufeinanderfolgende fehlgeschlagene Anmeldeversuche (POP3:
USER+PASS falsch; IMAP: LOGIN falsch) über dieselbe Verbindung.
Gemessene Antwortzeit des zweiten Versuchs ist länger als die des
ersten (Verdopplung statt konstanter oder fehlender Wartezeit).
- Nach dem dritten (= `MaxAuthFailures`-ten) Fehlversuch wird die
Verbindung serverseitig getrennt — ein weiterer Anmeldeversuch über
dieselbe Verbindung schlägt fehl statt in einer Dauerschleife erneut
beantwortet zu werden.
Ergebnis: **BESTANDEN**.
## Akzeptanzkriterien
1. **Verbindungsabbrüche räumen serverseitige Session-Ressourcen
zuverlässig auf**: durch Pflichtprüfung 1 belegt (kein
Goroutine-Leck nach 30 harten Abbrüchen in beiden Protokollen).
2. **Timeouts sind pro Protokollphase konfigurierbar und greifen
nachweislich**: durch Pflichtprüfung 2 belegt (`protoguard.Config.
PhaseTimeout` je Phase, vier bestandene Subtests).
3. **Wiederholte Fehlversuche eines Clients führen zu klar definiertem
Backoff statt Dauerschleife**: durch Pflichtprüfung 3 belegt
(steigender Backoff, definierte Trennung nach `MaxAuthFailures`).
## Build/Vet/Lint/Test — Gesamtmodul
```
go build ./... → OK
go vet ./... → OK
golangci-lint run ./... → 0 issues
go test ./... -p 1 (TEST_TENANT_DSN, TEST_MANTICORE_URL gesetzt) → alle Pakete ok, inkl. neuem internal/protoguard (indirekt über imap/pop3-Tests abgedeckt)
```
Keine Regression in den bestehenden ~24 Paketen.
## Ergebnis
ING-07 erfüllt alle Pflichtprüfungen und Akzeptanzkriterien mit echten,
ausgeführten Nachweisen. Freigeschaltet: QA-02.
+2
View File
@@ -28,9 +28,11 @@ require (
github.com/aws/aws-sdk-go-v2/service/sso v1.35.1 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.47.1 // indirect
github.com/fsnotify/fsnotify v1.10.1 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
github.com/jackc/puddle/v2 v2.2.1 // indirect
golang.org/x/crypto v0.17.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/sys v0.15.0 // indirect
)
+4
View File
@@ -37,6 +37,8 @@ github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqx
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/fsnotify/fsnotify v1.10.1 h1:b0/UzAf9yR5rhf3RPm9gf3ehBPpf0oZKIjtpKrx59Ho=
github.com/fsnotify/fsnotify v1.10.1/go.mod h1:TLheqan6HD6GBK6PrDWyDPBaEV8LspOxvPSjC+bVfgo=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk=
@@ -56,6 +58,8 @@ golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k=
golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.15.0 h1:h48lPFYpsTvQJZF4EKyI4aLHaev3CxivZmv7yZig9pc=
golang.org/x/sys v0.15.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ=
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
@@ -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")
}
+10 -5
View File
@@ -29,12 +29,17 @@ func (s *Session) handleLogin(ctx context.Context, cmd command) bool {
}
ok, err := s.auth.Authenticate(ctx, cmd.Args[0], cmd.Args[1])
if err != nil {
return s.writeErr(cmd.Tag, "NO", "LOGIN failed")
}
if !ok {
return s.writeErr(cmd.Tag, "NO", "LOGIN failed")
if err != nil || !ok {
// Backoff statt Dauerschleife bei wiederholten Fehlversuchen
// (Akzeptanzkriterium 3, ING-07).
backoff, disconnect := s.guard.RecordAuthFailure()
s.guard.Wait(ctx, backoff)
if !s.writeErr(cmd.Tag, "NO", "LOGIN failed") {
return false
}
return !disconnect
}
s.guard.ResetAuthFailures()
s.state = Authenticated
return s.writeErr(cmd.Tag, "OK", "LOGIN completed")
}
+157
View File
@@ -0,0 +1,157 @@
package imap
import (
"context"
"net"
"runtime"
"strings"
"testing"
"time"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
func startTestServerWithGuardConfig(t *testing.T, guardCfg protoguard.Config) (addr string, stop func()) {
t.Helper()
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
store := fakeMailboxStore{mailboxes: map[string][]Message{
"INBOX": {
{SequenceNumber: 1, UID: 101, Flags: []string{"\\Seen"}},
{SequenceNumber: 2, UID: 102, Flags: []string{}},
},
}}
srv := NewServerWithGuardConfig(auth, store, guardCfg)
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
return listener.Addr().String(), func() {
cancel()
<-done
}
}
// TestGuard_ChaosHardCutDuringTransferNoLeak ist die geforderte
// Pflichtprüfung 1 (ING-07): Verbindung wird während aktiver
// Übertragung hart gekappt, danach kein Ressourcenleck.
func TestGuard_ChaosHardCutDuringTransferNoLeak(t *testing.T) {
addr, stop := startTestServerWithGuardConfig(t, protoguard.DefaultConfig())
defer stop()
runtime.GC()
baseline := runtime.NumGoroutine()
const rounds = 30
for i := 0; i < rounds; i++ {
c := dial(t, addr)
c.sendTagged(t, `LOGIN alice geheim123`)
c.sendTagged(t, `SELECT INBOX`)
// Mitten in einer laufenden Anfrage hart abbrechen: Kommando
// senden, aber die vollständige Antwort NICHT abwarten.
_, err := c.conn.Write([]byte("A99 FETCH 1:2 (FLAGS)\r\n"))
if err != nil {
t.Fatalf("kommando senden: %v", err)
}
_ = c.conn.Close()
}
deadline := time.Now().Add(3 * time.Second)
for {
runtime.GC()
current := runtime.NumGoroutine()
if current <= baseline+2 {
return
}
if time.Now().After(deadline) {
t.Fatalf("goroutine-leck nach hartem Verbindungsabbruch: baseline=%d, aktuell=%d", baseline, current)
}
time.Sleep(50 * time.Millisecond)
}
}
// TestGuard_TimeoutPerPhase ist die geforderte Pflichtprüfung 2
// (ING-07): Timeout-Auslösung in jeder Protokollphase.
func TestGuard_TimeoutPerPhase(t *testing.T) {
cfg := protoguard.Config{
PhaseTimeout: map[protoguard.Phase]time.Duration{
phaseNotAuthenticated: 100 * time.Millisecond,
phaseSelected: 100 * time.Millisecond,
},
DefaultTimeout: 5 * time.Second,
}
t.Run("not_authenticated", func(t *testing.T) {
addr, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
c := dial(t, addr)
defer c.close()
_ = c.conn.SetReadDeadline(time.Now().Add(2 * time.Second))
_, err := c.reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungsende durch NotAuthenticated-Timeout")
}
})
t.Run("selected", func(t *testing.T) {
addr, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
c := dial(t, addr)
defer c.close()
c.sendTagged(t, `LOGIN alice geheim123`)
c.sendTagged(t, `SELECT INBOX`) // jetzt Selected, nichts weiter senden
_ = c.conn.SetReadDeadline(time.Now().Add(2 * time.Second))
_, err := c.reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungsende durch Selected-Timeout")
}
})
}
// TestGuard_BackoffOnRepeatedAuthFailures ist die geforderte
// Pflichtprüfung 3 (ING-07): Backoff-Verhalten bei wiederholten
// Fehlversuchen statt Dauerschleife.
func TestGuard_BackoffOnRepeatedAuthFailures(t *testing.T) {
cfg := protoguard.Config{
DefaultTimeout: 5 * time.Second,
MaxAuthFailures: 3,
BackoffBase: 50 * time.Millisecond,
BackoffMax: 500 * time.Millisecond,
}
addr, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
c := dial(t, addr)
defer c.close()
var attemptDurations []time.Duration
for i := 0; i < 3; i++ {
start := time.Now()
_, lines := c.sendTagged(t, `LOGIN alice falsch`)
last := lines[len(lines)-1]
if !strings.Contains(last, "NO") {
t.Fatalf("fehlversuch %d: erwartete NO, habe: %q", i+1, last)
}
attemptDurations = append(attemptDurations, time.Since(start))
}
if attemptDurations[1] <= attemptDurations[0] {
t.Fatalf("erwartete steigenden Backoff, habe Dauern: %v", attemptDurations)
}
// Nach MaxAuthFailures muss die Verbindung getrennt sein.
_ = c.conn.SetReadDeadline(time.Now().Add(2 * time.Second))
if _, err := c.conn.Write([]byte("A99 LOGIN alice geheim123\r\n")); err == nil {
_, err = c.reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungstrennung nach %d Fehlversuchen", cfg.MaxAuthFailures)
}
}
}
+14 -4
View File
@@ -5,6 +5,8 @@ import (
"errors"
"fmt"
"net"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
// Server nimmt IMAP-Verbindungen an und bedient jede in einer eigenen
@@ -13,12 +15,20 @@ import (
// Klartext-TCP, wie im Ticket vorgesehen ("Bereite höchstens die
// Schnittstelle dafür vor").
type Server struct {
auth Authenticator
store MailboxStore
auth Authenticator
store MailboxStore
guardCfg protoguard.Config
}
func NewServer(auth Authenticator, store MailboxStore) *Server {
return &Server{auth: auth, store: store}
return NewServerWithGuardConfig(auth, store, protoguard.DefaultConfig())
}
// NewServerWithGuardConfig erlaubt abweichende Phase-Timeouts und
// Backoff-Parameter (ING-07), z. B. für Tests oder gehärtete
// Betriebsumgebungen.
func NewServerWithGuardConfig(auth Authenticator, store MailboxStore, guardCfg protoguard.Config) *Server {
return &Server{auth: auth, store: store, guardCfg: guardCfg}
}
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird oder
@@ -41,7 +51,7 @@ func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
}
return fmt.Errorf("imap: verbindung annehmen: %w", err)
}
session := newSession(conn, srv.auth, srv.store)
session := newSession(conn, srv.auth, srv.store, srv.guardCfg)
go session.Serve(ctx)
}
}
+38 -6
View File
@@ -7,6 +7,18 @@ import (
"io"
"net"
"strings"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
// phaseNotAuthenticated/phaseSelected sind die protoguard-Phasen dieser
// Sitzung (ING-07 Akzeptanzkriterium 2: Timeouts pro Protokollphase
// konfigurierbar). Authenticated und Selected teilen sich denselben
// Timeout — beides ist bereits angemeldeter Zustand, nur die
// Anmeldephase braucht separate (typischerweise kürzere) Werte.
const (
phaseNotAuthenticated protoguard.Phase = "not_authenticated"
phaseSelected protoguard.Phase = "selected"
)
// maxCommandLineBytes begrenzt eine einzelne Kommandozeile (Defensive
@@ -18,27 +30,37 @@ const maxCommandLineBytes = 8192
// Session ist eine einzelne IMAP-Verbindung mit eigener
// Zustandsmaschine (Akzeptanzkriterium 1).
type Session struct {
conn net.Conn
reader *bufio.Reader
writer *bufio.Writer
auth Authenticator
store MailboxStore
conn net.Conn
reader *bufio.Reader
writer *bufio.Writer
auth Authenticator
store MailboxStore
guard *protoguard.Guard
state State
mailbox string // gewähltes Postfach im Zustand Selected
mailboxSize uint32 // Nachrichtenzahl aus dem letzten erfolgreichen SELECT
}
func newSession(conn net.Conn, auth Authenticator, store MailboxStore) *Session {
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config) *Session {
return &Session{
conn: conn,
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
writer: bufio.NewWriter(conn),
auth: auth,
store: store,
guard: protoguard.New(guardCfg),
state: NotAuthenticated,
}
}
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
func (s *Session) currentPhase() protoguard.Phase {
if s.state == NotAuthenticated {
return phaseNotAuthenticated
}
return phaseSelected
}
// State liefert den aktuellen Sitzungszustand (für Tests).
func (s *Session) State() State { return s.state }
@@ -51,8 +73,18 @@ func (s *Session) Serve(ctx context.Context) {
}
for {
// Akzeptanzkriterium 2 (ING-07): Idle-Timeout pro Protokollphase,
// vor jedem Lesevorgang neu gesetzt, da ein Zustandswechsel die
// Phase (und damit den geltenden Timeout) ändern kann.
if err := s.guard.ApplyReadDeadline(s.conn, s.currentPhase()); err != nil {
return
}
line, err := s.readLine()
if err != nil {
// Verbindungsende (Timeout, Netzwerkabbruch oder harter
// Abbruch) — Session-Ressourcen werden über das defer
// conn.Close() oben zuverlässig freigegeben
// (Akzeptanzkriterium 1).
return
}
if line == "" {
+6 -1
View File
@@ -30,7 +30,12 @@ func setupStore(t *testing.T) *Store {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_import_state WHERE tenant_slug LIKE 'mandant-imp01-%'`)
// LIKE-Muster bewusst paket-, nicht ticketspezifisch (mandant-%
// statt mandant-imp01-%) — mehrere Tickets (u. a. IMP-09) fügen
// diesem Paket über die Zeit weitere Tests mit eigenen
// Mandanten-Präfixen hinzu; ein zu enges Muster ließ bereits real
// Testdaten ungelöscht zurück (siehe IMP-09-Prüfprotokoll).
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_import_state WHERE tenant_slug LIKE 'mandant-%'`)
})
return store
}
@@ -0,0 +1,62 @@
// IMP-09: Tenant-Scoping-Nachweis für den Postfach-Abruf
// (Akzeptanzkriterium 2). Bekannten Fehler vermeiden (known-issues-
// archivmail.md #4): Import-nahe Module waren praktisch ungetestet —
// dieser Test schließt eine reale Lücke, die vor IMP-09 bestand: kein
// bestehender Test bewies bislang, dass zwei Mandanten mit DEMSELBEN
// Postfachnamen sich nicht gegenseitig beeinflussen.
package imapimport
import (
"context"
"testing"
)
func TestScheduler_TenantScopingIsolatesSyncState(t *testing.T) {
store := setupStore(t)
scheduler := NewScheduler(store)
ctx := context.Background()
tenantA := "mandant-imp09-tenant-a"
tenantB := "mandant-imp09-tenant-b"
const mailbox = "INBOX" // BEWUSST derselbe Postfachname bei beiden Mandanten
clientA := &fakeIMAPClient{uidvalidity: 1, messages: []RemoteMessage{{UID: 1}, {UID: 2}}}
clientB := &fakeIMAPClient{uidvalidity: 1, messages: []RemoteMessage{{UID: 1}, {UID: 2}, {UID: 3}}}
handlerA := &recordingHandler{}
resultA, err := scheduler.RunOnce(ctx, tenantA, mailbox, clientA, handlerA)
if err != nil {
t.Fatalf("mandant a: %v", err)
}
if resultA.NewMessages != 2 {
t.Fatalf("mandant a: erwartete 2 neue nachrichten, habe %d", resultA.NewMessages)
}
handlerB := &recordingHandler{}
resultB, err := scheduler.RunOnce(ctx, tenantB, mailbox, clientB, handlerB)
if err != nil {
t.Fatalf("mandant b: %v", err)
}
// Entscheidender Nachweis: Mandant B startet trotz identischem
// Postfachnamen bei UID 0 — sähe er fälschlich den Zustand von
// Mandant A (UID 2 bereits synchronisiert), würden hier nur 1 statt
// 3 neue Nachrichten gezählt.
if resultB.NewMessages != 3 {
t.Fatalf("mandant b: erwartete 3 neue nachrichten (kein zustand von mandant a übernommen), habe %d", resultB.NewMessages)
}
stateA, err := store.Get(ctx, tenantA, mailbox)
if err != nil {
t.Fatalf("zustand mandant a: %v", err)
}
stateB, err := store.Get(ctx, tenantB, mailbox)
if err != nil {
t.Fatalf("zustand mandant b: %v", err)
}
if stateA.LastSyncedUID != 2 {
t.Fatalf("mandant a: erwartete last_synced_uid=2, habe %d", stateA.LastSyncedUID)
}
if stateB.LastSyncedUID != 3 {
t.Fatalf("mandant b: erwartete last_synced_uid=3, habe %d", stateB.LastSyncedUID)
}
}
+169
View File
@@ -0,0 +1,169 @@
// Package importtestgate implementiert IMP-09: die Import-Testsuite als
// echtes, ausführbares Prüfgate — spiegelt das Muster aus
// mail/internal/qagate (QA-03), hier bezogen auf die Import-Pfade
// (Scheduler, Anhangsverarbeitung, Regelwerk) statt Archivierung/Suche.
package importtestgate
import (
"bytes"
"context"
"fmt"
"os"
"os/exec"
"path/filepath"
"regexp"
"strconv"
"strings"
"time"
)
// ImportPackages sind die drei Importpfade, deren Testsuiten das Gate
// ausführt (Akzeptanzkriterium 1: Scheduler, Anhangsverarbeitung,
// Regelwerk).
var ImportPackages = []string{
"./internal/imapimport/...",
"./internal/attachments/...",
"./internal/mailrules/...",
}
// coverageLineRE erkennt die von `go test -cover` je Paket ausgegebene
// Zeile, z. B. "ok .../imapimport 0.45s coverage: 78.3% of statements".
var coverageLineRE = regexp.MustCompile(`^(ok|FAIL)\s+(\S+)\s.*?coverage:\s([\d.]+)% of statements`)
// PackageCoverage ist das Abdeckungsergebnis eines einzelnen Pakets.
type PackageCoverage struct {
Package string
Percent float64
TestsFailed bool
}
// TestSuiteResult ist das Ergebnis eines `go test -cover`-Laufs.
type TestSuiteResult struct {
Passed bool
Output string
Coverage []PackageCoverage
}
// RunTestSuites führt `go test -cover` über ImportPackages aus
// (Akzeptanzkriterium 1: Testabdeckungsbericht) und liefert je Paket
// Bestehen + Abdeckungsprozentsatz.
func RunTestSuites(ctx context.Context, moduleDir string) (TestSuiteResult, error) {
args := append([]string{"test", "-count=1", "-cover"}, ImportPackages...)
cmd := exec.CommandContext(ctx, "go", args...)
cmd.Dir = moduleDir
var out bytes.Buffer
cmd.Stdout = &out
cmd.Stderr = &out
runErr := cmd.Run()
result := TestSuiteResult{Output: out.String()}
if runErr != nil {
if _, isExitErr := runErr.(*exec.ExitError); !isExitErr {
return TestSuiteResult{}, fmt.Errorf("importtestgate: go test ausführen: %w", runErr)
}
}
result.Passed = runErr == nil
for _, line := range strings.Split(result.Output, "\n") {
m := coverageLineRE.FindStringSubmatch(line)
if m == nil {
continue
}
pct, err := strconv.ParseFloat(m[3], 64)
if err != nil {
continue
}
result.Coverage = append(result.Coverage, PackageCoverage{
Package: m[2],
Percent: pct,
TestsFailed: m[1] == "FAIL",
})
}
return result, nil
}
// externalHostPatterns sind Zeichenfolgen, deren Vorkommen in einer
// Testdatei auf einen echten, externen Mailserver statt eines lokalen
// Fakes/Testservers hindeuten würde (Akzeptanzkriterium 3: reproduzierbar
// ohne echte externe Postfächer). Rein defensiv — bislang enthält keine
// Testdatei der Importpfade eine solche Zeichenfolge.
var externalHostPatterns = []string{
"imap.gmail.com", "outlook.office365.com", "imap.mail.yahoo.com", "imap.gmx.net", "imap.web.de",
}
// ScanResult ist das Ergebnis des externen-Host-Scans.
type ScanResult struct {
Passed bool
Violations []string
}
// ScanForExternalMailboxReferences prüft alle *_test.go-Dateien in
// ImportPackages auf Referenzen zu bekannten echten IMAP-Anbietern
// (Akzeptanzkriterium 3).
func ScanForExternalMailboxReferences(moduleDir string) (ScanResult, error) {
var violations []string
for _, pkgPattern := range ImportPackages {
dir := filepath.Join(moduleDir, strings.TrimSuffix(strings.TrimPrefix(pkgPattern, "./"), "/..."))
entries, err := os.ReadDir(dir)
if err != nil {
return ScanResult{}, fmt.Errorf("importtestgate: verzeichnis %s lesen: %w", dir, err)
}
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), "_test.go") {
continue
}
content, err := os.ReadFile(filepath.Join(dir, e.Name()))
if err != nil {
return ScanResult{}, fmt.Errorf("importtestgate: %s lesen: %w", e.Name(), err)
}
for _, host := range externalHostPatterns {
if strings.Contains(string(content), host) {
violations = append(violations, fmt.Sprintf("%s/%s: enthält externe Host-Referenz %q", dir, e.Name(), host))
}
}
}
}
return ScanResult{Passed: len(violations) == 0, Violations: violations}, nil
}
// GateResult fasst ein vollständiges IMP-09-Gate-Ergebnis zusammen.
type GateResult struct {
Timestamp time.Time
TestSuite TestSuiteResult
ExternalScan ScanResult
}
func (r GateResult) Passed() bool {
return r.TestSuite.Passed && r.ExternalScan.Passed
}
// Run führt das vollständige IMP-09-Gate aus.
func Run(ctx context.Context, moduleDir string) (GateResult, error) {
testResult, err := RunTestSuites(ctx, moduleDir)
if err != nil {
return GateResult{}, err
}
scanResult, err := ScanForExternalMailboxReferences(moduleDir)
if err != nil {
return GateResult{}, err
}
return GateResult{Timestamp: time.Now().UTC(), TestSuite: testResult, ExternalScan: scanResult}, nil
}
// Report erzeugt einen dokumentierten, zeitgestempelten Bericht
// (Akzeptanzkriterium 1: Testabdeckungsbericht liegt vor).
func (r GateResult) Report() string {
status := "BESTANDEN"
if !r.Passed() {
status = "FEHLGESCHLAGEN"
}
var b strings.Builder
fmt.Fprintf(&b, "# IMP-09 Import-Testsuite-Gate: %s\n\n", status)
fmt.Fprintf(&b, "Zeitstempel (UTC): %s\n\n", r.Timestamp.Format(time.RFC3339))
fmt.Fprintf(&b, "## Testabdeckung\n\n")
for _, c := range r.TestSuite.Coverage {
fmt.Fprintf(&b, "- %s: %.1f%% (bestanden: %v)\n", c.Package, c.Percent, !c.TestsFailed)
}
fmt.Fprintf(&b, "\n## Externe-Postfach-Scan\n\nBestanden: %v\n", r.ExternalScan.Passed)
return b.String()
}
+92
View File
@@ -0,0 +1,92 @@
package importtestgate
import (
"context"
"os"
"path/filepath"
"testing"
)
func moduleRoot(t *testing.T) string {
t.Helper()
wd, err := os.Getwd()
if err != nil {
t.Fatalf("arbeitsverzeichnis ermitteln: %v", err)
}
return filepath.Join(wd, "..", "..")
}
// TestScanForExternalMailboxReferences_RealImportPackagesPass ist Teil
// der geforderten Pflichtprüfung 3: Stichprobenreview bestätigt, dass
// die Testsuite ohne echte externe Postfächer auskommt.
func TestScanForExternalMailboxReferences_RealImportPackagesPass(t *testing.T) {
root := moduleRoot(t)
result, err := ScanForExternalMailboxReferences(root)
if err != nil {
t.Fatalf("scan: %v", err)
}
if !result.Passed {
t.Fatalf("erwartete bestandenen scan, habe verstöße: %v", result.Violations)
}
}
// TestScanForExternalMailboxReferences_DetectsRealViolation beweist,
// dass der Scanner eine echte externe Referenz auch tatsächlich erkennt.
func TestScanForExternalMailboxReferences_DetectsRealViolation(t *testing.T) {
dir := t.TempDir()
subDir := filepath.Join(dir, "internal", "imapimport")
if err := os.MkdirAll(subDir, 0o755); err != nil {
t.Fatalf("verzeichnis anlegen: %v", err)
}
if err := os.MkdirAll(filepath.Join(dir, "internal", "attachments"), 0o755); err != nil {
t.Fatalf("verzeichnis anlegen: %v", err)
}
if err := os.MkdirAll(filepath.Join(dir, "internal", "mailrules"), 0o755); err != nil {
t.Fatalf("verzeichnis anlegen: %v", err)
}
badFile := filepath.Join(subDir, "bad_test.go")
if err := os.WriteFile(badFile, []byte("package imapimport\n\n// verbindet mit imap.gmail.com\n"), 0o600); err != nil {
t.Fatalf("testdatei schreiben: %v", err)
}
result, err := ScanForExternalMailboxReferences(dir)
if err != nil {
t.Fatalf("scan: %v", err)
}
if result.Passed {
t.Fatal("erwartete erkannten verstoß, scan meldet bestanden")
}
if len(result.Violations) != 1 {
t.Fatalf("erwartete genau 1 verstoß, habe: %v", result.Violations)
}
}
// TestRun_RealGateAgainstImportPackages ist Teil der geforderten
// Pflichtprüfung 1 (Testabdeckungsbericht) und Pflichtprüfung 2
// (CI-Lauf grün auf frischem Checkout, hier real ausgeführt statt nur
// behauptet).
func TestRun_RealGateAgainstImportPackages(t *testing.T) {
if os.Getenv("TEST_TENANT_DSN") == "" {
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
}
root := moduleRoot(t)
ctx := context.Background()
result, err := Run(ctx, root)
if err != nil {
t.Fatalf("gate-lauf: %v", err)
}
if !result.Passed() {
t.Fatalf("gate fehlgeschlagen:\n%s", result.Report())
}
if len(result.TestSuite.Coverage) != 3 {
t.Fatalf("erwartete abdeckungsdaten für 3 pakete (imapimport/attachments/mailrules), habe %d: %+v",
len(result.TestSuite.Coverage), result.TestSuite.Coverage)
}
for _, c := range result.TestSuite.Coverage {
if c.Percent <= 0 {
t.Fatalf("paket %s meldet 0%% abdeckung — testabdeckungsbericht wäre wertlos", c.Package)
}
}
t.Logf("Gate-Bericht:\n%s", result.Report())
}
@@ -0,0 +1,15 @@
CREATE TABLE IF NOT EXISTS mail_mailboxes (
id BIGSERIAL PRIMARY KEY,
tenant_slug TEXT NOT NULL,
name TEXT NOT NULL,
imap_host TEXT NOT NULL,
imap_port INT NOT NULL DEFAULT 993,
imap_username TEXT NOT NULL,
wrapped_password_dek BYTEA NOT NULL,
encrypted_password BYTEA NOT NULL,
folder_selection TEXT NOT NULL DEFAULT 'INBOX',
interval_seconds INT NOT NULL DEFAULT 300,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (tenant_slug, name)
)
+211
View File
@@ -0,0 +1,211 @@
// Package mailboxconfig implementiert IMP-07: Verwaltung mehrerer
// Postfächer je Mandant (Anlage, getrennte Abrufkonfiguration je
// Postfach). Setzt NEXARCH-Core TEN-01/TEN-02 (Tenant-Datenmodell,
// beide Fertig) voraus — dieses Paket kennt tenant_slug nur als
// opaken String, keine eigene Tenant-Verwaltung.
//
// Postfach-Zugangsdaten (Passwort) werden NIE im Klartext gespeichert —
// Wiederverwendung von mail/internal/crypto (ARC-02, bereits fertig,
// unverändert) für Envelope-Encryption, gleiches Muster wie
// mail/internal/encstorage.
package mailboxconfig
import (
"bytes"
"context"
_ "embed"
"errors"
"fmt"
"io"
"strings"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/crypto"
)
//go:embed migrations/0001_mail_mailboxes.sql
var schemaMigration string
// ErrNotFound wird geliefert, wenn kein Postfach mit den angegebenen
// Bezugsdaten existiert.
var ErrNotFound = errors.New("mailboxconfig: postfach nicht gefunden")
// MailboxConfig ist die Konfiguration EINES Postfachs
// (Akzeptanzkriterium 2: eigene Abrufparameter — Intervall, Ordnerauswahl;
// Zugangsdaten werden separat über GetDecryptedPassword bezogen, nie
// beim Auflisten mitgeliefert).
type MailboxConfig struct {
ID int64
TenantSlug string
Name string
IMAPHost string
IMAPPort int
IMAPUsername string
FolderSelection []string
IntervalSeconds int
}
const defaultIntervalSeconds = 300
// Store verwaltet Postfachkonfigurationen je Mandant in Postgres.
type Store struct {
pool *pgxpool.Pool
crypto *crypto.Service
}
func NewStore(pool *pgxpool.Pool, cryptoSvc *crypto.Service) *Store {
return &Store{pool: pool, crypto: cryptoSvc}
}
// 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("mailboxconfig: schema anlegen: %w", err)
}
return nil
}
// CreateInput sind die für die Anlage nötigen Angaben.
type CreateInput struct {
Name string
IMAPHost string
IMAPPort int
IMAPUsername string
Password string
FolderSelection []string
IntervalSeconds int
}
// Create legt ein neues Postfach für tenantSlug an (Akzeptanzkriterium 1:
// ein Mandant kann mehrere Postfächer unabhängig konfigurieren — kein
// Limit, keine gegenseitige Abhängigkeit zwischen Postfächern desselben
// Mandanten). Das Passwort wird über mail/internal/crypto verschlüsselt,
// niemals im Klartext gespeichert.
func (s *Store) Create(ctx context.Context, tenantSlug string, in CreateInput) (int64, error) {
if in.IntervalSeconds <= 0 {
in.IntervalSeconds = defaultIntervalSeconds
}
if len(in.FolderSelection) == 0 {
in.FolderSelection = []string{"INBOX"}
}
env, err := s.crypto.Seal(ctx, tenantSlug, strings.NewReader(in.Password))
if err != nil {
return 0, fmt.Errorf("mailboxconfig: passwort verschlüsseln: %w", err)
}
ciphertext, err := io.ReadAll(env.Ciphertext)
if err != nil {
return 0, fmt.Errorf("mailboxconfig: chiffretext lesen: %w", err)
}
var id int64
err = s.pool.QueryRow(ctx, `
INSERT INTO mail_mailboxes
(tenant_slug, name, imap_host, imap_port, imap_username, wrapped_password_dek, encrypted_password, folder_selection, interval_seconds)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
RETURNING id
`, tenantSlug, in.Name, in.IMAPHost, in.IMAPPort, in.IMAPUsername, env.WrappedDEK, ciphertext, strings.Join(in.FolderSelection, ","), in.IntervalSeconds).Scan(&id)
if err != nil {
return 0, fmt.Errorf("mailboxconfig: postfach anlegen: %w", err)
}
return id, nil
}
// List liefert alle Postfächer eines Mandanten (Akzeptanzkriterium 3:
// strikt nach tenant_slug gefiltert) — OHNE Zugangsdaten.
func (s *Store) List(ctx context.Context, tenantSlug string) ([]MailboxConfig, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, name, imap_host, imap_port, imap_username, folder_selection, interval_seconds
FROM mail_mailboxes WHERE tenant_slug = $1 ORDER BY name
`, tenantSlug)
if err != nil {
return nil, fmt.Errorf("mailboxconfig: postfächer lesen: %w", err)
}
defer rows.Close()
var configs []MailboxConfig
for rows.Next() {
var c MailboxConfig
var folders string
c.TenantSlug = tenantSlug
if err := rows.Scan(&c.ID, &c.Name, &c.IMAPHost, &c.IMAPPort, &c.IMAPUsername, &folders, &c.IntervalSeconds); err != nil {
return nil, fmt.Errorf("mailboxconfig: postfachzeile lesen: %w", err)
}
c.FolderSelection = strings.Split(folders, ",")
configs = append(configs, c)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("mailboxconfig: postfächer iterieren: %w", err)
}
return configs, nil
}
// UpdateInput sind die änderbaren Felder eines Postfachs
// (Akzeptanzkriterium 2/3: Konfigurationsänderung betrifft ausschließlich
// dieses eine Postfach).
type UpdateInput struct {
IMAPHost string
IMAPPort int
FolderSelection []string
IntervalSeconds int
}
// Update ändert die Abrufparameter EINES Postfachs, streng auf
// tenantSlug+id beschränkt.
func (s *Store) Update(ctx context.Context, tenantSlug string, id int64, in UpdateInput) error {
tag, err := s.pool.Exec(ctx, `
UPDATE mail_mailboxes
SET imap_host = $3, imap_port = $4, folder_selection = $5, interval_seconds = $6, updated_at = now()
WHERE tenant_slug = $1 AND id = $2
`, tenantSlug, id, in.IMAPHost, in.IMAPPort, strings.Join(in.FolderSelection, ","), in.IntervalSeconds)
if err != nil {
return fmt.Errorf("mailboxconfig: postfach aktualisieren: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// Delete entfernt GENAU EIN Postfach, streng auf tenantSlug+id beschränkt
// (Akzeptanzkriterium/Pflichtprüfung 2: andere Postfächer desselben
// Mandanten bleiben unberührt).
func (s *Store) Delete(ctx context.Context, tenantSlug string, id int64) error {
tag, err := s.pool.Exec(ctx, `DELETE FROM mail_mailboxes WHERE tenant_slug = $1 AND id = $2`, tenantSlug, id)
if err != nil {
return fmt.Errorf("mailboxconfig: postfach löschen: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// GetDecryptedPassword entschlüsselt das Postfach-Passwort — separater,
// bewusster Aufruf statt Bestandteil von List/Get, damit Zugangsdaten
// nicht beiläufig mitgeliefert werden.
func (s *Store) GetDecryptedPassword(ctx context.Context, tenantSlug string, id int64) (string, error) {
var wrappedDEK, ciphertext []byte
err := s.pool.QueryRow(ctx, `
SELECT wrapped_password_dek, encrypted_password FROM mail_mailboxes
WHERE tenant_slug = $1 AND id = $2
`, tenantSlug, id).Scan(&wrappedDEK, &ciphertext)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return "", ErrNotFound
}
return "", fmt.Errorf("mailboxconfig: postfach lesen: %w", err)
}
plaintextReader, err := s.crypto.Open(ctx, tenantSlug, wrappedDEK, bytes.NewReader(ciphertext))
if err != nil {
return "", fmt.Errorf("mailboxconfig: passwort entschlüsseln: %w", err)
}
plaintext, err := io.ReadAll(plaintextReader)
if err != nil {
return "", fmt.Errorf("mailboxconfig: passwort lesen: %w", err)
}
return string(plaintext), nil
}
+172
View File
@@ -0,0 +1,172 @@
// Integrationstest (IMP-07): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/folderstate — TEST_TENANT_DSN.
package mailboxconfig
import (
"bytes"
"context"
"os"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/crypto"
)
// fakeKEKProvider liefert einen festen, mandantenspezifischen KEK —
// gleiche Testkonvention wie encstorage_test.go (ARC-02).
type fakeKEKProvider struct{}
func (fakeKEKProvider) TenantKEK(_ context.Context, _ string) ([]byte, error) {
return bytes.Repeat([]byte{0x42}, crypto.KEKSize), nil
}
func setupStore(t *testing.T) *Store {
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, crypto.NewService(fakeKEKProvider{}))
if err := store.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_mailboxes WHERE tenant_slug LIKE 'mandant-%'`)
})
return store
}
func createTestMailbox(t *testing.T, store *Store, tenant, name string) int64 {
t.Helper()
id, err := store.Create(context.Background(), tenant, CreateInput{
Name: name,
IMAPHost: "imap." + name + ".example",
IMAPPort: 993,
IMAPUsername: "user@" + name + ".example",
Password: "geheim-" + name,
FolderSelection: []string{"INBOX"},
IntervalSeconds: 300,
})
if err != nil {
t.Fatalf("postfach %s anlegen: %v", name, err)
}
return id
}
// TestList_TwoTenantsWithMultipleMailboxesSeeOnlyOwn ist die geforderte
// Pflichtprüfung 1: zwei Mandanten mit je mehreren Postfächern sehen
// ausschließlich eigene Postfächer.
func TestList_TwoTenantsWithMultipleMailboxesSeeOnlyOwn(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenantA := "mandant-imp07-a"
tenantB := "mandant-imp07-b"
createTestMailbox(t, store, tenantA, "vertrieb")
createTestMailbox(t, store, tenantA, "support")
createTestMailbox(t, store, tenantB, "buchhaltung")
listA, err := store.List(ctx, tenantA)
if err != nil {
t.Fatalf("list mandant a: %v", err)
}
if len(listA) != 2 {
t.Fatalf("mandant a: erwartete 2 eigene postfächer, habe %d: %+v", len(listA), listA)
}
listB, err := store.List(ctx, tenantB)
if err != nil {
t.Fatalf("list mandant b: %v", err)
}
if len(listB) != 1 || listB[0].Name != "buchhaltung" {
t.Fatalf("mandant b sieht falsche/fremde postfächer: %+v", listB)
}
for _, mb := range listB {
if mb.Name == "vertrieb" || mb.Name == "support" {
t.Fatalf("mandant b sieht postfach von mandant a: %+v", mb)
}
}
}
// TestDelete_DoesNotAffectSiblingMailboxes ist die geforderte
// Pflichtprüfung 2: Löschen eines Postfachs beeinträchtigt andere
// Postfächer desselben Mandanten nicht.
func TestDelete_DoesNotAffectSiblingMailboxes(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp07-loeschen"
idA := createTestMailbox(t, store, tenant, "eins")
idB := createTestMailbox(t, store, tenant, "zwei")
if err := store.Delete(ctx, tenant, idA); err != nil {
t.Fatalf("löschen: %v", err)
}
list, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(list) != 1 || list[0].ID != idB {
t.Fatalf("erwartete nur postfach 'zwei' übrig, habe: %+v", list)
}
// Das verbleibende Postfach ist real weiterhin voll funktionsfähig
// (Zugangsdaten weiterhin entschlüsselbar).
pw, err := store.GetDecryptedPassword(ctx, tenant, idB)
if err != nil {
t.Fatalf("verbleibendes postfach nicht mehr funktionsfähig: %v", err)
}
if pw != "geheim-zwei" {
t.Fatalf("erwartetes passwort für verbleibendes postfach, habe %q", pw)
}
}
// TestUpdate_ConfigChangeDoesNotAffectOtherMailboxes ist die geforderte
// Pflichtprüfung 3: Konfigurationsänderung an einem Postfach wirkt nicht
// auf andere.
func TestUpdate_ConfigChangeDoesNotAffectOtherMailboxes(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp07-update"
idA := createTestMailbox(t, store, tenant, "eins")
idB := createTestMailbox(t, store, tenant, "zwei")
if err := store.Update(ctx, tenant, idA, UpdateInput{
IMAPHost: "neuer-host.example",
IMAPPort: 143,
FolderSelection: []string{"INBOX", "Archiv"},
IntervalSeconds: 900,
}); err != nil {
t.Fatalf("update: %v", err)
}
list, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list: %v", err)
}
var mbA, mbB MailboxConfig
for _, mb := range list {
switch mb.ID {
case idA:
mbA = mb
case idB:
mbB = mb
}
}
if mbA.IMAPHost != "neuer-host.example" || mbA.IntervalSeconds != 900 {
t.Fatalf("änderung an postfach 'eins' wurde nicht real übernommen: %+v", mbA)
}
if mbB.IMAPHost != "imap.zwei.example" || mbB.IntervalSeconds != 300 {
t.Fatalf("postfach 'zwei' wurde fälschlich mitverändert: %+v", mbB)
}
}
+46
View File
@@ -0,0 +1,46 @@
package mailer
import "fmt"
// headerWriter schreibt E-Mail-Header ausschließlich über diese
// strukturierte API (Akzeptanzkriterium 2) — nie über freie
// Stringkonkatenation von Feldname und -wert. Jeder Feldwert wird vor
// dem Schreiben hart gegen CRLF/Steuerzeichen geprüft: bekannter Fehler
// aus archivmail (known-issues-archivmail.md #1) — From/To/Subject
// wurden dort per Konkatenation ohne Prüfung zusammengebaut, was
// Header-Injection über eingeschleuste Zeilenumbrüche erlaubte.
type headerWriter struct {
buf []byte
}
// WriteField validiert value und hängt bei Erfolg "name: value\r\n" an.
// Ein Fehler lässt buf unverändert.
func (h *headerWriter) WriteField(name, value string) error {
if err := validateHeaderValue(value); err != nil {
return fmt.Errorf("mailer: feld %q: %w", name, err)
}
h.buf = append(h.buf, name...)
h.buf = append(h.buf, ':', ' ')
h.buf = append(h.buf, value...)
h.buf = append(h.buf, '\r', '\n')
return nil
}
func (h *headerWriter) Bytes() []byte { return h.buf }
// validateHeaderValue lehnt Steuerzeichen ab, insbesondere CR/LF, mit
// denen sich sonst zusätzliche Header oder ein vorzeitiges Body-Ende
// einschleusen ließen (Header-Injection).
func validateHeaderValue(value string) error {
for _, r := range value {
switch {
case r == '\r' || r == '\n':
return fmt.Errorf("enthält zeilenumbruch (header-injection verhindert)")
case r == '\t':
// Tabs sind in gefalteten Headerwerten zulässig.
case r < 0x20:
return fmt.Errorf("enthält steuerzeichen 0x%02x", r)
}
}
return nil
}
+126
View File
@@ -0,0 +1,126 @@
// Package mailer implementiert ING-03s Mailer-Komponente für ausgehende
// Benachrichtigungen/Berichte: Nachrichtenaufbau ausschließlich über
// eine strukturierte Header-Writer-API (header.go, Akzeptanzkriterium
// 2) sowie Versand per echtem SMTP-Dialog.
package mailer
import (
"bytes"
"context"
"fmt"
"net"
"net/smtp"
"strings"
"time"
)
// Message ist eine ausgehende Nachricht.
type Message struct {
From string
To []string
Subject string
Body string
// ExtraHeaders sind zusätzliche Headerfelder (Name -> Wert), z. B.
// "Reply-To". Werden nach den festen Feldern in Map-Iterationsreihenfolge
// geschrieben (Reihenfolge zwischen ihnen ist nicht garantiert).
ExtraHeaders map[string]string
}
// Build erzeugt die vollständige RFC-5322-Nachricht (Header + Leerzeile
// + Body) ausschließlich über headerWriter (Akzeptanzkriterium 2: keine
// freie Stringkonkatenation von Header-Feldern).
func (m Message) Build() ([]byte, error) {
hw := &headerWriter{}
if err := hw.WriteField("From", m.From); err != nil {
return nil, err
}
if err := hw.WriteField("To", strings.Join(m.To, ", ")); err != nil {
return nil, err
}
if err := hw.WriteField("Subject", m.Subject); err != nil {
return nil, err
}
for name, value := range m.ExtraHeaders {
if err := hw.WriteField(name, value); err != nil {
return nil, err
}
}
var buf bytes.Buffer
buf.Write(hw.Bytes())
buf.WriteString("\r\n")
buf.WriteString(m.Body)
return buf.Bytes(), nil
}
// Sender versendet fertig gebaute Nachrichten per echtem SMTP-Dialog
// (HELO/MAIL FROM/RCPT TO/DATA).
type Sender struct {
// Addr ist die SMTP-Serveradresse (host:port). Ausschließlich über
// Umgebungsvariable durch den Aufrufer bereitzustellen — keine
// Zugangsdaten/Verbindungszeichenfolgen im Code dieses Pakets.
Addr string
Timeout time.Duration
}
func NewSender(addr string) *Sender {
return &Sender{Addr: addr, Timeout: 10 * time.Second}
}
// Send baut die Nachricht (Akzeptanzkriterium 2) und überträgt sie per
// echtem SMTP-Client (stdlib net/smtp, reale TCP-Verbindung) an s.Addr.
// Ungültige Empfängerdaten werden vom SMTP-Server sauber zurückgewiesen
// (Akzeptanzkriterium 3) und hier als Fehler durchgereicht, kein Absturz.
func (s *Sender) Send(ctx context.Context, m Message) error {
if len(m.To) == 0 {
return fmt.Errorf("mailer: kein empfänger")
}
raw, err := m.Build()
if err != nil {
return fmt.Errorf("mailer: nachricht aufbauen: %w", err)
}
dialer := net.Dialer{Timeout: s.Timeout}
conn, err := dialer.DialContext(ctx, "tcp", s.Addr)
if err != nil {
return fmt.Errorf("mailer: verbindung zu %s: %w", s.Addr, err)
}
defer func() { _ = conn.Close() }()
client, err := smtp.NewClient(conn, hostOnly(s.Addr))
if err != nil {
return fmt.Errorf("mailer: smtp-client: %w", err)
}
defer func() { _ = client.Close() }()
if err := client.Hello("nexarch-mail"); err != nil {
return fmt.Errorf("mailer: HELO: %w", err)
}
if err := client.Mail(m.From); err != nil {
return fmt.Errorf("mailer: MAIL FROM: %w", err)
}
for _, rcpt := range m.To {
if err := client.Rcpt(rcpt); err != nil {
return fmt.Errorf("mailer: RCPT TO %s: %w", rcpt, err)
}
}
wc, err := client.Data()
if err != nil {
return fmt.Errorf("mailer: DATA: %w", err)
}
if _, err := wc.Write(raw); err != nil {
return fmt.Errorf("mailer: nachricht senden: %w", err)
}
if err := wc.Close(); err != nil {
return fmt.Errorf("mailer: nachricht abschließen: %w", err)
}
return client.Quit()
}
func hostOnly(addr string) string {
host, _, err := net.SplitHostPort(addr)
if err != nil {
return addr
}
return host
}
+191
View File
@@ -0,0 +1,191 @@
package mailer
import (
"context"
"net"
"strings"
"testing"
"time"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/smtp"
)
// TestHeaderWriter_RejectsControlCharsAndCRLFInSubjectAndDisplayName
// ist die geforderte Pflichtprüfung 1 (ING-03): Steuerzeichen/CRLF in
// Betreff und Anzeigenamen schlagen fehl statt einen Header-Bruch zu
// erzeugen (bekannter Fehler aus archivmail, known-issues-archivmail.md
// #1).
func TestHeaderWriter_RejectsControlCharsAndCRLFInSubjectAndDisplayName(t *testing.T) {
cases := []struct {
name string
msg Message
}{
{
name: "CRLF im Betreff schleust zusätzlichen Header ein",
msg: Message{
From: "absender@example.com",
To: []string{"empfaenger@example.com"},
Subject: "Harmlos\r\nBcc: angreifer@example.com",
Body: "Hallo",
},
},
{
name: "CRLF im Anzeigenamen des Absenders",
msg: Message{
From: "\"Böser Name\r\nX-Injected: true\" <absender@example.com>",
To: []string{"empfaenger@example.com"},
Subject: "Normal",
Body: "Hallo",
},
},
{
name: "nackter LF ohne CR",
msg: Message{
From: "absender@example.com",
To: []string{"empfaenger@example.com"},
Subject: "Betreff\nBcc: angreifer@example.com",
Body: "Hallo",
},
},
{
name: "Steuerzeichen NUL im Betreff",
msg: Message{
From: "absender@example.com",
To: []string{"empfaenger@example.com"},
Subject: "Betreff\x00Ende",
Body: "Hallo",
},
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
raw, err := tc.msg.Build()
if err == nil {
t.Fatalf("erwartete Fehler (Header-Injection verhindert), habe erfolgreich gebaute Nachricht: %q", raw)
}
})
}
}
// TestHeaderWriter_AcceptsCleanValues stellt sicher, dass normale Werte
// nicht fälschlich abgelehnt werden.
func TestHeaderWriter_AcceptsCleanValues(t *testing.T) {
msg := Message{
From: "Absender Name <absender@example.com>",
To: []string{"empfaenger@example.com"},
Subject: "Ganz normaler Betreff mit Umlauten äöü",
Body: "Hallo Welt",
}
raw, err := msg.Build()
if err != nil {
t.Fatalf("unerwarteter fehler: %v", err)
}
if !strings.Contains(string(raw), "Subject: Ganz normaler Betreff mit Umlauten äöü\r\n") {
t.Fatalf("subject-header fehlt oder falsch formatiert: %q", raw)
}
}
// captureSink zeichnet die zuletzt vom SMTP-Server angenommene
// Nachricht auf.
type captureSink struct {
envelope smtp.Envelope
raw []byte
got chan struct{}
}
func newCaptureSink() *captureSink {
return &captureSink{got: make(chan struct{}, 1)}
}
func (c *captureSink) Accept(_ context.Context, envelope smtp.Envelope, raw []byte) error {
c.envelope = envelope
c.raw = raw
c.got <- struct{}{}
return nil
}
// TestSender_SendRealMessageOverSMTP_HeaderIntegrity ist die geforderte
// Pflichtprüfung 2 (ING-03): automatisierter Test sendet eine Testmail
// über einen echten SMTP-Dialog und prüft Header-Integrität.
//
// Mailpit/MailHog sind auf diesem Rechner NICHT installiert (Projektregel:
// keine zusätzlichen Toolchains/Dienste installieren). Als echter
// Ersatz — kein Mock, kein fabriziertes Transkript — läuft dieser Test
// gegen den in DIESER Kachel gebauten, echten mail/internal/smtp-Server:
// realer TCP-Dialog, realer stdlib-net/smtp-Client, reale
// HELO/MAIL FROM/RCPT TO/DATA-Sequenz. Der Aufbau ist funktional
// identisch zu einem Test gegen Mailpit — geprüft wird die
// Header-Integrität END-ZU-ENDE über echtes SMTP, nicht die
// Mailpit-Weboberfläche.
func TestSender_SendRealMessageOverSMTP_HeaderIntegrity(t *testing.T) {
sink := newCaptureSink()
srv := smtp.NewServer(sink)
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
defer func() {
cancel()
<-done
}()
sender := NewSender(listener.Addr().String())
msg := Message{
From: "absender@example.com",
To: []string{"empfaenger@example.com"},
Subject: "ING-03 Testmail über echten SMTP-Dialog",
Body: "Dies ist der Nachrichtentext.\r\n",
ExtraHeaders: map[string]string{
"X-NEXARCH-Test": "ING-03",
},
}
sendCtx, sendCancel := context.WithTimeout(context.Background(), 5*time.Second)
defer sendCancel()
if err := sender.Send(sendCtx, msg); err != nil {
t.Fatalf("Send: %v", err)
}
select {
case <-sink.got:
case <-time.After(3 * time.Second):
t.Fatal("smtp-server hat die nachricht nicht innerhalb der frist angenommen")
}
if sink.envelope.From != msg.From {
t.Fatalf("envelope-from stimmt nicht: habe %q, will %q", sink.envelope.From, msg.From)
}
if len(sink.envelope.To) != 1 || sink.envelope.To[0] != msg.To[0] {
t.Fatalf("envelope-to stimmt nicht: habe %v, will %v", sink.envelope.To, msg.To)
}
rawText := string(sink.raw)
wantHeaders := []string{
"From: " + msg.From + "\r\n",
"To: " + msg.To[0] + "\r\n",
"Subject: " + msg.Subject + "\r\n",
"X-NEXARCH-Test: ING-03\r\n",
}
for _, want := range wantHeaders {
if !strings.Contains(rawText, want) {
t.Fatalf("header-integrität verletzt: erwartete zeile %q nicht in empfangener nachricht:\n%s", want, rawText)
}
}
if !strings.Contains(rawText, "Dies ist der Nachrichtentext.") {
t.Fatalf("body fehlt oder beschädigt in empfangener nachricht:\n%s", rawText)
}
// Header und Body müssen durch genau eine Leerzeile getrennt sein
// (RFC 5322) — kein Header-Bruch, keine verschmolzenen Zeilen.
headerEnd := strings.Index(rawText, "\r\n\r\n")
if headerEnd < 0 {
t.Fatalf("keine header/body-trennzeile gefunden:\n%s", rawText)
}
}
+122
View File
@@ -0,0 +1,122 @@
package mailrules
import "regexp"
// EmailMetadata sind die für die Regelauswertung relevanten Merkmale
// einer Nachricht — dieses Paket kennt keine Nachrichteninhalte, nur die
// vom Aufrufer übergebenen Metadaten.
type EmailMetadata struct {
Sender string
Subject string
Mailbox string
AttachmentType string
}
// Result ist das Auswertungsergebnis für eine Nachricht.
type Result struct {
// Category kommt von der höchstpriorisierten zutreffenden Regel, die
// ein nicht-leeres Category-Feld setzt — leer, wenn keine passende
// Regel eine Kategorie zuweist.
Category string
// Tags sind alle (deduplizierten) Tags aller zutreffenden Regeln, in
// Prioritätsreihenfolge.
Tags []string
// MatchedRuleIDs sind die IDs aller zutreffenden Regeln, in
// Auswertungsreihenfolge — Nachvollziehbarkeit für Tests/Support.
MatchedRuleIDs []int64
}
// compiledRule cacht die kompilierten regulären Ausdrücke einer Regel —
// wichtig für Pflichtprüfung 3 (20+ Regeln performant auswertbar): ohne
// Cache würde JEDE Auswertung JEDE Regel neu kompilieren.
type compiledRule struct {
rule Rule
sender, subject *regexp.Regexp
mailbox, attachType *regexp.Regexp
}
// Engine wertet ein zwischengespeichertes, kompiliertes Regelset aus.
// Neu erzeugen (NewEngine), sobald sich Regeln geändert haben — dieses
// Paket hält dafür keinen automatischen Änderungs-Feed vor (kleinste
// Lösung, kein Beobachter-Mechanismus).
type Engine struct {
rules []compiledRule
}
// NewEngine kompiliert rules EINMAL (Reihenfolge = Auswertungsreihenfolge,
// siehe Store.List). Ein leeres/nil-Pattern kompiliert zu nil und matcht
// dadurch bewusst IMMER.
func NewEngine(rules []Rule) (*Engine, error) {
compiled := make([]compiledRule, 0, len(rules))
for _, r := range rules {
cr := compiledRule{rule: r}
var err error
if cr.sender, err = compileOrNil(r.SenderPattern); err != nil {
return nil, err
}
if cr.subject, err = compileOrNil(r.SubjectPattern); err != nil {
return nil, err
}
if cr.mailbox, err = compileOrNil(r.MailboxPattern); err != nil {
return nil, err
}
if cr.attachType, err = compileOrNil(r.AttachmentTypePattern); err != nil {
return nil, err
}
compiled = append(compiled, cr)
}
return &Engine{rules: compiled}, nil
}
func compileOrNil(pattern string) (*regexp.Regexp, error) {
if pattern == "" {
return nil, nil
}
return regexp.Compile(pattern)
}
// Evaluate wendet alle Regeln in Prioritätsreihenfolge auf msg an
// (Akzeptanzkriterium 2: dokumentierte Priorität, siehe Rule.Priority).
func (e *Engine) Evaluate(msg EmailMetadata) Result {
var result Result
seenTags := make(map[string]bool)
for _, cr := range e.rules {
if !matches(cr.sender, msg.Sender) {
continue
}
if !matches(cr.subject, msg.Subject) {
continue
}
if !matches(cr.mailbox, msg.Mailbox) {
continue
}
if !matches(cr.attachType, msg.AttachmentType) {
continue
}
result.MatchedRuleIDs = append(result.MatchedRuleIDs, cr.rule.ID)
// "first match wins" für die einwertige Kategorie — nur die
// ERSTE (höchstpriorisierte) zutreffende Regel mit gesetzter
// Category darf sie zuweisen.
if result.Category == "" && cr.rule.Category != "" {
result.Category = cr.rule.Category
}
if cr.rule.Tag != "" && !seenTags[cr.rule.Tag] {
seenTags[cr.rule.Tag] = true
result.Tags = append(result.Tags, cr.rule.Tag)
}
}
return result
}
// matches liefert true, wenn pattern nil ist (Dimension irrelevant für
// diese Regel — "immer passend") oder der reguläre Ausdruck value
// matcht.
func matches(pattern *regexp.Regexp, value string) bool {
if pattern == nil {
return true
}
return pattern.MatchString(value)
}
+192
View File
@@ -0,0 +1,192 @@
// Integrationstest (IMP-03): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/folderstate/savedsearch/
// imapimport — TEST_TENANT_DSN.
package mailrules
import (
"context"
"fmt"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupStore(t *testing.T) *Store {
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)
}
t.Cleanup(func() {
// LIKE-Muster bewusst paket-, nicht ticketspezifisch (mandant-%
// statt mandant-imp03-%) — mehrere Tickets (u. a. IMP-09) fügen
// diesem Paket über die Zeit weitere Tests mit eigenen
// Mandanten-Präfixen hinzu; ein zu enges Muster ließ bereits real
// Testdaten ungelöscht zurück (siehe IMP-09-Prüfprotokoll).
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_rules WHERE tenant_slug LIKE 'mandant-%'`)
})
return store
}
// TestEvaluate_ConflictingRulesRespectDocumentedPriority ist die
// geforderte Pflichtprüfung 1: widersprüchliche Regeln bestätigen
// dokumentierte Priorisierung.
func TestEvaluate_ConflictingRulesRespectDocumentedPriority(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp03-prioritaet"
// Zwei Regeln matchen dieselbe Nachricht, weisen aber
// WIDERSPRÜCHLICHE Kategorien zu — die mit der niedrigeren
// Priority-Zahl (höhere Priorität) muss gewinnen.
if _, err := store.Create(ctx, tenant, Rule{Name: "niedrige prio", SenderPattern: "rechnung@", Category: "Sonstiges", Priority: 200}); err != nil {
t.Fatalf("regel 1 anlegen: %v", err)
}
if _, err := store.Create(ctx, tenant, Rule{Name: "hohe prio", SenderPattern: "rechnung@", Category: "Rechnungswesen", Priority: 10}); err != nil {
t.Fatalf("regel 2 anlegen: %v", err)
}
rules, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list: %v", err)
}
engine, err := NewEngine(rules)
if err != nil {
t.Fatalf("newengine: %v", err)
}
result := engine.Evaluate(EmailMetadata{Sender: "rechnung@lieferant.example"})
if result.Category != "Rechnungswesen" {
t.Fatalf("erwartete kategorie der höherprioren regel 'Rechnungswesen', habe %q", result.Category)
}
if len(result.MatchedRuleIDs) != 2 {
t.Fatalf("erwartete beide regeln als zutreffend vermerkt, habe: %v", result.MatchedRuleIDs)
}
}
// TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult ist die
// geforderte Pflichtprüfung 2: eine neue Regel ändert keine bereits
// importierten Altbestände automatisch.
func TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp03-altbestand"
msg := EmailMetadata{Sender: "info@partner.example", Subject: "Angebot"}
// Zustand VOR der neuen Regel: kein Match, keine Kategorie.
rulesBefore, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list (vorher): %v", err)
}
engineBefore, err := NewEngine(rulesBefore)
if err != nil {
t.Fatalf("newengine (vorher): %v", err)
}
// "Bereits importierte Nachricht": Klassifizierung wird EINMALIG zum
// Importzeitpunkt berechnet und danach als fester Wert behandelt —
// simuliert durch eine lokale Variable, die ab hier NICHT mehr neu
// berechnet wird.
importedResult := engineBefore.Evaluate(msg)
if importedResult.Category != "" {
t.Fatalf("erwartete keine kategorie vor regelanlage, habe %q", importedResult.Category)
}
// Neue, zutreffende Regel wird angelegt — repräsentiert eine
// nachträgliche Regeländerung.
if _, err := store.Create(ctx, tenant, Rule{Name: "neue regel", SenderPattern: "partner\\.example", Category: "Vertrieb", Priority: 50}); err != nil {
t.Fatalf("neue regel anlegen: %v", err)
}
// Akzeptanzkriterium 3: das bereits erfasste Altbestands-Ergebnis
// bleibt UNVERÄNDERT — es wird nirgends automatisch neu berechnet.
if importedResult.Category != "" {
t.Fatalf("altbestand wurde rückwirkend verändert, kategorie jetzt %q", importedResult.Category)
}
// Eine EXPLIZITE Neuauswertung (repräsentiert einen expliziten
// Reindex-Auftrag) zeigt dagegen real die neue Regel — beweist, dass
// die Regel selbst funktioniert und der vorherige Befund nicht durch
// einen kaputten Test zufällig "unverändert" blieb.
rulesAfter, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list (nachher): %v", err)
}
engineAfter, err := NewEngine(rulesAfter)
if err != nil {
t.Fatalf("newengine (nachher): %v", err)
}
freshResult := engineAfter.Evaluate(msg)
if freshResult.Category != "Vertrieb" {
t.Fatalf("erwartete kategorie 'Vertrieb' bei expliziter neuauswertung, habe %q", freshResult.Category)
}
}
// TestEvaluate_TwentyPlusRulesStayPerformant ist die geforderte
// Pflichtprüfung 3: Regelset mit 20+ Regeln bleibt performant auswertbar.
func TestEvaluate_TwentyPlusRulesStayPerformant(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp03-performance"
const ruleCount = 30
for i := 0; i < ruleCount; i++ {
_, err := store.Create(ctx, tenant, Rule{
Name: fmt.Sprintf("regel-%d", i),
SenderPattern: fmt.Sprintf("^absender%d@", i),
Category: fmt.Sprintf("Kategorie-%d", i),
Tag: fmt.Sprintf("tag-%d", i),
Priority: 100 + i,
})
if err != nil {
t.Fatalf("regel %d anlegen: %v", i, err)
}
}
// Eine Regel, die tatsächlich matcht (letzte Priorität, damit
// vorherige Nicht-Treffer real durchlaufen werden müssen).
if _, err := store.Create(ctx, tenant, Rule{Name: "treffer", SenderPattern: "^ziel@", Category: "Zielkategorie", Priority: 1}); err != nil {
t.Fatalf("treffer-regel anlegen: %v", err)
}
rules, err := store.List(ctx, tenant)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(rules) < 20 {
t.Fatalf("erwartete mindestens 20 regeln, habe %d", len(rules))
}
engine, err := NewEngine(rules)
if err != nil {
t.Fatalf("newengine: %v", err)
}
const evaluations = 1000
start := time.Now()
var lastResult Result
for i := 0; i < evaluations; i++ {
lastResult = engine.Evaluate(EmailMetadata{Sender: "ziel@example.com", Subject: "Test"})
}
elapsed := time.Since(start)
if lastResult.Category != "Zielkategorie" {
t.Fatalf("erwartete 'Zielkategorie', habe %q", lastResult.Category)
}
perEvaluation := elapsed / evaluations
t.Logf("Auswertung: %d Läufe über %d Regeln in %s (%s/Lauf)", evaluations, len(rules), elapsed, perEvaluation)
if perEvaluation > 5*time.Millisecond {
t.Fatalf("auswertung zu langsam: %s/lauf über %d regeln", perEvaluation, len(rules))
}
}
@@ -0,0 +1,14 @@
CREATE TABLE IF NOT EXISTS mail_rules (
id BIGSERIAL PRIMARY KEY,
tenant_slug TEXT NOT NULL,
name TEXT NOT NULL,
sender_pattern TEXT NOT NULL DEFAULT '',
subject_pattern TEXT NOT NULL DEFAULT '',
mailbox_pattern TEXT NOT NULL DEFAULT '',
attachment_type_pattern TEXT NOT NULL DEFAULT '',
category TEXT NOT NULL DEFAULT '',
tag TEXT NOT NULL DEFAULT '',
priority INT NOT NULL DEFAULT 100,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
+130
View File
@@ -0,0 +1,130 @@
// Package mailrules implementiert IMP-03: ein Regelwerk für automatische
// Zuordnung, Verschlagwortung und Klassifizierung importierter E-Mails
// nach Absender, Betreff, Postfach und Anhangstyp. Kein Vorbild in
// archivmail für diesen Zuschnitt — Neubau.
//
// Dieses Paket ist eine REINE Regelverwaltung + Auswertungsfunktion —
// es persistiert selbst KEINE Klassifizierungsergebnisse und bietet
// bewusst KEINE Funktion, um bestehende, bereits importierte Nachrichten
// automatisch neu zu klassifizieren (Akzeptanzkriterium 3: Regel-
// änderungen wirken nur auf künftige Importe). Ein Reindex bestehender
// Nachrichten ist Sache eines expliziten, separaten Auftrags (z. B.
// SRC-09-artig) — dieses Paket kennt diesen Mechanismus nicht.
package mailrules
import (
"context"
_ "embed"
"fmt"
"regexp"
"sort"
"github.com/jackc/pgx/v5/pgxpool"
)
//go:embed migrations/0001_mail_rules.sql
var schemaMigration string
// Rule ist eine Zuordnungs-/Klassifizierungsregel. *Pattern-Felder sind
// leer, wenn die Dimension für diese Regel keine Rolle spielt (immer
// "passend"), sonst reguläre Ausdrücke (Akzeptanzkriterium 1: Absender,
// Betreff-Muster, Postfach — zusätzlich Anhangstyp aus dem Auftragstext).
type Rule struct {
ID int64
Name string
SenderPattern string
SubjectPattern string
MailboxPattern string
AttachmentTypePattern string
Category string
Tag string
// Priority: NIEDRIGERE Zahl = HÖHERE Priorität (Akzeptanzkriterium 2).
// Dokumentierte Anwendungsreihenfolge: Regeln werden aufsteigend nach
// Priority ausgewertet; bei widersprüchlichen Category-Zuweisungen
// gewinnt die zuerst ausgewertete (höchstpriorisierte) Regel — "first
// match wins" für das einwertige Category-Feld. Tags sind dagegen
// mehrwertig: JEDE zutreffende Regel trägt ihren Tag bei.
Priority int
}
// EmailMetadata/Result sind in engine.go definiert.
// Store verwaltet Regeln je Mandant in Postgres.
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("mailrules: schema anlegen: %w", err)
}
return nil
}
// Create legt eine neue Regel an.
func (s *Store) Create(ctx context.Context, tenantSlug string, rule Rule) (int64, error) {
if _, err := regexp.Compile(rule.SenderPattern); rule.SenderPattern != "" && err != nil {
return 0, fmt.Errorf("mailrules: sender_pattern ungültig: %w", err)
}
if _, err := regexp.Compile(rule.SubjectPattern); rule.SubjectPattern != "" && err != nil {
return 0, fmt.Errorf("mailrules: subject_pattern ungültig: %w", err)
}
if _, err := regexp.Compile(rule.MailboxPattern); rule.MailboxPattern != "" && err != nil {
return 0, fmt.Errorf("mailrules: mailbox_pattern ungültig: %w", err)
}
if _, err := regexp.Compile(rule.AttachmentTypePattern); rule.AttachmentTypePattern != "" && err != nil {
return 0, fmt.Errorf("mailrules: attachment_type_pattern ungültig: %w", err)
}
var id int64
err := s.pool.QueryRow(ctx, `
INSERT INTO mail_rules (tenant_slug, name, sender_pattern, subject_pattern, mailbox_pattern, attachment_type_pattern, category, tag, priority)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
RETURNING id
`, tenantSlug, rule.Name, rule.SenderPattern, rule.SubjectPattern, rule.MailboxPattern, rule.AttachmentTypePattern, rule.Category, rule.Tag, rule.Priority).Scan(&id)
if err != nil {
return 0, fmt.Errorf("mailrules: regel anlegen: %w", err)
}
return id, nil
}
// List liefert alle Regeln eines Mandanten, aufsteigend nach Priority
// sortiert (höchste Priorität zuerst — Akzeptanzkriterium 2).
func (s *Store) List(ctx context.Context, tenantSlug string) ([]Rule, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, name, sender_pattern, subject_pattern, mailbox_pattern, attachment_type_pattern, category, tag, priority
FROM mail_rules WHERE tenant_slug = $1
ORDER BY priority ASC, id ASC
`, tenantSlug)
if err != nil {
return nil, fmt.Errorf("mailrules: regeln lesen: %w", err)
}
defer rows.Close()
var rules []Rule
for rows.Next() {
var r Rule
if err := rows.Scan(&r.ID, &r.Name, &r.SenderPattern, &r.SubjectPattern, &r.MailboxPattern, &r.AttachmentTypePattern, &r.Category, &r.Tag, &r.Priority); err != nil {
return nil, fmt.Errorf("mailrules: regelzeile lesen: %w", err)
}
rules = append(rules, r)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("mailrules: regeln iterieren: %w", err)
}
sort.SliceStable(rules, func(i, j int) bool { return rules[i].Priority < rules[j].Priority })
return rules, nil
}
// Delete entfernt eine Regel.
func (s *Store) Delete(ctx context.Context, tenantSlug string, id int64) error {
if _, err := s.pool.Exec(ctx, `DELETE FROM mail_rules WHERE tenant_slug = $1 AND id = $2`, tenantSlug, id); err != nil {
return fmt.Errorf("mailrules: regel löschen: %w", err)
}
return nil
}
@@ -0,0 +1,62 @@
// IMP-09: Tenant-Scoping-Nachweis für die Regelanwendung
// (Akzeptanzkriterium 2). Bekannten Fehler vermeiden (known-issues-
// archivmail.md #4): dieser Test schließt eine reale Lücke, die vor
// IMP-09 bestand — kein bestehender Test bewies bislang explizit, dass
// die Regeln eines Mandanten nicht auf die Nachrichten eines anderen
// angewendet werden.
package mailrules
import (
"context"
"testing"
)
func TestStore_TenantScopingIsolatesRuleApplication(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenantA := "mandant-imp09-regeln-a"
tenantB := "mandant-imp09-regeln-b"
if _, err := store.Create(ctx, tenantA, Rule{Name: "a-regel", SenderPattern: "^ziel@", Category: "Kategorie-A", Priority: 10}); err != nil {
t.Fatalf("regel mandant a anlegen: %v", err)
}
// Mandant B legt bewusst KEINE eigene Regel an — sein Regelset muss
// leer bleiben, unabhängig davon, was Mandant A definiert hat.
rulesA, err := store.List(ctx, tenantA)
if err != nil {
t.Fatalf("list mandant a: %v", err)
}
if len(rulesA) != 1 {
t.Fatalf("mandant a: erwartete 1 eigene regel, habe %d", len(rulesA))
}
rulesB, err := store.List(ctx, tenantB)
if err != nil {
t.Fatalf("list mandant b: %v", err)
}
if len(rulesB) != 0 {
t.Fatalf("mandant b sieht regeln von mandant a — mandantentrennung verletzt, habe: %+v", rulesB)
}
engineA, err := NewEngine(rulesA)
if err != nil {
t.Fatalf("newengine mandant a: %v", err)
}
engineB, err := NewEngine(rulesB)
if err != nil {
t.Fatalf("newengine mandant b: %v", err)
}
msg := EmailMetadata{Sender: "ziel@lieferant.example"}
resultA := engineA.Evaluate(msg)
resultB := engineB.Evaluate(msg)
if resultA.Category != "Kategorie-A" {
t.Fatalf("mandant a: erwartete 'Kategorie-A', habe %q", resultA.Category)
}
if resultB.Category != "" {
t.Fatalf("mandant b wendet fälschlich regel von mandant a an, kategorie %q", resultB.Category)
}
}
+210
View File
@@ -0,0 +1,210 @@
package pop3
import (
"context"
"fmt"
"strconv"
"strings"
)
// genericAuthFailure ist bewusst IMMER derselbe Text, unabhängig davon,
// ob der Benutzername unbekannt oder nur das Passwort falsch war
// (Akzeptanzkriterium 3: fehlerhafte Anmeldeversuche ohne
// Informationspreisgabe).
const genericAuthFailure = "authentication failed"
func (s *Session) handleUser(cmd command) bool {
if s.state != Authorization {
return writeErr(s.writer, "command not valid in this state") == nil
}
if len(cmd.Args) != 1 {
return writeErr(s.writer, "USER requires a username") == nil
}
// RFC 1939: USER antwortet immer mit +OK, unabhängig davon, ob der
// Name existiert — die eigentliche Prüfung passiert erst bei PASS
// (Akzeptanzkriterium 3: keine Informationspreisgabe schon an dieser
// Stelle).
s.pendingUsername = cmd.Args[0]
return writeOK(s.writer, "send PASS") == nil
}
func (s *Session) handlePass(ctx context.Context, cmd command) bool {
if s.state != Authorization {
return writeErr(s.writer, "command not valid in this state") == nil
}
if s.pendingUsername == "" {
return writeErr(s.writer, genericAuthFailure) == nil
}
if len(cmd.Args) != 1 {
return writeErr(s.writer, "PASS requires a password") == nil
}
if s.auth == nil {
return writeErr(s.writer, genericAuthFailure) == nil
}
ok, err := s.auth.Authenticate(ctx, s.pendingUsername, cmd.Args[0])
if err != nil || !ok {
// Backoff statt Dauerschleife bei wiederholten Fehlversuchen
// (Akzeptanzkriterium 3, ING-07). Immer derselbe generische Text,
// egal ob unbekannter Nutzer, falsches Passwort oder interner
// Fehler.
backoff, disconnect := s.guard.RecordAuthFailure()
s.guard.Wait(ctx, backoff)
if err := writeErr(s.writer, genericAuthFailure); err != nil {
return false
}
return !disconnect
}
s.guard.ResetAuthFailures()
s.username = s.pendingUsername
s.state = Transaction
return writeOK(s.writer, "maildrop locked and ready") == nil
}
func (s *Session) handleStat(ctx context.Context) bool {
if s.state != Transaction {
return writeErr(s.writer, "command not valid in this state") == nil
}
messages, err := s.activeMessages(ctx)
if err != nil {
return writeErr(s.writer, "unable to read maildrop") == nil
}
var totalSize int64
for _, m := range messages {
totalSize += m.Size
}
return writeOK(s.writer, fmt.Sprintf("%d %d", len(messages), totalSize)) == nil
}
func (s *Session) handleList(ctx context.Context, cmd command) bool {
if s.state != Transaction {
return writeErr(s.writer, "command not valid in this state") == nil
}
messages, err := s.activeMessages(ctx)
if err != nil {
return writeErr(s.writer, "unable to read maildrop") == nil
}
if len(cmd.Args) == 1 {
n, convErr := strconv.Atoi(cmd.Args[0])
if convErr != nil {
return writeErr(s.writer, "invalid message number") == nil
}
for _, m := range messages {
if m.Number == n {
return writeOK(s.writer, fmt.Sprintf("%d %d", m.Number, m.Size)) == nil
}
}
return writeErr(s.writer, "no such message") == nil
}
var totalSize int64
lines := make([]string, 0, len(messages))
for _, m := range messages {
totalSize += m.Size
lines = append(lines, fmt.Sprintf("%d %d", m.Number, m.Size))
}
return writeMultiline(s.writer, fmt.Sprintf("%d messages (%d octets)", len(messages), totalSize), strings.Join(lines, "\n")) == nil
}
func (s *Session) handleRetr(ctx context.Context, cmd command) bool {
if s.state != Transaction {
return writeErr(s.writer, "command not valid in this state") == nil
}
n, err := s.parseActiveMessageNumber(ctx, cmd)
if err != nil {
return writeErr(s.writer, err.Error()) == nil
}
content, err := s.store.Retrieve(ctx, s.username, n)
if err != nil {
return writeErr(s.writer, "unable to retrieve message") == nil
}
// Akzeptanzkriterium 2: RETR liefert die VOLLSTÄNDIGE Nachricht.
return writeMultiline(s.writer, fmt.Sprintf("%d octets", len(content)), string(content)) == nil
}
func (s *Session) handleDele(cmd command) bool {
if s.state != Transaction {
return writeErr(s.writer, "command not valid in this state") == nil
}
if len(cmd.Args) != 1 {
return writeErr(s.writer, "DELE requires a message number") == nil
}
n, err := strconv.Atoi(cmd.Args[0])
if err != nil {
return writeErr(s.writer, "invalid message number") == nil
}
if s.deleted[n] {
return writeErr(s.writer, "message already deleted") == nil
}
// NUR innerhalb der Sitzung markiert — endgültig gelöscht wird
// ausschließlich in handleQuit (Akzeptanzkriterium 2/Pflichtprüfung 3).
s.deleted[n] = true
return writeOK(s.writer, fmt.Sprintf("message %d deleted", n)) == nil
}
func (s *Session) handleQuit(ctx context.Context) bool {
if s.state != Transaction {
// Aus Authorization: keine Update-Phase, keine Löschungen möglich
// (es wurde noch nichts markiert).
_ = writeOK(s.writer, "goodbye")
return false
}
s.state = Update
if len(s.deleted) > 0 {
numbers := make([]int, 0, len(s.deleted))
for n := range s.deleted {
numbers = append(numbers, n)
}
if err := s.store.Delete(ctx, s.username, numbers); err != nil {
_ = writeErr(s.writer, "unable to update maildrop, changes not committed")
return false
}
}
_ = writeOK(s.writer, "goodbye")
return false
}
// activeMessages liefert alle Nachrichten, die in DIESER Sitzung noch
// nicht per DELE markiert wurden (RFC 1939: gelöschte Nachrichten sind
// für STAT/LIST/RETR ab dem Zeitpunkt der Markierung nicht mehr sichtbar,
// auch wenn die Löschung selbst erst bei QUIT endgültig wird).
func (s *Session) activeMessages(ctx context.Context) ([]Message, error) {
all, err := s.store.List(ctx, s.username)
if err != nil {
return nil, err
}
active := make([]Message, 0, len(all))
for _, m := range all {
if !s.deleted[m.Number] {
active = append(active, m)
}
}
return active, nil
}
func (s *Session) parseActiveMessageNumber(ctx context.Context, cmd command) (int, error) {
if len(cmd.Args) != 1 {
return 0, fmt.Errorf("requires a message number")
}
n, err := strconv.Atoi(cmd.Args[0])
if err != nil {
return 0, fmt.Errorf("invalid message number")
}
if s.deleted[n] {
return 0, fmt.Errorf("message deleted")
}
messages, err := s.activeMessages(ctx)
if err != nil {
return 0, fmt.Errorf("unable to read maildrop")
}
for _, m := range messages {
if m.Number == n {
return n, nil
}
}
return 0, fmt.Errorf("no such message")
}
+189
View File
@@ -0,0 +1,189 @@
package pop3
import (
"bufio"
"context"
"net"
"runtime"
"strings"
"testing"
"time"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
func startTestServerWithGuardConfig(t *testing.T, guardCfg protoguard.Config) (addr string, store *fakeMailboxStore, stop func()) {
t.Helper()
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
store = newFakeMailboxStore()
srv := NewServerWithGuardConfig(auth, store, guardCfg)
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
return listener.Addr().String(), store, func() {
cancel()
<-done
}
}
// TestGuard_ChaosHardCutDuringTransferNoLeak ist die geforderte
// Pflichtprüfung 1 (ING-07): Verbindung wird während aktiver
// Übertragung hart gekappt, danach kein Ressourcenleck.
func TestGuard_ChaosHardCutDuringTransferNoLeak(t *testing.T) {
addr, _, stop := startTestServerWithGuardConfig(t, protoguard.DefaultConfig())
defer stop()
runtime.GC()
baseline := runtime.NumGoroutine()
const rounds = 30
for i := 0; i < rounds; i++ {
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
reader := bufio.NewReader(conn)
_, _ = reader.ReadString('\n') // Begrüßung
_, _ = conn.Write([]byte("USER alice\r\n"))
_, _ = reader.ReadString('\n')
_, _ = conn.Write([]byte("PASS geheim123\r\n"))
_, _ = reader.ReadString('\n')
// Mitten in der Multiline-Übertragung (RETR-Antwortkopf gelesen,
// Datenzeilen NICHT vollständig konsumiert) hart abbrechen.
_, _ = conn.Write([]byte("RETR 1\r\n"))
_, _ = reader.ReadString('\n') // nur die "+OK ... octets"-Kopfzeile
_ = conn.Close()
}
// Server braucht kurz, um die abgebrochenen Sessions abzuräumen.
deadline := time.Now().Add(3 * time.Second)
for {
runtime.GC()
current := runtime.NumGoroutine()
if current <= baseline+2 { // kleine Toleranz für Laufzeit-Jitter
return
}
if time.Now().After(deadline) {
t.Fatalf("goroutine-leck nach hartem Verbindungsabbruch: baseline=%d, aktuell=%d", baseline, current)
}
time.Sleep(50 * time.Millisecond)
}
}
// TestGuard_TimeoutPerPhase ist die geforderte Pflichtprüfung 2
// (ING-07): Timeout-Auslösung in jeder Protokollphase.
func TestGuard_TimeoutPerPhase(t *testing.T) {
cfg := protoguard.Config{
PhaseTimeout: map[protoguard.Phase]time.Duration{
phaseAuthorization: 100 * time.Millisecond,
phaseTransaction: 100 * time.Millisecond,
},
DefaultTimeout: 5 * time.Second,
}
t.Run("authorization", func(t *testing.T) {
addr, _, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
defer func() { _ = conn.Close() }()
reader := bufio.NewReader(conn)
_, _ = reader.ReadString('\n') // Begrüßung, aber nichts weiter senden
_ = conn.SetReadDeadline(time.Now().Add(2 * time.Second))
_, err = reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungsende durch Authorization-Timeout")
}
})
t.Run("transaction", func(t *testing.T) {
addr, _, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
defer func() { _ = conn.Close() }()
reader := bufio.NewReader(conn)
_, _ = reader.ReadString('\n')
_, _ = conn.Write([]byte("USER alice\r\n"))
_, _ = reader.ReadString('\n')
_, _ = conn.Write([]byte("PASS geheim123\r\n"))
_, _ = reader.ReadString('\n') // jetzt in Transaction, nichts weiter senden
_ = conn.SetReadDeadline(time.Now().Add(2 * time.Second))
_, err = reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungsende durch Transaction-Timeout")
}
})
}
// TestGuard_BackoffOnRepeatedAuthFailures ist die geforderte
// Pflichtprüfung 3 (ING-07): Backoff-Verhalten bei wiederholten
// Fehlversuchen statt Dauerschleife.
func TestGuard_BackoffOnRepeatedAuthFailures(t *testing.T) {
cfg := protoguard.Config{
DefaultTimeout: 5 * time.Second,
MaxAuthFailures: 3,
BackoffBase: 50 * time.Millisecond,
BackoffMax: 500 * time.Millisecond,
}
addr, _, stop := startTestServerWithGuardConfig(t, cfg)
defer stop()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
defer func() { _ = conn.Close() }()
reader := bufio.NewReader(conn)
_, _ = reader.ReadString('\n')
var attemptDurations []time.Duration
for i := 0; i < 3; i++ {
_, _ = conn.Write([]byte("USER alice\r\n"))
_, _ = reader.ReadString('\n')
start := time.Now()
_, _ = conn.Write([]byte("PASS falsch\r\n"))
_ = conn.SetReadDeadline(time.Now().Add(3 * time.Second))
resp, err := reader.ReadString('\n')
if err != nil {
if i < 2 {
t.Fatalf("fehlversuch %d: unerwarteter Verbindungsabbruch: %v", i+1, err)
}
// dritter Fehlversuch: Trennung nach der Antwort ist erlaubt.
} else if !strings.Contains(resp, "-ERR") {
t.Fatalf("fehlversuch %d: erwartete -ERR, habe: %q", i+1, resp)
}
attemptDurations = append(attemptDurations, time.Since(start))
}
// Backoff steigt: der zweite Fehlversuch muss spürbar länger dauern
// als der erste (Verdopplung statt konstanter/keiner Wartezeit).
if attemptDurations[1] <= attemptDurations[0] {
t.Fatalf("erwartete steigenden Backoff, habe Dauern: %v", attemptDurations)
}
// Nach MaxAuthFailures muss die Verbindung getrennt werden (kein
// endloses erneutes USER/PASS erlaubt) statt in Dauerschleife.
_ = conn.SetReadDeadline(time.Now().Add(2 * time.Second))
if _, err := conn.Write([]byte("USER alice\r\n")); err == nil {
_, err = reader.ReadString('\n')
if err == nil {
t.Fatalf("erwartete Verbindungstrennung nach %d Fehlversuchen", cfg.MaxAuthFailures)
}
}
}
+38
View File
@@ -0,0 +1,38 @@
// Package pop3 implementiert ING-02: den POP3-Server (RFC 1939) mit den
// Zuständen Authorization/Transaction/Update und den Kernbefehlen
// USER/PASS/STAT/LIST/RETR/DELE/QUIT. Bewusste Neuimplementierung nach
// NEXARCH-Techstack, kein 1:1-Übernehmen von archivmail — gleiche
// Konvention wie mail/internal/imap (ING-01): eigene, schmale
// Authenticator/MailboxStore-Schnittstellen statt geteilter Typen über
// Paketgrenzen hinweg, CRLF-sichere Antworten (response.go).
package pop3
import "context"
// Authenticator prüft Zugangsdaten für PASS.
type Authenticator interface {
Authenticate(ctx context.Context, username, password string) (ok bool, err error)
}
// Message ist eine Nachricht im Postfach (nur Nummer/Größe für STAT/
// LIST — Inhalt kommt separat über MailboxStore.Retrieve, damit LIST
// nicht unnötig alle Nachrichteninhalte laden muss).
type Message struct {
Number int
Size int64
}
// MailboxStore liefert Postfachzustand für STAT/LIST/RETR/DELE.
type MailboxStore interface {
// List liefert alle (noch nicht gelöschten) Nachrichten des Postfachs
// username.
List(ctx context.Context, username string) ([]Message, error)
// Retrieve liefert den vollständigen Inhalt einer Nachricht
// (Akzeptanzkriterium 2: RETR liefert vollständige Nachrichten).
Retrieve(ctx context.Context, username string, number int) ([]byte, error)
// Delete löscht die angegebenen Nachrichtennummern ENDGÜLTIG — wird
// AUSSCHLIESSLICH im Update-Zustand nach einem regulären QUIT
// aufgerufen (Akzeptanzkriterium 2/Pflichtprüfung 3: DELE markiert
// nur innerhalb der Sitzung, committet wird erst hier).
Delete(ctx context.Context, username string, numbers []int) error
}
+26
View File
@@ -0,0 +1,26 @@
package pop3
import "strings"
// command ist eine geparste POP3-Kommandozeile — POP3 hat (anders als
// IMAP) keine Tags, nur "KOMMANDO [Argumente]".
type command struct {
Name string // groß geschrieben (z. B. "USER")
Args []string
}
// parseCommandLine zerlegt eine Kommandozeile (bereits ohne CRLF) in
// Kommandoname und Leerzeichen-getrennte Argumente. POP3-Argumente
// (Benutzername/Passwort/Nachrichtennummern) enthalten in der Praxis
// keine Anführungszeichen-Syntax wie IMAP — ein einfacher Split genügt
// für die kleinste Lösung.
func parseCommandLine(line string) command {
fields := strings.Fields(line)
if len(fields) == 0 {
return command{}
}
return command{
Name: strings.ToUpper(fields[0]),
Args: fields[1:],
}
}
+301
View File
@@ -0,0 +1,301 @@
package pop3
import (
"bufio"
"context"
"errors"
"net"
"strings"
"sync"
"testing"
"time"
)
type fakeAuthenticator struct {
users map[string]string
}
func (f fakeAuthenticator) Authenticate(_ context.Context, username, password string) (bool, error) {
want, ok := f.users[username]
return ok && want == password, nil
}
// fakeMailboxStore hält Nachrichten im Prozessspeicher — Delete entfernt
// sie erst bei tatsächlichem Aufruf (durch handleQuit im Update-Zustand).
type fakeMailboxStore struct {
mu sync.Mutex
messages map[string]map[int]string // username -> nummer -> inhalt
}
func newFakeMailboxStore() *fakeMailboxStore {
return &fakeMailboxStore{messages: map[string]map[int]string{
"alice": {1: "Erste Testnachricht\nmit zwei Zeilen", 2: "Zweite Testnachricht"},
}}
}
func (f *fakeMailboxStore) List(_ context.Context, username string) ([]Message, error) {
f.mu.Lock()
defer f.mu.Unlock()
msgs := f.messages[username]
result := make([]Message, 0, len(msgs))
for n, content := range msgs {
result = append(result, Message{Number: n, Size: int64(len(content))})
}
return result, nil
}
func (f *fakeMailboxStore) Retrieve(_ context.Context, username string, number int) ([]byte, error) {
f.mu.Lock()
defer f.mu.Unlock()
content, ok := f.messages[username][number]
if !ok {
return nil, errors.New("keine solche nachricht")
}
return []byte(content), nil
}
func (f *fakeMailboxStore) Delete(_ context.Context, username string, numbers []int) error {
f.mu.Lock()
defer f.mu.Unlock()
for _, n := range numbers {
delete(f.messages[username], n)
}
return nil
}
func (f *fakeMailboxStore) count(username string) int {
f.mu.Lock()
defer f.mu.Unlock()
return len(f.messages[username])
}
func startTestServer(t *testing.T) (addr string, store *fakeMailboxStore, stop func()) {
t.Helper()
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
store = newFakeMailboxStore()
srv := NewServer(auth, store)
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
return listener.Addr().String(), store, func() {
cancel()
<-done
}
}
type pop3Client struct {
conn net.Conn
reader *bufio.Reader
}
func dial(t *testing.T, addr string) *pop3Client {
t.Helper()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
c := &pop3Client{conn: conn, reader: bufio.NewReader(conn)}
c.readLine(t) // Begrüßung
return c
}
func (c *pop3Client) readLine(t *testing.T) string {
t.Helper()
_ = c.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
line, err := c.reader.ReadString('\n')
if err != nil {
t.Fatalf("antwort lesen: %v", err)
}
return strings.TrimRight(line, "\r\n")
}
// send sendet EIN Kommando und liest EINE Antwortzeile (Statuszeile).
func (c *pop3Client) send(t *testing.T, cmd string) string {
t.Helper()
if _, err := c.conn.Write([]byte(cmd + "\r\n")); err != nil {
t.Fatalf("kommando senden: %v", err)
}
return c.readLine(t)
}
// sendMultiline sendet ein Kommando und liest bis zur "."-Abschlusszeile.
func (c *pop3Client) sendMultiline(t *testing.T, cmd string) (status string, dataLines []string) {
t.Helper()
status = c.send(t, cmd)
if !strings.HasPrefix(status, "+OK") {
return status, nil
}
for {
line := c.readLine(t)
if line == "." {
return status, dataLines
}
dataLines = append(dataLines, line)
}
}
func (c *pop3Client) close() { _ = c.conn.Close() }
func loginAsAlice(t *testing.T, c *pop3Client) {
t.Helper()
if resp := c.send(t, "USER alice"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("USER: %s", resp)
}
if resp := c.send(t, "PASS geheim123"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("PASS: %s", resp)
}
}
// TestSession_StateTransitions ist die geforderte Pflichtprüfung 1:
// automatisierter Test für jede Zustandsübergangs-Regel.
func TestSession_StateTransitions(t *testing.T) {
addr, _, stop := startTestServer(t)
defer stop()
c := dial(t, addr)
defer c.close()
// Verbotener Übergang: STAT/RETR/DELE in Authorization.
if resp := c.send(t, "STAT"); !strings.HasPrefix(resp, "-ERR") {
t.Fatalf("erwartete -ERR für STAT in Authorization, habe: %s", resp)
}
if resp := c.send(t, "RETR 1"); !strings.HasPrefix(resp, "-ERR") {
t.Fatalf("erwartete -ERR für RETR in Authorization, habe: %s", resp)
}
// PASS ohne vorheriges USER.
if resp := c.send(t, "PASS irgendwas"); !strings.HasPrefix(resp, "-ERR") {
t.Fatalf("erwartete -ERR für PASS ohne USER, habe: %s", resp)
}
// Authorization -> Transaction.
loginAsAlice(t, c)
// Verbotener Übergang: USER/PASS erneut in Transaction.
if resp := c.send(t, "USER alice"); !strings.HasPrefix(resp, "-ERR") {
t.Fatalf("erwartete -ERR für USER in Transaction, habe: %s", resp)
}
// In Transaction erlaubt: STAT.
if resp := c.send(t, "STAT"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("erwartete +OK für STAT in Transaction, habe: %s", resp)
}
// Transaction -> (Update, real durchlaufen) -> Verbindungsende.
if resp := c.send(t, "QUIT"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("erwartete +OK für QUIT, habe: %s", resp)
}
}
// TestCommands_RetrDeleFullCycle deckt Akzeptanzkriterium 2 ab: RETR
// liefert vollständige Nachrichten, DELE + QUIT löscht endgültig.
func TestCommands_RetrDeleFullCycle(t *testing.T) {
addr, store, stop := startTestServer(t)
defer stop()
c := dial(t, addr)
defer c.close()
loginAsAlice(t, c)
status, lines := c.sendMultiline(t, "RETR 1")
if !strings.HasPrefix(status, "+OK") {
t.Fatalf("RETR: %s", status)
}
full := strings.Join(lines, "\n")
if full != "Erste Testnachricht\nmit zwei Zeilen" {
t.Fatalf("RETR lieferte keine vollständige nachricht, habe: %q", full)
}
if resp := c.send(t, "DELE 1"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("DELE: %s", resp)
}
if resp := c.send(t, "QUIT"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("QUIT: %s", resp)
}
if store.count("alice") != 1 {
t.Fatalf("erwartete 1 verbleibende nachricht nach DELE+QUIT, habe %d", store.count("alice"))
}
}
// TestCommands_DeleWithoutQuitDeletesNothing ist die geforderte
// Pflichtprüfung 3: DELE ohne anschließendes QUIT löscht nichts
// endgültig.
func TestCommands_DeleWithoutQuitDeletesNothing(t *testing.T) {
addr, store, stop := startTestServer(t)
defer stop()
c := dial(t, addr)
loginAsAlice(t, c)
if resp := c.send(t, "DELE 1"); !strings.HasPrefix(resp, "+OK") {
t.Fatalf("DELE: %s", resp)
}
// Verbindung OHNE QUIT abrupt schließen.
c.close()
time.Sleep(100 * time.Millisecond) // server real verarbeiten lassen
if store.count("alice") != 2 {
t.Fatalf("erwartete weiterhin 2 nachrichten (kein QUIT, keine endgültige löschung), habe %d", store.count("alice"))
}
}
// TestPass_RejectsWithoutInformationLeak ist die geforderte
// Akzeptanzkriterium-3-Prüfung: fehlerhafte Anmeldeversuche ohne
// Informationspreisgabe.
func TestPass_RejectsWithoutInformationLeak(t *testing.T) {
addr, _, stop := startTestServer(t)
defer stop()
c1 := dial(t, addr)
defer c1.close()
c1.send(t, "USER unbekannter_nutzer")
respUnknownUser := c1.send(t, "PASS irgendwas")
c2 := dial(t, addr)
defer c2.close()
c2.send(t, "USER alice")
respWrongPassword := c2.send(t, "PASS falschespasswort")
if respUnknownUser != respWrongPassword {
t.Fatalf("unterschiedliche fehlermeldungen verraten, ob der nutzer existiert: %q vs %q", respUnknownUser, respWrongPassword)
}
if !strings.HasPrefix(respUnknownUser, "-ERR") {
t.Fatalf("erwartete -ERR, habe: %s", respUnknownUser)
}
}
// TestServer_ManyParallelSessions belegt Robustheit unter Last (Vorbild
// ING-01) — kein expliziter Lasttest im Ticket gefordert, aber sinnvolle
// Ergänzung zur Zustandsmaschinen-Testabdeckung.
func TestServer_ManyParallelSessions(t *testing.T) {
addr, _, stop := startTestServer(t)
defer stop()
const sessions = 20
var wg sync.WaitGroup
for i := 0; i < sessions; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
conn, err := net.DialTimeout("tcp", addr, 3*time.Second)
if err != nil {
t.Errorf("dial %d: %v", n, err)
return
}
defer func() { _ = conn.Close() }()
c := &pop3Client{conn: conn, reader: bufio.NewReader(conn)}
c.readLine(t)
loginAsAlice(t, c)
c.send(t, "STAT")
c.send(t, "QUIT")
}(i)
}
wg.Wait()
}
+58
View File
@@ -0,0 +1,58 @@
package pop3
import (
"bufio"
"strings"
)
// sanitizeResponseText entfernt eingebettete CR/LF aus text, BEVOR er in
// eine Antwortzeile eingebettet wird (Bekannter Fehler vermeiden — gleiche
// Konvention wie mail/internal/imap/response.go: archivmail erlaubte
// Header-/Zeilen-Injection durch Stringkonkatenation ohne CRLF-Prüfung).
func sanitizeResponseText(text string) string {
text = strings.ReplaceAll(text, "\r", "")
text = strings.ReplaceAll(text, "\n", "")
return text
}
func writeOK(w *bufio.Writer, text string) error {
_, err := w.WriteString("+OK " + sanitizeResponseText(text) + "\r\n")
if err != nil {
return err
}
return w.Flush()
}
func writeErr(w *bufio.Writer, text string) error {
_, err := w.WriteString("-ERR " + sanitizeResponseText(text) + "\r\n")
if err != nil {
return err
}
return w.Flush()
}
// writeMultiline schreibt eine mehrzeilige POP3-Antwort (LIST/RETR):
// "+OK ...\r\n" gefolgt von den Datenzeilen und einer abschließenden
// "." -Zeile (RFC 1939 §3). content wird an "\n" in Zeilen zerlegt; jede
// Zeile, die selbst mit "." beginnt, wird per "Byte-Stuffing" verdoppelt
// (RFC-Pflicht UND zusätzlicher Schutz gegen eine vorzeitig wirkende
// Terminierungszeile durch Nachrichteninhalt).
func writeMultiline(w *bufio.Writer, okText, content string) error {
if _, err := w.WriteString("+OK " + sanitizeResponseText(okText) + "\r\n"); err != nil {
return err
}
normalized := strings.ReplaceAll(content, "\r\n", "\n")
for _, line := range strings.Split(normalized, "\n") {
line = strings.TrimSuffix(line, "\r")
if strings.HasPrefix(line, ".") {
line = "." + line
}
if _, err := w.WriteString(line + "\r\n"); err != nil {
return err
}
}
if _, err := w.WriteString(".\r\n"); err != nil {
return err
}
return w.Flush()
}
+55
View File
@@ -0,0 +1,55 @@
package pop3
import (
"context"
"errors"
"fmt"
"net"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
// Server nimmt POP3-Verbindungen an und bedient jede in einer eigenen
// Goroutine (Akzeptanzkriterium 1) — gleiches Muster wie
// mail/internal/imap.Server. TLS/STARTTLS ist Sache von ING-06, nicht
// dieser Kachel.
type Server struct {
auth Authenticator
store MailboxStore
guardCfg protoguard.Config
}
func NewServer(auth Authenticator, store MailboxStore) *Server {
return NewServerWithGuardConfig(auth, store, protoguard.DefaultConfig())
}
// NewServerWithGuardConfig erlaubt abweichende Phase-Timeouts und
// Backoff-Parameter (ING-07), z. B. für Tests oder gehärtete
// Betriebsumgebungen.
func NewServerWithGuardConfig(auth Authenticator, store MailboxStore, guardCfg protoguard.Config) *Server {
return &Server{auth: auth, store: store, guardCfg: guardCfg}
}
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
go func() {
<-ctx.Done()
_ = listener.Close()
}()
for {
conn, err := listener.Accept()
if err != nil {
if ctx.Err() != nil {
return nil
}
var netErr net.Error
if errors.As(err, &netErr) && netErr.Timeout() {
continue
}
return fmt.Errorf("pop3: verbindung annehmen: %w", err)
}
session := newSession(conn, srv.auth, srv.store, srv.guardCfg)
go session.Serve(ctx)
}
}
+145
View File
@@ -0,0 +1,145 @@
package pop3
import (
"bufio"
"context"
"errors"
"io"
"net"
"strings"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
)
// phaseAuthorization/phaseTransaction sind die protoguard-Phasen dieser
// Sitzung (ING-07 Akzeptanzkriterium 2: Timeouts pro Protokollphase
// konfigurierbar).
const (
phaseAuthorization protoguard.Phase = "authorization"
phaseTransaction protoguard.Phase = "transaction"
)
// maxCommandLineBytes begrenzt eine einzelne Kommandozeile (defensive
// Fehlerbehandlung bei nicht-konformen Gegenstellen, gleiche Konvention
// wie mail/internal/imap).
const maxCommandLineBytes = 8192
// Session ist eine einzelne POP3-Verbindung mit eigener Zustandsmaschine
// (Akzeptanzkriterium 1).
type Session struct {
conn net.Conn
reader *bufio.Reader
writer *bufio.Writer
auth Authenticator
store MailboxStore
guard *protoguard.Guard
state State
pendingUsername string // nach USER, vor erfolgreichem PASS
username string // nach erfolgreichem PASS
deleted map[int]bool
}
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config) *Session {
return &Session{
conn: conn,
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
writer: bufio.NewWriter(conn),
auth: auth,
store: store,
guard: protoguard.New(guardCfg),
state: Authorization,
deleted: map[int]bool{},
}
}
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
func (s *Session) currentPhase() protoguard.Phase {
if s.state == Authorization {
return phaseAuthorization
}
return phaseTransaction
}
// State liefert den aktuellen Sitzungszustand (für Tests).
func (s *Session) State() State { return s.state }
// Serve führt die Sitzung bis QUIT oder Verbindungsende aus.
func (s *Session) Serve(ctx context.Context) {
defer func() { _ = s.conn.Close() }()
if err := writeOK(s.writer, "POP3 server ready"); err != nil {
return
}
for {
// Akzeptanzkriterium 2 (ING-07): Idle-Timeout pro Protokollphase,
// vor jedem Lesevorgang neu gesetzt, da ein Zustandswechsel die
// Phase (und damit den geltenden Timeout) ändern kann.
if err := s.guard.ApplyReadDeadline(s.conn, s.currentPhase()); err != nil {
return
}
line, err := s.readLine()
if err != nil {
// Verbindung endet OHNE QUIT (Timeout, Netzwerkabbruch oder
// harter Verbindungsabbruch) — Akzeptanzkriterium 1: die
// Session-Ressourcen (Verbindung, Reader/Writer) werden über
// das defer conn.Close() oben zuverlässig freigegeben.
// Zusätzlich Pflichtprüfung 3: als Deleted markierte
// Nachrichten dürfen dadurch NICHT gelöscht werden. Da
// store.Delete nur im regulären handleQuit aufgerufen wird,
// ist das hier bereits strukturell garantiert (kein Aufruf,
// keine Löschung).
return
}
if line == "" {
continue
}
cmd := parseCommandLine(line)
if cmd.Name == "" {
if err := writeErr(s.writer, "unrecognized command"); err != nil {
return
}
continue
}
if !s.dispatch(ctx, cmd) {
return
}
}
}
func (s *Session) readLine() (string, error) {
line, err := s.reader.ReadString('\n')
if err != nil {
if errors.Is(err, io.EOF) && line != "" {
return strings.TrimRight(line, "\r"), nil
}
return "", err
}
return strings.TrimRight(line, "\r\n"), nil
}
// dispatch verarbeitet EIN geparstes Kommando. false bedeutet: Sitzung
// beenden (QUIT abgeschlossen oder Schreibfehler).
func (s *Session) dispatch(ctx context.Context, cmd command) bool {
switch cmd.Name {
case "USER":
return s.handleUser(cmd)
case "PASS":
return s.handlePass(ctx, cmd)
case "STAT":
return s.handleStat(ctx)
case "LIST":
return s.handleList(ctx, cmd)
case "RETR":
return s.handleRetr(ctx, cmd)
case "DELE":
return s.handleDele(cmd)
case "QUIT":
return s.handleQuit(ctx)
default:
return writeErr(s.writer, "unknown command") == nil
}
}
+24
View File
@@ -0,0 +1,24 @@
package pop3
// State ist einer der drei POP3-Sitzungszustände (RFC 1939 §3),
// Akzeptanzkriterium 1.
type State int
const (
Authorization State = iota
Transaction
Update
)
func (s State) String() string {
switch s {
case Authorization:
return "AUTHORIZATION"
case Transaction:
return "TRANSACTION"
case Update:
return "UPDATE"
default:
return "unknown"
}
}
+119
View File
@@ -0,0 +1,119 @@
// Package protoguard bündelt die Fehlerbehandlungs- und
// Wiederverbindungslogik, die IMAP- und POP3-Sessions gemeinsam
// brauchen (ING-07): pro Protokollphase konfigurierbare Idle-Timeouts
// und Backoff statt Dauerschleife bei wiederholten Anmeldefehlern.
// Ressourcenaufräumung selbst passiert bereits strukturell durch
// defer conn.Close() in den Sessions — Guard sorgt dafür, dass dieser
// Pfad auch bei hängenden oder böswilligen Gegenstellen zuverlässig
// erreicht wird.
package protoguard
import (
"context"
"net"
"time"
)
// Phase identifiziert eine Protokollphase, für die ein eigener
// Idle-Timeout gilt.
type Phase string
// Config steuert Timeout- und Backoff-Verhalten einer Verbindung.
type Config struct {
// PhaseTimeout liefert den Idle-Timeout je Phase. Fehlt ein Eintrag,
// gilt DefaultTimeout.
PhaseTimeout map[Phase]time.Duration
// DefaultTimeout gilt, wenn für die aktuelle Phase kein eigener Wert
// gesetzt ist. 0 bedeutet: kein Timeout.
DefaultTimeout time.Duration
// MaxAuthFailures ist die Anzahl fehlgeschlagener Anmeldeversuche,
// nach der eine Verbindung getrennt wird. 0 bedeutet: unbegrenzt
// (kein Trennen, nur Backoff).
MaxAuthFailures int
// BackoffBase ist die Wartezeit vor der Antwort nach dem ersten
// Fehlversuch, verdoppelt sich je weiterem Fehlversuch bis
// BackoffMax.
BackoffBase time.Duration
BackoffMax time.Duration
}
// DefaultConfig liefert praxistaugliche Werte für Produktionsbetrieb.
func DefaultConfig() Config {
return Config{
DefaultTimeout: 5 * time.Minute,
MaxAuthFailures: 5,
BackoffBase: 200 * time.Millisecond,
BackoffMax: 5 * time.Second,
}
}
// Guard kapselt den Fehlerbehandlungszustand EINER Verbindung: aktuell
// angewandte Phase-Timeouts und Zahl der Anmeldefehlversuche.
type Guard struct {
cfg Config
authFailures int
}
// New erstellt einen Guard für eine einzelne Session.
func New(cfg Config) *Guard {
return &Guard{cfg: cfg}
}
// ApplyReadDeadline setzt die Lese-Deadline von conn passend zur
// angegebenen Protokollphase (Akzeptanzkriterium 2).
func (g *Guard) ApplyReadDeadline(conn net.Conn, phase Phase) error {
d := g.cfg.DefaultTimeout
if pd, ok := g.cfg.PhaseTimeout[phase]; ok {
d = pd
}
if d <= 0 {
return conn.SetReadDeadline(time.Time{})
}
return conn.SetReadDeadline(time.Now().Add(d))
}
// RecordAuthFailure zählt einen fehlgeschlagenen Anmeldeversuch dieser
// Verbindung und liefert die Backoff-Wartezeit vor der Fehlerantwort
// sowie ob die Verbindung danach getrennt werden muss (Akzeptanzkriterium
// 3: klar definierter Backoff statt Dauerschleife).
func (g *Guard) RecordAuthFailure() (backoff time.Duration, disconnect bool) {
g.authFailures++
backoff = g.backoffFor(g.authFailures)
disconnect = g.cfg.MaxAuthFailures > 0 && g.authFailures >= g.cfg.MaxAuthFailures
return backoff, disconnect
}
// ResetAuthFailures setzt den Fehlversuchszähler nach erfolgreicher
// Anmeldung zurück.
func (g *Guard) ResetAuthFailures() { g.authFailures = 0 }
func (g *Guard) backoffFor(failures int) time.Duration {
if g.cfg.BackoffBase <= 0 {
return 0
}
d := g.cfg.BackoffBase
for i := 1; i < failures; i++ {
d *= 2
if g.cfg.BackoffMax > 0 && d >= g.cfg.BackoffMax {
return g.cfg.BackoffMax
}
}
if g.cfg.BackoffMax > 0 && d > g.cfg.BackoffMax {
return g.cfg.BackoffMax
}
return d
}
// Wait wartet d, bricht aber bei ctx-Abbruch sofort ab, damit ein
// Server-Shutdown nicht auf eine laufende Backoff-Pause warten muss.
func (g *Guard) Wait(ctx context.Context, d time.Duration) {
if d <= 0 {
return
}
timer := time.NewTimer(d)
defer timer.Stop()
select {
case <-timer.C:
case <-ctx.Done():
}
}
+127
View File
@@ -0,0 +1,127 @@
package smtp
import (
"bytes"
"context"
"strings"
)
func (s *Session) handleHelo(verb, arg string) bool {
if strings.TrimSpace(arg) == "" {
return s.reply(501, verb+" requires a domain/address") == nil
}
// HELO/EHLO setzt den Envelope zurück, falls bereits einer im
// Aufbau war (RFC 5321 §4.1.1.1).
s.from = ""
s.to = nil
s.state = Ready
if verb == "EHLO" {
return s.replyMultiline(250, []string{"nexarch-mail greets " + arg, "8BITMIME"}) == nil
}
return s.reply(250, "nexarch-mail greets "+arg) == nil
}
// handleMailFrom ist Teil des Envelope-Aufbaus (Akzeptanzkriterium 1):
// die Absenderadresse wird vor der Annahme validiert.
func (s *Session) handleMailFrom(arg string) bool {
if s.state == Greeting {
return s.reply(503, "send HELO/EHLO first") == nil
}
addr, err := parseMailAddressArg(arg, "FROM")
if err != nil {
return s.reply(501, "invalid MAIL FROM syntax") == nil
}
if err := validateAddress(addr); err != nil {
// Akzeptanzkriterium 3: ungültige Absenderdaten -> saubere
// SMTP-Fehlermeldung statt Absturz oder Verbindungsabbruch.
return s.reply(553, "invalid sender address") == nil
}
s.from = addr
s.to = nil
s.state = MailFromSet
return s.reply(250, "OK") == nil
}
// handleRcptTo ist Teil des Envelope-Aufbaus (Akzeptanzkriterium 1):
// jede Empfängeradresse wird vor der Annahme validiert; mehrere RCPT TO
// sind erlaubt.
func (s *Session) handleRcptTo(arg string) bool {
if s.state != MailFromSet && s.state != RcptToSet {
return s.reply(503, "send MAIL FROM first") == nil
}
addr, err := parseMailAddressArg(arg, "TO")
if err != nil {
return s.reply(501, "invalid RCPT TO syntax") == nil
}
if err := validateAddress(addr); err != nil {
// Akzeptanzkriterium 3: ungültige Empfängerdaten -> saubere
// SMTP-Fehlermeldung statt Absturz oder Verbindungsabbruch.
return s.reply(553, "invalid recipient address") == nil
}
s.to = append(s.to, addr)
s.state = RcptToSet
return s.reply(250, "OK") == nil
}
func (s *Session) handleRset() bool {
s.from = ""
s.to = nil
if s.state != Greeting {
s.state = Ready
}
return s.reply(250, "OK") == nil
}
// handleData verlangt einen vollständig aufgebauten und validierten
// Envelope (Akzeptanzkriterium 1: Envelope UND Nachrichtengröße werden
// vor der Annahme geprüft) und liest die dot-gestuffte Nachricht bis zur
// Abschlusszeile ".".
func (s *Session) handleData(ctx context.Context) bool {
if s.state != RcptToSet {
return s.reply(503, "send MAIL FROM/RCPT TO first") == nil
}
if err := s.reply(354, "Start mail input; end with <CRLF>.<CRLF>"); err != nil {
return false
}
var buf bytes.Buffer
for {
line, err := s.readLine()
if err != nil {
return false
}
if line == "." {
break
}
// Byte-Stuffing rückgängig machen (RFC 5321 §4.5.2): eine Zeile,
// die mit "." beginnt, verliert genau diesen ersten Punkt.
line = strings.TrimPrefix(line, ".")
buf.WriteString(line)
buf.WriteString("\r\n")
if int64(buf.Len()) > s.maxMessageBytes {
// Akzeptanzkriterium 1: Nachrichtengröße wird VOR der
// endgültigen Annahme geprüft — sauberer Fehlercode statt
// unbegrenztem Pufferwachstum.
_ = s.drainUntilDot()
s.from = ""
s.to = nil
s.state = Ready
return s.reply(552, "message size exceeds fixed maximum message size") == nil
}
}
envelope := Envelope{From: s.from, To: s.to}
raw := buf.Bytes()
s.from = ""
s.to = nil
s.state = Ready
if s.sink != nil {
if err := s.sink.Accept(ctx, envelope, raw); err != nil {
return s.reply(451, "unable to accept message, try again later") == nil
}
}
return s.reply(250, "OK: message accepted") == nil
}
+23
View File
@@ -0,0 +1,23 @@
// Package smtp implementiert ING-03s SMTP-Server (RFC 5321) für
// eingehende Mails: HELO/EHLO, MAIL FROM, RCPT TO, DATA, RSET, QUIT.
// Bewusste Neuimplementierung nach NEXARCH-Techstack, gleiche Konvention
// wie mail/internal/imap und mail/internal/pop3 — eigene Session je
// Verbindung in eigener Goroutine, schmale Sink-Schnittstelle statt
// geteilter Typen über Paketgrenzen hinweg.
package smtp
import "context"
// Envelope ist der SMTP-Umschlag einer eingehenden Nachricht, wie er
// vor der DATA-Annahme validiert wurde (Akzeptanzkriterium 1).
type Envelope struct {
From string
To []string
}
// MessageSink nimmt eine vollständig empfangene, dot-entstuffte
// Nachricht entgegen — Speicherung/Weiterverarbeitung ist Sache
// anderer Kacheln.
type MessageSink interface {
Accept(ctx context.Context, envelope Envelope, raw []byte) error
}
+59
View File
@@ -0,0 +1,59 @@
package smtp
import (
"fmt"
"strings"
)
// parseCommand zerlegt eine Kommandozeile in Verb (großgeschrieben) und
// restliches Argument.
func parseCommand(line string) (verb, arg string) {
parts := strings.SplitN(strings.TrimSpace(line), " ", 2)
verb = strings.ToUpper(parts[0])
if len(parts) == 2 {
arg = strings.TrimSpace(parts[1])
}
return verb, arg
}
// parseMailAddressArg extrahiert die Adresse aus "FROM:<addr>" bzw.
// "TO:<addr>" (RFC 5321 §4.1.1.2/4.1.1.3). SMTP-Parameter wie SIZE=...
// werden für diese kleinste Lösung ignoriert.
func parseMailAddressArg(arg, keyword string) (string, error) {
trimmed := strings.TrimSpace(arg)
upper := strings.ToUpper(trimmed)
prefix := keyword + ":"
if !strings.HasPrefix(upper, prefix) {
return "", fmt.Errorf("smtp: erwartete %q am anfang von %q", prefix, arg)
}
rest := strings.TrimSpace(trimmed[len(prefix):])
if sp := strings.IndexByte(rest, ' '); sp >= 0 {
rest = rest[:sp]
}
rest = strings.TrimPrefix(rest, "<")
rest = strings.TrimSuffix(rest, ">")
if rest == "" {
return "", fmt.Errorf("smtp: leere adresse")
}
return rest, nil
}
// validateAddress prüft eine E-Mail-Adresse defensiv gegen
// Steuerzeichen und offensichtlich falsche Form (Akzeptanzkriterium 3:
// ungültige Empfänger-/Absenderdaten führen zu sauberer Fehlermeldung
// statt Absturz).
func validateAddress(addr string) error {
for _, r := range addr {
if r < 0x20 || r == 0x7f {
return fmt.Errorf("smtp: steuerzeichen in adresse")
}
}
at := strings.IndexByte(addr, '@')
if at <= 0 || at == len(addr)-1 {
return fmt.Errorf("smtp: ungültige adresse %q", addr)
}
if strings.IndexByte(addr[at+1:], '@') >= 0 {
return fmt.Errorf("smtp: ungültige adresse %q", addr)
}
return nil
}
+27
View File
@@ -0,0 +1,27 @@
package smtp
import "fmt"
// reply schreibt eine einzeilige SMTP-Antwort "code text\r\n".
func (s *Session) reply(code int, text string) error {
if _, err := fmt.Fprintf(s.writer, "%d %s\r\n", code, text); err != nil {
return err
}
return s.writer.Flush()
}
// replyMultiline schreibt eine mehrzeilige SMTP-Antwort (z. B. EHLO-
// Capability-Liste): alle Zeilen außer der letzten mit "-" statt " "
// nach dem Code (RFC 5321 §4.2.1).
func (s *Session) replyMultiline(code int, lines []string) error {
for i, line := range lines {
sep := "-"
if i == len(lines)-1 {
sep = " "
}
if _, err := fmt.Fprintf(s.writer, "%d%s%s\r\n", code, sep, line); err != nil {
return err
}
}
return s.writer.Flush()
}
+56
View File
@@ -0,0 +1,56 @@
package smtp
import (
"context"
"errors"
"fmt"
"net"
)
// defaultMaxMessageBytes ist die Standard-Höchstgröße einer
// angenommenen Nachricht (Akzeptanzkriterium 1).
const defaultMaxMessageBytes = 25 * 1024 * 1024 // 25 MiB
// Server nimmt SMTP-Verbindungen an und bedient jede in einer eigenen
// Goroutine — gleiches Muster wie mail/internal/imap.Server und
// mail/internal/pop3.Server. TLS/STARTTLS ist Sache von ING-06,
// Rate-Limiting Sache von ING-09, Protokoll-Logging Sache von ING-08 —
// keine dieser Kacheln.
type Server struct {
sink MessageSink
maxMessageBytes int64
}
func NewServer(sink MessageSink) *Server {
return NewServerWithMaxMessageBytes(sink, defaultMaxMessageBytes)
}
// NewServerWithMaxMessageBytes erlaubt eine abweichende
// Nachrichten-Höchstgröße, z. B. für Tests.
func NewServerWithMaxMessageBytes(sink MessageSink, maxMessageBytes int64) *Server {
return &Server{sink: sink, maxMessageBytes: maxMessageBytes}
}
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
go func() {
<-ctx.Done()
_ = listener.Close()
}()
for {
conn, err := listener.Accept()
if err != nil {
if ctx.Err() != nil {
return nil
}
var netErr net.Error
if errors.As(err, &netErr) && netErr.Timeout() {
continue
}
return fmt.Errorf("smtp: verbindung annehmen: %w", err)
}
session := newSession(conn, srv.sink, srv.maxMessageBytes)
go session.Serve(ctx)
}
}
+120
View File
@@ -0,0 +1,120 @@
package smtp
import (
"bufio"
"context"
"errors"
"io"
"net"
"strings"
)
// maxCommandLineBytes begrenzt eine einzelne Kommando-/DATA-Zeile
// (defensive Fehlerbehandlung bei nicht-konformen Gegenstellen statt
// optimistischem Parsing, gleiche Konvention wie mail/internal/imap und
// mail/internal/pop3).
const maxCommandLineBytes = 8192
// Session ist eine einzelne SMTP-Verbindung mit eigener
// Zustandsmaschine (Akzeptanzkriterium 1).
type Session struct {
conn net.Conn
reader *bufio.Reader
writer *bufio.Writer
sink MessageSink
maxMessageBytes int64
state State
from string
to []string
}
func newSession(conn net.Conn, sink MessageSink, maxMessageBytes int64) *Session {
return &Session{
conn: conn,
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
writer: bufio.NewWriter(conn),
sink: sink,
maxMessageBytes: maxMessageBytes,
state: Greeting,
}
}
// State liefert den aktuellen Sitzungszustand (für Tests).
func (s *Session) State() State { return s.state }
// Serve führt die Sitzung bis QUIT oder Verbindungsende aus.
func (s *Session) Serve(ctx context.Context) {
defer func() { _ = s.conn.Close() }()
if err := s.reply(220, "nexarch-mail SMTP server ready"); err != nil {
return
}
for {
line, err := s.readLine()
if err != nil {
return
}
if line == "" {
continue
}
verb, arg := parseCommand(line)
if !s.dispatch(ctx, verb, arg) {
return
}
}
}
func (s *Session) readLine() (string, error) {
line, err := s.reader.ReadString('\n')
if err != nil {
if errors.Is(err, io.EOF) && line != "" {
return strings.TrimRight(line, "\r"), nil
}
return "", err
}
return strings.TrimRight(line, "\r\n"), nil
}
// dispatch verarbeitet EIN geparstes Kommando. false bedeutet: Sitzung
// beenden (QUIT abgeschlossen oder nicht behebbarer Schreibfehler).
func (s *Session) dispatch(ctx context.Context, verb, arg string) bool {
switch verb {
case "HELO", "EHLO":
return s.handleHelo(verb, arg)
case "MAIL":
return s.handleMailFrom(arg)
case "RCPT":
return s.handleRcptTo(arg)
case "DATA":
return s.handleData(ctx)
case "RSET":
return s.handleRset()
case "NOOP":
return s.reply(250, "OK") == nil
case "QUIT":
_ = s.reply(221, "Bye")
return false
default:
return s.reply(500, "Command not recognized") == nil
}
}
// drainUntilDot liest Zeilen, ohne sie zu puffern, bis zur
// DATA-Abschlusszeile "." — hält das Protokoll nach einer wegen
// Größenüberschreitung abgelehnten Nachricht synchron, ohne den
// verworfenen Rest unbegrenzt im Speicher zu halten.
func (s *Session) drainUntilDot() error {
for {
line, err := s.readLine()
if err != nil {
return err
}
if line == "." {
return nil
}
}
}
+327
View File
@@ -0,0 +1,327 @@
package smtp
import (
"bufio"
"context"
"net"
"runtime"
"strings"
"sync"
"testing"
"time"
)
// fakeSink zeichnet angenommene Nachrichten im Prozessspeicher auf.
type fakeSink struct {
mu sync.Mutex
accepted []acceptedMessage
fail bool
}
type acceptedMessage struct {
envelope Envelope
raw []byte
}
func (f *fakeSink) Accept(_ context.Context, envelope Envelope, raw []byte) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.fail {
return errFakeSinkRejects
}
cp := make([]byte, len(raw))
copy(cp, raw)
f.accepted = append(f.accepted, acceptedMessage{envelope: envelope, raw: cp})
return nil
}
func (f *fakeSink) count() int {
f.mu.Lock()
defer f.mu.Unlock()
return len(f.accepted)
}
type sinkError string
func (e sinkError) Error() string { return string(e) }
const errFakeSinkRejects sinkError = "fake sink lehnt ab"
func startTestServer(t *testing.T, sink MessageSink, maxMessageBytes int64) (addr string, stop func()) {
t.Helper()
srv := NewServerWithMaxMessageBytes(sink, maxMessageBytes)
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
return listener.Addr().String(), func() {
cancel()
<-done
}
}
type smtpClient struct {
conn net.Conn
reader *bufio.Reader
}
func dial(t *testing.T, addr string) *smtpClient {
t.Helper()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial: %v", err)
}
c := &smtpClient{conn: conn, reader: bufio.NewReader(conn)}
c.readLine(t) // 220 Begrüßung
return c
}
func (c *smtpClient) readLine(t *testing.T) string {
t.Helper()
_ = c.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
line, err := c.reader.ReadString('\n')
if err != nil {
t.Fatalf("antwort lesen: %v", err)
}
return strings.TrimRight(line, "\r\n")
}
func (c *smtpClient) send(t *testing.T, cmd string) string {
t.Helper()
if _, err := c.conn.Write([]byte(cmd + "\r\n")); err != nil {
t.Fatalf("kommando senden: %v", err)
}
return c.readLine(t)
}
func (c *smtpClient) close() { _ = c.conn.Close() }
func code(line string) string {
if len(line) < 3 {
return line
}
return line[:3]
}
// TestSession_EnvelopeMustBeBuiltBeforeData ist die geforderte
// Zustandsmaschinen-Abdeckung für Akzeptanzkriterium 1: Envelope
// (HELO -> MAIL FROM -> RCPT TO) wird SCHRITTWEISE validiert, DATA ist
// erst nach vollständigem, gültigem Envelope erlaubt.
func TestSession_EnvelopeMustBeBuiltBeforeData(t *testing.T) {
sink := &fakeSink{}
addr, stop := startTestServer(t, sink, defaultMaxMessageBytes)
defer stop()
c := dial(t, addr)
defer c.close()
// MAIL FROM vor HELO -> 503.
if resp := c.send(t, "MAIL FROM:<a@example.com>"); code(resp) != "503" {
t.Fatalf("erwartete 503 für MAIL FROM vor HELO, habe: %s", resp)
}
if resp := c.send(t, "EHLO client.example.com"); code(resp) != "250" {
t.Fatalf("erwartete 250 für EHLO, habe: %s", resp)
}
// Mehrzeilige EHLO-Antwort vollständig lesen.
for {
line := c.readLine(t)
if strings.HasPrefix(line, "250 ") {
break
}
}
// RCPT TO vor MAIL FROM -> 503.
if resp := c.send(t, "RCPT TO:<b@example.com>"); code(resp) != "503" {
t.Fatalf("erwartete 503 für RCPT TO vor MAIL FROM, habe: %s", resp)
}
// DATA vor RCPT TO -> 503.
if resp := c.send(t, "DATA"); code(resp) != "503" {
t.Fatalf("erwartete 503 für DATA ohne RCPT TO, habe: %s", resp)
}
if resp := c.send(t, "MAIL FROM:<a@example.com>"); code(resp) != "250" {
t.Fatalf("erwartete 250 für MAIL FROM, habe: %s", resp)
}
if resp := c.send(t, "RCPT TO:<b@example.com>"); code(resp) != "250" {
t.Fatalf("erwartete 250 für RCPT TO, habe: %s", resp)
}
if resp := c.send(t, "DATA"); code(resp) != "354" {
t.Fatalf("erwartete 354 für DATA nach vollständigem Envelope, habe: %s", resp)
}
if resp := c.send(t, "Subject: test\r\n\r\nHallo\r\n."); code(resp) != "250" {
t.Fatalf("erwartete 250 nach abgeschlossener DATA, habe: %s", resp)
}
if sink.count() != 1 {
t.Fatalf("erwartete 1 angenommene nachricht, habe %d", sink.count())
}
}
// TestData_MessageSizeCheckedBeforeAcceptance ist die geforderte
// Pflichtprüfung für Akzeptanzkriterium 1 (Größenanteil): eine
// Nachricht über der konfigurierten Höchstgröße wird sauber
// zurückgewiesen, der Sink bekommt sie NICHT.
func TestData_MessageSizeCheckedBeforeAcceptance(t *testing.T) {
sink := &fakeSink{}
const tinyLimit = 32 // Bytes
addr, stop := startTestServer(t, sink, tinyLimit)
defer stop()
c := dial(t, addr)
defer c.close()
c.send(t, "EHLO client.example.com")
for {
line := c.readLine(t)
if strings.HasPrefix(line, "250 ") {
break
}
}
c.send(t, "MAIL FROM:<a@example.com>")
c.send(t, "RCPT TO:<b@example.com>")
if resp := c.send(t, "DATA"); code(resp) != "354" {
t.Fatalf("erwartete 354, habe: %s", resp)
}
// Die überlange Zeile überschreitet das Limit bereits selbst — der
// Server antwortet SOFORT mit 552, OHNE auf die Abschlusszeile "."
// zu warten (drainUntilDot liest sie erst danach weg, damit das
// Protokoll synchron bleibt). Deshalb hier NICHT auf eine Antwort
// zur ersten Zeile warten, sondern erst die Abschlusszeile senden
// und dann einmal lesen.
longBody := strings.Repeat("x", 200)
if _, err := c.conn.Write([]byte(longBody + "\r\n")); err != nil {
t.Fatalf("kommando senden: %v", err)
}
resp := c.send(t, ".")
if code(resp) != "552" {
t.Fatalf("erwartete 552 (nachricht zu groß), habe: %s", resp)
}
if sink.count() != 0 {
t.Fatalf("sink hätte die zu große nachricht nicht bekommen dürfen, habe %d", sink.count())
}
// Verbindung muss danach weiter benutzbar sein (kein Absturz/Hänger).
if resp := c.send(t, "NOOP"); code(resp) != "250" {
t.Fatalf("session nach größenfehler nicht mehr funktionsfähig: %s", resp)
}
}
// TestRcptTo_InvalidRecipientCleanError ist die geforderte
// Akzeptanzkriterium-3-Prüfung: ungültige Empfängerdaten führen zu
// sauberer SMTP-Fehlermeldung statt Absturz.
func TestRcptTo_InvalidRecipientCleanError(t *testing.T) {
sink := &fakeSink{}
addr, stop := startTestServer(t, sink, defaultMaxMessageBytes)
defer stop()
c := dial(t, addr)
defer c.close()
c.send(t, "EHLO client.example.com")
for {
line := c.readLine(t)
if strings.HasPrefix(line, "250 ") {
break
}
}
c.send(t, "MAIL FROM:<a@example.com>")
if resp := c.send(t, "RCPT TO:<keine-gueltige-adresse>"); code(resp) != "553" {
t.Fatalf("erwartete 553 für ungültigen empfänger, habe: %s", resp)
}
// Verbindung bleibt nutzbar — kein Absturz, kein Verbindungsabbruch.
if resp := c.send(t, "RCPT TO:<b@example.com>"); code(resp) != "250" {
t.Fatalf("erwartete 250 für gültigen empfänger nach vorherigem fehler, habe: %s", resp)
}
}
// TestMailFrom_InvalidSenderCleanError deckt Akzeptanzkriterium 3 auch
// für den Absender ab.
func TestMailFrom_InvalidSenderCleanError(t *testing.T) {
sink := &fakeSink{}
addr, stop := startTestServer(t, sink, defaultMaxMessageBytes)
defer stop()
c := dial(t, addr)
defer c.close()
c.send(t, "EHLO client.example.com")
for {
line := c.readLine(t)
if strings.HasPrefix(line, "250 ") {
break
}
}
if resp := c.send(t, "MAIL FROM:<keine-gueltige-adresse>"); code(resp) != "553" {
t.Fatalf("erwartete 553 für ungültigen absender, habe: %s", resp)
}
if resp := c.send(t, "NOOP"); code(resp) != "250" {
t.Fatalf("session nach ungültigem absender nicht mehr funktionsfähig: %s", resp)
}
}
// TestServer_ConcurrentConnectionsNoLeak ist die geforderte
// Pflichtprüfung 3: Lasttest mit gleichzeitigen Verbindungen ohne
// Verbindungsleck.
func TestServer_ConcurrentConnectionsNoLeak(t *testing.T) {
sink := &fakeSink{}
addr, stop := startTestServer(t, sink, defaultMaxMessageBytes)
defer stop()
runtime.GC()
baseline := runtime.NumGoroutine()
const concurrency = 50
var wg sync.WaitGroup
for i := 0; i < concurrency; i++ {
wg.Add(1)
go func() {
defer wg.Done()
conn, err := net.DialTimeout("tcp", addr, 3*time.Second)
if err != nil {
t.Errorf("dial: %v", err)
return
}
defer func() { _ = conn.Close() }()
c := &smtpClient{conn: conn, reader: bufio.NewReader(conn)}
c.readLine(t)
c.send(t, "EHLO client.example.com")
for {
line := c.readLine(t)
if strings.HasPrefix(line, "250 ") {
break
}
}
c.send(t, "MAIL FROM:<a@example.com>")
c.send(t, "RCPT TO:<b@example.com>")
c.send(t, "DATA")
c.send(t, "Subject: last\r\n\r\nHallo\r\n.")
c.send(t, "QUIT")
}()
}
wg.Wait()
if sink.count() != concurrency {
t.Fatalf("erwartete %d angenommene nachrichten, habe %d", concurrency, sink.count())
}
deadline := time.Now().Add(3 * time.Second)
for {
runtime.GC()
current := runtime.NumGoroutine()
if current <= baseline+2 {
return
}
if time.Now().After(deadline) {
t.Fatalf("verbindungs-/goroutine-leck nach lasttest: baseline=%d, aktuell=%d", baseline, current)
}
time.Sleep(50 * time.Millisecond)
}
}
+28
View File
@@ -0,0 +1,28 @@
package smtp
// State ist einer der vier SMTP-Sitzungszustände dieser Implementierung
// (RFC 5321 §3.3), Akzeptanzkriterium 1: Envelope wird schrittweise vor
// der DATA-Annahme aufgebaut und geprüft.
type State int
const (
Greeting State = iota // vor HELO/EHLO
Ready // nach HELO/EHLO, bereit für MAIL FROM
MailFromSet // nach gültigem MAIL FROM, wartet auf RCPT TO
RcptToSet // mind. ein gültiges RCPT TO, DATA erlaubt
)
func (s State) String() string {
switch s {
case Greeting:
return "GREETING"
case Ready:
return "READY"
case MailFromSet:
return "MAIL FROM SET"
case RcptToSet:
return "RCPT TO SET"
default:
return "unknown"
}
}
+80
View File
@@ -0,0 +1,80 @@
// Package syncalert implementiert IMP-08: Benachrichtigung bei
// wiederholtem Postfach-Sync-Ausfall, mit Eskalationsschwelle statt
// Einzel-Alarm pro Fehlversuch. Versand ausschließlich über den
// zentralen Core-Benachrichtigungs-Dispatcher (CFG-02, bereits Fertig)
// — dieses Paket baut KEINEN eigenen E-Mail-Versand, sondern ruft
// ausschließlich NotificationDispatcher.Enqueue auf (Akzeptanzkriterium
// 1), exakt einmal je Eskalation.
package syncalert
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
)
// NotificationDispatcher ist die schmale Schnittstelle zu Core CFG-02
// (internal/notify.Dispatcher.Enqueue) — Mail ruft ausschließlich diese
// EINE Methode auf, kein eigener Versandcode.
type NotificationDispatcher interface {
Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (id string, err error)
}
// HTTPNotificationDispatcher spricht CFG-02 über HTTP an — dieselbe
// Service-Credential-Konvention wie mail/internal/crypto.HTTPKEKProvider
// (API-02, X-Nexarch-Client-Id/Secret).
type HTTPNotificationDispatcher struct {
endpointURL string
clientID string
clientSecret string
httpClient *http.Client
}
func NewHTTPNotificationDispatcher(endpointURL, clientID, clientSecret string, httpClient *http.Client) *HTTPNotificationDispatcher {
if httpClient == nil {
httpClient = http.DefaultClient
}
return &HTTPNotificationDispatcher{endpointURL: endpointURL, clientID: clientID, clientSecret: clientSecret, httpClient: httpClient}
}
type enqueueRequest struct {
Channel string `json:"channel"`
Recipient string `json:"recipient"`
Payload map[string]any `json:"payload"`
}
type enqueueResponse struct {
ID string `json:"id"`
}
func (d *HTTPNotificationDispatcher) Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (string, error) {
body, err := json.Marshal(enqueueRequest{Channel: channel, Recipient: recipient, Payload: payload})
if err != nil {
return "", fmt.Errorf("syncalert: anfrage serialisieren: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, d.endpointURL, bytes.NewReader(body))
if err != nil {
return "", fmt.Errorf("syncalert: anfrage aufbauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Nexarch-Client-Id", d.clientID)
req.Header.Set("X-Nexarch-Client-Secret", d.clientSecret)
resp, err := d.httpClient.Do(req)
if err != nil {
return "", fmt.Errorf("syncalert: anfrage senden: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("syncalert: cfg-02 lehnte anfrage ab: status %d", resp.StatusCode)
}
var out enqueueResponse
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return "", fmt.Errorf("syncalert: antwort dekodieren: %w", err)
}
return out.ID, nil
}
@@ -0,0 +1,67 @@
package syncalert
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
)
// TestHTTPNotificationDispatcher_SendsCorrectRequestFormat beweist real
// über echtes HTTP, dass HTTPNotificationDispatcher Channel/Recipient/
// Payload sowie die Service-Credential-Header korrekt sendet — kein
// eigener E-Mail-Versand, nur ein einziger CFG-02-Aufruf
// (Akzeptanzkriterium 1).
func TestHTTPNotificationDispatcher_SendsCorrectRequestFormat(t *testing.T) {
var capturedBody map[string]any
var capturedClientID, capturedClientSecret string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
capturedClientID = r.Header.Get("X-Nexarch-Client-Id")
capturedClientSecret = r.Header.Get("X-Nexarch-Client-Secret")
if err := json.NewDecoder(r.Body).Decode(&capturedBody); err != nil {
t.Fatal(err)
}
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"id":"real-notification-id-123"}`))
}))
defer srv.Close()
dispatcher := NewHTTPNotificationDispatcher(srv.URL, "mail", "mail-service-secret", nil)
id, err := dispatcher.Enqueue(context.Background(), NotificationChannel, AdminRecipient, map[string]any{
"mailbox": "INBOX",
"reason": "verbindung abgelehnt",
})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
if id != "real-notification-id-123" {
t.Fatalf("erwartete reale id vom server, habe %q", id)
}
if capturedClientID != "mail" || capturedClientSecret != "mail-service-secret" {
t.Fatalf("service-credential-header fehlen/falsch: id=%q secret=%q", capturedClientID, capturedClientSecret)
}
if capturedBody["channel"] != NotificationChannel || capturedBody["recipient"] != AdminRecipient {
t.Fatalf("channel/recipient falsch übertragen: %+v", capturedBody)
}
payload, ok := capturedBody["payload"].(map[string]any)
if !ok || payload["mailbox"] != "INBOX" {
t.Fatalf("payload nicht korrekt übertragen: %+v", capturedBody)
}
}
// TestHTTPNotificationDispatcher_RejectedByServerReturnsError bestätigt,
// dass eine Ablehnung durch CFG-02 real als Fehler durchgereicht wird,
// statt stillschweigend zu verschwinden.
func TestHTTPNotificationDispatcher_RejectedByServerReturnsError(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusForbidden)
}))
defer srv.Close()
dispatcher := NewHTTPNotificationDispatcher(srv.URL, "mail", "falsch", nil)
if _, err := dispatcher.Enqueue(context.Background(), "c", "r", nil); err == nil {
t.Fatal("erwartete fehler bei abgelehnter anfrage, habe nil")
}
}
@@ -0,0 +1,11 @@
CREATE TABLE IF NOT EXISTS mail_sync_alert_state (
tenant_slug TEXT NOT NULL,
mailbox_name TEXT NOT NULL,
consecutive_failures INT NOT NULL DEFAULT 0,
alerted BOOLEAN NOT NULL DEFAULT false,
last_success_at TIMESTAMPTZ,
last_failure_reason TEXT,
last_failure_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (tenant_slug, mailbox_name)
)
+156
View File
@@ -0,0 +1,156 @@
package syncalert
import (
"context"
_ "embed"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
//go:embed migrations/0001_mail_sync_alert_state.sql
var schemaMigration string
const (
// NotificationChannel/AdminRecipient sind bewusst statisch (kleinste
// Lösung) — eine konfigurierbare Empfängerverwaltung ist Sache einer
// späteren Kachel, nicht Bestandteil von IMP-08.
NotificationChannel = "mail-sync-failure"
AdminRecipient = "mail-admins"
)
// Monitor verfolgt Sync-Fehlschläge je Mandant/Postfach und löst bei
// Überschreiten der Schwelle GENAU EINE Benachrichtigung aus
// (Akzeptanzkriterium 1).
type Monitor struct {
pool *pgxpool.Pool
dispatcher NotificationDispatcher
threshold int
now func() time.Time
}
// DefaultThreshold ist die Vorgabe-Eskalationsschwelle (konsekutive
// Fehlschläge), überschreibbar über WithThreshold.
const DefaultThreshold = 3
func NewMonitor(pool *pgxpool.Pool, dispatcher NotificationDispatcher) *Monitor {
return &Monitor{pool: pool, dispatcher: dispatcher, threshold: DefaultThreshold, now: time.Now}
}
// WithThreshold setzt eine abweichende Eskalationsschwelle.
func (m *Monitor) WithThreshold(threshold int) *Monitor {
m.threshold = threshold
return m
}
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
func (m *Monitor) EnsureSchema(ctx context.Context) error {
if _, err := m.pool.Exec(ctx, schemaMigration); err != nil {
return fmt.Errorf("syncalert: schema anlegen: %w", err)
}
return nil
}
type alertState struct {
consecutiveFailures int
alerted bool
lastSuccessAt *time.Time
}
func (m *Monitor) getOrCreate(ctx context.Context, tenantSlug, mailboxName string) (alertState, error) {
if _, err := m.pool.Exec(ctx, `
INSERT INTO mail_sync_alert_state (tenant_slug, mailbox_name)
VALUES ($1, $2)
ON CONFLICT (tenant_slug, mailbox_name) DO NOTHING
`, tenantSlug, mailboxName); err != nil {
return alertState{}, fmt.Errorf("syncalert: zustand anlegen: %w", err)
}
var st alertState
err := m.pool.QueryRow(ctx, `
SELECT consecutive_failures, alerted, last_success_at
FROM mail_sync_alert_state WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName).Scan(&st.consecutiveFailures, &st.alerted, &st.lastSuccessAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return alertState{}, fmt.Errorf("syncalert: gerade angelegten zustand nicht gefunden")
}
return alertState{}, fmt.Errorf("syncalert: zustand lesen: %w", err)
}
return st, nil
}
// RecordFailure verzeichnet einen fehlgeschlagenen Sync-Versuch. Erst
// wenn consecutive_failures die konfigurierte Schwelle ERSTMALIG
// erreicht (noch nicht "alerted"), wird GENAU EINE Benachrichtigung an
// CFG-02 ausgelöst (Akzeptanzkriterium 1) — weitere Fehlschläge danach
// lösen KEINE zusätzliche Benachrichtigung aus, solange der Alarmzustand
// nicht durch einen erfolgreichen Sync zurückgesetzt wurde (kein
// Einzel-Alarm pro Fehlversuch, keine Spam-Flut).
func (m *Monitor) RecordFailure(ctx context.Context, tenantSlug, mailboxName, reason string) error {
now := m.now()
st, err := m.getOrCreate(ctx, tenantSlug, mailboxName)
if err != nil {
return err
}
newFailures := st.consecutiveFailures + 1
if _, err := m.pool.Exec(ctx, `
UPDATE mail_sync_alert_state
SET consecutive_failures = $3, last_failure_reason = $4, last_failure_at = $5, updated_at = now()
WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName, newFailures, reason, now); err != nil {
return fmt.Errorf("syncalert: fehlschlag erfassen: %w", err)
}
if newFailures < m.threshold || st.alerted {
return nil
}
// Akzeptanzkriterium 2: Benachrichtigung enthält Postfach,
// Fehlerursache und Zeitpunkt des letzten erfolgreichen Abrufs.
payload := map[string]any{
"tenant_slug": tenantSlug,
"mailbox": mailboxName,
"reason": reason,
"consecutive_failures": newFailures,
"last_successful_sync": formatOptionalTime(st.lastSuccessAt),
}
if _, err := m.dispatcher.Enqueue(ctx, NotificationChannel, AdminRecipient, payload); err != nil {
return fmt.Errorf("syncalert: benachrichtigung auslösen: %w", err)
}
if _, err := m.pool.Exec(ctx, `
UPDATE mail_sync_alert_state SET alerted = true, updated_at = now()
WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName); err != nil {
return fmt.Errorf("syncalert: alarmzustand markieren: %w", err)
}
return nil
}
// RecordSuccess verzeichnet einen erfolgreichen Sync und setzt den
// Alarmzustand zurück (Akzeptanzkriterium 3) — der nächste Fehlschlag
// nach einem Erfolg beginnt wieder bei 0 konsekutiven Fehlschlägen.
func (m *Monitor) RecordSuccess(ctx context.Context, tenantSlug, mailboxName string) error {
now := m.now()
if _, err := m.pool.Exec(ctx, `
INSERT INTO mail_sync_alert_state (tenant_slug, mailbox_name, consecutive_failures, alerted, last_success_at)
VALUES ($1, $2, 0, false, $3)
ON CONFLICT (tenant_slug, mailbox_name) DO UPDATE
SET consecutive_failures = 0, alerted = false, last_success_at = $3, updated_at = now()
`, tenantSlug, mailboxName, now); err != nil {
return fmt.Errorf("syncalert: erfolg erfassen: %w", err)
}
return nil
}
func formatOptionalTime(t *time.Time) string {
if t == nil {
return ""
}
return t.UTC().Format(time.RFC3339)
}
+220
View File
@@ -0,0 +1,220 @@
// Integrationstest (IMP-08): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/folderstate/imapimport —
// TEST_TENANT_DSN.
package syncalert
import (
"context"
"fmt"
"os"
"sync"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
// fakeDispatcher zeichnet jeden Enqueue-Aufruf auf — echte HTTP-
// Anbindung ist Sache von dispatcher_http_test.go, hier wird die
// Eskalationslogik isoliert geprüft (gleiche Konvention wie
// fakeAuthenticator/fakeKEKProvider in anderen Mail-Paketen).
type fakeDispatcher struct {
mu sync.Mutex
calls []map[string]any
}
func (f *fakeDispatcher) Enqueue(_ context.Context, channel, recipient string, payload map[string]any) (string, error) {
f.mu.Lock()
defer f.mu.Unlock()
call := map[string]any{"channel": channel, "recipient": recipient}
for k, v := range payload {
call[k] = v
}
f.calls = append(f.calls, call)
return fmt.Sprintf("notif-%d", len(f.calls)), nil
}
func (f *fakeDispatcher) count() int {
f.mu.Lock()
defer f.mu.Unlock()
return len(f.calls)
}
func setupMonitor(t *testing.T, dispatcher NotificationDispatcher) *Monitor {
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() })
monitor := NewMonitor(pool, dispatcher).WithThreshold(3)
if err := monitor.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_sync_alert_state WHERE tenant_slug LIKE 'mandant-imp08-%'`)
})
return monitor
}
// TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification
// ist die geforderte Pflichtprüfung 1: N aufeinanderfolgende
// Fehlschläge lösen genau eine Benachrichtigung aus, keine Spam-Flut.
func TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification(t *testing.T) {
dispatcher := &fakeDispatcher{}
monitor := setupMonitor(t, dispatcher)
ctx := context.Background()
tenant := "mandant-imp08-schwelle"
// Schwelle ist 3 — die ersten 2 Fehlschläge dürfen NICHTS auslösen.
for i := 0; i < 2; i++ {
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
t.Fatalf("recordfailure %d: %v", i, err)
}
}
if dispatcher.count() != 0 {
t.Fatalf("erwartete 0 benachrichtigungen vor erreichen der schwelle, habe %d", dispatcher.count())
}
// Dritter Fehlschlag erreicht die Schwelle — GENAU EINE Benachrichtigung.
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
t.Fatalf("recordfailure 3: %v", err)
}
if dispatcher.count() != 1 {
t.Fatalf("erwartete genau 1 benachrichtigung bei erreichen der schwelle, habe %d", dispatcher.count())
}
// Weitere Fehlschläge DANACH dürfen KEINE zusätzliche Benachrichtigung
// auslösen (kein Einzel-Alarm pro Fehlversuch, keine Spam-Flut).
for i := 0; i < 5; i++ {
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
t.Fatalf("weiterer fehlschlag %d: %v", i, err)
}
}
if dispatcher.count() != 1 {
t.Fatalf("erwartete weiterhin genau 1 benachrichtigung nach 5 weiteren fehlschlägen, habe %d", dispatcher.count())
}
}
// TestRecordSuccess_EndsAlertStateVerifiably ist die geforderte
// Pflichtprüfung 2: erfolgreicher Lauf nach Ausfall beendet den
// Alarmzustand nachvollziehbar.
func TestRecordSuccess_EndsAlertStateVerifiably(t *testing.T) {
dispatcher := &fakeDispatcher{}
monitor := setupMonitor(t, dispatcher)
ctx := context.Background()
tenant := "mandant-imp08-reset"
for i := 0; i < 3; i++ {
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
t.Fatalf("recordfailure %d: %v", i, err)
}
}
if dispatcher.count() != 1 {
t.Fatalf("erwartete 1 benachrichtigung nach 3 fehlschlägen, habe %d", dispatcher.count())
}
if err := monitor.RecordSuccess(ctx, tenant, "INBOX"); err != nil {
t.Fatalf("recordsuccess: %v", err)
}
// Nachvollziehbar zurückgesetzt: der NÄCHSTE Fehlschlags-Zyklus muss
// real wieder bei 0 beginnen und erneut die volle Schwelle
// durchlaufen, bevor eine ZWEITE Benachrichtigung ausgelöst wird.
for i := 0; i < 2; i++ {
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
t.Fatalf("recordfailure nach reset %d: %v", i, err)
}
}
if dispatcher.count() != 1 {
t.Fatalf("erwartete weiterhin nur 1 benachrichtigung (schwelle nach reset noch nicht erreicht), habe %d", dispatcher.count())
}
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
t.Fatalf("dritter fehlschlag nach reset: %v", err)
}
if dispatcher.count() != 2 {
t.Fatalf("erwartete 2. benachrichtigung nach erneutem erreichen der schwelle, habe %d", dispatcher.count())
}
}
// TestRecordFailure_MultipleAffectedMailboxesStayIsolated ist die
// geforderte Pflichtprüfung 3: Test mit mehreren betroffenen
// Postfächern gleichzeitig bleibt übersichtlich (korrekt isoliert).
func TestRecordFailure_MultipleAffectedMailboxesStayIsolated(t *testing.T) {
dispatcher := &fakeDispatcher{}
monitor := setupMonitor(t, dispatcher)
ctx := context.Background()
tenant := "mandant-imp08-mehrere"
mailboxes := []string{"INBOX", "Archiv", "Vertrieb"}
var wg sync.WaitGroup
for _, mailbox := range mailboxes {
wg.Add(1)
go func(mb string) {
defer wg.Done()
for i := 0; i < 3; i++ {
_ = monitor.RecordFailure(ctx, tenant, mb, "gleichzeitiger ausfall")
}
}(mailbox)
}
wg.Wait()
if dispatcher.count() != len(mailboxes) {
t.Fatalf("erwartete genau 1 benachrichtigung je betroffenem postfach (%d), habe %d", len(mailboxes), dispatcher.count())
}
seenMailboxes := map[string]bool{}
dispatcher.mu.Lock()
for _, call := range dispatcher.calls {
mb, _ := call["mailbox"].(string)
if seenMailboxes[mb] {
t.Fatalf("postfach %q hat mehr als eine benachrichtigung erhalten", mb)
}
seenMailboxes[mb] = true
}
dispatcher.mu.Unlock()
for _, mb := range mailboxes {
if !seenMailboxes[mb] {
t.Fatalf("postfach %q fehlt unter den benachrichtigten, habe: %v", mb, seenMailboxes)
}
}
}
// TestRecordFailure_NotificationContainsRequiredFields deckt
// Akzeptanzkriterium 2 ab: Benachrichtigung enthält Postfach,
// Fehlerursache und Zeitpunkt des letzten erfolgreichen Abrufs.
func TestRecordFailure_NotificationContainsRequiredFields(t *testing.T) {
dispatcher := &fakeDispatcher{}
monitor := setupMonitor(t, dispatcher)
ctx := context.Background()
tenant := "mandant-imp08-inhalt"
if err := monitor.RecordSuccess(ctx, tenant, "INBOX"); err != nil {
t.Fatalf("initialer erfolg: %v", err)
}
for i := 0; i < 3; i++ {
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "authentifizierung fehlgeschlagen"); err != nil {
t.Fatalf("recordfailure %d: %v", i, err)
}
}
if dispatcher.count() != 1 {
t.Fatalf("erwartete 1 benachrichtigung, habe %d", dispatcher.count())
}
call := dispatcher.calls[0]
if call["mailbox"] != "INBOX" {
t.Fatalf("erwartete postfach 'INBOX' in der benachrichtigung, habe: %v", call["mailbox"])
}
if call["reason"] != "authentifizierung fehlgeschlagen" {
t.Fatalf("erwartete fehlerursache in der benachrichtigung, habe: %v", call["reason"])
}
lastSuccess, _ := call["last_successful_sync"].(string)
if lastSuccess == "" {
t.Fatal("erwartete zeitpunkt des letzten erfolgreichen abrufs in der benachrichtigung")
}
}
@@ -0,0 +1,70 @@
// fakeClamd implementiert das reale clamd-INSTREAM-Protokoll
// protokolltreu (kein echter ClamAV-Daemon auf dem Testhost installiert
// — siehe Paket-Dokumentation in scanner.go). Erkennt die offizielle
// EICAR-Testsignatur exakt wie ein echter Virenscanner es täte.
package virusscan
import (
"encoding/binary"
"io"
"net"
"strings"
"testing"
)
// eicarTestString ist die offizielle, von allen Antivirus-Herstellern
// gemeinsam definierte, VOLLKOMMEN UNGEFÄHRLICHE Testsignatur (EICAR
// Institute) — kein echter Schadcode, universeller Standardtest für
// Virenscanner-Integrationen.
const eicarTestString = `X5O!P%@AP[4\PZX54(P^)7CC)7}$EICAR-STANDARD-ANTIVIRUS-TEST-FILE!$H+H*`
func startFakeClamd(t *testing.T) (addr string) {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
go func() {
for {
conn, err := listener.Accept()
if err != nil {
return
}
go handleFakeClamdConn(conn)
}
}()
t.Cleanup(func() { _ = listener.Close() })
return listener.Addr().String()
}
func handleFakeClamdConn(conn net.Conn) {
defer func() { _ = conn.Close() }()
header := make([]byte, len("zINSTREAM\x00"))
if _, err := io.ReadFull(conn, header); err != nil {
return
}
var content []byte
for {
var lenBuf [4]byte
if _, err := io.ReadFull(conn, lenBuf[:]); err != nil {
return
}
chunkLen := binary.BigEndian.Uint32(lenBuf[:])
if chunkLen == 0 {
break
}
chunk := make([]byte, chunkLen)
if _, err := io.ReadFull(conn, chunk); err != nil {
return
}
content = append(content, chunk...)
}
if strings.Contains(string(content), "EICAR-STANDARD-ANTIVIRUS-TEST-FILE") {
_, _ = conn.Write([]byte("stream: Eicar-Test-Signature FOUND\x00"))
return
}
_, _ = conn.Write([]byte("stream: OK\x00"))
}
@@ -0,0 +1,8 @@
CREATE TABLE IF NOT EXISTS mail_quarantine (
id BIGSERIAL PRIMARY KEY,
tenant_slug TEXT NOT NULL,
filename TEXT NOT NULL,
content_hash TEXT NOT NULL,
signature_name TEXT NOT NULL,
quarantined_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
+127
View File
@@ -0,0 +1,127 @@
package virusscan
import (
"context"
"crypto/sha256"
_ "embed"
"encoding/hex"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
//go:embed migrations/0001_mail_quarantine.sql
var schemaMigration string
// Decision ist das Ergebnis der Scan-Entscheidung für einen Anhang
// (Akzeptanzkriterium 1/2/3).
type Decision int
const (
// DecisionArchive: sauber, darf archiviert werden.
DecisionArchive Decision = iota
// DecisionQuarantine: Fund, Archivierung unterbleibt, Anhang
// gequarantänt (Akzeptanzkriterium 2).
DecisionQuarantine
// DecisionError: Scanner nicht erreichbar/Fehler — definierter
// Fehlerzustand statt automatischer Archivierung ODER unbegrenzter
// Blockade (Akzeptanzkriterium 3).
DecisionError
)
// QuarantineStore persistiert Quarantänefälle je Mandant.
type QuarantineStore struct {
pool *pgxpool.Pool
}
func NewQuarantineStore(pool *pgxpool.Pool) *QuarantineStore {
return &QuarantineStore{pool: pool}
}
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
func (s *QuarantineStore) EnsureSchema(ctx context.Context) error {
if _, err := s.pool.Exec(ctx, schemaMigration); err != nil {
return fmt.Errorf("virusscan: schema anlegen: %w", err)
}
return nil
}
func (s *QuarantineStore) record(ctx context.Context, tenantSlug, filename, contentHash, signatureName string) error {
if _, err := s.pool.Exec(ctx, `
INSERT INTO mail_quarantine (tenant_slug, filename, content_hash, signature_name)
VALUES ($1, $2, $3, $4)
`, tenantSlug, filename, contentHash, signatureName); err != nil {
return fmt.Errorf("virusscan: quarantänefall speichern: %w", err)
}
return nil
}
// List liefert alle Quarantänefälle eines Mandanten — Nachvollziehbarkeit
// (klare Statusanzeige, Akzeptanzkriterium 1).
func (s *QuarantineStore) List(ctx context.Context, tenantSlug string) ([]QuarantineEntry, error) {
rows, err := s.pool.Query(ctx, `
SELECT filename, content_hash, signature_name, quarantined_at
FROM mail_quarantine WHERE tenant_slug = $1 ORDER BY quarantined_at DESC
`, tenantSlug)
if err != nil {
return nil, fmt.Errorf("virusscan: quarantänefälle lesen: %w", err)
}
defer rows.Close()
var entries []QuarantineEntry
for rows.Next() {
var e QuarantineEntry
if err := rows.Scan(&e.Filename, &e.ContentHash, &e.SignatureName, &e.QuarantinedAt); err != nil {
return nil, fmt.Errorf("virusscan: quarantänezeile lesen: %w", err)
}
entries = append(entries, e)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("virusscan: quarantänefälle iterieren: %w", err)
}
return entries, nil
}
// QuarantineEntry ist ein einzelner Quarantänefall.
type QuarantineEntry struct {
Filename string
ContentHash string
SignatureName string
QuarantinedAt time.Time
}
// Processor verbindet Scanner mit QuarantineStore
// (Akzeptanzkriterium 1: jeder Anhang wird vor Archivierung geprüft).
type Processor struct {
scanner Scanner
quarantine *QuarantineStore
}
func NewProcessor(scanner Scanner, quarantine *QuarantineStore) *Processor {
return &Processor{scanner: scanner, quarantine: quarantine}
}
// ScanAndDecide prüft content und liefert die Archivierungsentscheidung.
// Bei DecisionQuarantine wurde der Fall bereits real in QuarantineStore
// verzeichnet, bevor ScanAndDecide zurückkehrt.
func (p *Processor) ScanAndDecide(ctx context.Context, tenantSlug, filename string, content []byte) (Decision, Result, error) {
result, err := p.scanner.Scan(ctx, content)
if err != nil {
if errors.Is(err, ErrScannerUnavailable) {
return DecisionError, Result{}, err
}
return DecisionError, Result{}, fmt.Errorf("virusscan: scan fehlgeschlagen: %w", err)
}
if result.Clean {
return DecisionArchive, result, nil
}
hash := sha256.Sum256(content)
if err := p.quarantine.record(ctx, tenantSlug, filename, hex.EncodeToString(hash[:]), result.SignatureName); err != nil {
return DecisionError, result, err
}
return DecisionQuarantine, result, nil
}
+164
View File
@@ -0,0 +1,164 @@
// Integrationstest (IMP-06): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/folderstate —
// TEST_TENANT_DSN. Der Virenscanner selbst ist der protokolltreue
// fakeClamd (siehe fake_clamd_test.go), die Netzwerk-/Protokollschicht
// (ClamdScanner) ist vollständig real.
package virusscan
import (
"context"
"errors"
"fmt"
"net"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupProcessor(t *testing.T, scanner Scanner) (*Processor, *QuarantineStore, 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() })
quarantine := NewQuarantineStore(pool)
if err := quarantine.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
tenant := "mandant-imp06-virenscan"
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_quarantine WHERE tenant_slug LIKE 'mandant-%'`)
})
return NewProcessor(scanner, quarantine), quarantine, tenant
}
// TestScanAndDecide_EICARTriggersQuarantine ist die geforderte
// Pflichtprüfung 1: Test mit EICAR-Testdatei bestätigt
// Quarantäne-Verhalten.
func TestScanAndDecide_EICARTriggersQuarantine(t *testing.T) {
addr := startFakeClamd(t)
scanner := NewClamdScanner(addr)
processor, quarantine, tenant := setupProcessor(t, scanner)
ctx := context.Background()
decision, result, err := processor.ScanAndDecide(ctx, tenant, "eicar.txt", []byte(eicarTestString))
if err != nil {
t.Fatalf("scanandDecide: %v", err)
}
if decision != DecisionQuarantine {
t.Fatalf("erwartete DecisionQuarantine für EICAR, habe %v", decision)
}
if result.SignatureName == "" {
t.Fatal("erwartete gemeldeten signaturnamen bei fund")
}
entries, err := quarantine.List(ctx, tenant)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(entries) != 1 || entries[0].Filename != "eicar.txt" {
t.Fatalf("erwartete real verzeichneten quarantänefall für eicar.txt, habe: %+v", entries)
}
// Saubere Datei zum Vergleich: DARF archiviert werden.
decision2, _, err := processor.ScanAndDecide(ctx, tenant, "harmlos.txt", []byte("ganz normaler anhangsinhalt"))
if err != nil {
t.Fatalf("scanandDecide (harmlos): %v", err)
}
if decision2 != DecisionArchive {
t.Fatalf("erwartete DecisionArchive für harmlosen inhalt, habe %v", decision2)
}
}
// TestScan_ScannerUnreachableFailsFastNotHang ist die geforderte
// Pflichtprüfung 2: Scanner nicht erreichbar führt zu klar sichtbarem
// Fehlerzustand statt Hänger.
func TestScan_ScannerUnreachableFailsFastNotHang(t *testing.T) {
// Ein real geschlossener Port (nichts lauscht) — kein Hänger, sofortige
// Verbindungsablehnung durch das Betriebssystem.
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
unreachableAddr := listener.Addr().String()
_ = listener.Close() // sofort wieder geschlossen -> Verbindung wird real abgelehnt
scanner := NewClamdScanner(unreachableAddr).WithTimeout(2 * time.Second)
start := time.Now()
_, err = scanner.Scan(context.Background(), []byte("beliebiger inhalt"))
elapsed := time.Since(start)
if err == nil {
t.Fatal("erwartete fehler bei nicht erreichbarem scanner, habe nil")
}
if !errors.Is(err, ErrScannerUnavailable) {
t.Fatalf("erwartete ErrScannerUnavailable, habe: %v", err)
}
if elapsed > 2*time.Second {
t.Fatalf("scan hing über die konfigurierte frist hinaus: %s", elapsed)
}
t.Logf("nicht erreichbarer scanner meldete real nach %s: %v", elapsed, err)
}
// TestScanAndDecide_ScannerUnavailableYieldsDefinedErrorState ergänzt
// Pflichtprüfung 2 auf Processor-Ebene: ScanAndDecide liefert
// DecisionError statt automatischer Archivierung.
func TestScanAndDecide_ScannerUnavailableYieldsDefinedErrorState(t *testing.T) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
unreachableAddr := listener.Addr().String()
_ = listener.Close()
scanner := NewClamdScanner(unreachableAddr).WithTimeout(1 * time.Second)
processor, _, tenant := setupProcessor(t, scanner)
decision, _, err := processor.ScanAndDecide(context.Background(), tenant, "irgendwas.pdf", []byte("inhalt"))
if err == nil {
t.Fatal("erwartete fehler, habe nil")
}
if decision != DecisionError {
t.Fatalf("erwartete DecisionError (NICHT automatische archivierung) bei nicht erreichbarem scanner, habe %v", decision)
}
}
// TestScan_ThroughputWithManyAttachmentsIsAcceptable ist die geforderte
// Pflichtprüfung 3: Durchsatztest bestätigt akzeptable Verzögerung durch
// den Scan-Schritt.
func TestScan_ThroughputWithManyAttachmentsIsAcceptable(t *testing.T) {
addr := startFakeClamd(t)
scanner := NewClamdScanner(addr)
const attachments = 50
const targetPerScan = 100 * time.Millisecond
start := time.Now()
for i := 0; i < attachments; i++ {
content := []byte(fmt.Sprintf("anhangsinhalt nummer %d, harmlos", i))
result, err := scanner.Scan(context.Background(), content)
if err != nil {
t.Fatalf("scan %d: %v", i, err)
}
if !result.Clean {
t.Fatalf("scan %d: erwartete sauberes ergebnis, habe fund %q", i, result.SignatureName)
}
}
elapsed := time.Since(start)
perScan := elapsed / attachments
t.Logf("Durchsatz: %d Anhänge in %s (%s/Anhang, Ziel %s/Anhang)", attachments, elapsed, perScan, targetPerScan)
if perScan > targetPerScan {
t.Fatalf("scan zu langsam: %s/anhang, ziel %s/anhang", perScan, targetPerScan)
}
}
+146
View File
@@ -0,0 +1,146 @@
// Package virusscan implementiert IMP-06: Anbindung eines Virenscanners
// für importierte Anhänge, mit Quarantäne-Verhalten bei Fund und klarer
// Statusanzeige. Kein Vorbild in archivmail für diesen Zuschnitt — Neubau.
//
// ClamdScanner spricht das reale, dokumentierte clamd-INSTREAM-Protokoll
// (TCP, Längen-präfixierte Chunks) — kein ClamAV-Daemon wurde für diese
// Kachel auf dem Testhost installiert (ein Antivirus-Daemon samt
// Signaturdatenbank ist ein deutlich größerer, sicherheitsrelevanter
// Eingriff als ein einzelnes Go-Modul und wird nicht unaufgefordert
// vorgenommen). Stattdessen wird ein protokolltreuer Fake-Server für
// Tests verwendet (gleiches Prinzip wie IMP-08s
// HTTPNotificationDispatcher-Tests) — der reale Netzwerkpfad
// (ClamdScanner) ist vollständig echt und real getestet, nur die
// Gegenstelle ist ein Test-Double statt eines echten ClamAV-Daemons.
package virusscan
import (
"bufio"
"context"
"encoding/binary"
"errors"
"fmt"
"net"
"strings"
"time"
)
// Result ist das Ergebnis eines Scans (Akzeptanzkriterium 1).
type Result struct {
Clean bool
SignatureName string
}
// ErrScannerUnavailable wird geliefert, wenn der Virenscanner nicht
// erreichbar ist oder innerhalb der Frist nicht antwortet
// (Akzeptanzkriterium 3: definierter Fehlerzustand statt unbegrenzter
// Blockade).
var ErrScannerUnavailable = errors.New("virusscan: scanner nicht erreichbar")
// Scanner prüft Anhangsinhalte auf Schadsoftware.
type Scanner interface {
Scan(ctx context.Context, content []byte) (Result, error)
}
// ClamdScanner spricht das clamd-INSTREAM-Protokoll über TCP.
type ClamdScanner struct {
addr string
dialer net.Dialer
timeout time.Duration
}
// DefaultScanTimeout begrenzt einen einzelnen Scan-Vorgang
// (Akzeptanzkriterium 3).
const DefaultScanTimeout = 10 * time.Second
func NewClamdScanner(addr string) *ClamdScanner {
return &ClamdScanner{addr: addr, timeout: DefaultScanTimeout}
}
// WithTimeout überschreibt die Standard-Scan-Zeitüberschreitung (Tests
// nutzen eine kürzere Frist, um Nicht-Erreichbarkeit real zügig zu
// beweisen).
func (c *ClamdScanner) WithTimeout(d time.Duration) *ClamdScanner {
c.timeout = d
return c
}
const clamdChunkSize = 4096
// Scan überträgt content per INSTREAM (RFC-artiges, dokumentiertes
// clamd-Protokoll: "zINSTREAM\0" gefolgt von 4-Byte-Big-Endian-
// Längenpräfixen je Chunk, abgeschlossen durch ein Null-Längen-Chunk) und
// interpretiert die Antwortzeile.
func (c *ClamdScanner) Scan(ctx context.Context, content []byte) (Result, error) {
scanCtx := ctx
var cancel context.CancelFunc
if c.timeout > 0 {
scanCtx, cancel = context.WithTimeout(ctx, c.timeout)
defer cancel()
}
conn, err := c.dialer.DialContext(scanCtx, "tcp", c.addr)
if err != nil {
return Result{}, fmt.Errorf("%w: %v", ErrScannerUnavailable, err)
}
defer func() { _ = conn.Close() }()
if deadline, ok := scanCtx.Deadline(); ok {
_ = conn.SetDeadline(deadline)
}
if _, err := conn.Write([]byte("zINSTREAM\x00")); err != nil {
return Result{}, fmt.Errorf("%w: %v", ErrScannerUnavailable, err)
}
for offset := 0; offset < len(content); offset += clamdChunkSize {
end := offset + clamdChunkSize
if end > len(content) {
end = len(content)
}
chunk := content[offset:end]
var lenBuf [4]byte
binary.BigEndian.PutUint32(lenBuf[:], uint32(len(chunk)))
if _, err := conn.Write(lenBuf[:]); err != nil {
return Result{}, fmt.Errorf("%w: %v", ErrScannerUnavailable, err)
}
if _, err := conn.Write(chunk); err != nil {
return Result{}, fmt.Errorf("%w: %v", ErrScannerUnavailable, err)
}
}
// Null-Längen-Chunk signalisiert Ende des Streams.
var zero [4]byte
if _, err := conn.Write(zero[:]); err != nil {
return Result{}, fmt.Errorf("%w: %v", ErrScannerUnavailable, err)
}
reader := bufio.NewReader(conn)
line, err := reader.ReadString('\x00')
if err != nil {
return Result{}, fmt.Errorf("%w: antwort lesen: %v", ErrScannerUnavailable, err)
}
line = strings.TrimRight(line, "\x00\r\n")
return parseClamdResponse(line)
}
// parseClamdResponse interpretiert eine clamd-Antwortzeile, z. B.
// "stream: OK" oder "stream: Eicar-Test-Signature FOUND".
func parseClamdResponse(line string) (Result, error) {
switch {
case strings.HasSuffix(line, "OK"):
return Result{Clean: true}, nil
case strings.HasSuffix(line, "FOUND"):
// Format: "stream: <Signaturname> FOUND"
trimmed := strings.TrimSuffix(line, "FOUND")
trimmed = strings.TrimSpace(trimmed)
signature := trimmed
if idx := strings.LastIndex(trimmed, ":"); idx != -1 {
signature = strings.TrimSpace(trimmed[idx+1:])
}
return Result{Clean: false, SignatureName: signature}, nil
default:
return Result{}, fmt.Errorf("virusscan: unerwartete scanner-antwort: %q", line)
}
}