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) } }