Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5dcfa36f99 | ||
|
|
145a161f8a | ||
|
|
6e01cecca7 |
@@ -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.
|
||||||
@@ -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.
|
||||||
@@ -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.
|
||||||
@@ -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/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/ssooidc v1.40.1 // indirect
|
||||||
github.com/aws/aws-sdk-go-v2/service/sts v1.47.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/pgpassfile v1.0.0 // indirect
|
||||||
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
|
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
|
||||||
github.com/jackc/puddle/v2 v2.2.1 // indirect
|
github.com/jackc/puddle/v2 v2.2.1 // indirect
|
||||||
golang.org/x/crypto v0.17.0 // indirect
|
golang.org/x/crypto v0.17.0 // indirect
|
||||||
golang.org/x/sync v0.1.0 // indirect
|
golang.org/x/sync v0.1.0 // indirect
|
||||||
|
golang.org/x/sys v0.15.0 // indirect
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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.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 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
|
||||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
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 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
|
||||||
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
|
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=
|
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/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
|
||||||
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
|
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/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 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ=
|
||||||
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
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=
|
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)
|
||||||
|
)
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
)
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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()
|
||||||
|
)
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user