Compare commits

..
Author SHA1 Message Date
sysopsandClaude Sonnet 5 f7863fd4d2 API-05: verteilte-jwt-verifikation-rechte-feature-flag-cache-kontrakt
internal/moduletrust: asymmetrische JWT-Signatur (Ed25519) mit JWKS-
Verteilung, wie im Entscheidungsverlauf "Vertrauensstellung Core<->Module"
(nexarch-state.json) festgelegt. Getrennt von IAM-02s HS256-Session-Cookie
(Browser-Login bleibt unangetastet) — dies ist der Modul-zu-Core-
Vertrauensmechanismus.

KeyManager haelt ALLE noch gueltigen Schluesselpaare (nicht nur das aktuell
signierende); Rotate() erzeugt einen neuen Schluessel, alte bleiben in
PublicKeySet() erhalten — bereits ausgestellte Tokens bleiben dadurch nach
einer Rotation weiterhin verifizierbar (Akzeptanzkriterium 3, keine
Ausfallzeit). ServeJWKS/ParseJWKS sind der Verteilungsmechanismus.

StaleCache[T] ist der generische Rechte-/Feature-Flag-Cache-Kontrakt
(Akzeptanzkriterium 2), mit zwei explizit benannten und begruendeten
Verhalten: Get() ist FAIL-OPEN (nutzt bei Core-Ausfall einen vorhandenen,
abgelaufenen Stand weiter — ein bereits authentifiziertes Modul soll nicht
hart blockieren), RequireFresh() ist FAIL-CLOSED (nie zwischengespeichert,
schlaegt bei Core-Ausfall klar fehl — fuer sicherheitskritische Aktionen wie
einen neuen Login). LIC-02s internal/flag.Service implementiert bereits
denselben Kontrakt fuer Feature-Flags; StaleCache verallgemeinert dasselbe
Muster fuer JWT-Schluessel, damit beide Faelle derselben dokumentierten
Policy folgen statt zwei unterschiedlichen Ad-hoc-Loesungen.

Verifier.Verify ruft KeyFetchFunc nur bei abgelaufener TTL auf, nicht pro
Aufruf (Akzeptanzkriterium 1) — Signaturpruefung selbst ist immer lokal.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Core simuliert abgeschaltet, andere Module bleiben fuer bereits
   authentifizierte Nutzer funktionsfaehig bis TTL/Fail-Open greift —
   TestVerify_FailsOpenWhenCoreUnreachableButStaleKeysExist: Verify()
   funktioniert weiter mit letztbekanntem Schluesselstand. PASS.
2. Neue sicherheitskritische Aktion schlaegt bei Core-Ausfall klar fehl,
   statt andere Funktionen mitzureissen —
   TestRequireFreshKeys_FailsClosedWhenCoreUnreachable: Fehler trotz
   vorhandenem (aelterem) Cache-Stand. PASS.
3. Schluesselrotation ohne Downtime in einem simulierten zweiten Modul —
   TestRotate_NoDowntimeForAlreadyIssuedTokens: vor UND nach Rotation
   ausgestellte Tokens beide weiterhin gueltig fuer Modul B. PASS.

