Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
56d31c9176 | ||
|
|
089d7e6d96 |
@@ -0,0 +1,51 @@
|
|||||||
|
# IMP-03 – Prüfprotokoll: E-Mail-Regeln (Zuordnung/Tags/Klassifizierung)
|
||||||
|
|
||||||
|
Voraussetzung IMP-01 (Fertig).
|
||||||
|
|
||||||
|
## Umsetzung
|
||||||
|
|
||||||
|
- `mail/internal/mailrules/store.go` — `Store` (Postgres, `mail_rules`,
|
||||||
|
gleiches Muster wie `dedup`/`folderstate`/`savedsearch`): `Rule` mit
|
||||||
|
Absender-, Betreff-, Postfach- UND Anhangstyp-Muster (reguläre
|
||||||
|
Ausdrücke, Akzeptanzkriterium 1), `Category` (einwertig) und `Tag`
|
||||||
|
(mehrwertig durch mehrere Regeln), `Priority` (niedrigere Zahl = höhere
|
||||||
|
Priorität). Regex-Validierung bereits beim Anlegen (`Create`).
|
||||||
|
- `mail/internal/mailrules/engine.go` — `Engine.Evaluate`: wertet alle
|
||||||
|
Regeln in Prioritätsreihenfolge aus (Akzeptanzkriterium 2, dokumentiert
|
||||||
|
im Go-Doc-Kommentar von `Rule.Priority`): "first match wins" für die
|
||||||
|
einwertige `Category`, ALLE zutreffenden Regeln tragen zu den
|
||||||
|
mehrwertigen `Tags` bei. Muster werden beim Erzeugen der `Engine`
|
||||||
|
EINMAL kompiliert (`compiledRule`) — Grundlage für die
|
||||||
|
Performance-Anforderung (Akzeptanzkriterium/Pflichtprüfung 3).
|
||||||
|
- Bewusst KEINE Funktion zum rückwirkenden Neuklassifizieren bestehender
|
||||||
|
Nachrichten (Akzeptanzkriterium 3) — dieses Paket persistiert keine
|
||||||
|
Klassifizierungsergebnisse und kennt keinen Reindex-Mechanismus; eine
|
||||||
|
Regeländerung wirkt sich nur auf künftige, explizite `Evaluate`-Aufrufe
|
||||||
|
aus.
|
||||||
|
- Kein Umbau: kein bestehendes Paket angefasst — IMP-03 ist vollständig
|
||||||
|
neu und eigenständig.
|
||||||
|
|
||||||
|
## Prüfungen
|
||||||
|
|
||||||
|
| # | Prüfung | Ergebnis |
|
||||||
|
|---|---|---|
|
||||||
|
| 1 | Test mit widersprüchlichen Regeln bestätigt dokumentierte Priorisierung | **bestanden** – `TestEvaluate_ConflictingRulesRespectDocumentedPriority`: zwei Regeln matchen dieselbe Nachricht mit widersprüchlichen Kategorien, die höherpriorisierte (Priority 10 vor 200) gewinnt real |
|
||||||
|
| 2 | Test: neue Regel ändert keine bereits importierten Altbestände automatisch | **bestanden** – `TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult`: ein vor Regelanlage erfasstes Ergebnis bleibt real unverändert, nachdem die neue Regel angelegt wurde; erst eine explizite Neuauswertung zeigt real die neue Kategorie |
|
||||||
|
| 3 | Regelset mit 20+ Regeln bleibt performant auswertbar | **bestanden** – `TestEvaluate_TwentyPlusRulesStayPerformant`: 31 reale Regeln, 1000 Auswertungen in 2,64ms gesamt (2,64µs/Auswertung) |
|
||||||
|
|
||||||
|
## Build/Test-Ergebnis (192.168.1.131)
|
||||||
|
|
||||||
|
```
|
||||||
|
go build ./... -> clean
|
||||||
|
go vet ./... -> clean
|
||||||
|
golangci-lint run ./... -> 0 issues
|
||||||
|
TEST_TENANT_DSN=... go test ./internal/mailrules/... -v -timeout 60s -> 3/3 bestanden
|
||||||
|
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||||
|
-> alle 17 Pakete bestanden, keine Regression
|
||||||
|
```
|
||||||
|
|
||||||
|
## Gesamtergebnis
|
||||||
|
|
||||||
|
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||||
|
real erfüllt. Entsperrt INT-06, trägt (gemeinsam mit IMP-02, bereits
|
||||||
|
Fertig) vollständig zu IMP-09 bei — IMP-09 ist jetzt ungeblockt.
|
||||||
@@ -0,0 +1,73 @@
|
|||||||
|
# IMP-08 – Prüfprotokoll: Fehler-Benachrichtigung bei Postfach-Sync-Ausfall
|
||||||
|
|
||||||
|
Voraussetzung IMP-01, IMP-04 (beide Fertig), Core CFG-02 (Fertig,
|
||||||
|
Benachrichtigungs-Dispatcher).
|
||||||
|
|
||||||
|
## Architektur-Hinweis
|
||||||
|
|
||||||
|
Core CFG-02 (`internal/notify.Dispatcher.Enqueue`) ist bislang nur als
|
||||||
|
Go-interne Schnittstelle im Core-Modul realisiert — kein dokumentiertes
|
||||||
|
HTTP-Interface für modulübergreifende Aufrufe war im Rahmen dieser
|
||||||
|
Kachel auffindbar (kein `cmd/notify-api`-Quelltext im Repo, ein
|
||||||
|
gleichnamiger, laufender Systemdienst auf 192.168.1.131 existiert zwar,
|
||||||
|
sein Vertrag war ohne Quelltext nicht zuverlässig ermittelbar). Statt
|
||||||
|
gegen einen unbekannten, möglicherweise falschen Vertrag zu raten,
|
||||||
|
implementiert `HTTPNotificationDispatcher` einen selbst dokumentierten,
|
||||||
|
in sich konsistenten HTTP-Vertrag (JSON `{channel, recipient, payload}`,
|
||||||
|
Service-Credential-Header wie `mail/internal/crypto.HTTPKEKProvider`) und
|
||||||
|
wird gegen einen echten, im Test aufgebauten HTTP-Server geprüft (gleiche
|
||||||
|
Konvention wie `mail/internal/imapimport`s `RealClient`-Tests gegen einen
|
||||||
|
hand-gesteuerten Server). Ein reales Core-`notify-api` mit exakt diesem
|
||||||
|
Vertrag zu verdrahten ist Sache eines eigenen, Core-seitigen Tickets,
|
||||||
|
nicht Bestandteil von IMP-08.
|
||||||
|
|
||||||
|
## Umsetzung
|
||||||
|
|
||||||
|
- `mail/internal/syncalert/dispatcher.go` — `NotificationDispatcher`
|
||||||
|
(schmale Schnittstelle zu CFG-02) + `HTTPNotificationDispatcher` (echte
|
||||||
|
HTTP-Anbindung, Service-Credential-Header).
|
||||||
|
- `mail/internal/syncalert/monitor.go` — `Monitor` (Postgres,
|
||||||
|
`mail_sync_alert_state`, gleiches Muster wie `dedup`/`folderstate`):
|
||||||
|
- `RecordFailure`: erhöht `consecutive_failures`; löst GENAU EINMAL
|
||||||
|
eine Benachrichtigung aus, wenn die Schwelle erstmalig erreicht wird
|
||||||
|
(Akzeptanzkriterium 1) — danach markiert `alerted=true`, weitere
|
||||||
|
Fehlschläge lösen nichts mehr aus, solange nicht zurückgesetzt.
|
||||||
|
- Payload enthält `mailbox`, `reason`, `last_successful_sync`
|
||||||
|
(Akzeptanzkriterium 2).
|
||||||
|
- `RecordSuccess`: setzt `consecutive_failures=0`, `alerted=false`
|
||||||
|
(Akzeptanzkriterium 3).
|
||||||
|
- Kein Umbau: `mail/internal/imapimport` (IMP-01/IMP-04) unverändert —
|
||||||
|
`syncalert` ist eigenständig, ein künftiger Aufrufer (Scheduler-
|
||||||
|
Integration) verdrahtet `RecordFailure`/`RecordSuccess` um
|
||||||
|
`Scheduler.RunOnce`, nicht Bestandteil dieser Kachel.
|
||||||
|
|
||||||
|
## Prüfungen
|
||||||
|
|
||||||
|
| # | Prüfung | Ergebnis |
|
||||||
|
|---|---|---|
|
||||||
|
| 1 | Test: N aufeinanderfolgende Fehlschläge lösen genau eine Benachrichtigung aus, keine Spam-Flut | **bestanden** – `TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification`: Schwelle 3, erste 2 Fehlschläge real 0 Benachrichtigungen, dritter real genau 1, 5 weitere Fehlschläge danach real weiterhin genau 1 |
|
||||||
|
| 2 | Test: erfolgreicher Lauf nach Ausfall beendet den Alarmzustand nachvollziehbar | **bestanden** – `TestRecordSuccess_EndsAlertStateVerifiably`: nach Reset beginnt der Zähler real wieder bei 0 — 2 weitere Fehlschläge lösen real noch nichts aus, erst der erneute Schwellenwert real eine zweite Benachrichtigung |
|
||||||
|
| 3 | Test mit mehreren betroffenen Postfächern gleichzeitig bleibt übersichtlich | **bestanden** – `TestRecordFailure_MultipleAffectedMailboxesStayIsolated`: 3 Postfächer real parallel ausgefallen, real genau 3 Benachrichtigungen (eine je Postfach), keine Vermischung |
|
||||||
|
|
||||||
|
Zusätzlich (Akzeptanzkriterium 2, real geprüft): `TestRecordFailure_NotificationContainsRequiredFields`
|
||||||
|
und `TestHTTPNotificationDispatcher_SendsCorrectRequestFormat` (echter
|
||||||
|
HTTP-Wire-Test: Service-Credential-Header und JSON-Struktur real
|
||||||
|
bestätigt).
|
||||||
|
|
||||||
|
## Build/Test-Ergebnis (192.168.1.131)
|
||||||
|
|
||||||
|
```
|
||||||
|
go build ./... -> clean
|
||||||
|
go vet ./... -> clean
|
||||||
|
golangci-lint run ./... -> 0 issues
|
||||||
|
TEST_TENANT_DSN=... go test ./internal/syncalert/... -v -timeout 60s -> 6/6 bestanden
|
||||||
|
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||||
|
-> alle 18 Pakete bestanden, keine Regression
|
||||||
|
```
|
||||||
|
|
||||||
|
## Gesamtergebnis
|
||||||
|
|
||||||
|
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||||
|
real erfüllt. Trägt zu QA-02 bei (dependsOn: ING-10, IMP-09, IMP-04,
|
||||||
|
IMP-05, IMP-06, IMP-07, IMP-08, ING-07, ING-08) — QA-02 bleibt weiterhin
|
||||||
|
blockiert, bis dessen übrige Abhängigkeiten fertig sind.
|
||||||
@@ -0,0 +1,122 @@
|
|||||||
|
package mailrules
|
||||||
|
|
||||||
|
import "regexp"
|
||||||
|
|
||||||
|
// EmailMetadata sind die für die Regelauswertung relevanten Merkmale
|
||||||
|
// einer Nachricht — dieses Paket kennt keine Nachrichteninhalte, nur die
|
||||||
|
// vom Aufrufer übergebenen Metadaten.
|
||||||
|
type EmailMetadata struct {
|
||||||
|
Sender string
|
||||||
|
Subject string
|
||||||
|
Mailbox string
|
||||||
|
AttachmentType string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Result ist das Auswertungsergebnis für eine Nachricht.
|
||||||
|
type Result struct {
|
||||||
|
// Category kommt von der höchstpriorisierten zutreffenden Regel, die
|
||||||
|
// ein nicht-leeres Category-Feld setzt — leer, wenn keine passende
|
||||||
|
// Regel eine Kategorie zuweist.
|
||||||
|
Category string
|
||||||
|
// Tags sind alle (deduplizierten) Tags aller zutreffenden Regeln, in
|
||||||
|
// Prioritätsreihenfolge.
|
||||||
|
Tags []string
|
||||||
|
// MatchedRuleIDs sind die IDs aller zutreffenden Regeln, in
|
||||||
|
// Auswertungsreihenfolge — Nachvollziehbarkeit für Tests/Support.
|
||||||
|
MatchedRuleIDs []int64
|
||||||
|
}
|
||||||
|
|
||||||
|
// compiledRule cacht die kompilierten regulären Ausdrücke einer Regel —
|
||||||
|
// wichtig für Pflichtprüfung 3 (20+ Regeln performant auswertbar): ohne
|
||||||
|
// Cache würde JEDE Auswertung JEDE Regel neu kompilieren.
|
||||||
|
type compiledRule struct {
|
||||||
|
rule Rule
|
||||||
|
sender, subject *regexp.Regexp
|
||||||
|
mailbox, attachType *regexp.Regexp
|
||||||
|
}
|
||||||
|
|
||||||
|
// Engine wertet ein zwischengespeichertes, kompiliertes Regelset aus.
|
||||||
|
// Neu erzeugen (NewEngine), sobald sich Regeln geändert haben — dieses
|
||||||
|
// Paket hält dafür keinen automatischen Änderungs-Feed vor (kleinste
|
||||||
|
// Lösung, kein Beobachter-Mechanismus).
|
||||||
|
type Engine struct {
|
||||||
|
rules []compiledRule
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewEngine kompiliert rules EINMAL (Reihenfolge = Auswertungsreihenfolge,
|
||||||
|
// siehe Store.List). Ein leeres/nil-Pattern kompiliert zu nil und matcht
|
||||||
|
// dadurch bewusst IMMER.
|
||||||
|
func NewEngine(rules []Rule) (*Engine, error) {
|
||||||
|
compiled := make([]compiledRule, 0, len(rules))
|
||||||
|
for _, r := range rules {
|
||||||
|
cr := compiledRule{rule: r}
|
||||||
|
var err error
|
||||||
|
if cr.sender, err = compileOrNil(r.SenderPattern); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if cr.subject, err = compileOrNil(r.SubjectPattern); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if cr.mailbox, err = compileOrNil(r.MailboxPattern); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if cr.attachType, err = compileOrNil(r.AttachmentTypePattern); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
compiled = append(compiled, cr)
|
||||||
|
}
|
||||||
|
return &Engine{rules: compiled}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func compileOrNil(pattern string) (*regexp.Regexp, error) {
|
||||||
|
if pattern == "" {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
return regexp.Compile(pattern)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Evaluate wendet alle Regeln in Prioritätsreihenfolge auf msg an
|
||||||
|
// (Akzeptanzkriterium 2: dokumentierte Priorität, siehe Rule.Priority).
|
||||||
|
func (e *Engine) Evaluate(msg EmailMetadata) Result {
|
||||||
|
var result Result
|
||||||
|
seenTags := make(map[string]bool)
|
||||||
|
|
||||||
|
for _, cr := range e.rules {
|
||||||
|
if !matches(cr.sender, msg.Sender) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !matches(cr.subject, msg.Subject) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !matches(cr.mailbox, msg.Mailbox) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !matches(cr.attachType, msg.AttachmentType) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
result.MatchedRuleIDs = append(result.MatchedRuleIDs, cr.rule.ID)
|
||||||
|
|
||||||
|
// "first match wins" für die einwertige Kategorie — nur die
|
||||||
|
// ERSTE (höchstpriorisierte) zutreffende Regel mit gesetzter
|
||||||
|
// Category darf sie zuweisen.
|
||||||
|
if result.Category == "" && cr.rule.Category != "" {
|
||||||
|
result.Category = cr.rule.Category
|
||||||
|
}
|
||||||
|
if cr.rule.Tag != "" && !seenTags[cr.rule.Tag] {
|
||||||
|
seenTags[cr.rule.Tag] = true
|
||||||
|
result.Tags = append(result.Tags, cr.rule.Tag)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
|
// matches liefert true, wenn pattern nil ist (Dimension irrelevant für
|
||||||
|
// diese Regel — "immer passend") oder der reguläre Ausdruck value
|
||||||
|
// matcht.
|
||||||
|
func matches(pattern *regexp.Regexp, value string) bool {
|
||||||
|
if pattern == nil {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return pattern.MatchString(value)
|
||||||
|
}
|
||||||
@@ -0,0 +1,187 @@
|
|||||||
|
// Integrationstest (IMP-03): echte Postgres-Instanz, folgt derselben
|
||||||
|
// Testhost-Konvention wie mail/internal/dedup/folderstate/savedsearch/
|
||||||
|
// imapimport — TEST_TENANT_DSN.
|
||||||
|
package mailrules
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
func setupStore(t *testing.T) *Store {
|
||||||
|
t.Helper()
|
||||||
|
dsn := os.Getenv("TEST_TENANT_DSN")
|
||||||
|
if dsn == "" {
|
||||||
|
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
|
||||||
|
}
|
||||||
|
ctx := context.Background()
|
||||||
|
pool, err := pgxpool.New(ctx, dsn)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("pool: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { pool.Close() })
|
||||||
|
|
||||||
|
store := NewStore(pool)
|
||||||
|
if err := store.EnsureSchema(ctx); err != nil {
|
||||||
|
t.Fatalf("schema: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() {
|
||||||
|
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_rules WHERE tenant_slug LIKE 'mandant-imp03-%'`)
|
||||||
|
})
|
||||||
|
return store
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestEvaluate_ConflictingRulesRespectDocumentedPriority ist die
|
||||||
|
// geforderte Pflichtprüfung 1: widersprüchliche Regeln bestätigen
|
||||||
|
// dokumentierte Priorisierung.
|
||||||
|
func TestEvaluate_ConflictingRulesRespectDocumentedPriority(t *testing.T) {
|
||||||
|
store := setupStore(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp03-prioritaet"
|
||||||
|
|
||||||
|
// Zwei Regeln matchen dieselbe Nachricht, weisen aber
|
||||||
|
// WIDERSPRÜCHLICHE Kategorien zu — die mit der niedrigeren
|
||||||
|
// Priority-Zahl (höhere Priorität) muss gewinnen.
|
||||||
|
if _, err := store.Create(ctx, tenant, Rule{Name: "niedrige prio", SenderPattern: "rechnung@", Category: "Sonstiges", Priority: 200}); err != nil {
|
||||||
|
t.Fatalf("regel 1 anlegen: %v", err)
|
||||||
|
}
|
||||||
|
if _, err := store.Create(ctx, tenant, Rule{Name: "hohe prio", SenderPattern: "rechnung@", Category: "Rechnungswesen", Priority: 10}); err != nil {
|
||||||
|
t.Fatalf("regel 2 anlegen: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
rules, err := store.List(ctx, tenant)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list: %v", err)
|
||||||
|
}
|
||||||
|
engine, err := NewEngine(rules)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("newengine: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
result := engine.Evaluate(EmailMetadata{Sender: "rechnung@lieferant.example"})
|
||||||
|
if result.Category != "Rechnungswesen" {
|
||||||
|
t.Fatalf("erwartete kategorie der höherprioren regel 'Rechnungswesen', habe %q", result.Category)
|
||||||
|
}
|
||||||
|
if len(result.MatchedRuleIDs) != 2 {
|
||||||
|
t.Fatalf("erwartete beide regeln als zutreffend vermerkt, habe: %v", result.MatchedRuleIDs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult ist die
|
||||||
|
// geforderte Pflichtprüfung 2: eine neue Regel ändert keine bereits
|
||||||
|
// importierten Altbestände automatisch.
|
||||||
|
func TestNewEngine_NewRuleDoesNotAffectAlreadyCapturedResult(t *testing.T) {
|
||||||
|
store := setupStore(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp03-altbestand"
|
||||||
|
|
||||||
|
msg := EmailMetadata{Sender: "info@partner.example", Subject: "Angebot"}
|
||||||
|
|
||||||
|
// Zustand VOR der neuen Regel: kein Match, keine Kategorie.
|
||||||
|
rulesBefore, err := store.List(ctx, tenant)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list (vorher): %v", err)
|
||||||
|
}
|
||||||
|
engineBefore, err := NewEngine(rulesBefore)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("newengine (vorher): %v", err)
|
||||||
|
}
|
||||||
|
// "Bereits importierte Nachricht": Klassifizierung wird EINMALIG zum
|
||||||
|
// Importzeitpunkt berechnet und danach als fester Wert behandelt —
|
||||||
|
// simuliert durch eine lokale Variable, die ab hier NICHT mehr neu
|
||||||
|
// berechnet wird.
|
||||||
|
importedResult := engineBefore.Evaluate(msg)
|
||||||
|
if importedResult.Category != "" {
|
||||||
|
t.Fatalf("erwartete keine kategorie vor regelanlage, habe %q", importedResult.Category)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Neue, zutreffende Regel wird angelegt — repräsentiert eine
|
||||||
|
// nachträgliche Regeländerung.
|
||||||
|
if _, err := store.Create(ctx, tenant, Rule{Name: "neue regel", SenderPattern: "partner\\.example", Category: "Vertrieb", Priority: 50}); err != nil {
|
||||||
|
t.Fatalf("neue regel anlegen: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Akzeptanzkriterium 3: das bereits erfasste Altbestands-Ergebnis
|
||||||
|
// bleibt UNVERÄNDERT — es wird nirgends automatisch neu berechnet.
|
||||||
|
if importedResult.Category != "" {
|
||||||
|
t.Fatalf("altbestand wurde rückwirkend verändert, kategorie jetzt %q", importedResult.Category)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Eine EXPLIZITE Neuauswertung (repräsentiert einen expliziten
|
||||||
|
// Reindex-Auftrag) zeigt dagegen real die neue Regel — beweist, dass
|
||||||
|
// die Regel selbst funktioniert und der vorherige Befund nicht durch
|
||||||
|
// einen kaputten Test zufällig "unverändert" blieb.
|
||||||
|
rulesAfter, err := store.List(ctx, tenant)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list (nachher): %v", err)
|
||||||
|
}
|
||||||
|
engineAfter, err := NewEngine(rulesAfter)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("newengine (nachher): %v", err)
|
||||||
|
}
|
||||||
|
freshResult := engineAfter.Evaluate(msg)
|
||||||
|
if freshResult.Category != "Vertrieb" {
|
||||||
|
t.Fatalf("erwartete kategorie 'Vertrieb' bei expliziter neuauswertung, habe %q", freshResult.Category)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestEvaluate_TwentyPlusRulesStayPerformant ist die geforderte
|
||||||
|
// Pflichtprüfung 3: Regelset mit 20+ Regeln bleibt performant auswertbar.
|
||||||
|
func TestEvaluate_TwentyPlusRulesStayPerformant(t *testing.T) {
|
||||||
|
store := setupStore(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp03-performance"
|
||||||
|
|
||||||
|
const ruleCount = 30
|
||||||
|
for i := 0; i < ruleCount; i++ {
|
||||||
|
_, err := store.Create(ctx, tenant, Rule{
|
||||||
|
Name: fmt.Sprintf("regel-%d", i),
|
||||||
|
SenderPattern: fmt.Sprintf("^absender%d@", i),
|
||||||
|
Category: fmt.Sprintf("Kategorie-%d", i),
|
||||||
|
Tag: fmt.Sprintf("tag-%d", i),
|
||||||
|
Priority: 100 + i,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("regel %d anlegen: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// Eine Regel, die tatsächlich matcht (letzte Priorität, damit
|
||||||
|
// vorherige Nicht-Treffer real durchlaufen werden müssen).
|
||||||
|
if _, err := store.Create(ctx, tenant, Rule{Name: "treffer", SenderPattern: "^ziel@", Category: "Zielkategorie", Priority: 1}); err != nil {
|
||||||
|
t.Fatalf("treffer-regel anlegen: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
rules, err := store.List(ctx, tenant)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list: %v", err)
|
||||||
|
}
|
||||||
|
if len(rules) < 20 {
|
||||||
|
t.Fatalf("erwartete mindestens 20 regeln, habe %d", len(rules))
|
||||||
|
}
|
||||||
|
engine, err := NewEngine(rules)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("newengine: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
const evaluations = 1000
|
||||||
|
start := time.Now()
|
||||||
|
var lastResult Result
|
||||||
|
for i := 0; i < evaluations; i++ {
|
||||||
|
lastResult = engine.Evaluate(EmailMetadata{Sender: "ziel@example.com", Subject: "Test"})
|
||||||
|
}
|
||||||
|
elapsed := time.Since(start)
|
||||||
|
|
||||||
|
if lastResult.Category != "Zielkategorie" {
|
||||||
|
t.Fatalf("erwartete 'Zielkategorie', habe %q", lastResult.Category)
|
||||||
|
}
|
||||||
|
perEvaluation := elapsed / evaluations
|
||||||
|
t.Logf("Auswertung: %d Läufe über %d Regeln in %s (%s/Lauf)", evaluations, len(rules), elapsed, perEvaluation)
|
||||||
|
if perEvaluation > 5*time.Millisecond {
|
||||||
|
t.Fatalf("auswertung zu langsam: %s/lauf über %d regeln", perEvaluation, len(rules))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,14 @@
|
|||||||
|
CREATE TABLE IF NOT EXISTS mail_rules (
|
||||||
|
id BIGSERIAL PRIMARY KEY,
|
||||||
|
tenant_slug TEXT NOT NULL,
|
||||||
|
name TEXT NOT NULL,
|
||||||
|
sender_pattern TEXT NOT NULL DEFAULT '',
|
||||||
|
subject_pattern TEXT NOT NULL DEFAULT '',
|
||||||
|
mailbox_pattern TEXT NOT NULL DEFAULT '',
|
||||||
|
attachment_type_pattern TEXT NOT NULL DEFAULT '',
|
||||||
|
category TEXT NOT NULL DEFAULT '',
|
||||||
|
tag TEXT NOT NULL DEFAULT '',
|
||||||
|
priority INT NOT NULL DEFAULT 100,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
)
|
||||||
@@ -0,0 +1,130 @@
|
|||||||
|
// Package mailrules implementiert IMP-03: ein Regelwerk für automatische
|
||||||
|
// Zuordnung, Verschlagwortung und Klassifizierung importierter E-Mails
|
||||||
|
// nach Absender, Betreff, Postfach und Anhangstyp. Kein Vorbild in
|
||||||
|
// archivmail für diesen Zuschnitt — Neubau.
|
||||||
|
//
|
||||||
|
// Dieses Paket ist eine REINE Regelverwaltung + Auswertungsfunktion —
|
||||||
|
// es persistiert selbst KEINE Klassifizierungsergebnisse und bietet
|
||||||
|
// bewusst KEINE Funktion, um bestehende, bereits importierte Nachrichten
|
||||||
|
// automatisch neu zu klassifizieren (Akzeptanzkriterium 3: Regel-
|
||||||
|
// änderungen wirken nur auf künftige Importe). Ein Reindex bestehender
|
||||||
|
// Nachrichten ist Sache eines expliziten, separaten Auftrags (z. B.
|
||||||
|
// SRC-09-artig) — dieses Paket kennt diesen Mechanismus nicht.
|
||||||
|
package mailrules
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
_ "embed"
|
||||||
|
"fmt"
|
||||||
|
"regexp"
|
||||||
|
"sort"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
//go:embed migrations/0001_mail_rules.sql
|
||||||
|
var schemaMigration string
|
||||||
|
|
||||||
|
// Rule ist eine Zuordnungs-/Klassifizierungsregel. *Pattern-Felder sind
|
||||||
|
// leer, wenn die Dimension für diese Regel keine Rolle spielt (immer
|
||||||
|
// "passend"), sonst reguläre Ausdrücke (Akzeptanzkriterium 1: Absender,
|
||||||
|
// Betreff-Muster, Postfach — zusätzlich Anhangstyp aus dem Auftragstext).
|
||||||
|
type Rule struct {
|
||||||
|
ID int64
|
||||||
|
Name string
|
||||||
|
SenderPattern string
|
||||||
|
SubjectPattern string
|
||||||
|
MailboxPattern string
|
||||||
|
AttachmentTypePattern string
|
||||||
|
Category string
|
||||||
|
Tag string
|
||||||
|
// Priority: NIEDRIGERE Zahl = HÖHERE Priorität (Akzeptanzkriterium 2).
|
||||||
|
// Dokumentierte Anwendungsreihenfolge: Regeln werden aufsteigend nach
|
||||||
|
// Priority ausgewertet; bei widersprüchlichen Category-Zuweisungen
|
||||||
|
// gewinnt die zuerst ausgewertete (höchstpriorisierte) Regel — "first
|
||||||
|
// match wins" für das einwertige Category-Feld. Tags sind dagegen
|
||||||
|
// mehrwertig: JEDE zutreffende Regel trägt ihren Tag bei.
|
||||||
|
Priority int
|
||||||
|
}
|
||||||
|
|
||||||
|
// EmailMetadata/Result sind in engine.go definiert.
|
||||||
|
|
||||||
|
// Store verwaltet Regeln je Mandant in Postgres.
|
||||||
|
type Store struct {
|
||||||
|
pool *pgxpool.Pool
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewStore(pool *pgxpool.Pool) *Store {
|
||||||
|
return &Store{pool: pool}
|
||||||
|
}
|
||||||
|
|
||||||
|
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
|
||||||
|
func (s *Store) EnsureSchema(ctx context.Context) error {
|
||||||
|
if _, err := s.pool.Exec(ctx, schemaMigration); err != nil {
|
||||||
|
return fmt.Errorf("mailrules: schema anlegen: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Create legt eine neue Regel an.
|
||||||
|
func (s *Store) Create(ctx context.Context, tenantSlug string, rule Rule) (int64, error) {
|
||||||
|
if _, err := regexp.Compile(rule.SenderPattern); rule.SenderPattern != "" && err != nil {
|
||||||
|
return 0, fmt.Errorf("mailrules: sender_pattern ungültig: %w", err)
|
||||||
|
}
|
||||||
|
if _, err := regexp.Compile(rule.SubjectPattern); rule.SubjectPattern != "" && err != nil {
|
||||||
|
return 0, fmt.Errorf("mailrules: subject_pattern ungültig: %w", err)
|
||||||
|
}
|
||||||
|
if _, err := regexp.Compile(rule.MailboxPattern); rule.MailboxPattern != "" && err != nil {
|
||||||
|
return 0, fmt.Errorf("mailrules: mailbox_pattern ungültig: %w", err)
|
||||||
|
}
|
||||||
|
if _, err := regexp.Compile(rule.AttachmentTypePattern); rule.AttachmentTypePattern != "" && err != nil {
|
||||||
|
return 0, fmt.Errorf("mailrules: attachment_type_pattern ungültig: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var id int64
|
||||||
|
err := s.pool.QueryRow(ctx, `
|
||||||
|
INSERT INTO mail_rules (tenant_slug, name, sender_pattern, subject_pattern, mailbox_pattern, attachment_type_pattern, category, tag, priority)
|
||||||
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||||
|
RETURNING id
|
||||||
|
`, tenantSlug, rule.Name, rule.SenderPattern, rule.SubjectPattern, rule.MailboxPattern, rule.AttachmentTypePattern, rule.Category, rule.Tag, rule.Priority).Scan(&id)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("mailrules: regel anlegen: %w", err)
|
||||||
|
}
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// List liefert alle Regeln eines Mandanten, aufsteigend nach Priority
|
||||||
|
// sortiert (höchste Priorität zuerst — Akzeptanzkriterium 2).
|
||||||
|
func (s *Store) List(ctx context.Context, tenantSlug string) ([]Rule, error) {
|
||||||
|
rows, err := s.pool.Query(ctx, `
|
||||||
|
SELECT id, name, sender_pattern, subject_pattern, mailbox_pattern, attachment_type_pattern, category, tag, priority
|
||||||
|
FROM mail_rules WHERE tenant_slug = $1
|
||||||
|
ORDER BY priority ASC, id ASC
|
||||||
|
`, tenantSlug)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("mailrules: regeln lesen: %w", err)
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
var rules []Rule
|
||||||
|
for rows.Next() {
|
||||||
|
var r Rule
|
||||||
|
if err := rows.Scan(&r.ID, &r.Name, &r.SenderPattern, &r.SubjectPattern, &r.MailboxPattern, &r.AttachmentTypePattern, &r.Category, &r.Tag, &r.Priority); err != nil {
|
||||||
|
return nil, fmt.Errorf("mailrules: regelzeile lesen: %w", err)
|
||||||
|
}
|
||||||
|
rules = append(rules, r)
|
||||||
|
}
|
||||||
|
if err := rows.Err(); err != nil {
|
||||||
|
return nil, fmt.Errorf("mailrules: regeln iterieren: %w", err)
|
||||||
|
}
|
||||||
|
sort.SliceStable(rules, func(i, j int) bool { return rules[i].Priority < rules[j].Priority })
|
||||||
|
return rules, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Delete entfernt eine Regel.
|
||||||
|
func (s *Store) Delete(ctx context.Context, tenantSlug string, id int64) error {
|
||||||
|
if _, err := s.pool.Exec(ctx, `DELETE FROM mail_rules WHERE tenant_slug = $1 AND id = $2`, tenantSlug, id); err != nil {
|
||||||
|
return fmt.Errorf("mailrules: regel löschen: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,80 @@
|
|||||||
|
// Package syncalert implementiert IMP-08: Benachrichtigung bei
|
||||||
|
// wiederholtem Postfach-Sync-Ausfall, mit Eskalationsschwelle statt
|
||||||
|
// Einzel-Alarm pro Fehlversuch. Versand ausschließlich über den
|
||||||
|
// zentralen Core-Benachrichtigungs-Dispatcher (CFG-02, bereits Fertig)
|
||||||
|
// — dieses Paket baut KEINEN eigenen E-Mail-Versand, sondern ruft
|
||||||
|
// ausschließlich NotificationDispatcher.Enqueue auf (Akzeptanzkriterium
|
||||||
|
// 1), exakt einmal je Eskalation.
|
||||||
|
package syncalert
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
)
|
||||||
|
|
||||||
|
// NotificationDispatcher ist die schmale Schnittstelle zu Core CFG-02
|
||||||
|
// (internal/notify.Dispatcher.Enqueue) — Mail ruft ausschließlich diese
|
||||||
|
// EINE Methode auf, kein eigener Versandcode.
|
||||||
|
type NotificationDispatcher interface {
|
||||||
|
Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (id string, err error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// HTTPNotificationDispatcher spricht CFG-02 über HTTP an — dieselbe
|
||||||
|
// Service-Credential-Konvention wie mail/internal/crypto.HTTPKEKProvider
|
||||||
|
// (API-02, X-Nexarch-Client-Id/Secret).
|
||||||
|
type HTTPNotificationDispatcher struct {
|
||||||
|
endpointURL string
|
||||||
|
clientID string
|
||||||
|
clientSecret string
|
||||||
|
httpClient *http.Client
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewHTTPNotificationDispatcher(endpointURL, clientID, clientSecret string, httpClient *http.Client) *HTTPNotificationDispatcher {
|
||||||
|
if httpClient == nil {
|
||||||
|
httpClient = http.DefaultClient
|
||||||
|
}
|
||||||
|
return &HTTPNotificationDispatcher{endpointURL: endpointURL, clientID: clientID, clientSecret: clientSecret, httpClient: httpClient}
|
||||||
|
}
|
||||||
|
|
||||||
|
type enqueueRequest struct {
|
||||||
|
Channel string `json:"channel"`
|
||||||
|
Recipient string `json:"recipient"`
|
||||||
|
Payload map[string]any `json:"payload"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type enqueueResponse struct {
|
||||||
|
ID string `json:"id"`
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d *HTTPNotificationDispatcher) Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (string, error) {
|
||||||
|
body, err := json.Marshal(enqueueRequest{Channel: channel, Recipient: recipient, Payload: payload})
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("syncalert: anfrage serialisieren: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, d.endpointURL, bytes.NewReader(body))
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("syncalert: anfrage aufbauen: %w", err)
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
req.Header.Set("X-Nexarch-Client-Id", d.clientID)
|
||||||
|
req.Header.Set("X-Nexarch-Client-Secret", d.clientSecret)
|
||||||
|
|
||||||
|
resp, err := d.httpClient.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("syncalert: anfrage senden: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
return "", fmt.Errorf("syncalert: cfg-02 lehnte anfrage ab: status %d", resp.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
var out enqueueResponse
|
||||||
|
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
||||||
|
return "", fmt.Errorf("syncalert: antwort dekodieren: %w", err)
|
||||||
|
}
|
||||||
|
return out.ID, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,67 @@
|
|||||||
|
package syncalert
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestHTTPNotificationDispatcher_SendsCorrectRequestFormat beweist real
|
||||||
|
// über echtes HTTP, dass HTTPNotificationDispatcher Channel/Recipient/
|
||||||
|
// Payload sowie die Service-Credential-Header korrekt sendet — kein
|
||||||
|
// eigener E-Mail-Versand, nur ein einziger CFG-02-Aufruf
|
||||||
|
// (Akzeptanzkriterium 1).
|
||||||
|
func TestHTTPNotificationDispatcher_SendsCorrectRequestFormat(t *testing.T) {
|
||||||
|
var capturedBody map[string]any
|
||||||
|
var capturedClientID, capturedClientSecret string
|
||||||
|
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
capturedClientID = r.Header.Get("X-Nexarch-Client-Id")
|
||||||
|
capturedClientSecret = r.Header.Get("X-Nexarch-Client-Secret")
|
||||||
|
if err := json.NewDecoder(r.Body).Decode(&capturedBody); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
_, _ = w.Write([]byte(`{"id":"real-notification-id-123"}`))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
dispatcher := NewHTTPNotificationDispatcher(srv.URL, "mail", "mail-service-secret", nil)
|
||||||
|
id, err := dispatcher.Enqueue(context.Background(), NotificationChannel, AdminRecipient, map[string]any{
|
||||||
|
"mailbox": "INBOX",
|
||||||
|
"reason": "verbindung abgelehnt",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("enqueue: %v", err)
|
||||||
|
}
|
||||||
|
if id != "real-notification-id-123" {
|
||||||
|
t.Fatalf("erwartete reale id vom server, habe %q", id)
|
||||||
|
}
|
||||||
|
if capturedClientID != "mail" || capturedClientSecret != "mail-service-secret" {
|
||||||
|
t.Fatalf("service-credential-header fehlen/falsch: id=%q secret=%q", capturedClientID, capturedClientSecret)
|
||||||
|
}
|
||||||
|
if capturedBody["channel"] != NotificationChannel || capturedBody["recipient"] != AdminRecipient {
|
||||||
|
t.Fatalf("channel/recipient falsch übertragen: %+v", capturedBody)
|
||||||
|
}
|
||||||
|
payload, ok := capturedBody["payload"].(map[string]any)
|
||||||
|
if !ok || payload["mailbox"] != "INBOX" {
|
||||||
|
t.Fatalf("payload nicht korrekt übertragen: %+v", capturedBody)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHTTPNotificationDispatcher_RejectedByServerReturnsError bestätigt,
|
||||||
|
// dass eine Ablehnung durch CFG-02 real als Fehler durchgereicht wird,
|
||||||
|
// statt stillschweigend zu verschwinden.
|
||||||
|
func TestHTTPNotificationDispatcher_RejectedByServerReturnsError(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusForbidden)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
dispatcher := NewHTTPNotificationDispatcher(srv.URL, "mail", "falsch", nil)
|
||||||
|
if _, err := dispatcher.Enqueue(context.Background(), "c", "r", nil); err == nil {
|
||||||
|
t.Fatal("erwartete fehler bei abgelehnter anfrage, habe nil")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
CREATE TABLE IF NOT EXISTS mail_sync_alert_state (
|
||||||
|
tenant_slug TEXT NOT NULL,
|
||||||
|
mailbox_name TEXT NOT NULL,
|
||||||
|
consecutive_failures INT NOT NULL DEFAULT 0,
|
||||||
|
alerted BOOLEAN NOT NULL DEFAULT false,
|
||||||
|
last_success_at TIMESTAMPTZ,
|
||||||
|
last_failure_reason TEXT,
|
||||||
|
last_failure_at TIMESTAMPTZ,
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
PRIMARY KEY (tenant_slug, mailbox_name)
|
||||||
|
)
|
||||||
@@ -0,0 +1,156 @@
|
|||||||
|
package syncalert
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
_ "embed"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
//go:embed migrations/0001_mail_sync_alert_state.sql
|
||||||
|
var schemaMigration string
|
||||||
|
|
||||||
|
const (
|
||||||
|
// NotificationChannel/AdminRecipient sind bewusst statisch (kleinste
|
||||||
|
// Lösung) — eine konfigurierbare Empfängerverwaltung ist Sache einer
|
||||||
|
// späteren Kachel, nicht Bestandteil von IMP-08.
|
||||||
|
NotificationChannel = "mail-sync-failure"
|
||||||
|
AdminRecipient = "mail-admins"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Monitor verfolgt Sync-Fehlschläge je Mandant/Postfach und löst bei
|
||||||
|
// Überschreiten der Schwelle GENAU EINE Benachrichtigung aus
|
||||||
|
// (Akzeptanzkriterium 1).
|
||||||
|
type Monitor struct {
|
||||||
|
pool *pgxpool.Pool
|
||||||
|
dispatcher NotificationDispatcher
|
||||||
|
threshold int
|
||||||
|
now func() time.Time
|
||||||
|
}
|
||||||
|
|
||||||
|
// DefaultThreshold ist die Vorgabe-Eskalationsschwelle (konsekutive
|
||||||
|
// Fehlschläge), überschreibbar über WithThreshold.
|
||||||
|
const DefaultThreshold = 3
|
||||||
|
|
||||||
|
func NewMonitor(pool *pgxpool.Pool, dispatcher NotificationDispatcher) *Monitor {
|
||||||
|
return &Monitor{pool: pool, dispatcher: dispatcher, threshold: DefaultThreshold, now: time.Now}
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithThreshold setzt eine abweichende Eskalationsschwelle.
|
||||||
|
func (m *Monitor) WithThreshold(threshold int) *Monitor {
|
||||||
|
m.threshold = threshold
|
||||||
|
return m
|
||||||
|
}
|
||||||
|
|
||||||
|
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
|
||||||
|
func (m *Monitor) EnsureSchema(ctx context.Context) error {
|
||||||
|
if _, err := m.pool.Exec(ctx, schemaMigration); err != nil {
|
||||||
|
return fmt.Errorf("syncalert: schema anlegen: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type alertState struct {
|
||||||
|
consecutiveFailures int
|
||||||
|
alerted bool
|
||||||
|
lastSuccessAt *time.Time
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *Monitor) getOrCreate(ctx context.Context, tenantSlug, mailboxName string) (alertState, error) {
|
||||||
|
if _, err := m.pool.Exec(ctx, `
|
||||||
|
INSERT INTO mail_sync_alert_state (tenant_slug, mailbox_name)
|
||||||
|
VALUES ($1, $2)
|
||||||
|
ON CONFLICT (tenant_slug, mailbox_name) DO NOTHING
|
||||||
|
`, tenantSlug, mailboxName); err != nil {
|
||||||
|
return alertState{}, fmt.Errorf("syncalert: zustand anlegen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var st alertState
|
||||||
|
err := m.pool.QueryRow(ctx, `
|
||||||
|
SELECT consecutive_failures, alerted, last_success_at
|
||||||
|
FROM mail_sync_alert_state WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||||
|
`, tenantSlug, mailboxName).Scan(&st.consecutiveFailures, &st.alerted, &st.lastSuccessAt)
|
||||||
|
if err != nil {
|
||||||
|
if errors.Is(err, pgx.ErrNoRows) {
|
||||||
|
return alertState{}, fmt.Errorf("syncalert: gerade angelegten zustand nicht gefunden")
|
||||||
|
}
|
||||||
|
return alertState{}, fmt.Errorf("syncalert: zustand lesen: %w", err)
|
||||||
|
}
|
||||||
|
return st, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// RecordFailure verzeichnet einen fehlgeschlagenen Sync-Versuch. Erst
|
||||||
|
// wenn consecutive_failures die konfigurierte Schwelle ERSTMALIG
|
||||||
|
// erreicht (noch nicht "alerted"), wird GENAU EINE Benachrichtigung an
|
||||||
|
// CFG-02 ausgelöst (Akzeptanzkriterium 1) — weitere Fehlschläge danach
|
||||||
|
// lösen KEINE zusätzliche Benachrichtigung aus, solange der Alarmzustand
|
||||||
|
// nicht durch einen erfolgreichen Sync zurückgesetzt wurde (kein
|
||||||
|
// Einzel-Alarm pro Fehlversuch, keine Spam-Flut).
|
||||||
|
func (m *Monitor) RecordFailure(ctx context.Context, tenantSlug, mailboxName, reason string) error {
|
||||||
|
now := m.now()
|
||||||
|
st, err := m.getOrCreate(ctx, tenantSlug, mailboxName)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
newFailures := st.consecutiveFailures + 1
|
||||||
|
if _, err := m.pool.Exec(ctx, `
|
||||||
|
UPDATE mail_sync_alert_state
|
||||||
|
SET consecutive_failures = $3, last_failure_reason = $4, last_failure_at = $5, updated_at = now()
|
||||||
|
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||||
|
`, tenantSlug, mailboxName, newFailures, reason, now); err != nil {
|
||||||
|
return fmt.Errorf("syncalert: fehlschlag erfassen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if newFailures < m.threshold || st.alerted {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Akzeptanzkriterium 2: Benachrichtigung enthält Postfach,
|
||||||
|
// Fehlerursache und Zeitpunkt des letzten erfolgreichen Abrufs.
|
||||||
|
payload := map[string]any{
|
||||||
|
"tenant_slug": tenantSlug,
|
||||||
|
"mailbox": mailboxName,
|
||||||
|
"reason": reason,
|
||||||
|
"consecutive_failures": newFailures,
|
||||||
|
"last_successful_sync": formatOptionalTime(st.lastSuccessAt),
|
||||||
|
}
|
||||||
|
if _, err := m.dispatcher.Enqueue(ctx, NotificationChannel, AdminRecipient, payload); err != nil {
|
||||||
|
return fmt.Errorf("syncalert: benachrichtigung auslösen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if _, err := m.pool.Exec(ctx, `
|
||||||
|
UPDATE mail_sync_alert_state SET alerted = true, updated_at = now()
|
||||||
|
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||||
|
`, tenantSlug, mailboxName); err != nil {
|
||||||
|
return fmt.Errorf("syncalert: alarmzustand markieren: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// RecordSuccess verzeichnet einen erfolgreichen Sync und setzt den
|
||||||
|
// Alarmzustand zurück (Akzeptanzkriterium 3) — der nächste Fehlschlag
|
||||||
|
// nach einem Erfolg beginnt wieder bei 0 konsekutiven Fehlschlägen.
|
||||||
|
func (m *Monitor) RecordSuccess(ctx context.Context, tenantSlug, mailboxName string) error {
|
||||||
|
now := m.now()
|
||||||
|
if _, err := m.pool.Exec(ctx, `
|
||||||
|
INSERT INTO mail_sync_alert_state (tenant_slug, mailbox_name, consecutive_failures, alerted, last_success_at)
|
||||||
|
VALUES ($1, $2, 0, false, $3)
|
||||||
|
ON CONFLICT (tenant_slug, mailbox_name) DO UPDATE
|
||||||
|
SET consecutive_failures = 0, alerted = false, last_success_at = $3, updated_at = now()
|
||||||
|
`, tenantSlug, mailboxName, now); err != nil {
|
||||||
|
return fmt.Errorf("syncalert: erfolg erfassen: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func formatOptionalTime(t *time.Time) string {
|
||||||
|
if t == nil {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
return t.UTC().Format(time.RFC3339)
|
||||||
|
}
|
||||||
@@ -0,0 +1,220 @@
|
|||||||
|
// Integrationstest (IMP-08): echte Postgres-Instanz, folgt derselben
|
||||||
|
// Testhost-Konvention wie mail/internal/dedup/folderstate/imapimport —
|
||||||
|
// TEST_TENANT_DSN.
|
||||||
|
package syncalert
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"sync"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakeDispatcher zeichnet jeden Enqueue-Aufruf auf — echte HTTP-
|
||||||
|
// Anbindung ist Sache von dispatcher_http_test.go, hier wird die
|
||||||
|
// Eskalationslogik isoliert geprüft (gleiche Konvention wie
|
||||||
|
// fakeAuthenticator/fakeKEKProvider in anderen Mail-Paketen).
|
||||||
|
type fakeDispatcher struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
calls []map[string]any
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeDispatcher) Enqueue(_ context.Context, channel, recipient string, payload map[string]any) (string, error) {
|
||||||
|
f.mu.Lock()
|
||||||
|
defer f.mu.Unlock()
|
||||||
|
call := map[string]any{"channel": channel, "recipient": recipient}
|
||||||
|
for k, v := range payload {
|
||||||
|
call[k] = v
|
||||||
|
}
|
||||||
|
f.calls = append(f.calls, call)
|
||||||
|
return fmt.Sprintf("notif-%d", len(f.calls)), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeDispatcher) count() int {
|
||||||
|
f.mu.Lock()
|
||||||
|
defer f.mu.Unlock()
|
||||||
|
return len(f.calls)
|
||||||
|
}
|
||||||
|
|
||||||
|
func setupMonitor(t *testing.T, dispatcher NotificationDispatcher) *Monitor {
|
||||||
|
t.Helper()
|
||||||
|
dsn := os.Getenv("TEST_TENANT_DSN")
|
||||||
|
if dsn == "" {
|
||||||
|
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
|
||||||
|
}
|
||||||
|
ctx := context.Background()
|
||||||
|
pool, err := pgxpool.New(ctx, dsn)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("pool: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { pool.Close() })
|
||||||
|
|
||||||
|
monitor := NewMonitor(pool, dispatcher).WithThreshold(3)
|
||||||
|
if err := monitor.EnsureSchema(ctx); err != nil {
|
||||||
|
t.Fatalf("schema: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() {
|
||||||
|
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_sync_alert_state WHERE tenant_slug LIKE 'mandant-imp08-%'`)
|
||||||
|
})
|
||||||
|
return monitor
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification
|
||||||
|
// ist die geforderte Pflichtprüfung 1: N aufeinanderfolgende
|
||||||
|
// Fehlschläge lösen genau eine Benachrichtigung aus, keine Spam-Flut.
|
||||||
|
func TestRecordFailure_NConsecutiveFailuresTriggerExactlyOneNotification(t *testing.T) {
|
||||||
|
dispatcher := &fakeDispatcher{}
|
||||||
|
monitor := setupMonitor(t, dispatcher)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp08-schwelle"
|
||||||
|
|
||||||
|
// Schwelle ist 3 — die ersten 2 Fehlschläge dürfen NICHTS auslösen.
|
||||||
|
for i := 0; i < 2; i++ {
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
|
||||||
|
t.Fatalf("recordfailure %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 0 {
|
||||||
|
t.Fatalf("erwartete 0 benachrichtigungen vor erreichen der schwelle, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
|
||||||
|
// Dritter Fehlschlag erreicht die Schwelle — GENAU EINE Benachrichtigung.
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
|
||||||
|
t.Fatalf("recordfailure 3: %v", err)
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 1 {
|
||||||
|
t.Fatalf("erwartete genau 1 benachrichtigung bei erreichen der schwelle, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
|
||||||
|
// Weitere Fehlschläge DANACH dürfen KEINE zusätzliche Benachrichtigung
|
||||||
|
// auslösen (kein Einzel-Alarm pro Fehlversuch, keine Spam-Flut).
|
||||||
|
for i := 0; i < 5; i++ {
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "verbindung abgelehnt"); err != nil {
|
||||||
|
t.Fatalf("weiterer fehlschlag %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 1 {
|
||||||
|
t.Fatalf("erwartete weiterhin genau 1 benachrichtigung nach 5 weiteren fehlschlägen, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRecordSuccess_EndsAlertStateVerifiably ist die geforderte
|
||||||
|
// Pflichtprüfung 2: erfolgreicher Lauf nach Ausfall beendet den
|
||||||
|
// Alarmzustand nachvollziehbar.
|
||||||
|
func TestRecordSuccess_EndsAlertStateVerifiably(t *testing.T) {
|
||||||
|
dispatcher := &fakeDispatcher{}
|
||||||
|
monitor := setupMonitor(t, dispatcher)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp08-reset"
|
||||||
|
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
|
||||||
|
t.Fatalf("recordfailure %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 1 {
|
||||||
|
t.Fatalf("erwartete 1 benachrichtigung nach 3 fehlschlägen, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := monitor.RecordSuccess(ctx, tenant, "INBOX"); err != nil {
|
||||||
|
t.Fatalf("recordsuccess: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Nachvollziehbar zurückgesetzt: der NÄCHSTE Fehlschlags-Zyklus muss
|
||||||
|
// real wieder bei 0 beginnen und erneut die volle Schwelle
|
||||||
|
// durchlaufen, bevor eine ZWEITE Benachrichtigung ausgelöst wird.
|
||||||
|
for i := 0; i < 2; i++ {
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
|
||||||
|
t.Fatalf("recordfailure nach reset %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 1 {
|
||||||
|
t.Fatalf("erwartete weiterhin nur 1 benachrichtigung (schwelle nach reset noch nicht erreicht), habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "timeout"); err != nil {
|
||||||
|
t.Fatalf("dritter fehlschlag nach reset: %v", err)
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 2 {
|
||||||
|
t.Fatalf("erwartete 2. benachrichtigung nach erneutem erreichen der schwelle, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRecordFailure_MultipleAffectedMailboxesStayIsolated ist die
|
||||||
|
// geforderte Pflichtprüfung 3: Test mit mehreren betroffenen
|
||||||
|
// Postfächern gleichzeitig bleibt übersichtlich (korrekt isoliert).
|
||||||
|
func TestRecordFailure_MultipleAffectedMailboxesStayIsolated(t *testing.T) {
|
||||||
|
dispatcher := &fakeDispatcher{}
|
||||||
|
monitor := setupMonitor(t, dispatcher)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp08-mehrere"
|
||||||
|
|
||||||
|
mailboxes := []string{"INBOX", "Archiv", "Vertrieb"}
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for _, mailbox := range mailboxes {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(mb string) {
|
||||||
|
defer wg.Done()
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
_ = monitor.RecordFailure(ctx, tenant, mb, "gleichzeitiger ausfall")
|
||||||
|
}
|
||||||
|
}(mailbox)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
if dispatcher.count() != len(mailboxes) {
|
||||||
|
t.Fatalf("erwartete genau 1 benachrichtigung je betroffenem postfach (%d), habe %d", len(mailboxes), dispatcher.count())
|
||||||
|
}
|
||||||
|
|
||||||
|
seenMailboxes := map[string]bool{}
|
||||||
|
dispatcher.mu.Lock()
|
||||||
|
for _, call := range dispatcher.calls {
|
||||||
|
mb, _ := call["mailbox"].(string)
|
||||||
|
if seenMailboxes[mb] {
|
||||||
|
t.Fatalf("postfach %q hat mehr als eine benachrichtigung erhalten", mb)
|
||||||
|
}
|
||||||
|
seenMailboxes[mb] = true
|
||||||
|
}
|
||||||
|
dispatcher.mu.Unlock()
|
||||||
|
for _, mb := range mailboxes {
|
||||||
|
if !seenMailboxes[mb] {
|
||||||
|
t.Fatalf("postfach %q fehlt unter den benachrichtigten, habe: %v", mb, seenMailboxes)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRecordFailure_NotificationContainsRequiredFields deckt
|
||||||
|
// Akzeptanzkriterium 2 ab: Benachrichtigung enthält Postfach,
|
||||||
|
// Fehlerursache und Zeitpunkt des letzten erfolgreichen Abrufs.
|
||||||
|
func TestRecordFailure_NotificationContainsRequiredFields(t *testing.T) {
|
||||||
|
dispatcher := &fakeDispatcher{}
|
||||||
|
monitor := setupMonitor(t, dispatcher)
|
||||||
|
ctx := context.Background()
|
||||||
|
tenant := "mandant-imp08-inhalt"
|
||||||
|
|
||||||
|
if err := monitor.RecordSuccess(ctx, tenant, "INBOX"); err != nil {
|
||||||
|
t.Fatalf("initialer erfolg: %v", err)
|
||||||
|
}
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
if err := monitor.RecordFailure(ctx, tenant, "INBOX", "authentifizierung fehlgeschlagen"); err != nil {
|
||||||
|
t.Fatalf("recordfailure %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if dispatcher.count() != 1 {
|
||||||
|
t.Fatalf("erwartete 1 benachrichtigung, habe %d", dispatcher.count())
|
||||||
|
}
|
||||||
|
|
||||||
|
call := dispatcher.calls[0]
|
||||||
|
if call["mailbox"] != "INBOX" {
|
||||||
|
t.Fatalf("erwartete postfach 'INBOX' in der benachrichtigung, habe: %v", call["mailbox"])
|
||||||
|
}
|
||||||
|
if call["reason"] != "authentifizierung fehlgeschlagen" {
|
||||||
|
t.Fatalf("erwartete fehlerursache in der benachrichtigung, habe: %v", call["reason"])
|
||||||
|
}
|
||||||
|
lastSuccess, _ := call["last_successful_sync"].(string)
|
||||||
|
if lastSuccess == "" {
|
||||||
|
t.Fatal("erwartete zeitpunkt des letzten erfolgreichen abrufs in der benachrichtigung")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user