internal/usage: Store.Increment aktualisiert Zaehlerstaende ueber ein einziges atomares SQL-UPSERT (value = value + delta) statt Read-Modify-Write in Go — haelt Zaehlerstaende bei parallelen Schreibzugriffen konsistent (Akzeptanzkriterium 1), ganz ohne Anwendungs-Lock. Quotas sind Konfiguration (usage_quotas-Tabelle je Tenant+Metrik), kein Hardcode. Check/Enforce leiten aus Zaehlerstand + Quota eine definierte Reaktion ab (StatusOK/Warning bei 80%/Exceeded, Akzeptanzkriterium 2) — Enforce ruft eine uebergebene Reaction-Funktion auf, wenn der Status nicht OK ist; die konkrete Sperr-/Benachrichtigungslogik bleibt beim Aufrufer (z.B. TEN-02 vor Benutzeranlage), Enforce garantiert nur zuverlaessiges Ausloesen. Fehlende Quota-Konfiguration bedeutet unbegrenzt (StatusOK), kein Fehler. RunPeriodicAggregation ist das Aggregations-Grundgerüst (Akzeptanzkriterium 1: "periodisch aggregiert") — dieselbe In-Prozess-Worker-Goroutine-Konvention wie internal/tenant.Lifecycle.RunSweeper. Die konkrete Aggregationsquelle (Zeilen zaehlen in Modul-Tabellen) haengt vom jeweiligen Modul ab und ist nicht Teil dieser Kachel. Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS): 1. Quota-Ueberschreitung automatisiert erkannt, definierte Reaktion ausgeloest — TestEnforce_TriggersReactionOnExceeded: Reaction-Callback wird mit StatusExceeded aufgerufen. PASS. 2. Aggregationsjob liefert bei parallelen Schreibzugriffen konsistente Zaehlerstaende — TestIncrement_ConsistentUnderConcurrentWrites: 50 nebenlaeufige Increments, Endstand exakt 50 (kein Lost Update). PASS. 3. Zaehlerstand eines Tenants beeinflusst nicht den eines anderen — TestIncrement_IsolatedBetweenTenants. PASS. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
145 lines
4.7 KiB
Go
145 lines
4.7 KiB
Go
// 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
|
|
}
|