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>
207 lines
5.4 KiB
Go
207 lines
5.4 KiB
Go
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)
|
|
}
|
|
}
|