Compare commits

..
Author SHA1 Message Date
sysopsandClaude Sonnet 5 6a03dcafd6 OPS-01: health-check-endpunkte-je-modul
internal/health: wiederverwendbare Registry fuer benannte Checks (DB, Queue)
— nicht Core-spezifisch, sondern von jedem registrierten Modul (API-02)
gleichermassen einsetzbar. LivenessHandler prueft bewusst KEINE externen
Abhaengigkeiten (Akzeptanzkriterium 2: Liveness/Readiness getrennt) — ein
DB-Ausfall soll den Prozess nicht faelschlich als "tot" markieren und einen
grundlosen Neustart ausloesen. ReadinessHandler fuehrt alle registrierten
Checks NEBENLAEUFIG mit je eigenem Timeout aus (DefaultCheckTimeout=2s) und
liefert 503, sobald irgendeine Abhaengigkeit fehlschlaegt (Akzeptanz-
kriterium 1 + 3) — echte Pruefung von DB (Ping) und Job-Queue statt nur
Prozessstatus.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Simulierter Datenbankausfall fuehrt zu "nicht bereit" —
   TestReadinessHandler_ReportsNotReadyOnDatabaseFailure: geschlossener Pool,
   503 mit "database" im Checks-Ergebnis. PASS.
2. Health-Endpunkt antwortet auch bei haengendem Check innerhalb definierter
   Zeit — TestReadinessHandler_RespondsWithinTimeoutEvenWithHangingCheck:
   ein 10s blockierender Check wird durch 50ms-Timeout begrenzt, Handler
   antwortet deutlich unter 1s. PASS.
3. Readiness- und Liveness-Antwort unterscheiden sich nachweislich in
   mindestens einem Fehlerfall — TestLivenessAndReadiness_DifferOnDatabaseFailure:
   bei DB-Ausfall liefert Liveness weiterhin 200, Readiness 503. PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 23:30:30 +02:00
