Compare commits

..
Author SHA1 Message Date
sysops d6e9412043 API-11: wiederanlauf-nachsynchronisierungs-endpunkt-starten
- Merge API-02 (bereits Vorfahre), AUD-01, LIC-03 in API-06-Branch, um
  echte Produktionsimplementierungen fuer resync.Handler zu erhalten
  (moduleregistry.Registry, audit.Log, usage.Store) - go.mod/go.sum/
  scripts/reset-test-env.sh additiv zusammengefuehrt, keine Logikkonflikte
- cmd/resync-api: startet internal/resync.Handler (API-06) als
  eigenstaendigen HTTP-Dienst, kleiner auditAdapter fuer Signatur-
  Anpassung (kein neuer Fachcode)
- reines Wiring, kein Diff an internal/resync|audit|usage|moduleregistry
  (verifiziert)
- real deployed auf 131, end-zu-ende per curl mit echtem, ueber
  moduleregistry.Registry.Provision ausgestelltem Service-Credential:
  angewendetes usage-delta real in usage_counters bestaetigt, falsches
  Credential -> 401
- reale Grant-Luecke gefunden und behoben (nexarch_core auf modules/
  module_credentials/audit_events/feature_flags/usage_counters/
  resync_*_buffer), ueber information_schema verifiziert

Pruefungen siehe docs/API-11-PRUEFPROTOKOLL.md
2026-08-30 20:14:07 +02:00
sysops e6892535ed Merge branch 'feature/lic-03-nutzungszaehler-quotas' into feature/api-11-wiederanlauf-nachsynchronisierungs-endpunkt-starten
# Conflicts:
#	go.mod
#	scripts/reset-test-env.sh
2026-08-30 20:09:41 +02:00
sysops 1c41b65a47 Merge branch 'feature/aud-01-zentrales-audit-log-modell' into feature/api-11-wiederanlauf-nachsynchronisierungs-endpunkt-starten
# Conflicts:
#	go.mod
#	go.sum
#	scripts/reset-test-env.sh
2026-08-30 20:09:26 +02:00
sysops ec9bb27bb3 API-06: go.sum/go.mod aktualisieren (golang-jwt/jwt/v5 fuer moduletrust) 2026-08-28 09:22:56 +02:00
sysops 0fd9856b10 API-06: wiederanlauf-nachsynchronisierung-nach-core-ausfall (postgres-puffer, service-credential, sofort-invalidate) 2026-08-28 09:22:08 +02:00
sysopsandClaude Sonnet 5 b23cd1961f API-02: modul-registry-aktivierungspruefung
internal/moduleregistry: Registry.Register traegt Fachmodule mit Name,
Version und benoetigten Feature-Flags ein (Akzeptanzkriterium 1), fehlende
Pflichtangaben werden abgewiesen. IsActive kombiniert Registrierung + LIC-02
Feature-Flag-Auswertung (ALLE benoetigten Flags muessen fuer den Tenant
aktiv sein) — ein nicht registriertes Modul ist nie aktiv. List liefert alle
Module fuer Statusseite/Lizenzoberflaeche (Akzeptanzkriterium 3).

RequireActiveModule ist die zentrale Durchsetzungs-Middleware (Casbin-
Prinzip): weist Anfragen an ein deaktiviertes Modul ab, BEVOR der
Modul-Handler ueberhaupt aufgerufen wird (Akzeptanzkriterium 2) —
Pruefung per Test belegt, dass der Handler bei Deaktivierung nachweislich
nicht erreicht wird.

Service-Credentials (Akzeptanzkriterium 4): Registry.Provision stellt pro
Modul-Instanz Client-ID + Secret aus, gespeichert wird nur der SHA-256-Hash
des Secrets. Registry.Authenticate vergleicht timing-safe (dieselbe
subtle.ConstantTimeCompare-Referenzimplementierung wie AUD-02).
RequireServiceCredential-Middleware liest X-Client-Id/X-Client-Secret und
weist Aufrufe ohne gueltiges Credential mit 401 ab, bevor der Core-seitige
Endpunkt (z.B. Audit-Nachlieferung, Nutzungsmeldung) erreicht wird.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Anfrage an deaktiviertes Modul nachweislich vor Modul-Logik abgewiesen —
   TestRequireActiveModule_BlocksBeforeHandler: handlerReached bleibt false
   bei 403, wird true erst nach Aktivierung bei 200. PASS.
2. Registrierung mit fehlenden Pflichtangaben abgewiesen —
   TestRegister_RejectsMissingFields (leerer Name, leere Version). PASS.
3. Registry-Abfrage liefert konsistente Daten nach Aktivierung/Deaktivierung —
   TestIsActive_ReflectsFlagStateConsistently: aus/an/aus-Zyklus, IsActive
   folgt dem Flag-Zustand korrekt. PASS.
4. Aufruf mit ungueltigem/fehlendem Service-Credential abgewiesen, mit
   gueltigem angenommen — TestRequireServiceCredential_RejectsInvalidAcceptsValid
   und TestProvisionAndAuthenticate (falsches Secret, unbekannte Client-ID,
   korrektes Credential). PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 21:17:50 +02:00
sysopsandClaude Sonnet 5 d447869246 LIC-03: nutzungszaehler-quotas
internal/usage: Store.Increment aktualisiert Zaehlerstaende ueber ein
einziges atomares SQL-UPSERT (value = value + delta) statt Read-Modify-Write
in Go — haelt Zaehlerstaende bei parallelen Schreibzugriffen konsistent
(Akzeptanzkriterium 1), ganz ohne Anwendungs-Lock. Quotas sind Konfiguration
(usage_quotas-Tabelle je Tenant+Metrik), kein Hardcode.

Check/Enforce leiten aus Zaehlerstand + Quota eine definierte Reaktion ab
(StatusOK/Warning bei 80%/Exceeded, Akzeptanzkriterium 2) — Enforce ruft eine
uebergebene Reaction-Funktion auf, wenn der Status nicht OK ist; die
konkrete Sperr-/Benachrichtigungslogik bleibt beim Aufrufer (z.B. TEN-02 vor
Benutzeranlage), Enforce garantiert nur zuverlaessiges Ausloesen. Fehlende
Quota-Konfiguration bedeutet unbegrenzt (StatusOK), kein Fehler.

RunPeriodicAggregation ist das Aggregations-Grundgerüst (Akzeptanzkriterium 1:
"periodisch aggregiert") — dieselbe In-Prozess-Worker-Goroutine-Konvention
wie internal/tenant.Lifecycle.RunSweeper. Die konkrete Aggregationsquelle
(Zeilen zaehlen in Modul-Tabellen) haengt vom jeweiligen Modul ab und ist
nicht Teil dieser Kachel.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Quota-Ueberschreitung automatisiert erkannt, definierte Reaktion
   ausgeloest — TestEnforce_TriggersReactionOnExceeded: Reaction-Callback
   wird mit StatusExceeded aufgerufen. PASS.
2. Aggregationsjob liefert bei parallelen Schreibzugriffen konsistente
   Zaehlerstaende — TestIncrement_ConsistentUnderConcurrentWrites: 50
   nebenlaeufige Increments, Endstand exakt 50 (kein Lost Update). PASS.
3. Zaehlerstand eines Tenants beeinflusst nicht den eines anderen —
   TestIncrement_IsolatedBetweenTenants. PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 21:10:11 +02:00
sysopsandClaude Sonnet 5 b12d53f469 AUD-01: zentrales-audit-log-modell
internal/audit: eigenes, strukturiertes Audit-Datenmodell (Akteur, Aktion,
Zielobjekt, Zeitpunkt, Tenant) in der Registry-DB, getrennt von jedem
allgemeinen Anwendungs-Log (eigenes Paket, eigene Tabelle audit_events,
kein Logging-Framework). Log.Record ist der EINE zentrale Schreibpfad —
es gibt keine zweite Schreibmoeglichkeit, ueber die ein Handler die
Validierung umgehen koennte.

Fehlender Tenant-Bezug wird zweifach verhindert (Akzeptanzkriterium 2):
Log.Record weist leeren TenantSlug direkt ab (ErrMissingTenant), zusaetzlich
erzwingt eine CHECK-Constraint in der Migration dasselbe auf Datenbankebene,
selbst wenn Log.Record umgangen wuerde. Mandantenuebergreifende Ereignisse
(z.B. Superadmin-Aktionen) nutzen den reservierten Wert audit.SystemTenant
statt NULL oder leerem String — es gibt keinen Weg, ganz ohne Tenant-Bezug
zu schreiben.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Automatisierter Test belegt genau einen Audit-Eintrag pro
   sicherheitsrelevantem Vorgang — TestRecord_PersistsExactlyOneEventPerSecurityIncident
   (simulierter fehlgeschlagener Login), Feldinhalte verifiziert. PASS.
2. Fehlender Tenant-Bezug durch Constraint/Test verhindert —
   TestRecord_RejectsMissingTenant (App-Ebene) UND
   TestConstraint_RejectsMissingTenantAtDatabaseLevel (direkter INSERT unter
   Umgehung von Log.Record, durch CHECK-Constraint abgewiesen). PASS.
3. Datenmodell von zweiter Person gegen Dokumentation geprueft — NICHT
   durchgefuehrt (keine zweite Person in dieser Session verfuegbar). Offen.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:22:50 +02:00
sysopsandClaude Sonnet 5 4a30345e07 LIC-02: feature-flag-service-je-tenant
internal/flag: Store (Verwaltung) + Service (Auswertung mit TTL-Cache,
Default 5s) — Unleash-Prinzip Flag-Verwaltung vs. Flag-Auswertung getrennt,
als Kernfunktion des Core-Dienstes selbst statt separater Infrastruktur.

evaluate() wendet drei Strategien in fester Reihenfolge an: global an/aus,
Tenant-Zielgruppe, deterministischer Prozentsatz-Rollout (FNV-Hash aus
Tenant+Key, stabil pro Tenant). IsEnabled liefert IMMER nur bool (kein
Fehlerwert) — ein nicht erreichbarer Flag-Dienst kann damit keinen
Aufrufer zum Absturz bringen: bei DB-Fehler wird der zuletzt bekannte
Cache-Stand verwendet, ohne jeglichen Stand faellt der Dienst sicher auf
false zurueck. Service.Invalidate erzwingt sofortiges Neuladen fuer den
Schreiber selbst, andere Instanzen sehen Aenderungen spaetestens nach der
TTL (Akzeptanzkriterium 3, kein Neustart noetig).

Bugfix waehrend Tests: Store.Set uebergab ein nil-TargetTenantSlugs-Slice
als SQL NULL statt leerem Array (NOT-NULL-Verletzung) — auf leeres Slice
normalisiert.

Akzeptanzkriterium 4 (Deaktivierung loescht keine Daten): dieses Paket
besitzt ausschliesslich die eigene feature_flags-Zeile, hat keinerlei
Code-Pfad, der Modul-Geschaeftsdaten anfassen koennte — Loeschung bleibt
strukturell der Archive-Retention-Engine vorbehalten.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Cache-Invalidierungszeit automatisiert gemessen —
   TestService_CacheInvalidationTiming: Aenderung wirksam nach 153ms bei
   TTL=150ms (innerhalb Ziel+Toleranz), vorher nachweislich noch alter Stand. PASS.
2. Zielgruppen-Strategie liefert erwartete Auswertung —
   TestService_TargetTenantStrategy / TestEvaluate_TargetTenantStrategy. PASS.
3. Ausfall des Flag-Dienstes fuehrt zu dokumentiertem Fallback, kein Absturz —
   TestService_FallsBackOnStoreFailure (mit recover()-Absicherung): Fallback
   auf Cache-Stand bzw. sicheres false bei komplett unerreichbarer DB, geloggt. PASS.
4. Modul-Deaktivierung/Reaktivierung ohne Datenverlust — architektonisch durch
   fehlenden Code-Pfad sichergestellt (siehe oben), zusaetzlich durch
   TestService_InvalidateForcesImmediateRefresh (Toggle aus/an bleibt
   konsistent nachvollziehbar) mitabgedeckt. PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:19:59 +02:00
