diff --git a/internal/license/license.go b/internal/license/license.go new file mode 100644 index 0000000..34dab1b --- /dev/null +++ b/internal/license/license.go @@ -0,0 +1,106 @@ +// 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 +} diff --git a/internal/license/license_test.go b/internal/license/license_test.go new file mode 100644 index 0000000..fb9316a --- /dev/null +++ b/internal/license/license_test.go @@ -0,0 +1,106 @@ +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) + } + } +} diff --git a/internal/license/store.go b/internal/license/store.go new file mode 100644 index 0000000..a634738 --- /dev/null +++ b/internal/license/store.go @@ -0,0 +1,88 @@ +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 +} diff --git a/internal/license/store_test.go b/internal/license/store_test.go new file mode 100644 index 0000000..dec430b --- /dev/null +++ b/internal/license/store_test.go @@ -0,0 +1,177 @@ +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) + } +} diff --git a/internal/usage/aggregation.go b/internal/usage/aggregation.go new file mode 100644 index 0000000..e0e5bc6 --- /dev/null +++ b/internal/usage/aggregation.go @@ -0,0 +1,32 @@ +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) + } + } + } +} diff --git a/internal/usage/usage.go b/internal/usage/usage.go new file mode 100644 index 0000000..be6aeba --- /dev/null +++ b/internal/usage/usage.go @@ -0,0 +1,144 @@ +// 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 +} diff --git a/internal/usage/usage_test.go b/internal/usage/usage_test.go new file mode 100644 index 0000000..8906675 --- /dev/null +++ b/internal/usage/usage_test.go @@ -0,0 +1,206 @@ +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) + } +} diff --git a/migrations/0004_tenant_licenses.down.sql b/migrations/0004_tenant_licenses.down.sql new file mode 100644 index 0000000..d9728f1 --- /dev/null +++ b/migrations/0004_tenant_licenses.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS tenant_licenses; diff --git a/migrations/0004_tenant_licenses.up.sql b/migrations/0004_tenant_licenses.up.sql new file mode 100644 index 0000000..9872cbf --- /dev/null +++ b/migrations/0004_tenant_licenses.up.sql @@ -0,0 +1,12 @@ +-- 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() +); diff --git a/migrations/0005_usage_counters.down.sql b/migrations/0005_usage_counters.down.sql new file mode 100644 index 0000000..7e267ae --- /dev/null +++ b/migrations/0005_usage_counters.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS usage_quotas; +DROP TABLE IF EXISTS usage_counters; diff --git a/migrations/0005_usage_counters.up.sql b/migrations/0005_usage_counters.up.sql new file mode 100644 index 0000000..0658918 --- /dev/null +++ b/migrations/0005_usage_counters.up.sql @@ -0,0 +1,15 @@ +-- 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) +);