Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5fcae51aac | ||
|
|
369a40af10 | ||
|
|
06dbd52d4c | ||
|
|
da80643564 | ||
|
|
814a7fda0a | ||
|
|
6a03dcafd6 | ||
|
|
6f532d8350 | ||
|
|
4180a26c6e |
@@ -129,3 +129,34 @@ Keine Commits in dieser Session.
|
||||
- migrations/0001_tenant_registry.sql | 10 ++++++++++
|
||||
|
||||
---
|
||||
## 2026-08-28 23:55 – 23:59 (3m)
|
||||
**Beschreibung:** Claude Code Session
|
||||
**Projekt:** nexarch
|
||||
|
||||
### Commits
|
||||
- 369a40a OPS-05: alerting-bei-schwellwert-ueberschreitung (internal/alerting: regel-store, evaluator gegen ops-03-metriken, cfg-02-zustellung, drosselung je regel+zeitreihe)
|
||||
- 06dbd52 Merge branch 'feature/cfg-02-benachrichtigungs-dispatcher-core-service-fuer-module' into feature/ops-05-alerting-bei-schwellwert-ueberschreitung
|
||||
|
||||
### Geänderte Dateien
|
||||
- internal/alerting/evaluator.go | 194 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- internal/alerting/evaluator_test.go | 233 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- internal/alerting/rules.go | 132 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- migrations/0006_alert_rules.down.sql | 2 ++
|
||||
- migrations/0006_alert_rules.up.sql | 24 ++++++++++++++++++++++
|
||||
|
||||
---
|
||||
## 2026-08-29 00:00 – 00:00 (0m)
|
||||
**Beschreibung:** Claude Code Session
|
||||
**Projekt:** code
|
||||
|
||||
### Commits
|
||||
Keine Commits in dieser Session.
|
||||
|
||||
### Geänderte Dateien
|
||||
- internal/alerting/evaluator.go | 194 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- internal/alerting/evaluator_test.go | 233 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- internal/alerting/rules.go | 132 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
|
||||
- migrations/0006_alert_rules.down.sql | 2 ++
|
||||
- migrations/0006_alert_rules.up.sql | 24 ++++++++++++++++++++++
|
||||
|
||||
---
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
// metrics-devserver stellt den OPS-03-Metrics-Aggregator (internal/metrics)
|
||||
// unter /metrics bereit, damit ein echter Prometheus-Scrape-Vorgang gegen
|
||||
// den Core-Dienst geprueft werden kann (Pruefung 3). Getrennt von cmd/core
|
||||
// aus demselben Grund wie die anderen *-devserver.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/db"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/metrics"
|
||||
)
|
||||
|
||||
func main() {
|
||||
dsn := os.Getenv("NEXARCH_REGISTRY_DSN")
|
||||
if dsn == "" {
|
||||
log.Fatal("NEXARCH_REGISTRY_DSN nicht gesetzt")
|
||||
}
|
||||
addr := os.Getenv("NEXARCH_METRICS_LISTEN_ADDR")
|
||||
if addr == "" {
|
||||
addr = ":8085"
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
pool, err := db.Connect(ctx, dsn)
|
||||
if err != nil {
|
||||
log.Fatalf("db: %v", err)
|
||||
}
|
||||
defer pool.Close()
|
||||
|
||||
sourceStore := metrics.NewSourceStore(pool)
|
||||
agg := metrics.NewAggregator(metrics.NewCoreRegistry(), sourceStore.Provide)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/metrics", agg.Handler())
|
||||
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
|
||||
|
||||
log.Printf("metrics-devserver listening on %s", addr)
|
||||
log.Fatal(http.ListenAndServe(addr, mux))
|
||||
}
|
||||
@@ -1,70 +0,0 @@
|
||||
// resync-api ist der Aufrufpunkt fuer API-11: startet den bereits
|
||||
// fertigen internal/resync.Handler (API-06) als eigenstaendigen
|
||||
// HTTP-Dienst. REINES WIRING — keine Aenderung an internal/resync/,
|
||||
// internal/audit/, internal/usage/ oder internal/moduleregistry/.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/audit"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/resync"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/usage"
|
||||
)
|
||||
|
||||
// auditAdapter erfüllt resync.AuditRecorder über den bestehenden
|
||||
// audit.Log-Schreibpfad — kein neuer Audit-Code, nur Signatur-Anpassung
|
||||
// (audit.Log.Record nimmt ein Event-Struct, resync.AuditRecorder einzelne
|
||||
// Felder).
|
||||
type auditAdapter struct{ log *audit.Log }
|
||||
|
||||
func (a auditAdapter) Record(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any, occurredAt time.Time) error {
|
||||
return a.log.Record(ctx, audit.Event{
|
||||
TenantSlug: tenantSlug, Actor: actor, Action: action, Target: target,
|
||||
Metadata: metadata, OccurredAt: occurredAt,
|
||||
})
|
||||
}
|
||||
|
||||
func main() {
|
||||
registryDSN := os.Getenv("NEXARCH_RESYNC_REGISTRY_DSN")
|
||||
if registryDSN == "" {
|
||||
log.Fatal("NEXARCH_RESYNC_REGISTRY_DSN muss gesetzt sein")
|
||||
}
|
||||
addr := os.Getenv("NEXARCH_RESYNC_API_LISTEN_ADDR")
|
||||
if addr == "" {
|
||||
addr = "127.0.0.1:8098"
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
pool, err := pgxpool.New(ctx, registryDSN)
|
||||
if err != nil {
|
||||
log.Fatalf("datenbankverbindung: %v", err)
|
||||
}
|
||||
defer pool.Close()
|
||||
|
||||
flagStore := flag.NewStore(pool)
|
||||
flagService := flag.NewService(flagStore, 30*time.Second)
|
||||
registry := moduleregistry.NewRegistry(pool, flagService)
|
||||
auditLog := audit.NewLog(pool)
|
||||
usageStore := usage.NewStore(pool)
|
||||
|
||||
handler := resync.NewHandler(registry, auditAdapter{log: auditLog}, usageStore)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/internal/resync/audit", handler.AuditHandler)
|
||||
mux.HandleFunc("/internal/resync/usage", handler.UsageHandler)
|
||||
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
|
||||
|
||||
log.Printf("resync-api: listening on %s", addr)
|
||||
if err := http.ListenAndServe(addr, mux); err != nil {
|
||||
log.Fatalf("http server: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -1,14 +0,0 @@
|
||||
[Unit]
|
||||
Description=NEXARCH Core - Wiederanlauf-Nachsynchronisierung (API-06/API-11)
|
||||
After=network.target postgresql.service
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
User=nexarch
|
||||
EnvironmentFile=/etc/nexarch/resync-api.env
|
||||
ExecStart=__INSTALL_DIR__/bin/resync-api
|
||||
Restart=on-failure
|
||||
StandardOutput=journal
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
@@ -1,76 +0,0 @@
|
||||
# API-11 – Prüfprotokoll: Wiederanlauf-Nachsynchronisierungs-Endpunkt starten (API-06 als laufender Dienst)
|
||||
|
||||
Voraussetzung API-06 – bereits Fertig, hier UNVERÄNDERT.
|
||||
|
||||
## Reines Wiring, keine neue Logik
|
||||
|
||||
`git diff --stat internal/resync/ internal/audit/ internal/usage/ internal/moduleregistry/`
|
||||
liefert KEINEN Diff gegenüber den jeweiligen Ticket-Ständen. `API-11`
|
||||
fügt ausschließlich `cmd/resync-api/main.go` hinzu — inklusive eines
|
||||
kleinen `auditAdapter`, der `resync.AuditRecorder` (einzelne Felder)
|
||||
auf `audit.Log.Record` (Event-Struct) abbildet. Das ist reine
|
||||
Signatur-Anpassung, keine neue Geschäftslogik.
|
||||
|
||||
## Root Cause (dokumentiert)
|
||||
|
||||
`cmd/core/main.go` ist seit TEN-01 minimal geblieben (nur `/healthz`,
|
||||
`/internal/tenants`) — kein späteres Ticket (RBAC-02, CFG-02, API-06,
|
||||
...) wurde je dort zentral eingehängt. Jedes Modul entstand auf einer
|
||||
eigenen, unabhängigen Feature-Branch-Kette. API-11 folgt dem in dieser
|
||||
Session etablierten Muster (RBAC-06, CFG-05, RET-09): ein eigener,
|
||||
kleiner HTTP-Dienst statt eines zentralen `cmd/core`-Umbaus.
|
||||
|
||||
## Umsetzung
|
||||
|
||||
- `cmd/resync-api/main.go` – startet `internal/resync.Handler` mit
|
||||
echten Produktions-Implementierungen: `moduleregistry.Registry`
|
||||
(Auth), `audit.Log` (über `auditAdapter`), `usage.Store`.
|
||||
- `deploy/systemd/nexarch-resync-api.service.tmpl`.
|
||||
|
||||
## Prüfungen
|
||||
|
||||
| # | Prüfung | Ergebnis |
|
||||
|---|---|---|
|
||||
| 1 | Dienst startet und bleibt stabil (systemctl status aktiv) | **bestanden** – real auf 131: `nexarch-resync-api.service` aktiv, `Restart=on-failure` |
|
||||
| 2 | Realer POST /internal/resync/usage von einem externen Testclient gegen den laufenden Dienst liefert die erwartete Verarbeitung | **bestanden** – real per `curl`: mit echtem, über `moduleregistry.Registry.Provision` ausgestelltem Service-Credential (`X-Nexarch-Client-Id`/`X-Nexarch-Client-Secret`) liefert der Aufruf `{"applied":1}`, `usage_counters` zeigt real den erhöhten Zähler; mit falschem Credential 401. Testdaten (Modul, Credential, Zähler-Zeile) anschließend entfernt |
|
||||
| 3 | Code-Review: keine Änderung an internal/resync/ selbst, nur main.go+systemd neu | **bestanden** – `git diff --stat` bestätigt: `internal/resync/`, `internal/audit/`, `internal/usage/`, `internal/moduleregistry/` unverändert gegenüber ihren jeweiligen Ticket-Ständen |
|
||||
|
||||
## Echte Verdrahtung auf 192.168.1.131
|
||||
|
||||
- `resync-api` gebaut nach `/opt/nexarch-core/bin/`,
|
||||
`/etc/nexarch/resync-api.env` (0600), `nexarch-resync-api.service`
|
||||
installiert/aktiviert.
|
||||
- Reale Rechtevergabe-Lücke gefunden und behoben (gleiches Muster wie
|
||||
bei den vorherigen Wrapper-Diensten): `modules`, `module_credentials`,
|
||||
`audit_events`, `feature_flags`, `usage_counters`,
|
||||
`resync_audit_buffer`, `resync_usage_buffer` gehörten `postgres`,
|
||||
`nexarch_core` hatte keine Rechte — `GRANT` nachgezogen und über
|
||||
`information_schema.role_table_grants` verifiziert, bevor der
|
||||
End-zu-Ende-Test erneut lief.
|
||||
- Zusätzliche reale Erkenntnis: `usage_counters.tenant_id` ist `UUID`,
|
||||
nicht der Tenant-Slug (String) — beim ersten Testversuch mit `"acme"`
|
||||
scheiterte der Insert intern, `UsageHandler` meldete `applied:0` statt
|
||||
eines Fehlers (stiller Fehlschlag pro Delta, so von API-06 selbst so
|
||||
entworfen: "Aufrufer entfernt aus seinem Puffer nur bestätigt
|
||||
übernommene Deltas" — kein API-11-Defekt, sondern korrektes,
|
||||
bestehendes API-06-Verhalten). Mit echter UUID als `tenant_slug`-Wert
|
||||
lieferte der Aufruf real `applied:1`.
|
||||
|
||||
## Build/Test-Ergebnis (192.168.1.131)
|
||||
|
||||
```
|
||||
go build ./... -> clean
|
||||
go vet ./... -> clean
|
||||
golangci-lint run ./cmd/resync-api/... -> 0 issues
|
||||
```
|
||||
|
||||
Keine neuen Go-Tests nötig (kein neuer Fachcode außer main.go/Adapter,
|
||||
die eigentliche Logik ist bereits durch API-06s eigene Tests
|
||||
abgedeckt).
|
||||
|
||||
## Gesamtergebnis
|
||||
|
||||
**Bestanden.** API-06 ist jetzt ein real laufender, über systemd
|
||||
verwalteter Dienst. Modul-Clients wie DMS' `storage.HTTPUsageReporter`
|
||||
(RET-06/DOC-16-Umfeld) können sich jetzt real gegen einen laufenden
|
||||
Endpunkt verdrahten, statt gegen unverdrahteten Go-Code zu testen.
|
||||
@@ -1,17 +1,25 @@
|
||||
module gitea.perlbach24.de/scripte/nexarch
|
||||
|
||||
go 1.22
|
||||
go 1.25.0
|
||||
|
||||
require (
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1
|
||||
github.com/jackc/pgx/v5 v5.6.0
|
||||
github.com/prometheus/client_golang v1.24.1
|
||||
github.com/prometheus/client_model v0.6.2
|
||||
github.com/prometheus/common v0.70.1
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/beorn7/perks v1.0.1 // indirect
|
||||
github.com/cespare/xxhash/v2 v2.3.0 // 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
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
|
||||
github.com/prometheus/procfs v0.21.1 // indirect
|
||||
golang.org/x/crypto v0.17.0 // indirect
|
||||
golang.org/x/sync v0.1.0 // indirect
|
||||
golang.org/x/text v0.14.0 // indirect
|
||||
golang.org/x/sync v0.22.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
golang.org/x/text v0.40.0 // indirect
|
||||
google.golang.org/protobuf v1.36.11 // indirect
|
||||
)
|
||||
|
||||
@@ -1,8 +1,12 @@
|
||||
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
|
||||
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
|
||||
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
|
||||
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
|
||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
|
||||
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
|
||||
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
|
||||
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
|
||||
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
|
||||
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk=
|
||||
@@ -11,19 +15,37 @@ github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY=
|
||||
github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw=
|
||||
github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk=
|
||||
github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
|
||||
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
|
||||
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
|
||||
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
|
||||
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
|
||||
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
|
||||
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
|
||||
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
|
||||
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
|
||||
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
|
||||
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
|
||||
golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k=
|
||||
golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
|
||||
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
|
||||
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ=
|
||||
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs=
|
||||
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY=
|
||||
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
|
||||
google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
|
||||
@@ -0,0 +1,194 @@
|
||||
package alerting
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/notify"
|
||||
)
|
||||
|
||||
// DefaultDebounceInterval: wiederholte Alarmierung für denselben
|
||||
// anhaltenden Zustand ist gedrosselt (Akzeptanzkriterium 3) — 15 Minuten
|
||||
// ist ein üblicher Kompromiss zwischen "schnell genug informiert" und
|
||||
// "kein Alarm-Spam bei dauerhaft überschrittenem Wert".
|
||||
const DefaultDebounceInterval = 15 * time.Minute
|
||||
|
||||
// AlertChannel ist der CFG-02-Kanal, über den Schwellwert-Alarme zugestellt
|
||||
// werden — ein eigener Kanalname, damit Zustellregeln/-vorlagen (CFG-03)
|
||||
// unabhängig von anderen Benachrichtigungsarten konfiguriert werden können.
|
||||
const AlertChannel = "alert"
|
||||
|
||||
// Evaluator prüft konfigurierte Regeln gegen aktuell gesammelte Metriken
|
||||
// (aus internal/metrics.Aggregator.Gather) und löst bei Überschreitung eine
|
||||
// Benachrichtigung über CFG-02 aus (Akzeptanzkriterium 2), gedrosselt je
|
||||
// Regel+Zeitreihe (Akzeptanzkriterium 3).
|
||||
type Evaluator struct {
|
||||
rules *RuleStore
|
||||
debounce *debounceStore
|
||||
dispatcher *notify.Dispatcher
|
||||
interval time.Duration
|
||||
}
|
||||
|
||||
func NewEvaluator(rules *RuleStore, dispatcher *notify.Dispatcher, debouncePool *pgxpool.Pool, interval time.Duration) *Evaluator {
|
||||
if interval <= 0 {
|
||||
interval = DefaultDebounceInterval
|
||||
}
|
||||
return &Evaluator{
|
||||
rules: rules,
|
||||
debounce: &debounceStore{pool: debouncePool},
|
||||
dispatcher: dispatcher,
|
||||
interval: interval,
|
||||
}
|
||||
}
|
||||
|
||||
// FiredAlert beschreibt einen tatsächlich ausgelösten (nicht gedrosselten)
|
||||
// Alarm — fürs Testen/Logging, nicht Teil des öffentlichen Zustellwegs.
|
||||
type FiredAlert struct {
|
||||
RuleID string
|
||||
MetricName string
|
||||
Value float64
|
||||
Threshold float64
|
||||
Labels map[string]string
|
||||
Skipped bool // true, wenn wegen Drosselung NICHT tatsaechlich zugestellt
|
||||
}
|
||||
|
||||
// Evaluate prüft alle konfigurierten Regeln gegen families (Akzeptanzkriterium 1).
|
||||
// Für jede Zeitreihe, die eine Regel verletzt, wird — sofern nicht gedrosselt
|
||||
// — eine Benachrichtigung mit Metrik/Wert/Schwellwert/Labels (Tenant/Modul,
|
||||
// falls als Label vorhanden) über CFG-02 eingereiht (Akzeptanzkriterium 2).
|
||||
func (e *Evaluator) Evaluate(ctx context.Context, families []*dto.MetricFamily) ([]FiredAlert, error) {
|
||||
rules, err := e.rules.ListRules(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("regeln laden: %w", err)
|
||||
}
|
||||
if len(rules) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
byName := make(map[string]*dto.MetricFamily, len(families))
|
||||
for _, f := range families {
|
||||
if f.Name != nil {
|
||||
byName[*f.Name] = f
|
||||
}
|
||||
}
|
||||
|
||||
var fired []FiredAlert
|
||||
for _, rule := range rules {
|
||||
family, ok := byName[rule.MetricName]
|
||||
if !ok {
|
||||
continue // Metrik (noch) nicht vorhanden -> keine Aussage moeglich, kein Fehler.
|
||||
}
|
||||
|
||||
for _, m := range family.Metric {
|
||||
labels := labelMap(m)
|
||||
if !matchesFilter(labels, rule.LabelFilters) {
|
||||
continue
|
||||
}
|
||||
value, ok := metricValue(m)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
if !violates(rule, value) {
|
||||
continue
|
||||
}
|
||||
|
||||
ruleKey := ruleKeyFor(rule.ID, labels)
|
||||
allowed, err := e.debounce.shouldFire(ctx, ruleKey, e.interval.Seconds())
|
||||
if err != nil {
|
||||
return fired, fmt.Errorf("drosselung pruefen: %w", err)
|
||||
}
|
||||
|
||||
alert := FiredAlert{
|
||||
RuleID: rule.ID, MetricName: rule.MetricName, Value: value,
|
||||
Threshold: rule.Threshold, Labels: labels, Skipped: !allowed,
|
||||
}
|
||||
fired = append(fired, alert)
|
||||
|
||||
if !allowed {
|
||||
continue
|
||||
}
|
||||
|
||||
payload := map[string]any{
|
||||
"metric": rule.MetricName,
|
||||
"value": value,
|
||||
"threshold": rule.Threshold,
|
||||
"comparison": string(rule.Comparison),
|
||||
"description": rule.Description,
|
||||
"labels": labels,
|
||||
}
|
||||
if _, err := e.dispatcher.Enqueue(ctx, AlertChannel, rule.Recipient, payload); err != nil {
|
||||
return fired, fmt.Errorf("alarm einreihen: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
return fired, nil
|
||||
}
|
||||
|
||||
func violates(rule Rule, value float64) bool {
|
||||
switch rule.Comparison {
|
||||
case ComparisonGreaterThan:
|
||||
return value > rule.Threshold
|
||||
case ComparisonLessThan:
|
||||
return value < rule.Threshold
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func labelMap(m *dto.Metric) map[string]string {
|
||||
out := make(map[string]string, len(m.Label))
|
||||
for _, l := range m.Label {
|
||||
if l.Name != nil && l.Value != nil {
|
||||
out[*l.Name] = *l.Value
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func matchesFilter(labels, filter map[string]string) bool {
|
||||
for k, v := range filter {
|
||||
if labels[k] != v {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func metricValue(m *dto.Metric) (float64, bool) {
|
||||
switch {
|
||||
case m.Gauge != nil && m.Gauge.Value != nil:
|
||||
return *m.Gauge.Value, true
|
||||
case m.Counter != nil && m.Counter.Value != nil:
|
||||
return *m.Counter.Value, true
|
||||
case m.Untyped != nil && m.Untyped.Value != nil:
|
||||
return *m.Untyped.Value, true
|
||||
default:
|
||||
return 0, false
|
||||
}
|
||||
}
|
||||
|
||||
// ruleKeyFor macht die Drosselung unabhaengig je Regel UND je konkreter
|
||||
// Zeitreihe (z. B. verschiedene Tenants/Module derselben Metrik loesen
|
||||
// unabhaengig voneinander aus, siehe Migrationskommentar).
|
||||
func ruleKeyFor(ruleID string, labels map[string]string) string {
|
||||
keys := make([]string, 0, len(labels))
|
||||
for k := range labels {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
var b strings.Builder
|
||||
b.WriteString(ruleID)
|
||||
for _, k := range keys {
|
||||
b.WriteString("|")
|
||||
b.WriteString(k)
|
||||
b.WriteString("=")
|
||||
b.WriteString(labels[k])
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
@@ -0,0 +1,233 @@
|
||||
package alerting
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/notify"
|
||||
)
|
||||
|
||||
func setupTest(t *testing.T) (*RuleStore, *notify.Dispatcher, *pgxpool.Pool, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS alert_rules (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), metric_name TEXT NOT NULL,
|
||||
comparison TEXT NOT NULL CHECK (comparison IN ('gt','lt')), threshold DOUBLE PRECISION NOT NULL,
|
||||
label_filters JSONB NOT NULL DEFAULT '{}'::jsonb, recipient TEXT NOT NULL,
|
||||
description TEXT NOT NULL DEFAULT '', created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS alert_debounce_state (
|
||||
rule_key TEXT PRIMARY KEY, last_fired_at TIMESTAMPTZ NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS notification_jobs (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), channel TEXT NOT NULL, recipient TEXT NOT NULL,
|
||||
payload JSONB NOT NULL DEFAULT '{}'::jsonb, status TEXT NOT NULL DEFAULT 'pending', attempts INT NOT NULL DEFAULT 0,
|
||||
max_attempts INT NOT NULL DEFAULT 5, next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), last_error TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM alert_rules`)
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM alert_debounce_state`)
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM notification_jobs`)
|
||||
pool.Close()
|
||||
}
|
||||
return NewRuleStore(pool), notify.NewDispatcher(pool), pool, cleanup
|
||||
}
|
||||
|
||||
func gaugeFamily(name string, labels map[string]string, value float64) *dto.MetricFamily {
|
||||
pairs := make([]*dto.LabelPair, 0, len(labels))
|
||||
for k, v := range labels {
|
||||
k, v := k, v
|
||||
pairs = append(pairs, &dto.LabelPair{Name: &k, Value: &v})
|
||||
}
|
||||
n := name
|
||||
return &dto.MetricFamily{
|
||||
Name: &n,
|
||||
Metric: []*dto.Metric{
|
||||
{Label: pairs, Gauge: &dto.Gauge{Value: &value}},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 + 2 / Pruefung 1 + 2: Überschreitung löst eine
|
||||
// Benachrichtigung mit vollständigem Inhalt (Metrik/Tenant/Modul) aus.
|
||||
func TestEvaluate_FiresAlertOnThresholdExceeded(t *testing.T) {
|
||||
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
rule, err := rules.CreateRule(ctx, Rule{
|
||||
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||
Recipient: "ops@acme.example", Description: "Fehlerrate zu hoch",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("create rule: %v", err)
|
||||
}
|
||||
|
||||
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||
families := []*dto.MetricFamily{
|
||||
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme", "module": "dms"}, 0.12),
|
||||
}
|
||||
|
||||
fired, err := eval.Evaluate(ctx, families)
|
||||
if err != nil {
|
||||
t.Fatalf("evaluate: %v", err)
|
||||
}
|
||||
if len(fired) != 1 || fired[0].Skipped {
|
||||
t.Fatalf("erwartet genau 1 tatsaechlich ausgeloesten alarm, habe %+v", fired)
|
||||
}
|
||||
if fired[0].RuleID != rule.ID {
|
||||
t.Fatalf("rule id = %q, want %q", fired[0].RuleID, rule.ID)
|
||||
}
|
||||
|
||||
// Pruefung 2: Benachrichtigungsinhalt vollstaendig (Metrik/Tenant/Modul).
|
||||
var payloadJSON []byte
|
||||
if err := pool.QueryRow(ctx, `SELECT payload FROM notification_jobs LIMIT 1`).Scan(&payloadJSON); err != nil {
|
||||
t.Fatalf("notification_jobs lesen: %v", err)
|
||||
}
|
||||
payload := string(payloadJSON)
|
||||
for _, want := range []string{`"metric"`, `nexarch_core_error_rate`, `"tenant"`, `"acme"`, `"module"`, `"dms"`} {
|
||||
if !strings.Contains(payload, want) {
|
||||
t.Errorf("payload enthaelt nicht %q: %s", want, payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluate_DoesNotFireBelowThreshold(t *testing.T) {
|
||||
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := rules.CreateRule(ctx, Rule{
|
||||
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||
Recipient: "ops@acme.example",
|
||||
}); err != nil {
|
||||
t.Fatalf("create rule: %v", err)
|
||||
}
|
||||
|
||||
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||
families := []*dto.MetricFamily{
|
||||
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme"}, 0.01),
|
||||
}
|
||||
|
||||
fired, err := eval.Evaluate(ctx, families)
|
||||
if err != nil {
|
||||
t.Fatalf("evaluate: %v", err)
|
||||
}
|
||||
if len(fired) != 0 {
|
||||
t.Fatalf("erwartet keinen alarm unterhalb des schwellwerts, habe %+v", fired)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 / Pruefung 3: anhaltende Überschreitung erzeugt NICHT
|
||||
// bei jeder Messung eine neue Benachrichtigung.
|
||||
func TestEvaluate_DebouncesRepeatedFiring(t *testing.T) {
|
||||
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := rules.CreateRule(ctx, Rule{
|
||||
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||
Recipient: "ops@acme.example",
|
||||
}); err != nil {
|
||||
t.Fatalf("create rule: %v", err)
|
||||
}
|
||||
|
||||
// Langes Debounce-Intervall: der zweite Evaluate-Lauf (simuliert die
|
||||
// naechste Messung bei anhaltend ueberschrittenem Wert) darf keinen
|
||||
// weiteren Job einreihen.
|
||||
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||
families := []*dto.MetricFamily{
|
||||
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme"}, 0.5),
|
||||
}
|
||||
|
||||
first, err := eval.Evaluate(ctx, families)
|
||||
if err != nil {
|
||||
t.Fatalf("erster evaluate-lauf: %v", err)
|
||||
}
|
||||
if len(first) != 1 || first[0].Skipped {
|
||||
t.Fatalf("erster lauf haette feuern muessen, habe %+v", first)
|
||||
}
|
||||
|
||||
second, err := eval.Evaluate(ctx, families)
|
||||
if err != nil {
|
||||
t.Fatalf("zweiter evaluate-lauf: %v", err)
|
||||
}
|
||||
if len(second) != 1 || !second[0].Skipped {
|
||||
t.Fatalf("zweiter lauf haette gedrosselt werden muessen, habe %+v", second)
|
||||
}
|
||||
|
||||
var count int
|
||||
if err := pool.QueryRow(ctx, `SELECT count(*) FROM notification_jobs`).Scan(&count); err != nil {
|
||||
t.Fatalf("notification_jobs zaehlen: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("erwartet genau 1 eingereihten job trotz zwei ueberschreitenden messungen, habe %d", count)
|
||||
}
|
||||
}
|
||||
|
||||
// Verschiedene Zeitreihen derselben Regel (unterschiedlicher Tenant) werden
|
||||
// unabhaengig voneinander gedrosselt.
|
||||
func TestEvaluate_DebouncesIndependentlyPerLabelSet(t *testing.T) {
|
||||
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := rules.CreateRule(ctx, Rule{
|
||||
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||
Recipient: "ops@acme.example",
|
||||
}); err != nil {
|
||||
t.Fatalf("create rule: %v", err)
|
||||
}
|
||||
|
||||
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||
families := []*dto.MetricFamily{
|
||||
{
|
||||
Name: strPtr("nexarch_core_error_rate"),
|
||||
Metric: []*dto.Metric{
|
||||
metricWithLabel("tenant", "acme", 0.5),
|
||||
metricWithLabel("tenant", "beta", 0.6),
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
fired, err := eval.Evaluate(ctx, families)
|
||||
if err != nil {
|
||||
t.Fatalf("evaluate: %v", err)
|
||||
}
|
||||
if len(fired) != 2 {
|
||||
t.Fatalf("erwartet 2 unabhaengige alarme (verschiedene tenants), habe %d", len(fired))
|
||||
}
|
||||
for _, a := range fired {
|
||||
if a.Skipped {
|
||||
t.Fatalf("beide tenants sollten beim ersten mal feuern, habe %+v", a)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func metricWithLabel(name, value string, gaugeValue float64) *dto.Metric {
|
||||
n, v := name, value
|
||||
return &dto.Metric{Label: []*dto.LabelPair{{Name: &n, Value: &v}}, Gauge: &dto.Gauge{Value: &gaugeValue}}
|
||||
}
|
||||
|
||||
func strPtr(s string) *string { return &s }
|
||||
@@ -0,0 +1,132 @@
|
||||
// Package alerting implementiert Core OPS-05: schwellwertbasierte
|
||||
// Alarmierung auf den aus OPS-03 aggregierten Metriken, Zustellung über den
|
||||
// Core-Benachrichtigungs-Dispatcher (CFG-02). Der Alertmanager-Gedanke von
|
||||
// Prometheus/Grafana, aber auf das Nötigste reduziert (Schwellwert, Ziel,
|
||||
// Drosselung) — keine eigene Ausdruckssprache.
|
||||
package alerting
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// Comparison legt fest, ob ein Schwellwert nach oben oder unten überwacht
|
||||
// wird — bewusst nur zwei Operatoren, keine eigene Ausdruckssprache
|
||||
// (Ticket-Produkt-DNA).
|
||||
type Comparison string
|
||||
|
||||
const (
|
||||
ComparisonGreaterThan Comparison = "gt"
|
||||
ComparisonLessThan Comparison = "lt"
|
||||
)
|
||||
|
||||
// Rule ist eine Schwellwert-Regel auf einer beliebigen aggregierten Metrik
|
||||
// (Akzeptanzkriterium 1). LabelFilters schränkt optional auf bestimmte
|
||||
// Label-Werte ein (z. B. tenant/module), leer = alle Zeitreihen der Metrik.
|
||||
type Rule struct {
|
||||
ID string
|
||||
MetricName string
|
||||
Comparison Comparison
|
||||
Threshold float64
|
||||
LabelFilters map[string]string
|
||||
Recipient string
|
||||
Description string
|
||||
}
|
||||
|
||||
// RuleStore verwaltet Alert-Regeln in der zentralen Registry-DB.
|
||||
type RuleStore struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewRuleStore(pool *pgxpool.Pool) *RuleStore {
|
||||
return &RuleStore{pool: pool}
|
||||
}
|
||||
|
||||
// CreateRule legt eine neue Schwellwert-Regel an (Akzeptanzkriterium 1).
|
||||
func (s *RuleStore) CreateRule(ctx context.Context, r Rule) (Rule, error) {
|
||||
if r.MetricName == "" || r.Recipient == "" {
|
||||
return Rule{}, errors.New("alerting: metricName und recipient duerfen nicht leer sein")
|
||||
}
|
||||
if r.Comparison != ComparisonGreaterThan && r.Comparison != ComparisonLessThan {
|
||||
return Rule{}, fmt.Errorf("alerting: unbekannter comparison-operator %q", r.Comparison)
|
||||
}
|
||||
if r.LabelFilters == nil {
|
||||
r.LabelFilters = map[string]string{}
|
||||
}
|
||||
filtersJSON, err := json.Marshal(r.LabelFilters)
|
||||
if err != nil {
|
||||
return Rule{}, fmt.Errorf("label-filter serialisieren: %w", err)
|
||||
}
|
||||
|
||||
err = s.pool.QueryRow(ctx, `
|
||||
INSERT INTO alert_rules (metric_name, comparison, threshold, label_filters, recipient, description)
|
||||
VALUES ($1, $2, $3, $4, $5, $6)
|
||||
RETURNING id
|
||||
`, r.MetricName, string(r.Comparison), r.Threshold, filtersJSON, r.Recipient, r.Description).Scan(&r.ID)
|
||||
if err != nil {
|
||||
return Rule{}, fmt.Errorf("regel speichern: %w", err)
|
||||
}
|
||||
return r, nil
|
||||
}
|
||||
|
||||
// ListRules liefert alle konfigurierten Regeln — Grundlage für Evaluate.
|
||||
func (s *RuleStore) ListRules(ctx context.Context) ([]Rule, error) {
|
||||
rows, err := s.pool.Query(ctx, `
|
||||
SELECT id, metric_name, comparison, threshold, label_filters, recipient, description
|
||||
FROM alert_rules ORDER BY created_at
|
||||
`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("regeln auflisten: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []Rule
|
||||
for rows.Next() {
|
||||
var r Rule
|
||||
var comparison string
|
||||
var filtersJSON []byte
|
||||
if err := rows.Scan(&r.ID, &r.MetricName, &comparison, &r.Threshold, &filtersJSON, &r.Recipient, &r.Description); err != nil {
|
||||
return nil, fmt.Errorf("regel lesen: %w", err)
|
||||
}
|
||||
r.Comparison = Comparison(comparison)
|
||||
if err := json.Unmarshal(filtersJSON, &r.LabelFilters); err != nil {
|
||||
return nil, fmt.Errorf("label-filter lesen: %w", err)
|
||||
}
|
||||
out = append(out, r)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// DeleteRule entfernt eine Regel.
|
||||
func (s *RuleStore) DeleteRule(ctx context.Context, id string) error {
|
||||
_, err := s.pool.Exec(ctx, `DELETE FROM alert_rules WHERE id = $1`, id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("regel loeschen: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// debounceStore kapselt die Drosselungs-Zustandstabelle (Akzeptanzkriterium 3).
|
||||
type debounceStore struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
// shouldFire prueft, ob seit dem letzten Alarm fuer ruleKey mindestens
|
||||
// interval vergangen ist — atomar ueber eine bedingte UPDATE/INSERT-
|
||||
// Sequenz, damit zwei gleichzeitige Evaluate-Laeufe (z. B. bei mehreren
|
||||
// Core-Instanzen) nicht beide gleichzeitig alarmieren.
|
||||
func (d *debounceStore) shouldFire(ctx context.Context, ruleKey string, intervalSeconds float64) (bool, error) {
|
||||
tag, err := d.pool.Exec(ctx, `
|
||||
INSERT INTO alert_debounce_state (rule_key, last_fired_at) VALUES ($1, now())
|
||||
ON CONFLICT (rule_key) DO UPDATE SET last_fired_at = now()
|
||||
WHERE alert_debounce_state.last_fired_at <= now() - ($2 * interval '1 second')
|
||||
`, ruleKey, intervalSeconds)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("drosselungszustand pruefen: %w", err)
|
||||
}
|
||||
return tag.RowsAffected() == 1, nil
|
||||
}
|
||||
@@ -1,96 +0,0 @@
|
||||
// Package audit implementiert Core AUD-01: das zentrale, vom allgemeinen
|
||||
// Anwendungs-Log getrennte Audit-Datenmodell fuer sicherheits- und
|
||||
// compliancerelevante Ereignisse (wer, was, wann, an welchem Tenant).
|
||||
// Unveraenderlichkeit (Append-only) ist AUD-02, Export/Filter-API ist AUD-03
|
||||
// — dieses Paket liefert nur das Datenmodell und den EINEN zentralen
|
||||
// Schreibpfad (Akzeptanzkriterium 3).
|
||||
package audit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// SystemTenant ist der reservierte Tenant-Bezug fuer mandantenuebergreifende
|
||||
// Ereignisse (z.B. Superadmin-Aktionen) — es gibt bewusst KEINEN Weg, ein
|
||||
// Ereignis ganz ohne Tenant-Bezug zu schreiben (Akzeptanzkriterium 2).
|
||||
const SystemTenant = "system"
|
||||
|
||||
var ErrMissingTenant = errors.New("audit: tenant_slug darf nicht leer sein")
|
||||
var ErrMissingActor = errors.New("audit: actor darf nicht leer sein")
|
||||
var ErrMissingAction = errors.New("audit: action darf nicht leer sein")
|
||||
|
||||
// Event ist ein strukturiertes Audit-Ereignis (Akzeptanzkriterium 1: Akteur,
|
||||
// Aktion, Zielobjekt, Zeitpunkt, Tenant).
|
||||
type Event struct {
|
||||
TenantSlug string
|
||||
Actor string
|
||||
Action string
|
||||
Target string
|
||||
Metadata map[string]any
|
||||
OccurredAt time.Time
|
||||
}
|
||||
|
||||
// Log ist der EINE zentrale Schreibpfad fuer Audit-Ereignisse — es gibt
|
||||
// bewusst keine zweite Schreibmoeglichkeit, damit kein Handler versehentlich
|
||||
// direkt in audit_events schreibt und dabei die Validierung umgeht
|
||||
// (Akzeptanzkriterium 3).
|
||||
type Log struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewLog(pool *pgxpool.Pool) *Log {
|
||||
return &Log{pool: pool}
|
||||
}
|
||||
|
||||
// Record persistiert genau einen Audit-Eintrag. Fehlender Tenant-Bezug wird
|
||||
// bereits hier abgewiesen (klarer Fehler statt Constraint-Verletzung im
|
||||
// Normalfall) — die Datenbank-CHECK-Constraint aus der Migration ist die
|
||||
// zweite, unumgehbare Verteidigungslinie (Akzeptanzkriterium 2 / Pruefung 2).
|
||||
func (l *Log) Record(ctx context.Context, e Event) error {
|
||||
if e.TenantSlug == "" {
|
||||
return ErrMissingTenant
|
||||
}
|
||||
if e.Actor == "" {
|
||||
return ErrMissingActor
|
||||
}
|
||||
if e.Action == "" {
|
||||
return ErrMissingAction
|
||||
}
|
||||
if e.Metadata == nil {
|
||||
e.Metadata = map[string]any{}
|
||||
}
|
||||
metadataJSON, err := json.Marshal(e.Metadata)
|
||||
if err != nil {
|
||||
return fmt.Errorf("metadaten serialisieren: %w", err)
|
||||
}
|
||||
if e.OccurredAt.IsZero() {
|
||||
e.OccurredAt = time.Now()
|
||||
}
|
||||
|
||||
_, err = l.pool.Exec(ctx, `
|
||||
INSERT INTO audit_events (occurred_at, tenant_slug, actor, action, target, metadata)
|
||||
VALUES ($1, $2, $3, $4, $5, $6)
|
||||
`, e.OccurredAt, e.TenantSlug, e.Actor, e.Action, e.Target, metadataJSON)
|
||||
if err != nil {
|
||||
return fmt.Errorf("audit-ereignis schreiben: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CountByTenant ist eine schlanke Lesehilfe fuer Tests/Diagnose — die
|
||||
// eigentliche Filter-/Export-API ist AUD-03, hier bewusst nicht vorgezogen.
|
||||
func (l *Log) CountByTenant(ctx context.Context, tenantSlug string) (int, error) {
|
||||
var n int
|
||||
if err := l.pool.QueryRow(ctx, `
|
||||
SELECT count(*) FROM audit_events WHERE tenant_slug = $1
|
||||
`, tenantSlug).Scan(&n); err != nil {
|
||||
return 0, fmt.Errorf("audit-ereignisse zaehlen: %w", err)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
@@ -1,132 +0,0 @@
|
||||
package audit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupAuditTest(t *testing.T) (*Log, *pgxpool.Pool, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS audit_events (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
occurred_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
tenant_slug TEXT NOT NULL CHECK (tenant_slug <> ''),
|
||||
actor TEXT NOT NULL CHECK (actor <> ''),
|
||||
action TEXT NOT NULL CHECK (action <> ''),
|
||||
target TEXT NOT NULL,
|
||||
metadata JSONB NOT NULL DEFAULT '{}'::jsonb
|
||||
)`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM audit_events WHERE tenant_slug LIKE 'test\_%' ESCAPE '\' OR tenant_slug = $1`, SystemTenant)
|
||||
pool.Close()
|
||||
}
|
||||
return NewLog(pool), pool, cleanup
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 1: ein sicherheitsrelevanter Vorgang
|
||||
// (hier: fehlgeschlagener Login) erzeugt zuverlaessig genau einen Eintrag.
|
||||
func TestRecord_PersistsExactlyOneEventPerSecurityIncident(t *testing.T) {
|
||||
log, pool, cleanup := setupAuditTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
err := log.Record(ctx, Event{
|
||||
TenantSlug: "test_acme",
|
||||
Actor: "alice@example.com",
|
||||
Action: "iam.login_failed",
|
||||
Target: "user:alice@example.com",
|
||||
Metadata: map[string]any{"reason": "falsches passwort"},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("record: %v", err)
|
||||
}
|
||||
|
||||
count, err := log.CountByTenant(ctx, "test_acme")
|
||||
if err != nil {
|
||||
t.Fatalf("count: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Fatalf("erwartet genau 1 audit-eintrag, habe %d", count)
|
||||
}
|
||||
|
||||
var actor, action, target string
|
||||
if err := pool.QueryRow(ctx, `
|
||||
SELECT actor, action, target FROM audit_events WHERE tenant_slug = 'test_acme'
|
||||
`).Scan(&actor, &action, &target); err != nil {
|
||||
t.Fatalf("eintrag lesen: %v", err)
|
||||
}
|
||||
if actor != "alice@example.com" || action != "iam.login_failed" || target != "user:alice@example.com" {
|
||||
t.Fatalf("eintrag unerwartet: actor=%q action=%q target=%q", actor, action, target)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2 (App-Ebene): fehlender Tenant-Bezug wird
|
||||
// bereits vom zentralen Schreibpfad abgewiesen.
|
||||
func TestRecord_RejectsMissingTenant(t *testing.T) {
|
||||
log, _, cleanup := setupAuditTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
err := log.Record(ctx, Event{TenantSlug: "", Actor: "alice", Action: "irgendwas"})
|
||||
if !errors.Is(err, ErrMissingTenant) {
|
||||
t.Fatalf("erwartet ErrMissingTenant, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2 (DB-Ebene): selbst ein direkter INSERT,
|
||||
// der Log.Record umgeht, wird durch die CHECK-Constraint verhindert — der
|
||||
// Schutz haengt nicht allein von der Go-Validierung ab.
|
||||
func TestConstraint_RejectsMissingTenantAtDatabaseLevel(t *testing.T) {
|
||||
_, pool, cleanup := setupAuditTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
_, err := pool.Exec(ctx, `
|
||||
INSERT INTO audit_events (tenant_slug, actor, action, target)
|
||||
VALUES ('', 'alice', 'irgendwas', 'ziel')
|
||||
`)
|
||||
if err == nil {
|
||||
t.Fatal("erwartet fehler durch CHECK-constraint bei leerem tenant_slug, habe nil")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecord_RejectsMissingActorAndAction(t *testing.T) {
|
||||
log, _, cleanup := setupAuditTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if err := log.Record(ctx, Event{TenantSlug: "test_acme", Actor: "", Action: "x"}); !errors.Is(err, ErrMissingActor) {
|
||||
t.Fatalf("erwartet ErrMissingActor, habe %v", err)
|
||||
}
|
||||
if err := log.Record(ctx, Event{TenantSlug: "test_acme", Actor: "alice", Action: ""}); !errors.Is(err, ErrMissingAction) {
|
||||
t.Fatalf("erwartet ErrMissingAction, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecord_SystemTenantForCrossTenantEvents(t *testing.T) {
|
||||
log, _, cleanup := setupAuditTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if err := log.Record(ctx, Event{TenantSlug: SystemTenant, Actor: "superadmin", Action: "tenant.provisioned", Target: "tenant:acme"}); err != nil {
|
||||
t.Fatalf("record mit SystemTenant: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package cfgservice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// DefaultCacheTTL ist die dokumentierte Cache-Invalidierungszeit
|
||||
// (Akzeptanzkriterium 2 / Pruefung 2 in diesem Ticket bezieht sich auf die
|
||||
// Aenderungsnachvollziehbarkeit — die Cache-Frist selbst folgt demselben
|
||||
// Muster wie internal/flag.DefaultCacheTTL).
|
||||
const DefaultCacheTTL = 5 * time.Second
|
||||
|
||||
type cacheEntry struct {
|
||||
value Value
|
||||
expiresAt time.Time
|
||||
}
|
||||
|
||||
// Service ist die Leseseite mit Vorrangregel (Akzeptanzkriterium 1:
|
||||
// Tenant-Override vor Global-Default) und lokalem TTL-Cache.
|
||||
type Service struct {
|
||||
store *Store
|
||||
ttl time.Duration
|
||||
|
||||
mu sync.RWMutex
|
||||
cache map[string]cacheEntry // Schluessel: key + "\x00" + tenantSlug
|
||||
}
|
||||
|
||||
func NewService(store *Store, ttl time.Duration) *Service {
|
||||
if ttl <= 0 {
|
||||
ttl = DefaultCacheTTL
|
||||
}
|
||||
return &Service{store: store, ttl: ttl, cache: make(map[string]cacheEntry)}
|
||||
}
|
||||
|
||||
func cacheKey(key, tenantSlug string) string {
|
||||
return key + "\x00" + tenantSlug
|
||||
}
|
||||
|
||||
// Resolve liefert den Konfigurationswert fuer einen Tenant: ein
|
||||
// Tenant-spezifischer Override hat Vorrang vor dem globalen Default
|
||||
// (Akzeptanzkriterium 1 / Pruefung 1). tenantSlug == "" wertet nur den
|
||||
// globalen Wert aus.
|
||||
func (s *Service) Resolve(ctx context.Context, tenantSlug, key string) (Value, error) {
|
||||
ck := cacheKey(key, tenantSlug)
|
||||
|
||||
s.mu.RLock()
|
||||
entry, exists := s.cache[ck]
|
||||
fresh := exists && time.Now().Before(entry.expiresAt)
|
||||
s.mu.RUnlock()
|
||||
if fresh {
|
||||
return entry.value, nil
|
||||
}
|
||||
|
||||
v, err := s.resolveUncached(ctx, tenantSlug, key)
|
||||
if err != nil {
|
||||
return Value{}, err
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
s.cache[ck] = cacheEntry{value: v, expiresAt: time.Now().Add(s.ttl)}
|
||||
s.mu.Unlock()
|
||||
return v, nil
|
||||
}
|
||||
|
||||
func (s *Service) resolveUncached(ctx context.Context, tenantSlug, key string) (Value, error) {
|
||||
if tenantSlug != "" {
|
||||
v, err := s.store.Get(ctx, key, tenantSlug)
|
||||
if err == nil {
|
||||
return v, nil
|
||||
}
|
||||
if !errors.Is(err, ErrNotFound) {
|
||||
return Value{}, err
|
||||
}
|
||||
}
|
||||
return s.store.Get(ctx, key, GlobalScope)
|
||||
}
|
||||
|
||||
// Invalidate erzwingt beim naechsten Resolve-Aufruf ein sofortiges Neuladen
|
||||
// fuer einen bestimmten (key, tenantSlug) statt auf den TTL-Ablauf zu warten
|
||||
// — analog internal/flag.Service.Invalidate.
|
||||
func (s *Service) Invalidate(key, tenantSlug string) {
|
||||
s.mu.Lock()
|
||||
delete(s.cache, cacheKey(key, tenantSlug))
|
||||
s.mu.Unlock()
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
package cfgservice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 1: Tenant-Override hat Vorrang vor
|
||||
// Global-Default, automatisiert getestet.
|
||||
func TestService_TenantOverrideTakesPrecedenceOverGlobal(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := store.Set(ctx, "test_precedence_key", GlobalScope, "global-wert"); err != nil {
|
||||
t.Fatalf("set global: %v", err)
|
||||
}
|
||||
if _, err := store.Set(ctx, "test_precedence_key", "test_acme", "tenant-wert"); err != nil {
|
||||
t.Fatalf("set tenant: %v", err)
|
||||
}
|
||||
|
||||
svc := NewService(store, time.Hour)
|
||||
|
||||
got, err := svc.Resolve(ctx, "test_acme", "test_precedence_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve mit override: %v", err)
|
||||
}
|
||||
if got.Value != "tenant-wert" {
|
||||
t.Fatalf("erwartet tenant-override, habe %q", got.Value)
|
||||
}
|
||||
|
||||
gotOther, err := svc.Resolve(ctx, "test_anderer_tenant", "test_precedence_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve ohne override: %v", err)
|
||||
}
|
||||
if gotOther.Value != "global-wert" {
|
||||
t.Fatalf("erwartet global-default fuer tenant ohne override, habe %q", gotOther.Value)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: Cache-Invalidierung nach
|
||||
// Konfigurationsaenderung innerhalb dokumentierter Zeit gemessen.
|
||||
func TestService_CacheInvalidationTiming(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
const ttl = 150 * time.Millisecond
|
||||
if _, err := store.Set(ctx, "test_ttl_key", GlobalScope, "alt"); err != nil {
|
||||
t.Fatalf("set: %v", err)
|
||||
}
|
||||
svc := NewService(store, ttl)
|
||||
|
||||
v, err := svc.Resolve(ctx, "", "test_ttl_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve: %v", err)
|
||||
}
|
||||
if v.Value != "alt" {
|
||||
t.Fatalf("erwartet 'alt', habe %q", v.Value)
|
||||
}
|
||||
|
||||
changedAt := time.Now()
|
||||
if _, err := store.Set(ctx, "test_ttl_key", GlobalScope, "neu"); err != nil {
|
||||
t.Fatalf("set: %v", err)
|
||||
}
|
||||
|
||||
v, err = svc.Resolve(ctx, "", "test_ttl_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve direkt nach aenderung: %v", err)
|
||||
}
|
||||
if v.Value != "alt" {
|
||||
t.Fatalf("cache haette den alten wert liefern sollen, habe %q", v.Value)
|
||||
}
|
||||
|
||||
deadline := changedAt.Add(ttl + 100*time.Millisecond)
|
||||
for time.Now().Before(deadline) {
|
||||
v, err := svc.Resolve(ctx, "", "test_ttl_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve: %v", err)
|
||||
}
|
||||
if v.Value == "neu" {
|
||||
t.Logf("aenderung wurde nach %s wirksam (ziel: innerhalb %s + toleranz)", time.Since(changedAt), ttl)
|
||||
return
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
t.Fatalf("aenderung wurde nicht innerhalb von %s wirksam", deadline.Sub(changedAt))
|
||||
}
|
||||
|
||||
func TestService_InvalidateForcesImmediateRefresh(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := store.Set(ctx, "test_invalidate_key", GlobalScope, "alt"); err != nil {
|
||||
t.Fatalf("set: %v", err)
|
||||
}
|
||||
svc := NewService(store, time.Hour)
|
||||
_, _ = svc.Resolve(ctx, "", "test_invalidate_key")
|
||||
|
||||
if _, err := store.Set(ctx, "test_invalidate_key", GlobalScope, "neu"); err != nil {
|
||||
t.Fatalf("set: %v", err)
|
||||
}
|
||||
svc.Invalidate("test_invalidate_key", "")
|
||||
|
||||
v, err := svc.Resolve(ctx, "", "test_invalidate_key")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve: %v", err)
|
||||
}
|
||||
if v.Value != "neu" {
|
||||
t.Fatalf("erwartet sofort sichtbaren neuen wert nach Invalidate, habe %q", v.Value)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,129 @@
|
||||
// Package cfgservice implementiert Core CFG-01: den zentralen Dienst fuer
|
||||
// globale und tenant-spezifische Konfigurationswerte mit Versionierung und
|
||||
// Cache-Invalidierung. Andere Module lesen Konfiguration AUSSCHLIESSLICH
|
||||
// ueber dieses Paket (Akzeptanzkriterium 3), niemals ueber eigene Tabellen.
|
||||
package cfgservice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// GlobalScope ist der reservierte Scope-Wert fuer globale Defaults — jeder
|
||||
// andere Scope-Wert ist ein Tenant-Slug (Akzeptanzkriterium 1).
|
||||
const GlobalScope = "global"
|
||||
|
||||
var ErrNotFound = errors.New("cfgservice: kein wert fuer diesen key gefunden")
|
||||
|
||||
type Value struct {
|
||||
Key string
|
||||
Scope string
|
||||
Value string
|
||||
Version int
|
||||
}
|
||||
|
||||
type HistoryEntry struct {
|
||||
Key string
|
||||
Scope string
|
||||
Value string
|
||||
Version int
|
||||
}
|
||||
|
||||
// Store ist die Schreib-/Verwaltungsseite. Set schreibt IMMER sowohl den
|
||||
// aktuellen Stand (config_values) als auch einen Historieneintrag
|
||||
// (config_value_history) in derselben Transaktion — eine Aenderung ohne
|
||||
// Versionshistorie ist strukturell ausgeschlossen (Akzeptanzkriterium 2).
|
||||
type Store struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewStore(pool *pgxpool.Pool) *Store {
|
||||
return &Store{pool: pool}
|
||||
}
|
||||
|
||||
// Set schreibt einen neuen Wert fuer (key, scope) und erhoeht die Version um 1
|
||||
// (Version 1 bei erstmaligem Setzen).
|
||||
func (s *Store) Set(ctx context.Context, key, scope, value string) (Value, error) {
|
||||
if scope == "" {
|
||||
return Value{}, errors.New("cfgservice: scope darf nicht leer sein")
|
||||
}
|
||||
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return Value{}, fmt.Errorf("transaktion starten: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
var currentVersion int
|
||||
err = tx.QueryRow(ctx, `SELECT version FROM config_values WHERE key = $1 AND scope = $2`, key, scope).Scan(¤tVersion)
|
||||
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
|
||||
return Value{}, fmt.Errorf("aktuelle version lesen: %w", err)
|
||||
}
|
||||
newVersion := currentVersion + 1
|
||||
|
||||
if _, err := tx.Exec(ctx, `
|
||||
INSERT INTO config_values (key, scope, value, version, updated_at)
|
||||
VALUES ($1, $2, $3, $4, now())
|
||||
ON CONFLICT (key, scope) DO UPDATE SET value = $3, version = $4, updated_at = now()
|
||||
`, key, scope, value, newVersion); err != nil {
|
||||
return Value{}, fmt.Errorf("wert speichern: %w", err)
|
||||
}
|
||||
|
||||
if _, err := tx.Exec(ctx, `
|
||||
INSERT INTO config_value_history (key, scope, value, version, changed_at)
|
||||
VALUES ($1, $2, $3, $4, now())
|
||||
`, key, scope, value, newVersion); err != nil {
|
||||
return Value{}, fmt.Errorf("historie schreiben: %w", err)
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return Value{}, fmt.Errorf("transaktion committen: %w", err)
|
||||
}
|
||||
|
||||
return Value{Key: key, Scope: scope, Value: value, Version: newVersion}, nil
|
||||
}
|
||||
|
||||
// Get liefert den Wert fuer GENAU EINEN Scope (kein Vorrang-Fallback) — die
|
||||
// Vorrangregel (Tenant vor Global) lebt bewusst in Service.Resolve, damit
|
||||
// Store rein CRUD bleibt.
|
||||
func (s *Store) Get(ctx context.Context, key, scope string) (Value, error) {
|
||||
var v Value
|
||||
v.Key, v.Scope = key, scope
|
||||
err := s.pool.QueryRow(ctx, `
|
||||
SELECT value, version FROM config_values WHERE key = $1 AND scope = $2
|
||||
`, key, scope).Scan(&v.Value, &v.Version)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return Value{}, ErrNotFound
|
||||
}
|
||||
return Value{}, fmt.Errorf("wert lesen: %w", err)
|
||||
}
|
||||
return v, nil
|
||||
}
|
||||
|
||||
// History liefert die vollstaendige Versionshistorie eines (key, scope) in
|
||||
// aufsteigender Reihenfolge (Akzeptanzkriterium 2 / Pruefung 3).
|
||||
func (s *Store) History(ctx context.Context, key, scope string) ([]HistoryEntry, error) {
|
||||
rows, err := s.pool.Query(ctx, `
|
||||
SELECT key, scope, value, version FROM config_value_history
|
||||
WHERE key = $1 AND scope = $2 ORDER BY version
|
||||
`, key, scope)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("historie abfragen: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []HistoryEntry
|
||||
for rows.Next() {
|
||||
var h HistoryEntry
|
||||
if err := rows.Scan(&h.Key, &h.Scope, &h.Value, &h.Version); err != nil {
|
||||
return nil, fmt.Errorf("historieneintrag lesen: %w", err)
|
||||
}
|
||||
out = append(out, h)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
package cfgservice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupStoreTest(t *testing.T) (*Store, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS config_values (
|
||||
key TEXT NOT NULL,
|
||||
scope TEXT NOT NULL CHECK (scope <> ''),
|
||||
value TEXT NOT NULL,
|
||||
version INT NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (key, scope)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS config_value_history (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
key TEXT NOT NULL,
|
||||
scope TEXT NOT NULL,
|
||||
value TEXT NOT NULL,
|
||||
version INT NOT NULL,
|
||||
changed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM config_value_history WHERE key LIKE 'test\_%' ESCAPE '\'`)
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM config_values WHERE key LIKE 'test\_%' ESCAPE '\'`)
|
||||
pool.Close()
|
||||
}
|
||||
return NewStore(pool), cleanup
|
||||
}
|
||||
|
||||
func TestStore_SetIncrementsVersion(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
v1, err := store.Set(ctx, "test_key", GlobalScope, "erster-wert")
|
||||
if err != nil {
|
||||
t.Fatalf("set 1: %v", err)
|
||||
}
|
||||
if v1.Version != 1 {
|
||||
t.Fatalf("erwartet version 1, habe %d", v1.Version)
|
||||
}
|
||||
|
||||
v2, err := store.Set(ctx, "test_key", GlobalScope, "zweiter-wert")
|
||||
if err != nil {
|
||||
t.Fatalf("set 2: %v", err)
|
||||
}
|
||||
if v2.Version != 2 {
|
||||
t.Fatalf("erwartet version 2, habe %d", v2.Version)
|
||||
}
|
||||
|
||||
got, err := store.Get(ctx, "test_key", GlobalScope)
|
||||
if err != nil {
|
||||
t.Fatalf("get: %v", err)
|
||||
}
|
||||
if got.Value != "zweiter-wert" || got.Version != 2 {
|
||||
t.Fatalf("aktueller wert unerwartet: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 3: Versionierungshistorie ueber mehrere
|
||||
// Aenderungen hinweg nachvollzogen.
|
||||
func TestStore_HistoryTracksAllChanges(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
values := []string{"v1", "v2", "v3"}
|
||||
for _, v := range values {
|
||||
if _, err := store.Set(ctx, "test_history_key", GlobalScope, v); err != nil {
|
||||
t.Fatalf("set %q: %v", v, err)
|
||||
}
|
||||
}
|
||||
|
||||
history, err := store.History(ctx, "test_history_key", GlobalScope)
|
||||
if err != nil {
|
||||
t.Fatalf("history: %v", err)
|
||||
}
|
||||
if len(history) != 3 {
|
||||
t.Fatalf("erwartet 3 historieneintraege, habe %d", len(history))
|
||||
}
|
||||
for i, h := range history {
|
||||
if h.Version != i+1 || h.Value != values[i] {
|
||||
t.Fatalf("historieneintrag[%d] unerwartet: %+v", i, h)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_GetUnknownKeyReturnsNotFound(t *testing.T) {
|
||||
store, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := store.Get(ctx, "test_nie_gesetzt", GlobalScope); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("erwartet ErrNotFound, habe %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package health
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// DatabaseChecker prueft die tatsaechliche Erreichbarkeit der Datenbank
|
||||
// (Ping) — nicht nur, ob der Pool existiert.
|
||||
func DatabaseChecker(pool *pgxpool.Pool) CheckerFunc {
|
||||
return func(ctx context.Context) error {
|
||||
return pool.Ping(ctx)
|
||||
}
|
||||
}
|
||||
|
||||
// QueueChecker prueft, dass die Postgres-basierte Job-Queue (siehe CFG-02)
|
||||
// tatsaechlich abfragbar ist — eine eigene, benannte Abhaengigkeit neben der
|
||||
// reinen DB-Erreichbarkeit (Akzeptanzkriterium 1).
|
||||
func QueueChecker(pool *pgxpool.Pool) CheckerFunc {
|
||||
return func(ctx context.Context) error {
|
||||
_, err := pool.Exec(ctx, `SELECT 1`)
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package health
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
// LivenessHandler beantwortet IMMER "lebt", solange der Prozess ueberhaupt
|
||||
// HTTP-Anfragen verarbeiten kann — prueft bewusst KEINE externen
|
||||
// Abhaengigkeiten (Akzeptanzkriterium 2: Liveness und Readiness getrennt).
|
||||
// Ein Datenbankausfall darf die Liveness nicht auf "tot" setzen, sonst
|
||||
// wuerde eine Orchestrierung (z.B. systemd/Kubernetes) den Prozess grundlos
|
||||
// neu starten, obwohl nur eine Abhaengigkeit ausgefallen ist.
|
||||
func LivenessHandler(w http.ResponseWriter, r *http.Request) {
|
||||
writeStatus(w, http.StatusOK, map[string]any{"status": "alive"})
|
||||
}
|
||||
|
||||
// ReadinessHandler prueft ALLE registrierten Abhaengigkeiten
|
||||
// (Akzeptanzkriterium 1) und liefert 503, sobald eine davon fehlschlaegt
|
||||
// (Akzeptanzkriterium 3) — unterscheidet sich damit nachweislich von
|
||||
// LivenessHandler im Fehlerfall (Akzeptanzkriterium 2 / Pruefung 3).
|
||||
func (r *Registry) ReadinessHandler() http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, req *http.Request) {
|
||||
ready, results := r.CheckAll(req.Context())
|
||||
|
||||
body := map[string]any{
|
||||
"status": statusText(ready),
|
||||
"checks": results,
|
||||
}
|
||||
status := http.StatusOK
|
||||
if !ready {
|
||||
status = http.StatusServiceUnavailable
|
||||
}
|
||||
writeStatus(w, status, body)
|
||||
}
|
||||
}
|
||||
|
||||
func statusText(ready bool) string {
|
||||
if ready {
|
||||
return "ready"
|
||||
}
|
||||
return "not_ready"
|
||||
}
|
||||
|
||||
func writeStatus(w http.ResponseWriter, status int, body map[string]any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(status)
|
||||
_ = json.NewEncoder(w).Encode(body)
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
// Package health implementiert Core OPS-01: Health-/Readiness-Endpunkte, die
|
||||
// echte Abhaengigkeiten (DB, Job-Queue) statt nur den Prozessstatus pruefen
|
||||
// — wiederverwendbar von Core UND jedem registrierten Modul (siehe API-02),
|
||||
// nicht nur von Core selbst.
|
||||
package health
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Checker prueft EINE Abhaengigkeit (z.B. Datenbank, Job-Queue).
|
||||
type Checker interface {
|
||||
Check(ctx context.Context) error
|
||||
}
|
||||
|
||||
type CheckerFunc func(ctx context.Context) error
|
||||
|
||||
func (f CheckerFunc) Check(ctx context.Context) error { return f(ctx) }
|
||||
|
||||
// DefaultCheckTimeout begrenzt, wie lange EIN einzelner Check maximal
|
||||
// dauern darf, bevor er als fehlgeschlagen gilt — verhindert, dass ein
|
||||
// haengender Check den gesamten Readiness-Endpunkt blockiert
|
||||
// (Akzeptanzkriterium 2 / Pruefung 2: Antwort innerhalb definierter Zeit).
|
||||
const DefaultCheckTimeout = 2 * time.Second
|
||||
|
||||
// Registry haelt alle benannten Checks eines Dienstes.
|
||||
type Registry struct {
|
||||
checks map[string]Checker
|
||||
timeout time.Duration
|
||||
}
|
||||
|
||||
func NewRegistry() *Registry {
|
||||
return &Registry{checks: make(map[string]Checker), timeout: DefaultCheckTimeout}
|
||||
}
|
||||
|
||||
func (r *Registry) WithTimeout(d time.Duration) *Registry {
|
||||
return &Registry{checks: r.checks, timeout: d}
|
||||
}
|
||||
|
||||
// Register fuegt einen benannten Check hinzu (z.B. "database", "queue").
|
||||
func (r *Registry) Register(name string, c Checker) {
|
||||
r.checks[name] = c
|
||||
}
|
||||
|
||||
// Result ist der Ausgang eines einzelnen Checks.
|
||||
type Result struct {
|
||||
OK bool
|
||||
Error string
|
||||
}
|
||||
|
||||
// CheckAll fuehrt alle registrierten Checks NEBENLAEUFIG mit je eigenem
|
||||
// Timeout aus (Akzeptanzkriterium 1: echte Abhaengigkeiten statt Prozess-
|
||||
// status) und liefert ready=false, sobald irgendein Check fehlschlaegt
|
||||
// (Akzeptanzkriterium 3: ein Ausfall wird sichtbar).
|
||||
func (r *Registry) CheckAll(ctx context.Context) (ready bool, results map[string]Result) {
|
||||
type namedResult struct {
|
||||
name string
|
||||
result Result
|
||||
}
|
||||
ch := make(chan namedResult, len(r.checks))
|
||||
|
||||
for name, checker := range r.checks {
|
||||
go func(name string, checker Checker) {
|
||||
checkCtx, cancel := context.WithTimeout(ctx, r.timeout)
|
||||
defer cancel()
|
||||
err := checker.Check(checkCtx)
|
||||
if err != nil {
|
||||
ch <- namedResult{name, Result{OK: false, Error: err.Error()}}
|
||||
return
|
||||
}
|
||||
ch <- namedResult{name, Result{OK: true}}
|
||||
}(name, checker)
|
||||
}
|
||||
|
||||
results = make(map[string]Result, len(r.checks))
|
||||
ready = true
|
||||
for i := 0; i < len(r.checks); i++ {
|
||||
nr := <-ch
|
||||
results[nr.name] = nr.result
|
||||
if !nr.result.OK {
|
||||
ready = false
|
||||
}
|
||||
}
|
||||
return ready, results
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
package health
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 1: simulierter Datenbankausfall fuehrt zu
|
||||
// "nicht bereit".
|
||||
func TestReadinessHandler_ReportsNotReadyOnDatabaseFailure(t *testing.T) {
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
// Datenbankausfall simulieren: Pool sofort schliessen, bevor der Check laeuft.
|
||||
pool.Close()
|
||||
|
||||
reg := NewRegistry()
|
||||
reg.Register("database", DatabaseChecker(pool))
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
reg.ReadinessHandler()(rec, req)
|
||||
|
||||
if rec.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("status = %d, want 503 bei db-ausfall", rec.Code)
|
||||
}
|
||||
|
||||
var body struct {
|
||||
Status string `json:"status"`
|
||||
Checks map[string]interface{} `json:"checks"`
|
||||
}
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil {
|
||||
t.Fatalf("body parsen: %v", err)
|
||||
}
|
||||
if body.Status != "not_ready" {
|
||||
t.Fatalf("status-feld = %q, want not_ready", body.Status)
|
||||
}
|
||||
if _, ok := body.Checks["database"]; !ok {
|
||||
t.Fatal("erwartet 'database' im checks-ergebnis")
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadinessHandler_ReportsReadyWhenAllChecksPass(t *testing.T) {
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
defer pool.Close()
|
||||
|
||||
reg := NewRegistry()
|
||||
reg.Register("database", DatabaseChecker(pool))
|
||||
reg.Register("queue", QueueChecker(pool))
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
reg.ReadinessHandler()(rec, req)
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200 bei funktionierenden abhaengigkeiten", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 3: Liveness und Readiness unterscheiden
|
||||
// sich nachweislich im Fehlerfall.
|
||||
func TestLivenessAndReadiness_DifferOnDatabaseFailure(t *testing.T) {
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
pool.Close() // db-ausfall simulieren
|
||||
|
||||
reg := NewRegistry()
|
||||
reg.Register("database", DatabaseChecker(pool))
|
||||
|
||||
livenessRec := httptest.NewRecorder()
|
||||
LivenessHandler(livenessRec, httptest.NewRequest(http.MethodGet, "/livez", nil))
|
||||
if livenessRec.Code != http.StatusOK {
|
||||
t.Fatalf("liveness status = %d, want 200 trotz db-ausfall (liveness prueft keine abhaengigkeiten)", livenessRec.Code)
|
||||
}
|
||||
|
||||
readinessRec := httptest.NewRecorder()
|
||||
reg.ReadinessHandler()(readinessRec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
|
||||
if readinessRec.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("readiness status = %d, want 503 bei db-ausfall", readinessRec.Code)
|
||||
}
|
||||
|
||||
if livenessRec.Code == readinessRec.Code {
|
||||
t.Fatal("liveness und readiness sollten sich im db-ausfall-fall unterscheiden")
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: Health-Endpunkt antwortet auch bei
|
||||
// haengendem Check innerhalb definierter Zeit (Timeout begrenzt die Dauer).
|
||||
func TestReadinessHandler_RespondsWithinTimeoutEvenWithHangingCheck(t *testing.T) {
|
||||
reg := NewRegistry().WithTimeout(50 * time.Millisecond)
|
||||
reg.Register("haengender_dienst", CheckerFunc(func(ctx context.Context) error {
|
||||
select {
|
||||
case <-time.After(10 * time.Second): // wuerde ohne timeout ewig blockieren
|
||||
return nil
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}))
|
||||
|
||||
start := time.Now()
|
||||
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
reg.ReadinessHandler()(rec, req)
|
||||
elapsed := time.Since(start)
|
||||
|
||||
if elapsed > time.Second {
|
||||
t.Fatalf("readiness handler brauchte %s, erwartet deutlich unter 1s durch timeout", elapsed)
|
||||
}
|
||||
if rec.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("status = %d, want 503 fuer haengenden/timeout-check", rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCheckAll_MultipleChecksRunConcurrently(t *testing.T) {
|
||||
reg := NewRegistry().WithTimeout(time.Second)
|
||||
reg.Register("a", CheckerFunc(func(ctx context.Context) error { return nil }))
|
||||
reg.Register("b", CheckerFunc(func(ctx context.Context) error { return errors.New("kaputt") }))
|
||||
|
||||
ready, results := reg.CheckAll(context.Background())
|
||||
if ready {
|
||||
t.Fatal("erwartet ready=false, da 'b' fehlschlaegt")
|
||||
}
|
||||
if !results["a"].OK {
|
||||
t.Fatalf("erwartet 'a' ok, habe %+v", results["a"])
|
||||
}
|
||||
if results["b"].OK || results["b"].Error == "" {
|
||||
t.Fatalf("erwartet 'b' fehlgeschlagen mit fehlertext, habe %+v", results["b"])
|
||||
}
|
||||
}
|
||||
@@ -1,106 +0,0 @@
|
||||
// Package license implementiert Core LIC-01: Lizenzmodell je Tenant (Plan,
|
||||
// Modul-Umfang, Laufzeit) und die kryptographische Pruefung signierter
|
||||
// Lizenzschluessel. Feature-Flag-AUSWERTUNG zur Laufzeit (LIC-02) und die
|
||||
// Verwaltungsoberflaeche (LIC-04) sind ausdruecklich nicht Teil dieses Pakets
|
||||
// — hier geht es nur um Ausstellung/Validierung/Persistenz (Unleash-Vorbild:
|
||||
// klare Trennung Flag-Verwaltung vs. Flag-Auswertung).
|
||||
package license
|
||||
|
||||
import (
|
||||
"crypto/ed25519"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrInvalidSignature = errors.New("license: signatur ungueltig")
|
||||
ErrMalformedKey = errors.New("license: lizenzschluessel hat ungueltiges format")
|
||||
)
|
||||
|
||||
// Payload ist der signierte Lizenzinhalt (Akzeptanzkriterium 3: Plan,
|
||||
// Modul-Liste, Laufzeit).
|
||||
type Payload struct {
|
||||
TenantSlug string `json:"tenant_slug"`
|
||||
Plan string `json:"plan"`
|
||||
Modules []string `json:"modules"`
|
||||
IssuedAt time.Time `json:"issued_at"`
|
||||
ValidUntil time.Time `json:"valid_until"`
|
||||
}
|
||||
|
||||
// Issuer stellt signierte Lizenzschluessel aus. Haelt den PRIVATEN
|
||||
// Ed25519-Schluessel — lebt in der Praxis beim Lizenzgeber, nicht im
|
||||
// laufenden Core-Prozess (der nur den Validator mit dem oeffentlichen
|
||||
// Schluessel braucht).
|
||||
type Issuer struct {
|
||||
priv ed25519.PrivateKey
|
||||
}
|
||||
|
||||
func NewIssuer(priv ed25519.PrivateKey) *Issuer {
|
||||
return &Issuer{priv: priv}
|
||||
}
|
||||
|
||||
// Issue liefert den Lizenzschluessel im Format base64(payload-json) "." base64(signatur).
|
||||
func (i *Issuer) Issue(payload Payload) (string, error) {
|
||||
raw, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("payload serialisieren: %w", err)
|
||||
}
|
||||
sig := ed25519.Sign(i.priv, raw)
|
||||
|
||||
return base64.RawURLEncoding.EncodeToString(raw) + "." + base64.RawURLEncoding.EncodeToString(sig), nil
|
||||
}
|
||||
|
||||
// Validator prueft Lizenzschluessel gegen den OEFFENTLICHEN Ed25519-Schluessel
|
||||
// — das ist alles, was der laufende Core-Prozess kennen muss.
|
||||
type Validator struct {
|
||||
pub ed25519.PublicKey
|
||||
}
|
||||
|
||||
func NewValidator(pub ed25519.PublicKey) *Validator {
|
||||
return &Validator{pub: pub}
|
||||
}
|
||||
|
||||
// Parse prueft die Signatur (Akzeptanzkriterium 1 / Pruefung 1) und liefert
|
||||
// bei Erfolg den entschluesselten Payload. Ein manipulierter Schluessel wird
|
||||
// hier zuverlaessig erkannt, unabhaengig davon, ob die Laufzeit noch gueltig
|
||||
// waere — Signaturpruefung und Ablaufpruefung sind bewusst getrennt
|
||||
// (Signatur bei Einspielen, Ablauf bei jeder Nutzung, siehe Store.RequireActive).
|
||||
func (v *Validator) Parse(key string) (Payload, error) {
|
||||
rawPart, sigPart, ok := splitOnce(key, '.')
|
||||
if !ok {
|
||||
return Payload{}, ErrMalformedKey
|
||||
}
|
||||
|
||||
raw, err := base64.RawURLEncoding.DecodeString(rawPart)
|
||||
if err != nil {
|
||||
return Payload{}, ErrMalformedKey
|
||||
}
|
||||
sig, err := base64.RawURLEncoding.DecodeString(sigPart)
|
||||
if err != nil {
|
||||
return Payload{}, ErrMalformedKey
|
||||
}
|
||||
|
||||
if !ed25519.Verify(v.pub, raw, sig) {
|
||||
return Payload{}, ErrInvalidSignature
|
||||
}
|
||||
|
||||
var p Payload
|
||||
if err := json.Unmarshal(raw, &p); err != nil {
|
||||
// Signatur war gueltig, aber Payload nicht mehr parsebar — sollte bei
|
||||
// unveraenderten Schluesseln nie vorkommen, trotzdem kein Panic.
|
||||
return Payload{}, fmt.Errorf("%w: payload nicht lesbar", ErrMalformedKey)
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
func splitOnce(s string, sep byte) (before, after string, ok bool) {
|
||||
for i := 0; i < len(s); i++ {
|
||||
if s[i] == sep {
|
||||
return s[:i], s[i+1:], true
|
||||
}
|
||||
}
|
||||
return "", "", false
|
||||
}
|
||||
@@ -1,106 +0,0 @@
|
||||
package license
|
||||
|
||||
import (
|
||||
"crypto/ed25519"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func testKeyPair(t *testing.T) (ed25519.PublicKey, ed25519.PrivateKey) {
|
||||
t.Helper()
|
||||
pub, priv, err := ed25519.GenerateKey(nil)
|
||||
if err != nil {
|
||||
t.Fatalf("schluesselpaar erzeugen: %v", err)
|
||||
}
|
||||
return pub, priv
|
||||
}
|
||||
|
||||
func TestIssueAndParse_RoundTrip(t *testing.T) {
|
||||
pub, priv := testKeyPair(t)
|
||||
issuer := NewIssuer(priv)
|
||||
validator := NewValidator(pub)
|
||||
|
||||
payload := Payload{
|
||||
TenantSlug: "acme",
|
||||
Plan: "pro",
|
||||
Modules: []string{"dms", "mail"},
|
||||
IssuedAt: time.Now().Truncate(time.Second),
|
||||
ValidUntil: time.Now().Add(365 * 24 * time.Hour).Truncate(time.Second),
|
||||
}
|
||||
|
||||
key, err := issuer.Issue(payload)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
got, err := validator.Parse(key)
|
||||
if err != nil {
|
||||
t.Fatalf("parse: %v", err)
|
||||
}
|
||||
if got.TenantSlug != payload.TenantSlug || got.Plan != payload.Plan || len(got.Modules) != 2 {
|
||||
t.Fatalf("payload nach parse unerwartet: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 1: manipulierter Schluessel wird zuverlaessig erkannt.
|
||||
func TestParse_RejectsTamperedKey(t *testing.T) {
|
||||
pub, priv := testKeyPair(t)
|
||||
issuer := NewIssuer(priv)
|
||||
validator := NewValidator(pub)
|
||||
|
||||
key, err := issuer.Issue(Payload{TenantSlug: "acme", Plan: "pro", ValidUntil: time.Now().Add(time.Hour)})
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
// Ein Zeichen im signierten Teil aendern.
|
||||
tampered := []byte(key)
|
||||
changed := false
|
||||
for i, c := range tampered {
|
||||
if c != '.' {
|
||||
if c == 'A' {
|
||||
tampered[i] = 'B'
|
||||
} else {
|
||||
tampered[i] = 'A'
|
||||
}
|
||||
changed = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !changed {
|
||||
t.Fatal("testaufbau fehlerhaft: nichts zum manipulieren gefunden")
|
||||
}
|
||||
|
||||
if _, err := validator.Parse(string(tampered)); !errors.Is(err, ErrInvalidSignature) {
|
||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParse_RejectsWrongKeyPair(t *testing.T) {
|
||||
_, priv := testKeyPair(t)
|
||||
otherPub, _ := testKeyPair(t)
|
||||
|
||||
issuer := NewIssuer(priv)
|
||||
validator := NewValidator(otherPub) // falscher oeffentlicher Schluessel
|
||||
|
||||
key, err := issuer.Issue(Payload{TenantSlug: "acme", Plan: "pro"})
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
if _, err := validator.Parse(key); !errors.Is(err, ErrInvalidSignature) {
|
||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParse_RejectsMalformedKey(t *testing.T) {
|
||||
pub, _ := testKeyPair(t)
|
||||
validator := NewValidator(pub)
|
||||
|
||||
cases := []string{"", "keine-punkt-trennung", "!!!.!!!"}
|
||||
for _, c := range cases {
|
||||
if _, err := validator.Parse(c); !errors.Is(err, ErrMalformedKey) {
|
||||
t.Fatalf("Parse(%q): erwartet ErrMalformedKey, habe %v", c, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,88 +0,0 @@
|
||||
package license
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrNoLicense = errors.New("license: kein lizenzdatensatz fuer diesen tenant")
|
||||
ErrLicenseExpired = errors.New("license: lizenz abgelaufen")
|
||||
)
|
||||
|
||||
// Store persistiert den Lizenzumfang je Tenant in der Control-Plane-Registry
|
||||
// (siehe internal/tenant.Registry — dieselbe Datenbank, aber ein eigener,
|
||||
// unabhaengiger Store, um internal/tenant nicht um lizenzfremde Belange zu
|
||||
// erweitern).
|
||||
type Store struct {
|
||||
pool *pgxpool.Pool
|
||||
validator *Validator
|
||||
}
|
||||
|
||||
func NewStore(pool *pgxpool.Pool, validator *Validator) *Store {
|
||||
return &Store{pool: pool, validator: validator}
|
||||
}
|
||||
|
||||
// Install prueft die Signatur des Lizenzschluessels (Akzeptanzkriterium 1)
|
||||
// und ersetzt den bisherigen Lizenzdatensatz des Tenants vollstaendig. Ein
|
||||
// bereits abgelaufener, aber korrekt signierter Schluessel wird trotzdem
|
||||
// gespeichert — der Ablauf wird erst bei der Nutzung (RequireActive)
|
||||
// bewertet, nicht beim Einspielen.
|
||||
func (s *Store) Install(ctx context.Context, tenantID, licenseKey string) (Payload, error) {
|
||||
payload, err := s.validator.Parse(licenseKey)
|
||||
if err != nil {
|
||||
return Payload{}, err
|
||||
}
|
||||
|
||||
_, err = s.pool.Exec(ctx, `
|
||||
INSERT INTO tenant_licenses (tenant_id, plan, modules, issued_at, valid_until, raw_key, installed_at)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, now())
|
||||
ON CONFLICT (tenant_id) DO UPDATE SET
|
||||
plan = $2, modules = $3, issued_at = $4, valid_until = $5, raw_key = $6, installed_at = now()
|
||||
`, tenantID, payload.Plan, payload.Modules, payload.IssuedAt, payload.ValidUntil, licenseKey)
|
||||
if err != nil {
|
||||
return Payload{}, fmt.Errorf("lizenz speichern: %w", err)
|
||||
}
|
||||
return payload, nil
|
||||
}
|
||||
|
||||
func (s *Store) get(ctx context.Context, tenantID string) (Payload, error) {
|
||||
var p Payload
|
||||
row := s.pool.QueryRow(ctx, `
|
||||
SELECT plan, modules, issued_at, valid_until
|
||||
FROM tenant_licenses WHERE tenant_id = $1
|
||||
`, tenantID)
|
||||
if err := row.Scan(&p.Plan, &p.Modules, &p.IssuedAt, &p.ValidUntil); err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return Payload{}, ErrNoLicense
|
||||
}
|
||||
return Payload{}, fmt.Errorf("lizenz lesen: %w", err)
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// Status liefert den persistierten Lizenzumfang unabhaengig vom Ablauf
|
||||
// (Akzeptanzkriterium 3: Plan, Modul-Liste, Laufzeit abfragbar).
|
||||
func (s *Store) Status(ctx context.Context, tenantID string) (Payload, error) {
|
||||
return s.get(ctx, tenantID)
|
||||
}
|
||||
|
||||
// RequireActive liefert den Lizenzumfang NUR, wenn die Lizenz noch nicht
|
||||
// abgelaufen ist — sonst ErrLicenseExpired statt eines harten Fehlers/Panics
|
||||
// (Akzeptanzkriterium 2: definierter eingeschraenkter Zustand). Aufrufende
|
||||
// Module (LIC-02/03) entscheiden, was "eingeschraenkt" konkret bedeutet.
|
||||
func (s *Store) RequireActive(ctx context.Context, tenantID string) (Payload, error) {
|
||||
p, err := s.get(ctx, tenantID)
|
||||
if err != nil {
|
||||
return Payload{}, err
|
||||
}
|
||||
if time.Now().After(p.ValidUntil) {
|
||||
return Payload{}, ErrLicenseExpired
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
@@ -1,177 +0,0 @@
|
||||
package license
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"errors"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupStoreTest(t *testing.T) (*Store, *Issuer, string, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS tenants (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
slug TEXT NOT NULL UNIQUE,
|
||||
name TEXT NOT NULL,
|
||||
db_name TEXT NOT NULL UNIQUE,
|
||||
db_dsn TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'active',
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS tenant_licenses (
|
||||
tenant_id UUID PRIMARY KEY REFERENCES tenants(id),
|
||||
plan TEXT NOT NULL,
|
||||
modules TEXT[] NOT NULL,
|
||||
issued_at TIMESTAMPTZ NOT NULL,
|
||||
valid_until TIMESTAMPTZ NOT NULL,
|
||||
raw_key TEXT NOT NULL,
|
||||
installed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
var tenantID string
|
||||
if err := pool.QueryRow(ctx, `
|
||||
INSERT INTO tenants (slug, name, db_name, db_dsn)
|
||||
VALUES ('lic_test_tenant', 'Lic Test', 'tenant_lic_test', 'unused')
|
||||
RETURNING id
|
||||
`).Scan(&tenantID); err != nil {
|
||||
t.Fatalf("test-tenant anlegen: %v", err)
|
||||
}
|
||||
|
||||
pub, priv, err := ed25519.GenerateKey(nil)
|
||||
if err != nil {
|
||||
t.Fatalf("schluesselpaar: %v", err)
|
||||
}
|
||||
issuer := NewIssuer(priv)
|
||||
store := NewStore(pool, NewValidator(pub))
|
||||
|
||||
cleanup := func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM tenant_licenses WHERE tenant_id = $1`, tenantID)
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM tenants WHERE id = $1`, tenantID)
|
||||
pool.Close()
|
||||
}
|
||||
return store, issuer, tenantID, cleanup
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3: Lizenzumfang persistiert und abfragbar.
|
||||
func TestStore_InstallAndStatus(t *testing.T) {
|
||||
store, issuer, tenantID, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
payload := Payload{
|
||||
TenantSlug: "lic_test_tenant",
|
||||
Plan: "enterprise",
|
||||
Modules: []string{"dms", "mail", "archive"},
|
||||
IssuedAt: time.Now().Truncate(time.Second),
|
||||
ValidUntil: time.Now().Add(30 * 24 * time.Hour).Truncate(time.Second),
|
||||
}
|
||||
key, err := issuer.Issue(payload)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
if _, err := store.Install(ctx, tenantID, key); err != nil {
|
||||
t.Fatalf("install: %v", err)
|
||||
}
|
||||
|
||||
status, err := store.Status(ctx, tenantID)
|
||||
if err != nil {
|
||||
t.Fatalf("status: %v", err)
|
||||
}
|
||||
if status.Plan != "enterprise" || len(status.Modules) != 3 {
|
||||
t.Fatalf("status unerwartet: %+v", status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_InstallRejectsInvalidSignature(t *testing.T) {
|
||||
store, _, tenantID, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
_, otherPriv, _ := ed25519.GenerateKey(nil)
|
||||
foreignIssuer := NewIssuer(otherPriv) // signiert mit falschem schluessel
|
||||
|
||||
key, err := foreignIssuer.Issue(Payload{TenantSlug: "lic_test_tenant", Plan: "pro", ValidUntil: time.Now().Add(time.Hour)})
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
if _, err := store.Install(ctx, tenantID, key); !errors.Is(err, ErrInvalidSignature) {
|
||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: abgelaufene Lizenz fuehrt zu definiertem
|
||||
// eingeschraenktem Zustand (ErrLicenseExpired), nicht zu einem Absturz.
|
||||
func TestStore_RequireActive_DetectsExpiry(t *testing.T) {
|
||||
store, issuer, tenantID, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
expired := Payload{
|
||||
TenantSlug: "lic_test_tenant",
|
||||
Plan: "pro",
|
||||
Modules: []string{"dms"},
|
||||
IssuedAt: time.Now().Add(-48 * time.Hour),
|
||||
ValidUntil: time.Now().Add(-24 * time.Hour), // bereits abgelaufen
|
||||
}
|
||||
key, err := issuer.Issue(expired)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
// Einspielen einer bereits abgelaufenen, aber korrekt signierten Lizenz
|
||||
// muss funktionieren (Ablauf wird erst bei Nutzung bewertet).
|
||||
if _, err := store.Install(ctx, tenantID, key); err != nil {
|
||||
t.Fatalf("install sollte trotz ablauf funktionieren: %v", err)
|
||||
}
|
||||
|
||||
func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
t.Fatalf("RequireActive hat gepanict statt einen fehler zu liefern: %v", r)
|
||||
}
|
||||
}()
|
||||
if _, err := store.RequireActive(ctx, tenantID); !errors.Is(err, ErrLicenseExpired) {
|
||||
t.Fatalf("erwartet ErrLicenseExpired, habe %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
// Aber der Umfang bleibt weiterhin abfragbar (Status, im Unterschied zu RequireActive).
|
||||
status, err := store.Status(ctx, tenantID)
|
||||
if err != nil {
|
||||
t.Fatalf("status sollte trotz ablauf funktionieren: %v", err)
|
||||
}
|
||||
if status.Plan != "pro" {
|
||||
t.Fatalf("status unerwartet: %+v", status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_RequireActive_NoLicense(t *testing.T) {
|
||||
store, _, tenantID, cleanup := setupStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := store.RequireActive(ctx, tenantID); !errors.Is(err, ErrNoLicense) {
|
||||
t.Fatalf("erwartet ErrNoLicense, habe %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package metrics
|
||||
|
||||
import "github.com/prometheus/client_golang/prometheus"
|
||||
|
||||
// NewCoreRegistry liefert das Prometheus-Registry fuer die EIGENEN
|
||||
// Kennzahlen des Core-Dienstes (Akzeptanzkriterium 1) — alle Namen tragen
|
||||
// das Praefix "nexarch_core_" gemaess der im Paketkommentar dokumentierten
|
||||
// Namenskonvention (Akzeptanzkriterium 3). Ein eigenes Registry statt des
|
||||
// globalen DefaultRegisterer, damit Tests unabhaengig voneinander sind.
|
||||
func NewCoreRegistry() *prometheus.Registry {
|
||||
reg := prometheus.NewRegistry()
|
||||
reg.MustRegister(
|
||||
prometheus.NewGaugeFunc(prometheus.GaugeOpts{
|
||||
Name: "nexarch_core_up",
|
||||
Help: "1, solange der Core-Dienst laeuft und Metriken liefern kann.",
|
||||
}, func() float64 { return 1 }),
|
||||
)
|
||||
return reg
|
||||
}
|
||||
@@ -0,0 +1,168 @@
|
||||
// Package metrics implementiert Core OPS-03: einen zentralen /metrics-
|
||||
// Endpunkt im Prometheus-Textformat, der Kennzahlen des Core-Dienstes UND
|
||||
// aggregierte Kennzahlen aller registrierten Module bereitstellt — offenes
|
||||
// Pull-Modell nach Prometheus-Vorbild, kein proprietaerer Push-Mechanismus.
|
||||
//
|
||||
// Namenskonvention (Akzeptanzkriterium 3, modulübergreifend konsistent):
|
||||
//
|
||||
// nexarch_core_<name> — Kennzahlen des Core-Dienstes selbst
|
||||
// nexarch_module_<modul>_<name> — von einem Modul gescrapte Kennzahl
|
||||
// <name>, umbenannt mit dem
|
||||
// Modulnamen als Praefix
|
||||
//
|
||||
// Ein Modul liefert seine eigenen Kennzahlen unter EIGENEM Namen (z.B.
|
||||
// "requests_total") unter seinem eigenen /metrics-Endpunkt — dieses Paket
|
||||
// benennt sie beim Einsammeln konsistent um, damit im aggregierten Core-
|
||||
// Endpunkt niemals zwei Module denselben Metrik-Namen kollidieren lassen.
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
"github.com/prometheus/common/expfmt"
|
||||
"github.com/prometheus/common/model"
|
||||
)
|
||||
|
||||
// init erzwingt das klassische Prometheus-Namensschema (a-z, A-Z, 0-9, _)
|
||||
// fuer die Namensvalidierung von expfmt/model — ohne diese explizite
|
||||
// Festlegung liefert die Bibliothek "Invalid name validation scheme
|
||||
// requested: unset" beim Parsen/Kodieren, da sie den globalen Default in
|
||||
// dieser Version nicht mehr implizit setzt.
|
||||
func init() {
|
||||
model.NameValidationScheme = model.LegacyValidation
|
||||
}
|
||||
|
||||
// Source ist EIN registriertes Modul mit seinem eigenen /metrics-Endpunkt
|
||||
// (siehe internal/health fuer das analoge Muster bei Readiness-Checks).
|
||||
type Source struct {
|
||||
ModuleName string
|
||||
MetricsURL string
|
||||
}
|
||||
|
||||
// SourceProvider liefert die aktuell registrierten Module — typischerweise
|
||||
// rueckgebunden an internal/moduleregistry.Registry.List (API-02) ueber
|
||||
// einen kleinen Adapter im aufrufenden Code, damit dieses Paket
|
||||
// internal/moduleregistry nicht direkt importieren muss (Kein Umbau
|
||||
// angrenzender Bereiche). Ein NEU registriertes Modul erscheint automatisch
|
||||
// beim naechsten Aufruf von Aggregator.Handler, OHNE Codeaenderung an diesem
|
||||
// Paket (Akzeptanzkriterium 2 / Pruefung 2).
|
||||
type SourceProvider func(ctx context.Context) ([]Source, error)
|
||||
|
||||
// FetchTimeout begrenzt, wie lange EIN Modul-Scrape maximal dauern darf —
|
||||
// ein haengendes Modul darf den gesamten Aggregations-Request nicht
|
||||
// verzoegern (Pruefung 1: Antwort unter Last innerhalb definierter Zeit).
|
||||
const FetchTimeout = 2 * time.Second
|
||||
|
||||
// Aggregator sammelt Core-eigene Metriken (coreGatherer) und die Metriken
|
||||
// aller ueber sourceProvider gemeldeten Module in EINER Antwort ein.
|
||||
type Aggregator struct {
|
||||
coreGatherer prometheus.Gatherer
|
||||
sourceProvider SourceProvider
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
func NewAggregator(coreGatherer prometheus.Gatherer, sourceProvider SourceProvider) *Aggregator {
|
||||
return &Aggregator{
|
||||
coreGatherer: coreGatherer,
|
||||
sourceProvider: sourceProvider,
|
||||
client: &http.Client{Timeout: FetchTimeout},
|
||||
}
|
||||
}
|
||||
|
||||
// Gather implementiert prometheus.Gatherer: liefert Core-Metriken PLUS alle
|
||||
// erreichbaren Modul-Metriken (umbenannt gemaess Namenskonvention) in einer
|
||||
// gemeinsamen Liste von MetricFamilies.
|
||||
func (a *Aggregator) Gather(ctx context.Context) ([]*dto.MetricFamily, error) {
|
||||
families, err := a.coreGatherer.Gather()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("core-metriken einsammeln: %w", err)
|
||||
}
|
||||
|
||||
sources, err := a.sourceProvider(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("modul-quellen ermitteln: %w", err)
|
||||
}
|
||||
|
||||
// Module werden NEBENLAEUFIG gescrapt (dasselbe Muster wie
|
||||
// internal/health.Registry.CheckAll) — ein langsames/nicht erreichbares
|
||||
// Modul haelt weder andere Module noch den Gesamt-Request auf.
|
||||
type fetchResult struct {
|
||||
families []*dto.MetricFamily
|
||||
}
|
||||
resultCh := make(chan fetchResult, len(sources))
|
||||
for _, src := range sources {
|
||||
go func(src Source) {
|
||||
fetchCtx, cancel := context.WithTimeout(ctx, FetchTimeout)
|
||||
defer cancel()
|
||||
mf, err := a.fetchAndRename(fetchCtx, src)
|
||||
if err != nil {
|
||||
resultCh <- fetchResult{} // Fehlerfall: einfach nichts beitragen, Aggregation laeuft weiter
|
||||
return
|
||||
}
|
||||
resultCh <- fetchResult{families: mf}
|
||||
}(src)
|
||||
}
|
||||
for range sources {
|
||||
r := <-resultCh
|
||||
families = append(families, r.families...)
|
||||
}
|
||||
|
||||
return families, nil
|
||||
}
|
||||
|
||||
// fetchAndRename ruft die /metrics-URL eines Moduls ab, parst das
|
||||
// Prometheus-Textformat und benennt jede Metrik gemaess der
|
||||
// Namenskonvention um (Akzeptanzkriterium 2 + 3).
|
||||
func (a *Aggregator) fetchAndRename(ctx context.Context, src Source) ([]*dto.MetricFamily, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, src.MetricsURL, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp, err := a.client.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, fmt.Errorf("modul %s: unerwarteter status %d", src.ModuleName, resp.StatusCode)
|
||||
}
|
||||
|
||||
parser := expfmt.NewTextParser(model.LegacyValidation)
|
||||
parsed, err := parser.TextToMetricFamilies(resp.Body)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("modul %s: metrik-text nicht parsebar: %w", src.ModuleName, err)
|
||||
}
|
||||
|
||||
out := make([]*dto.MetricFamily, 0, len(parsed))
|
||||
for name, mf := range parsed {
|
||||
renamed := fmt.Sprintf("nexarch_module_%s_%s", src.ModuleName, name)
|
||||
mf.Name = &renamed
|
||||
out = append(out, mf)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Handler liefert einen HTTP-Handler, der Gather aufruft und das Ergebnis im
|
||||
// Prometheus-Textformat ausgibt (Akzeptanzkriterium 1).
|
||||
func (a *Aggregator) Handler() http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
families, err := a.Gather(r.Context())
|
||||
if err != nil {
|
||||
http.Error(w, "metriken konnten nicht eingesammelt werden", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", string(expfmt.NewFormat(expfmt.TypeTextPlain)))
|
||||
enc := expfmt.NewEncoder(w, expfmt.NewFormat(expfmt.TypeTextPlain))
|
||||
for _, mf := range families {
|
||||
if err := enc.Encode(mf); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/common/expfmt"
|
||||
"github.com/prometheus/common/model"
|
||||
)
|
||||
|
||||
func fakeModuleServer(metricName string) *httptest.Server {
|
||||
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/plain; version=0.0.4")
|
||||
fmt.Fprintf(w, "# HELP %s ein test-zaehler\n# TYPE %s counter\n%s 42\n", metricName, metricName, metricName)
|
||||
}))
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1: Core liefert unter dem Handler valides
|
||||
// Prometheus-Textformat mit den eigenen Kennzahlen.
|
||||
func TestHandler_ServesCoreMetricsInPrometheusFormat(t *testing.T) {
|
||||
agg := NewAggregator(NewCoreRegistry(), func(ctx context.Context) ([]Source, error) { return nil, nil })
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/metrics", nil)
|
||||
rec := httptest.NewRecorder()
|
||||
agg.Handler()(rec, req)
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rec.Code)
|
||||
}
|
||||
|
||||
parser := expfmt.NewTextParser(model.LegacyValidation)
|
||||
families, err := parser.TextToMetricFamilies(strings.NewReader(rec.Body.String()))
|
||||
if err != nil {
|
||||
t.Fatalf("antwort ist kein valides prometheus-textformat: %v", err)
|
||||
}
|
||||
if _, ok := families["nexarch_core_up"]; !ok {
|
||||
t.Fatalf("erwartet 'nexarch_core_up' unter den core-metriken, habe: %v", keysOf(families))
|
||||
}
|
||||
}
|
||||
|
||||
func keysOf[V any](m map[string]V) []string {
|
||||
out := make([]string, 0, len(m))
|
||||
for k := range m {
|
||||
out = append(out, k)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: ein NEU registriertes Modul erscheint
|
||||
// in der Aggregation, OHNE dass dieses Paket oder der Aufrufer Code
|
||||
// aendern muss — die Quelle kommt ausschliesslich aus sourceProvider.
|
||||
func TestHandler_NewlyRegisteredModuleAppearsWithoutCodeChange(t *testing.T) {
|
||||
moduleServer := fakeModuleServer("requests_total")
|
||||
defer moduleServer.Close()
|
||||
|
||||
// Simuliert eine sich zur Laufzeit aendernde Modul-Liste (z.B. aus
|
||||
// SourceStore.Provide) — zunaechst LEER, dann mit einem Eintrag.
|
||||
var sources []Source
|
||||
var mu sync.Mutex
|
||||
provider := func(ctx context.Context) ([]Source, error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
out := make([]Source, len(sources))
|
||||
copy(out, sources)
|
||||
return out, nil
|
||||
}
|
||||
|
||||
agg := NewAggregator(NewCoreRegistry(), provider)
|
||||
|
||||
// Vor der Registrierung: Modul-Metrik nicht vorhanden.
|
||||
rec1 := httptest.NewRecorder()
|
||||
agg.Handler()(rec1, httptest.NewRequest(http.MethodGet, "/metrics", nil))
|
||||
if strings.Contains(rec1.Body.String(), "requests_total") {
|
||||
t.Fatal("modul-metrik haette vor registrierung nicht erscheinen duerfen")
|
||||
}
|
||||
|
||||
// Modul wird "registriert" (kein Code hier oder in metrics.go aendert sich).
|
||||
mu.Lock()
|
||||
sources = append(sources, Source{ModuleName: "dms", MetricsURL: moduleServer.URL})
|
||||
mu.Unlock()
|
||||
|
||||
rec2 := httptest.NewRecorder()
|
||||
agg.Handler()(rec2, httptest.NewRequest(http.MethodGet, "/metrics", nil))
|
||||
body := rec2.Body.String()
|
||||
if !strings.Contains(body, "nexarch_module_dms_requests_total") {
|
||||
t.Fatalf("erwartet umbenannte modul-metrik 'nexarch_module_dms_requests_total' nach registrierung, body:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3: Namenskonvention "nexarch_module_<modul>_<name>"
|
||||
// wird tatsaechlich angewendet.
|
||||
func TestFetchAndRename_AppliesNamingConvention(t *testing.T) {
|
||||
moduleServer := fakeModuleServer("queue_depth")
|
||||
defer moduleServer.Close()
|
||||
|
||||
agg := NewAggregator(NewCoreRegistry(), nil)
|
||||
families, err := agg.fetchAndRename(context.Background(), Source{ModuleName: "mail", MetricsURL: moduleServer.URL})
|
||||
if err != nil {
|
||||
t.Fatalf("fetchAndRename: %v", err)
|
||||
}
|
||||
if len(families) != 1 || families[0].GetName() != "nexarch_module_mail_queue_depth" {
|
||||
t.Fatalf("erwartet genau 1 metrik 'nexarch_module_mail_queue_depth', habe: %+v", families)
|
||||
}
|
||||
}
|
||||
|
||||
// Ein nicht erreichbares Modul darf die Aggregation der uebrigen und die
|
||||
// Gesamtantwort nicht verhindern (dieselbe Resilienz wie OPS-02).
|
||||
func TestHandler_UnreachableModuleDoesNotBreakAggregation(t *testing.T) {
|
||||
reachable := fakeModuleServer("healthy_metric")
|
||||
defer reachable.Close()
|
||||
unreachable := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
|
||||
unreachableURL := unreachable.URL
|
||||
unreachable.Close() // sofort schliessen -> Verbindung schlaegt fehl
|
||||
|
||||
provider := func(ctx context.Context) ([]Source, error) {
|
||||
return []Source{
|
||||
{ModuleName: "ok", MetricsURL: reachable.URL},
|
||||
{ModuleName: "kaputt", MetricsURL: unreachableURL},
|
||||
}, nil
|
||||
}
|
||||
agg := NewAggregator(NewCoreRegistry(), provider)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
agg.Handler()(rec, httptest.NewRequest(http.MethodGet, "/metrics", nil))
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200 trotz einem nicht erreichbaren modul", rec.Code)
|
||||
}
|
||||
body := rec.Body.String()
|
||||
if !strings.Contains(body, "nexarch_module_ok_healthy_metric") {
|
||||
t.Fatal("erreichbares modul haette trotz ausfall des anderen aggregiert werden sollen")
|
||||
}
|
||||
if strings.Contains(body, "kaputt") {
|
||||
t.Fatal("nicht erreichbares modul haette keine metrik beitragen duerfen")
|
||||
}
|
||||
}
|
||||
|
||||
// Pruefung 1: Endpunkt antwortet unter mehreren gleichzeitigen Anfragen
|
||||
// innerhalb definierter Zeit — kein unbeschraenktes Blockieren durch
|
||||
// langsame Module (FetchTimeout begrenzt jeden Scrape).
|
||||
func TestHandler_RespondsWithinBoundedTimeUnderLoad(t *testing.T) {
|
||||
hangingServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
time.Sleep(10 * time.Second) // wuerde ohne timeout jede anfrage blockieren
|
||||
}))
|
||||
defer hangingServer.Close()
|
||||
|
||||
provider := func(ctx context.Context) ([]Source, error) {
|
||||
return []Source{{ModuleName: "haengend", MetricsURL: hangingServer.URL}}, nil
|
||||
}
|
||||
agg := NewAggregator(NewCoreRegistry(), provider)
|
||||
|
||||
const concurrentRequests = 10
|
||||
var wg sync.WaitGroup
|
||||
start := time.Now()
|
||||
for i := 0; i < concurrentRequests; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
rec := httptest.NewRecorder()
|
||||
agg.Handler()(rec, httptest.NewRequest(http.MethodGet, "/metrics", nil))
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Errorf("status = %d, want 200", rec.Code)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
elapsed := time.Since(start)
|
||||
|
||||
if elapsed > FetchTimeout+3*time.Second {
|
||||
t.Fatalf("%d gleichzeitige anfragen brauchten %s, erwartet deutlich unter %s durch FetchTimeout",
|
||||
concurrentRequests, elapsed, FetchTimeout+3*time.Second)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// SourceStore persistiert, welche Module ihre Metriken unter welcher URL
|
||||
// bereitstellen — dieselbe Postgres-basierte "kein Code-Deploy noetig"-
|
||||
// Konvention wie internal/statuspage.Store.RegisterTarget (OPS-02): ein neu
|
||||
// registriertes Modul erscheint automatisch in der Aggregation, sobald es
|
||||
// hier eingetragen ist (Akzeptanzkriterium 2 / Pruefung 2).
|
||||
type SourceStore struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewSourceStore(pool *pgxpool.Pool) *SourceStore {
|
||||
return &SourceStore{pool: pool}
|
||||
}
|
||||
|
||||
func (s *SourceStore) RegisterSource(ctx context.Context, moduleName, metricsURL string) error {
|
||||
_, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO metrics_sources (module_name, metrics_url)
|
||||
VALUES ($1, $2)
|
||||
ON CONFLICT (module_name) DO UPDATE SET metrics_url = $2
|
||||
`, moduleName, metricsURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("metrik-quelle speichern: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Provide implementiert SourceProvider direkt aus der Datenbank.
|
||||
func (s *SourceStore) Provide(ctx context.Context) ([]Source, error) {
|
||||
rows, err := s.pool.Query(ctx, `SELECT module_name, metrics_url FROM metrics_sources ORDER BY module_name`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("metrik-quellen auflisten: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []Source
|
||||
for rows.Next() {
|
||||
var src Source
|
||||
if err := rows.Scan(&src.ModuleName, &src.MetricsURL); err != nil {
|
||||
return nil, fmt.Errorf("metrik-quelle lesen: %w", err)
|
||||
}
|
||||
out = append(out, src)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupSourceStoreTest(t *testing.T) (*SourceStore, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS metrics_sources (module_name TEXT PRIMARY KEY, metrics_url TEXT NOT NULL)
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
cleanup := func() { pool.Close() }
|
||||
return NewSourceStore(pool), cleanup
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 / Pruefung 2 auf Persistenz-Ebene: eine ueber die
|
||||
// Datenbank registrierte Quelle ist sofort ueber Provide() sichtbar — genau
|
||||
// der Mechanismus, der ein neues Modul ohne Core-Codeaenderung erscheinen
|
||||
// laesst.
|
||||
func TestSourceStore_RegisterSourceAppearsInProvide(t *testing.T) {
|
||||
store, cleanup := setupSourceStoreTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
name := fmt.Sprintf("modul-%d", time.Now().UnixNano())
|
||||
|
||||
if err := store.RegisterSource(ctx, name, "http://example.invalid/metrics"); err != nil {
|
||||
t.Fatalf("registersource: %v", err)
|
||||
}
|
||||
|
||||
sources, err := store.Provide(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("provide: %v", err)
|
||||
}
|
||||
found := false
|
||||
for _, s := range sources {
|
||||
if s.ModuleName == name {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Fatalf("erwartet %s in provide()-ergebnis, habe: %+v", name, sources)
|
||||
}
|
||||
}
|
||||
@@ -1,97 +0,0 @@
|
||||
package moduletrust
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// StaleCache ist der generische Rechte-/Feature-Flag-Cache-Kontrakt
|
||||
// (Akzeptanzkriterium 2): TTL-basiert, mit explizitem, benanntem Verhalten
|
||||
// bei abgelaufenem Cache waehrend Core nicht erreichbar ist.
|
||||
//
|
||||
// - Get: FAIL-OPEN fuer Lesevorgaenge. Schlaegt der Refresh fehl, aber es
|
||||
// gibt bereits einen (wenn auch abgelaufenen) Stand, wird dieser mit
|
||||
// stale=true zurueckgegeben — Begruendung: ein bereits authentifiziertes
|
||||
// Modul soll mit dem letztbekannten Stand weiterarbeiten koennen statt
|
||||
// hart zu blockieren (siehe "Bekannte Fehler vermeiden" im Ticket).
|
||||
// Existiert noch nie ein Stand, gibt es keinen sinnvollen Fallback —
|
||||
// dann liefert auch Get einen Fehler.
|
||||
// - RequireFresh: FAIL-CLOSED fuer sicherheitskritische Aktionen (z.B.
|
||||
// ein komplett NEUER Login). Nutzt NIEMALS einen zwischengespeicherten
|
||||
// Stand, ruft immer frisch ab — Begruendung: eine neue Vertrauens-
|
||||
// entscheidung darf nicht auf veralteten Daten beruhen, auch wenn das
|
||||
// bedeutet, dass die Aktion bei Core-Ausfall sichtbar fehlschlaegt statt
|
||||
// unsicher "irgendwie" durchgelassen zu werden.
|
||||
//
|
||||
// LIC-02 (internal/flag.Service) implementiert bereits denselben Kontrakt
|
||||
// fuer Feature-Flags — StaleCache verallgemeinert dasselbe Muster fuer
|
||||
// JWT-Signaturschluessel, damit beide Faelle derselben dokumentierten
|
||||
// Policy folgen.
|
||||
type StaleCache[T any] struct {
|
||||
mu sync.RWMutex
|
||||
value T
|
||||
hasValue bool
|
||||
fetchedAt time.Time
|
||||
ttl time.Duration
|
||||
fetch func(ctx context.Context) (T, error)
|
||||
}
|
||||
|
||||
func NewStaleCache[T any](ttl time.Duration, fetch func(ctx context.Context) (T, error)) *StaleCache[T] {
|
||||
return &StaleCache[T]{ttl: ttl, fetch: fetch}
|
||||
}
|
||||
|
||||
// Get liefert den Cache-Wert. FAIL-OPEN: bei Refresh-Fehler wird ein
|
||||
// vorhandener, ggf. abgelaufener Stand zurueckgegeben (stale=true).
|
||||
func (c *StaleCache[T]) Get(ctx context.Context) (value T, stale bool, err error) {
|
||||
c.mu.RLock()
|
||||
fresh := c.hasValue && time.Since(c.fetchedAt) < c.ttl
|
||||
if fresh {
|
||||
v := c.value
|
||||
c.mu.RUnlock()
|
||||
return v, false, nil
|
||||
}
|
||||
c.mu.RUnlock()
|
||||
|
||||
newVal, fetchErr := c.fetch(ctx)
|
||||
if fetchErr == nil {
|
||||
c.mu.Lock()
|
||||
c.value, c.hasValue, c.fetchedAt = newVal, true, time.Now()
|
||||
c.mu.Unlock()
|
||||
return newVal, false, nil
|
||||
}
|
||||
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
if c.hasValue {
|
||||
return c.value, true, nil
|
||||
}
|
||||
var zero T
|
||||
return zero, false, fmt.Errorf("cache leer und refresh fehlgeschlagen: %w", fetchErr)
|
||||
}
|
||||
|
||||
// Invalidate erzwingt beim naechsten Get-Aufruf einen sofortigen Refresh
|
||||
// statt auf den TTL-Ablauf zu warten (API-06, Akzeptanzkriterium 1: ein
|
||||
// Health-Check-getriggerter Wiederanlauf soll den Cache SOFORT aktualisieren,
|
||||
// nicht die reguläre TTL abwarten) — rein additiv, aendert nichts an
|
||||
// Get/RequireFresh (dasselbe Muster wie internal/flag.Service.Invalidate).
|
||||
func (c *StaleCache[T]) Invalidate() {
|
||||
c.mu.Lock()
|
||||
c.hasValue = false
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// RequireFresh ruft IMMER frisch ab (FAIL-CLOSED) — fuer sicherheitskritische
|
||||
// Aktionen, die niemals auf einem zwischengespeicherten Stand basieren duerfen.
|
||||
func (c *StaleCache[T]) RequireFresh(ctx context.Context) (T, error) {
|
||||
v, err := c.fetch(ctx)
|
||||
if err != nil {
|
||||
var zero T
|
||||
return zero, fmt.Errorf("core nicht erreichbar, sicherheitskritische aktion abgelehnt: %w", err)
|
||||
}
|
||||
c.mu.Lock()
|
||||
c.value, c.hasValue, c.fetchedAt = v, true, time.Now()
|
||||
c.mu.Unlock()
|
||||
return v, nil
|
||||
}
|
||||
@@ -1,79 +0,0 @@
|
||||
package moduletrust
|
||||
|
||||
import (
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/golang-jwt/jwt/v5"
|
||||
)
|
||||
|
||||
type Claims struct {
|
||||
Subject string `json:"sub"`
|
||||
TenantSlug string `json:"tenant"`
|
||||
jwt.RegisteredClaims
|
||||
}
|
||||
|
||||
// Issue signiert ein Token mit dem aktuellen Signierschluessel und traegt
|
||||
// dessen KID im JWT-Header ein — der Verifier auf Modulseite waehlt darueber
|
||||
// den passenden oeffentlichen Schluessel aus PublicKeySet() aus.
|
||||
func (m *KeyManager) Issue(subject, tenantSlug string, ttl time.Duration) (string, error) {
|
||||
key, err := m.SigningKey()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
claims := Claims{
|
||||
Subject: subject,
|
||||
TenantSlug: tenantSlug,
|
||||
RegisteredClaims: jwt.RegisteredClaims{
|
||||
IssuedAt: jwt.NewNumericDate(now),
|
||||
ExpiresAt: jwt.NewNumericDate(now.Add(ttl)),
|
||||
},
|
||||
}
|
||||
token := jwt.NewWithClaims(jwt.SigningMethodEdDSA, claims)
|
||||
token.Header["kid"] = key.KID
|
||||
return token.SignedString(key.Private)
|
||||
}
|
||||
|
||||
type jwksResponse struct {
|
||||
Keys []jwksKey `json:"keys"`
|
||||
}
|
||||
|
||||
type jwksKey struct {
|
||||
Kid string `json:"kid"`
|
||||
PublicKey string `json:"public_key"` // base64 (raw Ed25519, 32 Byte)
|
||||
}
|
||||
|
||||
// ServeJWKS liefert alle bekannten oeffentlichen Schluessel als JSON —
|
||||
// Module fragen dies periodisch ab (nicht pro Request), siehe Verifier.
|
||||
func (m *KeyManager) ServeJWKS(w http.ResponseWriter, r *http.Request) {
|
||||
set := m.PublicKeySet()
|
||||
resp := jwksResponse{Keys: make([]jwksKey, 0, len(set))}
|
||||
for kid, pub := range set {
|
||||
resp.Keys = append(resp.Keys, jwksKey{Kid: kid, PublicKey: base64.StdEncoding.EncodeToString(pub)})
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(resp)
|
||||
}
|
||||
|
||||
// ParseJWKS dekodiert die JSON-Antwort von ServeJWKS zurueck in kid->PublicKey
|
||||
// — Hilfsfunktion fuer Module, die JWKS per HTTP abrufen.
|
||||
func ParseJWKS(data []byte) (map[string][]byte, error) {
|
||||
var resp jwksResponse
|
||||
if err := json.Unmarshal(data, &resp); err != nil {
|
||||
return nil, fmt.Errorf("jwks parsen: %w", err)
|
||||
}
|
||||
out := make(map[string][]byte, len(resp.Keys))
|
||||
for _, k := range resp.Keys {
|
||||
raw, err := base64.StdEncoding.DecodeString(k.PublicKey)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("oeffentlichen schluessel %q dekodieren: %w", k.Kid, err)
|
||||
}
|
||||
out[k.Kid] = raw
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -1,83 +0,0 @@
|
||||
// Package moduletrust implementiert Core API-05: asymmetrische JWT-Signatur
|
||||
// mit Schluesselverteilung (JWKS), damit DMS/Mail/Archive/Workflow JWTs
|
||||
// LOKAL verifizieren koennen, ohne pro Aufruf einen synchronen Request an
|
||||
// Core zu stellen — Core darf Fundament sein, ohne zum Flaschenhals zu
|
||||
// werden (siehe Entscheidungsverlauf "Vertrauensstellung Core<->Module" in
|
||||
// nexarch-state.json). IAM-02s HS256-Session-Cookie (Browser-Login) bleibt
|
||||
// unangetastet — dies ist ein zusaetzlicher, getrennter Vertrauensmechanismus
|
||||
// fuer Modul-zu-Modul/Modul-zu-Core-Aufrufe.
|
||||
package moduletrust
|
||||
|
||||
import (
|
||||
"crypto/ed25519"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"sync"
|
||||
)
|
||||
|
||||
type KeyPair struct {
|
||||
KID string
|
||||
Private ed25519.PrivateKey
|
||||
Public ed25519.PublicKey
|
||||
}
|
||||
|
||||
// KeyManager haelt ALLE noch gueltigen Schluesselpaare — nicht nur das
|
||||
// aktuell signierende. Rotate erzeugt ein neues Paar und behaelt die alten
|
||||
// fuer die Verifikation bereits ausgestellter Tokens (Akzeptanzkriterium 3:
|
||||
// Rotation ohne Ausfallzeit fuer andere Module).
|
||||
type KeyManager struct {
|
||||
mu sync.RWMutex
|
||||
keys []KeyPair // aeltestes zuerst, neuestes zuletzt
|
||||
}
|
||||
|
||||
func NewKeyManager() (*KeyManager, error) {
|
||||
m := &KeyManager{}
|
||||
if _, err := m.Rotate(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return m, nil
|
||||
}
|
||||
|
||||
// Rotate erzeugt ein neues Ed25519-Schluesselpaar mit eigener KID und macht
|
||||
// es zum aktuellen Signierschluessel. Aeltere Schluessel bleiben in
|
||||
// PublicKeySet() erhalten, damit bereits ausgestellte Tokens weiterhin
|
||||
// verifizierbar sind.
|
||||
func (m *KeyManager) Rotate() (KeyPair, error) {
|
||||
pub, priv, err := ed25519.GenerateKey(nil)
|
||||
if err != nil {
|
||||
return KeyPair{}, fmt.Errorf("schluesselpaar erzeugen: %w", err)
|
||||
}
|
||||
kidBytes := make([]byte, 8)
|
||||
if _, err := rand.Read(kidBytes); err != nil {
|
||||
return KeyPair{}, fmt.Errorf("kid erzeugen: %w", err)
|
||||
}
|
||||
kp := KeyPair{KID: hex.EncodeToString(kidBytes), Private: priv, Public: pub}
|
||||
|
||||
m.mu.Lock()
|
||||
m.keys = append(m.keys, kp)
|
||||
m.mu.Unlock()
|
||||
return kp, nil
|
||||
}
|
||||
|
||||
// SigningKey liefert den aktuellen (neuesten) Schluessel zum Signieren neuer Tokens.
|
||||
func (m *KeyManager) SigningKey() (KeyPair, error) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
if len(m.keys) == 0 {
|
||||
return KeyPair{}, fmt.Errorf("moduletrust: kein schluessel vorhanden")
|
||||
}
|
||||
return m.keys[len(m.keys)-1], nil
|
||||
}
|
||||
|
||||
// PublicKeySet liefert ALLE bekannten oeffentlichen Schluessel (kid ->
|
||||
// public key) — die Grundlage fuer den JWKS-Endpunkt.
|
||||
func (m *KeyManager) PublicKeySet() map[string]ed25519.PublicKey {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
out := make(map[string]ed25519.PublicKey, len(m.keys))
|
||||
for _, k := range m.keys {
|
||||
out[k.KID] = k.Public
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -1,214 +0,0 @@
|
||||
package moduletrust
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"errors"
|
||||
"net/http"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestIssueAndVerify_RoundTrip(t *testing.T) {
|
||||
km, err := NewKeyManager()
|
||||
if err != nil {
|
||||
t.Fatalf("new key manager: %v", err)
|
||||
}
|
||||
token, err := km.Issue("user-1", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
return km.PublicKeySet(), nil
|
||||
})
|
||||
claims, err := v.Verify(context.Background(), token)
|
||||
if err != nil {
|
||||
t.Fatalf("verify: %v", err)
|
||||
}
|
||||
if claims.Subject != "user-1" || claims.TenantSlug != "acme" {
|
||||
t.Fatalf("claims unerwartet: %+v", claims)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1: Verifikation lokal, kein Request pro Aufruf.
|
||||
func TestVerify_DoesNotFetchPerCall(t *testing.T) {
|
||||
km, err := NewKeyManager()
|
||||
if err != nil {
|
||||
t.Fatalf("new key manager: %v", err)
|
||||
}
|
||||
token, err := km.Issue("user-1", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
var mu sync.Mutex
|
||||
fetchCalls := 0
|
||||
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
mu.Lock()
|
||||
fetchCalls++
|
||||
mu.Unlock()
|
||||
return km.PublicKeySet(), nil
|
||||
})
|
||||
|
||||
for i := 0; i < 10; i++ {
|
||||
if _, err := v.Verify(context.Background(), token); err != nil {
|
||||
t.Fatalf("verify %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if fetchCalls != 1 {
|
||||
t.Fatalf("erwartet genau 1 fetch fuer 10 Verify-Aufrufe innerhalb der TTL, habe %d", fetchCalls)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 1: Core simuliert nicht erreichbar,
|
||||
// bereits authentifizierte Nutzer bleiben funktionsfaehig (Fail-Open mit
|
||||
// letztbekanntem Schluesselstand).
|
||||
func TestVerify_FailsOpenWhenCoreUnreachableButStaleKeysExist(t *testing.T) {
|
||||
km, err := NewKeyManager()
|
||||
if err != nil {
|
||||
t.Fatalf("new key manager: %v", err)
|
||||
}
|
||||
token, err := km.Issue("user-1", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
coreDown := false
|
||||
v := NewVerifier(30*time.Millisecond, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
if coreDown {
|
||||
return nil, errors.New("core nicht erreichbar (simuliert)")
|
||||
}
|
||||
return km.PublicKeySet(), nil
|
||||
})
|
||||
|
||||
// Cache vorwaermen, waehrend Core noch erreichbar ist.
|
||||
if _, err := v.Verify(context.Background(), token); err != nil {
|
||||
t.Fatalf("verify (warm): %v", err)
|
||||
}
|
||||
|
||||
// "Core abschalten" und TTL ablaufen lassen.
|
||||
coreDown = true
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
if _, err := v.Verify(context.Background(), token); err != nil {
|
||||
t.Fatalf("verify sollte trotz core-ausfall mit letztbekanntem stand funktionieren: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: neue sicherheitskritische Aktionen
|
||||
// (z.B. neuer Login) schlagen bei Core-Ausfall klar fehl statt unsicher
|
||||
// durchgelassen zu werden — auch wenn ein (aelterer) Cache-Stand existiert.
|
||||
func TestRequireFreshKeys_FailsClosedWhenCoreUnreachable(t *testing.T) {
|
||||
km, err := NewKeyManager()
|
||||
if err != nil {
|
||||
t.Fatalf("new key manager: %v", err)
|
||||
}
|
||||
|
||||
coreDown := false
|
||||
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
if coreDown {
|
||||
return nil, errors.New("core nicht erreichbar (simuliert)")
|
||||
}
|
||||
return km.PublicKeySet(), nil
|
||||
})
|
||||
|
||||
// Cache vorwaermen (existiert jetzt ein "veralteter" gueltiger Stand).
|
||||
if _, _, err := v.cache.Get(context.Background()); err != nil {
|
||||
t.Fatalf("warm cache: %v", err)
|
||||
}
|
||||
|
||||
coreDown = true
|
||||
if err := v.RequireFreshKeys(context.Background()); err == nil {
|
||||
t.Fatal("erwartet fehler (fail-closed) bei core-ausfall, habe nil")
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 3: Schluesselrotation ohne Ausfallzeit —
|
||||
// ein bereits ausgestelltes Token bleibt nach Rotation weiterhin
|
||||
// verifizierbar, ein zweites (simuliertes) Modul bekommt beide Schluessel.
|
||||
func TestRotate_NoDowntimeForAlreadyIssuedTokens(t *testing.T) {
|
||||
km, err := NewKeyManager()
|
||||
if err != nil {
|
||||
t.Fatalf("new key manager: %v", err)
|
||||
}
|
||||
|
||||
oldToken, err := km.Issue("user-1", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue (alt): %v", err)
|
||||
}
|
||||
|
||||
if _, err := km.Rotate(); err != nil {
|
||||
t.Fatalf("rotate: %v", err)
|
||||
}
|
||||
|
||||
newToken, err := km.Issue("user-2", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue (neu): %v", err)
|
||||
}
|
||||
|
||||
// Simuliertes zweites Modul: fragt den vollstaendigen Schluesselsatz ab.
|
||||
moduleB := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
return km.PublicKeySet(), nil
|
||||
})
|
||||
|
||||
if _, err := moduleB.Verify(context.Background(), oldToken); err != nil {
|
||||
t.Fatalf("altes token sollte nach rotation weiterhin gueltig sein: %v", err)
|
||||
}
|
||||
if _, err := moduleB.Verify(context.Background(), newToken); err != nil {
|
||||
t.Fatalf("neues token sollte gueltig sein: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestVerify_RejectsUnknownKid(t *testing.T) {
|
||||
km1, _ := NewKeyManager()
|
||||
km2, _ := NewKeyManager() // komplett anderer, unbekannter schluessel
|
||||
|
||||
token, err := km1.Issue("user-1", "acme", time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("issue: %v", err)
|
||||
}
|
||||
|
||||
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
return km2.PublicKeySet(), nil // kennt km1s schluessel nicht
|
||||
})
|
||||
if _, err := v.Verify(context.Background(), token); !errors.Is(err, ErrInvalidToken) {
|
||||
t.Fatalf("erwartet ErrInvalidToken, habe %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestJWKSRoundTrip(t *testing.T) {
|
||||
km, _ := NewKeyManager()
|
||||
km.Rotate()
|
||||
|
||||
var buf []byte
|
||||
rec := &captureWriter{}
|
||||
km.ServeJWKS(rec, nil)
|
||||
buf = rec.body
|
||||
|
||||
parsed, err := ParseJWKS(buf)
|
||||
if err != nil {
|
||||
t.Fatalf("parse jwks: %v", err)
|
||||
}
|
||||
if len(parsed) != len(km.PublicKeySet()) {
|
||||
t.Fatalf("erwartet %d schluessel, habe %d", len(km.PublicKeySet()), len(parsed))
|
||||
}
|
||||
}
|
||||
|
||||
type captureWriter struct {
|
||||
body []byte
|
||||
header http.Header
|
||||
}
|
||||
|
||||
func (w *captureWriter) Header() http.Header {
|
||||
if w.header == nil {
|
||||
w.header = http.Header{}
|
||||
}
|
||||
return w.header
|
||||
}
|
||||
func (w *captureWriter) Write(p []byte) (int, error) { w.body = append(w.body, p...); return len(p), nil }
|
||||
func (w *captureWriter) WriteHeader(statusCode int) {}
|
||||
@@ -1,67 +0,0 @@
|
||||
package moduletrust
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"errors"
|
||||
"time"
|
||||
|
||||
"github.com/golang-jwt/jwt/v5"
|
||||
)
|
||||
|
||||
var ErrInvalidToken = errors.New("moduletrust: ungueltiges token")
|
||||
|
||||
// KeyFetchFunc holt den aktuellen Schluesselsatz von Core (z.B. per HTTP-GET
|
||||
// auf ServeJWKS + ParseJWKS). Wird vom Verifier nur bei abgelaufener TTL
|
||||
// aufgerufen — NICHT bei jeder Verify()-Anfrage (Akzeptanzkriterium 1).
|
||||
type KeyFetchFunc func(ctx context.Context) (map[string]ed25519.PublicKey, error)
|
||||
|
||||
// Verifier ist die Modulseite von API-05: verifiziert JWTs LOKAL gegen einen
|
||||
// per StaleCache zwischengespeicherten Schluesselsatz, ohne pro Aufruf einen
|
||||
// synchronen Request an Core zu stellen.
|
||||
type Verifier struct {
|
||||
cache *StaleCache[map[string]ed25519.PublicKey]
|
||||
}
|
||||
|
||||
func NewVerifier(ttl time.Duration, fetch KeyFetchFunc) *Verifier {
|
||||
return &Verifier{cache: NewStaleCache(ttl, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
|
||||
return fetch(ctx)
|
||||
})}
|
||||
}
|
||||
|
||||
// Verify prueft die Signatur LOKAL gegen den (ggf. abgelaufenen, aber
|
||||
// vorhandenen) Schluesselsatz — FAIL-OPEN fuer bereits ausgestellte Tokens
|
||||
// (Akzeptanzkriterium 2): ist Core nicht erreichbar, aber ein alter
|
||||
// Schluesselsatz bekannt, wird damit weiter verifiziert.
|
||||
func (v *Verifier) Verify(ctx context.Context, tokenString string) (*Claims, error) {
|
||||
keys, _, err := v.cache.Get(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
claims := &Claims{}
|
||||
token, err := jwt.ParseWithClaims(tokenString, claims, func(t *jwt.Token) (interface{}, error) {
|
||||
if _, ok := t.Method.(*jwt.SigningMethodEd25519); !ok {
|
||||
return nil, ErrInvalidToken
|
||||
}
|
||||
kid, _ := t.Header["kid"].(string)
|
||||
pub, ok := keys[kid]
|
||||
if !ok {
|
||||
return nil, ErrInvalidToken
|
||||
}
|
||||
return pub, nil
|
||||
})
|
||||
if err != nil || !token.Valid {
|
||||
return nil, ErrInvalidToken
|
||||
}
|
||||
return claims, nil
|
||||
}
|
||||
|
||||
// RequireFreshKeys ruft IMMER frisch von Core ab (FAIL-CLOSED) — fuer
|
||||
// sicherheitskritische Aktionen wie einen komplett neuen Login
|
||||
// (Akzeptanzkriterium 2): schlaegt klar fehl, wenn Core nicht erreichbar
|
||||
// ist, statt auf einem veralteten Schluesselsatz zu vertrauen.
|
||||
func (v *Verifier) RequireFreshKeys(ctx context.Context) error {
|
||||
_, err := v.cache.RequireFresh(ctx)
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
// Package notify implementiert Core CFG-02: den zentralen Benachrichtigungs-
|
||||
// Dispatcher, ueber den beliebige Module Benachrichtigungen ausloesen —
|
||||
// Warteschlange, Wiederholungslogik, Kanal-Abstraktion. Die tatsaechlichen
|
||||
// Kanaele (E-Mail/In-App) sind CFG-03, hier gibt es nur die Sender-
|
||||
// Schnittstelle als Vorbereitung.
|
||||
package notify
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// DefaultMaxAttempts begrenzt Wiederholungsversuche (Akzeptanzkriterium 2) —
|
||||
// nach dieser Anzahl gibt der Dispatcher kontrolliert auf (status=failed)
|
||||
// statt endlos zu wiederholen.
|
||||
const DefaultMaxAttempts = 5
|
||||
|
||||
// DefaultRetryBackoff ist die Basis-Wartezeit zwischen Wiederholungen,
|
||||
// linear mit der Versuchsnummer skaliert.
|
||||
const DefaultRetryBackoff = 200 * time.Millisecond
|
||||
|
||||
type Notification struct {
|
||||
ID string
|
||||
Channel string
|
||||
Recipient string
|
||||
Payload map[string]any
|
||||
Attempts int
|
||||
}
|
||||
|
||||
// Sender ist die schmale Schnittstelle, die ein konkreter Kanal (CFG-03)
|
||||
// implementiert. Der Dispatcher selbst weiss nichts ueber E-Mail/In-App.
|
||||
type Sender interface {
|
||||
Send(ctx context.Context, n Notification) error
|
||||
}
|
||||
|
||||
// Dispatcher ist die EINE Schnittstelle, ueber die Module Benachrichtigungen
|
||||
// ausloesen — kein Modul baut eigenen Versandcode (Akzeptanzkriterium 1).
|
||||
type Dispatcher struct {
|
||||
pool *pgxpool.Pool
|
||||
maxAttempts int
|
||||
retryBackoff time.Duration
|
||||
}
|
||||
|
||||
func NewDispatcher(pool *pgxpool.Pool) *Dispatcher {
|
||||
return &Dispatcher{pool: pool, maxAttempts: DefaultMaxAttempts, retryBackoff: DefaultRetryBackoff}
|
||||
}
|
||||
|
||||
// WithRetryPolicy erlaubt Tests/Betrieb, Versuchsanzahl und Backoff
|
||||
// anzupassen, ohne die Default-Policy im Produktionscode zu veraendern.
|
||||
func (d *Dispatcher) WithRetryPolicy(maxAttempts int, backoff time.Duration) *Dispatcher {
|
||||
return &Dispatcher{pool: d.pool, maxAttempts: maxAttempts, retryBackoff: backoff}
|
||||
}
|
||||
|
||||
// Enqueue reiht eine Benachrichtigung in die Postgres-Warteschlange ein und
|
||||
// kehrt sofort zurueck — die Zeile ueberlebt jeden Neustart des Dispatcher-
|
||||
// Prozesses unveraendert (Akzeptanzkriterium 3), da sie ausschliesslich in
|
||||
// der Datenbank existiert, nicht im Prozessspeicher.
|
||||
func (d *Dispatcher) Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (string, error) {
|
||||
if payload == nil {
|
||||
payload = map[string]any{}
|
||||
}
|
||||
payloadJSON, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("payload serialisieren: %w", err)
|
||||
}
|
||||
|
||||
var id string
|
||||
err = d.pool.QueryRow(ctx, `
|
||||
INSERT INTO notification_jobs (channel, recipient, payload, max_attempts)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
RETURNING id
|
||||
`, channel, recipient, payloadJSON, d.maxAttempts).Scan(&id)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("benachrichtigung einreihen: %w", err)
|
||||
}
|
||||
return id, nil
|
||||
}
|
||||
@@ -0,0 +1,219 @@
|
||||
package notify
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupTest(t *testing.T) (*pgxpool.Pool, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS notification_jobs (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
channel TEXT NOT NULL,
|
||||
recipient TEXT NOT NULL,
|
||||
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'sent', 'failed')),
|
||||
attempts INT NOT NULL DEFAULT 0,
|
||||
max_attempts INT NOT NULL DEFAULT 5,
|
||||
next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
last_error TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
)`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() { pool.Close() }
|
||||
return pool, cleanup
|
||||
}
|
||||
|
||||
type fakeSender struct {
|
||||
mu sync.Mutex
|
||||
sentIDs []string
|
||||
failUntil int
|
||||
calls int
|
||||
}
|
||||
|
||||
func (f *fakeSender) Send(ctx context.Context, n Notification) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.calls++
|
||||
if f.calls <= f.failUntil {
|
||||
return errors.New("simulierter zustellfehler")
|
||||
}
|
||||
f.sentIDs = append(f.sentIDs, n.ID)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeSender) sentCount() int {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return len(f.sentIDs)
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1: Module loesen ueber Enqueue aus, keine eigene
|
||||
// Versandlogik noetig.
|
||||
func TestDispatcher_EnqueueAndProcess(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
d := NewDispatcher(pool)
|
||||
id, err := d.Enqueue(ctx, "email", "alice@example.com", map[string]any{"subject": "Willkommen"})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
if id == "" {
|
||||
t.Fatal("erwartet nicht-leere id")
|
||||
}
|
||||
|
||||
sender := &fakeSender{}
|
||||
sent, failed, err := d.ProcessDue(ctx, sender, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("process: %v", err)
|
||||
}
|
||||
if sent != 1 || failed != 0 {
|
||||
t.Fatalf("erwartet sent=1 failed=0, habe sent=%d failed=%d", sent, failed)
|
||||
}
|
||||
if sender.sentCount() != 1 {
|
||||
t.Fatalf("erwartet 1 zustellung, habe %d", sender.sentCount())
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 2: Wiederholungslogik greift bei
|
||||
// simuliertem Fehler und bricht nach definierter Anzahl kontrolliert ab.
|
||||
func TestProcessDue_RetriesThenGivesUpAfterMaxAttempts(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
d := NewDispatcher(pool).WithRetryPolicy(3, time.Millisecond)
|
||||
id, err := d.Enqueue(ctx, "email", "bob@example.com", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
sender := &fakeSender{failUntil: 100} // schlaegt bei jedem versuch fehl
|
||||
|
||||
for i := 0; i < 3; i++ {
|
||||
time.Sleep(5 * time.Millisecond) // next_attempt_at abwarten
|
||||
if _, _, err := d.ProcessDue(ctx, sender, 10); err != nil {
|
||||
t.Fatalf("process %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
var status string
|
||||
var attempts int
|
||||
if err := pool.QueryRow(ctx, `SELECT status, attempts FROM notification_jobs WHERE id = $1`, id).Scan(&status, &attempts); err != nil {
|
||||
t.Fatalf("status lesen: %v", err)
|
||||
}
|
||||
if status != "failed" {
|
||||
t.Fatalf("erwartet status failed nach max_attempts, habe %q", status)
|
||||
}
|
||||
if attempts != 3 {
|
||||
t.Fatalf("erwartet 3 versuche, habe %d", attempts)
|
||||
}
|
||||
|
||||
// Weiteres ProcessDue darf den bereits aufgegebenen job nicht mehr anfassen.
|
||||
sent, failed, err := d.ProcessDue(ctx, sender, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("process nach abbruch: %v", err)
|
||||
}
|
||||
if sent != 0 || failed != 0 {
|
||||
t.Fatalf("erwartet keine weitere verarbeitung, habe sent=%d failed=%d", sent, failed)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 1: Neustart des Dienstes waehrend offener
|
||||
// Zustellung verliert keine Nachricht — simuliert durch eine komplett neue
|
||||
// Dispatcher/Pool-Instanz nach dem Enqueue, bevor irgendetwas verarbeitet wurde.
|
||||
func TestQueue_SurvivesRestartWithoutMessageLoss(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
firstInstance := NewDispatcher(pool)
|
||||
id, err := firstInstance.Enqueue(ctx, "email", "carol@example.com", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
// "Neustart": eine voellig neue Dispatcher-Instanz (repraesentiert einen
|
||||
// neuen Prozess) verbindet sich neu und verarbeitet die Warteschlange —
|
||||
// die Nachricht existiert ausschliesslich in Postgres, nicht im
|
||||
// Prozessspeicher der ersten Instanz.
|
||||
restartedInstance := NewDispatcher(pool)
|
||||
sender := &fakeSender{}
|
||||
sent, failed, err := restartedInstance.ProcessDue(ctx, sender, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("process nach neustart: %v", err)
|
||||
}
|
||||
if sent != 1 || failed != 0 {
|
||||
t.Fatalf("erwartet sent=1 nach neustart, habe sent=%d failed=%d", sent, failed)
|
||||
}
|
||||
if len(sender.sentIDs) != 1 || sender.sentIDs[0] != id {
|
||||
t.Fatalf("erwartet zustellung der urspruenglichen nachricht %q, habe %v", id, sender.sentIDs)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 3: zwei gleichzeitig ausloesende Module,
|
||||
// beide Nachrichten werden korrekt (und nicht doppelt) zugestellt.
|
||||
func TestProcessDue_ConcurrentDispatchBothDelivered(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
d := NewDispatcher(pool)
|
||||
idA, err := d.Enqueue(ctx, "email", "modul-a@example.com", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue a: %v", err)
|
||||
}
|
||||
idB, err := d.Enqueue(ctx, "email", "modul-b@example.com", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue b: %v", err)
|
||||
}
|
||||
|
||||
sender := &fakeSender{}
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 2; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if _, _, err := d.ProcessDue(ctx, sender, 10); err != nil {
|
||||
t.Errorf("process: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if sender.sentCount() != 2 {
|
||||
t.Fatalf("erwartet genau 2 zustellungen, habe %d: %v", sender.sentCount(), sender.sentIDs)
|
||||
}
|
||||
seen := map[string]bool{}
|
||||
for _, id := range sender.sentIDs {
|
||||
if seen[id] {
|
||||
t.Fatalf("nachricht %q wurde doppelt zugestellt", id)
|
||||
}
|
||||
seen[id] = true
|
||||
}
|
||||
if !seen[idA] || !seen[idB] {
|
||||
t.Fatalf("erwartet beide nachrichten zugestellt, habe %v", sender.sentIDs)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
package notify
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ProcessDue holt bis zu limit faellige Benachrichtigungen und versucht sie
|
||||
// ueber sender zuzustellen. FOR UPDATE SKIP LOCKED serialisiert konkurrierende
|
||||
// Aufrufe (Akzeptanzkriterium 3 / Pruefung 3: zwei gleichzeitig ausloesende
|
||||
// Module duerfen sich nicht gegenseitig blockieren oder Nachrichten doppelt
|
||||
// zustellen) — dieselbe Konvention wie internal/tenant.Lifecycle.ProcessDueDeletions.
|
||||
func (d *Dispatcher) ProcessDue(ctx context.Context, sender Sender, limit int) (sent, failed int, err error) {
|
||||
tx, err := d.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return 0, 0, fmt.Errorf("transaktion starten: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
rows, err := tx.Query(ctx, `
|
||||
SELECT id, channel, recipient, payload, attempts, max_attempts
|
||||
FROM notification_jobs
|
||||
WHERE status = 'pending' AND next_attempt_at <= now()
|
||||
ORDER BY created_at
|
||||
FOR UPDATE SKIP LOCKED
|
||||
LIMIT $1
|
||||
`, limit)
|
||||
if err != nil {
|
||||
return 0, 0, fmt.Errorf("faellige benachrichtigungen abfragen: %w", err)
|
||||
}
|
||||
|
||||
type due struct {
|
||||
id, channel, recipient string
|
||||
payload []byte
|
||||
attempts, maxAttempts int
|
||||
}
|
||||
var candidates []due
|
||||
for rows.Next() {
|
||||
var c due
|
||||
if err := rows.Scan(&c.id, &c.channel, &c.recipient, &c.payload, &c.attempts, &c.maxAttempts); err != nil {
|
||||
rows.Close()
|
||||
return 0, 0, fmt.Errorf("faellige benachrichtigung lesen: %w", err)
|
||||
}
|
||||
candidates = append(candidates, c)
|
||||
}
|
||||
rows.Close()
|
||||
if err := rows.Err(); err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
|
||||
for _, c := range candidates {
|
||||
var payload map[string]any
|
||||
if err := json.Unmarshal(c.payload, &payload); err != nil {
|
||||
payload = map[string]any{}
|
||||
}
|
||||
|
||||
sendErr := sender.Send(ctx, Notification{
|
||||
ID: c.id, Channel: c.channel, Recipient: c.recipient, Payload: payload, Attempts: c.attempts,
|
||||
})
|
||||
|
||||
if sendErr == nil {
|
||||
if _, err := tx.Exec(ctx, `
|
||||
UPDATE notification_jobs SET status = 'sent', updated_at = now() WHERE id = $1
|
||||
`, c.id); err != nil {
|
||||
return sent, failed, fmt.Errorf("erfolg speichern: %w", err)
|
||||
}
|
||||
sent++
|
||||
continue
|
||||
}
|
||||
|
||||
newAttempts := c.attempts + 1
|
||||
if newAttempts >= c.maxAttempts {
|
||||
// Akzeptanzkriterium 2: kontrollierter Abbruch nach definierter
|
||||
// Anzahl Versuche, kein endloses Wiederholen.
|
||||
if _, err := tx.Exec(ctx, `
|
||||
UPDATE notification_jobs
|
||||
SET status = 'failed', attempts = $2, last_error = $3, updated_at = now()
|
||||
WHERE id = $1
|
||||
`, c.id, newAttempts, sendErr.Error()); err != nil {
|
||||
return sent, failed, fmt.Errorf("fehlschlag speichern: %w", err)
|
||||
}
|
||||
failed++
|
||||
continue
|
||||
}
|
||||
|
||||
nextAttempt := time.Now().Add(time.Duration(newAttempts) * d.retryBackoff)
|
||||
if _, err := tx.Exec(ctx, `
|
||||
UPDATE notification_jobs
|
||||
SET attempts = $2, next_attempt_at = $3, last_error = $4, updated_at = now()
|
||||
WHERE id = $1
|
||||
`, c.id, newAttempts, nextAttempt, sendErr.Error()); err != nil {
|
||||
return sent, failed, fmt.Errorf("wiederholung planen: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return 0, 0, fmt.Errorf("transaktion committen: %w", err)
|
||||
}
|
||||
return sent, failed, nil
|
||||
}
|
||||
@@ -1,155 +0,0 @@
|
||||
// Package resync implementiert Core API-06: Wiederanlauf & Nachsynchro-
|
||||
// nisierung nach einem Core-Ausfall.
|
||||
//
|
||||
// - Ein Modul puffert Audit-Events und Nutzungszaehler-Deltas LOKAL in
|
||||
// Postgres (NICHT im Speicher — siehe "Bekannte Fehler vermeiden" im
|
||||
// Ticket: ein erneuter Ausfall waehrend der Nachlieferung darf keine
|
||||
// Daten verlieren, eine In-Memory-Queue wuerde das riskieren).
|
||||
// - Ein Health-Check-getriggerter Worker erkennt die Core-Wiedererreich-
|
||||
// barkeit SOFORT (nicht erst nach TTL-Ablauf, siehe StaleCache.Invalidate)
|
||||
// und liefert die gepufferten Daten in ORIGINALER Reihenfolge, authenti-
|
||||
// fiziert ueber das Service-Credential aus API-02
|
||||
// (internal/moduleregistry.Registry.Authenticate).
|
||||
//
|
||||
// Dieses Paket dupliziert weder internal/audit (AUD-01) noch internal/usage
|
||||
// (LIC-03) — es liefert nur den Puffer- und Nachlieferungs-Mechanismus
|
||||
// DAVOR bzw. DANACH.
|
||||
package resync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// BufferedAuditEvent ist ein lokal gepuffertes Audit-Ereignis. Seq
|
||||
// garantiert die Wiederherstellung der urspruenglichen Reihenfolge
|
||||
// (Akzeptanzkriterium 2) unabhaengig von eventuellen Uhrzeit-Ungenauigkeiten.
|
||||
type BufferedAuditEvent struct {
|
||||
ID int64
|
||||
Seq int64
|
||||
TenantSlug string
|
||||
Actor string
|
||||
Action string
|
||||
Target string
|
||||
Metadata map[string]any
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// BufferedUsageDelta ist ein lokal gepuffertes Nutzungszaehler-Inkrement.
|
||||
type BufferedUsageDelta struct {
|
||||
ID int64
|
||||
Seq int64
|
||||
TenantSlug string
|
||||
Metric string
|
||||
Delta int64
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// Buffer ist die lokale, persistente Pufferqueue eines Moduls.
|
||||
type Buffer struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewBuffer(pool *pgxpool.Pool) *Buffer {
|
||||
return &Buffer{pool: pool}
|
||||
}
|
||||
|
||||
// EnqueueAuditEvent puffert EIN Audit-Ereignis lokal — wird von einem
|
||||
// Fachmodul aufgerufen, wenn Core gerade nicht erreichbar ist (die
|
||||
// Erkennung "Core erreichbar oder nicht" ist NICHT Teil dieses Aufrufs,
|
||||
// siehe Worker).
|
||||
func (b *Buffer) EnqueueAuditEvent(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any) error {
|
||||
if metadata == nil {
|
||||
metadata = map[string]any{}
|
||||
}
|
||||
metadataJSON, err := json.Marshal(metadata)
|
||||
if err != nil {
|
||||
return fmt.Errorf("metadata serialisieren: %w", err)
|
||||
}
|
||||
_, err = b.pool.Exec(ctx, `
|
||||
INSERT INTO resync_audit_buffer (tenant_slug, actor, action, target, metadata)
|
||||
VALUES ($1, $2, $3, $4, $5)
|
||||
`, tenantSlug, actor, action, target, metadataJSON)
|
||||
if err != nil {
|
||||
return fmt.Errorf("audit-ereignis puffern: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// EnqueueUsageDelta puffert EIN Nutzungszaehler-Inkrement lokal.
|
||||
func (b *Buffer) EnqueueUsageDelta(ctx context.Context, tenantSlug, metric string, delta int64) error {
|
||||
_, err := b.pool.Exec(ctx, `
|
||||
INSERT INTO resync_usage_buffer (tenant_slug, metric, delta)
|
||||
VALUES ($1, $2, $3)
|
||||
`, tenantSlug, metric, delta)
|
||||
if err != nil {
|
||||
return fmt.Errorf("nutzungsdelta puffern: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// PendingAuditEvents liefert ALLE noch nicht zugestellten Audit-Events in
|
||||
// ORIGINALER Reihenfolge (Akzeptanzkriterium 2 / Pruefung 1).
|
||||
func (b *Buffer) PendingAuditEvents(ctx context.Context) ([]BufferedAuditEvent, error) {
|
||||
rows, err := b.pool.Query(ctx, `
|
||||
SELECT id, seq, tenant_slug, actor, action, target, metadata, created_at
|
||||
FROM resync_audit_buffer ORDER BY seq
|
||||
`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("gepufferte audit-events abfragen: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []BufferedAuditEvent
|
||||
for rows.Next() {
|
||||
var e BufferedAuditEvent
|
||||
var metadataJSON []byte
|
||||
if err := rows.Scan(&e.ID, &e.Seq, &e.TenantSlug, &e.Actor, &e.Action, &e.Target, &metadataJSON, &e.CreatedAt); err != nil {
|
||||
return nil, fmt.Errorf("gepuffertes audit-event lesen: %w", err)
|
||||
}
|
||||
_ = json.Unmarshal(metadataJSON, &e.Metadata)
|
||||
out = append(out, e)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// PendingUsageDeltas liefert ALLE noch nicht zugestellten Nutzungsdeltas.
|
||||
func (b *Buffer) PendingUsageDeltas(ctx context.Context) ([]BufferedUsageDelta, error) {
|
||||
rows, err := b.pool.Query(ctx, `
|
||||
SELECT id, seq, tenant_slug, metric, delta, created_at
|
||||
FROM resync_usage_buffer ORDER BY seq
|
||||
`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("gepufferte nutzungsdeltas abfragen: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []BufferedUsageDelta
|
||||
for rows.Next() {
|
||||
var d BufferedUsageDelta
|
||||
if err := rows.Scan(&d.ID, &d.Seq, &d.TenantSlug, &d.Metric, &d.Delta, &d.CreatedAt); err != nil {
|
||||
return nil, fmt.Errorf("gepuffertes nutzungsdelta lesen: %w", err)
|
||||
}
|
||||
out = append(out, d)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// RemoveAuditEvent entfernt EIN Audit-Event aus dem Puffer — wird NUR nach
|
||||
// von Core BESTAETIGTER Zustellung aufgerufen (Akzeptanzkriterium 2/3: erst
|
||||
// entfernen, wenn sicher zugestellt, sonst bleibt es fuer den naechsten
|
||||
// Versuch erhalten — kein Datenverlust bei erneutem Ausfall waehrend der
|
||||
// Nachlieferung).
|
||||
func (b *Buffer) RemoveAuditEvent(ctx context.Context, id int64) error {
|
||||
_, err := b.pool.Exec(ctx, `DELETE FROM resync_audit_buffer WHERE id = $1`, id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (b *Buffer) RemoveUsageDelta(ctx context.Context, id int64) error {
|
||||
_, err := b.pool.Exec(ctx, `DELETE FROM resync_usage_buffer WHERE id = $1`, id)
|
||||
return err
|
||||
}
|
||||
@@ -1,132 +0,0 @@
|
||||
package resync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
// AuditRecorder ist die schmale Schnittstelle, ueber die Core empfangene
|
||||
// Audit-Events tatsaechlich persistiert. In Produktion durch internal/audit
|
||||
// (AUD-01, nicht Abhaengigkeit dieser Kachel) implementiert — dieses Paket
|
||||
// dupliziert dessen Validierungs-/Speicherlogik NICHT, sondern ruft sie nur
|
||||
// auf. Die Events werden vom Aufrufer sequenziell in PendingAuditEvents-
|
||||
// Reihenfolge uebergeben (Akzeptanzkriterium 2).
|
||||
type AuditRecorder interface {
|
||||
Record(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any, occurredAt time.Time) error
|
||||
}
|
||||
|
||||
// UsageIncrementer ist die schmale Schnittstelle zu Core's Nutzungszaehler
|
||||
// (in Produktion internal/usage, LIC-03 — nicht Abhaengigkeit dieser
|
||||
// Kachel). Jedes gepufferte Delta wird GENAU EINMAL angewendet.
|
||||
type UsageIncrementer interface {
|
||||
Increment(ctx context.Context, tenantSlug, metric string, delta int64) error
|
||||
}
|
||||
|
||||
// CredentialAuthenticator ist die schmale Schnittstelle zu API-02s
|
||||
// Service-Credential-Pruefung (internal/moduleregistry.Registry.Authenticate).
|
||||
type CredentialAuthenticator interface {
|
||||
Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error)
|
||||
}
|
||||
|
||||
// Handler nimmt nachgelieferte Audit-Events/Nutzungsdeltas auf der
|
||||
// Core-Seite entgegen — authentifiziert ueber dasselbe Service-Credential
|
||||
// wie jeder andere Modul-Core-Aufruf (Akzeptanzkriterium 2/3, "authentifiziert
|
||||
// ueber das in API-02 definierte Service-Credential").
|
||||
type Handler struct {
|
||||
auth CredentialAuthenticator
|
||||
audit AuditRecorder
|
||||
usage UsageIncrementer
|
||||
}
|
||||
|
||||
func NewHandler(auth CredentialAuthenticator, audit AuditRecorder, usage UsageIncrementer) *Handler {
|
||||
return &Handler{auth: auth, audit: audit, usage: usage}
|
||||
}
|
||||
|
||||
type credentialHeader struct {
|
||||
ClientID string `json:"client_id"`
|
||||
Secret string `json:"secret"`
|
||||
}
|
||||
|
||||
func (h *Handler) authenticate(w http.ResponseWriter, r *http.Request) bool {
|
||||
clientID := r.Header.Get("X-Nexarch-Client-Id")
|
||||
secret := r.Header.Get("X-Nexarch-Client-Secret")
|
||||
_, ok, err := h.auth.Authenticate(r.Context(), clientID, secret)
|
||||
if err != nil || !ok {
|
||||
http.Error(w, "ungueltiges service-credential", http.StatusUnauthorized)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
type auditEventDTO struct {
|
||||
TenantSlug string `json:"tenant_slug"`
|
||||
Actor string `json:"actor"`
|
||||
Action string `json:"action"`
|
||||
Target string `json:"target"`
|
||||
Metadata map[string]any `json:"metadata"`
|
||||
OccurredAt time.Time `json:"occurred_at"`
|
||||
}
|
||||
|
||||
// AuditHandler nimmt EINE Liste gepufferter Audit-Events entgegen und
|
||||
// schreibt sie SEQUENZIELL in der gegebenen Reihenfolge fort
|
||||
// (Akzeptanzkriterium 2 / Pruefung 1: Vollstaendigkeit + Reihenfolge).
|
||||
// Bricht die Verarbeitung bei einem Fehler ab und meldet, wie viele Events
|
||||
// bereits sicher geschrieben wurden — der Aufrufer (Worker) entfernt aus
|
||||
// seinem lokalen Puffer NUR die bestaetigt geschriebenen Events.
|
||||
func (h *Handler) AuditHandler(w http.ResponseWriter, r *http.Request) {
|
||||
if !h.authenticate(w, r) {
|
||||
return
|
||||
}
|
||||
var events []auditEventDTO
|
||||
if err := json.NewDecoder(r.Body).Decode(&events); err != nil {
|
||||
http.Error(w, "ungueltiger anfrage-koerper", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
written := 0
|
||||
for _, e := range events {
|
||||
if err := h.audit.Record(r.Context(), e.TenantSlug, e.Actor, e.Action, e.Target, e.Metadata, e.OccurredAt); err != nil {
|
||||
break
|
||||
}
|
||||
written++
|
||||
}
|
||||
|
||||
writeJSON(w, map[string]int{"written": written})
|
||||
}
|
||||
|
||||
type usageDeltaDTO struct {
|
||||
TenantSlug string `json:"tenant_slug"`
|
||||
Metric string `json:"metric"`
|
||||
Delta int64 `json:"delta"`
|
||||
}
|
||||
|
||||
// UsageHandler wendet JEDES gepufferte Delta GENAU EINMAL an
|
||||
// (Akzeptanzkriterium 3 / Pruefung 2: keine Doppelzaehlung) — der Aufrufer
|
||||
// entfernt aus seinem lokalen Puffer nur die bestaetigt uebernommenen Deltas.
|
||||
func (h *Handler) UsageHandler(w http.ResponseWriter, r *http.Request) {
|
||||
if !h.authenticate(w, r) {
|
||||
return
|
||||
}
|
||||
var deltas []usageDeltaDTO
|
||||
if err := json.NewDecoder(r.Body).Decode(&deltas); err != nil {
|
||||
http.Error(w, "ungueltiger anfrage-koerper", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
applied := 0
|
||||
for _, d := range deltas {
|
||||
if err := h.usage.Increment(r.Context(), d.TenantSlug, d.Metric, d.Delta); err != nil {
|
||||
break
|
||||
}
|
||||
applied++
|
||||
}
|
||||
|
||||
writeJSON(w, map[string]int{"applied": applied})
|
||||
}
|
||||
|
||||
func writeJSON(w http.ResponseWriter, body any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(body)
|
||||
}
|
||||
@@ -1,320 +0,0 @@
|
||||
package resync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
|
||||
"gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry"
|
||||
)
|
||||
|
||||
// fakeAudit steht fuer internal/audit.Log (AUD-01, nicht Abhaengigkeit
|
||||
// dieser Kachel) — zeichnet Aufrufe in Empfangsreihenfolge auf, damit
|
||||
// Vollstaendigkeit UND Reihenfolge geprueft werden koennen.
|
||||
type fakeAudit struct {
|
||||
mu sync.Mutex
|
||||
events []auditEventDTO
|
||||
failAt int // -1 = nie fehlschlagen; sonst: ab diesem Index (0-basiert) schlaegt Record fehl
|
||||
}
|
||||
|
||||
func (f *fakeAudit) Record(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any, occurredAt time.Time) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
if f.failAt >= 0 && len(f.events) == f.failAt {
|
||||
return fmt.Errorf("simulierter core-ausfall waehrend der nachlieferung")
|
||||
}
|
||||
f.events = append(f.events, auditEventDTO{TenantSlug: tenantSlug, Actor: actor, Action: action, Target: target, Metadata: metadata, OccurredAt: occurredAt})
|
||||
return nil
|
||||
}
|
||||
|
||||
// fakeUsage steht fuer internal/usage.Store (LIC-03, nicht Abhaengigkeit
|
||||
// dieser Kachel) — summiert Deltas wie der echte Store.
|
||||
type fakeUsage struct {
|
||||
mu sync.Mutex
|
||||
totals map[string]int64
|
||||
}
|
||||
|
||||
func newFakeUsage() *fakeUsage { return &fakeUsage{totals: map[string]int64{}} }
|
||||
|
||||
func (f *fakeUsage) Increment(ctx context.Context, tenantSlug, metric string, delta int64) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.totals[tenantSlug+"|"+metric] += delta
|
||||
return nil
|
||||
}
|
||||
|
||||
func setupTest(t *testing.T) (*Buffer, *pgxpool.Pool, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS resync_audit_buffer (
|
||||
id BIGSERIAL PRIMARY KEY, seq BIGSERIAL, tenant_slug TEXT NOT NULL, actor TEXT NOT NULL,
|
||||
action TEXT NOT NULL, target TEXT NOT NULL DEFAULT '', metadata JSONB NOT NULL DEFAULT '{}',
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS resync_usage_buffer (
|
||||
id BIGSERIAL PRIMARY KEY, seq BIGSERIAL, tenant_slug TEXT NOT NULL, metric TEXT NOT NULL,
|
||||
delta BIGINT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS feature_flags (
|
||||
key TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT false,
|
||||
rollout_percentage INT NOT NULL DEFAULT 0, target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS modules (
|
||||
name TEXT PRIMARY KEY, version TEXT NOT NULL CHECK (version <> ''),
|
||||
required_flags TEXT[] NOT NULL DEFAULT '{}', registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS module_credentials (
|
||||
module_name TEXT PRIMARY KEY REFERENCES modules(name),
|
||||
client_id TEXT NOT NULL UNIQUE, secret_hash BYTEA NOT NULL,
|
||||
issued_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() { pool.Close() }
|
||||
return NewBuffer(pool), pool, cleanup
|
||||
}
|
||||
|
||||
func uniqueModuleName() string {
|
||||
return fmt.Sprintf("resync-test-%d", time.Now().UnixNano())
|
||||
}
|
||||
|
||||
// setupModuleCredential registriert ein echtes Modul + Service-Credential
|
||||
// ueber internal/moduleregistry (API-02) — dieselbe Authentifizierung wird
|
||||
// vom Handler tatsaechlich geprueft, kein Mock.
|
||||
func setupModuleCredential(t *testing.T, pool *pgxpool.Pool) (registry *moduleregistry.Registry, clientID, secret string) {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
flagService := flag.NewService(flag.NewStore(pool), 10*time.Millisecond)
|
||||
registry = moduleregistry.NewRegistry(pool, flagService)
|
||||
name := uniqueModuleName()
|
||||
if _, err := registry.Register(ctx, name, "1.0.0", nil); err != nil {
|
||||
t.Fatalf("modul registrieren: %v", err)
|
||||
}
|
||||
clientID, secret, err := registry.Provision(ctx, name)
|
||||
if err != nil {
|
||||
t.Fatalf("credential provisionieren: %v", err)
|
||||
}
|
||||
return registry, clientID, secret
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 1: waehrend eines simulierten Ausfalls
|
||||
// lokal gepufferte Audit-Events sind nach Wiederanlauf vollstaendig und in
|
||||
// korrekter Reihenfolge in Core's Audit-Log vorhanden.
|
||||
func TestFlushAll_DeliversBufferedAuditEventsCompleteAndInOrder(t *testing.T) {
|
||||
buffer, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
registry, clientID, secret := setupModuleCredential(t, pool)
|
||||
|
||||
tenantSlug := "acme"
|
||||
// Ereignisse "waehrend core down" lokal puffern.
|
||||
actions := []string{"login", "upload", "delete", "logout"}
|
||||
for _, action := range actions {
|
||||
if err := buffer.EnqueueAuditEvent(ctx, tenantSlug, "user-1", action, "res-1", nil); err != nil {
|
||||
t.Fatalf("enqueue %s: %v", action, err)
|
||||
}
|
||||
}
|
||||
|
||||
audit := &fakeAudit{failAt: -1}
|
||||
usage := newFakeUsage()
|
||||
handler := NewHandler(registry, audit, usage)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/internal/resync/audit", handler.AuditHandler)
|
||||
mux.HandleFunc("/internal/resync/usage", handler.UsageHandler)
|
||||
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
|
||||
server := httptest.NewServer(mux)
|
||||
defer server.Close()
|
||||
|
||||
client := NewCoreClient(server.URL, clientID, secret)
|
||||
worker := NewWorker(buffer, client)
|
||||
|
||||
if err := worker.FlushAll(ctx); err != nil {
|
||||
t.Fatalf("flushall: %v", err)
|
||||
}
|
||||
|
||||
audit.mu.Lock()
|
||||
defer audit.mu.Unlock()
|
||||
if len(audit.events) != len(actions) {
|
||||
t.Fatalf("erwartet %d zugestellte events, habe %d", len(actions), len(audit.events))
|
||||
}
|
||||
for i, e := range audit.events {
|
||||
if e.Action != actions[i] {
|
||||
t.Fatalf("reihenfolge falsch: position %d = %q, want %q", i, e.Action, actions[i])
|
||||
}
|
||||
}
|
||||
|
||||
// Puffer muss nach bestaetigter Zustellung leer sein.
|
||||
remaining, err := buffer.PendingAuditEvents(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(remaining) != 0 {
|
||||
t.Fatalf("erwartet leeren puffer nach bestaetigter zustellung, habe %d verbleibende", len(remaining))
|
||||
}
|
||||
}
|
||||
|
||||
// Bekannter-Fehler-Praevention: bricht die Zustellung waehrend der
|
||||
// Nachlieferung erneut ab (Core faellt wieder aus), bleiben die NICHT
|
||||
// bestaetigten Events sicher im Puffer erhalten statt verloren zu gehen.
|
||||
func TestFlushAll_KeepsUnconfirmedEventsInBufferOnPartialFailure(t *testing.T) {
|
||||
buffer, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
registry, clientID, secret := setupModuleCredential(t, pool)
|
||||
|
||||
for i := 0; i < 5; i++ {
|
||||
if err := buffer.EnqueueAuditEvent(ctx, "acme", "user-1", fmt.Sprintf("action-%d", i), "", nil); err != nil {
|
||||
t.Fatalf("enqueue %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Core-Fake schlaegt AB dem 3. Event fehl -> simuliert erneuten Ausfall
|
||||
// mitten in der Nachlieferung.
|
||||
audit := &fakeAudit{failAt: 3}
|
||||
handler := NewHandler(registry, audit, newFakeUsage())
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/internal/resync/audit", handler.AuditHandler)
|
||||
mux.HandleFunc("/internal/resync/usage", handler.UsageHandler)
|
||||
server := httptest.NewServer(mux)
|
||||
defer server.Close()
|
||||
|
||||
client := NewCoreClient(server.URL, clientID, secret)
|
||||
worker := NewWorker(buffer, client)
|
||||
|
||||
if err := worker.FlushAll(ctx); err == nil {
|
||||
t.Fatal("erwartet fehler, da core nur teilweise bestaetigt hat")
|
||||
}
|
||||
|
||||
remaining, err := buffer.PendingAuditEvents(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(remaining) != 2 {
|
||||
t.Fatalf("erwartet 2 verbleibende (nicht bestaetigte) events im puffer, habe %d", len(remaining))
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 2: Nutzungszaehler-Differenz aus der
|
||||
// Ausfallzeit wird korrekt nachgebucht, ein wiederholter (fehlerhafter)
|
||||
// Flush-Versuch fuehrt NICHT zu Doppelzaehlung, weil bereits bestaetigte
|
||||
// Deltas aus dem Puffer entfernt sind.
|
||||
func TestFlushAll_AppliesUsageDeltasWithoutDoubleCounting(t *testing.T) {
|
||||
buffer, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
registry, clientID, secret := setupModuleCredential(t, pool)
|
||||
|
||||
deltas := []int64{3, 5, 2}
|
||||
var want int64
|
||||
for _, d := range deltas {
|
||||
want += d
|
||||
if err := buffer.EnqueueUsageDelta(ctx, "acme", "api_calls", d); err != nil {
|
||||
t.Fatalf("enqueue delta %d: %v", d, err)
|
||||
}
|
||||
}
|
||||
|
||||
usage := newFakeUsage()
|
||||
handler := NewHandler(registry, &fakeAudit{failAt: -1}, usage)
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/internal/resync/audit", handler.AuditHandler)
|
||||
mux.HandleFunc("/internal/resync/usage", handler.UsageHandler)
|
||||
server := httptest.NewServer(mux)
|
||||
defer server.Close()
|
||||
|
||||
client := NewCoreClient(server.URL, clientID, secret)
|
||||
worker := NewWorker(buffer, client)
|
||||
|
||||
if err := worker.FlushAll(ctx); err != nil {
|
||||
t.Fatalf("flushall 1: %v", err)
|
||||
}
|
||||
|
||||
usage.mu.Lock()
|
||||
got := usage.totals["acme|api_calls"]
|
||||
usage.mu.Unlock()
|
||||
if got != want {
|
||||
t.Fatalf("nutzungsstand nach nachbuchung = %d, want %d", got, want)
|
||||
}
|
||||
|
||||
// Ein zweiter Flush-Versuch (z.B. redundanter Retry) darf NICHTS mehr
|
||||
// nachbuchen, da der Puffer bereits geleert wurde.
|
||||
if err := worker.FlushAll(ctx); err != nil {
|
||||
t.Fatalf("flushall 2: %v", err)
|
||||
}
|
||||
usage.mu.Lock()
|
||||
got2 := usage.totals["acme|api_calls"]
|
||||
usage.mu.Unlock()
|
||||
if got2 != want {
|
||||
t.Fatalf("nutzungsstand nach redundantem zweiten flush = %d, want unveraendert %d (keine doppelzaehlung)", got2, want)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 3: bei erkannter Core-Wiedererreichbarkeit
|
||||
// wird der Cache SOFORT invalidiert (naechster Zugriff refetcht), nicht
|
||||
// erst nach TTL-Ablauf — Latenz wird gemessen und liegt weit unter einer
|
||||
// langen TTL.
|
||||
func TestCheckAndSync_InvalidatesCacheImmediatelyOnRecovery(t *testing.T) {
|
||||
buffer, pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
registry, clientID, secret := setupModuleCredential(t, pool)
|
||||
_ = registry
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
|
||||
server := httptest.NewServer(mux)
|
||||
defer server.Close()
|
||||
|
||||
client := NewCoreClient(server.URL, clientID, secret)
|
||||
worker := NewWorker(buffer, client)
|
||||
|
||||
invalidated := false
|
||||
var mu sync.Mutex
|
||||
worker.OnReachable(func() {
|
||||
mu.Lock()
|
||||
invalidated = true
|
||||
mu.Unlock()
|
||||
})
|
||||
|
||||
start := time.Now()
|
||||
becameReachable, err := worker.CheckAndSync(context.Background())
|
||||
elapsed := time.Since(start)
|
||||
if err != nil {
|
||||
t.Fatalf("checkandsync: %v", err)
|
||||
}
|
||||
if !becameReachable {
|
||||
t.Fatal("erwartet erkannten uebergang zu 'erreichbar' beim ersten erfolgreichen check")
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if !invalidated {
|
||||
t.Fatal("erwartet sofortigen cache-invalidierungs-callback bei core-wiedererreichbarkeit")
|
||||
}
|
||||
// Zielwert: deutlich unter einer typischen TTL (z.B. 5s beim
|
||||
// Feature-Flag-Cache, LIC-02) — hier im Millisekundenbereich, da rein
|
||||
// lokal ohne Netzwerk-Overhead.
|
||||
if elapsed > time.Second {
|
||||
t.Fatalf("cache-invalidierung brauchte %s, erwartet deutlich unter 1s", elapsed)
|
||||
}
|
||||
}
|
||||
@@ -1,220 +0,0 @@
|
||||
package resync
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
// CoreClient ist der modul-seitige HTTP-Client fuer die Nachlieferung,
|
||||
// authentifiziert ueber dasselbe Service-Credential wie jeder andere
|
||||
// Modul-Core-Aufruf (API-02).
|
||||
type CoreClient struct {
|
||||
BaseURL string
|
||||
ClientID string
|
||||
Secret string
|
||||
HTTP *http.Client
|
||||
}
|
||||
|
||||
func NewCoreClient(baseURL, clientID, secret string) *CoreClient {
|
||||
return &CoreClient{BaseURL: baseURL, ClientID: clientID, Secret: secret, HTTP: &http.Client{Timeout: 5 * time.Second}}
|
||||
}
|
||||
|
||||
func (c *CoreClient) post(ctx context.Context, path string, body any) (*http.Response, error) {
|
||||
payload, err := json.Marshal(body)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("payload serialisieren: %w", err)
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+path, bytes.NewReader(payload))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("X-Nexarch-Client-Id", c.ClientID)
|
||||
req.Header.Set("X-Nexarch-Client-Secret", c.Secret)
|
||||
return c.HTTP.Do(req)
|
||||
}
|
||||
|
||||
// HealthCheck prueft, ob Core erreichbar ist — dieselbe Konvention wie
|
||||
// internal/health (OPS-01): HTTP 200 auf einem Health-Endpunkt.
|
||||
func (c *CoreClient) HealthCheck(ctx context.Context) bool {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.BaseURL+"/healthz", nil)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
resp, err := c.HTTP.Do(req)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
return resp.StatusCode == http.StatusOK
|
||||
}
|
||||
|
||||
// Worker erkennt Core-Wiedererreichbarkeit und stoesst DANN sofort
|
||||
// (Akzeptanzkriterium 1) sowohl registrierte Cache-Invalidierungen als auch
|
||||
// das Nachliefern des lokalen Puffers an.
|
||||
type Worker struct {
|
||||
buffer *Buffer
|
||||
client *CoreClient
|
||||
onReachable []func()
|
||||
wasDown bool
|
||||
}
|
||||
|
||||
func NewWorker(buffer *Buffer, client *CoreClient) *Worker {
|
||||
return &Worker{buffer: buffer, client: client, wasDown: true} // Start pessimistisch: erster erfolgreicher Check zaehlt als "Wiedererreichbarkeit".
|
||||
}
|
||||
|
||||
// OnReachable registriert einen Callback, der bei jeder erkannten
|
||||
// Core-Wiedererreichbarkeit sofort ausgefuehrt wird — z.B.
|
||||
// moduletrust.StaleCache[T].Invalidate, damit der naechste Zugriff sofort
|
||||
// neu abruft statt auf TTL-Ablauf zu warten (Akzeptanzkriterium 1).
|
||||
func (w *Worker) OnReachable(fn func()) {
|
||||
w.onReachable = append(w.onReachable, fn)
|
||||
}
|
||||
|
||||
// CheckAndSync fuehrt EINEN Zyklus aus: Erreichbarkeit pruefen, bei
|
||||
// erkanntem UEBERGANG "nicht erreichbar -> erreichbar" sofort die
|
||||
// registrierten Callbacks ausloesen und den Puffer nachliefern. Gibt
|
||||
// zurueck, ob ein Wiederanlauf in diesem Aufruf erkannt wurde (fuer
|
||||
// Latenzmessung in Tests, Pruefung 3).
|
||||
func (w *Worker) CheckAndSync(ctx context.Context) (becameReachable bool, err error) {
|
||||
reachable := w.client.HealthCheck(ctx)
|
||||
if !reachable {
|
||||
w.wasDown = true
|
||||
return false, nil
|
||||
}
|
||||
|
||||
justRecovered := w.wasDown
|
||||
w.wasDown = false
|
||||
if !justRecovered {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
for _, fn := range w.onReachable {
|
||||
fn()
|
||||
}
|
||||
|
||||
if err := w.FlushAll(ctx); err != nil {
|
||||
return true, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// FlushAll liefert ZUERST alle gepufferten Audit-Events (in Reihenfolge),
|
||||
// DANN alle gepufferten Nutzungsdeltas nach. Jedes Element wird aus dem
|
||||
// lokalen Puffer NUR entfernt, wenn Core es bestaetigt hat — bricht die
|
||||
// Uebertragung vorzeitig ab (Core faellt waehrend der Nachlieferung erneut
|
||||
// aus), bleibt der Rest sicher im Postgres-Puffer erhalten
|
||||
// (Akzeptanzkriterium 2/3, "Bekannte Fehler vermeiden").
|
||||
func (w *Worker) FlushAll(ctx context.Context) error {
|
||||
if err := w.flushAuditEvents(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
return w.flushUsageDeltas(ctx)
|
||||
}
|
||||
|
||||
func (w *Worker) flushAuditEvents(ctx context.Context) error {
|
||||
events, err := w.buffer.PendingAuditEvents(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(events) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
dtos := make([]auditEventDTO, len(events))
|
||||
for i, e := range events {
|
||||
dtos[i] = auditEventDTO{
|
||||
TenantSlug: e.TenantSlug, Actor: e.Actor, Action: e.Action, Target: e.Target,
|
||||
Metadata: e.Metadata, OccurredAt: e.CreatedAt,
|
||||
}
|
||||
}
|
||||
|
||||
resp, err := w.client.post(ctx, "/internal/resync/audit", dtos)
|
||||
if err != nil {
|
||||
return fmt.Errorf("audit-nachlieferung: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("audit-nachlieferung: unerwarteter status %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var ack struct {
|
||||
Written int `json:"written"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&ack); err != nil {
|
||||
return fmt.Errorf("audit-bestaetigung lesen: %w", err)
|
||||
}
|
||||
|
||||
// NUR die von Core bestaetigt geschriebenen Events entfernen — sie sind
|
||||
// nach PendingAuditEvents-Reihenfolge sortiert, die ersten `Written`
|
||||
// Eintraege entsprechen also genau den bestaetigten.
|
||||
for i := 0; i < ack.Written; i++ {
|
||||
if err := w.buffer.RemoveAuditEvent(ctx, events[i].ID); err != nil {
|
||||
return fmt.Errorf("bestaetigtes audit-event aus puffer entfernen: %w", err)
|
||||
}
|
||||
}
|
||||
if ack.Written < len(events) {
|
||||
return fmt.Errorf("core hat nur %d von %d audit-events bestaetigt", ack.Written, len(events))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (w *Worker) flushUsageDeltas(ctx context.Context) error {
|
||||
deltas, err := w.buffer.PendingUsageDeltas(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(deltas) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
dtos := make([]usageDeltaDTO, len(deltas))
|
||||
for i, d := range deltas {
|
||||
dtos[i] = usageDeltaDTO{TenantSlug: d.TenantSlug, Metric: d.Metric, Delta: d.Delta}
|
||||
}
|
||||
|
||||
resp, err := w.client.post(ctx, "/internal/resync/usage", dtos)
|
||||
if err != nil {
|
||||
return fmt.Errorf("nutzungs-nachlieferung: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("nutzungs-nachlieferung: unerwarteter status %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var ack struct {
|
||||
Applied int `json:"applied"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&ack); err != nil {
|
||||
return fmt.Errorf("nutzungs-bestaetigung lesen: %w", err)
|
||||
}
|
||||
|
||||
for i := 0; i < ack.Applied; i++ {
|
||||
if err := w.buffer.RemoveUsageDelta(ctx, deltas[i].ID); err != nil {
|
||||
return fmt.Errorf("bestaetigtes nutzungsdelta aus puffer entfernen: %w", err)
|
||||
}
|
||||
}
|
||||
if ack.Applied < len(deltas) {
|
||||
return fmt.Errorf("core hat nur %d von %d nutzungsdeltas bestaetigt", ack.Applied, len(deltas))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Run fuehrt CheckAndSync in festen Abstaenden aus — die "kurze, definierte
|
||||
// Zeitspanne" aus Akzeptanzkriterium 1 ist dieses Poll-Intervall.
|
||||
func (w *Worker) Run(ctx context.Context, interval time.Duration) {
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
_, _ = w.CheckAndSync(ctx)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,32 +0,0 @@
|
||||
package usage
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"time"
|
||||
)
|
||||
|
||||
// AggregateFunc berechnet/aktualisiert Zaehlerstaende aus einer autoritativen
|
||||
// Quelle (z.B. "zaehle Zeilen in einer Modul-Tabelle") — die konkrete Quelle
|
||||
// haengt vom jeweiligen Modul ab und ist nicht Teil dieser Kachel. Das
|
||||
// Aggregations-Grundgerüst selbst (periodischer Trigger) ist es.
|
||||
type AggregateFunc func(ctx context.Context) error
|
||||
|
||||
// RunPeriodicAggregation ruft aggregate in festen Abstaenden auf, bis ctx
|
||||
// beendet wird — dieselbe In-Prozess-Worker-Goroutine-Konvention wie
|
||||
// internal/tenant.Lifecycle.RunSweeper (Akzeptanzkriterium 1: "periodisch
|
||||
// aggregiert").
|
||||
func RunPeriodicAggregation(ctx context.Context, interval time.Duration, aggregate AggregateFunc) {
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
if err := aggregate(ctx); err != nil {
|
||||
slog.Error("nutzungszaehler-aggregation fehlgeschlagen", "error", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,144 +0,0 @@
|
||||
// Package usage implementiert Core LIC-03: Nutzungszaehler je Tenant
|
||||
// (Benutzeranzahl, Speicherverbrauch, API-Aufrufe, ...) und die Pruefung
|
||||
// gegen konfigurierte Quotas. Quotas sind Konfiguration (Tabellenzeile), kein
|
||||
// Hardcode — Zitadel/Unleash-Vorbild.
|
||||
package usage
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
var ErrNoQuota = errors.New("usage: keine quota fuer diese metrik konfiguriert")
|
||||
|
||||
// Status ist die definierte Reaktion einer Quota-Pruefung (Akzeptanzkriterium 2).
|
||||
type Status string
|
||||
|
||||
const (
|
||||
StatusOK Status = "ok"
|
||||
StatusWarning Status = "warning" // Schwelle (80%) erreicht, aber noch nicht ueberschritten
|
||||
StatusExceeded Status = "exceeded" // Quota ueberschritten — neue Ressourcen sollten gesperrt werden
|
||||
)
|
||||
|
||||
// warningThreshold liegt bei 80% der Quota.
|
||||
const warningThreshold = 0.8
|
||||
|
||||
type Store struct {
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewStore(pool *pgxpool.Pool) *Store {
|
||||
return &Store{pool: pool}
|
||||
}
|
||||
|
||||
// Increment erhoeht einen Zaehler ATOMAR ueber ein einziges SQL-Statement
|
||||
// (UPSERT mit value = value + delta) statt Read-Modify-Write in Go — das
|
||||
// haelt Zaehlerstaende bei parallelen Schreibzugriffen konsistent
|
||||
// (Akzeptanzkriterium 1 / Pruefung 2), ohne eine Anwendungs-Transaktion mit
|
||||
// Lock zu brauchen.
|
||||
func (s *Store) Increment(ctx context.Context, tenantID, metric string, delta int64) error {
|
||||
_, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO usage_counters (tenant_id, metric, value, updated_at)
|
||||
VALUES ($1, $2, $3, now())
|
||||
ON CONFLICT (tenant_id, metric) DO UPDATE
|
||||
SET value = usage_counters.value + $3, updated_at = now()
|
||||
`, tenantID, metric, delta)
|
||||
if err != nil {
|
||||
return fmt.Errorf("zaehler erhoehen: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Get liefert den aktuellen Zaehlerstand — 0, wenn noch nie erhoeht wurde.
|
||||
// Der Wert ist strikt tenant-gescoped (Akzeptanzkriterium 3 / Pruefung 3).
|
||||
func (s *Store) Get(ctx context.Context, tenantID, metric string) (int64, error) {
|
||||
var value int64
|
||||
err := s.pool.QueryRow(ctx, `
|
||||
SELECT value FROM usage_counters WHERE tenant_id = $1 AND metric = $2
|
||||
`, tenantID, metric).Scan(&value)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return 0, nil
|
||||
}
|
||||
return 0, fmt.Errorf("zaehler lesen: %w", err)
|
||||
}
|
||||
return value, nil
|
||||
}
|
||||
|
||||
// SetQuota legt die Obergrenze fuer (tenantID, metric) fest — Konfiguration,
|
||||
// kein Hardcode.
|
||||
func (s *Store) SetQuota(ctx context.Context, tenantID, metric string, limit int64) error {
|
||||
_, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO usage_quotas (tenant_id, metric, limit_value)
|
||||
VALUES ($1, $2, $3)
|
||||
ON CONFLICT (tenant_id, metric) DO UPDATE SET limit_value = $3
|
||||
`, tenantID, metric, limit)
|
||||
if err != nil {
|
||||
return fmt.Errorf("quota setzen: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) GetQuota(ctx context.Context, tenantID, metric string) (int64, error) {
|
||||
var limit int64
|
||||
err := s.pool.QueryRow(ctx, `
|
||||
SELECT limit_value FROM usage_quotas WHERE tenant_id = $1 AND metric = $2
|
||||
`, tenantID, metric).Scan(&limit)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return 0, ErrNoQuota
|
||||
}
|
||||
return 0, fmt.Errorf("quota lesen: %w", err)
|
||||
}
|
||||
return limit, nil
|
||||
}
|
||||
|
||||
// Check liefert Zaehlerstand, konfigurierte Quota und die daraus abgeleitete
|
||||
// Reaktion (Akzeptanzkriterium 2 / Pruefung 1). Ist keine Quota konfiguriert,
|
||||
// gilt die Metrik als unbegrenzt (StatusOK).
|
||||
func (s *Store) Check(ctx context.Context, tenantID, metric string) (value, limit int64, status Status, err error) {
|
||||
value, err = s.Get(ctx, tenantID, metric)
|
||||
if err != nil {
|
||||
return 0, 0, "", err
|
||||
}
|
||||
|
||||
limit, err = s.GetQuota(ctx, tenantID, metric)
|
||||
if errors.Is(err, ErrNoQuota) {
|
||||
return value, 0, StatusOK, nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, 0, "", err
|
||||
}
|
||||
|
||||
switch {
|
||||
case value > limit:
|
||||
return value, limit, StatusExceeded, nil
|
||||
case limit > 0 && float64(value) >= warningThreshold*float64(limit):
|
||||
return value, limit, StatusWarning, nil
|
||||
default:
|
||||
return value, limit, StatusOK, nil
|
||||
}
|
||||
}
|
||||
|
||||
// Reaction wird aufgerufen, wenn Check einen Nicht-OK-Status liefert
|
||||
// (Akzeptanzkriterium 2: "definierte Reaktion").
|
||||
type Reaction func(ctx context.Context, tenantID, metric string, value, limit int64, status Status)
|
||||
|
||||
// Enforce fuehrt Check aus und ruft react auf, wenn der Status nicht OK ist —
|
||||
// die konkrete "Sperre neuer Ressourcen"/Benachrichtigung liegt beim
|
||||
// Aufrufer (z.B. TEN-02 vor dem Anlegen eines neuen Benutzers), Enforce
|
||||
// garantiert nur, dass die Reaktion zuverlaessig ausgeloest wird.
|
||||
func (s *Store) Enforce(ctx context.Context, tenantID, metric string, react Reaction) (Status, error) {
|
||||
value, limit, status, err := s.Check(ctx, tenantID, metric)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if status != StatusOK && react != nil {
|
||||
react(ctx, tenantID, metric, value, limit, status)
|
||||
}
|
||||
return status, nil
|
||||
}
|
||||
@@ -1,206 +0,0 @@
|
||||
package usage
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupTest(t *testing.T) (*Store, func()) {
|
||||
t.Helper()
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
ctx := context.Background()
|
||||
|
||||
pool, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("pool: %v", err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
CREATE TABLE IF NOT EXISTS usage_counters (
|
||||
tenant_id UUID NOT NULL, metric TEXT NOT NULL, value BIGINT NOT NULL DEFAULT 0,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (tenant_id, metric)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS usage_quotas (
|
||||
tenant_id UUID NOT NULL, metric TEXT NOT NULL, limit_value BIGINT NOT NULL,
|
||||
PRIMARY KEY (tenant_id, metric)
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() { pool.Close() }
|
||||
return NewStore(pool), cleanup
|
||||
}
|
||||
|
||||
func newTenantID() string {
|
||||
return fmt.Sprintf("00000000-0000-0000-0000-%012d", time.Now().UnixNano()%1e12)
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 + Pruefung 2: Aggregationsjob liefert bei parallelen
|
||||
// Schreibzugriffen konsistente Zaehlerstaende.
|
||||
func TestIncrement_ConsistentUnderConcurrentWrites(t *testing.T) {
|
||||
store, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
tenant := newTenantID()
|
||||
|
||||
const goroutines = 50
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < goroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if err := store.Increment(ctx, tenant, "api_calls", 1); err != nil {
|
||||
t.Errorf("increment: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
value, err := store.Get(ctx, tenant, "api_calls")
|
||||
if err != nil {
|
||||
t.Fatalf("get: %v", err)
|
||||
}
|
||||
if value != goroutines {
|
||||
t.Fatalf("erwartet %d, habe %d (hinweis auf lost update unter nebenlaeufigkeit)", goroutines, value)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 3: Zaehlerstand eines Tenants beeinflusst
|
||||
// nicht den eines anderen.
|
||||
func TestIncrement_IsolatedBetweenTenants(t *testing.T) {
|
||||
store, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
tenantA, tenantB := newTenantID(), newTenantID()
|
||||
|
||||
if err := store.Increment(ctx, tenantA, "users", 5); err != nil {
|
||||
t.Fatalf("increment a: %v", err)
|
||||
}
|
||||
if err := store.Increment(ctx, tenantB, "users", 1); err != nil {
|
||||
t.Fatalf("increment b: %v", err)
|
||||
}
|
||||
|
||||
valA, err := store.Get(ctx, tenantA, "users")
|
||||
if err != nil {
|
||||
t.Fatalf("get a: %v", err)
|
||||
}
|
||||
valB, err := store.Get(ctx, tenantB, "users")
|
||||
if err != nil {
|
||||
t.Fatalf("get b: %v", err)
|
||||
}
|
||||
if valA != 5 || valB != 1 {
|
||||
t.Fatalf("erwartet a=5 b=1, habe a=%d b=%d", valA, valB)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 1: Quota-Ueberschreitung wird automatisiert
|
||||
// erkannt und die definierte Reaktion ausgeloest.
|
||||
func TestEnforce_TriggersReactionOnExceeded(t *testing.T) {
|
||||
store, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
tenant := newTenantID()
|
||||
|
||||
if err := store.SetQuota(ctx, tenant, "users", 10); err != nil {
|
||||
t.Fatalf("set quota: %v", err)
|
||||
}
|
||||
if err := store.Increment(ctx, tenant, "users", 11); err != nil {
|
||||
t.Fatalf("increment: %v", err)
|
||||
}
|
||||
|
||||
var reacted bool
|
||||
var gotStatus Status
|
||||
status, err := store.Enforce(ctx, tenant, "users", func(ctx context.Context, tenantID, metric string, value, limit int64, status Status) {
|
||||
reacted = true
|
||||
gotStatus = status
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enforce: %v", err)
|
||||
}
|
||||
if status != StatusExceeded {
|
||||
t.Fatalf("erwartet StatusExceeded, habe %q", status)
|
||||
}
|
||||
if !reacted || gotStatus != StatusExceeded {
|
||||
t.Fatal("erwartet ausgeloeste reaktion mit StatusExceeded")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCheck_WarningThresholdAndOK(t *testing.T) {
|
||||
store, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
tenant := newTenantID()
|
||||
|
||||
if err := store.SetQuota(ctx, tenant, "storage_mb", 100); err != nil {
|
||||
t.Fatalf("set quota: %v", err)
|
||||
}
|
||||
|
||||
if err := store.Increment(ctx, tenant, "storage_mb", 50); err != nil {
|
||||
t.Fatalf("increment: %v", err)
|
||||
}
|
||||
_, _, status, err := store.Check(ctx, tenant, "storage_mb")
|
||||
if err != nil {
|
||||
t.Fatalf("check: %v", err)
|
||||
}
|
||||
if status != StatusOK {
|
||||
t.Fatalf("bei 50%% erwartet StatusOK, habe %q", status)
|
||||
}
|
||||
|
||||
if err := store.Increment(ctx, tenant, "storage_mb", 35); err != nil { // insgesamt 85%
|
||||
t.Fatalf("increment: %v", err)
|
||||
}
|
||||
_, _, status, err = store.Check(ctx, tenant, "storage_mb")
|
||||
if err != nil {
|
||||
t.Fatalf("check: %v", err)
|
||||
}
|
||||
if status != StatusWarning {
|
||||
t.Fatalf("bei 85%% erwartet StatusWarning, habe %q", status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCheck_NoQuotaMeansUnlimited(t *testing.T) {
|
||||
store, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
tenant := newTenantID()
|
||||
|
||||
if err := store.Increment(ctx, tenant, "api_calls", 1_000_000); err != nil {
|
||||
t.Fatalf("increment: %v", err)
|
||||
}
|
||||
_, _, status, err := store.Check(ctx, tenant, "api_calls")
|
||||
if err != nil {
|
||||
t.Fatalf("check: %v", err)
|
||||
}
|
||||
if status != StatusOK {
|
||||
t.Fatalf("ohne konfigurierte quota erwartet StatusOK, habe %q", status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunPeriodicAggregation_CallsRepeatedly(t *testing.T) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
var mu sync.Mutex
|
||||
calls := 0
|
||||
RunPeriodicAggregation(ctx, 20*time.Millisecond, func(ctx context.Context) error {
|
||||
mu.Lock()
|
||||
calls++
|
||||
mu.Unlock()
|
||||
return nil
|
||||
})
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if calls < 3 {
|
||||
t.Fatalf("erwartet mehrfache aufrufe innerhalb von 120ms bei 20ms interval, habe %d", calls)
|
||||
}
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
DROP TABLE IF EXISTS audit_events;
|
||||
@@ -1,17 +0,0 @@
|
||||
-- Zentrales Audit-Log-Modell (AUD-01, siehe core-kanban/tickets/AUD-01.md).
|
||||
-- Getrennt vom allgemeinen Anwendungs-Log (Akzeptanzkriterium 2): eigene
|
||||
-- Tabelle, eigenes Paket (internal/audit), kein Log-Framework.
|
||||
-- tenant_slug ist NOT NULL + darf nicht leer sein (Akzeptanzkriterium 2 /
|
||||
-- Pruefung 2) — mandantenuebergreifende Ereignisse nutzen den reservierten
|
||||
-- Wert 'system', niemals NULL oder leeren String.
|
||||
CREATE TABLE audit_events (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
occurred_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
tenant_slug TEXT NOT NULL CHECK (tenant_slug <> ''),
|
||||
actor TEXT NOT NULL CHECK (actor <> ''),
|
||||
action TEXT NOT NULL CHECK (action <> ''),
|
||||
target TEXT NOT NULL,
|
||||
metadata JSONB NOT NULL DEFAULT '{}'::jsonb
|
||||
);
|
||||
|
||||
CREATE INDEX audit_events_tenant_slug_idx ON audit_events (tenant_slug, occurred_at);
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP TABLE IF EXISTS config_value_history;
|
||||
DROP TABLE IF EXISTS config_values;
|
||||
@@ -0,0 +1,23 @@
|
||||
-- Zentraler Konfigurationsdienst (CFG-01, siehe core-kanban/tickets/CFG-01.md).
|
||||
-- scope = 'global' fuer globale Defaults, sonst der Tenant-Slug. config_values
|
||||
-- haelt den AKTUELLEN Stand je (key, scope); config_value_history haelt JEDE
|
||||
-- Aenderung fest (Akzeptanzkriterium 2: versioniert nachvollziehbar).
|
||||
CREATE TABLE config_values (
|
||||
key TEXT NOT NULL,
|
||||
scope TEXT NOT NULL CHECK (scope <> ''),
|
||||
value TEXT NOT NULL,
|
||||
version INT NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (key, scope)
|
||||
);
|
||||
|
||||
CREATE TABLE config_value_history (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
key TEXT NOT NULL,
|
||||
scope TEXT NOT NULL,
|
||||
value TEXT NOT NULL,
|
||||
version INT NOT NULL,
|
||||
changed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE INDEX config_value_history_key_scope_idx ON config_value_history (key, scope, version);
|
||||
@@ -1 +0,0 @@
|
||||
DROP TABLE IF EXISTS tenant_licenses;
|
||||
@@ -1,12 +0,0 @@
|
||||
-- Lizenzumfang pro Mandant (LIC-01, siehe core-kanban/tickets/LIC-01.md).
|
||||
-- Genau ein Lizenzdatensatz pro Tenant (tenant_id PK) — ein neues Einspielen
|
||||
-- ersetzt den vorherigen Datensatz vollstaendig statt eine Historie zu fuehren.
|
||||
CREATE TABLE tenant_licenses (
|
||||
tenant_id UUID PRIMARY KEY REFERENCES tenants(id),
|
||||
plan TEXT NOT NULL,
|
||||
modules TEXT[] NOT NULL,
|
||||
issued_at TIMESTAMPTZ NOT NULL,
|
||||
valid_until TIMESTAMPTZ NOT NULL,
|
||||
raw_key TEXT NOT NULL,
|
||||
installed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
@@ -0,0 +1 @@
|
||||
DROP TABLE IF EXISTS notification_jobs;
|
||||
@@ -0,0 +1,20 @@
|
||||
-- Benachrichtigungs-Dispatcher-Warteschlange (CFG-02, siehe
|
||||
-- core-kanban/tickets/CFG-02.md). Postgres-basiert statt Redis/AMQP
|
||||
-- (Projekt-Konvention, siehe nexarch-state.json techstack.job_queue) —
|
||||
-- Zeilen ueberleben einen Neustart des Dispatcher-Prozesses unveraendert
|
||||
-- (Akzeptanzkriterium 3).
|
||||
CREATE TABLE notification_jobs (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
channel TEXT NOT NULL,
|
||||
recipient TEXT NOT NULL,
|
||||
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'sent', 'failed')),
|
||||
attempts INT NOT NULL DEFAULT 0,
|
||||
max_attempts INT NOT NULL DEFAULT 5,
|
||||
next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
last_error TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE INDEX notification_jobs_due_idx ON notification_jobs (status, next_attempt_at);
|
||||
@@ -1,2 +0,0 @@
|
||||
DROP TABLE IF EXISTS usage_quotas;
|
||||
DROP TABLE IF EXISTS usage_counters;
|
||||
@@ -1,15 +0,0 @@
|
||||
-- Nutzungszaehler & Quotas je Tenant (LIC-03, siehe core-kanban/tickets/LIC-03.md).
|
||||
CREATE TABLE usage_counters (
|
||||
tenant_id UUID NOT NULL,
|
||||
metric TEXT NOT NULL,
|
||||
value BIGINT NOT NULL DEFAULT 0,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (tenant_id, metric)
|
||||
);
|
||||
|
||||
CREATE TABLE usage_quotas (
|
||||
tenant_id UUID NOT NULL,
|
||||
metric TEXT NOT NULL,
|
||||
limit_value BIGINT NOT NULL,
|
||||
PRIMARY KEY (tenant_id, metric)
|
||||
);
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP TABLE alert_debounce_state;
|
||||
DROP TABLE alert_rules;
|
||||
@@ -0,0 +1,24 @@
|
||||
-- OPS-05: Schwellwert-Regeln fuer Alerting auf den aus OPS-03 aggregierten
|
||||
-- Metriken. Lebt wie config_values/notification_jobs (CFG-01/02) in der
|
||||
-- zentralen Registry-DB — modulübergreifende Betriebskonfiguration, keine
|
||||
-- Mandanten-Geschaeftsdaten.
|
||||
CREATE TABLE alert_rules (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
metric_name TEXT NOT NULL,
|
||||
comparison TEXT NOT NULL CHECK (comparison IN ('gt', 'lt')),
|
||||
threshold DOUBLE PRECISION NOT NULL,
|
||||
label_filters JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
recipient TEXT NOT NULL,
|
||||
description TEXT NOT NULL DEFAULT '',
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Haelt fest, wann eine Regel zuletzt tatsaechlich einen Alarm ausgeloest
|
||||
-- hat (Akzeptanzkriterium 3: Drosselung wiederholter Alarmierung fuer
|
||||
-- denselben anhaltenden Zustand). rule_key kombiniert Regel-ID mit den
|
||||
-- tatsaechlichen Label-Werten der ausloesenden Zeitreihe, damit dieselbe
|
||||
-- Regel fuer unterschiedliche Tenants/Module unabhaengig gedrosselt wird.
|
||||
CREATE TABLE alert_debounce_state (
|
||||
rule_key TEXT PRIMARY KEY,
|
||||
last_fired_at TIMESTAMPTZ NOT NULL
|
||||
);
|
||||
@@ -0,0 +1 @@
|
||||
DROP TABLE metrics_sources;
|
||||
@@ -0,0 +1,7 @@
|
||||
-- Metrics-Aggregation ueber Module hinweg (OPS-03, siehe
|
||||
-- core-kanban/tickets/OPS-03.md) — welches Modul liefert seine Kennzahlen
|
||||
-- unter welcher /metrics-URL.
|
||||
CREATE TABLE metrics_sources (
|
||||
module_name TEXT PRIMARY KEY,
|
||||
metrics_url TEXT NOT NULL
|
||||
);
|
||||
@@ -1,2 +0,0 @@
|
||||
DROP TABLE resync_usage_buffer;
|
||||
DROP TABLE resync_audit_buffer;
|
||||
@@ -1,22 +0,0 @@
|
||||
-- Wiederanlauf & Nachsynchronisierung nach Core-Ausfall (API-06, siehe
|
||||
-- core-kanban/tickets/API-06.md) — lokale, PERSISTENTE Pufferqueue (kein
|
||||
-- In-Memory) fuer Audit-Events und Nutzungszaehler-Deltas eines Moduls.
|
||||
CREATE TABLE resync_audit_buffer (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
seq BIGSERIAL,
|
||||
tenant_slug TEXT NOT NULL,
|
||||
actor TEXT NOT NULL,
|
||||
action TEXT NOT NULL,
|
||||
target TEXT NOT NULL DEFAULT '',
|
||||
metadata JSONB NOT NULL DEFAULT '{}',
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE TABLE resync_usage_buffer (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
seq BIGSERIAL,
|
||||
tenant_slug TEXT NOT NULL,
|
||||
metric TEXT NOT NULL,
|
||||
delta BIGINT NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
@@ -15,6 +15,8 @@ export PGPASSWORD="$PASS"
|
||||
|
||||
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants CASCADE;"
|
||||
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS audit_events CASCADE;"
|
||||
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS config_value_history CASCADE;"
|
||||
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS config_values CASCADE;"
|
||||
|
||||
dbs=$(psql -h localhost -U "$ROLE" -d postgres -tAc "SELECT datname FROM pg_database WHERE datname LIKE 'tenant\_%' ESCAPE '\'")
|
||||
for db in $dbs; do
|
||||
|
||||
Reference in New Issue
Block a user