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/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) +);