sysopsandClaude Sonnet 5 d63db93a4c LIC-01: lizenzmodell-lizenzschluessel-pruefung
internal/license: Ed25519-signierte Lizenzschluessel (stdlib crypto/ed25519,
keine neue Abhaengigkeit). Issuer haelt den privaten Schluessel (lebt beim
Lizenzgeber), Validator nur den oeffentlichen (lebt im Core-Prozess) — klare
Trennung Ausstellung/Pruefung nach Unleash-Vorbild (Flag-Verwaltung vs.
Flag-Auswertung).

Store.Install prueft NUR die Signatur und persistiert den Lizenzumfang
(Plan, Modul-Liste, Laufzeit) in tenant_licenses (Registry-DB, 1:1 zu
tenants). Eine bereits abgelaufene, aber korrekt signierte Lizenz laesst
sich trotzdem einspielen — der Ablauf wird erst bei Store.RequireActive
bewertet (liefert ErrLicenseExpired statt Panic/Absturz), waehrend
Store.Status den Umfang unabhaengig vom Ablauf weiterhin liefert.

Neu: scripts/run-checks.sh buendelt reset-test-env.sh + go build/vet/test
(-p 1) zu einem Ein-Kommando-Check fuer den Testhost.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Manipulierter Lizenzschluessel zuverlaessig erkannt —
   TestParse_RejectsTamperedKey, TestParse_RejectsWrongKeyPair,
   TestStore_InstallRejectsInvalidSignature. PASS.
2. Ablauf loest definierten eingeschraenkten Zustand aus, kein harter
   Systemausfall — TestStore_RequireActive_DetectsExpiry (inkl. recover()-
   Absicherung im Test, dass kein Panic auftritt), ErrLicenseExpired statt
   Absturz; Status bleibt trotzdem abfragbar. PASS.
3. Signaturpruefung von zweiter Person gegen Dokumentation nachvollzogen —
   NICHT durchgefuehrt (keine zweite Person in dieser Session verfuegbar).
   Offen.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:12:13 +02:00