Zusaetzlich: TestVerify_DoesNotFetchPerCall belegt Akzeptanzkriterium 1
direkt (10 Verify-Aufrufe, genau 1 Fetch). PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 21:38:57 +02:00
15 changed files with 529 additions and 835 deletions
-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)
}
}
+86
View File
@@ -0,0 +1,86 @@
package moduletrust
import (
"context"
"fmt"
"sync"
"time"
)
// StaleCache ist der generische Rechte-/Feature-Flag-Cache-Kontrakt
// (Akzeptanzkriterium 2): TTL-basiert, mit explizitem, benanntem Verhalten
// bei abgelaufenem Cache waehrend Core nicht erreichbar ist.
//
// - Get: FAIL-OPEN fuer Lesevorgaenge. Schlaegt der Refresh fehl, aber es
// gibt bereits einen (wenn auch abgelaufenen) Stand, wird dieser mit
// stale=true zurueckgegeben — Begruendung: ein bereits authentifiziertes
// Modul soll mit dem letztbekannten Stand weiterarbeiten koennen statt
// hart zu blockieren (siehe "Bekannte Fehler vermeiden" im Ticket).
// Existiert noch nie ein Stand, gibt es keinen sinnvollen Fallback —
// dann liefert auch Get einen Fehler.
// - RequireFresh: FAIL-CLOSED fuer sicherheitskritische Aktionen (z.B.
// ein komplett NEUER Login). Nutzt NIEMALS einen zwischengespeicherten
// Stand, ruft immer frisch ab — Begruendung: eine neue Vertrauens-
// entscheidung darf nicht auf veralteten Daten beruhen, auch wenn das
// bedeutet, dass die Aktion bei Core-Ausfall sichtbar fehlschlaegt statt
// unsicher "irgendwie" durchgelassen zu werden.
//
// LIC-02 (internal/flag.Service) implementiert bereits denselben Kontrakt
// fuer Feature-Flags — StaleCache verallgemeinert dasselbe Muster fuer
// JWT-Signaturschluessel, damit beide Faelle derselben dokumentierten
// Policy folgen.
type StaleCache[T any] struct {
mu sync.RWMutex
value T
hasValue bool
fetchedAt time.Time
ttl time.Duration
fetch func(ctx context.Context) (T, error)
}
func NewStaleCache[T any](ttl time.Duration, fetch func(ctx context.Context) (T, error)) *StaleCache[T] {
return &StaleCache[T]{ttl: ttl, fetch: fetch}
}
// Get liefert den Cache-Wert. FAIL-OPEN: bei Refresh-Fehler wird ein
// vorhandener, ggf. abgelaufener Stand zurueckgegeben (stale=true).
func (c *StaleCache[T]) Get(ctx context.Context) (value T, stale bool, err error) {
c.mu.RLock()
fresh := c.hasValue && time.Since(c.fetchedAt) < c.ttl
if fresh {
v := c.value
c.mu.RUnlock()
return v, false, nil
}
c.mu.RUnlock()
newVal, fetchErr := c.fetch(ctx)
if fetchErr == nil {
c.mu.Lock()
c.value, c.hasValue, c.fetchedAt = newVal, true, time.Now()
c.mu.Unlock()
return newVal, false, nil
}
c.mu.RLock()
defer c.mu.RUnlock()
if c.hasValue {
return c.value, true, nil
}
var zero T
return zero, false, fmt.Errorf("cache leer und refresh fehlgeschlagen: %w", fetchErr)
}
// RequireFresh ruft IMMER frisch ab (FAIL-CLOSED) — fuer sicherheitskritische
// Aktionen, die niemals auf einem zwischengespeicherten Stand basieren duerfen.
func (c *StaleCache[T]) RequireFresh(ctx context.Context) (T, error) {
v, err := c.fetch(ctx)
if err != nil {
var zero T
return zero, fmt.Errorf("core nicht erreichbar, sicherheitskritische aktion abgelehnt: %w", err)
}
c.mu.Lock()
c.value, c.hasValue, c.fetchedAt = v, true, time.Now()
c.mu.Unlock()
return v, nil
}
+79
View File
@@ -0,0 +1,79 @@
package moduletrust
import (
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/golang-jwt/jwt/v5"
)
type Claims struct {
Subject string `json:"sub"`
TenantSlug string `json:"tenant"`
jwt.RegisteredClaims
}
// Issue signiert ein Token mit dem aktuellen Signierschluessel und traegt
// dessen KID im JWT-Header ein — der Verifier auf Modulseite waehlt darueber
// den passenden oeffentlichen Schluessel aus PublicKeySet() aus.
func (m *KeyManager) Issue(subject, tenantSlug string, ttl time.Duration) (string, error) {
key, err := m.SigningKey()
if err != nil {
return "", err
}
now := time.Now()
claims := Claims{
Subject: subject,
TenantSlug: tenantSlug,
RegisteredClaims: jwt.RegisteredClaims{
IssuedAt: jwt.NewNumericDate(now),
ExpiresAt: jwt.NewNumericDate(now.Add(ttl)),
},
}
token := jwt.NewWithClaims(jwt.SigningMethodEdDSA, claims)
token.Header["kid"] = key.KID
return token.SignedString(key.Private)
}
type jwksResponse struct {
Keys []jwksKey `json:"keys"`
}
type jwksKey struct {
Kid string `json:"kid"`
PublicKey string `json:"public_key"` // base64 (raw Ed25519, 32 Byte)
}
// ServeJWKS liefert alle bekannten oeffentlichen Schluessel als JSON —
// Module fragen dies periodisch ab (nicht pro Request), siehe Verifier.
func (m *KeyManager) ServeJWKS(w http.ResponseWriter, r *http.Request) {
set := m.PublicKeySet()
resp := jwksResponse{Keys: make([]jwksKey, 0, len(set))}
for kid, pub := range set {
resp.Keys = append(resp.Keys, jwksKey{Kid: kid, PublicKey: base64.StdEncoding.EncodeToString(pub)})
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(resp)
}
// ParseJWKS dekodiert die JSON-Antwort von ServeJWKS zurueck in kid->PublicKey
// — Hilfsfunktion fuer Module, die JWKS per HTTP abrufen.
func ParseJWKS(data []byte) (map[string][]byte, error) {
var resp jwksResponse
if err := json.Unmarshal(data, &resp); err != nil {
return nil, fmt.Errorf("jwks parsen: %w", err)
}
out := make(map[string][]byte, len(resp.Keys))
for _, k := range resp.Keys {
raw, err := base64.StdEncoding.DecodeString(k.PublicKey)
if err != nil {
return nil, fmt.Errorf("oeffentlichen schluessel %q dekodieren: %w", k.Kid, err)
}
out[k.Kid] = raw
}
return out, nil
}
+83
View File
@@ -0,0 +1,83 @@
// Package moduletrust implementiert Core API-05: asymmetrische JWT-Signatur
// mit Schluesselverteilung (JWKS), damit DMS/Mail/Archive/Workflow JWTs
// LOKAL verifizieren koennen, ohne pro Aufruf einen synchronen Request an
// Core zu stellen — Core darf Fundament sein, ohne zum Flaschenhals zu
// werden (siehe Entscheidungsverlauf "Vertrauensstellung Core<->Module" in
// nexarch-state.json). IAM-02s HS256-Session-Cookie (Browser-Login) bleibt
// unangetastet — dies ist ein zusaetzlicher, getrennter Vertrauensmechanismus
// fuer Modul-zu-Modul/Modul-zu-Core-Aufrufe.
package moduletrust
import (
"crypto/ed25519"
"crypto/rand"
"encoding/hex"
"fmt"
"sync"
)
type KeyPair struct {
KID string
Private ed25519.PrivateKey
Public ed25519.PublicKey
}
// KeyManager haelt ALLE noch gueltigen Schluesselpaare — nicht nur das
// aktuell signierende. Rotate erzeugt ein neues Paar und behaelt die alten
// fuer die Verifikation bereits ausgestellter Tokens (Akzeptanzkriterium 3:
// Rotation ohne Ausfallzeit fuer andere Module).
type KeyManager struct {
mu sync.RWMutex
keys []KeyPair // aeltestes zuerst, neuestes zuletzt
}
func NewKeyManager() (*KeyManager, error) {
m := &KeyManager{}
if _, err := m.Rotate(); err != nil {
return nil, err
}
return m, nil
}
// Rotate erzeugt ein neues Ed25519-Schluesselpaar mit eigener KID und macht
// es zum aktuellen Signierschluessel. Aeltere Schluessel bleiben in
// PublicKeySet() erhalten, damit bereits ausgestellte Tokens weiterhin
// verifizierbar sind.
func (m *KeyManager) Rotate() (KeyPair, error) {
pub, priv, err := ed25519.GenerateKey(nil)
if err != nil {
return KeyPair{}, fmt.Errorf("schluesselpaar erzeugen: %w", err)
}
kidBytes := make([]byte, 8)
if _, err := rand.Read(kidBytes); err != nil {
return KeyPair{}, fmt.Errorf("kid erzeugen: %w", err)
}
kp := KeyPair{KID: hex.EncodeToString(kidBytes), Private: priv, Public: pub}
m.mu.Lock()
m.keys = append(m.keys, kp)
m.mu.Unlock()
return kp, nil
}
// SigningKey liefert den aktuellen (neuesten) Schluessel zum Signieren neuer Tokens.
func (m *KeyManager) SigningKey() (KeyPair, error) {
m.mu.RLock()
defer m.mu.RUnlock()
if len(m.keys) == 0 {
return KeyPair{}, fmt.Errorf("moduletrust: kein schluessel vorhanden")
}
return m.keys[len(m.keys)-1], nil
}
// PublicKeySet liefert ALLE bekannten oeffentlichen Schluessel (kid ->
// public key) — die Grundlage fuer den JWKS-Endpunkt.
func (m *KeyManager) PublicKeySet() map[string]ed25519.PublicKey {
m.mu.RLock()
defer m.mu.RUnlock()
out := make(map[string]ed25519.PublicKey, len(m.keys))
for _, k := range m.keys {
out[k.KID] = k.Public
}
return out
}
+214
View File
@@ -0,0 +1,214 @@
package moduletrust
import (
"context"
"crypto/ed25519"
"errors"
"net/http"
"sync"
"testing"
"time"
)
func TestIssueAndVerify_RoundTrip(t *testing.T) {
km, err := NewKeyManager()
if err != nil {
t.Fatalf("new key manager: %v", err)
}
token, err := km.Issue("user-1", "acme", time.Hour)
if err != nil {
t.Fatalf("issue: %v", err)
}
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
return km.PublicKeySet(), nil
})
claims, err := v.Verify(context.Background(), token)
if err != nil {
t.Fatalf("verify: %v", err)
}
if claims.Subject != "user-1" || claims.TenantSlug != "acme" {
t.Fatalf("claims unerwartet: %+v", claims)
}
}
// Akzeptanzkriterium 1: Verifikation lokal, kein Request pro Aufruf.
func TestVerify_DoesNotFetchPerCall(t *testing.T) {
km, err := NewKeyManager()
if err != nil {
t.Fatalf("new key manager: %v", err)
}
token, err := km.Issue("user-1", "acme", time.Hour)
if err != nil {
t.Fatalf("issue: %v", err)
}
var mu sync.Mutex
fetchCalls := 0
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
mu.Lock()
fetchCalls++
mu.Unlock()
return km.PublicKeySet(), nil
})
for i := 0; i < 10; i++ {
if _, err := v.Verify(context.Background(), token); err != nil {
t.Fatalf("verify %d: %v", i, err)
}
}
mu.Lock()
defer mu.Unlock()
if fetchCalls != 1 {
t.Fatalf("erwartet genau 1 fetch fuer 10 Verify-Aufrufe innerhalb der TTL, habe %d", fetchCalls)
}
}
// Akzeptanzkriterium 2 + Pruefung 1: Core simuliert nicht erreichbar,
// bereits authentifizierte Nutzer bleiben funktionsfaehig (Fail-Open mit
// letztbekanntem Schluesselstand).
func TestVerify_FailsOpenWhenCoreUnreachableButStaleKeysExist(t *testing.T) {
km, err := NewKeyManager()
if err != nil {
t.Fatalf("new key manager: %v", err)
}
token, err := km.Issue("user-1", "acme", time.Hour)
if err != nil {
t.Fatalf("issue: %v", err)
}
coreDown := false
v := NewVerifier(30*time.Millisecond, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
if coreDown {
return nil, errors.New("core nicht erreichbar (simuliert)")
}
return km.PublicKeySet(), nil
})
// Cache vorwaermen, waehrend Core noch erreichbar ist.
if _, err := v.Verify(context.Background(), token); err != nil {
t.Fatalf("verify (warm): %v", err)
}
// "Core abschalten" und TTL ablaufen lassen.
coreDown = true
time.Sleep(50 * time.Millisecond)
if _, err := v.Verify(context.Background(), token); err != nil {
t.Fatalf("verify sollte trotz core-ausfall mit letztbekanntem stand funktionieren: %v", err)
}
}
// Akzeptanzkriterium 2 + Pruefung 2: neue sicherheitskritische Aktionen
// (z.B. neuer Login) schlagen bei Core-Ausfall klar fehl statt unsicher
// durchgelassen zu werden — auch wenn ein (aelterer) Cache-Stand existiert.
func TestRequireFreshKeys_FailsClosedWhenCoreUnreachable(t *testing.T) {
km, err := NewKeyManager()
if err != nil {
t.Fatalf("new key manager: %v", err)
}
coreDown := false
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
if coreDown {
return nil, errors.New("core nicht erreichbar (simuliert)")
}
return km.PublicKeySet(), nil
})
// Cache vorwaermen (existiert jetzt ein "veralteter" gueltiger Stand).
if _, _, err := v.cache.Get(context.Background()); err != nil {
t.Fatalf("warm cache: %v", err)
}
coreDown = true
if err := v.RequireFreshKeys(context.Background()); err == nil {
t.Fatal("erwartet fehler (fail-closed) bei core-ausfall, habe nil")
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Schluesselrotation ohne Ausfallzeit —
// ein bereits ausgestelltes Token bleibt nach Rotation weiterhin
// verifizierbar, ein zweites (simuliertes) Modul bekommt beide Schluessel.
func TestRotate_NoDowntimeForAlreadyIssuedTokens(t *testing.T) {
km, err := NewKeyManager()
if err != nil {
t.Fatalf("new key manager: %v", err)
}
oldToken, err := km.Issue("user-1", "acme", time.Hour)
if err != nil {
t.Fatalf("issue (alt): %v", err)
}
if _, err := km.Rotate(); err != nil {
t.Fatalf("rotate: %v", err)
}
newToken, err := km.Issue("user-2", "acme", time.Hour)
if err != nil {
t.Fatalf("issue (neu): %v", err)
}
// Simuliertes zweites Modul: fragt den vollstaendigen Schluesselsatz ab.
moduleB := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
return km.PublicKeySet(), nil
})
if _, err := moduleB.Verify(context.Background(), oldToken); err != nil {
t.Fatalf("altes token sollte nach rotation weiterhin gueltig sein: %v", err)
}
if _, err := moduleB.Verify(context.Background(), newToken); err != nil {
t.Fatalf("neues token sollte gueltig sein: %v", err)
}
}
func TestVerify_RejectsUnknownKid(t *testing.T) {
km1, _ := NewKeyManager()
km2, _ := NewKeyManager() // komplett anderer, unbekannter schluessel
token, err := km1.Issue("user-1", "acme", time.Hour)
if err != nil {
t.Fatalf("issue: %v", err)
}
v := NewVerifier(time.Hour, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
return km2.PublicKeySet(), nil // kennt km1s schluessel nicht
})
if _, err := v.Verify(context.Background(), token); !errors.Is(err, ErrInvalidToken) {
t.Fatalf("erwartet ErrInvalidToken, habe %v", err)
}
}
func TestJWKSRoundTrip(t *testing.T) {
km, _ := NewKeyManager()
km.Rotate()
var buf []byte
rec := &captureWriter{}
km.ServeJWKS(rec, nil)
buf = rec.body
parsed, err := ParseJWKS(buf)
if err != nil {
t.Fatalf("parse jwks: %v", err)
}
if len(parsed) != len(km.PublicKeySet()) {
t.Fatalf("erwartet %d schluessel, habe %d", len(km.PublicKeySet()), len(parsed))
}
}
type captureWriter struct {
body []byte
header http.Header
}
func (w *captureWriter) Header() http.Header {
if w.header == nil {
w.header = http.Header{}
}
return w.header
}
func (w *captureWriter) Write(p []byte) (int, error) { w.body = append(w.body, p...); return len(p), nil }
func (w *captureWriter) WriteHeader(statusCode int) {}
+67
View File
@@ -0,0 +1,67 @@
package moduletrust
import (
"context"
"crypto/ed25519"
"errors"
"time"
"github.com/golang-jwt/jwt/v5"
)
var ErrInvalidToken = errors.New("moduletrust: ungueltiges token")
// KeyFetchFunc holt den aktuellen Schluesselsatz von Core (z.B. per HTTP-GET
// auf ServeJWKS + ParseJWKS). Wird vom Verifier nur bei abgelaufener TTL
// aufgerufen — NICHT bei jeder Verify()-Anfrage (Akzeptanzkriterium 1).
type KeyFetchFunc func(ctx context.Context) (map[string]ed25519.PublicKey, error)
// Verifier ist die Modulseite von API-05: verifiziert JWTs LOKAL gegen einen
// per StaleCache zwischengespeicherten Schluesselsatz, ohne pro Aufruf einen
// synchronen Request an Core zu stellen.
type Verifier struct {
cache *StaleCache[map[string]ed25519.PublicKey]
}
func NewVerifier(ttl time.Duration, fetch KeyFetchFunc) *Verifier {
return &Verifier{cache: NewStaleCache(ttl, func(ctx context.Context) (map[string]ed25519.PublicKey, error) {
return fetch(ctx)
})}
}
// Verify prueft die Signatur LOKAL gegen den (ggf. abgelaufenen, aber
// vorhandenen) Schluesselsatz — FAIL-OPEN fuer bereits ausgestellte Tokens
// (Akzeptanzkriterium 2): ist Core nicht erreichbar, aber ein alter
// Schluesselsatz bekannt, wird damit weiter verifiziert.
func (v *Verifier) Verify(ctx context.Context, tokenString string) (*Claims, error) {
keys, _, err := v.cache.Get(ctx)
if err != nil {
return nil, err
}
claims := &Claims{}
token, err := jwt.ParseWithClaims(tokenString, claims, func(t *jwt.Token) (interface{}, error) {
if _, ok := t.Method.(*jwt.SigningMethodEd25519); !ok {
return nil, ErrInvalidToken
}
kid, _ := t.Header["kid"].(string)
pub, ok := keys[kid]
if !ok {
return nil, ErrInvalidToken
}
return pub, nil
})
if err != nil || !token.Valid {
return nil, ErrInvalidToken
}
return claims, nil
}
// RequireFreshKeys ruft IMMER frisch von Core ab (FAIL-CLOSED) — fuer
// sicherheitskritische Aktionen wie einen komplett neuen Login
// (Akzeptanzkriterium 2): schlaegt klar fehl, wenn Core nicht erreichbar
// ist, statt auf einem veralteten Schluesselsatz zu vertrauen.
func (v *Verifier) RequireFreshKeys(ctx context.Context) error {
_, err := v.cache.RequireFresh(ctx)
return err
}
-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")
}
}
-2
View File
@@ -1,2 +0,0 @@
DROP TABLE webhook_deliveries;
DROP TABLE webhook_subscriptions;
-27
View File
@@ -1,27 +0,0 @@
-- Zentrale Webhook-Registry & Zustellung (API-07, siehe
-- core-kanban/tickets/API-07.md) — Postgres-basierte Jobqueue, kein
-- Redis/AMQP (Ticket-Vorgabe).
CREATE TABLE 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 INDEX webhook_subscriptions_event_type_idx ON webhook_subscriptions (event_type);
CREATE TABLE 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', -- pending|delivered|failed
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
);
CREATE INDEX webhook_deliveries_due_idx ON webhook_deliveries (status, next_attempt_at);