169 lines
4.3 KiB
Go
169 lines
4.3 KiB
Go
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)
|
|
}
|
|
}
|