API-03: zentrales-rate-limiting-api-gateway-schicht (postgres-basierter shared state)
This commit is contained in:
@@ -0,0 +1,168 @@
|
||||
package ratelimit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func setupTest(t *testing.T) (*pgxpool.Pool, 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 rate_limit_configs (
|
||||
key TEXT PRIMARY KEY, limit_value INT NOT NULL, window_seconds INT NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS rate_limit_counters (
|
||||
key TEXT NOT NULL, window_start TIMESTAMPTZ NOT NULL, count INT NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (key, window_start)
|
||||
);
|
||||
`); err != nil {
|
||||
t.Fatalf("schema: %v", err)
|
||||
}
|
||||
|
||||
cleanup := func() { pool.Close() }
|
||||
return pool, cleanup
|
||||
}
|
||||
|
||||
func newTestKey() string {
|
||||
return fmt.Sprintf("test-key-%d", time.Now().UnixNano())
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 3 + Pruefung 2: Ueberschreitung des Limits wird
|
||||
// nachweislich erkannt.
|
||||
func TestAllow_RejectsAfterLimitExceeded(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
store := NewStore(pool)
|
||||
ctx := context.Background()
|
||||
key := newTestKey()
|
||||
|
||||
if err := store.SetLimit(ctx, key, 3, time.Minute); err != nil {
|
||||
t.Fatalf("setlimit: %v", err)
|
||||
}
|
||||
|
||||
for i := 0; i < 3; i++ {
|
||||
res, err := store.Allow(ctx, key)
|
||||
if err != nil {
|
||||
t.Fatalf("allow %d: %v", i, err)
|
||||
}
|
||||
if !res.Allowed {
|
||||
t.Fatalf("anfrage %d haette erlaubt sein muessen, limit=3", i)
|
||||
}
|
||||
}
|
||||
|
||||
res, err := store.Allow(ctx, key)
|
||||
if err != nil {
|
||||
t.Fatalf("allow (4.): %v", err)
|
||||
}
|
||||
if res.Allowed {
|
||||
t.Fatal("4. anfrage haette abgelehnt werden muessen (limit=3)")
|
||||
}
|
||||
if res.RetryAfter <= 0 {
|
||||
t.Fatalf("erwartet positive retry-after-dauer, habe %v", res.RetryAfter)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 1 / Pruefung 3: Konfigurationsaenderung wirkt sich
|
||||
// NICHT auf andere Schluessel (Tenants) aus.
|
||||
func TestSetLimit_DoesNotAffectOtherKeys(t *testing.T) {
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
store := NewStore(pool)
|
||||
ctx := context.Background()
|
||||
keyA := newTestKey()
|
||||
keyB := newTestKey()
|
||||
|
||||
if err := store.SetLimit(ctx, keyA, 1, time.Minute); err != nil {
|
||||
t.Fatalf("setlimit a: %v", err)
|
||||
}
|
||||
|
||||
resA1, _ := store.Allow(ctx, keyA)
|
||||
resA2, _ := store.Allow(ctx, keyA)
|
||||
if !resA1.Allowed || resA2.Allowed {
|
||||
t.Fatalf("key a: erwartet erlaubt,abgelehnt, habe %v,%v", resA1.Allowed, resA2.Allowed)
|
||||
}
|
||||
|
||||
// keyB hat keine eigene Konfiguration -> DefaultLimit, unbeeinflusst von keyA.
|
||||
resB, err := store.Allow(ctx, keyB)
|
||||
if err != nil {
|
||||
t.Fatalf("allow b: %v", err)
|
||||
}
|
||||
if !resB.Allowed {
|
||||
t.Fatal("key b haette trotz limit-ueberschreitung bei key a erlaubt sein muessen")
|
||||
}
|
||||
if resB.Limit != DefaultLimit {
|
||||
t.Fatalf("key b limit = %d, want default %d", resB.Limit, DefaultLimit)
|
||||
}
|
||||
}
|
||||
|
||||
// Akzeptanzkriterium 2 + Pruefung 1: zwei "Dienstinstanzen" (zwei
|
||||
// unabhaengige Pools/Store-Objekte) gegen dieselbe Datenbank setzen
|
||||
// dasselbe Limit korrekt durch — kein In-Process-Zustand kann das leisten.
|
||||
func TestAllow_TwoInstancesShareStateCorrectly(t *testing.T) {
|
||||
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||
if adminDSN == "" {
|
||||
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||
}
|
||||
pool, cleanup := setupTest(t)
|
||||
defer cleanup()
|
||||
|
||||
ctx := context.Background()
|
||||
pool2, err := pgxpool.New(ctx, adminDSN)
|
||||
if err != nil {
|
||||
t.Fatalf("zweiter pool (simuliert zweite instanz): %v", err)
|
||||
}
|
||||
defer pool2.Close()
|
||||
|
||||
storeInstance1 := NewStore(pool)
|
||||
storeInstance2 := NewStore(pool2)
|
||||
|
||||
key := newTestKey()
|
||||
const limit = 20
|
||||
if err := storeInstance1.SetLimit(ctx, key, limit, time.Minute); err != nil {
|
||||
t.Fatalf("setlimit: %v", err)
|
||||
}
|
||||
|
||||
const totalRequests = 50
|
||||
var allowedCount int64
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < totalRequests; i++ {
|
||||
wg.Add(1)
|
||||
store := storeInstance1
|
||||
if i%2 == 0 {
|
||||
store = storeInstance2 // Haelfte der Anfragen "kommt" von der zweiten Instanz
|
||||
}
|
||||
go func(s *Store) {
|
||||
defer wg.Done()
|
||||
res, err := s.Allow(ctx, key)
|
||||
if err != nil {
|
||||
t.Errorf("allow: %v", err)
|
||||
return
|
||||
}
|
||||
if res.Allowed {
|
||||
atomic.AddInt64(&allowedCount, 1)
|
||||
}
|
||||
}(store)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if allowedCount != limit {
|
||||
t.Fatalf("erlaubte anfragen ueber beide instanzen = %d, want exakt %d (geteilter zustand)", allowedCount, limit)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user