From ee98efb51e3f5bb58a4523fc11d8512ba07b6075 Mon Sep 17 00:00:00 2001 From: sysops Date: Sun, 30 Aug 2026 23:43:12 +0200 Subject: [PATCH] ARC-01: objekt-speicher-anbindung-fuer-mails-anhaenge - mail/internal/storage: LocalDriver/S3Driver (bewaehrtes Muster aus DMS FDN-03, bewusste Neuimplementierung - Mail kann DMS nicht importieren), ObjectKey mit festem Pfadschema - Service.Put/GetVerified: Pruefsummenverifikation AN DIESER SCHICHT (Erweiterung gegenueber FDN-03) - SHA-256-Sidecar, sofortige Ruecklese-Verifikation beim Schreiben, Erkennung manipulierter Objekte beim Lesen - HTTPUsageReporter: meldet an Core API-11 (resync-api/LIC-05), identisches Muster wie DMS FDN-03 - 4 Tests real bestanden: byteidentischer Read-back, manipuliertes Objekt erkannt, Lasttest (500 Objekte, 105.8us/Objekt), Nutzungsmeldung bei Schreiben+Loeschen - zusaetzlich echter End-zu-Ende-Beweis gegen den laufenden nexarch-resync-api.service: reales Service-Credential provisioniert, Put->GetVerified->Delete komplett durchlaufen, usage_counters zeigt reales +29/-29-Delta (beide Meldungen real angewendet) Pruefungen siehe mail/docs/ARC-01-PRUEFPROTOKOLL.md --- mail/docs/ARC-01-PRUEFPROTOKOLL.md | 54 +++++++++++ mail/go.mod | 22 ++++- mail/go.sum | 36 +++++++ mail/internal/storage/driver.go | 49 ++++++++++ mail/internal/storage/localdriver.go | 63 ++++++++++++ mail/internal/storage/s3driver.go | 110 +++++++++++++++++++++ mail/internal/storage/service.go | 127 ++++++++++++++++++++++++ mail/internal/storage/service_test.go | 134 ++++++++++++++++++++++++++ mail/internal/storage/usagereport.go | 75 ++++++++++++++ 9 files changed, 669 insertions(+), 1 deletion(-) create mode 100644 mail/docs/ARC-01-PRUEFPROTOKOLL.md create mode 100644 mail/internal/storage/driver.go create mode 100644 mail/internal/storage/localdriver.go create mode 100644 mail/internal/storage/s3driver.go create mode 100644 mail/internal/storage/service.go create mode 100644 mail/internal/storage/service_test.go create mode 100644 mail/internal/storage/usagereport.go diff --git a/mail/docs/ARC-01-PRUEFPROTOKOLL.md b/mail/docs/ARC-01-PRUEFPROTOKOLL.md new file mode 100644 index 0000000..88c29c1 --- /dev/null +++ b/mail/docs/ARC-01-PRUEFPROTOKOLL.md @@ -0,0 +1,54 @@ +# ARC-01 – Prüfprotokoll: Objekt-Speicher-Anbindung für Mails/Anhänge + +Voraussetzung ING-04 – bereits Fertig. ARC-01 ist der Startpunkt der +Foundation-Kette (analog DMS FDN-03), nicht nur eine Ergänzung — es +entsperrt ARC-02 bis ARC-10 sowie mehrere Ingestion-Tickets. + +## Umsetzung + +Bewährtes Muster aus DMS FDN-03 (LocalDriver/S3Driver-Abstraktion) +übernommen — bewusste Neuimplementierung statt Cross-Modul-Import +(Mail ist eigenständiges Go-Modul, kann DMS' `internal/` nicht +importieren): + +- `mail/internal/storage.Driver` — `Put`/`Get`/`Delete`, zwei + Implementierungen (`LocalDriver`, `S3Driver`). +- `ObjectKey(messageID, partIndex)` — festes, dokumentiertes + Pfadschema `messages//parts/` (Akzeptanzkriterium 1). + Lesezugriff hängt NUR von `messageID`+`partIndex` ab, nicht vom + ursprünglichen Importpfad (Akzeptanzkriterium 3). +- **Erweiterung gegenüber FDN-03** — Prüfsummenverifikation AN DIESER + SCHICHT (Akzeptanzkriterium 2, von ARC-01 explizit gefordert, anders + als FDN-03): `Service.Put` schreibt Inhalt + SHA-256-Sidecar-Objekt, + liest SOFORT zurück und verifiziert — ein fehlgeschlagener + Rücklese-Vergleich lässt `Put` selbst fehlschlagen, keine unbemerkt + fehlerhafte Ablage. `Service.GetVerified` wiederholt die Prüfung bei + jedem späteren Lesezugriff. +- `HTTPUsageReporter` — identisches Muster wie DMS FDN-03, meldet über + Core API-11 (`resync-api`, `internal/resync.Handler.UsageHandler`, + Service-Credential wie API-02) an LIC-05 (Akzeptanzkriterium 4). + +## Prüfungen + +| # | Prüfung | Ergebnis | +|---|---|---| +| 1 | Test: geschriebenes Objekt liefert beim Lesen byteidentischen Inhalt | **bestanden** – `TestPut_ReadBackIsByteIdentical`: `GetVerified` liefert exakt den geschriebenen Inhalt | +| 2 | Test: absichtlich beschädigtes Objekt wird bei Prüfsummenvergleich erkannt | **bestanden** – `TestGetVerified_DetectsTamperedObject`: Objekt direkt am Dateisystem manipuliert (umgeht `Service` vollständig), `GetVerified` liefert real `ErrChecksumMismatch` | +| 3 | Lasttest mit vielen kleinen Objekten bestätigt akzeptable Latenz | **bestanden** – `TestPut_ManySmallObjectsAcceptableLatency`: 500 reale `Put`-Aufrufe (inkl. Schreiben+Sidecar+Rücklese-Verifikation) in 52,9 ms — **105,8 µs/Objekt**, weit unter der 10-ms-Grenze | +| 4 | Melde-Aufruf an Core LIC-05 bei Schreib- und Löschvorgang nachweislich ausgelöst, mit korrekter Größenangabe | **bestanden** – `TestPut_ReportsUsageOnWriteAndDelete` (Fake-Reporter, exakte Delta-Werte); ZUSÄTZLICH real auf 131 gegen den laufenden `nexarch-resync-api.service` (API-11) bewiesen: echtes Service-Credential provisioniert, `Put`→`GetVerified`→`Delete` komplett durchlaufen, `usage_counters` zeigt reales Delta `+29` dann `-29` (Nettosumme 0 — beide Meldungen real angewendet, nicht nur eine) | + +## Build/Test-Ergebnis (192.168.1.131) + +``` +go build ./... -> clean +go vet ./... -> clean +golangci-lint run ./... -> 0 issues +go test ./... -p 1 -> alle Mail-Pakete bestanden (storage, mimeparse, example, pflichttestgate) +``` + +## Gesamtergebnis + +**Bestanden.** Alle vier Akzeptanzkriterien und alle vier +Pflichtprüfungen real erfüllt, inklusive eines echten End-zu-Ende-Laufs +gegen den live laufenden Core-API-11-Dienst (nicht nur einen Fake). +Entsperrt ARC-02–ARC-10 sowie mehrere Ingestion-Tickets. diff --git a/mail/go.mod b/mail/go.mod index c29d056..aa02577 100644 --- a/mail/go.mod +++ b/mail/go.mod @@ -1,13 +1,33 @@ module gitea.perlbach24.de/scripte/nexarch/mail -go 1.22 +go 1.24 + +toolchain go1.24.4 require ( + github.com/aws/aws-sdk-go-v2 v1.45.1 + github.com/aws/aws-sdk-go-v2/config v1.33.1 + github.com/aws/aws-sdk-go-v2/credentials v1.20.1 + github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1 + github.com/aws/smithy-go v1.28.1 github.com/jackc/pgx/v5 v5.6.0 golang.org/x/text v0.14.0 ) require ( + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.7.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/sts v1.47.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 diff --git a/mail/go.sum b/mail/go.sum index 5c39671..995792c 100644 --- a/mail/go.sum +++ b/mail/go.sum @@ -1,3 +1,39 @@ +github.com/aws/aws-sdk-go-v2 v1.45.1 h1:iIoG3NaLhV6UZpPXyPXlDj2I9oS8tV/nMcMnITCC6Ks= +github.com/aws/aws-sdk-go-v2 v1.45.1/go.mod h1:bttEH6JqnUL8LepvDVfdrds/fZ5bCIxzpe3abyUrhDU= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 h1:GPRlPwz40I2B2VrBEASOA3Bi77NyeqejNLkifosX0rs= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20/go.mod h1:g7PNzKcsOKWb4fkSRBA7BZVAS6Y8IcxzN+nRohhQ1Q8= +github.com/aws/aws-sdk-go-v2/config v1.33.1 h1:bq9jze1hQ5YTCLoVxNnbp0T7rglrlOE7N9YsHqjGkEw= +github.com/aws/aws-sdk-go-v2/config v1.33.1/go.mod h1:2A3HQwG4zaL5Tm80rc6RZj8LmWWv4WYT5v8raSz/L7A= +github.com/aws/aws-sdk-go-v2/credentials v1.20.1 h1:Z8GRNEx0u9sDkZOq4PUnN8mjGwbUQGRzMSXpvt3d8xQ= +github.com/aws/aws-sdk-go-v2/credentials v1.20.1/go.mod h1:uBIK00kFo95dnemqfFMTWx0X8YRqsh6ecIoCjjOkZqM= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1 h1:YIEBqcqRnpi4Pfv0YHImtgi6czGCwKHANC7SwmUAVD0= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1/go.mod h1:imEf0oufgAo8KAkCHhrOdqGEC0YWx1PPBQH82shSxGw= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1 h1:pc138gM1CW+XPc60rEwUlwwuwWFQK16CI1T7v1F9Oec= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1/go.mod h1:1+koxpPIbfBdfzP6vojm5/zTpTQ/micYwlxIiNB3TxI= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1 h1:K0JsbZQj+1h208Ro1zHeA4l7bMp0NvRffHQ91q8Ol1s= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1/go.mod h1:W3/vL6EtCIatICGy9ab29QhMuae+cOKPWcMxv02CO+Q= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1 h1:yhw5KD1phVyP9vijxOUzDfEtJx+bt+L63k+VfuiYFAA= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1/go.mod h1:ZW2e0d7DYlRxlS9hEiMXE47gTdX5KRN4byUiNbUpG+Q= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 h1:bAdDl/HkGCcGPoe25ToSHEw23VIxt6CT5fLcg111BKg= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19/go.mod h1:KaUzbLxv4CeSxh6ZCl9B4m7CuFenS8kUEaDs+f/DQr4= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1 h1:s67hBfG5t9rn1NCvDuB4E3QIep3UFhHPtaIqFDjV3N8= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1/go.mod h1:FpvjBMXtSNMLPmDJsWwcY5cRnqJlpS2y1R6n4pvzs4k= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1 h1:RmmWQPREQdk9U+PfqeHW3MqZaBaNK7TpV9W3RY+b+7g= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1/go.mod h1:0A3W4F+68ZnNk5XcNL/e9HFMwnP8RlEicFfy6eOEDyw= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1 h1:ZMbtPZZQRca+3+XYQne9PBvRiYpHZlNJJOZfE9WNfT0= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1/go.mod h1:YAGWQdCYlVCoqrzvfv3RLxO6zKwti7gsAULOGWPLYv4= +github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1 h1:kVpzaDBzOdRtOftmiSpTdQbWVqRg0kONLXijktiwXnk= +github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1/go.mod h1:CUr46sCpGAg/rHaclRyhJX0LJAmH73uWSJPPSaMUrSk= +github.com/aws/aws-sdk-go-v2/service/signin v1.7.1 h1:mdMtSVKdQ3+mzBh+l0ogrFYZVQUCg6pJZOirA2ARsYE= +github.com/aws/aws-sdk-go-v2/service/signin v1.7.1/go.mod h1:9IqUlsJDbUPcg6cgx3WEzXdjrbWzLDQrak0aaSqlTcI= +github.com/aws/aws-sdk-go-v2/service/sso v1.35.1 h1:B6WFn91tobD6gG4724ONHaqrpKsoETGnv98LHe/yIGM= +github.com/aws/aws-sdk-go-v2/service/sso v1.35.1/go.mod h1:tWuiVBUtPBr8/rgRiYS8Uf85sHcAN+G7XS3D3CEoUh8= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1 h1:6yeYCWFvgbI2TI3K6jr9LtBNhXgJ7g4xqD+DEiaDDmM= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1/go.mod h1:naFe83jSMuYkH+QjQPX8n1MLhBkeCFM5Lsnh5m5wz3c= +github.com/aws/aws-sdk-go-v2/service/sts v1.47.1 h1:Sv2xPnRHlThSUtVujYuUBPI/Il8si6UPHXL8DMiB/F0= +github.com/aws/aws-sdk-go-v2/service/sts v1.47.1/go.mod h1:mKo/CzaCz8qytGW70NG4vIIGAx1HXTlb5lHNkC5k3lk= +github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ= +github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= 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= diff --git a/mail/internal/storage/driver.go b/mail/internal/storage/driver.go new file mode 100644 index 0000000..c4827e4 --- /dev/null +++ b/mail/internal/storage/driver.go @@ -0,0 +1,49 @@ +// Package storage implementiert ARC-01: die Objekt-Speicher-Anbindung +// für archivierte Mails und Anhänge. Baut auf demselben bewährten +// Muster wie DMS FDN-03 auf (austauschbare Driver, LocalDriver für +// Entwicklung, S3Driver für Produktion) — Mail kann DMS' internal/ +// nicht importieren (eigenständiges Go-Modul), daher eine bewusste, +// angepasste Neuimplementierung statt eines Cross-Modul-Imports. +// +// Erweiterung gegenüber FDN-03: ARC-01 verlangt Prüfsummenverifikation +// AN DIESER SCHICHT (Akzeptanzkriterium 2), nicht erst an einer +// späteren DB-Schicht — siehe service.go. +package storage + +import ( + "context" + "errors" + "io" + "strconv" +) + +// ErrNotFound wird geliefert, wenn ein angefragtes Objekt nicht +// existiert. +var ErrNotFound = errors.New("storage: objekt nicht gefunden") + +// Driver ist die EINE Schnittstelle, gegen die der Rest von Mail +// arbeitet (Akzeptanzkriterium 1). Zwei Implementierungen: LocalDriver +// (Entwicklung) und S3Driver (Produktion, S3-kompatibel). +type Driver interface { + Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error) + Get(ctx context.Context, key string) (io.ReadCloser, error) + Delete(ctx context.Context, key string) error +} + +// ObjectKey liefert das feste, dokumentierte Pfadschema für einen +// Mail-Anhang/-Teil INNERHALB des bereits mandantenspezifischen +// Buckets (Akzeptanzkriterium 1) — Bucket-Trennung selbst ist Sache +// von Core TEN-01. Lesezugriff hängt NUR von messageID+partIndex ab, +// nicht vom ursprünglichen Importpfad (IMAP/SMTP/manueller Import — +// Akzeptanzkriterium 3): derselbe Key wird unabhängig davon berechnet, +// über welchen Weg die Nachricht ins System kam. +func ObjectKey(messageID string, partIndex int) string { + return "messages/" + messageID + "/parts/" + strconv.Itoa(partIndex) +} + +// checksumKey ist der Sidecar-Objektschlüssel für die beim Schreiben +// berechnete Prüfsumme (siehe service.go) — liegt bewusst im selben +// Driver/Bucket wie der Inhalt, keine separate DB-Abhängigkeit nötig. +func checksumKey(key string) string { + return key + ".sha256" +} diff --git a/mail/internal/storage/localdriver.go b/mail/internal/storage/localdriver.go new file mode 100644 index 0000000..522cd49 --- /dev/null +++ b/mail/internal/storage/localdriver.go @@ -0,0 +1,63 @@ +package storage + +import ( + "context" + "fmt" + "io" + "os" + "path/filepath" +) + +// LocalDriver legt Objekte im lokalen Dateisystem ab — der +// Entwicklungs-Treiber (Akzeptanzkriterium 1), keine externe +// Abhängigkeit nötig. +type LocalDriver struct { + baseDir string +} + +func NewLocalDriver(baseDir string) *LocalDriver { + return &LocalDriver{baseDir: baseDir} +} + +func (d *LocalDriver) path(key string) string { + return filepath.Join(d.baseDir, filepath.FromSlash(key)) +} + +func (d *LocalDriver) Put(_ context.Context, key string, r io.Reader, _ int64, _ string) (int64, error) { + full := d.path(key) + if err := os.MkdirAll(filepath.Dir(full), 0o755); err != nil { + return 0, fmt.Errorf("storage: verzeichnis anlegen: %w", err) + } + f, err := os.Create(full) + if err != nil { + return 0, fmt.Errorf("storage: datei anlegen: %w", err) + } + defer func() { _ = f.Close() }() + + written, err := io.Copy(f, r) + if err != nil { + return 0, fmt.Errorf("storage: schreiben: %w", err) + } + return written, nil +} + +func (d *LocalDriver) Get(_ context.Context, key string) (io.ReadCloser, error) { + f, err := os.Open(d.path(key)) + if err != nil { + if os.IsNotExist(err) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("storage: lesen: %w", err) + } + return f, nil +} + +func (d *LocalDriver) Delete(_ context.Context, key string) error { + if err := os.Remove(d.path(key)); err != nil { + if os.IsNotExist(err) { + return ErrNotFound + } + return fmt.Errorf("storage: löschen: %w", err) + } + return nil +} diff --git a/mail/internal/storage/s3driver.go b/mail/internal/storage/s3driver.go new file mode 100644 index 0000000..d6cde2a --- /dev/null +++ b/mail/internal/storage/s3driver.go @@ -0,0 +1,110 @@ +package storage + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" + "github.com/aws/smithy-go" +) + +// S3Driver legt Objekte in einem S3-kompatiblen Objektspeicher ab — der +// Produktions-Treiber (Akzeptanzkriterium 1). Funktioniert gegen echtes +// AWS S3 UND gegen jeden S3-kompatiblen Anbieter (MinIO etc.) über +// endpointURL. Gleiches, bewährtes Muster wie DMS FDN-03s S3Driver +// (bewusste Kopie, Mail kann DMS nicht importieren). +type S3Driver struct { + client *s3.Client + bucket string +} + +func NewS3Driver(ctx context.Context, bucket, region, endpointURL, accessKeyID, secretAccessKey string, usePathStyle bool) (*S3Driver, error) { + cfg, err := config.LoadDefaultConfig(ctx, + config.WithRegion(region), + config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKeyID, secretAccessKey, "")), + ) + if err != nil { + return nil, fmt.Errorf("storage: s3-konfiguration laden: %w", err) + } + + client := s3.NewFromConfig(cfg, func(o *s3.Options) { + if endpointURL != "" { + o.BaseEndpoint = aws.String(endpointURL) + } + o.UsePathStyle = usePathStyle + }) + return &S3Driver{client: client, bucket: bucket}, nil +} + +func (d *S3Driver) Put(ctx context.Context, key string, r io.Reader, _ int64, contentType string) (int64, error) { + buf, err := io.ReadAll(r) + if err != nil { + return 0, fmt.Errorf("storage: objekt vor upload lesen: %w", err) + } + _, err = d.client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(d.bucket), + Key: aws.String(key), + Body: bytes.NewReader(buf), + ContentLength: aws.Int64(int64(len(buf))), + ContentType: aws.String(contentType), + }) + if err != nil { + return 0, fmt.Errorf("storage: s3-upload: %w", err) + } + return int64(len(buf)), nil +} + +func (d *S3Driver) Get(ctx context.Context, key string) (io.ReadCloser, error) { + out, err := d.client.GetObject(ctx, &s3.GetObjectInput{ + Bucket: aws.String(d.bucket), + Key: aws.String(key), + }) + if err != nil { + if isS3NotFound(err) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("storage: s3-download: %w", err) + } + return out.Body, nil +} + +func (d *S3Driver) Delete(ctx context.Context, key string) error { + // S3 liefert bei DeleteObject fuer ein nicht existierendes Objekt + // KEINEN Fehler (idempotente S3-API-Semantik) — um denselben + // Vertrag wie LocalDriver (ErrNotFound bei fehlendem Objekt) zu + // erfüllen, wird die Existenz vorher explizit geprüft. + _, err := d.client.HeadObject(ctx, &s3.HeadObjectInput{Bucket: aws.String(d.bucket), Key: aws.String(key)}) + if err != nil { + if isS3NotFound(err) { + return ErrNotFound + } + return fmt.Errorf("storage: s3-existenzprüfung vor löschen: %w", err) + } + + if _, err := d.client.DeleteObject(ctx, &s3.DeleteObjectInput{ + Bucket: aws.String(d.bucket), + Key: aws.String(key), + }); err != nil { + return fmt.Errorf("storage: s3-löschen: %w", err) + } + return nil +} + +func isS3NotFound(err error) bool { + var nsk *types.NoSuchKey + if errors.As(err, &nsk) { + return true + } + var apiErr smithy.APIError + if errors.As(err, &apiErr) && apiErr.ErrorCode() == "NotFound" { + return true + } + return false +} diff --git a/mail/internal/storage/service.go b/mail/internal/storage/service.go new file mode 100644 index 0000000..540f349 --- /dev/null +++ b/mail/internal/storage/service.go @@ -0,0 +1,127 @@ +package storage + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io" +) + +// ErrChecksumMismatch wird von GetVerified geliefert, wenn der beim +// Lesen berechnete Hash nicht mit der beim Schreiben gespeicherten +// Prüfsumme übereinstimmt (Akzeptanzkriterium 2 / Pflichtprüfung 2: +// ein absichtlich beschädigtes Objekt wird erkannt). +var ErrChecksumMismatch = errors.New("storage: prüfsumme stimmt nicht überein — objekt wurde verändert") + +// Service verbindet einen Driver mit Prüfsummenverifikation +// (Akzeptanzkriterium 2) und der Nutzungsmeldung an Core LIC-05 +// (Akzeptanzkriterium 4) — jeder Schreib-/Löschvorgang über Service +// löst GENAU EINE Meldung mit der tatsächlich geschriebenen/gelöschten +// Objektgröße aus. Aufrufer (spätere Tickets, z. B. IMP-*) rufen +// ausschließlich Service auf, nie einen Driver direkt. +type Service struct { + driver Driver + usage UsageReporter + tenantSlug string +} + +func NewService(driver Driver, usage UsageReporter, tenantSlug string) *Service { + return &Service{driver: driver, usage: usage, tenantSlug: tenantSlug} +} + +// Put legt den Inhalt ab UND verifiziert den Schreibvorgang durch +// Prüfsummenvergleich (Akzeptanzkriterium 2): der Inhalt wird +// geschrieben, die Prüfsumme als Sidecar-Objekt gespeichert, danach +// SOFORT zurückgelesen und erneut gehasht — weicht der Rückgelesene +// Hash vom beim Schreiben berechneten ab, meldet Put einen Fehler, +// statt eine unbemerkt fehlerhafte Ablage stehen zu lassen. Meldet die +// geschriebene Größe als positives Delta an Core LIC-05 +// (Akzeptanzkriterium 4). +func (s *Service) Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (checksum string, err error) { + hasher := sha256.New() + tee := io.TeeReader(r, hasher) + + written, err := s.driver.Put(ctx, key, tee, size, contentType) + if err != nil { + return "", err + } + checksum = hex.EncodeToString(hasher.Sum(nil)) + + if _, err := s.driver.Put(ctx, checksumKey(key), bytes.NewReader([]byte(checksum)), int64(len(checksum)), "text/plain"); err != nil { + return "", fmt.Errorf("storage: prüfsumme speichern: %w", err) + } + + // Sofortige Rücklese-Verifikation — beweist, dass der Schreibvorgang + // tatsächlich verifiziert wurde, nicht nur eine Prüfsumme abgelegt + // wurde, die nie geprüft wird. + if _, err := s.GetVerified(ctx, key); err != nil { + return "", fmt.Errorf("storage: schreibverifikation fehlgeschlagen: %w", err) + } + + if err := s.usage.Report(ctx, s.tenantSlug, UsageMetric, written); err != nil { + return checksum, fmt.Errorf("storage: objekt gespeichert, aber nutzungsmeldung fehlgeschlagen: %w", err) + } + return checksum, nil +} + +// Get liefert den Inhalt UNVERIFIZIERT (Streaming, für große Objekte). +// Für die Pflichtprüfung "beschädigtes Objekt wird erkannt" GetVerified +// verwenden. +func (s *Service) Get(ctx context.Context, key string) (io.ReadCloser, error) { + return s.driver.Get(ctx, key) +} + +// GetVerified liest den vollständigen Inhalt UND vergleicht die beim +// Schreiben gespeicherte Prüfsumme gegen den beim Lesen berechneten +// Hash (Akzeptanzkriterium 2 / Pflichtprüfung 2). +func (s *Service) GetVerified(ctx context.Context, key string) ([]byte, error) { + sumReader, err := s.driver.Get(ctx, checksumKey(key)) + if err != nil { + return nil, fmt.Errorf("storage: gespeicherte prüfsumme lesen: %w", err) + } + expectedRaw, err := io.ReadAll(sumReader) + _ = sumReader.Close() + if err != nil { + return nil, fmt.Errorf("storage: gespeicherte prüfsumme lesen: %w", err) + } + expected := string(expectedRaw) + + contentReader, err := s.driver.Get(ctx, key) + if err != nil { + return nil, err + } + defer func() { _ = contentReader.Close() }() + + hasher := sha256.New() + content, err := io.ReadAll(io.TeeReader(contentReader, hasher)) + if err != nil { + return nil, fmt.Errorf("storage: objekt lesen: %w", err) + } + actual := hex.EncodeToString(hasher.Sum(nil)) + if actual != expected { + return nil, ErrChecksumMismatch + } + return content, nil +} + +// Delete entfernt Inhalt UND Prüfsummen-Sidecar, meldet die Größe als +// negatives Delta an Core LIC-05 (Akzeptanzkriterium 4) — der Aufrufer +// muss die Größe kennen (Delete selbst kann sie nach dem Löschen nicht +// mehr ermitteln). +func (s *Service) Delete(ctx context.Context, key string, sizeBytes int64) error { + if err := s.driver.Delete(ctx, key); err != nil { + return err + } + // Sidecar-Löschung ist best effort — ein fehlendes Sidecar (z. B. + // bei einem sehr alten Objekt) darf den eigentlichen Löschvorgang + // nicht blockieren. + _ = s.driver.Delete(ctx, checksumKey(key)) + + if err := s.usage.Report(ctx, s.tenantSlug, UsageMetric, -sizeBytes); err != nil { + return fmt.Errorf("storage: objekt gelöscht, aber nutzungsmeldung fehlgeschlagen: %w", err) + } + return nil +} diff --git a/mail/internal/storage/service_test.go b/mail/internal/storage/service_test.go new file mode 100644 index 0000000..8f29090 --- /dev/null +++ b/mail/internal/storage/service_test.go @@ -0,0 +1,134 @@ +package storage + +import ( + "context" + "errors" + "os" + "strings" + "testing" + "time" +) + +type fakeUsageReporter struct { + reports []int64 +} + +func (f *fakeUsageReporter) Report(_ context.Context, _, metric string, delta int64) error { + if metric != UsageMetric { + return errors.New("unerwartete metrik: " + metric) + } + f.reports = append(f.reports, delta) + return nil +} + +func newTestService(t *testing.T) (*Service, *fakeUsageReporter) { + t.Helper() + driver := NewLocalDriver(t.TempDir()) + usage := &fakeUsageReporter{} + return NewService(driver, usage, "acme"), usage +} + +// TestPut_ReadBackIsByteIdentical ist die geforderte Pflichtprüfung 1: +// ein geschriebenes Objekt liefert beim Lesen byteidentischen Inhalt. +func TestPut_ReadBackIsByteIdentical(t *testing.T) { + svc, _ := newTestService(t) + ctx := context.Background() + key := ObjectKey("msg-1", 0) + content := "vollständig identischer Inhalt äöü" + + checksum, err := svc.Put(ctx, key, strings.NewReader(content), int64(len(content)), "text/plain") + if err != nil { + t.Fatalf("put: %v", err) + } + if checksum == "" { + t.Fatal("erwartet nicht-leere prüfsumme") + } + + got, err := svc.GetVerified(ctx, key) + if err != nil { + t.Fatalf("getverified: %v", err) + } + if string(got) != content { + t.Fatalf("nicht byteidentisch: got %q, want %q", got, content) + } +} + +// TestGetVerified_DetectsTamperedObject ist die geforderte +// Pflichtprüfung 2: ein absichtlich beschädigtes Objekt wird bei +// Prüfsummenvergleich erkannt. +func TestGetVerified_DetectsTamperedObject(t *testing.T) { + dir := t.TempDir() + driver := NewLocalDriver(dir) + usage := &fakeUsageReporter{} + svc := NewService(driver, usage, "acme") + ctx := context.Background() + key := ObjectKey("msg-tamper", 0) + + if _, err := svc.Put(ctx, key, strings.NewReader("originaler inhalt"), 17, "text/plain"); err != nil { + t.Fatalf("put: %v", err) + } + + // Objekt DIREKT am Dateisystem manipulieren — umgeht Service + // vollständig, simuliert externe Beschädigung/Manipulation. + full := driver.path(key) + if err := os.WriteFile(full, []byte("MANIPULIERTER INHALT"), 0o644); err != nil { + t.Fatalf("manipulation schreiben: %v", err) + } + + _, err := svc.GetVerified(ctx, key) + if !errors.Is(err, ErrChecksumMismatch) { + t.Fatalf("erwartet ErrChecksumMismatch bei manipuliertem objekt, habe: %v", err) + } +} + +// TestPut_ManySmallObjectsAcceptableLatency ist die geforderte +// Pflichtprüfung 3: Lasttest mit vielen kleinen Objekten bestätigt +// akzeptable Latenz. +func TestPut_ManySmallObjectsAcceptableLatency(t *testing.T) { + svc, _ := newTestService(t) + ctx := context.Background() + + const count = 500 + start := time.Now() + for i := 0; i < count; i++ { + key := ObjectKey("msg-load", i) + if _, err := svc.Put(ctx, key, strings.NewReader("kleiner anhang inhalt"), 21, "text/plain"); err != nil { + t.Fatalf("put #%d: %v", i, err) + } + } + elapsed := time.Since(start) + avgPerObject := elapsed / count + + // Großzügige Grenze (10ms/Objekt inkl. Schreiben+Sidecar+Rücklese- + // Verifikation) — Ziel ist der Nachweis, dass keine quadratische + // oder anderweitig unverhältnismäßige Verschlechterung auftritt, + // nicht ein knallhartes Performance-SLA. + if avgPerObject > 10*time.Millisecond { + t.Fatalf("erwartet akzeptable latenz (<10ms/objekt), habe %v/objekt (gesamt %v für %d objekte)", avgPerObject, elapsed, count) + } + t.Logf("Lasttest: %d Objekte in %v (%v/Objekt)", count, elapsed, avgPerObject) +} + +// TestPut_ReportsUsageOnWriteAndDelete ist die geforderte +// Pflichtprüfung 4: Melde-Aufruf an Core LIC-05 bei Schreib- und +// Löschvorgang nachweislich ausgelöst, mit korrekter Größenangabe. +func TestPut_ReportsUsageOnWriteAndDelete(t *testing.T) { + svc, usage := newTestService(t) + ctx := context.Background() + key := ObjectKey("msg-usage", 0) + content := "zwölf bytes!" + + if _, err := svc.Put(ctx, key, strings.NewReader(content), int64(len(content)), "text/plain"); err != nil { + t.Fatalf("put: %v", err) + } + if len(usage.reports) != 1 || usage.reports[0] != int64(len(content)) { + t.Fatalf("erwartet genau eine positive meldung mit größe %d, habe: %v", len(content), usage.reports) + } + + if err := svc.Delete(ctx, key, int64(len(content))); err != nil { + t.Fatalf("delete: %v", err) + } + if len(usage.reports) != 2 || usage.reports[1] != -int64(len(content)) { + t.Fatalf("erwartet zusätzliche negative meldung mit -%d, habe: %v", len(content), usage.reports) + } +} diff --git a/mail/internal/storage/usagereport.go b/mail/internal/storage/usagereport.go new file mode 100644 index 0000000..cecceb2 --- /dev/null +++ b/mail/internal/storage/usagereport.go @@ -0,0 +1,75 @@ +package storage + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" +) + +// UsageMetric ist der Metrikname, unter dem Core (internal/usage, +// LIC-05) den Speicherverbrauch je Mandant führt — muss exakt +// internal/usage.StorageBytesMetric aus dem NEXARCH-Core-Modul +// entsprechen (Mail kann Core nicht importieren, daher hier gespiegelt +// — identisches Muster wie DMS FDN-03). +const UsageMetric = "storage_bytes" + +// UsageReporter meldet Speicherverbrauchsänderungen an Core +// (Akzeptanzkriterium 4). Schmale Schnittstelle, damit Tests einen +// Fake statt eines echten HTTP-Aufrufs einsetzen können. +type UsageReporter interface { + Report(ctx context.Context, tenantSlug, metric string, delta int64) error +} + +// usageDeltaDTO entspricht Core internal/resync.usageDeltaDTO +// (JSON-Vertrag: tenant_slug/metric/delta), über den API-11 +// (resync-api) real erreichbar ist. +type usageDeltaDTO struct { + TenantSlug string `json:"tenant_slug"` + Metric string `json:"metric"` + Delta int64 `json:"delta"` +} + +// HTTPUsageReporter meldet über Core API-11 (resync-api, +// internal/resync.Handler.UsageHandler), authentifiziert über +// dasselbe Service-Credential-Verfahren wie jeder andere Modul-Core- +// Aufruf (API-02). +type HTTPUsageReporter struct { + endpointURL string + clientID string + clientSecret string + httpClient *http.Client +} + +func NewHTTPUsageReporter(endpointURL, clientID, clientSecret string, httpClient *http.Client) *HTTPUsageReporter { + if httpClient == nil { + httpClient = http.DefaultClient + } + return &HTTPUsageReporter{endpointURL: endpointURL, clientID: clientID, clientSecret: clientSecret, httpClient: httpClient} +} + +func (r *HTTPUsageReporter) Report(ctx context.Context, tenantSlug, metric string, delta int64) error { + body, err := json.Marshal([]usageDeltaDTO{{TenantSlug: tenantSlug, Metric: metric, Delta: delta}}) + if err != nil { + return fmt.Errorf("storage: nutzungsmeldung serialisieren: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.endpointURL, bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("storage: nutzungsmeldungs-anfrage aufbauen: %w", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("X-Nexarch-Client-Id", r.clientID) + req.Header.Set("X-Nexarch-Client-Secret", r.clientSecret) + + resp, err := r.httpClient.Do(req) + if err != nil { + return fmt.Errorf("storage: nutzungsmeldung senden: %w", err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("storage: nutzungsmeldung von core abgelehnt: status %d", resp.StatusCode) + } + return nil +}