Compare commits

..
9 changed files with 389 additions and 511 deletions
+158
View File
@@ -0,0 +1,158 @@
// Package mtls implementiert Core API-09: gegenseitige TLS-Authentifizierung
// zwischen Core/DMS/Mail/Archive/Workflow/AI/Connect-Instanzen fuer
// Stufe-2/3-Installationen (siehe SKALIERUNGSKONZEPT.md) — schuetzt die
// VERTRAULICHKEIT des Transportwegs, dort wo API-05 (asymmetrische
// JWT-Signatur) und API-02 (Service-Credential) nur Authentizitaet/
// Integritaet, aber keine Transportverschluesselung zwischen getrennten
// Hosts absichern.
//
// Zertifikatsverteilung (Akzeptanzkriterium 1): Authority.IssueCert liefert
// PEM-kodiertes Zertifikat + privaten Schluessel, die ueber denselben Weg
// wie andere Secrets verteilt werden (Umgebungsvariablen/Secret-Store,
// niemals im Code) — automatisierbar, da IssueCert ein reiner
// Funktionsaufruf ohne manuelle Schritte ist.
package mtls
import (
"crypto/ed25519"
"crypto/rand"
"crypto/tls"
"crypto/x509"
"crypto/x509/pkix"
"encoding/pem"
"fmt"
"math/big"
"net"
"time"
)
// Authority ist eine interne Zertifizierungsstelle fuer Modul-zu-Modul-
// mTLS. Haelt — analog zu internal/moduletrust.KeyManager (API-05) — ALLE
// noch gueltigen historischen CA-Zertifikate im Vertrauensspeicher, damit
// eine Rotation bereits ausgestellte Leaf-Zertifikate nicht ungueltig macht
// (Akzeptanzkriterium 3: Rotation ohne Ausfallzeit).
type Authority struct {
current caGeneration
trustPool *x509.CertPool
generations []caGeneration
}
type caGeneration struct {
cert *x509.Certificate
key ed25519.PrivateKey
}
// NewAuthority erzeugt eine frische interne CA.
func NewAuthority() (*Authority, error) {
gen, err := newCAGeneration()
if err != nil {
return nil, err
}
pool := x509.NewCertPool()
pool.AddCert(gen.cert)
return &Authority{current: gen, trustPool: pool, generations: []caGeneration{gen}}, nil
}
func newCAGeneration() (caGeneration, error) {
pub, priv, err := ed25519.GenerateKey(rand.Reader)
if err != nil {
return caGeneration{}, fmt.Errorf("ca-schluesselpaar erzeugen: %w", err)
}
serial, err := rand.Int(rand.Reader, big.NewInt(1<<62))
if err != nil {
return caGeneration{}, fmt.Errorf("seriennummer erzeugen: %w", err)
}
template := &x509.Certificate{
SerialNumber: serial,
Subject: pkix.Name{CommonName: "nexarch-internal-ca"},
NotBefore: time.Now().Add(-time.Minute),
NotAfter: time.Now().Add(5 * 365 * 24 * time.Hour),
KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageCRLSign,
BasicConstraintsValid: true,
IsCA: true,
}
der, err := x509.CreateCertificate(rand.Reader, template, template, pub, priv)
if err != nil {
return caGeneration{}, fmt.Errorf("ca-zertifikat erstellen: %w", err)
}
cert, err := x509.ParseCertificate(der)
if err != nil {
return caGeneration{}, fmt.Errorf("ca-zertifikat parsen: %w", err)
}
return caGeneration{cert: cert, key: priv}, nil
}
// Rotate erzeugt eine NEUE CA-Generation fuer zukuenftige IssueCert-Aufrufe,
// behaelt aber ALLE bisherigen CA-Zertifikate im Vertrauensspeicher —
// bereits ausgestellte Leaf-Zertifikate bleiben dadurch gueltig, eine
// laufende mTLS-Verbindung wird durch Rotate NICHT unterbrochen
// (Akzeptanzkriterium 3 / Pruefung 2).
func (a *Authority) Rotate() error {
gen, err := newCAGeneration()
if err != nil {
return err
}
a.trustPool.AddCert(gen.cert)
a.generations = append(a.generations, gen)
a.current = gen
return nil
}
// TrustPool liefert den Vertrauensspeicher mit ALLEN (auch historischen,
// noch nicht abgelaufenen) CA-Zertifikaten — Grundlage der
// Server-seitigen Client-Zertifikatspruefung (Akzeptanzkriterium 2).
func (a *Authority) TrustPool() *x509.CertPool {
return a.trustPool
}
// IssuedCert ist ein ausgestelltes Leaf-Zertifikat inklusive privatem
// Schluessel, PEM-kodiert zur Verteilung (Akzeptanzkriterium 1).
type IssuedCert struct {
CertPEM []byte
KeyPEM []byte
}
// IssueCert stellt ein Leaf-Zertifikat fuer EINE Modul-Instanz aus, signiert
// mit der AKTUELLEN CA-Generation (Akzeptanzkriterium 1).
func (a *Authority) IssueCert(commonName string, validity time.Duration) (IssuedCert, error) {
pub, priv, err := ed25519.GenerateKey(rand.Reader)
if err != nil {
return IssuedCert{}, fmt.Errorf("leaf-schluesselpaar erzeugen: %w", err)
}
serial, err := rand.Int(rand.Reader, big.NewInt(1<<62))
if err != nil {
return IssuedCert{}, fmt.Errorf("seriennummer erzeugen: %w", err)
}
template := &x509.Certificate{
SerialNumber: serial,
Subject: pkix.Name{CommonName: commonName},
NotBefore: time.Now().Add(-time.Minute),
NotAfter: time.Now().Add(validity),
KeyUsage: x509.KeyUsageDigitalSignature,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth, x509.ExtKeyUsageServerAuth},
// DNSNames/IPAddresses decken lokale Testumgebungen ab (127.0.0.1,
// localhost) — in echten Stufe-2/3-Installationen entspricht
// commonName dem tatsaechlichen internen Hostnamen der Instanz.
DNSNames: []string{commonName, "localhost"},
IPAddresses: []net.IP{net.IPv4(127, 0, 0, 1), net.IPv6loopback},
}
der, err := x509.CreateCertificate(rand.Reader, template, a.current.cert, pub, a.current.key)
if err != nil {
return IssuedCert{}, fmt.Errorf("leaf-zertifikat erstellen: %w", err)
}
certPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})
keyBytes, err := x509.MarshalPKCS8PrivateKey(priv)
if err != nil {
return IssuedCert{}, fmt.Errorf("leaf-schluessel serialisieren: %w", err)
}
keyPEM := pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: keyBytes})
return IssuedCert{CertPEM: certPEM, KeyPEM: keyPEM}, nil
}
// TLSCertificate wandelt ein IssuedCert in ein tls.Certificate um, wie es
// tls.Config.Certificates erwartet.
func (c IssuedCert) TLSCertificate() (tls.Certificate, error) {
return tls.X509KeyPair(c.CertPEM, c.KeyPEM)
}
+46
View File
@@ -0,0 +1,46 @@
package mtls
import "crypto/tls"
// Mode legt fest, ob mTLS erzwungen wird. Fuer Stufe-1-Einzel-LXC ist mTLS
// laut Ticket "optional/nicht zwingend" (lokale Kommunikation) — ModeOff
// liefert einen ganz normalen TLS-Server OHNE Client-Zertifikatspruefung,
// damit Stufe-1-Betrieb unveraendert weiterlaeuft (Akzeptanzkriterium 3 /
// Pruefung 3).
type Mode string
const (
ModeOff Mode = "off" // Stufe 1: kein mTLS-Zwang
ModeRequired Mode = "required" // Stufe 2/3: Client-Zertifikat zwingend
)
// ServerTLSConfig liefert die tls.Config fuer eine Modul-Instanz, die
// eingehende Verbindungen ANDERER Modul-Instanzen annimmt.
//
// - ModeRequired: verlangt UND verifiziert ein Client-Zertifikat gegen
// authority.TrustPool() (Akzeptanzkriterium 2) — eine Verbindung ohne
// gueltiges Zertifikat wird vom TLS-Handshake selbst abgelehnt, bevor
// irgendein Anwendungscode erreicht wird.
// - ModeOff: normales TLS ohne Client-Zertifikatspruefung
// (Akzeptanzkriterium 3 / Pruefung 3: Stufe-1-Betrieb funktioniert
// weiterhin ohne mTLS-Zwang).
func ServerTLSConfig(mode Mode, serverCert tls.Certificate, authority *Authority) *tls.Config {
cfg := &tls.Config{Certificates: []tls.Certificate{serverCert}}
if mode == ModeRequired {
cfg.ClientAuth = tls.RequireAndVerifyClientCert
cfg.ClientCAs = authority.TrustPool()
}
return cfg
}
// ClientTLSConfig liefert die tls.Config fuer eine Modul-Instanz, die eine
// Verbindung zu einer ANDEREN Modul-Instanz aufbaut — praesentiert das
// eigene Zertifikat und vertraut Gegenstellen, die von derselben Authority
// signiert wurden.
func ClientTLSConfig(mode Mode, clientCert tls.Certificate, authority *Authority) *tls.Config {
cfg := &tls.Config{RootCAs: authority.TrustPool()}
if mode == ModeRequired {
cfg.Certificates = []tls.Certificate{clientCert}
}
return cfg
}
+185
View File
@@ -0,0 +1,185 @@
package mtls
import (
"crypto/tls"
"io"
"net/http"
"net/http/httptest"
"testing"
"time"
)
func newTLSServer(t *testing.T, cfg *tls.Config) *httptest.Server {
t.Helper()
server := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("ok"))
}))
server.TLS = cfg
server.StartTLS()
return server
}
func fetch(t *testing.T, url string, clientCfg *tls.Config) (int, error) {
t.Helper()
client := &http.Client{
Transport: &http.Transport{TLSClientConfig: clientCfg},
Timeout: 3 * time.Second,
}
resp, err := client.Get(url)
if err != nil {
return 0, err
}
defer resp.Body.Close()
_, _ = io.ReadAll(resp.Body)
return resp.StatusCode, nil
}
// Akzeptanzkriterium 1: Zertifikatsverteilung ist ein reiner
// Funktionsaufruf, das Ergebnis ist gueltiges, wiederverwendbares
// PEM-Material.
func TestIssueCert_ProducesValidPEMRoundTrip(t *testing.T) {
authority, err := NewAuthority()
if err != nil {
t.Fatalf("newauthority: %v", err)
}
issued, err := authority.IssueCert("dms", time.Hour)
if err != nil {
t.Fatalf("issuecert: %v", err)
}
if _, err := issued.TLSCertificate(); err != nil {
t.Fatalf("ausgestelltes zertifikat nicht als tls.Certificate ladbar: %v", err)
}
}
// Akzeptanzkriterium 2 + Pruefung 1: eine Verbindung ohne gueltiges
// Client-Zertifikat wird abgelehnt, eine mit gueltigem angenommen.
func TestServerTLSConfig_RejectsMissingClientCert(t *testing.T) {
authority, err := NewAuthority()
if err != nil {
t.Fatalf("newauthority: %v", err)
}
serverIssued, err := authority.IssueCert("core", time.Hour)
if err != nil {
t.Fatalf("server-zertifikat ausstellen: %v", err)
}
serverCert, err := serverIssued.TLSCertificate()
if err != nil {
t.Fatalf("server-zertifikat laden: %v", err)
}
server := newTLSServer(t, ServerTLSConfig(ModeRequired, serverCert, authority))
defer server.Close()
// Client OHNE Zertifikat, aber mit korrektem RootCA-Vertrauen fuer den
// Server — der Handshake muss trotzdem an der fehlenden Client-Auth
// scheitern.
_, err = fetch(t, server.URL, &tls.Config{RootCAs: authority.TrustPool()})
if err == nil {
t.Fatal("erwartet fehlschlagenden handshake ohne client-zertifikat")
}
clientIssued, err := authority.IssueCert("mail", time.Hour)
if err != nil {
t.Fatalf("client-zertifikat ausstellen: %v", err)
}
clientCert, err := clientIssued.TLSCertificate()
if err != nil {
t.Fatalf("client-zertifikat laden: %v", err)
}
status, err := fetch(t, server.URL, ClientTLSConfig(ModeRequired, clientCert, authority))
if err != nil {
t.Fatalf("erwartet erfolgreiche verbindung mit gueltigem client-zertifikat: %v", err)
}
if status != http.StatusOK {
t.Fatalf("status = %d, want 200", status)
}
}
// Akzeptanzkriterium 3 + Pruefung 2: Rotation der CA unterbricht laufenden
// Betrieb nicht — ein VOR der Rotation ausgestelltes Leaf-Zertifikat
// funktioniert danach weiter, UND ein NACH der Rotation neu ausgestelltes
// funktioniert ebenfalls, gegen denselben, weiterlaufenden Server (kein
// Neustart noetig, da TrustPool() denselben *x509.CertPool zurueckliefert,
// den Rotate() erweitert).
func TestAuthority_RotateWithoutDowntime(t *testing.T) {
authority, err := NewAuthority()
if err != nil {
t.Fatalf("newauthority: %v", err)
}
serverIssued, err := authority.IssueCert("core", time.Hour)
if err != nil {
t.Fatalf("server-zertifikat ausstellen: %v", err)
}
serverCert, err := serverIssued.TLSCertificate()
if err != nil {
t.Fatalf("server-zertifikat laden: %v", err)
}
server := newTLSServer(t, ServerTLSConfig(ModeRequired, serverCert, authority))
defer server.Close() // EIN Server-Prozess ueber die gesamte Rotation hinweg — kein Neustart.
preRotationIssued, err := authority.IssueCert("dms", time.Hour)
if err != nil {
t.Fatalf("client-zertifikat (vor rotation) ausstellen: %v", err)
}
preRotationCert, _ := preRotationIssued.TLSCertificate()
if err := authority.Rotate(); err != nil {
t.Fatalf("rotate: %v", err)
}
// VOR der Rotation ausgestelltes Zertifikat funktioniert WEITERHIN, ohne
// dass der Server neu gestartet wurde.
status, err := fetch(t, server.URL, ClientTLSConfig(ModeRequired, preRotationCert, authority))
if err != nil {
t.Fatalf("vor-rotation-zertifikat haette weiterhin funktionieren sollen: %v", err)
}
if status != http.StatusOK {
t.Fatalf("status = %d, want 200 (vor-rotation-zertifikat)", status)
}
// NACH der Rotation neu ausgestelltes Zertifikat funktioniert ebenfalls,
// gegen DENSELBEN laufenden Server.
postRotationIssued, err := authority.IssueCert("mail", time.Hour)
if err != nil {
t.Fatalf("client-zertifikat (nach rotation) ausstellen: %v", err)
}
postRotationCert, _ := postRotationIssued.TLSCertificate()
status, err = fetch(t, server.URL, ClientTLSConfig(ModeRequired, postRotationCert, authority))
if err != nil {
t.Fatalf("nach-rotation-zertifikat haette funktionieren sollen: %v", err)
}
if status != http.StatusOK {
t.Fatalf("status = %d, want 200 (nach-rotation-zertifikat)", status)
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Stufe-1-Betrieb (ModeOff) funktioniert
// weiterhin OHNE mTLS-Zwang — ein Client ganz ohne Zertifikat kommt durch.
func TestServerTLSConfig_ModeOffWorksWithoutClientCert(t *testing.T) {
authority, err := NewAuthority()
if err != nil {
t.Fatalf("newauthority: %v", err)
}
serverIssued, err := authority.IssueCert("core", time.Hour)
if err != nil {
t.Fatalf("server-zertifikat ausstellen: %v", err)
}
serverCert, err := serverIssued.TLSCertificate()
if err != nil {
t.Fatalf("server-zertifikat laden: %v", err)
}
server := newTLSServer(t, ServerTLSConfig(ModeOff, serverCert, authority))
defer server.Close()
status, err := fetch(t, server.URL, ClientTLSConfig(ModeOff, tls.Certificate{}, authority))
if err != nil {
t.Fatalf("stufe-1-betrieb (ModeOff) haette ohne client-zertifikat funktionieren sollen: %v", err)
}
if status != http.StatusOK {
t.Fatalf("status = %d, want 200", status)
}
}
-64
View File
@@ -1,64 +0,0 @@
package ratelimit
import (
"fmt"
"math"
"net/http"
"gitea.perlbach24.de/scripte/nexarch/internal/apiserver"
)
// KeyFunc bestimmt den Rate-Limit-Schluessel fuer eine Anfrage — typischerweise
// der Tenant-Slug aus apiserver.RequestContext, siehe KeyByTenant. Als
// eigenstaendiger Typ austauschbar (z.B. spaeter API-Token-basiert), ohne
// internal/apiserver aendern zu muessen (Kein Umbau angrenzender Bereiche).
type KeyFunc func(r *http.Request) (string, bool)
// KeyByTenant liest den Tenant aus dem von internal/apiserver.authAndTenantContext
// gesetzten RequestContext (IAM-02/API-01) — dieselbe Middleware-Kette,
// keine zweite Authentifizierung.
func KeyByTenant(r *http.Request) (string, bool) {
rc, ok := apiserver.FromContext(r.Context())
if !ok || rc.TenantSlug == "" {
return "", false
}
return rc.TenantSlug, true
}
// Middleware setzt sich VOR den eigentlichen Handler (nach Auth+Tenant-
// Kontext, siehe KeyByTenant) und lehnt Anfragen ueber dem konfigurierten
// Limit mit 429 + Retry-After ab (Akzeptanzkriterium 3).
func Middleware(store *Store, keyFn KeyFunc, next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
key, ok := keyFn(r)
if !ok {
// Kein Schluessel ermittelbar (z.B. kein Tenant-Kontext) — die
// Auth-Pruefung selbst ist Sache von authAndTenantContext, hier
// wird nur nicht limitiert, wenn schon kein Kontext vorliegt.
next(w, r)
return
}
result, err := store.Allow(r.Context(), key)
if err != nil {
apiserver.WriteError(w, http.StatusInternalServerError, "rate_limit_error", "rate-limit-pruefung fehlgeschlagen")
return
}
w.Header().Set("X-RateLimit-Limit", fmt.Sprintf("%d", result.Limit))
w.Header().Set("X-RateLimit-Remaining", fmt.Sprintf("%d", result.Remaining))
if !result.Allowed {
retryAfterSeconds := int(math.Ceil(result.RetryAfter.Seconds()))
if retryAfterSeconds < 1 {
retryAfterSeconds = 1
}
w.Header().Set("Retry-After", fmt.Sprintf("%d", retryAfterSeconds))
apiserver.WriteError(w, http.StatusTooManyRequests, "rate_limit_exceeded",
"zu viele anfragen, bitte spaeter erneut versuchen")
return
}
next(w, r)
}
}
-136
View File
@@ -1,136 +0,0 @@
package ratelimit
import (
"context"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/apiserver"
"gitea.perlbach24.de/scripte/nexarch/internal/auth"
)
func setupMiddlewareTest(t *testing.T) (*apiserver.Server, *auth.TokenIssuer, *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 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)
}
issuer := auth.NewTokenIssuer("test-secret-nur-fuer-tests")
srv := apiserver.NewServer(issuer)
store := NewStore(pool)
cleanup := func() { pool.Close() }
return srv, issuer, store, cleanup
}
// Akzeptanzkriterium 3 + Pruefung 2: Ueberschreitung liefert 429 mit
// Retry-After-Header, ueber die vollstaendige Middleware-Kette (Auth ->
// Tenant-Kontext -> Rate-Limit -> Handler) hinweg.
func TestMiddleware_Returns429WithRetryAfterOnExceeded(t *testing.T) {
srv, issuer, store, cleanup := setupMiddlewareTest(t)
defer cleanup()
ctx := context.Background()
tenantSlug := newTestKey()
if err := store.SetLimit(ctx, tenantSlug, 1, time.Minute); err != nil {
t.Fatalf("setlimit: %v", err)
}
called := 0
srv.HandleV1("/limited", Middleware(store, KeyByTenant, func(w http.ResponseWriter, r *http.Request) {
called++
w.WriteHeader(http.StatusOK)
}))
token, err := issuer.Issue("user-1", tenantSlug)
if err != nil {
t.Fatalf("issue: %v", err)
}
req1 := httptest.NewRequest(http.MethodGet, "/api/v1/limited", nil)
req1.AddCookie(&http.Cookie{Name: auth.CookieName, Value: token})
rec1 := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec1, req1)
if rec1.Code != http.StatusOK {
t.Fatalf("1. anfrage: status = %d, want 200", rec1.Code)
}
req2 := httptest.NewRequest(http.MethodGet, "/api/v1/limited", nil)
req2.AddCookie(&http.Cookie{Name: auth.CookieName, Value: token})
rec2 := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec2, req2)
if rec2.Code != http.StatusTooManyRequests {
t.Fatalf("2. anfrage: status = %d, want 429", rec2.Code)
}
if rec2.Header().Get("Retry-After") == "" {
t.Fatal("erwartet gesetzten Retry-After-Header bei 429")
}
if called != 1 {
t.Fatalf("handler wurde %d mal aufgerufen, want genau 1 (2. anfrage haette nicht durchgereicht werden duerfen)", called)
}
}
// Akzeptanzkriterium 1: das Limit ist pro Tenant konfigurierbar — ein
// anderer Tenant (anderer Key) bleibt von der Ausschoepfung unberuehrt.
func TestMiddleware_DifferentTenantsHaveIndependentLimits(t *testing.T) {
srv, issuer, store, cleanup := setupMiddlewareTest(t)
defer cleanup()
ctx := context.Background()
tenantA := newTestKey()
tenantB := newTestKey()
if err := store.SetLimit(ctx, tenantA, 1, time.Minute); err != nil {
t.Fatalf("setlimit a: %v", err)
}
if err := store.SetLimit(ctx, tenantB, 1, time.Minute); err != nil {
t.Fatalf("setlimit b: %v", err)
}
srv.HandleV1("/limited2", Middleware(store, KeyByTenant, func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
}))
tokenA, _ := issuer.Issue("user-a", tenantA)
tokenB, _ := issuer.Issue("user-b", tenantB)
// Tenant A schoepft sein Limit aus.
reqA := httptest.NewRequest(http.MethodGet, "/api/v1/limited2", nil)
reqA.AddCookie(&http.Cookie{Name: auth.CookieName, Value: tokenA})
recA := httptest.NewRecorder()
srv.Handler().ServeHTTP(recA, reqA)
if recA.Code != http.StatusOK {
t.Fatalf("tenant a, 1. anfrage: status = %d, want 200", recA.Code)
}
// Tenant B ist von der Ausschoepfung bei A unberuehrt.
reqB := httptest.NewRequest(http.MethodGet, "/api/v1/limited2", nil)
reqB.AddCookie(&http.Cookie{Name: auth.CookieName, Value: tokenB})
recB := httptest.NewRecorder()
srv.Handler().ServeHTTP(recB, reqB)
if recB.Code != http.StatusOK {
t.Fatalf("tenant b haette trotz ausgeschoepftem limit von tenant a erlaubt sein muessen, status = %d", recB.Code)
}
}
-126
View File
@@ -1,126 +0,0 @@
// Package ratelimit implementiert Core API-03: eine zentrale Rate-Limiting-
// Schicht fuer die gesamte Core-API mit GETEILTEM, EXTERNEM Zustand in
// Postgres — kein In-Process-Zaehler (siehe "Bekannte Fehler vermeiden" im
// Ticket: dasselbe Risiko wie bei IAM-07s In-Memory-Lockout). Fixed-Window-
// Algorithmus: einfach, korrekt unter nebenlaeufigem Zugriff (atomares
// UPSERT wie internal/usage.Store.Increment), und fuer ein API-Gateway
// ausreichend praezise — kein Sliding-Window-Overhead noetig.
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// DefaultLimit/DefaultWindow gelten fuer jeden Schluessel (Tenant oder
// API-Token), fuer den keine eigene Konfiguration existiert
// (Akzeptanzkriterium 1: "konfigurierbar", nicht "verpflichtend konfiguriert").
const (
DefaultLimit = 100
DefaultWindow = time.Minute
)
// Config ist das je Schluessel konfigurierbare Limit.
type Config struct {
Limit int
Window time.Duration
}
// Store haelt Konfiguration UND Zaehlerstand in Postgres — beides ueber
// dieselbe Verbindung erreichbar, damit mehrere Dienstinstanzen denselben
// Zustand sehen (Akzeptanzkriterium 2).
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
// SetLimit setzt die Konfiguration fuer EINEN Schluessel (z.B. einen
// Tenant-Slug oder eine API-Token-ID) — wirkt sich nicht auf andere
// Schluessel aus (Akzeptanzkriterium 1 / Pruefung 3).
func (s *Store) SetLimit(ctx context.Context, key string, limit int, window time.Duration) error {
_, err := s.pool.Exec(ctx, `
INSERT INTO rate_limit_configs (key, limit_value, window_seconds)
VALUES ($1, $2, $3)
ON CONFLICT (key) DO UPDATE SET limit_value = $2, window_seconds = $3
`, key, limit, int(window.Seconds()))
if err != nil {
return fmt.Errorf("rate-limit-konfiguration speichern: %w", err)
}
return nil
}
func (s *Store) getConfig(ctx context.Context, key string) (Config, error) {
var limit, windowSeconds int
err := s.pool.QueryRow(ctx, `
SELECT limit_value, window_seconds FROM rate_limit_configs WHERE key = $1
`, key).Scan(&limit, &windowSeconds)
if err != nil {
if err == pgx.ErrNoRows {
return Config{Limit: DefaultLimit, Window: DefaultWindow}, nil
}
return Config{}, fmt.Errorf("rate-limit-konfiguration lesen: %w", err)
}
return Config{Limit: limit, Window: time.Duration(windowSeconds) * time.Second}, nil
}
// Result ist der Ausgang einer Allow-Pruefung (Akzeptanzkriterium 3).
type Result struct {
Allowed bool
Limit int
Remaining int
RetryAfter time.Duration
}
// Allow erhoeht den Zaehler fuer (key, aktuelles Zeitfenster) ATOMAR ueber
// ein einziges UPSERT (dasselbe Muster wie internal/usage.Store.Increment)
// und vergleicht das Ergebnis gegen die konfigurierte Grenze — kein
// Lesen-Erhoehen-Schreiben in Go, damit zwei Dienstinstanzen, die
// gleichzeitig gegen dieselbe Datenbank inkrementieren, sich niemals
// gegenseitig ueberschreiben (Akzeptanzkriterium 2 / Pruefung 1).
func (s *Store) Allow(ctx context.Context, key string) (Result, error) {
cfg, err := s.getConfig(ctx, key)
if err != nil {
return Result{}, err
}
windowSeconds := int64(cfg.Window.Seconds())
if windowSeconds <= 0 {
windowSeconds = int64(DefaultWindow.Seconds())
}
now := time.Now().UTC()
windowStart := time.Unix((now.Unix()/windowSeconds)*windowSeconds, 0).UTC()
windowEnd := windowStart.Add(time.Duration(windowSeconds) * time.Second)
var count int
err = s.pool.QueryRow(ctx, `
INSERT INTO rate_limit_counters (key, window_start, count)
VALUES ($1, $2, 1)
ON CONFLICT (key, window_start) DO UPDATE
SET count = rate_limit_counters.count + 1
RETURNING count
`, key, windowStart).Scan(&count)
if err != nil {
return Result{}, fmt.Errorf("rate-limit-zaehler erhoehen: %w", err)
}
remaining := cfg.Limit - count
if remaining < 0 {
remaining = 0
}
if count > cfg.Limit {
return Result{
Allowed: false,
Limit: cfg.Limit,
Remaining: 0,
RetryAfter: windowEnd.Sub(now),
}, nil
}
return Result{Allowed: true, Limit: cfg.Limit, Remaining: remaining}, nil
}
-168
View File
@@ -1,168 +0,0 @@
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)
}
}
-2
View File
@@ -1,2 +0,0 @@
DROP TABLE rate_limit_counters;
DROP TABLE rate_limit_configs;
-15
View File
@@ -1,15 +0,0 @@
-- Rate-Limiting mit geteiltem, externem Zustand (API-03, siehe
-- core-kanban/tickets/API-03.md) — bewusst NICHT In-Process, damit mehrere
-- Dienstinstanzen dasselbe Limit korrekt durchsetzen.
CREATE TABLE rate_limit_configs (
key TEXT PRIMARY KEY,
limit_value INT NOT NULL,
window_seconds INT NOT NULL
);
CREATE TABLE rate_limit_counters (
key TEXT NOT NULL,
window_start TIMESTAMPTZ NOT NULL,
count INT NOT NULL DEFAULT 0,
PRIMARY KEY (key, window_start)
);