Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b23cd1961f | ||
|
|
4a30345e07 |
@@ -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
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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()
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
@@ -1,106 +0,0 @@
|
|||||||
// Package license implementiert Core LIC-01: Lizenzmodell je Tenant (Plan,
|
|
||||||
// Modul-Umfang, Laufzeit) und die kryptographische Pruefung signierter
|
|
||||||
// Lizenzschluessel. Feature-Flag-AUSWERTUNG zur Laufzeit (LIC-02) und die
|
|
||||||
// Verwaltungsoberflaeche (LIC-04) sind ausdruecklich nicht Teil dieses Pakets
|
|
||||||
// — hier geht es nur um Ausstellung/Validierung/Persistenz (Unleash-Vorbild:
|
|
||||||
// klare Trennung Flag-Verwaltung vs. Flag-Auswertung).
|
|
||||||
package license
|
|
||||||
|
|
||||||
import (
|
|
||||||
"crypto/ed25519"
|
|
||||||
"encoding/base64"
|
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
var (
|
|
||||||
ErrInvalidSignature = errors.New("license: signatur ungueltig")
|
|
||||||
ErrMalformedKey = errors.New("license: lizenzschluessel hat ungueltiges format")
|
|
||||||
)
|
|
||||||
|
|
||||||
// Payload ist der signierte Lizenzinhalt (Akzeptanzkriterium 3: Plan,
|
|
||||||
// Modul-Liste, Laufzeit).
|
|
||||||
type Payload struct {
|
|
||||||
TenantSlug string `json:"tenant_slug"`
|
|
||||||
Plan string `json:"plan"`
|
|
||||||
Modules []string `json:"modules"`
|
|
||||||
IssuedAt time.Time `json:"issued_at"`
|
|
||||||
ValidUntil time.Time `json:"valid_until"`
|
|
||||||
}
|
|
||||||
|
|
||||||
// Issuer stellt signierte Lizenzschluessel aus. Haelt den PRIVATEN
|
|
||||||
// Ed25519-Schluessel — lebt in der Praxis beim Lizenzgeber, nicht im
|
|
||||||
// laufenden Core-Prozess (der nur den Validator mit dem oeffentlichen
|
|
||||||
// Schluessel braucht).
|
|
||||||
type Issuer struct {
|
|
||||||
priv ed25519.PrivateKey
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewIssuer(priv ed25519.PrivateKey) *Issuer {
|
|
||||||
return &Issuer{priv: priv}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Issue liefert den Lizenzschluessel im Format base64(payload-json) "." base64(signatur).
|
|
||||||
func (i *Issuer) Issue(payload Payload) (string, error) {
|
|
||||||
raw, err := json.Marshal(payload)
|
|
||||||
if err != nil {
|
|
||||||
return "", fmt.Errorf("payload serialisieren: %w", err)
|
|
||||||
}
|
|
||||||
sig := ed25519.Sign(i.priv, raw)
|
|
||||||
|
|
||||||
return base64.RawURLEncoding.EncodeToString(raw) + "." + base64.RawURLEncoding.EncodeToString(sig), nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Validator prueft Lizenzschluessel gegen den OEFFENTLICHEN Ed25519-Schluessel
|
|
||||||
// — das ist alles, was der laufende Core-Prozess kennen muss.
|
|
||||||
type Validator struct {
|
|
||||||
pub ed25519.PublicKey
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewValidator(pub ed25519.PublicKey) *Validator {
|
|
||||||
return &Validator{pub: pub}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Parse prueft die Signatur (Akzeptanzkriterium 1 / Pruefung 1) und liefert
|
|
||||||
// bei Erfolg den entschluesselten Payload. Ein manipulierter Schluessel wird
|
|
||||||
// hier zuverlaessig erkannt, unabhaengig davon, ob die Laufzeit noch gueltig
|
|
||||||
// waere — Signaturpruefung und Ablaufpruefung sind bewusst getrennt
|
|
||||||
// (Signatur bei Einspielen, Ablauf bei jeder Nutzung, siehe Store.RequireActive).
|
|
||||||
func (v *Validator) Parse(key string) (Payload, error) {
|
|
||||||
rawPart, sigPart, ok := splitOnce(key, '.')
|
|
||||||
if !ok {
|
|
||||||
return Payload{}, ErrMalformedKey
|
|
||||||
}
|
|
||||||
|
|
||||||
raw, err := base64.RawURLEncoding.DecodeString(rawPart)
|
|
||||||
if err != nil {
|
|
||||||
return Payload{}, ErrMalformedKey
|
|
||||||
}
|
|
||||||
sig, err := base64.RawURLEncoding.DecodeString(sigPart)
|
|
||||||
if err != nil {
|
|
||||||
return Payload{}, ErrMalformedKey
|
|
||||||
}
|
|
||||||
|
|
||||||
if !ed25519.Verify(v.pub, raw, sig) {
|
|
||||||
return Payload{}, ErrInvalidSignature
|
|
||||||
}
|
|
||||||
|
|
||||||
var p Payload
|
|
||||||
if err := json.Unmarshal(raw, &p); err != nil {
|
|
||||||
// Signatur war gueltig, aber Payload nicht mehr parsebar — sollte bei
|
|
||||||
// unveraenderten Schluesseln nie vorkommen, trotzdem kein Panic.
|
|
||||||
return Payload{}, fmt.Errorf("%w: payload nicht lesbar", ErrMalformedKey)
|
|
||||||
}
|
|
||||||
return p, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func splitOnce(s string, sep byte) (before, after string, ok bool) {
|
|
||||||
for i := 0; i < len(s); i++ {
|
|
||||||
if s[i] == sep {
|
|
||||||
return s[:i], s[i+1:], true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return "", "", false
|
|
||||||
}
|
|
||||||
@@ -1,106 +0,0 @@
|
|||||||
package license
|
|
||||||
|
|
||||||
import (
|
|
||||||
"crypto/ed25519"
|
|
||||||
"errors"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
func testKeyPair(t *testing.T) (ed25519.PublicKey, ed25519.PrivateKey) {
|
|
||||||
t.Helper()
|
|
||||||
pub, priv, err := ed25519.GenerateKey(nil)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("schluesselpaar erzeugen: %v", err)
|
|
||||||
}
|
|
||||||
return pub, priv
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestIssueAndParse_RoundTrip(t *testing.T) {
|
|
||||||
pub, priv := testKeyPair(t)
|
|
||||||
issuer := NewIssuer(priv)
|
|
||||||
validator := NewValidator(pub)
|
|
||||||
|
|
||||||
payload := Payload{
|
|
||||||
TenantSlug: "acme",
|
|
||||||
Plan: "pro",
|
|
||||||
Modules: []string{"dms", "mail"},
|
|
||||||
IssuedAt: time.Now().Truncate(time.Second),
|
|
||||||
ValidUntil: time.Now().Add(365 * 24 * time.Hour).Truncate(time.Second),
|
|
||||||
}
|
|
||||||
|
|
||||||
key, err := issuer.Issue(payload)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := validator.Parse(key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("parse: %v", err)
|
|
||||||
}
|
|
||||||
if got.TenantSlug != payload.TenantSlug || got.Plan != payload.Plan || len(got.Modules) != 2 {
|
|
||||||
t.Fatalf("payload nach parse unerwartet: %+v", got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 1 + Pruefung 1: manipulierter Schluessel wird zuverlaessig erkannt.
|
|
||||||
func TestParse_RejectsTamperedKey(t *testing.T) {
|
|
||||||
pub, priv := testKeyPair(t)
|
|
||||||
issuer := NewIssuer(priv)
|
|
||||||
validator := NewValidator(pub)
|
|
||||||
|
|
||||||
key, err := issuer.Issue(Payload{TenantSlug: "acme", Plan: "pro", ValidUntil: time.Now().Add(time.Hour)})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Ein Zeichen im signierten Teil aendern.
|
|
||||||
tampered := []byte(key)
|
|
||||||
changed := false
|
|
||||||
for i, c := range tampered {
|
|
||||||
if c != '.' {
|
|
||||||
if c == 'A' {
|
|
||||||
tampered[i] = 'B'
|
|
||||||
} else {
|
|
||||||
tampered[i] = 'A'
|
|
||||||
}
|
|
||||||
changed = true
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if !changed {
|
|
||||||
t.Fatal("testaufbau fehlerhaft: nichts zum manipulieren gefunden")
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := validator.Parse(string(tampered)); !errors.Is(err, ErrInvalidSignature) {
|
|
||||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestParse_RejectsWrongKeyPair(t *testing.T) {
|
|
||||||
_, priv := testKeyPair(t)
|
|
||||||
otherPub, _ := testKeyPair(t)
|
|
||||||
|
|
||||||
issuer := NewIssuer(priv)
|
|
||||||
validator := NewValidator(otherPub) // falscher oeffentlicher Schluessel
|
|
||||||
|
|
||||||
key, err := issuer.Issue(Payload{TenantSlug: "acme", Plan: "pro"})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
if _, err := validator.Parse(key); !errors.Is(err, ErrInvalidSignature) {
|
|
||||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestParse_RejectsMalformedKey(t *testing.T) {
|
|
||||||
pub, _ := testKeyPair(t)
|
|
||||||
validator := NewValidator(pub)
|
|
||||||
|
|
||||||
cases := []string{"", "keine-punkt-trennung", "!!!.!!!"}
|
|
||||||
for _, c := range cases {
|
|
||||||
if _, err := validator.Parse(c); !errors.Is(err, ErrMalformedKey) {
|
|
||||||
t.Fatalf("Parse(%q): erwartet ErrMalformedKey, habe %v", c, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,88 +0,0 @@
|
|||||||
package license
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
|
||||||
)
|
|
||||||
|
|
||||||
var (
|
|
||||||
ErrNoLicense = errors.New("license: kein lizenzdatensatz fuer diesen tenant")
|
|
||||||
ErrLicenseExpired = errors.New("license: lizenz abgelaufen")
|
|
||||||
)
|
|
||||||
|
|
||||||
// Store persistiert den Lizenzumfang je Tenant in der Control-Plane-Registry
|
|
||||||
// (siehe internal/tenant.Registry — dieselbe Datenbank, aber ein eigener,
|
|
||||||
// unabhaengiger Store, um internal/tenant nicht um lizenzfremde Belange zu
|
|
||||||
// erweitern).
|
|
||||||
type Store struct {
|
|
||||||
pool *pgxpool.Pool
|
|
||||||
validator *Validator
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewStore(pool *pgxpool.Pool, validator *Validator) *Store {
|
|
||||||
return &Store{pool: pool, validator: validator}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Install prueft die Signatur des Lizenzschluessels (Akzeptanzkriterium 1)
|
|
||||||
// und ersetzt den bisherigen Lizenzdatensatz des Tenants vollstaendig. Ein
|
|
||||||
// bereits abgelaufener, aber korrekt signierter Schluessel wird trotzdem
|
|
||||||
// gespeichert — der Ablauf wird erst bei der Nutzung (RequireActive)
|
|
||||||
// bewertet, nicht beim Einspielen.
|
|
||||||
func (s *Store) Install(ctx context.Context, tenantID, licenseKey string) (Payload, error) {
|
|
||||||
payload, err := s.validator.Parse(licenseKey)
|
|
||||||
if err != nil {
|
|
||||||
return Payload{}, err
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = s.pool.Exec(ctx, `
|
|
||||||
INSERT INTO tenant_licenses (tenant_id, plan, modules, issued_at, valid_until, raw_key, installed_at)
|
|
||||||
VALUES ($1, $2, $3, $4, $5, $6, now())
|
|
||||||
ON CONFLICT (tenant_id) DO UPDATE SET
|
|
||||||
plan = $2, modules = $3, issued_at = $4, valid_until = $5, raw_key = $6, installed_at = now()
|
|
||||||
`, tenantID, payload.Plan, payload.Modules, payload.IssuedAt, payload.ValidUntil, licenseKey)
|
|
||||||
if err != nil {
|
|
||||||
return Payload{}, fmt.Errorf("lizenz speichern: %w", err)
|
|
||||||
}
|
|
||||||
return payload, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *Store) get(ctx context.Context, tenantID string) (Payload, error) {
|
|
||||||
var p Payload
|
|
||||||
row := s.pool.QueryRow(ctx, `
|
|
||||||
SELECT plan, modules, issued_at, valid_until
|
|
||||||
FROM tenant_licenses WHERE tenant_id = $1
|
|
||||||
`, tenantID)
|
|
||||||
if err := row.Scan(&p.Plan, &p.Modules, &p.IssuedAt, &p.ValidUntil); err != nil {
|
|
||||||
if errors.Is(err, pgx.ErrNoRows) {
|
|
||||||
return Payload{}, ErrNoLicense
|
|
||||||
}
|
|
||||||
return Payload{}, fmt.Errorf("lizenz lesen: %w", err)
|
|
||||||
}
|
|
||||||
return p, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Status liefert den persistierten Lizenzumfang unabhaengig vom Ablauf
|
|
||||||
// (Akzeptanzkriterium 3: Plan, Modul-Liste, Laufzeit abfragbar).
|
|
||||||
func (s *Store) Status(ctx context.Context, tenantID string) (Payload, error) {
|
|
||||||
return s.get(ctx, tenantID)
|
|
||||||
}
|
|
||||||
|
|
||||||
// RequireActive liefert den Lizenzumfang NUR, wenn die Lizenz noch nicht
|
|
||||||
// abgelaufen ist — sonst ErrLicenseExpired statt eines harten Fehlers/Panics
|
|
||||||
// (Akzeptanzkriterium 2: definierter eingeschraenkter Zustand). Aufrufende
|
|
||||||
// Module (LIC-02/03) entscheiden, was "eingeschraenkt" konkret bedeutet.
|
|
||||||
func (s *Store) RequireActive(ctx context.Context, tenantID string) (Payload, error) {
|
|
||||||
p, err := s.get(ctx, tenantID)
|
|
||||||
if err != nil {
|
|
||||||
return Payload{}, err
|
|
||||||
}
|
|
||||||
if time.Now().After(p.ValidUntil) {
|
|
||||||
return Payload{}, ErrLicenseExpired
|
|
||||||
}
|
|
||||||
return p, nil
|
|
||||||
}
|
|
||||||
@@ -1,177 +0,0 @@
|
|||||||
package license
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"crypto/ed25519"
|
|
||||||
"errors"
|
|
||||||
"os"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
|
||||||
)
|
|
||||||
|
|
||||||
func setupStoreTest(t *testing.T) (*Store, *Issuer, string, func()) {
|
|
||||||
t.Helper()
|
|
||||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
|
||||||
if adminDSN == "" {
|
|
||||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
|
||||||
}
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
pool, err := pgxpool.New(ctx, adminDSN)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("pool: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := pool.Exec(ctx, `
|
|
||||||
CREATE TABLE IF NOT EXISTS tenants (
|
|
||||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
||||||
slug TEXT NOT NULL UNIQUE,
|
|
||||||
name TEXT NOT NULL,
|
|
||||||
db_name TEXT NOT NULL UNIQUE,
|
|
||||||
db_dsn TEXT NOT NULL,
|
|
||||||
status TEXT NOT NULL DEFAULT 'active',
|
|
||||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
||||||
);
|
|
||||||
CREATE TABLE IF NOT EXISTS tenant_licenses (
|
|
||||||
tenant_id UUID PRIMARY KEY REFERENCES tenants(id),
|
|
||||||
plan TEXT NOT NULL,
|
|
||||||
modules TEXT[] NOT NULL,
|
|
||||||
issued_at TIMESTAMPTZ NOT NULL,
|
|
||||||
valid_until TIMESTAMPTZ NOT NULL,
|
|
||||||
raw_key TEXT NOT NULL,
|
|
||||||
installed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
||||||
);
|
|
||||||
`); err != nil {
|
|
||||||
t.Fatalf("schema: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var tenantID string
|
|
||||||
if err := pool.QueryRow(ctx, `
|
|
||||||
INSERT INTO tenants (slug, name, db_name, db_dsn)
|
|
||||||
VALUES ('lic_test_tenant', 'Lic Test', 'tenant_lic_test', 'unused')
|
|
||||||
RETURNING id
|
|
||||||
`).Scan(&tenantID); err != nil {
|
|
||||||
t.Fatalf("test-tenant anlegen: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
pub, priv, err := ed25519.GenerateKey(nil)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("schluesselpaar: %v", err)
|
|
||||||
}
|
|
||||||
issuer := NewIssuer(priv)
|
|
||||||
store := NewStore(pool, NewValidator(pub))
|
|
||||||
|
|
||||||
cleanup := func() {
|
|
||||||
_, _ = pool.Exec(ctx, `DELETE FROM tenant_licenses WHERE tenant_id = $1`, tenantID)
|
|
||||||
_, _ = pool.Exec(ctx, `DELETE FROM tenants WHERE id = $1`, tenantID)
|
|
||||||
pool.Close()
|
|
||||||
}
|
|
||||||
return store, issuer, tenantID, cleanup
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 3: Lizenzumfang persistiert und abfragbar.
|
|
||||||
func TestStore_InstallAndStatus(t *testing.T) {
|
|
||||||
store, issuer, tenantID, cleanup := setupStoreTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
payload := Payload{
|
|
||||||
TenantSlug: "lic_test_tenant",
|
|
||||||
Plan: "enterprise",
|
|
||||||
Modules: []string{"dms", "mail", "archive"},
|
|
||||||
IssuedAt: time.Now().Truncate(time.Second),
|
|
||||||
ValidUntil: time.Now().Add(30 * 24 * time.Hour).Truncate(time.Second),
|
|
||||||
}
|
|
||||||
key, err := issuer.Issue(payload)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := store.Install(ctx, tenantID, key); err != nil {
|
|
||||||
t.Fatalf("install: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
status, err := store.Status(ctx, tenantID)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("status: %v", err)
|
|
||||||
}
|
|
||||||
if status.Plan != "enterprise" || len(status.Modules) != 3 {
|
|
||||||
t.Fatalf("status unerwartet: %+v", status)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestStore_InstallRejectsInvalidSignature(t *testing.T) {
|
|
||||||
store, _, tenantID, cleanup := setupStoreTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
_, otherPriv, _ := ed25519.GenerateKey(nil)
|
|
||||||
foreignIssuer := NewIssuer(otherPriv) // signiert mit falschem schluessel
|
|
||||||
|
|
||||||
key, err := foreignIssuer.Issue(Payload{TenantSlug: "lic_test_tenant", Plan: "pro", ValidUntil: time.Now().Add(time.Hour)})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := store.Install(ctx, tenantID, key); !errors.Is(err, ErrInvalidSignature) {
|
|
||||||
t.Fatalf("erwartet ErrInvalidSignature, habe %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 2 + Pruefung 2: abgelaufene Lizenz fuehrt zu definiertem
|
|
||||||
// eingeschraenktem Zustand (ErrLicenseExpired), nicht zu einem Absturz.
|
|
||||||
func TestStore_RequireActive_DetectsExpiry(t *testing.T) {
|
|
||||||
store, issuer, tenantID, cleanup := setupStoreTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
expired := Payload{
|
|
||||||
TenantSlug: "lic_test_tenant",
|
|
||||||
Plan: "pro",
|
|
||||||
Modules: []string{"dms"},
|
|
||||||
IssuedAt: time.Now().Add(-48 * time.Hour),
|
|
||||||
ValidUntil: time.Now().Add(-24 * time.Hour), // bereits abgelaufen
|
|
||||||
}
|
|
||||||
key, err := issuer.Issue(expired)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("issue: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Einspielen einer bereits abgelaufenen, aber korrekt signierten Lizenz
|
|
||||||
// muss funktionieren (Ablauf wird erst bei Nutzung bewertet).
|
|
||||||
if _, err := store.Install(ctx, tenantID, key); err != nil {
|
|
||||||
t.Fatalf("install sollte trotz ablauf funktionieren: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
func() {
|
|
||||||
defer func() {
|
|
||||||
if r := recover(); r != nil {
|
|
||||||
t.Fatalf("RequireActive hat gepanict statt einen fehler zu liefern: %v", r)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
if _, err := store.RequireActive(ctx, tenantID); !errors.Is(err, ErrLicenseExpired) {
|
|
||||||
t.Fatalf("erwartet ErrLicenseExpired, habe %v", err)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
// Aber der Umfang bleibt weiterhin abfragbar (Status, im Unterschied zu RequireActive).
|
|
||||||
status, err := store.Status(ctx, tenantID)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("status sollte trotz ablauf funktionieren: %v", err)
|
|
||||||
}
|
|
||||||
if status.Plan != "pro" {
|
|
||||||
t.Fatalf("status unerwartet: %+v", status)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestStore_RequireActive_NoLicense(t *testing.T) {
|
|
||||||
store, _, tenantID, cleanup := setupStoreTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
if _, err := store.RequireActive(ctx, tenantID); !errors.Is(err, ErrNoLicense) {
|
|
||||||
t.Fatalf("erwartet ErrNoLicense, habe %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,32 +0,0 @@
|
|||||||
package usage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"log/slog"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
// AggregateFunc berechnet/aktualisiert Zaehlerstaende aus einer autoritativen
|
|
||||||
// Quelle (z.B. "zaehle Zeilen in einer Modul-Tabelle") — die konkrete Quelle
|
|
||||||
// haengt vom jeweiligen Modul ab und ist nicht Teil dieser Kachel. Das
|
|
||||||
// Aggregations-Grundgerüst selbst (periodischer Trigger) ist es.
|
|
||||||
type AggregateFunc func(ctx context.Context) error
|
|
||||||
|
|
||||||
// RunPeriodicAggregation ruft aggregate in festen Abstaenden auf, bis ctx
|
|
||||||
// beendet wird — dieselbe In-Prozess-Worker-Goroutine-Konvention wie
|
|
||||||
// internal/tenant.Lifecycle.RunSweeper (Akzeptanzkriterium 1: "periodisch
|
|
||||||
// aggregiert").
|
|
||||||
func RunPeriodicAggregation(ctx context.Context, interval time.Duration, aggregate AggregateFunc) {
|
|
||||||
ticker := time.NewTicker(interval)
|
|
||||||
defer ticker.Stop()
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
return
|
|
||||||
case <-ticker.C:
|
|
||||||
if err := aggregate(ctx); err != nil {
|
|
||||||
slog.Error("nutzungszaehler-aggregation fehlgeschlagen", "error", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,47 +0,0 @@
|
|||||||
package usage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
)
|
|
||||||
|
|
||||||
// StorageBytesMetric ist der feste Metrikname, unter dem der belegte
|
|
||||||
// Speicherplatz je Tenant gefuehrt wird (LIC-05, siehe
|
|
||||||
// core-kanban/tickets/LIC-05.md). LIC-03 fragt genau diese Metrik ueber
|
|
||||||
// Store.Get/Store.Check ab — kein zweiter, paralleler Speicher-Zaehler.
|
|
||||||
const StorageBytesMetric = "storage_bytes"
|
|
||||||
|
|
||||||
// ReportStorageWrite wird von den Objekt-Storage-Treibern der Module (DMS
|
|
||||||
// FDN-03, Mail ARC-01 — existieren als Code noch nicht) bei jedem
|
|
||||||
// Schreibvorgang aufgerufen. Nutzt Store.Increment, das bereits atomar ist
|
|
||||||
// (Akzeptanzkriterium 2, siehe LIC-03) — kein zweiter Inkrement-Mechanismus
|
|
||||||
// nur fuer Speicher.
|
|
||||||
func (s *Store) ReportStorageWrite(ctx context.Context, tenantID string, sizeBytes int64) error {
|
|
||||||
if sizeBytes < 0 {
|
|
||||||
return errors.New("usage: sizeBytes darf bei einem schreibvorgang nicht negativ sein")
|
|
||||||
}
|
|
||||||
if err := s.Increment(ctx, tenantID, StorageBytesMetric, sizeBytes); err != nil {
|
|
||||||
return fmt.Errorf("speicherverbrauch (schreiben) melden: %w", err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// ReportStorageDelete wird bei jedem Loeschvorgang aufgerufen — dekrementiert
|
|
||||||
// denselben Zaehler ueber ein negatives Delta desselben atomaren UPSERT.
|
|
||||||
func (s *Store) ReportStorageDelete(ctx context.Context, tenantID string, sizeBytes int64) error {
|
|
||||||
if sizeBytes < 0 {
|
|
||||||
return errors.New("usage: sizeBytes darf bei einem loeschvorgang nicht negativ sein")
|
|
||||||
}
|
|
||||||
if err := s.Increment(ctx, tenantID, StorageBytesMetric, -sizeBytes); err != nil {
|
|
||||||
return fmt.Errorf("speicherverbrauch (loeschen) melden: %w", err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// CurrentStorageUsage liefert den aktuellen Speicherverbrauch eines Tenants
|
|
||||||
// (Akzeptanzkriterium 3) — ein einfaches Get auf den bereits gefuehrten
|
|
||||||
// Zaehler, kein Scan des Objekt-Storage.
|
|
||||||
func (s *Store) CurrentStorageUsage(ctx context.Context, tenantID string) (int64, error) {
|
|
||||||
return s.Get(ctx, tenantID, StorageBytesMetric)
|
|
||||||
}
|
|
||||||
@@ -1,130 +0,0 @@
|
|||||||
package usage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"sync"
|
|
||||||
"testing"
|
|
||||||
)
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 1 + 2 + Pruefung 1: paralleler Schreib-Test (viele
|
|
||||||
// gleichzeitige Uploads) ergibt korrekten Endstand ohne verlorene Updates.
|
|
||||||
func TestReportStorageWrite_ConcurrentUploadsSumCorrectly(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
sizes := []int64{1024, 2048, 4096, 8192, 512, 256, 1000, 999, 1, 7000}
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
for _, size := range sizes {
|
|
||||||
wg.Add(1)
|
|
||||||
go func(sz int64) {
|
|
||||||
defer wg.Done()
|
|
||||||
if err := store.ReportStorageWrite(ctx, tenant, sz); err != nil {
|
|
||||||
t.Errorf("report write: %v", err)
|
|
||||||
}
|
|
||||||
}(size)
|
|
||||||
}
|
|
||||||
wg.Wait()
|
|
||||||
|
|
||||||
var expected int64
|
|
||||||
for _, s := range sizes {
|
|
||||||
expected += s
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := store.CurrentStorageUsage(ctx, tenant)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("current usage: %v", err)
|
|
||||||
}
|
|
||||||
if got != expected {
|
|
||||||
t.Fatalf("erwartet %d bytes (unabhaengige kontrollsumme), habe %d — hinweis auf verlorene updates", expected, got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 1 + Pruefung 2: Loeschvorgang dekrementiert korrekt.
|
|
||||||
func TestReportStorageDelete_Decrements(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
if err := store.ReportStorageWrite(ctx, tenant, 10_000); err != nil {
|
|
||||||
t.Fatalf("write: %v", err)
|
|
||||||
}
|
|
||||||
if err := store.ReportStorageDelete(ctx, tenant, 3_000); err != nil {
|
|
||||||
t.Fatalf("delete: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := store.CurrentStorageUsage(ctx, tenant)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("current usage: %v", err)
|
|
||||||
}
|
|
||||||
if got != 7_000 {
|
|
||||||
t.Fatalf("erwartet 7000 nach schreiben(10000)+loeschen(3000), habe %d", got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 3 + Pruefung 3: Abfrage liefert konsistenten Wert mit
|
|
||||||
// einer unabhaengigen Kontrollzaehlung ueber gemischte Schreib-/Loeschvorgaenge.
|
|
||||||
func TestCurrentStorageUsage_MatchesIndependentTally(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
type op struct {
|
|
||||||
write bool
|
|
||||||
size int64
|
|
||||||
}
|
|
||||||
ops := []op{
|
|
||||||
{true, 5000}, {true, 3000}, {false, 1000}, {true, 2000}, {false, 4000}, {true, 500},
|
|
||||||
}
|
|
||||||
|
|
||||||
var tally int64
|
|
||||||
for _, o := range ops {
|
|
||||||
if o.write {
|
|
||||||
if err := store.ReportStorageWrite(ctx, tenant, o.size); err != nil {
|
|
||||||
t.Fatalf("write: %v", err)
|
|
||||||
}
|
|
||||||
tally += o.size
|
|
||||||
} else {
|
|
||||||
if err := store.ReportStorageDelete(ctx, tenant, o.size); err != nil {
|
|
||||||
t.Fatalf("delete: %v", err)
|
|
||||||
}
|
|
||||||
tally -= o.size
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := store.CurrentStorageUsage(ctx, tenant)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("current usage: %v", err)
|
|
||||||
}
|
|
||||||
if got != tally {
|
|
||||||
t.Fatalf("erwartet %d (unabhaengige kontrollzaehlung), habe %d", tally, got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 3: LIC-03s generischer Store.Get liefert denselben Wert
|
|
||||||
// wie CurrentStorageUsage — kein zweiter, abweichender Zaehlmechanismus.
|
|
||||||
func TestCurrentStorageUsage_MatchesGenericStoreGet(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
if err := store.ReportStorageWrite(ctx, tenant, 42); err != nil {
|
|
||||||
t.Fatalf("write: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
viaStorage, err := store.CurrentStorageUsage(ctx, tenant)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("current usage: %v", err)
|
|
||||||
}
|
|
||||||
viaGeneric, err := store.Get(ctx, tenant, StorageBytesMetric)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("generic get: %v", err)
|
|
||||||
}
|
|
||||||
if viaStorage != viaGeneric || viaStorage != 42 {
|
|
||||||
t.Fatalf("erwartet beide wege liefern 42, habe storage=%d generic=%d", viaStorage, viaGeneric)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,144 +0,0 @@
|
|||||||
// Package usage implementiert Core LIC-03: Nutzungszaehler je Tenant
|
|
||||||
// (Benutzeranzahl, Speicherverbrauch, API-Aufrufe, ...) und die Pruefung
|
|
||||||
// gegen konfigurierte Quotas. Quotas sind Konfiguration (Tabellenzeile), kein
|
|
||||||
// Hardcode — Zitadel/Unleash-Vorbild.
|
|
||||||
package usage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
|
||||||
)
|
|
||||||
|
|
||||||
var ErrNoQuota = errors.New("usage: keine quota fuer diese metrik konfiguriert")
|
|
||||||
|
|
||||||
// Status ist die definierte Reaktion einer Quota-Pruefung (Akzeptanzkriterium 2).
|
|
||||||
type Status string
|
|
||||||
|
|
||||||
const (
|
|
||||||
StatusOK Status = "ok"
|
|
||||||
StatusWarning Status = "warning" // Schwelle (80%) erreicht, aber noch nicht ueberschritten
|
|
||||||
StatusExceeded Status = "exceeded" // Quota ueberschritten — neue Ressourcen sollten gesperrt werden
|
|
||||||
)
|
|
||||||
|
|
||||||
// warningThreshold liegt bei 80% der Quota.
|
|
||||||
const warningThreshold = 0.8
|
|
||||||
|
|
||||||
type Store struct {
|
|
||||||
pool *pgxpool.Pool
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewStore(pool *pgxpool.Pool) *Store {
|
|
||||||
return &Store{pool: pool}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Increment erhoeht einen Zaehler ATOMAR ueber ein einziges SQL-Statement
|
|
||||||
// (UPSERT mit value = value + delta) statt Read-Modify-Write in Go — das
|
|
||||||
// haelt Zaehlerstaende bei parallelen Schreibzugriffen konsistent
|
|
||||||
// (Akzeptanzkriterium 1 / Pruefung 2), ohne eine Anwendungs-Transaktion mit
|
|
||||||
// Lock zu brauchen.
|
|
||||||
func (s *Store) Increment(ctx context.Context, tenantID, metric string, delta int64) error {
|
|
||||||
_, err := s.pool.Exec(ctx, `
|
|
||||||
INSERT INTO usage_counters (tenant_id, metric, value, updated_at)
|
|
||||||
VALUES ($1, $2, $3, now())
|
|
||||||
ON CONFLICT (tenant_id, metric) DO UPDATE
|
|
||||||
SET value = usage_counters.value + $3, updated_at = now()
|
|
||||||
`, tenantID, metric, delta)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("zaehler erhoehen: %w", err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Get liefert den aktuellen Zaehlerstand — 0, wenn noch nie erhoeht wurde.
|
|
||||||
// Der Wert ist strikt tenant-gescoped (Akzeptanzkriterium 3 / Pruefung 3).
|
|
||||||
func (s *Store) Get(ctx context.Context, tenantID, metric string) (int64, error) {
|
|
||||||
var value int64
|
|
||||||
err := s.pool.QueryRow(ctx, `
|
|
||||||
SELECT value FROM usage_counters WHERE tenant_id = $1 AND metric = $2
|
|
||||||
`, tenantID, metric).Scan(&value)
|
|
||||||
if err != nil {
|
|
||||||
if errors.Is(err, pgx.ErrNoRows) {
|
|
||||||
return 0, nil
|
|
||||||
}
|
|
||||||
return 0, fmt.Errorf("zaehler lesen: %w", err)
|
|
||||||
}
|
|
||||||
return value, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// SetQuota legt die Obergrenze fuer (tenantID, metric) fest — Konfiguration,
|
|
||||||
// kein Hardcode.
|
|
||||||
func (s *Store) SetQuota(ctx context.Context, tenantID, metric string, limit int64) error {
|
|
||||||
_, err := s.pool.Exec(ctx, `
|
|
||||||
INSERT INTO usage_quotas (tenant_id, metric, limit_value)
|
|
||||||
VALUES ($1, $2, $3)
|
|
||||||
ON CONFLICT (tenant_id, metric) DO UPDATE SET limit_value = $3
|
|
||||||
`, tenantID, metric, limit)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("quota setzen: %w", err)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *Store) GetQuota(ctx context.Context, tenantID, metric string) (int64, error) {
|
|
||||||
var limit int64
|
|
||||||
err := s.pool.QueryRow(ctx, `
|
|
||||||
SELECT limit_value FROM usage_quotas WHERE tenant_id = $1 AND metric = $2
|
|
||||||
`, tenantID, metric).Scan(&limit)
|
|
||||||
if err != nil {
|
|
||||||
if errors.Is(err, pgx.ErrNoRows) {
|
|
||||||
return 0, ErrNoQuota
|
|
||||||
}
|
|
||||||
return 0, fmt.Errorf("quota lesen: %w", err)
|
|
||||||
}
|
|
||||||
return limit, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Check liefert Zaehlerstand, konfigurierte Quota und die daraus abgeleitete
|
|
||||||
// Reaktion (Akzeptanzkriterium 2 / Pruefung 1). Ist keine Quota konfiguriert,
|
|
||||||
// gilt die Metrik als unbegrenzt (StatusOK).
|
|
||||||
func (s *Store) Check(ctx context.Context, tenantID, metric string) (value, limit int64, status Status, err error) {
|
|
||||||
value, err = s.Get(ctx, tenantID, metric)
|
|
||||||
if err != nil {
|
|
||||||
return 0, 0, "", err
|
|
||||||
}
|
|
||||||
|
|
||||||
limit, err = s.GetQuota(ctx, tenantID, metric)
|
|
||||||
if errors.Is(err, ErrNoQuota) {
|
|
||||||
return value, 0, StatusOK, nil
|
|
||||||
}
|
|
||||||
if err != nil {
|
|
||||||
return 0, 0, "", err
|
|
||||||
}
|
|
||||||
|
|
||||||
switch {
|
|
||||||
case value > limit:
|
|
||||||
return value, limit, StatusExceeded, nil
|
|
||||||
case limit > 0 && float64(value) >= warningThreshold*float64(limit):
|
|
||||||
return value, limit, StatusWarning, nil
|
|
||||||
default:
|
|
||||||
return value, limit, StatusOK, nil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Reaction wird aufgerufen, wenn Check einen Nicht-OK-Status liefert
|
|
||||||
// (Akzeptanzkriterium 2: "definierte Reaktion").
|
|
||||||
type Reaction func(ctx context.Context, tenantID, metric string, value, limit int64, status Status)
|
|
||||||
|
|
||||||
// Enforce fuehrt Check aus und ruft react auf, wenn der Status nicht OK ist —
|
|
||||||
// die konkrete "Sperre neuer Ressourcen"/Benachrichtigung liegt beim
|
|
||||||
// Aufrufer (z.B. TEN-02 vor dem Anlegen eines neuen Benutzers), Enforce
|
|
||||||
// garantiert nur, dass die Reaktion zuverlaessig ausgeloest wird.
|
|
||||||
func (s *Store) Enforce(ctx context.Context, tenantID, metric string, react Reaction) (Status, error) {
|
|
||||||
value, limit, status, err := s.Check(ctx, tenantID, metric)
|
|
||||||
if err != nil {
|
|
||||||
return "", err
|
|
||||||
}
|
|
||||||
if status != StatusOK && react != nil {
|
|
||||||
react(ctx, tenantID, metric, value, limit, status)
|
|
||||||
}
|
|
||||||
return status, nil
|
|
||||||
}
|
|
||||||
@@ -1,206 +0,0 @@
|
|||||||
package usage
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"fmt"
|
|
||||||
"os"
|
|
||||||
"sync"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
|
||||||
)
|
|
||||||
|
|
||||||
func setupTest(t *testing.T) (*Store, func()) {
|
|
||||||
t.Helper()
|
|
||||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
|
||||||
if adminDSN == "" {
|
|
||||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
|
||||||
}
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
pool, err := pgxpool.New(ctx, adminDSN)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("pool: %v", err)
|
|
||||||
}
|
|
||||||
if _, err := pool.Exec(ctx, `
|
|
||||||
CREATE TABLE IF NOT EXISTS usage_counters (
|
|
||||||
tenant_id UUID NOT NULL, metric TEXT NOT NULL, value BIGINT NOT NULL DEFAULT 0,
|
|
||||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (tenant_id, metric)
|
|
||||||
);
|
|
||||||
CREATE TABLE IF NOT EXISTS usage_quotas (
|
|
||||||
tenant_id UUID NOT NULL, metric TEXT NOT NULL, limit_value BIGINT NOT NULL,
|
|
||||||
PRIMARY KEY (tenant_id, metric)
|
|
||||||
);
|
|
||||||
`); err != nil {
|
|
||||||
t.Fatalf("schema: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
cleanup := func() { pool.Close() }
|
|
||||||
return NewStore(pool), cleanup
|
|
||||||
}
|
|
||||||
|
|
||||||
func newTenantID() string {
|
|
||||||
return fmt.Sprintf("00000000-0000-0000-0000-%012d", time.Now().UnixNano()%1e12)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 1 + Pruefung 2: Aggregationsjob liefert bei parallelen
|
|
||||||
// Schreibzugriffen konsistente Zaehlerstaende.
|
|
||||||
func TestIncrement_ConsistentUnderConcurrentWrites(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
const goroutines = 50
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
for i := 0; i < goroutines; i++ {
|
|
||||||
wg.Add(1)
|
|
||||||
go func() {
|
|
||||||
defer wg.Done()
|
|
||||||
if err := store.Increment(ctx, tenant, "api_calls", 1); err != nil {
|
|
||||||
t.Errorf("increment: %v", err)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
|
||||||
wg.Wait()
|
|
||||||
|
|
||||||
value, err := store.Get(ctx, tenant, "api_calls")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("get: %v", err)
|
|
||||||
}
|
|
||||||
if value != goroutines {
|
|
||||||
t.Fatalf("erwartet %d, habe %d (hinweis auf lost update unter nebenlaeufigkeit)", goroutines, value)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 3 + Pruefung 3: Zaehlerstand eines Tenants beeinflusst
|
|
||||||
// nicht den eines anderen.
|
|
||||||
func TestIncrement_IsolatedBetweenTenants(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenantA, tenantB := newTenantID(), newTenantID()
|
|
||||||
|
|
||||||
if err := store.Increment(ctx, tenantA, "users", 5); err != nil {
|
|
||||||
t.Fatalf("increment a: %v", err)
|
|
||||||
}
|
|
||||||
if err := store.Increment(ctx, tenantB, "users", 1); err != nil {
|
|
||||||
t.Fatalf("increment b: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
valA, err := store.Get(ctx, tenantA, "users")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("get a: %v", err)
|
|
||||||
}
|
|
||||||
valB, err := store.Get(ctx, tenantB, "users")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("get b: %v", err)
|
|
||||||
}
|
|
||||||
if valA != 5 || valB != 1 {
|
|
||||||
t.Fatalf("erwartet a=5 b=1, habe a=%d b=%d", valA, valB)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Akzeptanzkriterium 2 + Pruefung 1: Quota-Ueberschreitung wird automatisiert
|
|
||||||
// erkannt und die definierte Reaktion ausgeloest.
|
|
||||||
func TestEnforce_TriggersReactionOnExceeded(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
if err := store.SetQuota(ctx, tenant, "users", 10); err != nil {
|
|
||||||
t.Fatalf("set quota: %v", err)
|
|
||||||
}
|
|
||||||
if err := store.Increment(ctx, tenant, "users", 11); err != nil {
|
|
||||||
t.Fatalf("increment: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var reacted bool
|
|
||||||
var gotStatus Status
|
|
||||||
status, err := store.Enforce(ctx, tenant, "users", func(ctx context.Context, tenantID, metric string, value, limit int64, status Status) {
|
|
||||||
reacted = true
|
|
||||||
gotStatus = status
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("enforce: %v", err)
|
|
||||||
}
|
|
||||||
if status != StatusExceeded {
|
|
||||||
t.Fatalf("erwartet StatusExceeded, habe %q", status)
|
|
||||||
}
|
|
||||||
if !reacted || gotStatus != StatusExceeded {
|
|
||||||
t.Fatal("erwartet ausgeloeste reaktion mit StatusExceeded")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestCheck_WarningThresholdAndOK(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
if err := store.SetQuota(ctx, tenant, "storage_mb", 100); err != nil {
|
|
||||||
t.Fatalf("set quota: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := store.Increment(ctx, tenant, "storage_mb", 50); err != nil {
|
|
||||||
t.Fatalf("increment: %v", err)
|
|
||||||
}
|
|
||||||
_, _, status, err := store.Check(ctx, tenant, "storage_mb")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("check: %v", err)
|
|
||||||
}
|
|
||||||
if status != StatusOK {
|
|
||||||
t.Fatalf("bei 50%% erwartet StatusOK, habe %q", status)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := store.Increment(ctx, tenant, "storage_mb", 35); err != nil { // insgesamt 85%
|
|
||||||
t.Fatalf("increment: %v", err)
|
|
||||||
}
|
|
||||||
_, _, status, err = store.Check(ctx, tenant, "storage_mb")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("check: %v", err)
|
|
||||||
}
|
|
||||||
if status != StatusWarning {
|
|
||||||
t.Fatalf("bei 85%% erwartet StatusWarning, habe %q", status)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestCheck_NoQuotaMeansUnlimited(t *testing.T) {
|
|
||||||
store, cleanup := setupTest(t)
|
|
||||||
defer cleanup()
|
|
||||||
ctx := context.Background()
|
|
||||||
tenant := newTenantID()
|
|
||||||
|
|
||||||
if err := store.Increment(ctx, tenant, "api_calls", 1_000_000); err != nil {
|
|
||||||
t.Fatalf("increment: %v", err)
|
|
||||||
}
|
|
||||||
_, _, status, err := store.Check(ctx, tenant, "api_calls")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("check: %v", err)
|
|
||||||
}
|
|
||||||
if status != StatusOK {
|
|
||||||
t.Fatalf("ohne konfigurierte quota erwartet StatusOK, habe %q", status)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestRunPeriodicAggregation_CallsRepeatedly(t *testing.T) {
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
var mu sync.Mutex
|
|
||||||
calls := 0
|
|
||||||
RunPeriodicAggregation(ctx, 20*time.Millisecond, func(ctx context.Context) error {
|
|
||||||
mu.Lock()
|
|
||||||
calls++
|
|
||||||
mu.Unlock()
|
|
||||||
return nil
|
|
||||||
})
|
|
||||||
|
|
||||||
mu.Lock()
|
|
||||||
defer mu.Unlock()
|
|
||||||
if calls < 3 {
|
|
||||||
t.Fatalf("erwartet mehrfache aufrufe innerhalb von 120ms bei 20ms interval, habe %d", calls)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
DROP TABLE IF EXISTS feature_flags;
|
||||||
@@ -0,0 +1,10 @@
|
|||||||
|
-- Feature-Flags zentral je Mandant/Zielgruppe (LIC-02, siehe core-kanban/tickets/LIC-02.md).
|
||||||
|
-- Lebt in der Registry-DB, nicht pro Tenant-Datenbank — Flags sind eine
|
||||||
|
-- Core-weite Konfiguration, keine Mandanten-Geschaeftsdaten.
|
||||||
|
CREATE TABLE feature_flags (
|
||||||
|
key TEXT PRIMARY KEY,
|
||||||
|
enabled BOOLEAN NOT NULL DEFAULT false,
|
||||||
|
rollout_percentage INT NOT NULL DEFAULT 0 CHECK (rollout_percentage BETWEEN 0 AND 100),
|
||||||
|
target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
@@ -1 +0,0 @@
|
|||||||
DROP TABLE IF EXISTS tenant_licenses;
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
-- Lizenzumfang pro Mandant (LIC-01, siehe core-kanban/tickets/LIC-01.md).
|
|
||||||
-- Genau ein Lizenzdatensatz pro Tenant (tenant_id PK) — ein neues Einspielen
|
|
||||||
-- ersetzt den vorherigen Datensatz vollstaendig statt eine Historie zu fuehren.
|
|
||||||
CREATE TABLE tenant_licenses (
|
|
||||||
tenant_id UUID PRIMARY KEY REFERENCES tenants(id),
|
|
||||||
plan TEXT NOT NULL,
|
|
||||||
modules TEXT[] NOT NULL,
|
|
||||||
issued_at TIMESTAMPTZ NOT NULL,
|
|
||||||
valid_until TIMESTAMPTZ NOT NULL,
|
|
||||||
raw_key TEXT NOT NULL,
|
|
||||||
installed_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
||||||
);
|
|
||||||
@@ -0,0 +1,2 @@
|
|||||||
|
DROP TABLE IF EXISTS module_credentials;
|
||||||
|
DROP TABLE IF EXISTS modules;
|
||||||
@@ -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,2 +0,0 @@
|
|||||||
DROP TABLE IF EXISTS usage_quotas;
|
|
||||||
DROP TABLE IF EXISTS usage_counters;
|
|
||||||
@@ -1,15 +0,0 @@
|
|||||||
-- Nutzungszaehler & Quotas je Tenant (LIC-03, siehe core-kanban/tickets/LIC-03.md).
|
|
||||||
CREATE TABLE usage_counters (
|
|
||||||
tenant_id UUID NOT NULL,
|
|
||||||
metric TEXT NOT NULL,
|
|
||||||
value BIGINT NOT NULL DEFAULT 0,
|
|
||||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
||||||
PRIMARY KEY (tenant_id, metric)
|
|
||||||
);
|
|
||||||
|
|
||||||
CREATE TABLE usage_quotas (
|
|
||||||
tenant_id UUID NOT NULL,
|
|
||||||
metric TEXT NOT NULL,
|
|
||||||
limit_value BIGINT NOT NULL,
|
|
||||||
PRIMARY KEY (tenant_id, metric)
|
|
||||||
);
|
|
||||||
Reference in New Issue
Block a user