sysopsandClaude Sonnet 5 b23cd1961f API-02: modul-registry-aktivierungspruefung
internal/moduleregistry: Registry.Register traegt Fachmodule mit Name,
Version und benoetigten Feature-Flags ein (Akzeptanzkriterium 1), fehlende
Pflichtangaben werden abgewiesen. IsActive kombiniert Registrierung + LIC-02
Feature-Flag-Auswertung (ALLE benoetigten Flags muessen fuer den Tenant
aktiv sein) — ein nicht registriertes Modul ist nie aktiv. List liefert alle
Module fuer Statusseite/Lizenzoberflaeche (Akzeptanzkriterium 3).

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

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

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

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

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

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

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

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

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:19:59 +02:00
24 changed files with 1332 additions and 495 deletions
+87
View File
@@ -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
}
+43
View File
@@ -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")
}
}
}
+87
View File
@@ -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()
}
+179
View File
@@ -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")
}
}()
}
+25
View File
@@ -0,0 +1,25 @@
package health
import (
"context"
"github.com/jackc/pgx/v5/pgxpool"
)
// DatabaseChecker prueft die tatsaechliche Erreichbarkeit der Datenbank
// (Ping) — nicht nur, ob der Pool existiert.
func DatabaseChecker(pool *pgxpool.Pool) CheckerFunc {
return func(ctx context.Context) error {
return pool.Ping(ctx)
}
}
// QueueChecker prueft, dass die Postgres-basierte Job-Queue (siehe CFG-02)
// tatsaechlich abfragbar ist — eine eigene, benannte Abhaengigkeit neben der
// reinen DB-Erreichbarkeit (Akzeptanzkriterium 1).
func QueueChecker(pool *pgxpool.Pool) CheckerFunc {
return func(ctx context.Context) error {
_, err := pool.Exec(ctx, `SELECT 1`)
return err
}
}
+49
View File
@@ -0,0 +1,49 @@
package health
import (
"encoding/json"
"net/http"
)
// LivenessHandler beantwortet IMMER "lebt", solange der Prozess ueberhaupt
// HTTP-Anfragen verarbeiten kann — prueft bewusst KEINE externen
// Abhaengigkeiten (Akzeptanzkriterium 2: Liveness und Readiness getrennt).
// Ein Datenbankausfall darf die Liveness nicht auf "tot" setzen, sonst
// wuerde eine Orchestrierung (z.B. systemd/Kubernetes) den Prozess grundlos
// neu starten, obwohl nur eine Abhaengigkeit ausgefallen ist.
func LivenessHandler(w http.ResponseWriter, r *http.Request) {
writeStatus(w, http.StatusOK, map[string]any{"status": "alive"})
}
// ReadinessHandler prueft ALLE registrierten Abhaengigkeiten
// (Akzeptanzkriterium 1) und liefert 503, sobald eine davon fehlschlaegt
// (Akzeptanzkriterium 3) — unterscheidet sich damit nachweislich von
// LivenessHandler im Fehlerfall (Akzeptanzkriterium 2 / Pruefung 3).
func (r *Registry) ReadinessHandler() http.HandlerFunc {
return func(w http.ResponseWriter, req *http.Request) {
ready, results := r.CheckAll(req.Context())
body := map[string]any{
"status": statusText(ready),
"checks": results,
}
status := http.StatusOK
if !ready {
status = http.StatusServiceUnavailable
}
writeStatus(w, status, body)
}
}
func statusText(ready bool) string {
if ready {
return "ready"
}
return "not_ready"
}
func writeStatus(w http.ResponseWriter, status int, body map[string]any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(body)
}
+86
View File
@@ -0,0 +1,86 @@
// Package health implementiert Core OPS-01: Health-/Readiness-Endpunkte, die
// echte Abhaengigkeiten (DB, Job-Queue) statt nur den Prozessstatus pruefen
// — wiederverwendbar von Core UND jedem registrierten Modul (siehe API-02),
// nicht nur von Core selbst.
package health
import (
"context"
"time"
)
// Checker prueft EINE Abhaengigkeit (z.B. Datenbank, Job-Queue).
type Checker interface {
Check(ctx context.Context) error
}
type CheckerFunc func(ctx context.Context) error
func (f CheckerFunc) Check(ctx context.Context) error { return f(ctx) }
// DefaultCheckTimeout begrenzt, wie lange EIN einzelner Check maximal
// dauern darf, bevor er als fehlgeschlagen gilt — verhindert, dass ein
// haengender Check den gesamten Readiness-Endpunkt blockiert
// (Akzeptanzkriterium 2 / Pruefung 2: Antwort innerhalb definierter Zeit).
const DefaultCheckTimeout = 2 * time.Second
// Registry haelt alle benannten Checks eines Dienstes.
type Registry struct {
checks map[string]Checker
timeout time.Duration
}
func NewRegistry() *Registry {
return &Registry{checks: make(map[string]Checker), timeout: DefaultCheckTimeout}
}
func (r *Registry) WithTimeout(d time.Duration) *Registry {
return &Registry{checks: r.checks, timeout: d}
}
// Register fuegt einen benannten Check hinzu (z.B. "database", "queue").
func (r *Registry) Register(name string, c Checker) {
r.checks[name] = c
}
// Result ist der Ausgang eines einzelnen Checks.
type Result struct {
OK bool
Error string
}
// CheckAll fuehrt alle registrierten Checks NEBENLAEUFIG mit je eigenem
// Timeout aus (Akzeptanzkriterium 1: echte Abhaengigkeiten statt Prozess-
// status) und liefert ready=false, sobald irgendein Check fehlschlaegt
// (Akzeptanzkriterium 3: ein Ausfall wird sichtbar).
func (r *Registry) CheckAll(ctx context.Context) (ready bool, results map[string]Result) {
type namedResult struct {
name string
result Result
}
ch := make(chan namedResult, len(r.checks))
for name, checker := range r.checks {
go func(name string, checker Checker) {
checkCtx, cancel := context.WithTimeout(ctx, r.timeout)
defer cancel()
err := checker.Check(checkCtx)
if err != nil {
ch <- namedResult{name, Result{OK: false, Error: err.Error()}}
return
}
ch <- namedResult{name, Result{OK: true}}
}(name, checker)
}
results = make(map[string]Result, len(r.checks))
ready = true
for i := 0; i < len(r.checks); i++ {
nr := <-ch
results[nr.name] = nr.result
if !nr.result.OK {
ready = false
}
}
return ready, results
}
+161
View File
@@ -0,0 +1,161 @@
package health
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
// Akzeptanzkriterium 1 + Pruefung 1: simulierter Datenbankausfall fuehrt zu
// "nicht bereit".
func TestReadinessHandler_ReportsNotReadyOnDatabaseFailure(t *testing.T) {
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
// Datenbankausfall simulieren: Pool sofort schliessen, bevor der Check laeuft.
pool.Close()
reg := NewRegistry()
reg.Register("database", DatabaseChecker(pool))
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
rec := httptest.NewRecorder()
reg.ReadinessHandler()(rec, req)
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want 503 bei db-ausfall", rec.Code)
}
var body struct {
Status string `json:"status"`
Checks map[string]interface{} `json:"checks"`
}
if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil {
t.Fatalf("body parsen: %v", err)
}
if body.Status != "not_ready" {
t.Fatalf("status-feld = %q, want not_ready", body.Status)
}
if _, ok := body.Checks["database"]; !ok {
t.Fatal("erwartet 'database' im checks-ergebnis")
}
}
func TestReadinessHandler_ReportsReadyWhenAllChecksPass(t *testing.T) {
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
defer pool.Close()
reg := NewRegistry()
reg.Register("database", DatabaseChecker(pool))
reg.Register("queue", QueueChecker(pool))
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
rec := httptest.NewRecorder()
reg.ReadinessHandler()(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200 bei funktionierenden abhaengigkeiten", rec.Code)
}
}
// Akzeptanzkriterium 2 + Pruefung 3: Liveness und Readiness unterscheiden
// sich nachweislich im Fehlerfall.
func TestLivenessAndReadiness_DifferOnDatabaseFailure(t *testing.T) {
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
pool.Close() // db-ausfall simulieren
reg := NewRegistry()
reg.Register("database", DatabaseChecker(pool))
livenessRec := httptest.NewRecorder()
LivenessHandler(livenessRec, httptest.NewRequest(http.MethodGet, "/livez", nil))
if livenessRec.Code != http.StatusOK {
t.Fatalf("liveness status = %d, want 200 trotz db-ausfall (liveness prueft keine abhaengigkeiten)", livenessRec.Code)
}
readinessRec := httptest.NewRecorder()
reg.ReadinessHandler()(readinessRec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
if readinessRec.Code != http.StatusServiceUnavailable {
t.Fatalf("readiness status = %d, want 503 bei db-ausfall", readinessRec.Code)
}
if livenessRec.Code == readinessRec.Code {
t.Fatal("liveness und readiness sollten sich im db-ausfall-fall unterscheiden")
}
}
// Akzeptanzkriterium 2 + Pruefung 2: Health-Endpunkt antwortet auch bei
// haengendem Check innerhalb definierter Zeit (Timeout begrenzt die Dauer).
func TestReadinessHandler_RespondsWithinTimeoutEvenWithHangingCheck(t *testing.T) {
reg := NewRegistry().WithTimeout(50 * time.Millisecond)
reg.Register("haengender_dienst", CheckerFunc(func(ctx context.Context) error {
select {
case <-time.After(10 * time.Second): // wuerde ohne timeout ewig blockieren
return nil
case <-ctx.Done():
return ctx.Err()
}
}))
start := time.Now()
req := httptest.NewRequest(http.MethodGet, "/readyz", nil)
rec := httptest.NewRecorder()
reg.ReadinessHandler()(rec, req)
elapsed := time.Since(start)
if elapsed > time.Second {
t.Fatalf("readiness handler brauchte %s, erwartet deutlich unter 1s durch timeout", elapsed)
}
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want 503 fuer haengenden/timeout-check", rec.Code)
}
}
func TestCheckAll_MultipleChecksRunConcurrently(t *testing.T) {
reg := NewRegistry().WithTimeout(time.Second)
reg.Register("a", CheckerFunc(func(ctx context.Context) error { return nil }))
reg.Register("b", CheckerFunc(func(ctx context.Context) error { return errors.New("kaputt") }))
ready, results := reg.CheckAll(context.Background())
if ready {
t.Fatal("erwartet ready=false, da 'b' fehlschlaegt")
}
if !results["a"].OK {
t.Fatalf("erwartet 'a' ok, habe %+v", results["a"])
}
if results["b"].OK || results["b"].Error == "" {
t.Fatalf("erwartet 'b' fehlgeschlagen mit fehlertext, habe %+v", results["b"])
}
}
+99
View File
@@ -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
}
+46
View File
@@ -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)
}
}
+120
View File
@@ -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
}
+304
View File
@@ -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)
}
}
-206
View File
@@ -1,206 +0,0 @@
package tenant
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var (
ErrTenantNotFound = errors.New("tenant: nicht gefunden")
ErrInvalidTransition = errors.New("tenant: ungueltiger zustandsuebergang")
// ErrTenantNotActive wird von Lifecycle.CheckActive verwendet — bewusst
// EIN Fehler fuer suspendiert/zur-Loeschung-vorgemerkt/geloescht, da der
// Aufrufer (z.B. Login) nur wissen muss "kein Zugriff", nicht welcher der
// Nicht-aktiv-Zustaende genau vorliegt.
ErrTenantNotActive = errors.New("tenant: nicht aktiv")
)
func scanTenantWithLifecycle(row pgx.Row) (Tenant, error) {
var t Tenant
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status,
&t.CreatedAt, &t.PreviousStatus, &t.DeletionScheduledAt); err != nil {
return Tenant{}, err
}
return t, nil
}
// transition fuehrt einen bewachten Zustandsuebergang aus: das UPDATE greift
// nur, wenn der aktuelle Status einer von allowedFrom ist (atomarer
// Check-and-Set, kein Race zwischen Lesen und Schreiben). Greift es nicht,
// wird zwischen "Tenant existiert nicht" und "Uebergang nicht erlaubt"
// unterschieden, damit AC1 ("ungueltige Uebergaenge werden abgewiesen") einen
// sprechenden Fehler liefert statt eines stillen No-Ops.
func (r *Registry) transition(ctx context.Context, slug string, allowedFrom []Status, to Status, previousStatus *string, deletionAt *time.Time) (Tenant, error) {
from := make([]string, len(allowedFrom))
for i, s := range allowedFrom {
from[i] = string(s)
}
row := r.pool.QueryRow(ctx, `
UPDATE tenants
SET status = $2, previous_status = $3, deletion_scheduled_at = $4
WHERE slug = $1 AND status = ANY($5)
RETURNING id, slug, name, db_name, db_dsn, status, created_at, previous_status, deletion_scheduled_at
`, slug, string(to), previousStatus, deletionAt, from)
t, err := scanTenantWithLifecycle(row)
if err == nil {
return t, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return Tenant{}, fmt.Errorf("zustandsuebergang: %w", err)
}
existing, getErr := r.GetBySlug(ctx, slug)
if getErr != nil {
return Tenant{}, ErrTenantNotFound
}
return Tenant{}, fmt.Errorf("%w: von %q nach %q (aktuell: %q)", ErrInvalidTransition, allowedFrom, to, existing.Status)
}
// Suspend haelt die Daten des Mandanten unveraendert, sperrt aber den Zugriff
// (Akzeptanzkriterium 1) — es findet keine Loeschung/Migration statt.
func (r *Registry) Suspend(ctx context.Context, slug string) (Tenant, error) {
return r.transition(ctx, slug, []Status{StatusActive}, StatusSuspended, nil, nil)
}
// Reactivate stellt den Zustand vor der Suspendierung vollstaendig wieder her
// (Akzeptanzkriterium 2) — da Suspend keine weiteren Daten veraendert, genuegt
// die Rueckkehr nach StatusActive.
func (r *Registry) Reactivate(ctx context.Context, slug string) (Tenant, error) {
return r.transition(ctx, slug, []Status{StatusSuspended}, StatusActive, nil, nil)
}
// ScheduleDeletion merkt den Mandanten zur Loeschung vor und startet die
// Karenzzeit (Akzeptanzkriterium 3). previous_status wird festgehalten, damit
// CancelDeletion exakt dorthin zurueckkehren kann (aktiv ODER suspendiert).
func (r *Registry) ScheduleDeletion(ctx context.Context, slug string, grace time.Duration) (Tenant, error) {
existing, err := r.GetBySlug(ctx, slug)
if err != nil {
return Tenant{}, ErrTenantNotFound
}
prev := string(existing.Status)
deletionAt := time.Now().Add(grace)
return r.transition(ctx, slug, []Status{StatusActive, StatusSuspended}, StatusPendingDeletion, &prev, &deletionAt)
}
// CancelDeletion widerruft eine Loeschvormerkung innerhalb der Karenzzeit und
// stellt exakt den zuvor gesicherten Zustand wieder her.
func (r *Registry) CancelDeletion(ctx context.Context, slug string) (Tenant, error) {
existing, err := r.GetBySlug(ctx, slug)
if err != nil {
return Tenant{}, ErrTenantNotFound
}
if existing.Status != StatusPendingDeletion || existing.PreviousStatus == nil {
return Tenant{}, fmt.Errorf("%w: von %q nach aktiv/suspendiert (aktuell: %q)", ErrInvalidTransition, StatusPendingDeletion, existing.Status)
}
restoreTo := Status(*existing.PreviousStatus)
return r.transition(ctx, slug, []Status{StatusPendingDeletion}, restoreTo, nil, nil)
}
// Lifecycle fuehrt die tatsaechliche, physische Loeschung nach Ablauf der
// Karenzzeit aus (Datenbank-Drop) und stellt die Zugriffsschutz-Pruefung
// bereit. Getrennt von Registry, weil hierfuer zusaetzlich der adminPool
// (fuer DROP DATABASE) noetig ist, siehe internal/tenant.Provisioner.
type Lifecycle struct {
registry *Registry
adminPool *pgxpool.Pool
}
func NewLifecycle(registry *Registry, adminPool *pgxpool.Pool) *Lifecycle {
return &Lifecycle{registry: registry, adminPool: adminPool}
}
// CheckActive verweigert Zugriff fuer jeden Nicht-aktiv-Zustand und loggt den
// Vorgang strukturiert (Akzeptanzkriterium 1 / Pruefung 2).
func (l *Lifecycle) CheckActive(ctx context.Context, slug string) error {
t, err := l.registry.GetBySlug(ctx, slug)
if err != nil {
return ErrTenantNotFound
}
if t.Status != StatusActive {
slog.Warn("zugriff auf nicht-aktiven mandanten verweigert",
"tenant_slug", slug, "tenant_status", t.Status)
return ErrTenantNotActive
}
return nil
}
// ProcessDueDeletions loescht alle Mandanten-Datenbanken, deren Karenzzeit
// abgelaufen ist (Akzeptanzkriterium 3 / Pruefung 3). FOR UPDATE SKIP LOCKED
// folgt der projektweiten Postgres-Jobqueue-Konvention (siehe
// SKALIERUNGSKONZEPT.md) und macht die Funktion sicher fuer mehrere parallel
// laufende Core-Instanzen.
func (l *Lifecycle) ProcessDueDeletions(ctx context.Context) (int, error) {
tx, err := l.registry.pool.Begin(ctx)
if err != nil {
return 0, fmt.Errorf("sweep-transaktion starten: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
rows, err := tx.Query(ctx, `
SELECT id, db_name FROM tenants
WHERE status = $1 AND deletion_scheduled_at <= now()
FOR UPDATE SKIP LOCKED
`, string(StatusPendingDeletion))
if err != nil {
return 0, fmt.Errorf("faellige loeschungen abfragen: %w", err)
}
type due struct{ id, dbName string }
var candidates []due
for rows.Next() {
var d due
if err := rows.Scan(&d.id, &d.dbName); err != nil {
rows.Close()
return 0, fmt.Errorf("faellige loeschung lesen: %w", err)
}
candidates = append(candidates, d)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, err
}
processed := 0
for _, c := range candidates {
if _, err := l.adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, c.dbName)); err != nil {
return processed, fmt.Errorf("tenant-datenbank %q loeschen: %w", c.dbName, err)
}
if _, err := tx.Exec(ctx, `
UPDATE tenants SET status = $2, previous_status = NULL, deletion_scheduled_at = NULL
WHERE id = $1
`, c.id, string(StatusDeleted)); err != nil {
return processed, fmt.Errorf("tenant %q als geloescht markieren: %w", c.id, err)
}
processed++
}
if err := tx.Commit(ctx); err != nil {
return 0, fmt.Errorf("sweep-transaktion committen: %w", err)
}
return processed, nil
}
// RunSweeper triggert ProcessDueDeletions periodisch, bis ctx beendet wird —
// die "In-Prozess-Worker-Goroutine" aus der projektweiten Jobqueue-Konvention.
func (l *Lifecycle) RunSweeper(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if _, err := l.ProcessDueDeletions(ctx); err != nil {
slog.Error("tenant-loeschung-sweep fehlgeschlagen", "error", err)
}
}
}
}
-253
View File
@@ -1,253 +0,0 @@
package tenant
import (
"context"
"errors"
"fmt"
"os"
"strings"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func newLifecycleTestSetup(t *testing.T) (*Registry, *Lifecycle, *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()
adminPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("admin pool: %v", err)
}
registryPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("registry pool: %v", err)
}
if _, err := registryPool.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(),
previous_status TEXT,
deletion_scheduled_at TIMESTAMPTZ
)`); err != nil {
t.Fatalf("registry-schema: %v", err)
}
registry := NewRegistry(registryPool)
dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1)
provisioner := NewProvisioner(adminPool, registry, dsnTemplate)
lifecycle := NewLifecycle(registry, adminPool)
cleanup := func() {
registryPool.Close()
adminPool.Close()
}
_ = provisioner
return registry, lifecycle, adminPool, cleanup
}
func provisionTestTenant(t *testing.T, registry *Registry, adminPool *pgxpool.Pool, slug string) {
t.Helper()
dsnTemplate := strings.Replace(os.Getenv("TEST_ADMIN_DSN"), "/postgres?", "/%s?", 1)
provisioner := NewProvisioner(adminPool, registry, dsnTemplate)
if _, err := provisioner.Provision(context.Background(), slug, slug); err != nil {
t.Fatalf("provision %s: %v", slug, err)
}
t.Cleanup(func() {
ctx := context.Background()
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbNameForSlug(slug)))
_, _ = registry.pool.Exec(ctx, `DELETE FROM tenants WHERE slug = $1`, slug)
})
}
// Akzeptanzkriterium 1 (Suspend) + 2 (Reactivate) + Pruefung 1 (Uebergaenge).
func TestLifecycle_SuspendAndReactivate(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_suspend")
ctx := context.Background()
suspended, err := registry.Suspend(ctx, "lc_suspend")
if err != nil {
t.Fatalf("suspend: %v", err)
}
if suspended.Status != StatusSuspended {
t.Fatalf("status = %q, want suspended", suspended.Status)
}
reactivated, err := registry.Reactivate(ctx, "lc_suspend")
if err != nil {
t.Fatalf("reactivate: %v", err)
}
if reactivated.Status != StatusActive {
t.Fatalf("status = %q, want active", reactivated.Status)
}
}
// Pruefung 1: ungueltige Uebergaenge werden abgewiesen.
func TestLifecycle_RejectsInvalidTransitions(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_invalid")
ctx := context.Background()
// Reactivate auf einem bereits aktiven Tenant ist kein gueltiger Uebergang.
if _, err := registry.Reactivate(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
if _, err := registry.Suspend(ctx, "lc_invalid"); err != nil {
t.Fatalf("suspend: %v", err)
}
// Suspend auf einem bereits suspendierten Tenant ist ebenfalls ungueltig.
if _, err := registry.Suspend(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
// CancelDeletion ohne vorherige Loeschvormerkung ist ungueltig.
if _, err := registry.CancelDeletion(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
if _, err := registry.Suspend(ctx, "unbekannter-slug-xyz"); !errors.Is(err, ErrTenantNotFound) {
t.Fatalf("erwartet ErrTenantNotFound, habe %v", err)
}
}
// Akzeptanzkriterium 3: Loeschung zweistufig mit Karenzzeit, innerhalb der
// Frist widerrufbar — sowohl aus 'active' als auch aus 'suspended' heraus,
// mit exakter Wiederherstellung des jeweiligen Vorzustands.
func TestLifecycle_ScheduleAndCancelDeletion_RestoresExactPreviousState(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_cancel_active")
provisionTestTenant(t, registry, adminPool, "lc_cancel_suspended")
ctx := context.Background()
// Fall 1: aus 'active' heraus vorgemerkt und widerrufen.
scheduled, err := registry.ScheduleDeletion(ctx, "lc_cancel_active", time.Hour)
if err != nil {
t.Fatalf("schedule deletion: %v", err)
}
if scheduled.Status != StatusPendingDeletion {
t.Fatalf("status = %q, want pending_deletion", scheduled.Status)
}
if scheduled.DeletionScheduledAt == nil {
t.Fatal("erwartet gesetzte deletion_scheduled_at")
}
restored, err := registry.CancelDeletion(ctx, "lc_cancel_active")
if err != nil {
t.Fatalf("cancel deletion: %v", err)
}
if restored.Status != StatusActive {
t.Fatalf("status = %q, want active (vorheriger zustand)", restored.Status)
}
// Fall 2: aus 'suspended' heraus vorgemerkt und widerrufen — muss zu
// 'suspended' zurueckkehren, NICHT zu 'active'.
if _, err := registry.Suspend(ctx, "lc_cancel_suspended"); err != nil {
t.Fatalf("suspend: %v", err)
}
if _, err := registry.ScheduleDeletion(ctx, "lc_cancel_suspended", time.Hour); err != nil {
t.Fatalf("schedule deletion: %v", err)
}
restoredSuspended, err := registry.CancelDeletion(ctx, "lc_cancel_suspended")
if err != nil {
t.Fatalf("cancel deletion: %v", err)
}
if restoredSuspended.Status != StatusSuspended {
t.Fatalf("status = %q, want suspended (vorheriger zustand)", restoredSuspended.Status)
}
}
// Akzeptanzkriterium 1 + Pruefung 2: suspendierter Tenant erzeugt bei jedem
// Zugriffsversuch einen klaren Fehler.
func TestLifecycle_CheckActive_RejectsNonActive(t *testing.T) {
registry, lifecycle, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_checkactive")
ctx := context.Background()
if err := lifecycle.CheckActive(ctx, "lc_checkactive"); err != nil {
t.Fatalf("aktiver tenant sollte durchgehen, habe %v", err)
}
if _, err := registry.Suspend(ctx, "lc_checkactive"); err != nil {
t.Fatalf("suspend: %v", err)
}
for i := 0; i < 3; i++ {
if err := lifecycle.CheckActive(ctx, "lc_checkactive"); !errors.Is(err, ErrTenantNotActive) {
t.Fatalf("versuch %d: erwartet ErrTenantNotActive, habe %v", i, err)
}
}
if err := lifecycle.CheckActive(ctx, "nie-registriert"); !errors.Is(err, ErrTenantNotFound) {
t.Fatalf("erwartet ErrTenantNotFound, habe %v", err)
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Loeschvorgang nach Ablauf der Karenzzeit
// automatisch ausgeloest (hier durch direkten Aufruf von ProcessDueDeletions,
// das RunSweeper periodisch aufruft).
func TestLifecycle_ProcessDueDeletions(t *testing.T) {
registry, lifecycle, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_due")
provisionTestTenant(t, registry, adminPool, "lc_not_due")
ctx := context.Background()
// lc_due: Karenzzeit liegt bereits in der Vergangenheit -> faellig.
if _, err := registry.ScheduleDeletion(ctx, "lc_due", -time.Minute); err != nil {
t.Fatalf("schedule deletion (due): %v", err)
}
// lc_not_due: Karenzzeit liegt weit in der Zukunft -> nicht faellig.
if _, err := registry.ScheduleDeletion(ctx, "lc_not_due", time.Hour); err != nil {
t.Fatalf("schedule deletion (not due): %v", err)
}
processed, err := lifecycle.ProcessDueDeletions(ctx)
if err != nil {
t.Fatalf("process due deletions: %v", err)
}
if processed != 1 {
t.Fatalf("erwartet genau 1 verarbeitete loeschung, habe %d", processed)
}
due, err := registry.GetBySlug(ctx, "lc_due")
if err != nil {
t.Fatalf("get lc_due: %v", err)
}
if due.Status != StatusDeleted {
t.Fatalf("lc_due status = %q, want deleted", due.Status)
}
notDue, err := registry.GetBySlug(ctx, "lc_not_due")
if err != nil {
t.Fatalf("get lc_not_due: %v", err)
}
if notDue.Status != StatusPendingDeletion {
t.Fatalf("lc_not_due status = %q, want pending_deletion (noch nicht faellig)", notDue.Status)
}
// Datenbank von lc_due wurde tatsaechlich physisch entfernt.
var exists bool
if err := adminPool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM pg_database WHERE datname = $1)`,
dbNameForSlug("lc_due")).Scan(&exists); err != nil {
t.Fatalf("pg_database pruefen: %v", err)
}
if exists {
t.Fatal("erwartet, dass die tenant-datenbank von lc_due geloescht wurde")
}
}
+2 -6
View File
@@ -35,17 +35,13 @@ func (r *Registry) insertTx(ctx context.Context, tx pgx.Tx, t Tenant) (Tenant, e
}
func (r *Registry) GetBySlug(ctx context.Context, slug string) (Tenant, error) {
// previous_status/deletion_scheduled_at werden mitgelesen, damit TEN-04
// (internal/tenant/lifecycle.go) den vollstaendigen Lebenszyklus-Zustand
// ueber GetBySlug ansehen kann, statt eine eigene Abfrage zu duplizieren.
var t Tenant
row := r.pool.QueryRow(ctx, `
SELECT id, slug, name, db_name, db_dsn, status, created_at, previous_status, deletion_scheduled_at
SELECT id, slug, name, db_name, db_dsn, status, created_at
FROM tenants WHERE slug = $1
`, slug)
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status, &t.CreatedAt,
&t.PreviousStatus, &t.DeletionScheduledAt); err != nil {
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status, &t.CreatedAt); err != nil {
return Tenant{}, fmt.Errorf("tenant laden: %w", err)
}
return t, nil
-9
View File
@@ -12,10 +12,6 @@ type Status string
const (
StatusActive Status = "active"
// Lebenszyklus-Zustaende aus TEN-04 (siehe internal/tenant/lifecycle.go).
StatusSuspended Status = "suspended"
StatusPendingDeletion Status = "pending_deletion"
StatusDeleted Status = "deleted"
)
type Tenant struct {
@@ -26,11 +22,6 @@ type Tenant struct {
DBDSN string
Status Status
CreatedAt time.Time
// PreviousStatus und DeletionScheduledAt sind nur waehrend
// StatusPendingDeletion gesetzt (TEN-04) — sie halten fest, in welchen
// Zustand CancelDeletion zurueckkehrt und wann die Karenzzeit ablaeuft.
PreviousStatus *string
DeletionScheduledAt *time.Time
}
// slugPattern erzwingt sichere, als SQL-Identifier verwendbare Slugs, damit
@@ -1,2 +0,0 @@
ALTER TABLE tenants DROP COLUMN previous_status;
ALTER TABLE tenants DROP COLUMN deletion_scheduled_at;
-6
View File
@@ -1,6 +0,0 @@
-- Lebenszyklus-Zustaende fuer Mandanten (TEN-04, siehe core-kanban/tickets/TEN-04.md).
-- previous_status haelt den Zustand VOR einer Loeschvormerkung, damit
-- CancelDeletion "den vorherigen Zustand vollstaendig wiederherstellt"
-- (aktiv ODER suspendiert), statt hart auf 'active' zurueckzusetzen.
ALTER TABLE tenants ADD COLUMN previous_status TEXT;
ALTER TABLE tenants ADD COLUMN deletion_scheduled_at TIMESTAMPTZ;
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS feature_flags;
+10
View File
@@ -0,0 +1,10 @@
-- Feature-Flags zentral je Mandant/Zielgruppe (LIC-02, siehe core-kanban/tickets/LIC-02.md).
-- Lebt in der Registry-DB, nicht pro Tenant-Datenbank — Flags sind eine
-- Core-weite Konfiguration, keine Mandanten-Geschaeftsdaten.
CREATE TABLE feature_flags (
key TEXT PRIMARY KEY,
enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0 CHECK (rollout_percentage BETWEEN 0 AND 100),
target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+2
View File
@@ -0,0 +1,2 @@
DROP TABLE IF EXISTS module_credentials;
DROP TABLE IF EXISTS modules;
+18
View File
@@ -0,0 +1,18 @@
-- Modul-Registry & Aktivierungspruefung (API-02, siehe core-kanban/tickets/API-02.md).
CREATE TABLE modules (
name TEXT PRIMARY KEY,
version TEXT NOT NULL CHECK (version <> ''),
required_flags TEXT[] NOT NULL DEFAULT '{}',
registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Service-Credential je Modul-Instanz, bei Provisionierung ausgestellt
-- (Akzeptanzkriterium 4). secret_hash enthaelt NIEMALS das Secret im
-- Klartext, nur dessen SHA-256-Hash (Timing-safe-Vergleich beim Login,
-- Referenzmuster siehe AUD-02).
CREATE TABLE module_credentials (
module_name TEXT PRIMARY KEY REFERENCES modules(name),
client_id TEXT NOT NULL UNIQUE,
secret_hash BYTEA NOT NULL,
issued_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+1 -13
View File
@@ -1,23 +1,11 @@
#!/usr/bin/env bash
# Setzt die nexarch-Testumgebung zurueck: loescht die geteilte
# Registry-Tabelle "tenants" in der postgres-Wartungsdatenbank sowie alle
# tenant_*-Datenbanken. Noetig, weil verschiedene Feature-Branches
# unterschiedliche Registry-Schemata erwarten, aber dieselbe physische
# Postgres-Instanz auf dem Testhost teilen (siehe [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/reset-test-env.sh
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"
ROLE="nexarch_test"
export PGPASSWORD="$PASS"
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants;"
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants 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
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP DATABASE IF EXISTS \"${db}\";"
done
echo "Testumgebung zurueckgesetzt: registry-tabelle + $(echo "$dbs" | grep -c . || true) tenant-datenbank(en) entfernt."
+12
View File
@@ -0,0 +1,12 @@
#!/usr/bin/env bash
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"
cd "$(dirname "$0")/.."
NEXARCH_TEST_DB_PASSWORD="$PASS" bash scripts/reset-test-env.sh
export TEST_ADMIN_DSN="postgresql://nexarch_test:${PASS}@localhost:5432/postgres?sslmode=disable"
echo "== go build =="
go build ./...
echo "== go vet =="
go vet ./...
echo "== go test (-p 1) =="
go test ./... -p 1 -count=1