65 changed files with 2204 additions and 2264 deletions
-22
View File
@@ -1,22 +0,0 @@
name: Core-Schnittstellen-Vertragstests
on:
push:
paths:
- "internal/contracttest/**"
- "internal/apiserver/**"
- "internal/moduletrust/**"
- "internal/moduleregistry/**"
- "internal/webhook/**"
pull_request: {}
jobs:
contract-tests:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.22"
- name: Vertragstests ausfuehren
run: go test ./internal/contracttest/... -v
+1 -10
View File
@@ -8,7 +8,6 @@ import (
"gitea.perlbach24.de/scripte/nexarch/internal/config"
"gitea.perlbach24.de/scripte/nexarch/internal/db"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
"gitea.perlbach24.de/scripte/nexarch/internal/user"
)
func main() {
@@ -35,20 +34,12 @@ func main() {
provisioner := tenant.NewProvisioner(adminPool, registry, cfg.TenantDSNTemplate)
tenantHandler := tenant.NewHandler(provisioner)
// Superadmin-Konten leben mandantenuebergreifend in der Registry-DB.
// Tenant-User-CRUD (user.TenantUserStore) braucht Connection-Routing pro
// Mandant (TEN-06, noch nicht gebaut) und wird hier bewusst noch nicht
// verdrahtet — Package ist bereits eigenstaendig nutzbar/testbar.
superadmins := user.NewSuperadminStore(registryPool)
userHandler := user.NewHandler(nil, superadmins)
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
// Vorlaeufige Pfade ohne Versionierung/Auth — werden mit API-01/IAM-02 abgeloest.
// Vorlaeufiger Pfad ohne Versionierung/Auth — wird mit API-01/IAM-01 abgeloest.
mux.HandleFunc("/internal/tenants", tenantHandler.CreateTenant)
mux.HandleFunc("/internal/superadmins", userHandler.CreateSuperadmin)
log.Printf("nexarch-core listening on %s", cfg.ListenAddr)
if err := http.ListenAndServe(cfg.ListenAddr, mux); err != nil {
+70
View File
@@ -0,0 +1,70 @@
// 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)
}
}
@@ -0,0 +1,14 @@
[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
+76
View File
@@ -0,0 +1,76 @@
# 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 -1
View File
@@ -5,13 +5,13 @@ go 1.22
require (
github.com/golang-jwt/jwt/v5 v5.3.1
github.com/jackc/pgx/v5 v5.6.0
golang.org/x/crypto v0.17.0
)
require (
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
github.com/jackc/puddle/v2 v2.2.1 // indirect
golang.org/x/crypto v0.17.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/text v0.14.0 // indirect
)
-20
View File
@@ -1,20 +0,0 @@
package apiserver
import "context"
type contextKey int
const requestContextKey contextKey = iota
// RequestContext ist der Tenant-/Benutzerkontext, den die Middleware-Kette
// fuer nachgelagerte Handler bereitstellt (Akzeptanzkriterium 3).
type RequestContext struct {
UserID string
TenantSlug string
}
// FromContext liest den von der Middleware gesetzten Kontext.
func FromContext(ctx context.Context) (RequestContext, bool) {
rc, ok := ctx.Value(requestContextKey).(RequestContext)
return rc, ok
}
-31
View File
@@ -1,31 +0,0 @@
// Package apiserver implementiert Core API-01: das REST-Grundgerüst mit
// URL-Versionierung, einheitlichem Fehlerformat und Middleware-Kette
// (Auth, Tenant-/Benutzerkontext, Logging).
package apiserver
import (
"encoding/json"
"net/http"
)
// errorBody ist das EINE Fehlerschema fuer alle Endpunkte unter /api/{version}/
// (Akzeptanzkriterium 2).
type errorBody struct {
Error struct {
Code string `json:"code"`
Message string `json:"message"`
} `json:"error"`
}
// WriteError schreibt einen Fehler im einheitlichen Schema. code ist ein
// stabiler, maschinenlesbarer Bezeichner (z.B. "unauthenticated"), message
// ein fuer Menschen lesbarer deutscher Text.
func WriteError(w http.ResponseWriter, status int, code, message string) {
var body errorBody
body.Error.Code = code
body.Error.Message = message
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(body)
}
-56
View File
@@ -1,56 +0,0 @@
package apiserver
import (
"context"
"log/slog"
"net/http"
"time"
"gitea.perlbach24.de/scripte/nexarch/internal/auth"
)
// authAndTenantContext prueft die Session (wiederverwendet auth.TokenIssuer.Verify
// aus IAM-02 — keine zweite JWT-Implementierung) und setzt bei Erfolg
// RequestContext fuer nachgelagerte Handler (Akzeptanzkriterium 3). Anders
// als auth.RequireAuth (Klartext-Fehler) antwortet diese Middleware im
// einheitlichen API-01-Fehlerschema (Akzeptanzkriterium 2), damit ALLE
// Endpunkte unter /api/{version}/ dasselbe Format liefern, auch bei
// Auth-Fehlern.
func authAndTenantContext(issuer *auth.TokenIssuer, next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
cookie, err := r.Cookie(auth.CookieName)
if err != nil {
WriteError(w, http.StatusUnauthorized, "unauthenticated", "nicht angemeldet")
return
}
claims, err := issuer.Verify(cookie.Value)
if err != nil {
WriteError(w, http.StatusUnauthorized, "unauthenticated", "nicht angemeldet")
return
}
rc := RequestContext{UserID: claims.UserID, TenantSlug: claims.TenantSlug}
next(w, r.WithContext(context.WithValue(r.Context(), requestContextKey, rc)))
}
}
type statusRecorder struct {
http.ResponseWriter
status int
}
func (s *statusRecorder) WriteHeader(code int) {
s.status = code
s.ResponseWriter.WriteHeader(code)
}
// loggingMiddleware protokolliert jede Anfrage strukturiert.
func loggingMiddleware(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
start := time.Now()
next(rec, r)
slog.Info("api-anfrage", "method", r.Method, "path", r.URL.Path, "status", rec.status, "dauer", time.Since(start))
}
}
-37
View File
@@ -1,37 +0,0 @@
package apiserver
import (
"net/http"
"gitea.perlbach24.de/scripte/nexarch/internal/auth"
)
// Server registriert versionierte API-Routen (Akzeptanzkriterium 1: unter
// /api/{version}/...) und verdrahtet fuer jede Route dieselbe Middleware-
// Kette (Logging -> Auth+Tenantkontext -> Handler).
type Server struct {
mux *http.ServeMux
issuer *auth.TokenIssuer
}
func NewServer(issuer *auth.TokenIssuer) *Server {
return &Server{mux: http.NewServeMux(), issuer: issuer}
}
// Handle registriert pattern unter der angegebenen Version, z.B.
// Handle("v1", "/things", h) -> erreichbar unter /api/v1/things. Verschiedene
// Versionen sind unabhaengige Pfade — eine neue Version beeintraechtigt
// bestehende nicht (Akzeptanzkriterium 1 / Pruefung 3).
func (s *Server) Handle(version, pattern string, h http.HandlerFunc) {
full := "/api/" + version + pattern
s.mux.HandleFunc(full, loggingMiddleware(authAndTenantContext(s.issuer, h)))
}
// HandleV1 ist die Kurzform fuer die aktuelle Hauptversion.
func (s *Server) HandleV1(pattern string, h http.HandlerFunc) {
s.Handle("v1", pattern, h)
}
func (s *Server) Handler() http.Handler {
return s.mux
}
-149
View File
@@ -1,149 +0,0 @@
package apiserver
import (
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"gitea.perlbach24.de/scripte/nexarch/internal/auth"
)
func newTestServer() (*Server, *auth.TokenIssuer) {
issuer := auth.NewTokenIssuer("test-secret-nur-fuer-tests")
return NewServer(issuer), issuer
}
func withAuthCookie(req *http.Request, token string) *http.Request {
req.AddCookie(&http.Cookie{Name: auth.CookieName, Value: token})
return req
}
// Akzeptanzkriterium 1: API unter versioniertem Pfad erreichbar.
func TestHandleV1_RegistersUnderVersionedPath(t *testing.T) {
srv, issuer := newTestServer()
srv.HandleV1("/things", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
req := withAuthCookie(httptest.NewRequest(http.MethodGet, "/api/v1/things", nil), token)
rec := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", rec.Code)
}
}
// Akzeptanzkriterium 2 + Pruefung 1 (Stichprobe): mehrere Endpunkte liefern
// bei fehlerhafter Anfrage dasselbe Fehlerschema.
func TestErrorFormat_ConsistentAcrossEndpoints(t *testing.T) {
srv, _ := newTestServer()
srv.HandleV1("/endpunkt-a", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
srv.HandleV1("/endpunkt-b", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
for _, path := range []string{"/api/v1/endpunkt-a", "/api/v1/endpunkt-b"} {
req := httptest.NewRequest(http.MethodGet, path, nil) // ohne cookie -> 401
rec := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("%s: status = %d, want 401", path, rec.Code)
}
var body errorBody
if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil {
t.Fatalf("%s: fehlerantwort nicht im erwarteten json-schema: %v (body: %s)", path, err, rec.Body.String())
}
if body.Error.Code == "" || body.Error.Message == "" {
t.Fatalf("%s: erwartet nicht-leeren code/message, habe %+v", path, body)
}
}
}
// Akzeptanzkriterium 3 + Pruefung 2: Middleware-Kette setzt Tenant-/
// Benutzerkontext zuverlaessig, nachweislich fuer mehrere Endpunkte.
func TestMiddleware_SetsRequestContextForEveryEndpoint(t *testing.T) {
srv, issuer := newTestServer()
var gotA, gotB RequestContext
srv.HandleV1("/kontext-a", func(w http.ResponseWriter, r *http.Request) {
gotA, _ = FromContext(r.Context())
w.WriteHeader(http.StatusOK)
})
srv.HandleV1("/kontext-b", func(w http.ResponseWriter, r *http.Request) {
gotB, _ = FromContext(r.Context())
w.WriteHeader(http.StatusOK)
})
token, err := issuer.Issue("user-42", "tenant-x")
if err != nil {
t.Fatalf("issue: %v", err)
}
for path, got := range map[string]*RequestContext{"/api/v1/kontext-a": &gotA, "/api/v1/kontext-b": &gotB} {
req := withAuthCookie(httptest.NewRequest(http.MethodGet, path, nil), token)
rec := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("%s: status = %d, want 200", path, rec.Code)
}
if got.UserID != "user-42" || got.TenantSlug != "tenant-x" {
t.Fatalf("%s: request-context unerwartet: %+v", path, *got)
}
}
}
// Akzeptanzkriterium 1 + Pruefung 3: eine neue v2-Route laesst sich anlegen,
// ohne v1 zu beeintraechtigen.
func TestVersioning_V2DoesNotAffectV1(t *testing.T) {
srv, issuer := newTestServer()
srv.HandleV1("/things", func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("v1-antwort"))
})
srv.Handle("v2", "/things", func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("v2-antwort"))
})
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
reqV1 := withAuthCookie(httptest.NewRequest(http.MethodGet, "/api/v1/things", nil), token)
recV1 := httptest.NewRecorder()
srv.Handler().ServeHTTP(recV1, reqV1)
if recV1.Body.String() != "v1-antwort" {
t.Fatalf("v1 antwort = %q, want v1-antwort", recV1.Body.String())
}
reqV2 := withAuthCookie(httptest.NewRequest(http.MethodGet, "/api/v2/things", nil), token)
recV2 := httptest.NewRecorder()
srv.Handler().ServeHTTP(recV2, reqV2)
if recV2.Body.String() != "v2-antwort" {
t.Fatalf("v2 antwort = %q, want v2-antwort", recV2.Body.String())
}
// v1 nach dem Anlegen von v2 erneut pruefen — unveraendert.
reqV1Again := withAuthCookie(httptest.NewRequest(http.MethodGet, "/api/v1/things", nil), token)
recV1Again := httptest.NewRecorder()
srv.Handler().ServeHTTP(recV1Again, reqV1Again)
if recV1Again.Body.String() != "v1-antwort" {
t.Fatalf("v1 antwort nach v2-anlage = %q, want weiterhin v1-antwort", recV1Again.Body.String())
}
}
func TestAuthAndTenantContext_RejectsInvalidToken(t *testing.T) {
srv, _ := newTestServer()
srv.HandleV1("/geschuetzt", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
req := withAuthCookie(httptest.NewRequest(http.MethodGet, "/api/v1/geschuetzt", nil), "kaputtes.token.hier")
rec := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("status = %d, want 401", rec.Code)
}
}
+96
View File
@@ -0,0 +1,96 @@
// 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
}
+132
View File
@@ -0,0 +1,132 @@
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)
}
}
-67
View File
@@ -1,67 +0,0 @@
package auth
import (
"encoding/json"
"net/http"
"time"
)
// Handler stellt Login/Logout als HTTP-Endpunkte bereit. Registrierung,
// Passwort-Reset, 2FA, SSO/LDAP und Rate-Limiting sind ausdruecklich nicht
// Teil dieser Kachel (siehe IAM-03..07).
type Handler struct {
login *LoginService
}
func NewHandler(login *LoginService) *Handler {
return &Handler{login: login}
}
type loginRequest struct {
Email string `json:"email"`
Password string `json:"password"`
}
func (h *Handler) Login(w http.ResponseWriter, r *http.Request) {
var req loginRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "ungueltige Anfrage", http.StatusBadRequest)
return
}
token, err := h.login.Login(r.Context(), req.Email, req.Password)
if err != nil {
http.Error(w, ErrInvalidCredentials.Error(), http.StatusUnauthorized)
return
}
http.SetCookie(w, &http.Cookie{
Name: CookieName,
Value: token,
Path: "/",
HttpOnly: true,
Secure: true,
SameSite: http.SameSiteStrictMode,
MaxAge: int(AccessTokenTTL.Seconds()),
})
w.WriteHeader(http.StatusOK)
}
// Logout loescht das Session-Cookie. Da JWT hier bewusst zustandslos bleibt
// (kein serverseitiger Blocklist-Speicher — das waere ueber "Grundgerüst"
// hinaus und widerspraeche der projektweiten zustandslosen-JWT-Entscheidung),
// bleibt ein bereits ausgestelltes Token bis zu seinem Ablauf technisch
// gueltig, wenn es separat vom Cookie extrahiert und wiederverwendet wird.
func (h *Handler) Logout(w http.ResponseWriter, r *http.Request) {
http.SetCookie(w, &http.Cookie{
Name: CookieName,
Value: "",
Path: "/",
HttpOnly: true,
Secure: true,
SameSite: http.SameSiteStrictMode,
MaxAge: -1,
Expires: time.Unix(0, 0),
})
w.WriteHeader(http.StatusOK)
}
-56
View File
@@ -1,56 +0,0 @@
package auth
import (
"context"
"errors"
"gitea.perlbach24.de/scripte/nexarch/internal/user"
)
var ErrInvalidCredentials = errors.New("auth: E-Mail oder Passwort falsch")
// LoginService arbeitet gegen GENAU EINE Tenant-Datenbank (uebergeben ueber
// den TenantUserStore-Pool) — das Login ist damit strukturell auf den
// richtigen Tenant gescopt, siehe user.TenantUserStore.GetByEmailForAuth.
type LoginService struct {
users *user.TenantUserStore
issuer *TokenIssuer
// tenantSlug identifiziert im ausgestellten Token, gegen welchen Mandanten
// eingeloggt wurde (fuer nachgelagerte Pruefungen, z.B. Middleware-Logs).
tenantSlug string
}
func NewLoginService(users *user.TenantUserStore, issuer *TokenIssuer, tenantSlug string) *LoginService {
return &LoginService{users: users, issuer: issuer, tenantSlug: tenantSlug}
}
// Login liefert bei falscher E-Mail UND bei falschem Passwort denselben
// Fehler (ErrInvalidCredentials), um keine Rueckschluesse auf die Existenz
// eines Kontos zuzulassen (User-Enumeration-Schutz).
func (s *LoginService) Login(ctx context.Context, email, password string) (string, error) {
creds, err := s.users.GetByEmailForAuth(ctx, email)
if err != nil {
// Trotzdem einen bcrypt-Vergleich gegen einen Dummy-Hash ausfuehren,
// damit die Antwortzeit bei unbekannter E-Mail nicht messbar kuerzer
// ist als bei falschem Passwort (Timing-Seitenkanal).
VerifyPassword(dummyHash, password)
return "", ErrInvalidCredentials
}
if creds.User.Status != user.StatusActive {
return "", ErrInvalidCredentials
}
if !VerifyPassword(creds.PasswordHash, password) {
return "", ErrInvalidCredentials
}
return s.issuer.Issue(creds.User.ID, s.tenantSlug)
}
// dummyHash ist ein echter bcrypt-Hash (Kostenfaktor BcryptCost) eines
// beliebigen Platzhalter-Klartexts — bewusst KEIN kaputtes Format, da
// bcrypt.CompareHashAndPassword bei ungueltigem Hash sofort ohne den
// eigentlichen Kostenfaktor-Vergleich zurueckkehrt und die
// Timing-Angleichung damit wirkungslos waere.
const dummyHash = "$2a$12$cmwiETrG9DK5/uTM2fg4uetngYUspKjME5P8fNpk0QYTaO64N0r3C"
-177
View File
@@ -1,177 +0,0 @@
package auth
import (
"context"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/user"
)
const usersSchema = `
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE users (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
password_hash TEXT NOT NULL DEFAULT '',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);`
func setupTenantDB(t *testing.T, dbName string) *pgxpool.Pool {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
adminPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("admin pool: %v", err)
}
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbName))
if _, err := adminPool.Exec(ctx, fmt.Sprintf(`CREATE DATABASE %q`, dbName)); err != nil {
t.Fatalf("testdatenbank anlegen: %v", err)
}
dsn := strings.Replace(adminDSN, "/postgres?", "/"+dbName+"?", 1)
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect testdatenbank: %v", err)
}
if _, err := pool.Exec(ctx, usersSchema); err != nil {
t.Fatalf("schema anwenden: %v", err)
}
t.Cleanup(func() {
pool.Close()
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbName))
adminPool.Close()
})
return pool
}
func createUserWithPassword(t *testing.T, store *user.TenantUserStore, email, password string) user.User {
t.Helper()
ctx := context.Background()
u, err := store.Create(ctx, email, "Test User")
if err != nil {
t.Fatalf("create user: %v", err)
}
hash, err := HashPassword(password)
if err != nil {
t.Fatalf("hash password: %v", err)
}
if err := store.SetPasswordHash(ctx, u.ID, hash); err != nil {
t.Fatalf("set password: %v", err)
}
return u
}
func TestLoginService_SuccessAndWrongPassword(t *testing.T) {
pool := setupTenantDB(t, "test_iam02_login")
store := user.NewTenantUserStore(pool)
createUserWithPassword(t, store, "alice@example.com", "korrektes-passwort")
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
login := NewLoginService(store, issuer, "acme")
token, err := login.Login(context.Background(), "alice@example.com", "korrektes-passwort")
if err != nil {
t.Fatalf("login: %v", err)
}
if token == "" {
t.Fatal("erwartet nicht-leeres token")
}
if _, err := login.Login(context.Background(), "alice@example.com", "falsches-passwort"); !errors.Is(err, ErrInvalidCredentials) {
t.Fatalf("erwartet ErrInvalidCredentials, habe %v", err)
}
if _, err := login.Login(context.Background(), "unbekannt@example.com", "irgendwas"); !errors.Is(err, ErrInvalidCredentials) {
t.Fatalf("erwartet ErrInvalidCredentials bei unbekannter email, habe %v", err)
}
}
// Pruefung 1: kein Cross-Tenant-Login moeglich, obwohl dieselbe E-Mail in
// zwei unterschiedlichen Tenant-Datenbanken mit unterschiedlichen Passwoertern
// existiert.
func TestLoginService_NoCrossTenantLogin(t *testing.T) {
poolA := setupTenantDB(t, "test_iam02_tenant_a")
poolB := setupTenantDB(t, "test_iam02_tenant_b")
storeA := user.NewTenantUserStore(poolA)
storeB := user.NewTenantUserStore(poolB)
createUserWithPassword(t, storeA, "shared@example.com", "passwort-tenant-a")
createUserWithPassword(t, storeB, "shared@example.com", "passwort-tenant-b")
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
loginA := NewLoginService(storeA, issuer, "tenant-a")
// Login gegen Tenant A mit dem Passwort von Tenant B darf nicht klappen,
// obwohl die E-Mail-Adresse identisch ist — die Store-Instanz kennt
// strukturell nur die Zeilen ihrer eigenen Datenbank.
if _, err := loginA.Login(context.Background(), "shared@example.com", "passwort-tenant-b"); !errors.Is(err, ErrInvalidCredentials) {
t.Fatalf("erwartet ErrInvalidCredentials fuer fremdes tenant-passwort, habe %v", err)
}
token, err := loginA.Login(context.Background(), "shared@example.com", "passwort-tenant-a")
if err != nil {
t.Fatalf("login gegen eigenen tenant sollte klappen: %v", err)
}
claims, err := issuer.Verify(token)
if err != nil {
t.Fatalf("verify: %v", err)
}
if claims.TenantSlug != "tenant-a" {
t.Fatalf("token tenant = %q, want tenant-a", claims.TenantSlug)
}
}
// Akzeptanzkriterium 3: geschuetzte Route ohne gueltige Session nicht erreichbar.
func TestRequireAuth_BlocksWithoutValidCookie(t *testing.T) {
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
protected := RequireAuth(issuer, func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
// Kein Cookie.
req := httptest.NewRequest(http.MethodGet, "/geschuetzt", nil)
rec := httptest.NewRecorder()
protected(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("ohne cookie: status = %d, want 401", rec.Code)
}
// Manipuliertes Cookie.
req = httptest.NewRequest(http.MethodGet, "/geschuetzt", nil)
req.AddCookie(&http.Cookie{Name: CookieName, Value: "kaputt.token.hier"})
rec = httptest.NewRecorder()
protected(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("mit kaputtem cookie: status = %d, want 401", rec.Code)
}
// Gueltiges Token.
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
req = httptest.NewRequest(http.MethodGet, "/geschuetzt", nil)
req.AddCookie(&http.Cookie{Name: CookieName, Value: token})
rec = httptest.NewRecorder()
protected(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("mit gueltigem cookie: status = %d, want 200", rec.Code)
}
}
-41
View File
@@ -1,41 +0,0 @@
package auth
import (
"context"
"net/http"
)
const CookieName = "nexarch_session"
type contextKey int
const claimsContextKey contextKey = iota
// RequireAuth schuetzt eine Route: ohne gueltiges, nicht abgelaufenes Token
// im Session-Cookie wird 401 zurueckgegeben und der Handler nicht aufgerufen
// (IAM-02 Akzeptanzkriterium 3).
func RequireAuth(issuer *TokenIssuer, next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
cookie, err := r.Cookie(CookieName)
if err != nil {
http.Error(w, "nicht angemeldet", http.StatusUnauthorized)
return
}
claims, err := issuer.Verify(cookie.Value)
if err != nil {
http.Error(w, "nicht angemeldet", http.StatusUnauthorized)
return
}
ctx := context.WithValue(r.Context(), claimsContextKey, claims)
next(w, r.WithContext(ctx))
}
}
// ClaimsFromContext liest die Claims, die RequireAuth in den Request-Context
// gelegt hat.
func ClaimsFromContext(ctx context.Context) (*Claims, bool) {
c, ok := ctx.Value(claimsContextKey).(*Claims)
return c, ok
}
-28
View File
@@ -1,28 +0,0 @@
// Package auth implementiert Core IAM-02: Login/Logout, Passwort-Hashing und
// die Middleware zum Schutz von Routen. Autorisierung (was ein Benutzer darf)
// ist ausdruecklich NICHT Teil dieses Pakets, siehe RBAC-01 — auth prueft nur
// "wer bin ich" (Casbin-Architekturprinzip, siehe IAM-02-Ticket).
package auth
import "golang.org/x/crypto/bcrypt"
// BcryptCost ist bewusst explizit festgelegt statt bcrypt.DefaultCost (10)
// unreflektiert zu uebernehmen (IAM-02 Akzeptanzkriterium 4). Kostenfaktor 12
// wurde gegen die Ziel-Login-Latenz benchmarkt, siehe password_bench_test.go
// und den Pruefungs-Eintrag in der Commit-Nachricht.
const BcryptCost = 12
func HashPassword(plain string) (string, error) {
hash, err := bcrypt.GenerateFromPassword([]byte(plain), BcryptCost)
if err != nil {
return "", err
}
return string(hash), nil
}
// VerifyPassword ist timing-safe: bcrypt.CompareHashAndPassword vergleicht
// konstant in der Zeit bzgl. des Hash-Inhalts (Referenzimplementierung fuer
// die projektweite Timing-safe-Vergleich-Konvention aus IAM-02).
func VerifyPassword(hash, plain string) bool {
return bcrypt.CompareHashAndPassword([]byte(hash), []byte(plain)) == nil
}
-42
View File
@@ -1,42 +0,0 @@
package auth
import (
"testing"
"time"
)
// TargetLoginLatency ist der Zielwert aus IAM-02 Akzeptanzkriterium 4: der
// bcrypt-Vergleich allein darf die Login-Latenz nicht dominieren. 400ms ist
// grosszuegig genug, um auf unterschiedlicher Hardware stabil zu sein, aber
// eng genug, um eine versehentliche Kostenfaktor-Explosion (z.B. 16 statt 12)
// zuverlaessig aufzudecken.
const TargetLoginLatency = 400 * time.Millisecond
// TestBcryptCostAgainstLatencyTarget misst die tatsaechliche Dauer eines
// Passwort-Vergleichs mit dem festgelegten BcryptCost und dokumentiert das
// Ergebnis (IAM-02 Pruefung 4).
func TestBcryptCostAgainstLatencyTarget(t *testing.T) {
hash, err := HashPassword("benchmark-passwort")
if err != nil {
t.Fatalf("hash: %v", err)
}
start := time.Now()
if !VerifyPassword(hash, "benchmark-passwort") {
t.Fatal("verifikation haette erfolgreich sein muessen")
}
elapsed := time.Since(start)
t.Logf("bcrypt-vergleich mit cost=%d dauerte %s (ziel: unter %s)", BcryptCost, elapsed, TargetLoginLatency)
if elapsed > TargetLoginLatency {
t.Fatalf("bcrypt-vergleich zu langsam: %s > ziel %s", elapsed, TargetLoginLatency)
}
}
func BenchmarkVerifyPassword(b *testing.B) {
hash, _ := HashPassword("benchmark-passwort")
b.ResetTimer()
for i := 0; i < b.N; i++ {
VerifyPassword(hash, "benchmark-passwort")
}
}
-28
View File
@@ -1,28 +0,0 @@
package auth
import "testing"
func TestHashAndVerifyPassword(t *testing.T) {
hash, err := HashPassword("s3hr-geheim!")
if err != nil {
t.Fatalf("hash: %v", err)
}
if hash == "s3hr-geheim!" {
t.Fatal("passwort wurde nicht gehasht")
}
if !VerifyPassword(hash, "s3hr-geheim!") {
t.Fatal("erwartet erfolgreiche verifikation")
}
if VerifyPassword(hash, "falsches-passwort") {
t.Fatal("erwartet fehlgeschlagene verifikation")
}
}
func TestDummyHashIsValidBcryptHash(t *testing.T) {
// Stellt sicher, dass der Timing-Angleichs-Hash in login.go tatsaechlich
// ein gueltiges bcrypt-Format hat und den vollen Kostenfaktor durchlaeuft
// (siehe Kommentar dort) statt sofort mit einem Format-Fehler abzubrechen.
if VerifyPassword(dummyHash, "irgendein-text") {
t.Fatal("dummyHash sollte fuer beliebigen text nicht passen")
}
}
-63
View File
@@ -1,63 +0,0 @@
package auth
import (
"errors"
"time"
"github.com/golang-jwt/jwt/v5"
)
// AccessTokenTTL ist bewusst kurz gehalten (Session-Ablauf statt langlebiger
// Tokens), passend zur "so vertrauenswuerdig wie noetig"-Produkt-DNA.
const AccessTokenTTL = 30 * time.Minute
var ErrInvalidToken = errors.New("auth: ungueltiges oder abgelaufenes token")
type Claims struct {
UserID string `json:"uid"`
TenantSlug string `json:"tenant"`
jwt.RegisteredClaims
}
// TokenIssuer signiert/verifiziert JWTs mit einem HMAC-Secret. Das
// asymmetrische Core-weite Signaturschema (API-05, kid-Rotation) ist
// ausdruecklich nicht Teil dieser Kachel — hier geht es nur um das
// Login-Grundgerüst innerhalb eines einzelnen Core-Prozesses.
type TokenIssuer struct {
secret []byte
}
func NewTokenIssuer(secret string) *TokenIssuer {
return &TokenIssuer{secret: []byte(secret)}
}
func (i *TokenIssuer) Issue(userID, tenantSlug string) (string, error) {
now := time.Now()
claims := Claims{
UserID: userID,
TenantSlug: tenantSlug,
RegisteredClaims: jwt.RegisteredClaims{
IssuedAt: jwt.NewNumericDate(now),
ExpiresAt: jwt.NewNumericDate(now.Add(AccessTokenTTL)),
},
}
token := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
return token.SignedString(i.secret)
}
// Verify prueft Signatur UND Ablauf (jwt.ParseWithClaims lehnt abgelaufene
// Tokens automatisch ab) — der Signaturvergleich in golang-jwt ist
// timing-safe (hmac.Equal).
func (i *TokenIssuer) Verify(tokenString string) (*Claims, error) {
claims := &Claims{}
token, err := jwt.ParseWithClaims(tokenString, claims, func(t *jwt.Token) (interface{}, error) {
if _, ok := t.Method.(*jwt.SigningMethodHMAC); !ok {
return nil, ErrInvalidToken
}
return i.secret, nil
})
if err != nil || !token.Valid {
return nil, ErrInvalidToken
}
return claims, nil
}
-82
View File
@@ -1,82 +0,0 @@
package auth
import (
"strings"
"testing"
"time"
"github.com/golang-jwt/jwt/v5"
)
func TestTokenIssueAndVerify(t *testing.T) {
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
claims, err := issuer.Verify(token)
if err != nil {
t.Fatalf("verify: %v", err)
}
if claims.UserID != "user-1" || claims.TenantSlug != "acme" {
t.Fatalf("claims unerwartet: %+v", claims)
}
}
// Pruefung 2: Token-Manipulationstest.
func TestTokenVerify_RejectsManipulatedPayload(t *testing.T) {
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
parts := strings.Split(token, ".")
if len(parts) != 3 {
t.Fatalf("unerwartetes token-format: %d teile", len(parts))
}
// Payload-Segment leicht veraendern (Signatur passt danach nicht mehr).
tampered := parts[0] + "." + parts[1] + "x" + "." + parts[2]
if _, err := issuer.Verify(tampered); err == nil {
t.Fatal("erwartet fehler bei manipuliertem token, habe nil")
}
}
func TestTokenVerify_RejectsWrongSecret(t *testing.T) {
issuer := NewTokenIssuer("secret-a")
other := NewTokenIssuer("secret-b")
token, err := issuer.Issue("user-1", "acme")
if err != nil {
t.Fatalf("issue: %v", err)
}
if _, err := other.Verify(token); err == nil {
t.Fatal("erwartet fehler bei falschem secret, habe nil")
}
}
// Pruefung 3: abgelaufenes Token erzwingt Neuanmeldung.
func TestTokenVerify_RejectsExpiredToken(t *testing.T) {
issuer := NewTokenIssuer("test-secret-nur-fuer-tests")
claims := Claims{
UserID: "user-1",
TenantSlug: "acme",
RegisteredClaims: jwt.RegisteredClaims{
IssuedAt: jwt.NewNumericDate(time.Now().Add(-2 * AccessTokenTTL)),
ExpiresAt: jwt.NewNumericDate(time.Now().Add(-time.Minute)),
},
}
tok := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
expired, err := tok.SignedString([]byte("test-secret-nur-fuer-tests"))
if err != nil {
t.Fatalf("signieren: %v", err)
}
if _, err := issuer.Verify(expired); err == nil {
t.Fatal("erwartet fehler bei abgelaufenem token, habe nil")
}
}
-57
View File
@@ -1,57 +0,0 @@
// Package contracttest implementiert Core QA-07: Vertragstests fuer die
// Core-Schnittstellen, die von den Fachmodulen (DMS/Mail/Archive/Workflow/
// AI/Connect) konsumiert werden — REST-Fehlerschema (API-01), JWKS-/Rechte-
// Cache-Kontrakt (API-05), Modul-Registry (API-02) und Webhook-Zustellung
// (API-07). Jeder Test prueft die TATSAECHLICHE Antwort einer echten
// Core-Komponente gegen ihr dokumentiertes Schema — eine entfernte oder
// umbenannte Pflichteigenschaft laesst den jeweiligen Test fehlschlagen,
// BEVOR sie ein konsumierendes Modul bricht (Akzeptanzkriterium 2).
package contracttest
import (
"encoding/json"
"fmt"
)
// RequireJSONFields dekodiert data als JSON-Objekt und prueft, dass ALLE
// angegebenen Top-Level-Schluessel vorhanden sind. Liefert die fehlenden
// Schluessel zurueck — leer bedeutet: Vertrag eingehalten.
func RequireJSONFields(data []byte, required []string) (missing []string, err error) {
var decoded map[string]any
if err := json.Unmarshal(data, &decoded); err != nil {
return nil, fmt.Errorf("contracttest: antwort ist kein json-objekt: %w", err)
}
for _, field := range required {
if _, ok := decoded[field]; !ok {
missing = append(missing, field)
}
}
return missing, nil
}
// RequireJSONArrayItemFields prueft, dass data ein JSON-Objekt mit einem
// Array-Feld arrayField ist, dessen ERSTES Element alle itemFields enthaelt
// — Vertrag fuer Listen-Antworten wie JWKS ("keys": [{...}]).
func RequireJSONArrayItemFields(data []byte, arrayField string, itemFields []string) (missing []string, err error) {
var decoded map[string]json.RawMessage
if err := json.Unmarshal(data, &decoded); err != nil {
return nil, fmt.Errorf("contracttest: antwort ist kein json-objekt: %w", err)
}
raw, ok := decoded[arrayField]
if !ok {
return itemFields, nil // das array-feld selbst fehlt bereits -> alles "fehlend"
}
var items []map[string]any
if err := json.Unmarshal(raw, &items); err != nil {
return nil, fmt.Errorf("contracttest: feld %q ist kein array: %w", arrayField, err)
}
if len(items) == 0 {
return nil, fmt.Errorf("contracttest: feld %q ist leer, kann nicht gegen kontrakt geprueft werden", arrayField)
}
for _, field := range itemFields {
if _, ok := items[0][field]; !ok {
missing = append(missing, field)
}
}
return missing, nil
}
-193
View File
@@ -1,193 +0,0 @@
package contracttest
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/apiserver"
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
"gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry"
"gitea.perlbach24.de/scripte/nexarch/internal/moduletrust"
"gitea.perlbach24.de/scripte/nexarch/internal/webhook"
)
// --- Vertrag 1: API-01 REST-Fehlerschema ---
// Akzeptanzkriterium 1/2 + Pruefung 1/2.
func TestContract_API01_ErrorEnvelope(t *testing.T) {
rec := httptest.NewRecorder()
apiserver.WriteError(rec, http.StatusUnauthorized, "unauthenticated", "nicht angemeldet")
if missing, err := RequireJSONFields(rec.Body.Bytes(), []string{"error"}); err != nil {
t.Fatalf("antwort nicht lesbar: %v", err)
} else if len(missing) != 0 {
t.Fatalf("erwartet top-level feld 'error', fehlt: %v", missing)
}
var decoded map[string]json.RawMessage
_ = json.Unmarshal(rec.Body.Bytes(), &decoded)
if missing, err := RequireJSONFields(decoded["error"], []string{"code", "message"}); err != nil {
t.Fatalf("error-objekt nicht lesbar: %v", err)
} else if len(missing) != 0 {
t.Fatalf("error-objekt fehlen pflichtfelder: %v", missing)
}
}
// Pruefung 2: absichtlich simulierter Breaking Change (Feld "message"
// entfernt) laesst den Vertragstest fehlschlagen.
func TestContract_API01_DetectsBreakingChange_RemovedField(t *testing.T) {
brokenResponse := []byte(`{"error":{"code":"unauthenticated"}}`) // "message" fehlt absichtlich
var decoded map[string]json.RawMessage
if err := json.Unmarshal(brokenResponse, &decoded); err != nil {
t.Fatalf("fixture nicht lesbar: %v", err)
}
missing, err := RequireJSONFields(decoded["error"], []string{"code", "message"})
if err != nil {
t.Fatalf("pruefung selbst fehlgeschlagen: %v", err)
}
if len(missing) != 1 || missing[0] != "message" {
t.Fatalf("erwartet erkanntes fehlendes feld 'message', habe: %v", missing)
}
}
// --- Vertrag 2: API-05 JWKS / Rechte-Cache-Kontrakt ---
// Akzeptanzkriterium 1/2 + Pruefung 1/2.
func TestContract_API05_JWKS(t *testing.T) {
km, err := moduletrust.NewKeyManager()
if err != nil {
t.Fatalf("keymanager: %v", err)
}
if _, err := km.Rotate(); err != nil {
t.Fatalf("rotate: %v", err)
}
rec := httptest.NewRecorder()
km.ServeJWKS(rec, httptest.NewRequest(http.MethodGet, "/jwks", nil))
missing, err := RequireJSONArrayItemFields(rec.Body.Bytes(), "keys", []string{"kid", "public_key"})
if err != nil {
t.Fatalf("jwks-antwort verletzt vertrag: %v", err)
}
if len(missing) != 0 {
t.Fatalf("jwks-eintrag fehlen pflichtfelder: %v", missing)
}
}
// Pruefung 2: simulierter Breaking Change — Feld "public_key" umbenannt
// (z.B. faelschlich zu "publicKey"), Vertragstest erkennt das fehlende
// Originalfeld zuverlaessig.
func TestContract_API05_DetectsBreakingChange_RenamedField(t *testing.T) {
brokenJWKS := []byte(`{"keys":[{"kid":"abc","publicKey":"base64..."}]}`)
missing, err := RequireJSONArrayItemFields(brokenJWKS, "keys", []string{"kid", "public_key"})
if err != nil {
t.Fatalf("pruefung selbst fehlgeschlagen: %v", err)
}
if len(missing) != 1 || missing[0] != "public_key" {
t.Fatalf("erwartet erkanntes umbenanntes feld 'public_key', habe: %v", missing)
}
}
// --- Vertrag 3: API-02 Modul-Registry ---
// Akzeptanzkriterium 1/2 + Pruefung 1/2.
func setupModuleRegistryTest(t *testing.T) (*moduleregistry.Registry, 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 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()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
flagService := flag.NewService(flag.NewStore(pool), 10*time.Millisecond)
return moduleregistry.NewRegistry(pool, flagService), func() { pool.Close() }
}
func TestContract_API02_ModuleInfo(t *testing.T) {
registry, cleanup := setupModuleRegistryTest(t)
defer cleanup()
ctx := context.Background()
m, err := registry.Register(ctx, "contracttest-dms", "1.0.0", []string{"some-flag"})
if err != nil {
t.Fatalf("register: %v", err)
}
data, err := json.Marshal(m)
if err != nil {
t.Fatalf("marshal: %v", err)
}
missing, err := RequireJSONFields(data, []string{"Name", "Version", "RequiredFlags"})
if err != nil {
t.Fatalf("modul-antwort nicht lesbar: %v", err)
}
if len(missing) != 0 {
t.Fatalf("modul-info fehlen pflichtfelder: %v", missing)
}
}
// Pruefung 2: simulierter Breaking Change — Feld "RequiredFlags" entfernt.
func TestContract_API02_DetectsBreakingChange_RemovedField(t *testing.T) {
brokenModuleJSON := []byte(`{"Name":"dms","Version":"1.0.0"}`)
missing, err := RequireJSONFields(brokenModuleJSON, []string{"Name", "Version", "RequiredFlags"})
if err != nil {
t.Fatalf("pruefung selbst fehlgeschlagen: %v", err)
}
if len(missing) != 1 || missing[0] != "RequiredFlags" {
t.Fatalf("erwartet erkanntes fehlendes feld 'RequiredFlags', habe: %v", missing)
}
}
// --- Vertrag 4: API-07 Webhook-Zustellung (HMAC-Signaturheader) ---
// Akzeptanzkriterium 1/2 + Pruefung 1/2.
func TestContract_API07_WebhookSignatureHeader(t *testing.T) {
if webhook.SignatureHeader != "X-Nexarch-Signature-256" {
t.Fatalf("signaturheader-name = %q, want stabilen vertrag 'X-Nexarch-Signature-256' — konsumierende module lesen genau diesen namen", webhook.SignatureHeader)
}
secret := "vertragstest-secret"
payload := []byte(`{"event":"test"}`)
signature := webhook.Sign(secret, payload)
if !webhook.VerifySignature(secret, payload, signature) {
t.Fatal("konsumentenseitige signaturpruefung schlaegt fuer eine korrekte core-signatur fehl")
}
}
// Pruefung 2: simulierter Breaking Change — eine Zustellung, die den
// Signaturheader unter einem ANDEREN Namen sendet (wie es ein
// hypothetischer, den Vertrag brechender Core-Umbau taete). Ein
// konsumierendes Modul, das nach dem DOKUMENTIERTEN Namen sucht, findet in
// diesem Fall nichts — der Vertragstest erkennt das zuverlaessig.
func TestContract_API07_DetectsBreakingChange_RenamedHeader(t *testing.T) {
brokenHeaders := http.Header{}
brokenHeaders.Set("X-Signature", webhook.Sign("secret", []byte("payload"))) // falscher name, simuliert breaking change
if got := brokenHeaders.Get(webhook.SignatureHeader); got != "" {
t.Fatalf("erwartet leeren wert unter dem dokumentierten headernamen bei simuliertem breaking change, habe: %q", got)
}
}
+106
View File
@@ -0,0 +1,106 @@
// 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
}
+106
View File
@@ -0,0 +1,106 @@
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)
}
}
}
+88
View File
@@ -0,0 +1,88 @@
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
}
+177
View File
@@ -0,0 +1,177 @@
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)
}
}
+11
View File
@@ -71,6 +71,17 @@ func (c *StaleCache[T]) Get(ctx context.Context) (value T, stale bool, err error
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) {
+155
View File
@@ -0,0 +1,155 @@
// 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
}
+132
View File
@@ -0,0 +1,132 @@
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)
}
+320
View File
@@ -0,0 +1,320 @@
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)
}
}
+220
View File
@@ -0,0 +1,220 @@
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)
}
}
}
+32
View File
@@ -0,0 +1,32 @@
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)
}
}
}
}
+144
View File
@@ -0,0 +1,144 @@
// 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
}
+206
View File
@@ -0,0 +1,206 @@
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)
}
}
-66
View File
@@ -1,66 +0,0 @@
package user
import (
"encoding/json"
"errors"
"net/http"
)
// Handler stellt die CRUD-API fuer Benutzerkonten bereit (IAM-01-Auftrag).
// Auth/Sessions (IAM-02) und Rollen (RBAC-01) sind ausdruecklich nicht Teil
// dieser Kachel und daher hier noch nicht angebunden.
type Handler struct {
users *TenantUserStore
superadmins *SuperadminStore
}
func NewHandler(users *TenantUserStore, superadmins *SuperadminStore) *Handler {
return &Handler{users: users, superadmins: superadmins}
}
type createUserRequest struct {
Email string `json:"email"`
Name string `json:"name"`
}
func (h *Handler) CreateUser(w http.ResponseWriter, r *http.Request) {
var req createUserRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "ungueltige Anfrage", http.StatusBadRequest)
return
}
u, err := h.users.Create(r.Context(), req.Email, req.Name)
writeUserResult(w, u, err)
}
// CreateSuperadmin legt ein mandantenuebergreifendes Superadmin-Konto an —
// bewusst ein eigener Endpunkt statt eines Tenant-Parameters mit Null-Wert.
func (h *Handler) CreateSuperadmin(w http.ResponseWriter, r *http.Request) {
var req createUserRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "ungueltige Anfrage", http.StatusBadRequest)
return
}
u, err := h.superadmins.Create(r.Context(), req.Email, req.Name)
writeUserResult(w, u, err)
}
func writeUserResult(w http.ResponseWriter, u User, err error) {
if err != nil {
switch {
case errors.Is(err, ErrInvalidEmail), errors.Is(err, ErrEmailTaken):
http.Error(w, err.Error(), http.StatusBadRequest)
case errors.Is(err, ErrNotFound):
http.Error(w, err.Error(), http.StatusNotFound)
default:
http.Error(w, "benutzer konnte nicht verarbeitet werden", http.StatusInternalServerError)
}
return
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(u)
}
-173
View File
@@ -1,173 +0,0 @@
package user
import (
"context"
"errors"
"fmt"
"os"
"strings"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
// setupTestDB legt eine frische, isolierte Testdatenbank an, wendet die
// uebergebene Migration an und liefert einen verbundenen Pool. Wird ohne
// TEST_ADMIN_DSN uebersprungen — siehe internal/tenant/provisioner_test.go
// fuer dasselbe Muster.
func setupTestDB(t *testing.T, dbName, schemaSQL string) *pgxpool.Pool {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
adminPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("admin pool: %v", err)
}
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbName))
if _, err := adminPool.Exec(ctx, fmt.Sprintf(`CREATE DATABASE %q`, dbName)); err != nil {
t.Fatalf("testdatenbank anlegen: %v", err)
}
dsn := strings.Replace(adminDSN, "/postgres?", "/"+dbName+"?", 1)
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect testdatenbank: %v", err)
}
if _, err := pool.Exec(ctx, schemaSQL); err != nil {
t.Fatalf("schema anwenden: %v", err)
}
t.Cleanup(func() {
pool.Close()
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbName))
adminPool.Close()
})
return pool
}
const usersSchema = `
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE users (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);`
const superadminsSchema = `
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE superadmins (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);`
// Akzeptanzkriterium 1 + 3, Pruefung 1 (inkl. Negativfall doppelte E-Mail).
func TestTenantUserStore_CRUD(t *testing.T) {
pool := setupTestDB(t, "test_iam01_users", usersSchema)
store := NewTenantUserStore(pool)
ctx := context.Background()
created, err := store.Create(ctx, "alice@example.com", "Alice")
if err != nil {
t.Fatalf("create: %v", err)
}
if created.Status != StatusActive {
t.Fatalf("erwartet status active, hat %q", created.Status)
}
got, err := store.Get(ctx, created.ID)
if err != nil {
t.Fatalf("get: %v", err)
}
if got.Email != "alice@example.com" {
t.Fatalf("get email = %q", got.Email)
}
updated, err := store.Update(ctx, created.ID, "", "Alice A.")
if err != nil {
t.Fatalf("update: %v", err)
}
if updated.Name != "Alice A." || updated.Email != "alice@example.com" {
t.Fatalf("update ergebnis unerwartet: %+v", updated)
}
list, err := store.List(ctx)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(list) != 1 {
t.Fatalf("erwartet 1 benutzer, habe %d", len(list))
}
deactivated, err := store.Deactivate(ctx, created.ID)
if err != nil {
t.Fatalf("deactivate: %v", err)
}
if deactivated.Status != StatusInactive {
t.Fatalf("erwartet status inactive, hat %q", deactivated.Status)
}
// Negativfall: doppelte E-Mail-Adresse.
if _, err := store.Create(ctx, "second@example.com", "Bob"); err != nil {
t.Fatalf("create second: %v", err)
}
if _, err := store.Create(ctx, "second@example.com", "Bob Zwei"); !errors.Is(err, ErrEmailTaken) {
t.Fatalf("erwartet ErrEmailTaken, habe %v", err)
}
// Negativfall: fehlender Benutzer.
if _, err := store.Get(ctx, created.ID+"-nicht-vorhanden"); err == nil {
t.Fatalf("erwartet fehler bei unbekannter/ungueltiger id")
}
}
// Akzeptanzkriterium 2 + Pruefung 2: Superadmin-Anlage ohne Tenant-Kontext.
// SuperadminStore.Create hat keinen Tenant-Parameter — es gibt syntaktisch
// keine Moeglichkeit, hier versehentlich einen Tenant-Sonderfall zu vergessen.
func TestSuperadminStore_CreateWithoutTenantContext(t *testing.T) {
pool := setupTestDB(t, "test_iam01_superadmins", superadminsSchema)
store := NewSuperadminStore(pool)
ctx := context.Background()
created, err := store.Create(ctx, "root@nexarch.internal", "Root")
if err != nil {
t.Fatalf("create superadmin: %v", err)
}
if created.Status != StatusActive {
t.Fatalf("erwartet status active, hat %q", created.Status)
}
got, err := store.Get(ctx, created.ID)
if err != nil {
t.Fatalf("get: %v", err)
}
if got.Email != "root@nexarch.internal" {
t.Fatalf("get email = %q", got.Email)
}
if _, err := store.Create(ctx, "root@nexarch.internal", "Root Zwei"); !errors.Is(err, ErrEmailTaken) {
t.Fatalf("erwartet ErrEmailTaken (globale eindeutigkeit), habe %v", err)
}
deactivated, err := store.Deactivate(ctx, created.ID)
if err != nil {
t.Fatalf("deactivate: %v", err)
}
if deactivated.Status != StatusInactive {
t.Fatalf("erwartet status inactive, hat %q", deactivated.Status)
}
}
-77
View File
@@ -1,77 +0,0 @@
package user
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
// SuperadminStore verwaltet mandantenuebergreifende Superadmin-Konten in der
// Control-Plane-Registry (siehe internal/tenant.Registry). Superadmin-ohne-
// Tenant ist dadurch ein eigener Typ statt eines Sonderfalls von User/
// TenantUserStore — es gibt keinen Tenant-Parameter, den man weglassen
// koennte (IAM-01, "ohne Sonderbehandlung im Code").
type SuperadminStore struct {
pool *pgxpool.Pool
}
func NewSuperadminStore(pool *pgxpool.Pool) *SuperadminStore {
return &SuperadminStore{pool: pool}
}
func (s *SuperadminStore) Create(ctx context.Context, email, name string) (User, error) {
if err := ValidateEmail(email); err != nil {
return User{}, err
}
var u User
u.Email, u.Name, u.Status = email, name, StatusActive
row := s.pool.QueryRow(ctx, `
INSERT INTO superadmins (email, name, status)
VALUES ($1, $2, $3)
RETURNING id, created_at, updated_at
`, u.Email, u.Name, u.Status)
if err := row.Scan(&u.ID, &u.CreatedAt, &u.UpdatedAt); err != nil {
return User{}, mapWriteErr(err)
}
return u, nil
}
func (s *SuperadminStore) Get(ctx context.Context, id string) (User, error) {
return scanUser(s.pool.QueryRow(ctx, `
SELECT id, email, name, status, created_at, updated_at
FROM superadmins WHERE id = $1
`, id))
}
func (s *SuperadminStore) List(ctx context.Context) ([]User, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, email, name, status, created_at, updated_at
FROM superadmins ORDER BY created_at
`)
if err != nil {
return nil, fmt.Errorf("superadmins auflisten: %w", err)
}
defer rows.Close()
var out []User
for rows.Next() {
var u User
if err := rows.Scan(&u.ID, &u.Email, &u.Name, &u.Status, &u.CreatedAt, &u.UpdatedAt); err != nil {
return nil, fmt.Errorf("superadmin lesen: %w", err)
}
out = append(out, u)
}
return out, rows.Err()
}
func (s *SuperadminStore) Deactivate(ctx context.Context, id string) (User, error) {
return scanUser(s.pool.QueryRow(ctx, `
UPDATE superadmins SET status = $2, updated_at = now()
WHERE id = $1
RETURNING id, email, name, status, created_at, updated_at
`, id, StatusInactive))
}
-173
View File
@@ -1,173 +0,0 @@
package user
import (
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
// TenantUserStore verwaltet Benutzer innerhalb GENAU EINER Tenant-Datenbank.
// Welcher Mandant gemeint ist, ergibt sich ausschliesslich aus dem
// uebergebenen Pool — es gibt keine tenant_id-Spalte (siehe migrations/tenant/0001_users.up.sql).
type TenantUserStore struct {
pool *pgxpool.Pool
}
func NewTenantUserStore(pool *pgxpool.Pool) *TenantUserStore {
return &TenantUserStore{pool: pool}
}
func (s *TenantUserStore) Create(ctx context.Context, email, name string) (User, error) {
if err := ValidateEmail(email); err != nil {
return User{}, err
}
var u User
u.Email, u.Name, u.Status = email, name, StatusActive
row := s.pool.QueryRow(ctx, `
INSERT INTO users (email, name, status)
VALUES ($1, $2, $3)
RETURNING id, created_at, updated_at
`, u.Email, u.Name, u.Status)
if err := row.Scan(&u.ID, &u.CreatedAt, &u.UpdatedAt); err != nil {
return User{}, mapWriteErr(err)
}
return u, nil
}
func (s *TenantUserStore) Get(ctx context.Context, id string) (User, error) {
return scanUser(s.pool.QueryRow(ctx, `
SELECT id, email, name, status, created_at, updated_at
FROM users WHERE id = $1
`, id))
}
func (s *TenantUserStore) List(ctx context.Context) ([]User, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, email, name, status, created_at, updated_at
FROM users ORDER BY created_at
`)
if err != nil {
return nil, fmt.Errorf("benutzer auflisten: %w", err)
}
defer rows.Close()
var out []User
for rows.Next() {
var u User
if err := rows.Scan(&u.ID, &u.Email, &u.Name, &u.Status, &u.CreatedAt, &u.UpdatedAt); err != nil {
return nil, fmt.Errorf("benutzer lesen: %w", err)
}
out = append(out, u)
}
return out, rows.Err()
}
// Update aendert Name und E-Mail. Eine leere email/name laesst das jeweilige
// Feld unveraendert.
func (s *TenantUserStore) Update(ctx context.Context, id, email, name string) (User, error) {
if email != "" {
if err := ValidateEmail(email); err != nil {
return User{}, err
}
}
row := s.pool.QueryRow(ctx, `
UPDATE users
SET email = COALESCE(NULLIF($2, ''), email),
name = COALESCE(NULLIF($3, ''), name),
updated_at = now()
WHERE id = $1
RETURNING id, email, name, status, created_at, updated_at
`, id, email, name)
u, err := scanUser(row)
if err != nil {
return User{}, mapWriteErr(err)
}
return u, nil
}
// Deactivate setzt den Benutzer auf inaktiv statt ihn zu loeschen.
func (s *TenantUserStore) Deactivate(ctx context.Context, id string) (User, error) {
return scanUser(s.pool.QueryRow(ctx, `
UPDATE users SET status = $2, updated_at = now()
WHERE id = $1
RETURNING id, email, name, status, created_at, updated_at
`, id, StatusInactive))
}
// SetPasswordHash schreibt einen bereits berechneten bcrypt-Hash (siehe
// internal/auth, IAM-02). Der Store selbst kennt kein Klartext-Passwort.
func (s *TenantUserStore) SetPasswordHash(ctx context.Context, id, hash string) error {
tag, err := s.pool.Exec(ctx, `
UPDATE users SET password_hash = $2, updated_at = now() WHERE id = $1
`, id, hash)
if err != nil {
return fmt.Errorf("passwort setzen: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// AuthCredentials wird ausschliesslich fuer den Login-Pfad (internal/auth)
// verwendet und traegt bewusst den password_hash, damit er nicht ueber den
// regulaeren User-Typ/JSON-Serialisierungspfad nach aussen dringen kann.
type AuthCredentials struct {
User User
PasswordHash string
}
// GetByEmailForAuth liefert Benutzer + Passwort-Hash zu einer E-Mail-Adresse
// aus GENAU DIESER Tenant-Datenbank — der Tenant-Scope ergibt sich damit
// zwingend aus dem verwendeten Pool, es gibt keine Moeglichkeit, versehentlich
// ueber Tenant-Grenzen hinweg zu suchen (bekannter archivmail-Fehler, siehe
// IAM-02 "Bekannte Fehler vermeiden").
func (s *TenantUserStore) GetByEmailForAuth(ctx context.Context, email string) (AuthCredentials, error) {
var c AuthCredentials
row := s.pool.QueryRow(ctx, `
SELECT id, email, name, status, created_at, updated_at, password_hash
FROM users WHERE email = $1
`, email)
if err := row.Scan(&c.User.ID, &c.User.Email, &c.User.Name, &c.User.Status,
&c.User.CreatedAt, &c.User.UpdatedAt, &c.PasswordHash); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return AuthCredentials{}, ErrNotFound
}
return AuthCredentials{}, fmt.Errorf("anmeldedaten lesen: %w", err)
}
return c, nil
}
func scanUser(row pgx.Row) (User, error) {
var u User
if err := row.Scan(&u.ID, &u.Email, &u.Name, &u.Status, &u.CreatedAt, &u.UpdatedAt); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return User{}, ErrNotFound
}
return User{}, fmt.Errorf("benutzer lesen: %w", err)
}
return u, nil
}
// mapWriteErr uebersetzt den Unique-Constraint-Verstoss der E-Mail-Spalte in
// einen sprechenden Fehler statt der rohen Postgres-Fehlermeldung.
func mapWriteErr(err error) error {
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) && pgErr.Code == "23505" {
return ErrEmailTaken
}
if errors.Is(err, pgx.ErrNoRows) {
return ErrNotFound
}
return fmt.Errorf("benutzer schreiben: %w", err)
}
-42
View File
@@ -1,42 +0,0 @@
// Package user implementiert Core IAM-01: das Benutzer-Datenmodell und die
// CRUD-Operationen. Tenant-Zugehoerigkeit ist ueber die Zieldatenbank
// gegeben (Modell C, siehe internal/tenant) — Superadmin-Konten leben
// dagegen mandantenuebergreifend in der Registry und sind ueber
// SuperadminStore als eigener, First-Class-Typ modelliert, nicht als
// tenant_id-NULL-Sonderfall in User.
package user
import (
"errors"
"regexp"
"time"
)
type Status string
const (
StatusActive Status = "active"
StatusInactive Status = "inactive"
)
type User struct {
ID string
Email string
Name string
Status Status
CreatedAt time.Time
UpdatedAt time.Time
}
var emailPattern = regexp.MustCompile(`^[^\s@]+@[^\s@]+\.[^\s@]+$`)
var ErrInvalidEmail = errors.New("user: ungueltige E-Mail-Adresse")
var ErrEmailTaken = errors.New("user: E-Mail-Adresse bereits vergeben")
var ErrNotFound = errors.New("user: nicht gefunden")
func ValidateEmail(email string) error {
if !emailPattern.MatchString(email) {
return ErrInvalidEmail
}
return nil
}
-24
View File
@@ -1,24 +0,0 @@
package user
import "testing"
func TestValidateEmail(t *testing.T) {
cases := []struct {
email string
wantErr bool
}{
{"a@b.de", false},
{"a.b+c@sub.example.com", false},
{"", true},
{"keine-email", true},
{"a@b", true},
{"@b.de", true},
}
for _, c := range cases {
err := ValidateEmail(c.email)
if (err != nil) != c.wantErr {
t.Errorf("ValidateEmail(%q) error = %v, wantErr %v", c.email, err, c.wantErr)
}
}
}
-177
View File
@@ -1,177 +0,0 @@
package webhook
import (
"bytes"
"context"
"crypto/subtle"
"encoding/hex"
"fmt"
"net/http"
"time"
"github.com/jackc/pgx/v5"
)
// DefaultMaxAttempts ist die konfigurierbare Obergrenze, ab der eine
// Zustellung endgueltig als fehlgeschlagen gilt (Akzeptanzkriterium 2).
const DefaultMaxAttempts = 5
// DefaultBaseBackoff ist die Basisdauer fuer exponentielles Backoff:
// naechster Versuch nach BaseBackoff * 2^attempt (Akzeptanzkriterium 2).
const DefaultBaseBackoff = 2 * time.Second
// Dispatcher liefert faellige Zustellungen aus. Konfigurierbar in Tests
// (kleine BaseBackoff, kleine MaxAttempts), damit Retry/Backoff/Obergrenze
// ohne minutenlange Wartezeit real durchlaufen werden koennen.
type Dispatcher struct {
pool pgxIface
client *http.Client
MaxAttempts int
BaseBackoff time.Duration
}
// pgxIface ist die schmale Teilmenge von *pgxpool.Pool, die der Dispatcher
// braucht — als Interface, damit Tests keine echte Verbindung fuer reine
// Signatur-/Backoff-Logik brauchen (wird hier aber durchgehend mit echten
// Integrationstests gegen Postgres verwendet, siehe dispatcher_test.go).
type pgxIface interface {
Begin(ctx context.Context) (pgx.Tx, error)
}
func NewDispatcher(pool pgxIface, client *http.Client) *Dispatcher {
if client == nil {
client = &http.Client{Timeout: 5 * time.Second}
}
return &Dispatcher{pool: pool, client: client, MaxAttempts: DefaultMaxAttempts, BaseBackoff: DefaultBaseBackoff}
}
// backoffFor berechnet die Wartezeit vor dem naechsten Versuch: exponentiell
// wachsend mit der Anzahl bereits unternommener Versuche.
func (d *Dispatcher) backoffFor(attempt int) time.Duration {
return d.BaseBackoff * time.Duration(1<<uint(attempt))
}
// ProcessDue liefert ALLE derzeit faelligen Zustellungen aus — dasselbe
// SELECT ... FOR UPDATE SKIP LOCKED-Muster wie
// internal/tenant.Lifecycle.ProcessDueDeletions, damit mehrere Dispatcher-
// Instanzen dieselbe Zustellung nie doppelt bearbeiten.
func (d *Dispatcher) ProcessDue(ctx context.Context) (int, error) {
tx, err := d.pool.Begin(ctx)
if err != nil {
return 0, fmt.Errorf("transaktion starten: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
rows, err := tx.Query(ctx, `
SELECT wd.id, wd.payload, wd.attempt, ws.target_url, ws.secret
FROM webhook_deliveries wd
JOIN webhook_subscriptions ws ON ws.id = wd.subscription_id
WHERE wd.status = $1 AND wd.next_attempt_at <= now()
FOR UPDATE OF wd SKIP LOCKED
`, StatusPending)
if err != nil {
return 0, fmt.Errorf("faellige zustellungen abfragen: %w", err)
}
var deliveries []delivery
for rows.Next() {
var del delivery
if err := rows.Scan(&del.ID, &del.Payload, &del.Attempt, &del.TargetURL, &del.Secret); err != nil {
rows.Close()
return 0, fmt.Errorf("zustellung lesen: %w", err)
}
deliveries = append(deliveries, del)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, err
}
for _, del := range deliveries {
d.attemptOne(ctx, tx, del)
}
if err := tx.Commit(ctx); err != nil {
return 0, fmt.Errorf("transaktion committen: %w", err)
}
return len(deliveries), nil
}
// attemptOne fuehrt GENAU EINEN Zustellversuch aus und aktualisiert den
// Zustellungsdatensatz entsprechend — Erfolg (Akzeptanzkriterium 1/Pruefung
// 1), erneuter Fehlversuch mit Backoff, oder endgueltiges Scheitern nach
// DefaultMaxAttempts (Akzeptanzkriterium 2/Pruefung 2). Ein Fehler bei
// GENAU EINER Zustellung darf die anderen in diesem Batch nicht verhindern.
func (d *Dispatcher) attemptOne(ctx context.Context, tx pgx.Tx, del delivery) {
signature := Sign(del.Secret, del.Payload)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, del.TargetURL, bytes.NewReader(del.Payload))
deliveryErr := err
var statusCode int
if err == nil {
req.Header.Set("Content-Type", "application/json")
req.Header.Set(SignatureHeader, signature)
resp, err := d.client.Do(req)
if err != nil {
deliveryErr = err
} else {
statusCode = resp.StatusCode
resp.Body.Close()
if statusCode < 200 || statusCode >= 300 {
deliveryErr = fmt.Errorf("unerwarteter statuscode %d", statusCode)
}
}
}
if deliveryErr == nil {
_, _ = tx.Exec(ctx, `
UPDATE webhook_deliveries SET status = $2, delivered_at = now(), attempt = attempt + 1
WHERE id = $1
`, del.ID, StatusDelivered)
return
}
nextAttempt := del.Attempt + 1
if nextAttempt >= d.MaxAttempts {
_, _ = tx.Exec(ctx, `
UPDATE webhook_deliveries SET status = $2, attempt = $3, last_error = $4
WHERE id = $1
`, del.ID, StatusFailed, nextAttempt, deliveryErr.Error())
return
}
nextAttemptAt := time.Now().Add(d.backoffFor(nextAttempt))
_, _ = tx.Exec(ctx, `
UPDATE webhook_deliveries SET attempt = $2, next_attempt_at = $3, last_error = $4
WHERE id = $1
`, del.ID, nextAttempt, nextAttemptAt, deliveryErr.Error())
}
// Run ruft ProcessDue in festen Abstaenden auf, bis ctx beendet wird —
// dieselbe Konvention wie internal/tenant.Lifecycle.RunSweeper.
func (d *Dispatcher) Run(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
_, _ = d.ProcessDue(ctx)
}
}
}
// VerifySignature prueft empfaengerseitig, ob signature zu payload und
// secret passt — timing-safe (dasselbe Muster wie internal/audit.timingsafe),
// damit ein Empfaenger die Authentizitaet einer Zustellung pruefen kann
// (Akzeptanzkriterium 3).
func VerifySignature(secret string, payload []byte, signature string) bool {
expected := Sign(secret, payload)
expectedBytes, err1 := hex.DecodeString(expected)
gotBytes, err2 := hex.DecodeString(signature)
if err1 != nil || err2 != nil {
return false
}
return subtle.ConstantTimeCompare(expectedBytes, gotBytes) == 1
}
-131
View File
@@ -1,131 +0,0 @@
// Package webhook implementiert Core API-07: eine zentrale Webhook-Registry
// und Zustellungs-Engine fuer alle Fachmodule (DMS/Mail/Archive/Workflow/AI).
// Module reichen Ereignisse EINMAL zur Zustellung ein und implementieren
// selbst KEINE eigene Retry-/Signatur-Logik — das ist der zentrale Zweck
// dieser Kachel ("Bewusst vermeiden: jedes Modul baut seine eigene
// Webhook-Zustellungs-Engine"). Zustellung laeuft ueber dieselbe
// Postgres-Jobqueue-Konvention (SELECT ... FOR UPDATE SKIP LOCKED) wie
// internal/tenant.Lifecycle.ProcessDueDeletions (TEN-04) und
// internal/notify.Dispatcher (CFG-02) — kein Redis/AMQP.
package webhook
import (
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
const (
StatusPending = "pending"
StatusDelivered = "delivered"
StatusFailed = "failed"
)
// Subscription ist EIN externer Abonnent fuer einen Ereignistyp
// (Akzeptanzkriterium 1: Module registrieren Ereignistypen, externe
// Abonnenten registrieren Ziel-URLs — dieses Paket modelliert die
// Abonnenten-Seite; welche Ereignistypen ein Modul anbietet, ist bewusst
// NICHT Teil dieser Kachel).
type Subscription struct {
ID string
EventType string
TargetURL string
Secret string
CreatedAt time.Time
}
// Store persistiert Abonnements und Zustellversuche.
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
// Subscribe registriert einen Abonnenten fuer einen Ereignistyp. secret wird
// spaeter zur HMAC-Signierung jeder Zustellung an diesen Abonnenten
// verwendet (Akzeptanzkriterium 3).
func (s *Store) Subscribe(ctx context.Context, eventType, targetURL, secret string) (Subscription, error) {
var sub Subscription
sub.EventType, sub.TargetURL, sub.Secret = eventType, targetURL, secret
row := s.pool.QueryRow(ctx, `
INSERT INTO webhook_subscriptions (event_type, target_url, secret)
VALUES ($1, $2, $3)
RETURNING id, created_at
`, eventType, targetURL, secret)
if err := row.Scan(&sub.ID, &sub.CreatedAt); err != nil {
return Subscription{}, fmt.Errorf("abonnement anlegen: %w", err)
}
return sub, nil
}
// Enqueue reicht EIN Ereignis zur Zustellung an ALLE Abonnenten des
// angegebenen Ereignistyps ein — dies ist die EINZIGE Schnittstelle, die
// ein Fachmodul braucht (Akzeptanzkriterium 1). Jeder Abonnent erhaelt
// einen eigenen, unabhaengigen Zustellversuch-Datensatz.
func (s *Store) Enqueue(ctx context.Context, eventType string, payload any) (int, error) {
payloadJSON, err := json.Marshal(payload)
if err != nil {
return 0, fmt.Errorf("payload serialisieren: %w", err)
}
rows, err := s.pool.Query(ctx, `
SELECT id FROM webhook_subscriptions WHERE event_type = $1
`, eventType)
if err != nil {
return 0, fmt.Errorf("abonnenten ermitteln: %w", err)
}
var subscriptionIDs []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
rows.Close()
return 0, fmt.Errorf("abonnent lesen: %w", err)
}
subscriptionIDs = append(subscriptionIDs, id)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, err
}
for _, subID := range subscriptionIDs {
if _, err := s.pool.Exec(ctx, `
INSERT INTO webhook_deliveries (subscription_id, event_type, payload, status, next_attempt_at)
VALUES ($1, $2, $3, $4, now())
`, subID, eventType, payloadJSON, StatusPending); err != nil {
return 0, fmt.Errorf("zustellung einreihen: %w", err)
}
}
return len(subscriptionIDs), nil
}
// delivery ist ein interner Datensatz fuer EINEN Zustellversuch, inklusive
// der zugehoerigen Abonnentendaten (per JOIN geladen).
type delivery struct {
ID string
TargetURL string
Secret string
Payload []byte
Attempt int
}
// Sign berechnet die HMAC-SHA256-Signatur des Payloads (Akzeptanzkriterium
// 3) — hex-kodiert, damit sie problemlos als HTTP-Header uebertragen werden
// kann.
func Sign(secret string, payload []byte) string {
mac := hmac.New(sha256.New, []byte(secret))
mac.Write(payload)
return hex.EncodeToString(mac.Sum(nil))
}
// SignatureHeader ist der HTTP-Header, unter dem die Signatur uebertragen
// wird — dokumentierter Vertrag fuer Empfaenger (Akzeptanzkriterium 3).
const SignatureHeader = "X-Nexarch-Signature-256"
-205
View File
@@ -1,205 +0,0 @@
package webhook
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"os"
"sync/atomic"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupTest(t *testing.T) (*Store, *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 webhook_subscriptions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), event_type TEXT NOT NULL,
target_url TEXT NOT NULL, secret TEXT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS webhook_deliveries (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), subscription_id UUID NOT NULL REFERENCES webhook_subscriptions(id),
event_type TEXT NOT NULL, payload JSONB NOT NULL, status TEXT NOT NULL DEFAULT 'pending',
attempt INT NOT NULL DEFAULT 0, next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(),
last_error TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), delivered_at TIMESTAMPTZ
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
cleanup := func() { pool.Close() }
return NewStore(pool), pool, cleanup
}
func uniqueEventType(prefix string) string {
return prefix + "-" + time.Now().Format("150405.000000000")
}
// Akzeptanzkriterium 1 + Pruefung 1: ein Modul reicht ein Ereignis ein
// (Enqueue) ohne eigene Zustellungslogik, der zentrale Dispatcher liefert
// zuverlaessig aus.
func TestEnqueueAndProcessDue_DeliversSuccessfully(t *testing.T) {
store, pool, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
var receivedBody []byte
var receivedSignature string
target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
receivedBody = body
receivedSignature = r.Header.Get(SignatureHeader)
w.WriteHeader(http.StatusOK)
}))
defer target.Close()
eventType := uniqueEventType("dms.file.created")
sub, err := store.Subscribe(ctx, eventType, target.URL, "geheimes-secret")
if err != nil {
t.Fatalf("subscribe: %v", err)
}
n, err := store.Enqueue(ctx, eventType, map[string]string{"file_id": "42"})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
if n != 1 {
t.Fatalf("erwartet 1 eingereihte zustellung, habe %d", n)
}
dispatcher := NewDispatcher(pool, target.Client())
processed, err := dispatcher.ProcessDue(ctx)
if err != nil {
t.Fatalf("processdue: %v", err)
}
if processed != 1 {
t.Fatalf("erwartet 1 verarbeitete zustellung, habe %d", processed)
}
var status string
if err := pool.QueryRow(ctx, `SELECT status FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status); err != nil {
t.Fatalf("status lesen: %v", err)
}
if status != StatusDelivered {
t.Fatalf("status = %q, want %q", status, StatusDelivered)
}
// Postgres' JSONB-Spalte kann die Byte-Repraesentation des Payloads
// gegenueber dem urspruenglichen json.Marshal kanonisieren (z.B.
// Leerzeichen) — das ist unschaedlich, denn der Dispatcher signiert
// IMMER exakt die Bytes, die er auch sendet. Die Pruefung vergleicht
// deshalb Signatur gegen tatsaechlich empfangene Bytes (Selbstkonsistenz),
// nicht gegen eine unabhaengig neu marshalte Referenz.
if !VerifySignature("geheimes-secret", receivedBody, receivedSignature) {
t.Fatalf("empfangene signatur %q passt nicht zum empfangenen payload %q", receivedSignature, receivedBody)
}
var decoded map[string]string
if err := json.Unmarshal(receivedBody, &decoded); err != nil {
t.Fatalf("empfangener payload nicht als json lesbar: %v", err)
}
if decoded["file_id"] != "42" {
t.Fatalf("empfangener payload = %v, want file_id=42", decoded)
}
}
// Akzeptanzkriterium 2 + Pruefung 2: fehlschlagendes Ziel loest Retry mit
// wachsendem Backoff aus und endet nach der konfigurierten Obergrenze in
// "failed".
func TestProcessDue_RetriesWithBackoffThenMarksFailed(t *testing.T) {
store, pool, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
var callCount int32
target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&callCount, 1)
w.WriteHeader(http.StatusInternalServerError)
}))
defer target.Close()
eventType := uniqueEventType("mail.send.failed")
sub, err := store.Subscribe(ctx, eventType, target.URL, "secret")
if err != nil {
t.Fatalf("subscribe: %v", err)
}
if _, err := store.Enqueue(ctx, eventType, map[string]string{"x": "y"}); err != nil {
t.Fatalf("enqueue: %v", err)
}
dispatcher := NewDispatcher(pool, target.Client())
dispatcher.MaxAttempts = 2
dispatcher.BaseBackoff = 30 * time.Millisecond
// 1. Versuch: schlaegt fehl, ist aber noch nicht die Obergrenze.
if _, err := dispatcher.ProcessDue(ctx); err != nil {
t.Fatalf("processdue 1: %v", err)
}
var status string
var attempt int
if err := pool.QueryRow(ctx, `SELECT status, attempt FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status, &attempt); err != nil {
t.Fatalf("status lesen 1: %v", err)
}
if status != StatusPending || attempt != 1 {
t.Fatalf("nach 1. fehlschlag: status=%q attempt=%d, want pending/1", status, attempt)
}
// Sofort erneut verarbeiten: Backoff ist noch nicht abgelaufen -> nichts faellig.
processedTooEarly, err := dispatcher.ProcessDue(ctx)
if err != nil {
t.Fatalf("processdue (zu frueh): %v", err)
}
if processedTooEarly != 0 {
t.Fatal("erwartet 0 verarbeitete zustellungen, solange backoff nicht abgelaufen ist")
}
time.Sleep(dispatcher.backoffFor(1) + 20*time.Millisecond)
// 2. Versuch: erreicht MaxAttempts=2 -> endgueltig fehlgeschlagen.
if _, err := dispatcher.ProcessDue(ctx); err != nil {
t.Fatalf("processdue 2: %v", err)
}
if err := pool.QueryRow(ctx, `SELECT status, attempt FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status, &attempt); err != nil {
t.Fatalf("status lesen 2: %v", err)
}
if status != StatusFailed || attempt != 2 {
t.Fatalf("nach 2. fehlschlag: status=%q attempt=%d, want failed/2", status, attempt)
}
if atomic.LoadInt32(&callCount) != 2 {
t.Fatalf("erwartet genau 2 tatsaechliche zustellversuche, habe %d", callCount)
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Signaturpruefung erkennt eine
// manipulierte Payload zuverlaessig.
func TestVerifySignature_DetectsTamperedPayload(t *testing.T) {
secret := "geteiltes-geheimnis"
payload := []byte(`{"file_id":"42"}`)
signature := Sign(secret, payload)
if !VerifySignature(secret, payload, signature) {
t.Fatal("erwartet gueltige signatur fuer unveraenderte payload")
}
tampered := []byte(`{"file_id":"99"}`)
if VerifySignature(secret, tampered, signature) {
t.Fatal("erwartet ungueltige signatur fuer manipulierte payload")
}
if VerifySignature("falsches-secret", payload, signature) {
t.Fatal("erwartet ungueltige signatur bei falschem secret")
}
}
-1
View File
@@ -1 +0,0 @@
DROP TABLE IF EXISTS superadmins;
-14
View File
@@ -1,14 +0,0 @@
-- Superadmin-Konten arbeiten mandantenuebergreifend und leben deshalb in der
-- Control-Plane-Registry (siehe TEN-01), nicht in einer Tenant-Datenbank.
-- Das bildet "Superadmin ohne Tenant" strukturell als First-Class-Zustand ab,
-- statt ihn als Sonderfall in der Tenant-users-Tabelle zu behandeln
-- (IAM-01, siehe core-kanban/tickets/IAM-01.md — bekannte Fehler vermeiden).
-- E-Mail-Eindeutigkeit ist hier global, da die Registry-DB einmalig existiert.
CREATE TABLE superadmins (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS audit_events;
+17
View File
@@ -0,0 +1,17 @@
-- 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);
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS feature_flags;
+10
View File
@@ -0,0 +1,10 @@
-- Feature-Flags zentral je Mandant/Zielgruppe (LIC-02, siehe core-kanban/tickets/LIC-02.md).
-- Lebt in der Registry-DB, nicht pro Tenant-Datenbank — Flags sind eine
-- Core-weite Konfiguration, keine Mandanten-Geschaeftsdaten.
CREATE TABLE feature_flags (
key TEXT PRIMARY KEY,
enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0 CHECK (rollout_percentage BETWEEN 0 AND 100),
target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS tenant_licenses;
+12
View File
@@ -0,0 +1,12 @@
-- 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()
);
+2
View File
@@ -0,0 +1,2 @@
DROP TABLE IF EXISTS module_credentials;
DROP TABLE IF EXISTS modules;
+18
View File
@@ -0,0 +1,18 @@
-- Modul-Registry & Aktivierungspruefung (API-02, siehe core-kanban/tickets/API-02.md).
CREATE TABLE modules (
name TEXT PRIMARY KEY,
version TEXT NOT NULL CHECK (version <> ''),
required_flags TEXT[] NOT NULL DEFAULT '{}',
registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Service-Credential je Modul-Instanz, bei Provisionierung ausgestellt
-- (Akzeptanzkriterium 4). secret_hash enthaelt NIEMALS das Secret im
-- Klartext, nur dessen SHA-256-Hash (Timing-safe-Vergleich beim Login,
-- Referenzmuster siehe AUD-02).
CREATE TABLE 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()
);
+2
View File
@@ -0,0 +1,2 @@
DROP TABLE IF EXISTS usage_quotas;
DROP TABLE IF EXISTS usage_counters;
+15
View File
@@ -0,0 +1,15 @@
-- 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)
);
+2
View File
@@ -0,0 +1,2 @@
DROP TABLE resync_usage_buffer;
DROP TABLE resync_audit_buffer;
+22
View File
@@ -0,0 +1,22 @@
-- 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()
);
-1
View File
@@ -1 +0,0 @@
DROP TABLE IF EXISTS users;
-16
View File
@@ -1,16 +0,0 @@
-- Benutzer-Datenmodell (IAM-01, siehe core-kanban/tickets/IAM-01.md).
-- Diese Migration laeuft in der DB EINES Mandanten (Modell C, siehe TEN-01) —
-- die Tenant-Zugehoerigkeit ist implizit durch die Datenbankverbindung
-- gegeben, es gibt daher bewusst KEINE tenant_id-Spalte.
-- E-Mail-Eindeutigkeit ist hier tenant-scoped: der UNIQUE-Constraint gilt
-- nur innerhalb dieser einen Tenant-Datenbank.
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE users (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
@@ -1 +0,0 @@
ALTER TABLE users DROP COLUMN password_hash;
@@ -1,3 +0,0 @@
-- Passwort-Hash-Spalte fuer Login (IAM-02, siehe core-kanban/tickets/IAM-02.md).
-- Enthaelt AUSSCHLIESSLICH den bcrypt-Hash, niemals das Klartext-Passwort.
ALTER TABLE users ADD COLUMN password_hash TEXT NOT NULL DEFAULT '';
+8
View File
@@ -1,4 +1,11 @@
#!/usr/bin/env bash
# Setzt die nexarch-Testumgebung zurueck: loescht die geteilte
# Registry-Tabelle "tenants" in der postgres-Wartungsdatenbank sowie alle
# tenant_*-Datenbanken. Noetig, weil verschiedene Feature-Branches
# unterschiedliche Registry-Schemata erwarten, aber dieselbe physische
# Postgres-Instanz auf dem Testhost teilen (siehe [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/reset-test-env.sh
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"
@@ -7,6 +14,7 @@ ROLE="nexarch_test"
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;"
dbs=$(psql -h localhost -U "$ROLE" -d postgres -tAc "SELECT datname FROM pg_database WHERE datname LIKE 'tenant\_%' ESCAPE '\'")
for db in $dbs; do
+6
View File
@@ -1,4 +1,10 @@
#!/usr/bin/env bash
# Ein-Kommando-Pruefung fuer den aktuellen Code-Stand auf dem Testhost:
# Registry+Tenant-DBs zuruecksetzen, dann build/vet/test in einem Rutsch.
# -p 1 ist Pflicht, da mehrere Pakete dieselbe physische Registry-Tabelle auf
# dem Testhost teilen (siehe [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/run-checks.sh
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"