From e0c82b5d6381ab4d2ab1d4efe240b922e40d332e Mon Sep 17 00:00:00 2001 From: sysops Date: Fri, 28 Aug 2026 10:30:12 +0200 Subject: [PATCH 1/2] QA-07: apiserver+moduleregistry+flag+webhook auf api-05-basis portiert (fuer vertragstests benoetigt) --- internal/apiserver/context.go | 20 ++ internal/apiserver/errors.go | 31 +++ internal/apiserver/middleware.go | 56 +++++ internal/apiserver/server.go | 37 +++ internal/apiserver/server_test.go | 149 +++++++++++ internal/flag/flag.go | 87 +++++++ internal/flag/flag_test.go | 43 ++++ internal/flag/service.go | 87 +++++++ internal/flag/service_test.go | 179 +++++++++++++ internal/moduleregistry/credentials.go | 99 ++++++++ internal/moduleregistry/middleware.go | 46 ++++ internal/moduleregistry/registry.go | 120 +++++++++ internal/moduleregistry/registry_test.go | 304 +++++++++++++++++++++++ internal/webhook/dispatcher.go | 177 +++++++++++++ internal/webhook/webhook.go | 131 ++++++++++ internal/webhook/webhook_test.go | 205 +++++++++++++++ 16 files changed, 1771 insertions(+) create mode 100644 internal/apiserver/context.go create mode 100644 internal/apiserver/errors.go create mode 100644 internal/apiserver/middleware.go create mode 100644 internal/apiserver/server.go create mode 100644 internal/apiserver/server_test.go create mode 100644 internal/flag/flag.go create mode 100644 internal/flag/flag_test.go create mode 100644 internal/flag/service.go create mode 100644 internal/flag/service_test.go create mode 100644 internal/moduleregistry/credentials.go create mode 100644 internal/moduleregistry/middleware.go create mode 100644 internal/moduleregistry/registry.go create mode 100644 internal/moduleregistry/registry_test.go create mode 100644 internal/webhook/dispatcher.go create mode 100644 internal/webhook/webhook.go create mode 100644 internal/webhook/webhook_test.go diff --git a/internal/apiserver/context.go b/internal/apiserver/context.go new file mode 100644 index 0000000..9ed7ca9 --- /dev/null +++ b/internal/apiserver/context.go @@ -0,0 +1,20 @@ +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 +} diff --git a/internal/apiserver/errors.go b/internal/apiserver/errors.go new file mode 100644 index 0000000..2d952bf --- /dev/null +++ b/internal/apiserver/errors.go @@ -0,0 +1,31 @@ +// 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) +} diff --git a/internal/apiserver/middleware.go b/internal/apiserver/middleware.go new file mode 100644 index 0000000..1587071 --- /dev/null +++ b/internal/apiserver/middleware.go @@ -0,0 +1,56 @@ +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)) + } +} diff --git a/internal/apiserver/server.go b/internal/apiserver/server.go new file mode 100644 index 0000000..bd0c0a7 --- /dev/null +++ b/internal/apiserver/server.go @@ -0,0 +1,37 @@ +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 +} diff --git a/internal/apiserver/server_test.go b/internal/apiserver/server_test.go new file mode 100644 index 0000000..6375976 --- /dev/null +++ b/internal/apiserver/server_test.go @@ -0,0 +1,149 @@ +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) + } +} diff --git a/internal/flag/flag.go b/internal/flag/flag.go new file mode 100644 index 0000000..76f23cb --- /dev/null +++ b/internal/flag/flag.go @@ -0,0 +1,87 @@ +// Package flag implementiert Core LIC-02: einen Feature-Flag-Dienst mit +// Strategien (global an/aus, Prozentsatz, Tenant-Zielgruppe) als Kernfunktion +// des Core-Dienstes selbst — keine zusaetzliche Infrastruktur (Unleash-Server +// + eigene DB), siehe "bewusst vermeiden" im LIC-02-Ticket. +package flag + +import ( + "context" + "errors" + "fmt" + "hash/fnv" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +var ErrNotFound = errors.New("flag: nicht gefunden") + +// Flag ist die zentrale Definition — Auswertung (Evaluate) ist bewusst davon +// getrennt (Unleash-Prinzip: Flag-Verwaltung vs. Flag-Auswertung). +type Flag struct { + Key string + Enabled bool + RolloutPercentage int + TargetTenantSlugs []string +} + +// Store ist die Verwaltungsseite (Admin): Flags definieren/lesen. +type Store struct { + pool *pgxpool.Pool +} + +func NewStore(pool *pgxpool.Pool) *Store { + return &Store{pool: pool} +} + +func (s *Store) Set(ctx context.Context, f Flag) error { + if f.TargetTenantSlugs == nil { + f.TargetTenantSlugs = []string{} // pgx uebertraegt ein nil-Slice sonst als SQL NULL statt leerem Array. + } + _, err := s.pool.Exec(ctx, ` + INSERT INTO feature_flags (key, enabled, rollout_percentage, target_tenant_slugs, updated_at) + VALUES ($1, $2, $3, $4, now()) + ON CONFLICT (key) DO UPDATE SET + enabled = $2, rollout_percentage = $3, target_tenant_slugs = $4, updated_at = now() + `, f.Key, f.Enabled, f.RolloutPercentage, f.TargetTenantSlugs) + if err != nil { + return fmt.Errorf("flag speichern: %w", err) + } + return nil +} + +func (s *Store) Get(ctx context.Context, key string) (Flag, error) { + var f Flag + row := s.pool.QueryRow(ctx, ` + SELECT key, enabled, rollout_percentage, target_tenant_slugs + FROM feature_flags WHERE key = $1 + `, key) + if err := row.Scan(&f.Key, &f.Enabled, &f.RolloutPercentage, &f.TargetTenantSlugs); err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return Flag{}, ErrNotFound + } + return Flag{}, fmt.Errorf("flag lesen: %w", err) + } + return f, nil +} + +// evaluate wendet die Strategien in fester Reihenfolge an: globaler +// An/Aus-Schalter zuerst, dann Tenant-Zielgruppe, dann Prozentsatz-Rollout. +// Ein unbekannter/nicht getroffener Fall ergibt false — Fail-Safe-Default, +// kein Feature wird versehentlich aktiv. +func evaluate(f Flag, tenantSlug string) bool { + if f.Enabled { + return true + } + for _, target := range f.TargetTenantSlugs { + if target == tenantSlug { + return true + } + } + if f.RolloutPercentage > 0 { + h := fnv.New32a() + _, _ = h.Write([]byte(f.Key + "|" + tenantSlug)) + return int(h.Sum32()%100) < f.RolloutPercentage + } + return false +} diff --git a/internal/flag/flag_test.go b/internal/flag/flag_test.go new file mode 100644 index 0000000..772a336 --- /dev/null +++ b/internal/flag/flag_test.go @@ -0,0 +1,43 @@ +package flag + +import "testing" + +func TestEvaluate_GlobalEnabled(t *testing.T) { + f := Flag{Key: "k", Enabled: true} + if !evaluate(f, "irgendein-tenant") { + t.Fatal("global aktiviertes flag sollte fuer jeden tenant true liefern") + } +} + +// Akzeptanzkriterium 1 + Pruefung 2: Zielgruppen-Strategie. +func TestEvaluate_TargetTenantStrategy(t *testing.T) { + f := Flag{Key: "k", Enabled: false, TargetTenantSlugs: []string{"acme"}} + if !evaluate(f, "acme") { + t.Fatal("erwartet true fuer tenant in zielgruppe") + } + if evaluate(f, "globex") { + t.Fatal("erwartet false fuer tenant ausserhalb der zielgruppe") + } +} + +func TestEvaluate_RolloutPercentageBoundaries(t *testing.T) { + full := Flag{Key: "k", RolloutPercentage: 100} + if !evaluate(full, "beliebiger-tenant-1") || !evaluate(full, "beliebiger-tenant-2") { + t.Fatal("100% rollout sollte immer true liefern") + } + + none := Flag{Key: "k", RolloutPercentage: 0} + if evaluate(none, "beliebiger-tenant") { + t.Fatal("0% rollout ohne enabled/zielgruppe sollte false liefern") + } +} + +func TestEvaluate_RolloutIsDeterministicPerTenant(t *testing.T) { + f := Flag{Key: "k", RolloutPercentage: 50} + first := evaluate(f, "stabiler-tenant") + for i := 0; i < 5; i++ { + if evaluate(f, "stabiler-tenant") != first { + t.Fatal("rollout-auswertung sollte fuer denselben tenant/key stabil sein") + } + } +} diff --git a/internal/flag/service.go b/internal/flag/service.go new file mode 100644 index 0000000..6718d28 --- /dev/null +++ b/internal/flag/service.go @@ -0,0 +1,87 @@ +package flag + +import ( + "context" + "log/slog" + "sync" + "time" +) + +// DefaultCacheTTL ist die dokumentierte Cache-Invalidierungszeit +// (Akzeptanzkriterium 2/3): eine Aenderung wirkt spaetestens nach dieser +// Zeit auf allen Core-Instanzen, ohne dass ein Dienst neu gestartet werden +// muss (Akzeptanzkriterium 3). +const DefaultCacheTTL = 5 * time.Second + +type cacheEntry struct { + flag Flag + expiresAt time.Time +} + +// Service ist die Auswertungsseite (SDK/Client-Analogon zu Unleash) mit +// lokalem TTL-Cache. Bewusst getrennt von Store (Verwaltung). +type Service struct { + store *Store + ttl time.Duration + + mu sync.RWMutex + cache map[string]cacheEntry +} + +func NewService(store *Store, ttl time.Duration) *Service { + if ttl <= 0 { + ttl = DefaultCacheTTL + } + return &Service{store: store, ttl: ttl, cache: make(map[string]cacheEntry)} +} + +// IsEnabled wertet ein Flag fuer einen Tenant aus. Liefert IMMER einen +// bool ohne Fehlerwert — ein nicht erreichbarer Flag-Dienst darf abhaengige +// Aufrufer nicht zum Absturz bringen oder zu Fehlerbehandlungscode zwingen, +// der leicht vergessen wird (Akzeptanzkriterium 3 / Pruefung 3: dokumentiertes +// Fallback-Verhalten = false, ggf. aus dem zuletzt bekannten Zwischenspeicher). +func (s *Service) IsEnabled(ctx context.Context, tenantSlug, key string) bool { + f, ok := s.resolve(ctx, key) + if !ok { + return false + } + return evaluate(f, tenantSlug) +} + +func (s *Service) resolve(ctx context.Context, key string) (Flag, bool) { + s.mu.RLock() + entry, exists := s.cache[key] + fresh := exists && time.Now().Before(entry.expiresAt) + s.mu.RUnlock() + if fresh { + return entry.flag, true + } + + f, err := s.store.Get(ctx, key) + if err != nil { + if exists { + slog.Warn("feature-flag-dienst nicht erreichbar, nutze zwischengespeicherten stand", + "flag_key", key, "error", err) + return entry.flag, true + } + slog.Warn("feature-flag-dienst nicht erreichbar, kein zwischengespeicherter stand vorhanden, fallback: deaktiviert", + "flag_key", key, "error", err) + return Flag{}, false + } + + s.mu.Lock() + s.cache[key] = cacheEntry{flag: f, expiresAt: time.Now().Add(s.ttl)} + s.mu.Unlock() + return f, true +} + +// Invalidate erzwingt beim naechsten IsEnabled-Aufruf ein sofortiges Neuladen +// aus der Datenbank statt auf den TTL-Ablauf zu warten — wird nach Store.Set +// auf derselben Instanz aufgerufen, damit der Schreiber die eigene Aenderung +// ohne Wartezeit sieht. Andere Core-Instanzen sehen sie spaetestens nach +// DefaultCacheTTL (siehe Akzeptanzkriterium 3). +func (s *Service) Invalidate(key string) { + s.mu.Lock() + delete(s.cache, key) + s.mu.Unlock() +} diff --git a/internal/flag/service_test.go b/internal/flag/service_test.go new file mode 100644 index 0000000..0ddd5b9 --- /dev/null +++ b/internal/flag/service_test.go @@ -0,0 +1,179 @@ +package flag + +import ( + "context" + "os" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +func setupFlagStoreTest(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 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() + )`); err != nil { + t.Fatalf("schema: %v", err) + } + + cleanup := func() { + _, _ = pool.Exec(ctx, `DELETE FROM feature_flags WHERE key LIKE 'test\_%' ESCAPE '\'`) + pool.Close() + } + return NewStore(pool), cleanup +} + +// Akzeptanzkriterium 1 + Pruefung 2: Zielgruppen-Strategie liefert im Test +// die erwartete Auswertung. +func TestService_TargetTenantStrategy(t *testing.T) { + store, cleanup := setupFlagStoreTest(t) + defer cleanup() + ctx := context.Background() + + if err := store.Set(ctx, Flag{Key: "test_target_flag", TargetTenantSlugs: []string{"acme"}}); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, time.Hour) + + if !svc.IsEnabled(ctx, "acme", "test_target_flag") { + t.Fatal("erwartet true fuer tenant in zielgruppe") + } + if svc.IsEnabled(ctx, "globex", "test_target_flag") { + t.Fatal("erwartet false fuer tenant ausserhalb der zielgruppe") + } +} + +// Akzeptanzkriterium 2 + 3 + Pruefung 1: Flag-Aenderung wirkt innerhalb der +// dokumentierten Cache-Invalidierungszeit, automatisiert gemessen. +func TestService_CacheInvalidationTiming(t *testing.T) { + store, cleanup := setupFlagStoreTest(t) + defer cleanup() + ctx := context.Background() + + const ttl = 150 * time.Millisecond + if err := store.Set(ctx, Flag{Key: "test_ttl_flag", Enabled: false}); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, ttl) + + if svc.IsEnabled(ctx, "acme", "test_ttl_flag") { + t.Fatal("erwartet false vor der aenderung") + } + + // Aenderung "auf einer anderen instanz" simulieren: direkt ueber den + // Store, ohne svc.Invalidate aufzurufen. + changedAt := time.Now() + if err := store.Set(ctx, Flag{Key: "test_ttl_flag", Enabled: true}); err != nil { + t.Fatalf("set: %v", err) + } + + // Sofort danach sollte der Cache noch den alten Stand liefern. + if svc.IsEnabled(ctx, "acme", "test_ttl_flag") { + t.Fatal("cache haette den alten (false) stand liefern sollen, direkt nach der aenderung") + } + + deadline := changedAt.Add(ttl + 100*time.Millisecond) + for time.Now().Before(deadline) { + if svc.IsEnabled(ctx, "acme", "test_ttl_flag") { + elapsed := time.Since(changedAt) + t.Logf("aenderung wurde nach %s wirksam (ziel: innerhalb %s + toleranz)", elapsed, ttl) + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("aenderung wurde nicht innerhalb von %s wirksam", deadline.Sub(changedAt)) +} + +func TestService_InvalidateForcesImmediateRefresh(t *testing.T) { + store, cleanup := setupFlagStoreTest(t) + defer cleanup() + ctx := context.Background() + + if err := store.Set(ctx, Flag{Key: "test_invalidate_flag", Enabled: false}); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, time.Hour) // lange TTL, damit Invalidate den unterschied macht + _ = svc.IsEnabled(ctx, "acme", "test_invalidate_flag") + + if err := store.Set(ctx, Flag{Key: "test_invalidate_flag", Enabled: true}); err != nil { + t.Fatalf("set: %v", err) + } + svc.Invalidate("test_invalidate_flag") + + if !svc.IsEnabled(ctx, "acme", "test_invalidate_flag") { + t.Fatal("erwartet sofort sichtbaren neuen stand nach Invalidate") + } +} + +// Akzeptanzkriterium 3 + Pruefung 3: Ausfall des Flag-Dienstes fuehrt zu +// dokumentiertem Fallback-Verhalten, nicht zum Absturz. +func TestService_FallsBackOnStoreFailure(t *testing.T) { + store, cleanup := setupFlagStoreTest(t) + defer cleanup() + ctx := context.Background() + + if err := store.Set(ctx, Flag{Key: "test_fallback_flag", Enabled: true}); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, time.Hour) + + // Cache vorwaermen, waehrend die DB noch erreichbar ist. + if !svc.IsEnabled(ctx, "acme", "test_fallback_flag") { + t.Fatal("erwartet true bei funktionierender db") + } + + brokenPool, err := pgxpool.New(ctx, "postgresql://nonexistent-host-fuer-test:5432/x?connect_timeout=1") + if err != nil { + t.Fatalf("broken pool erstellen (sollte nicht sofort verbinden): %v", err) + } + brokenStore := NewStore(brokenPool) + + svcWithCache := NewService(brokenStore, time.Nanosecond) // TTL sofort abgelaufen, erzwingt reload-versuch + svcWithCache.mu.Lock() + svcWithCache.cache["test_fallback_flag"] = cacheEntry{ + flag: Flag{Key: "test_fallback_flag", Enabled: true}, + expiresAt: time.Now().Add(-time.Hour), // bereits abgelaufen + } + svcWithCache.mu.Unlock() + + func() { + defer func() { + if r := recover(); r != nil { + t.Fatalf("IsEnabled hat gepanict statt einen fallback zu liefern: %v", r) + } + }() + if !svcWithCache.IsEnabled(ctx, "acme", "test_fallback_flag") { + t.Fatal("erwartet fallback auf zwischengespeicherten (true) stand bei db-ausfall") + } + }() + + // Voellig frischer Dienst ohne jeglichen cache + kaputte db -> sicherer + // default false, kein absturz. + freshSvc := NewService(brokenStore, time.Hour) + func() { + defer func() { + if r := recover(); r != nil { + t.Fatalf("IsEnabled hat gepanict: %v", r) + } + }() + if freshSvc.IsEnabled(ctx, "acme", "test_fallback_flag") { + t.Fatal("erwartet fail-safe false ohne cache und mit kaputter db") + } + }() +} diff --git a/internal/moduleregistry/credentials.go b/internal/moduleregistry/credentials.go new file mode 100644 index 0000000..c941096 --- /dev/null +++ b/internal/moduleregistry/credentials.go @@ -0,0 +1,99 @@ +package moduleregistry + +import ( + "context" + "crypto/rand" + "crypto/sha256" + "crypto/subtle" + "encoding/hex" + "errors" + "fmt" + + "github.com/jackc/pgx/v5" +) + +var ( + ErrModuleNotRegistered = errors.New("moduleregistry: modul muss vor provisionierung registriert sein") + ErrInvalidCredential = errors.New("moduleregistry: ungueltiges oder fehlendes service-credential") +) + +// Provision stellt ein Service-Credential (Client-ID + Secret) fuer eine +// Modul-Instanz aus (Akzeptanzkriterium 4). Das Secret wird NUR beim +// Ausstellen im Klartext zurueckgegeben, gespeichert wird ausschliesslich +// dessen SHA-256-Hash. +func (r *Registry) Provision(ctx context.Context, moduleName string) (clientID, secret string, err error) { + if _, err := r.Get(ctx, moduleName); err != nil { + if errors.Is(err, ErrModuleNotFound) { + return "", "", ErrModuleNotRegistered + } + return "", "", err + } + + clientID, err = randomToken(16) + if err != nil { + return "", "", fmt.Errorf("client-id erzeugen: %w", err) + } + secret, err = randomToken(32) + if err != nil { + return "", "", fmt.Errorf("secret erzeugen: %w", err) + } + hash := hashSecret(secret) + + _, err = r.pool.Exec(ctx, ` + INSERT INTO module_credentials (module_name, client_id, secret_hash, issued_at) + VALUES ($1, $2, $3, now()) + ON CONFLICT (module_name) DO UPDATE SET client_id = $2, secret_hash = $3, issued_at = now() + `, moduleName, clientID, hash) + if err != nil { + return "", "", fmt.Errorf("credential speichern: %w", err) + } + return clientID, secret, nil +} + +// Authenticate prueft ein Service-Credential timing-safe (Referenzmuster +// siehe AUD-02) — Aufrufe ohne gueltiges Credential werden abgelehnt +// (Akzeptanzkriterium 4 / Pruefung 4). +func (r *Registry) Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error) { + if clientID == "" || secret == "" { + return "", false, nil + } + + var storedHash []byte + err = r.pool.QueryRow(ctx, ` + SELECT module_name, secret_hash FROM module_credentials WHERE client_id = $1 + `, clientID).Scan(&moduleName, &storedHash) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return "", false, nil + } + return "", false, fmt.Errorf("credential lesen: %w", err) + } + + if !timingSafeEqual(hashSecret(secret), storedHash) { + return "", false, nil + } + return moduleName, true, nil +} + +func randomToken(n int) (string, error) { + buf := make([]byte, n) + if _, err := rand.Read(buf); err != nil { + return "", err + } + return hex.EncodeToString(buf), nil +} + +func hashSecret(secret string) []byte { + sum := sha256.Sum256([]byte(secret)) + return sum[:] +} + +// timingSafeEqual folgt derselben Referenzimplementierung wie AUD-02 +// (subtle.ConstantTimeCompare) — projektweite Konvention fuer jeden +// sicherheitsrelevanten Vergleich. +func timingSafeEqual(a, b []byte) bool { + if len(a) != len(b) { + return false + } + return subtle.ConstantTimeCompare(a, b) == 1 +} diff --git a/internal/moduleregistry/middleware.go b/internal/moduleregistry/middleware.go new file mode 100644 index 0000000..24ccf8c --- /dev/null +++ b/internal/moduleregistry/middleware.go @@ -0,0 +1,46 @@ +package moduleregistry + +import "net/http" + +// RequireActiveModule weist Anfragen an ein nicht aktiviertes Modul ZENTRAL +// ab, bevor der eigentliche Modul-Handler erreicht wird (Akzeptanzkriterium 2 / +// Pruefung 1) — Casbin-Prinzip: Durchsetzung als Middleware statt verstreuter +// Pruefungen in jedem Handler. tenantSlug/moduleName werden hier ueber +// Query-Parameter gelesen (echte Extraktion aus JWT/Tenant-Kontext ist +// API-05/TEN-06, nicht Teil dieser Kachel). +func (r *Registry) RequireActiveModule(moduleName string, next http.HandlerFunc) http.HandlerFunc { + return func(w http.ResponseWriter, req *http.Request) { + tenantSlug := req.URL.Query().Get("tenant") + active, err := r.IsActive(req.Context(), tenantSlug, moduleName) + if err != nil { + http.Error(w, "aktivierungspruefung fehlgeschlagen", http.StatusInternalServerError) + return + } + if !active { + http.Error(w, "modul nicht aktiviert", http.StatusForbidden) + return + } + next(w, req) + } +} + +// RequireServiceCredential authentifiziert eine Modul-Instanz ueber ihr +// Service-Credential (X-Client-Id/X-Client-Secret-Header) BEVOR der +// eigentliche Handler erreicht wird (Akzeptanzkriterium 4 / Pruefung 4). +func (r *Registry) RequireServiceCredential(next http.HandlerFunc) http.HandlerFunc { + return func(w http.ResponseWriter, req *http.Request) { + clientID := req.Header.Get("X-Client-Id") + secret := req.Header.Get("X-Client-Secret") + + _, ok, err := r.Authenticate(req.Context(), clientID, secret) + if err != nil { + http.Error(w, "authentifizierung fehlgeschlagen", http.StatusInternalServerError) + return + } + if !ok { + http.Error(w, ErrInvalidCredential.Error(), http.StatusUnauthorized) + return + } + next(w, req) + } +} diff --git a/internal/moduleregistry/registry.go b/internal/moduleregistry/registry.go new file mode 100644 index 0000000..dbda148 --- /dev/null +++ b/internal/moduleregistry/registry.go @@ -0,0 +1,120 @@ +// Package moduleregistry implementiert Core API-02: die Registry, in der +// sich Fachmodule (DMS, Mail, weitere) mit Metadaten eintragen, gekoppelt an +// die Aktivierungspruefung aus LIC-02 (Feature-Flags). Zusaetzlich +// authentifiziert die Registry Modul-Instanzen selbst ueber ein bei +// Provisionierung ausgestelltes Service-Credential. +package moduleregistry + +import ( + "context" + "errors" + "fmt" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/flag" +) + +var ( + ErrMissingName = errors.New("moduleregistry: name darf nicht leer sein") + ErrMissingVersion = errors.New("moduleregistry: version darf nicht leer sein") + ErrModuleNotFound = errors.New("moduleregistry: modul nicht registriert") +) + +type Module struct { + Name string + Version string + RequiredFlags []string +} + +type Registry struct { + pool *pgxpool.Pool + flags *flag.Service +} + +func NewRegistry(pool *pgxpool.Pool, flags *flag.Service) *Registry { + return &Registry{pool: pool, flags: flags} +} + +// Register traegt ein Modul mit Name, Version und benoetigten Feature-Flags +// ein (Akzeptanzkriterium 1). Fehlende Pflichtangaben werden abgewiesen +// (Akzeptanzkriterium 1 / Pruefung 2). Erneutes Register desselben Namens +// aktualisiert Version/Flags (Redeploy-Fall). +func (r *Registry) Register(ctx context.Context, name, version string, requiredFlags []string) (Module, error) { + if name == "" { + return Module{}, ErrMissingName + } + if version == "" { + return Module{}, ErrMissingVersion + } + if requiredFlags == nil { + requiredFlags = []string{} + } + + _, err := r.pool.Exec(ctx, ` + INSERT INTO modules (name, version, required_flags, registered_at) + VALUES ($1, $2, $3, now()) + ON CONFLICT (name) DO UPDATE SET version = $2, required_flags = $3, registered_at = now() + `, name, version, requiredFlags) + if err != nil { + return Module{}, fmt.Errorf("modul registrieren: %w", err) + } + return Module{Name: name, Version: version, RequiredFlags: requiredFlags}, nil +} + +func (r *Registry) Get(ctx context.Context, name string) (Module, error) { + var m Module + m.Name = name + err := r.pool.QueryRow(ctx, ` + SELECT version, required_flags FROM modules WHERE name = $1 + `, name).Scan(&m.Version, &m.RequiredFlags) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return Module{}, ErrModuleNotFound + } + return Module{}, fmt.Errorf("modul lesen: %w", err) + } + return m, nil +} + +// List liefert alle registrierten Module (Akzeptanzkriterium 3: ueber API +// abfragbar, z.B. fuer Statusseite/Lizenzoberflaeche). +func (r *Registry) List(ctx context.Context) ([]Module, error) { + rows, err := r.pool.Query(ctx, `SELECT name, version, required_flags FROM modules ORDER BY name`) + if err != nil { + return nil, fmt.Errorf("module auflisten: %w", err) + } + defer rows.Close() + + var out []Module + for rows.Next() { + var m Module + if err := rows.Scan(&m.Name, &m.Version, &m.RequiredFlags); err != nil { + return nil, fmt.Errorf("modul lesen: %w", err) + } + out = append(out, m) + } + return out, rows.Err() +} + +// IsActive prueft, ob ein registriertes Modul fuer einen Tenant aktiviert +// ist: registriert UND alle benoetigten Feature-Flags sind fuer diesen +// Tenant aktiv (Akzeptanzkriterium 2). Ein nicht registriertes Modul gilt +// immer als nicht aktiv. +func (r *Registry) IsActive(ctx context.Context, tenantSlug, moduleName string) (bool, error) { + m, err := r.Get(ctx, moduleName) + if err != nil { + if errors.Is(err, ErrModuleNotFound) { + return false, nil + } + return false, err + } + + for _, flagKey := range m.RequiredFlags { + if !r.flags.IsEnabled(ctx, tenantSlug, flagKey) { + return false, nil + } + } + return true, nil +} diff --git a/internal/moduleregistry/registry_test.go b/internal/moduleregistry/registry_test.go new file mode 100644 index 0000000..6d602a2 --- /dev/null +++ b/internal/moduleregistry/registry_test.go @@ -0,0 +1,304 @@ +package moduleregistry + +import ( + "context" + "errors" + "fmt" + "net/http" + "net/http/httptest" + "os" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/flag" +) + +func setupTest(t *testing.T) (*Registry, *flag.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 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) + } + + flagStore := flag.NewStore(pool) + // Kurze TTL, damit Tests, die den Flag-Store direkt aendern (an + // Registry.IsActive vorbei), den neuen Stand ohne manuelles Invalidate + // zuverlaessig sehen. + flagService := flag.NewService(flagStore, 10*time.Millisecond) + registry := NewRegistry(pool, flagService) + + cleanup := func() { pool.Close() } + return registry, flagStore, cleanup +} + +func uniqueModuleName(t *testing.T) string { + return fmt.Sprintf("dms_%d", time.Now().UnixNano()) +} + +// Akzeptanzkriterium 1 + Pruefung 2: fehlende Pflichtangaben abgewiesen. +func TestRegister_RejectsMissingFields(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + if _, err := registry.Register(ctx, "", "1.0", nil); !errors.Is(err, ErrMissingName) { + t.Fatalf("erwartet ErrMissingName, habe %v", err) + } + if _, err := registry.Register(ctx, "dms", "", nil); !errors.Is(err, ErrMissingVersion) { + t.Fatalf("erwartet ErrMissingVersion, habe %v", err) + } +} + +func TestRegister_AndGet(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + name := uniqueModuleName(t) + + m, err := registry.Register(ctx, name, "1.2.0", []string{"dms_enabled"}) + if err != nil { + t.Fatalf("register: %v", err) + } + if m.Version != "1.2.0" || len(m.RequiredFlags) != 1 { + t.Fatalf("unerwartet: %+v", m) + } + + got, err := registry.Get(ctx, name) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.Version != "1.2.0" { + t.Fatalf("get version = %q", got.Version) + } +} + +// Akzeptanzkriterium 2 + 3 + Pruefung 3: konsistente Daten nach +// Aktivierung/Deaktivierung eines Moduls. +func TestIsActive_ReflectsFlagStateConsistently(t *testing.T) { + registry, flagStore, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + name := uniqueModuleName(t) + flagKey := name + "_enabled" + + if _, err := registry.Register(ctx, name, "1.0", []string{flagKey}); err != nil { + t.Fatalf("register: %v", err) + } + + active, err := registry.IsActive(ctx, "acme", name) + if err != nil { + t.Fatalf("is active (vor flag): %v", err) + } + if active { + t.Fatal("erwartet nicht aktiv, solange flag nicht gesetzt ist") + } + + if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: true}); err != nil { + t.Fatalf("flag setzen: %v", err) + } + time.Sleep(20 * time.Millisecond) // TTL abwarten + + active, err = registry.IsActive(ctx, "acme", name) + if err != nil { + t.Fatalf("is active (nach flag an): %v", err) + } + if !active { + t.Fatal("erwartet aktiv, nachdem flag aktiviert wurde") + } + + if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: false}); err != nil { + t.Fatalf("flag zuruecksetzen: %v", err) + } + time.Sleep(20 * time.Millisecond) // TTL abwarten + active, err = registry.IsActive(ctx, "acme", name) + if err != nil { + t.Fatalf("is active (nach flag aus): %v", err) + } + if active { + t.Fatal("erwartet wieder nicht aktiv, nachdem flag deaktiviert wurde") + } +} + +func TestIsActive_UnregisteredModuleIsNeverActive(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + active, err := registry.IsActive(ctx, "acme", "nie-registriert") + if err != nil { + t.Fatalf("is active: %v", err) + } + if active { + t.Fatal("unregistriertes modul darf nie aktiv sein") + } +} + +// Akzeptanzkriterium 2 + Pruefung 1: Anfrage an deaktiviertes Modul wird +// zentral abgewiesen, BEVOR die Modul-Logik erreicht wird. +func TestRequireActiveModule_BlocksBeforeHandler(t *testing.T) { + registry, flagStore, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + name := uniqueModuleName(t) + flagKey := name + "_enabled" + + if _, err := registry.Register(ctx, name, "1.0", []string{flagKey}); err != nil { + t.Fatalf("register: %v", err) + } + + handlerReached := false + handler := registry.RequireActiveModule(name, func(w http.ResponseWriter, r *http.Request) { + handlerReached = true + w.WriteHeader(http.StatusOK) + }) + + req := httptest.NewRequest(http.MethodGet, "/modul?tenant=acme", nil) + rec := httptest.NewRecorder() + handler(rec, req) + if rec.Code != http.StatusForbidden { + t.Fatalf("status = %d, want 403", rec.Code) + } + if handlerReached { + t.Fatal("handler haette bei deaktiviertem modul NICHT erreicht werden duerfen") + } + + if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: true}); err != nil { + t.Fatalf("flag setzen: %v", err) + } + req2 := httptest.NewRequest(http.MethodGet, "/modul?tenant=acme", nil) + rec2 := httptest.NewRecorder() + handler(rec2, req2) + if rec2.Code != http.StatusOK { + t.Fatalf("status nach aktivierung = %d, want 200", rec2.Code) + } + if !handlerReached { + t.Fatal("handler haette bei aktiviertem modul erreicht werden muessen") + } +} + +// Akzeptanzkriterium 4 + Pruefung 4: gueltiges/ungueltiges Service-Credential. +func TestProvisionAndAuthenticate(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + name := uniqueModuleName(t) + + if _, err := registry.Register(ctx, name, "1.0", nil); err != nil { + t.Fatalf("register: %v", err) + } + + clientID, secret, err := registry.Provision(ctx, name) + if err != nil { + t.Fatalf("provision: %v", err) + } + if clientID == "" || secret == "" { + t.Fatal("erwartet nicht-leere client-id/secret") + } + + moduleName, ok, err := registry.Authenticate(ctx, clientID, secret) + if err != nil { + t.Fatalf("authenticate (korrekt): %v", err) + } + if !ok || moduleName != name { + t.Fatalf("erwartet erfolgreiche authentifizierung fuer %q, habe ok=%v moduleName=%q", name, ok, moduleName) + } + + _, ok, err = registry.Authenticate(ctx, clientID, "falsches-secret") + if err != nil { + t.Fatalf("authenticate (falsch): %v", err) + } + if ok { + t.Fatal("erwartet fehlschlag bei falschem secret") + } + + _, ok, err = registry.Authenticate(ctx, "unbekannte-client-id", secret) + if err != nil { + t.Fatalf("authenticate (unbekannt): %v", err) + } + if ok { + t.Fatal("erwartet fehlschlag bei unbekannter client-id") + } +} + +func TestProvision_RequiresRegisteredModule(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + if _, _, err := registry.Provision(ctx, "nie-registriert"); !errors.Is(err, ErrModuleNotRegistered) { + t.Fatalf("erwartet ErrModuleNotRegistered, habe %v", err) + } +} + +func TestRequireServiceCredential_RejectsInvalidAcceptsValid(t *testing.T) { + registry, _, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + name := uniqueModuleName(t) + + if _, err := registry.Register(ctx, name, "1.0", nil); err != nil { + t.Fatalf("register: %v", err) + } + clientID, secret, err := registry.Provision(ctx, name) + if err != nil { + t.Fatalf("provision: %v", err) + } + + handler := registry.RequireServiceCredential(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + }) + + // Fehlendes Credential. + req := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil) + rec := httptest.NewRecorder() + handler(rec, req) + if rec.Code != http.StatusUnauthorized { + t.Fatalf("ohne credential: status = %d, want 401", rec.Code) + } + + // Falsches Secret. + req2 := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil) + req2.Header.Set("X-Client-Id", clientID) + req2.Header.Set("X-Client-Secret", "falsch") + rec2 := httptest.NewRecorder() + handler(rec2, req2) + if rec2.Code != http.StatusUnauthorized { + t.Fatalf("falsches secret: status = %d, want 401", rec2.Code) + } + + // Gueltiges Credential. + req3 := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil) + req3.Header.Set("X-Client-Id", clientID) + req3.Header.Set("X-Client-Secret", secret) + rec3 := httptest.NewRecorder() + handler(rec3, req3) + if rec3.Code != http.StatusOK { + t.Fatalf("gueltiges credential: status = %d, want 200", rec3.Code) + } +} diff --git a/internal/webhook/dispatcher.go b/internal/webhook/dispatcher.go new file mode 100644 index 0000000..cf03419 --- /dev/null +++ b/internal/webhook/dispatcher.go @@ -0,0 +1,177 @@ +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<= 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 +} diff --git a/internal/webhook/webhook.go b/internal/webhook/webhook.go new file mode 100644 index 0000000..f533bee --- /dev/null +++ b/internal/webhook/webhook.go @@ -0,0 +1,131 @@ +// 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" diff --git a/internal/webhook/webhook_test.go b/internal/webhook/webhook_test.go new file mode 100644 index 0000000..2f78674 --- /dev/null +++ b/internal/webhook/webhook_test.go @@ -0,0 +1,205 @@ +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") + } +} From 196b48ce09ad59b8eef68e2235194975db3b75db Mon Sep 17 00:00:00 2001 From: sysops Date: Fri, 28 Aug 2026 10:34:25 +0200 Subject: [PATCH 2/2] QA-07: schnittstellen-vertragstests fuer api-01/api-05/api-02/api-07 + gitea-actions-workflow --- .gitea/workflows/contract-tests.yml | 22 +++ internal/contracttest/contracttest.go | 57 ++++++ internal/contracttest/contracttest_test.go | 193 +++++++++++++++++++++ 3 files changed, 272 insertions(+) create mode 100644 .gitea/workflows/contract-tests.yml create mode 100644 internal/contracttest/contracttest.go create mode 100644 internal/contracttest/contracttest_test.go diff --git a/.gitea/workflows/contract-tests.yml b/.gitea/workflows/contract-tests.yml new file mode 100644 index 0000000..fecd52d --- /dev/null +++ b/.gitea/workflows/contract-tests.yml @@ -0,0 +1,22 @@ +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 diff --git a/internal/contracttest/contracttest.go b/internal/contracttest/contracttest.go new file mode 100644 index 0000000..f821644 --- /dev/null +++ b/internal/contracttest/contracttest.go @@ -0,0 +1,57 @@ +// 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 +} diff --git a/internal/contracttest/contracttest_test.go b/internal/contracttest/contracttest_test.go new file mode 100644 index 0000000..88ac943 --- /dev/null +++ b/internal/contracttest/contracttest_test.go @@ -0,0 +1,193 @@ +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) + } +}