Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a67adcefa1 |
@@ -1,158 +0,0 @@
|
||||
// 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)
|
||||
}
|
||||
@@ -1,46 +0,0 @@
|
||||
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
|
||||
}
|
||||
@@ -1,185 +0,0 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,136 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
// 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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP TABLE rate_limit_counters;
|
||||
DROP TABLE rate_limit_configs;
|
||||
@@ -0,0 +1,15 @@
|
||||
-- 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)
|
||||
);
|
||||
Reference in New Issue
Block a user