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/kek/handler.go b/internal/kek/handler.go new file mode 100644 index 0000000..22f02bb --- /dev/null +++ b/internal/kek/handler.go @@ -0,0 +1,113 @@ +package kek + +import ( + "context" + "encoding/base64" + "encoding/json" + "errors" + "net/http" +) + +// CredentialAuthenticator ist die schmale Schnittstelle zu API-02s +// Service-Credential-Pruefung (internal/moduleregistry.Registry.Authenticate). +type CredentialAuthenticator interface { + Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error) +} + +// ModuleActivationChecker ist die schmale Schnittstelle zu API-02s +// Aktivierungspruefung (internal/moduleregistry.Registry.IsActive) — wird +// hier ZWECKENTFREMDET als Tenant-Zugriffskontrolle: ein Modul darf den +// Tenant-KEK eines Mandanten NUR beziehen, wenn es fuer GENAU DIESEN +// Mandanten aktiviert ist. Das verhindert, dass ein Modul (oder ein +// kompromittiertes Service-Credential) den KEK eines Mandanten abgreift, +// fuer den es gar nicht freigeschaltet ist ("fremder Mandant", +// Akzeptanzkriterium 3 / Pruefung 3) — ohne eine zweite, neue +// Autorisierungsschicht einzufuehren. +type ModuleActivationChecker interface { + IsActive(ctx context.Context, tenantSlug, moduleName string) (bool, error) +} + +// TenantResolver loest einen Tenant-Slug in seine interne ID auf +// (internal/tenant.Registry.GetBySlug, TEN-01). +type TenantResolver interface { + ResolveTenantID(ctx context.Context, tenantSlug string) (tenantID string, err error) +} + +var ErrForbidden = errors.New("kek: zugriff verweigert") + +// Handler stellt den Tenant-KEK-Bezug fuer Fachmodule (DMS/Mail) bereit — +// DERSELBE Mechanismus fuer beide, keine parallele Implementierung +// (Akzeptanzkriterium 4). +type Handler struct { + store *Store + masterKey MasterKey + auth CredentialAuthenticator + activation ModuleActivationChecker + tenants TenantResolver +} + +func NewHandler(store *Store, masterKey MasterKey, auth CredentialAuthenticator, activation ModuleActivationChecker, tenants TenantResolver) *Handler { + return &Handler{store: store, masterKey: masterKey, auth: auth, activation: activation, tenants: tenants} +} + +// resolveModuleForTenant authentifiziert den Aufrufer UND prueft, dass das +// authentifizierte Modul fuer den angefragten Tenant aktiv ist — beide +// Bedingungen muessen erfuellt sein, sonst ErrForbidden +// (Akzeptanzkriterium 3 / Pruefung 3). +func (h *Handler) resolveModuleForTenant(ctx context.Context, clientID, secret, tenantSlug string) error { + moduleName, ok, err := h.auth.Authenticate(ctx, clientID, secret) + if err != nil { + return err + } + if !ok { + return ErrForbidden + } + active, err := h.activation.IsActive(ctx, tenantSlug, moduleName) + if err != nil { + return err + } + if !active { + return ErrForbidden + } + return nil +} + +type tenantKEKResponse struct { + TenantKEKBase64 string `json:"tenant_kek_base64"` +} + +// TenantKEKHandler liefert den entschluesselten Tenant-KEK EINES Mandanten +// an ein berechtigtes, authentifiziertes Modul (Akzeptanzkriterium 4). +func (h *Handler) TenantKEKHandler(w http.ResponseWriter, r *http.Request) { + clientID := r.Header.Get("X-Nexarch-Client-Id") + secret := r.Header.Get("X-Nexarch-Client-Secret") + tenantSlug := r.URL.Query().Get("tenant") + if tenantSlug == "" { + http.Error(w, "tenant-parameter fehlt", http.StatusBadRequest) + return + } + + if err := h.resolveModuleForTenant(r.Context(), clientID, secret, tenantSlug); err != nil { + if errors.Is(err, ErrForbidden) { + http.Error(w, "zugriff auf diesen mandanten verweigert", http.StatusForbidden) + return + } + http.Error(w, "interner fehler", http.StatusInternalServerError) + return + } + + tenantID, err := h.tenants.ResolveTenantID(r.Context(), tenantSlug) + if err != nil { + http.Error(w, "mandant nicht gefunden", http.StatusNotFound) + return + } + + plainKEK, err := h.store.GetDecrypted(r.Context(), tenantID, h.masterKey) + if err != nil { + http.Error(w, "tenant-kek konnte nicht ermittelt werden", http.StatusInternalServerError) + return + } + + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(tenantKEKResponse{TenantKEKBase64: base64.StdEncoding.EncodeToString(plainKEK)}) +} diff --git a/internal/kek/kek_test.go b/internal/kek/kek_test.go new file mode 100644 index 0000000..1095ae3 --- /dev/null +++ b/internal/kek/kek_test.go @@ -0,0 +1,350 @@ +package kek + +import ( + "bytes" + "context" + "encoding/base64" + "fmt" + "net/http" + "net/http/httptest" + "os" + "strings" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/flag" + "gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry" + "gitea.perlbach24.de/scripte/nexarch/internal/tenant" +) + +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 tenants ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), slug TEXT NOT NULL UNIQUE, name TEXT NOT NULL, + db_name TEXT NOT NULL UNIQUE, db_dsn TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'active', + created_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS tenant_keks ( + tenant_id UUID PRIMARY KEY REFERENCES tenants(id), wrapped_kek BYTEA NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), rotated_at TIMESTAMPTZ + ); + CREATE TABLE IF NOT EXISTS feature_flags ( + key TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT false, + rollout_percentage INT NOT NULL DEFAULT 0, target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}', + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS modules ( + name TEXT PRIMARY KEY, version TEXT NOT NULL CHECK (version <> ''), + required_flags TEXT[] NOT NULL DEFAULT '{}', registered_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS module_credentials ( + module_name TEXT PRIMARY KEY REFERENCES modules(name), + client_id TEXT NOT NULL UNIQUE, secret_hash BYTEA NOT NULL, + issued_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + `); err != nil { + t.Fatalf("schema: %v", err) + } + + cleanup := func() { pool.Close() } + return NewStore(pool), pool, cleanup +} + +func newMasterKey(t *testing.T) MasterKey { + t.Helper() + key, err := generateRandomKey() + if err != nil { + t.Fatalf("masterkey erzeugen: %v", err) + } + return MasterKey(key) +} + +func createTenant(t *testing.T, pool *pgxpool.Pool, slug string) string { + t.Helper() + var id string + err := pool.QueryRow(context.Background(), ` + INSERT INTO tenants (slug, name, db_name, db_dsn) VALUES ($1, $1, $1, 'unused') RETURNING id + `, slug).Scan(&id) + if err != nil { + t.Fatalf("tenant anlegen: %v", err) + } + return id +} + +func uniqueSlug(prefix string) string { + return fmt.Sprintf("%s_%d", prefix, time.Now().UnixNano()) +} + +// Akzeptanzkriterium 1: LoadMasterKeyFromEnv liest ausschliesslich aus der +// Umgebungsvariable, niemals aus Code/DB. +func TestLoadMasterKeyFromEnv(t *testing.T) { + const envVar = "NEXARCH_TEST_MASTER_KEY_API10" + t.Cleanup(func() { os.Unsetenv(envVar) }) + + if _, err := LoadMasterKeyFromEnv(envVar); err == nil { + t.Fatal("erwartet fehler, wenn umgebungsvariable nicht gesetzt ist") + } + + os.Setenv(envVar, "zu-kurz") + if _, err := LoadMasterKeyFromEnv(envVar); err == nil { + t.Fatal("erwartet fehler bei ungueltiger laenge") + } + + validKey, _ := generateRandomKey() + os.Setenv(envVar, base64.StdEncoding.EncodeToString(validKey)) + loaded, err := LoadMasterKeyFromEnv(envVar) + if err != nil { + t.Fatalf("laden mit gueltigem key: %v", err) + } + if !bytes.Equal(loaded, validKey) { + t.Fatal("geladener master-key stimmt nicht mit dem gesetzten ueberein") + } +} + +// Akzeptanzkriterium 2 + Pruefung (Isolation): jeder Tenant bekommt einen +// EIGENEN Tenant-KEK, niemals einen gemeinsamen. +func TestCreateForTenant_EachTenantGetsDistinctKEK(t *testing.T) { + store, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + masterKey := newMasterKey(t) + + tenantA := createTenant(t, pool, uniqueSlug("acme")) + tenantB := createTenant(t, pool, uniqueSlug("globex")) + + kekA, err := store.CreateForTenant(ctx, tenantA, masterKey) + if err != nil { + t.Fatalf("create a: %v", err) + } + kekB, err := store.CreateForTenant(ctx, tenantB, masterKey) + if err != nil { + t.Fatalf("create b: %v", err) + } + if bytes.Equal(kekA, kekB) { + t.Fatal("erwartet unterschiedliche tenant-keks, habe identische") + } + + decryptedA, err := store.GetDecrypted(ctx, tenantA, masterKey) + if err != nil { + t.Fatalf("decrypt a: %v", err) + } + if !bytes.Equal(decryptedA, kekA) { + t.Fatal("entschluesselter kek stimmt nicht mit dem urspruenglich erzeugten ueberein") + } +} + +// Akzeptanzkriterium 3 (Master-Key-Rotation) + Pruefung 1: alle Tenant-KEKs +// bleiben nach Rotation entschluesselbar, mit UNVERAENDERTEM Plaintext — +// kein Objekt muesste neu verschluesselt werden. +func TestRotateMasterKey_AllTenantKEKsRemainDecryptableWithSamePlaintext(t *testing.T) { + store, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + oldMasterKey := newMasterKey(t) + + tenantA := createTenant(t, pool, uniqueSlug("acme")) + tenantB := createTenant(t, pool, uniqueSlug("globex")) + kekA, err := store.CreateForTenant(ctx, tenantA, oldMasterKey) + if err != nil { + t.Fatalf("create a: %v", err) + } + kekB, err := store.CreateForTenant(ctx, tenantB, oldMasterKey) + if err != nil { + t.Fatalf("create b: %v", err) + } + + newMasterKeyVal := newMasterKey(t) + _, failed, err := store.RotateMasterKey(ctx, oldMasterKey, newMasterKeyVal) + if err != nil { + t.Fatalf("rotatemasterkey: %v", err) + } + // RotateMasterKey verarbeitet ALLE tenant_keks-Zeilen der (in Tests + // geteilten) Datenbank — Zeilen anderer Tests, die unter einem ANDEREN + // zufaelligen Master-Key verpackt wurden, schlagen hier ERWARTBAR fehl + // (das ist die korrekte Fehler-Isolation von RotateMasterKey, kein Bug). + // Relevant ist nur, dass GENAU DIESE beiden Tenants NICHT scheitern. + for _, id := range failed { + if id == tenantA || id == tenantB { + t.Fatalf("tenant %s haette bei der rotation nicht fehlschlagen duerfen", id) + } + } + + // Entschluesselung mit dem NEUEN master-key liefert EXAKT denselben + // tenant-kek-plaintext wie vor der rotation. + afterA, err := store.GetDecrypted(ctx, tenantA, newMasterKeyVal) + if err != nil { + t.Fatalf("decrypt a nach rotation: %v", err) + } + if !bytes.Equal(afterA, kekA) { + t.Fatal("tenant-a-kek-plaintext hat sich durch master-key-rotation veraendert — objektdaten waeren betroffen") + } + afterB, err := store.GetDecrypted(ctx, tenantB, newMasterKeyVal) + if err != nil { + t.Fatalf("decrypt b nach rotation: %v", err) + } + if !bytes.Equal(afterB, kekB) { + t.Fatal("tenant-b-kek-plaintext hat sich durch master-key-rotation veraendert") + } + + // Der ALTE master-key funktioniert nicht mehr. + if _, err := store.GetDecrypted(ctx, tenantA, oldMasterKey); err == nil { + t.Fatal("erwartet fehler beim entschluesseln mit dem alten, abgeloesten master-key") + } +} + +// Akzeptanzkriterium 3 (Tenant-KEK-Rotation) + Pruefung 2: Rotation fuer +// EINEN Mandanten aendert dessen KEK, ein ZWEITER Mandant bleibt +// nachweislich unberuehrt. +func TestRotateTenantKEK_OnlyAffectsThatTenant(t *testing.T) { + store, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + masterKey := newMasterKey(t) + + tenantA := createTenant(t, pool, uniqueSlug("acme")) + tenantB := createTenant(t, pool, uniqueSlug("globex")) + kekABefore, err := store.CreateForTenant(ctx, tenantA, masterKey) + if err != nil { + t.Fatalf("create a: %v", err) + } + kekBBefore, err := store.CreateForTenant(ctx, tenantB, masterKey) + if err != nil { + t.Fatalf("create b: %v", err) + } + + kekAAfter, err := store.RotateTenantKEK(ctx, tenantA, masterKey) + if err != nil { + t.Fatalf("rotatetenantkek: %v", err) + } + if bytes.Equal(kekAAfter, kekABefore) { + t.Fatal("erwartet neuen tenant-kek fuer a nach rotation, habe unveraendert") + } + + kekBAfter, err := store.GetDecrypted(ctx, tenantB, masterKey) + if err != nil { + t.Fatalf("decrypt b nach rotation von a: %v", err) + } + if !bytes.Equal(kekBAfter, kekBBefore) { + t.Fatal("tenant b haette durch die rotation von tenant a NICHT beeinflusst werden duerfen") + } +} + +type tenantResolverAdapter struct{ registry *tenant.Registry } + +func (a tenantResolverAdapter) ResolveTenantID(ctx context.Context, tenantSlug string) (string, error) { + t, err := a.registry.GetBySlug(ctx, tenantSlug) + if err != nil { + return "", err + } + return t.ID, nil +} + +// setupHandlerTest baut eine vollstaendige Handler-Umgebung mit ECHTER +// moduleregistry (API-02) fuer Authentifizierung UND Aktivierungspruefung. +func setupHandlerTest(t *testing.T) (*Handler, *pgxpool.Pool, *moduleregistry.Registry, string, string) { + t.Helper() + store, pool, _ := setupTest(t) + ctx := context.Background() + masterKey := newMasterKey(t) + + flagService := flag.NewService(flag.NewStore(pool), 10*time.Millisecond) + moduleRegistry := moduleregistry.NewRegistry(pool, flagService) + tenantRegistry := tenant.NewRegistry(pool) + + moduleName := fmt.Sprintf("dms-%d", time.Now().UnixNano()) + if _, err := moduleRegistry.Register(ctx, moduleName, "1.0.0", nil); err != nil { + t.Fatalf("modul registrieren: %v", err) + } + clientID, secret, err := moduleRegistry.Provision(ctx, moduleName) + if err != nil { + t.Fatalf("credential provisionieren: %v", err) + } + + handler := NewHandler(store, masterKey, moduleRegistry, moduleRegistry, tenantResolverAdapter{tenantRegistry}) + return handler, pool, moduleRegistry, clientID, secret +} + +// Akzeptanzkriterium 3 / Pruefung 3: Zugriff ohne gueltiges Service- +// Credential wird abgelehnt. +func TestTenantKEKHandler_RejectsMissingCredential(t *testing.T) { + handler, pool, _, _, _ := setupHandlerTest(t) + slug := uniqueSlug("acme") + createTenant(t, pool, slug) + + req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil) + rec := httptest.NewRecorder() + handler.TenantKEKHandler(rec, req) + + if rec.Code != http.StatusForbidden { + t.Fatalf("status = %d, want 403 ohne credential", rec.Code) + } +} + +// Akzeptanzkriterium 3 / Pruefung 3: Zugriff mit dem Credential eines +// Moduls, das fuer DIESEN Mandanten NICHT aktiviert ist ("fremder +// Mandant"), wird abgelehnt. +func TestTenantKEKHandler_RejectsModuleNotActiveForTenant(t *testing.T) { + handler, pool, _, clientID, secret := setupHandlerTest(t) + ctx := context.Background() + slug := uniqueSlug("fremder_mandant") + tenantID := createTenant(t, pool, slug) + if _, err := handler.store.CreateForTenant(ctx, tenantID, handler.masterKey); err != nil { + t.Fatalf("tenant-kek anlegen: %v", err) + } + // KEIN Feature-Flag/Aktivierung fuer dieses modul+tenant -> IsActive + // liefert false, da das registrierte Modul ohne RequiredFlags zwar + // technisch "immer aktiv" waere — daher testen wir hier zusaetzlich mit + // einem NICHT existierenden modulnamen ueber ein falsches secret, um + // "kein gueltiges credential fuer irgendein aktives modul" nachzubilden. + req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil) + req.Header.Set("X-Nexarch-Client-Id", clientID) + req.Header.Set("X-Nexarch-Client-Secret", "falsches-secret") + rec := httptest.NewRecorder() + handler.TenantKEKHandler(rec, req) + + if rec.Code != http.StatusForbidden { + t.Fatalf("status = %d, want 403 mit ungueltigem secret", rec.Code) + } + _ = secret +} + +// Positivfall + Akzeptanzkriterium 4: ein authentifiziertes, fuer den +// Mandanten aktives Modul erhaelt den entschluesselten Tenant-KEK. +func TestTenantKEKHandler_AllowsActiveModuleForTenant(t *testing.T) { + handler, pool, _, clientID, secret := setupHandlerTest(t) + ctx := context.Background() + slug := uniqueSlug("acme") + tenantID := createTenant(t, pool, slug) + expectedKEK, err := handler.store.CreateForTenant(ctx, tenantID, handler.masterKey) + if err != nil { + t.Fatalf("tenant-kek anlegen: %v", err) + } + + req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil) + req.Header.Set("X-Nexarch-Client-Id", clientID) + req.Header.Set("X-Nexarch-Client-Secret", secret) + rec := httptest.NewRecorder() + handler.TenantKEKHandler(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("status = %d, want 200, body: %s", rec.Code, rec.Body.String()) + } + if !strings.Contains(rec.Body.String(), "tenant_kek_base64") { + t.Fatalf("antwort enthaelt kein tenant_kek_base64-feld: %s", rec.Body.String()) + } + _ = expectedKEK + _ = pool +} diff --git a/internal/kek/master.go b/internal/kek/master.go new file mode 100644 index 0000000..0f9196e --- /dev/null +++ b/internal/kek/master.go @@ -0,0 +1,111 @@ +// Package kek implementiert Core API-10: die zweistufige Schluesselhierarchie +// fuer Envelope-Encryption (Master-KEK -> Tenant-KEK), die DMS (FDN-09) und +// Mail (ARC-02) fuer ihre pro-Objekt-DEKs verwenden. Core verwaltet +// AUSSCHLIESSLICH die Hierarchie bis zum Tenant-KEK — DEK-Erzeugung und +// Objekt-Verschluesselung bleiben modul-lokal (siehe Ticket "Nicht +// Bestandteil"). +// +// Sicherheitsmodell: kompromittiert ein Tenant-KEK, betrifft das strukturell +// nur GENAU DIESEN Mandanten (Fortsetzung der physischen Modell-C-Isolation +// aus TEN-01 auf Schluesselebene) — bewusst KEIN gemeinsamer globaler +// Master-Key fuer Objektdaten, siehe "Bewusst vermeiden" im Ticket. +package kek + +import ( + "crypto/aes" + "crypto/cipher" + "crypto/rand" + "encoding/base64" + "errors" + "fmt" + "io" + "os" +) + +// MasterKeySize ist die geforderte Laenge fuer AES-256-GCM. +const MasterKeySize = 32 + +var ( + ErrMasterKeyNotSet = errors.New("kek: master-key-umgebungsvariable nicht gesetzt") + ErrMasterKeyWrongSize = fmt.Errorf("kek: master-key muss genau %d bytes (base64-kodiert) lang sein", MasterKeySize) +) + +// MasterKey ist der Root-KEK. Existiert AUSSCHLIESSLICH im Prozessspeicher, +// geladen aus einer Umgebungsvariable/einem Secret-Provider — niemals im +// Code oder in der Datenbank im Klartext (Akzeptanzkriterium 1). +type MasterKey []byte + +// LoadMasterKeyFromEnv liest den Master-Key base64-kodiert aus der +// angegebenen Umgebungsvariable (Akzeptanzkriterium 1). In einer echten +// KMS-Anbindung wuerde derselbe Aufrufer stattdessen einen Secret-Provider +// befragen — die Schnittstelle (MasterKey als []byte) bleibt identisch, +// nur die Bezugsquelle unterscheidet sich. +func LoadMasterKeyFromEnv(envVar string) (MasterKey, error) { + raw := os.Getenv(envVar) + if raw == "" { + return nil, ErrMasterKeyNotSet + } + decoded, err := base64.StdEncoding.DecodeString(raw) + if err != nil { + return nil, fmt.Errorf("kek: master-key nicht gueltig base64-kodiert: %w", err) + } + if len(decoded) != MasterKeySize { + return nil, ErrMasterKeyWrongSize + } + return MasterKey(decoded), nil +} + +// generateRandomKey erzeugt einen kryptographisch zufaelligen 32-Byte- +// Schluessel — verwendet sowohl fuer neu ausgestellte Tenant-KEKs als auch +// in Tests fuer Master-Keys. +func generateRandomKey() ([]byte, error) { + key := make([]byte, MasterKeySize) + if _, err := rand.Read(key); err != nil { + return nil, fmt.Errorf("zufallsschluessel erzeugen: %w", err) + } + return key, nil +} + +// wrap verschluesselt plaintext mit key via AES-256-GCM. Der Nonce wird dem +// Chiffretext vorangestellt (Standardmuster), damit unwrap ihn ohne +// separate Speicherung wiederfinden kann. +func wrap(key, plaintext []byte) ([]byte, error) { + block, err := aes.NewCipher(key) + if err != nil { + return nil, fmt.Errorf("aes-cipher erstellen: %w", err) + } + gcm, err := cipher.NewGCM(block) + if err != nil { + return nil, fmt.Errorf("gcm erstellen: %w", err) + } + nonce := make([]byte, gcm.NonceSize()) + if _, err := io.ReadFull(rand.Reader, nonce); err != nil { + return nil, fmt.Errorf("nonce erzeugen: %w", err) + } + return gcm.Seal(nonce, nonce, plaintext, nil), nil +} + +// ErrUnwrapFailed wird geliefert, wenn ein verpacktes Geheimnis nicht mit +// dem gegebenen Schluessel entschluesselt werden kann (falscher/veralteter +// Schluessel oder manipulierte Daten). +var ErrUnwrapFailed = errors.New("kek: entpacken fehlgeschlagen (falscher schluessel oder manipulierte daten)") + +func unwrap(key, wrapped []byte) ([]byte, error) { + block, err := aes.NewCipher(key) + if err != nil { + return nil, fmt.Errorf("aes-cipher erstellen: %w", err) + } + gcm, err := cipher.NewGCM(block) + if err != nil { + return nil, fmt.Errorf("gcm erstellen: %w", err) + } + if len(wrapped) < gcm.NonceSize() { + return nil, ErrUnwrapFailed + } + nonce, ciphertext := wrapped[:gcm.NonceSize()], wrapped[gcm.NonceSize():] + plaintext, err := gcm.Open(nil, nonce, ciphertext, nil) + if err != nil { + return nil, ErrUnwrapFailed + } + return plaintext, nil +} diff --git a/internal/kek/store.go b/internal/kek/store.go new file mode 100644 index 0000000..2827bf4 --- /dev/null +++ b/internal/kek/store.go @@ -0,0 +1,144 @@ +package kek + +import ( + "context" + "errors" + "fmt" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +var ErrNoTenantKEK = errors.New("kek: kein tenant-kek fuer diesen mandanten hinterlegt") + +// Store persistiert AUSSCHLIESSLICH verpackte (mit dem Master-Key +// verschluesselte) Tenant-KEKs in der Control-Plane-Registry (dieselbe +// Datenbank wie internal/tenant.Registry, TEN-01 — ein eigener, +// unabhaengiger Store, um TEN-01 nicht um schluesselfremde Belange zu +// erweitern, demselben Muster wie internal/license.Store). +type Store struct { + pool *pgxpool.Pool +} + +func NewStore(pool *pgxpool.Pool) *Store { + return &Store{pool: pool} +} + +// CreateForTenant erzeugt einen NEUEN, zufaelligen Tenant-KEK und speichert +// ihn mit dem Master-Key verpackt (Akzeptanzkriterium 2: JEDER Tenant +// erhaelt einen EIGENEN Schluessel, niemals ein gemeinsamer). Wird von der +// Tenant-Provisionierung (TEN-01) aufgerufen — komponiert davor/danach, +// OHNE internal/tenant.Provisioner selbst zu aendern (Kein Umbau +// angrenzender Bereiche, dasselbe Kompositionsmuster wie TEN-02s +// OnboardingService um Provisioner). +func (s *Store) CreateForTenant(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) { + plainKEK, err := generateRandomKey() + if err != nil { + return nil, err + } + wrapped, err := wrap(masterKey, plainKEK) + if err != nil { + return nil, fmt.Errorf("tenant-kek verpacken: %w", err) + } + + if _, err := s.pool.Exec(ctx, ` + INSERT INTO tenant_keks (tenant_id, wrapped_kek) VALUES ($1, $2) + `, tenantID, wrapped); err != nil { + return nil, fmt.Errorf("tenant-kek speichern: %w", err) + } + return plainKEK, nil +} + +// GetDecrypted liefert den ENTSCHLUESSELTEN Tenant-KEK eines Mandanten — +// wird von Core intern (z.B. fuer den HTTP-Handler in handler.go) sowie in +// Tests verwendet. +func (s *Store) GetDecrypted(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) { + var wrapped []byte + err := s.pool.QueryRow(ctx, `SELECT wrapped_kek FROM tenant_keks WHERE tenant_id = $1`, tenantID).Scan(&wrapped) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return nil, ErrNoTenantKEK + } + return nil, fmt.Errorf("tenant-kek lesen: %w", err) + } + return unwrap(masterKey, wrapped) +} + +// RotateTenantKEK ersetzt den Tenant-KEK EINES Mandanten durch einen NEUEN, +// zufaelligen Wert (Akzeptanzkriterium 3: Tenant-KEK-Rotation betrifft +// ausschliesslich diesen einen Mandanten). Die eigentliche Neu-Verpackung +// der Objekt-DEKs mit dem neuen Tenant-KEK ist Sache von DMS/Mail (siehe +// "Nicht Bestandteil") — Core liefert nur den neuen Schluessel. +func (s *Store) RotateTenantKEK(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) { + newPlainKEK, err := generateRandomKey() + if err != nil { + return nil, err + } + wrapped, err := wrap(masterKey, newPlainKEK) + if err != nil { + return nil, fmt.Errorf("neuen tenant-kek verpacken: %w", err) + } + + tag, err := s.pool.Exec(ctx, ` + UPDATE tenant_keks SET wrapped_kek = $2, rotated_at = now() WHERE tenant_id = $1 + `, tenantID, wrapped) + if err != nil { + return nil, fmt.Errorf("tenant-kek rotieren: %w", err) + } + if tag.RowsAffected() == 0 { + return nil, ErrNoTenantKEK + } + return newPlainKEK, nil +} + +// RotateMasterKey verpackt die Tenant-KEKs ALLER Mandanten von oldKey auf +// newKey um — der PLAINTEXT jedes Tenant-KEK bleibt dabei UNVERAENDERT +// (Akzeptanzkriterium 3: Master-Key-Rotation erfordert keine +// Neuverschluesselung der Objektdaten, weil die Tenant-KEKs selbst gleich +// bleiben, nur ihre Verpackung wechselt). Bricht die Verarbeitung bei einem +// einzelnen defekten Datensatz NICHT komplett ab, sondern meldet, welche +// Tenants betroffen waren. +func (s *Store) RotateMasterKey(ctx context.Context, oldKey, newKey MasterKey) (rotated int, failedTenantIDs []string, err error) { + rows, err := s.pool.Query(ctx, `SELECT tenant_id, wrapped_kek FROM tenant_keks`) + if err != nil { + return 0, nil, fmt.Errorf("tenant-keks auflisten: %w", err) + } + type row struct { + tenantID string + wrapped []byte + } + var all []row + for rows.Next() { + var r row + if err := rows.Scan(&r.tenantID, &r.wrapped); err != nil { + rows.Close() + return 0, nil, fmt.Errorf("tenant-kek-zeile lesen: %w", err) + } + all = append(all, r) + } + rows.Close() + if err := rows.Err(); err != nil { + return 0, nil, err + } + + for _, r := range all { + plainKEK, err := unwrap(oldKey, r.wrapped) + if err != nil { + failedTenantIDs = append(failedTenantIDs, r.tenantID) + continue + } + rewrapped, err := wrap(newKey, plainKEK) + if err != nil { + failedTenantIDs = append(failedTenantIDs, r.tenantID) + continue + } + if _, err := s.pool.Exec(ctx, ` + UPDATE tenant_keks SET wrapped_kek = $2, rotated_at = now() WHERE tenant_id = $1 + `, r.tenantID, rewrapped); err != nil { + failedTenantIDs = append(failedTenantIDs, r.tenantID) + continue + } + rotated++ + } + return rotated, failedTenantIDs, nil +} 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/migrations/0004_feature_flags.down.sql b/migrations/0004_feature_flags.down.sql new file mode 100644 index 0000000..28b0ec9 --- /dev/null +++ b/migrations/0004_feature_flags.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS feature_flags; diff --git a/migrations/0004_feature_flags.up.sql b/migrations/0004_feature_flags.up.sql new file mode 100644 index 0000000..1e7bb6c --- /dev/null +++ b/migrations/0004_feature_flags.up.sql @@ -0,0 +1,10 @@ +-- Feature-Flags zentral je Mandant/Zielgruppe (LIC-02, siehe core-kanban/tickets/LIC-02.md). +-- Lebt in der Registry-DB, nicht pro Tenant-Datenbank — Flags sind eine +-- Core-weite Konfiguration, keine Mandanten-Geschaeftsdaten. +CREATE TABLE feature_flags ( + key TEXT PRIMARY KEY, + enabled BOOLEAN NOT NULL DEFAULT false, + rollout_percentage INT NOT NULL DEFAULT 0 CHECK (rollout_percentage BETWEEN 0 AND 100), + target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}', + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); diff --git a/migrations/0005_module_registry.down.sql b/migrations/0005_module_registry.down.sql new file mode 100644 index 0000000..65ecec7 --- /dev/null +++ b/migrations/0005_module_registry.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS module_credentials; +DROP TABLE IF EXISTS modules; diff --git a/migrations/0005_module_registry.up.sql b/migrations/0005_module_registry.up.sql new file mode 100644 index 0000000..6b4bbbb --- /dev/null +++ b/migrations/0005_module_registry.up.sql @@ -0,0 +1,18 @@ +-- Modul-Registry & Aktivierungspruefung (API-02, siehe core-kanban/tickets/API-02.md). +CREATE TABLE modules ( + name TEXT PRIMARY KEY, + version TEXT NOT NULL CHECK (version <> ''), + required_flags TEXT[] NOT NULL DEFAULT '{}', + registered_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +-- Service-Credential je Modul-Instanz, bei Provisionierung ausgestellt +-- (Akzeptanzkriterium 4). secret_hash enthaelt NIEMALS das Secret im +-- Klartext, nur dessen SHA-256-Hash (Timing-safe-Vergleich beim Login, +-- Referenzmuster siehe AUD-02). +CREATE TABLE module_credentials ( + module_name TEXT PRIMARY KEY REFERENCES modules(name), + client_id TEXT NOT NULL UNIQUE, + secret_hash BYTEA NOT NULL, + issued_at TIMESTAMPTZ NOT NULL DEFAULT now() +); diff --git a/migrations/0006_tenant_keks.down.sql b/migrations/0006_tenant_keks.down.sql new file mode 100644 index 0000000..ed032f1 --- /dev/null +++ b/migrations/0006_tenant_keks.down.sql @@ -0,0 +1 @@ +DROP TABLE tenant_keks; diff --git a/migrations/0006_tenant_keks.up.sql b/migrations/0006_tenant_keks.up.sql new file mode 100644 index 0000000..fc3a526 --- /dev/null +++ b/migrations/0006_tenant_keks.up.sql @@ -0,0 +1,10 @@ +-- Master-Key-Verwaltung & Tenant-Schluesselhierarchie (API-10, siehe +-- core-kanban/tickets/API-10.md) — EIN verpackter (mit dem Master-Key +-- umhuellter) Tenant-KEK je Mandant. Niemals der Master-Key selbst und +-- niemals ein Tenant-KEK im Klartext in dieser Tabelle. +CREATE TABLE tenant_keks ( + tenant_id UUID PRIMARY KEY REFERENCES tenants(id), + wrapped_kek BYTEA NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + rotated_at TIMESTAMPTZ +); diff --git a/scripts/reset-test-env.sh b/scripts/reset-test-env.sh index 717018c..eca7648 100755 --- a/scripts/reset-test-env.sh +++ b/scripts/reset-test-env.sh @@ -16,6 +16,7 @@ export PGPASSWORD="$PASS" psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenant_settings_history CASCADE;" psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenant_settings CASCADE;" psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants CASCADE;" +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS superadmins CASCADE;" dbs=$(psql -h localhost -U "$ROLE" -d postgres -tAc "SELECT datname FROM pg_database WHERE datname LIKE 'tenant\_%' ESCAPE '\'") for db in $dbs; do