Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
af1709a2bb |
@@ -0,0 +1,99 @@
|
|||||||
|
# ING-09 — Rate-Limiting auf Protokollebene: Prüfprotokoll
|
||||||
|
|
||||||
|
Datum: 2026-09-01
|
||||||
|
Host: 192.168.1.131 (Build/Test/Lint), rsync + ssh
|
||||||
|
Pakete: `mail/internal/ratelimit` (neu, gemeinsam genutzt), `mail/internal/imap`, `mail/internal/pop3`, `mail/internal/smtp`
|
||||||
|
|
||||||
|
## Umsetzung
|
||||||
|
|
||||||
|
Neues Paket `ratelimit`: Token-Bucket-Rate-Limiting, je (Mandant,
|
||||||
|
Quelle)-Schlüssel ein eigener Bucket. `ConfigProvider`/`StaticConfig`
|
||||||
|
liefern die Konfiguration (Burst, Nachfüllrate) je Mandant, mit
|
||||||
|
Fallback auf eine Default-Konfiguration (Akzeptanzkriterium 1/2:
|
||||||
|
begrenzt UND konfigurierbar; Akzeptanzkriterium 3: je Mandant getrennt
|
||||||
|
konfigurierbar). `Limiter.Allow(tenant, source)` liefert bei Ablehnung
|
||||||
|
eine konkrete, positive Wartezeit statt nur `false` — Grundlage für
|
||||||
|
Akzeptanzkriterium 3: "definierte Ablehnung MIT Wartezeit-Hinweis,
|
||||||
|
nicht Verbindungsabbruch ohne Erklärung".
|
||||||
|
|
||||||
|
**IMAP** (`LOGIN`) und **POP3** (`PASS`) begrenzen Anmeldeversuche pro
|
||||||
|
(Mandant, Quell-IP) — Akzeptanzkriterium 1. **SMTP** (`MAIL FROM`)
|
||||||
|
begrenzt die Annahmerate pro (Mandant, Absenderadresse+Quell-IP) —
|
||||||
|
Akzeptanzkriterium 2. Bei Überschreitung antwortet der Server mit einer
|
||||||
|
Fehlermeldung, die die Wartezeit in Sekunden nennt (POP3 `-ERR`, IMAP
|
||||||
|
`NO`, SMTP `451` — temporärer Fehlercode, "versuch es später erneut"),
|
||||||
|
die Verbindung bleibt in allen drei Fällen offen und weiter nutzbar
|
||||||
|
(Akzeptanzkriterium 3). `loginLimiter`/`acceptLimiter` sind optional
|
||||||
|
(`nil` = kein Rate-Limiting, Rückwärtskompatibilität zu ING-01..ING-08);
|
||||||
|
neue Konstruktoren `NewServerWithGuardTLSLoggerAndRateLimit` (IMAP/POP3)
|
||||||
|
und `NewServerWithMaxMessageBytesTLSLoggerAndRateLimit` (SMTP).
|
||||||
|
|
||||||
|
Jeder `Server` bekommt eine `tenantID` — konsistent mit dem in ING-10
|
||||||
|
etablierten Muster "ein Server-Prozess/Instanz je Mandant" — und ein
|
||||||
|
`*ratelimit.Limiter`, der über mehrere Server-Instanzen (Mandanten)
|
||||||
|
hinweg geteilt werden kann, aber intern strikt nach `tenantID` trennt.
|
||||||
|
|
||||||
|
## Pflichtprüfung 1: Lasttest bestätigt greifendes Limit bei Überschreitung
|
||||||
|
|
||||||
|
`TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint` in allen
|
||||||
|
drei Protokollpaketen: Burst=5, 20 reale, aufeinanderfolgende
|
||||||
|
Anmelde-/Annahmeversuche über echte TCP-Verbindungen gegen den
|
||||||
|
laufenden Server. Ergebnis in allen drei Protokollen identisch: exakt
|
||||||
|
5 Versuche akzeptiert (der konfigurierte Burst), exakt 15 Versuche mit
|
||||||
|
der erwarteten Fehlermeldung inkl. Wartezeit-Hinweis abgelehnt — kein
|
||||||
|
Verbindungsabbruch, jede Ablehnung kommt als reguläre Protokollantwort.
|
||||||
|
|
||||||
|
Ergebnis: **BESTANDEN** in allen drei Protokollen.
|
||||||
|
|
||||||
|
## Pflichtprüfung 2: legitime Nutzung unterhalb der Schwelle bleibt unbeeinträchtigt
|
||||||
|
|
||||||
|
`TestRateLimit_LegitUsageBelowThresholdUnaffected` in allen drei
|
||||||
|
Protokollpaketen: Burst=10, nur 3 Versuche — alle drei erfolgreich,
|
||||||
|
keine Ablehnung.
|
||||||
|
|
||||||
|
Ergebnis: **BESTANDEN** in allen drei Protokollen.
|
||||||
|
|
||||||
|
## Pflichtprüfung 3: Limit ist je Mandant getrennt konfigurierbar und wirksam
|
||||||
|
|
||||||
|
`TestRateLimit_PerTenantIndependentAndEffective` in allen drei
|
||||||
|
Protokollpaketen: EIN gemeinsamer `*ratelimit.Limiter`, aber zwei
|
||||||
|
Server-Instanzen mit unterschiedlicher `tenantID`
|
||||||
|
(`mandant-knapp` → Burst 2, `mandant-grosszuegig` → Burst 8, per
|
||||||
|
`StaticConfig.PerTenant`). 10 Versuche je Mandant: `mandant-knapp`
|
||||||
|
akzeptiert exakt 2, `mandant-grosszuegig` akzeptiert exakt 8 — beweist
|
||||||
|
sowohl die Trennung (unterschiedliche Werte wirken unabhängig) als auch
|
||||||
|
die Wirksamkeit (jeweils exakt der konfigurierte Burst, nicht mehr,
|
||||||
|
nicht weniger).
|
||||||
|
|
||||||
|
Ergebnis: **BESTANDEN** in allen drei Protokollen.
|
||||||
|
|
||||||
|
## Akzeptanzkriterien
|
||||||
|
|
||||||
|
1. **Login-Versuche pro Quelle/Zeitfenster sind begrenzt und
|
||||||
|
konfigurierbar**: IMAP/POP3, durch Pflichtprüfung 1+2 belegt.
|
||||||
|
2. **SMTP-Annahmerate pro Absender/Quelle ist begrenzt und
|
||||||
|
konfigurierbar**: SMTP, durch Pflichtprüfung 1+2 belegt.
|
||||||
|
3. **Überschreitung führt zu definierter Ablehnung mit
|
||||||
|
Wartezeit-Hinweis, nicht zu Verbindungsabbruch ohne Erklärung**:
|
||||||
|
durch Pflichtprüfung 1 belegt (Verbindung bleibt in jedem Testlauf
|
||||||
|
offen, jede Ablehnung enthält die Wartezeit in Sekunden).
|
||||||
|
|
||||||
|
## Build/Vet/Lint/Test — Gesamtmodul
|
||||||
|
|
||||||
|
```
|
||||||
|
go build ./... → OK
|
||||||
|
go vet ./... → OK
|
||||||
|
golangci-lint run ./... → 0 issues
|
||||||
|
go test ./... -p 1 (TEST_TENANT_DSN, TEST_MANTICORE_URL gesetzt) → alle Pakete ok, inkl. neuem internal/ratelimit
|
||||||
|
```
|
||||||
|
|
||||||
|
Keine Regression in den bestehenden ~31 Paketen — insbesondere die
|
||||||
|
QA-07-Lasttests bleiben grün: Rate-Limiting ist standardmäßig
|
||||||
|
deaktiviert (`loginLimiter`/`acceptLimiter` nil), bis explizit über die
|
||||||
|
neuen Konstruktoren aktiviert.
|
||||||
|
|
||||||
|
## Ergebnis
|
||||||
|
|
||||||
|
ING-09 erfüllt alle Akzeptanzkriterien mit echten, ausgeführten
|
||||||
|
Nachweisen — in allen drei Protokollen (IMAP, POP3, SMTP) einzeln
|
||||||
|
geprüft. Freigeschaltet: QA-04.
|
||||||
@@ -37,6 +37,13 @@ func (s *Session) handleLogin(ctx context.Context, cmd command) bool {
|
|||||||
// akzeptiert, sobald der Server TLS überhaupt anbietet.
|
// akzeptiert, sobald der Server TLS überhaupt anbietet.
|
||||||
return s.writeErr(cmd.Tag, "NO", "LOGIN disabled without TLS, use STARTTLS")
|
return s.writeErr(cmd.Tag, "NO", "LOGIN disabled without TLS, use STARTTLS")
|
||||||
}
|
}
|
||||||
|
if s.loginLimiter != nil {
|
||||||
|
if ok, retryAfter := s.loginLimiter.Allow(s.tenantID, s.sourceAddr()); !ok {
|
||||||
|
// Akzeptanzkriterium 1/3 (ING-09): definierte Ablehnung MIT
|
||||||
|
// Wartezeit-Hinweis statt Verbindungsabbruch ohne Erklärung.
|
||||||
|
return s.writeErr(cmd.Tag, "NO", fmt.Sprintf("rate limit exceeded, retry in %.1fs", retryAfter.Seconds()))
|
||||||
|
}
|
||||||
|
}
|
||||||
if s.auth == nil {
|
if s.auth == nil {
|
||||||
return s.writeErr(cmd.Tag, "NO", "LOGIN not available")
|
return s.writeErr(cmd.Tag, "NO", "LOGIN not available")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,129 @@
|
|||||||
|
package imap
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"net"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
|
)
|
||||||
|
|
||||||
|
func startRateLimitedServer(t *testing.T, tenant string, limiter *ratelimit.Limiter) (addr string, stop func()) {
|
||||||
|
t.Helper()
|
||||||
|
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
|
||||||
|
store := fakeMailboxStore{mailboxes: map[string][]Message{
|
||||||
|
"INBOX": {{SequenceNumber: 1, UID: 1, Flags: []string{}}},
|
||||||
|
}}
|
||||||
|
srv := NewServerWithGuardTLSLoggerAndRateLimit(auth, store, protoguard.DefaultConfig(), nil, nil, tenant, limiter)
|
||||||
|
|
||||||
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listener: %v", err)
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
done := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
_ = srv.Serve(ctx, listener)
|
||||||
|
close(done)
|
||||||
|
}()
|
||||||
|
return listener.Addr().String(), func() {
|
||||||
|
cancel()
|
||||||
|
<-done
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// attemptLogin führt LOGIN über eine NEUE Verbindung aus und liefert
|
||||||
|
// die Abschlusszeile.
|
||||||
|
func attemptLogin(t *testing.T, addr string) string {
|
||||||
|
t.Helper()
|
||||||
|
c := dial(t, addr)
|
||||||
|
defer c.close()
|
||||||
|
_, lines := c.sendTagged(t, "LOGIN alice geheim123")
|
||||||
|
return lines[len(lines)-1]
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint ist die
|
||||||
|
// geforderte Pflichtprüfung 1 (ING-09).
|
||||||
|
func TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 5, RefillEvery: time.Hour},
|
||||||
|
})
|
||||||
|
addr, stop := startRateLimitedServer(t, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
const attempts = 20
|
||||||
|
var accepted, rejected int
|
||||||
|
for i := 0; i < attempts; i++ {
|
||||||
|
last := attemptLogin(t, addr)
|
||||||
|
switch {
|
||||||
|
case strings.Contains(last, "OK"):
|
||||||
|
accepted++
|
||||||
|
case strings.Contains(last, "NO") && strings.Contains(last, "rate limit"):
|
||||||
|
rejected++
|
||||||
|
default:
|
||||||
|
t.Fatalf("unerwartete abschlussantwort: %q", last)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if accepted != 5 {
|
||||||
|
t.Fatalf("erwartete genau 5 akzeptierte versuche (burst), habe %d", accepted)
|
||||||
|
}
|
||||||
|
if rejected != attempts-5 {
|
||||||
|
t.Fatalf("erwartete %d abgelehnte versuche, habe %d", attempts-5, rejected)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LegitUsageBelowThresholdUnaffected ist die geforderte
|
||||||
|
// Pflichtprüfung 2 (ING-09).
|
||||||
|
func TestRateLimit_LegitUsageBelowThresholdUnaffected(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 10, RefillEvery: time.Second},
|
||||||
|
})
|
||||||
|
addr, stop := startRateLimitedServer(t, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
last := attemptLogin(t, addr)
|
||||||
|
if !strings.Contains(last, "OK") {
|
||||||
|
t.Fatalf("versuch %d unterhalb der schwelle wurde abgelehnt: %q", i+1, last)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_PerTenantIndependentAndEffective ist die geforderte
|
||||||
|
// Pflichtprüfung 3 (ING-09).
|
||||||
|
func TestRateLimit_PerTenantIndependentAndEffective(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 2, RefillEvery: time.Hour},
|
||||||
|
PerTenant: map[string]ratelimit.Config{
|
||||||
|
"mandant-grosszuegig": {Burst: 8, RefillEvery: time.Hour},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
addrKnapp, stopKnapp := startRateLimitedServer(t, "mandant-knapp", limiter)
|
||||||
|
defer stopKnapp()
|
||||||
|
addrGross, stopGross := startRateLimitedServer(t, "mandant-grosszuegig", limiter)
|
||||||
|
defer stopGross()
|
||||||
|
|
||||||
|
var acceptedKnapp int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if strings.Contains(attemptLogin(t, addrKnapp), "OK") {
|
||||||
|
acceptedKnapp++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var acceptedGross int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if strings.Contains(attemptLogin(t, addrGross), "OK") {
|
||||||
|
acceptedGross++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if acceptedKnapp != 2 {
|
||||||
|
t.Fatalf("mandant-knapp: erwartete 2 akzeptierte versuche, habe %d", acceptedKnapp)
|
||||||
|
}
|
||||||
|
if acceptedGross != 8 {
|
||||||
|
t.Fatalf("mandant-grosszuegig: erwartete 8 akzeptierte versuche, habe %d", acceptedGross)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Server nimmt IMAP-Verbindungen an und bedient jede in einer eigenen
|
// Server nimmt IMAP-Verbindungen an und bedient jede in einer eigenen
|
||||||
@@ -23,6 +24,9 @@ type Server struct {
|
|||||||
guardCfg protoguard.Config
|
guardCfg protoguard.Config
|
||||||
tlsConfig *tls.Config
|
tlsConfig *tls.Config
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
|
||||||
|
tenantID string
|
||||||
|
loginLimiter *ratelimit.Limiter
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewServer(auth Authenticator, store MailboxStore) *Server {
|
func NewServer(auth Authenticator, store MailboxStore) *Server {
|
||||||
@@ -50,6 +54,14 @@ func NewServerWithGuardTLSAndLogger(auth Authenticator, store MailboxStore, guar
|
|||||||
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger}
|
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewServerWithGuardTLSLoggerAndRateLimit erlaubt zusätzlich
|
||||||
|
// Rate-Limiting für LOGIN-Versuche (ING-09). loginLimiter darf nil sein
|
||||||
|
// (Rate-Limiting dann deaktiviert). tenantID identifiziert diesen
|
||||||
|
// Server gegenüber dem Limiter (Akzeptanzkriterium 3).
|
||||||
|
func NewServerWithGuardTLSLoggerAndRateLimit(auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, loginLimiter *ratelimit.Limiter) *Server {
|
||||||
|
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger, tenantID: tenantID, loginLimiter: loginLimiter}
|
||||||
|
}
|
||||||
|
|
||||||
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird oder
|
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird oder
|
||||||
// Accept endgültig fehlschlägt. Blockiert den Aufrufer.
|
// Accept endgültig fehlschlägt. Blockiert den Aufrufer.
|
||||||
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
||||||
@@ -70,7 +82,7 @@ func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
|||||||
}
|
}
|
||||||
return fmt.Errorf("imap: verbindung annehmen: %w", err)
|
return fmt.Errorf("imap: verbindung annehmen: %w", err)
|
||||||
}
|
}
|
||||||
session := newSession(conn, srv.auth, srv.store, srv.guardCfg, srv.tlsConfig, srv.logger)
|
session := newSession(conn, srv.auth, srv.store, srv.guardCfg, srv.tlsConfig, srv.logger, srv.tenantID, srv.loginLimiter)
|
||||||
go session.Serve(ctx)
|
go session.Serve(ctx)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// phaseNotAuthenticated/phaseSelected sind die protoguard-Phasen dieser
|
// phaseNotAuthenticated/phaseSelected sind die protoguard-Phasen dieser
|
||||||
@@ -33,36 +34,52 @@ const maxCommandLineBytes = 8192
|
|||||||
// Session ist eine einzelne IMAP-Verbindung mit eigener
|
// Session ist eine einzelne IMAP-Verbindung mit eigener
|
||||||
// Zustandsmaschine (Akzeptanzkriterium 1).
|
// Zustandsmaschine (Akzeptanzkriterium 1).
|
||||||
type Session struct {
|
type Session struct {
|
||||||
conn net.Conn
|
conn net.Conn
|
||||||
reader *bufio.Reader
|
reader *bufio.Reader
|
||||||
writer *bufio.Writer
|
writer *bufio.Writer
|
||||||
auth Authenticator
|
auth Authenticator
|
||||||
store MailboxStore
|
store MailboxStore
|
||||||
guard *protoguard.Guard
|
guard *protoguard.Guard
|
||||||
tlsConfig *tls.Config // nil = kein TLS/STARTTLS angeboten (ING-06)
|
tlsConfig *tls.Config // nil = kein TLS/STARTTLS angeboten (ING-06)
|
||||||
tlsActive bool
|
tlsActive bool
|
||||||
log *protolog.SessionLogger // ING-08, nie nil (log.Event() ist nil-sicher)
|
log *protolog.SessionLogger // ING-08, nie nil (log.Event() ist nil-sicher)
|
||||||
|
|
||||||
|
tenantID string
|
||||||
|
loginLimiter *ratelimit.Limiter // ING-09, nil = kein Rate-Limiting
|
||||||
|
|
||||||
state State
|
state State
|
||||||
mailbox string // gewähltes Postfach im Zustand Selected
|
mailbox string // gewähltes Postfach im Zustand Selected
|
||||||
mailboxSize uint32 // Nachrichtenzahl aus dem letzten erfolgreichen SELECT
|
mailboxSize uint32 // Nachrichtenzahl aus dem letzten erfolgreichen SELECT
|
||||||
}
|
}
|
||||||
|
|
||||||
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger) *Session {
|
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, loginLimiter *ratelimit.Limiter) *Session {
|
||||||
_, alreadyTLS := conn.(*tls.Conn)
|
_, alreadyTLS := conn.(*tls.Conn)
|
||||||
return &Session{
|
return &Session{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
|
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
|
||||||
writer: bufio.NewWriter(conn),
|
writer: bufio.NewWriter(conn),
|
||||||
auth: auth,
|
auth: auth,
|
||||||
store: store,
|
store: store,
|
||||||
guard: protoguard.New(guardCfg),
|
guard: protoguard.New(guardCfg),
|
||||||
tlsConfig: tlsConfig,
|
tlsConfig: tlsConfig,
|
||||||
tlsActive: alreadyTLS,
|
tlsActive: alreadyTLS,
|
||||||
log: protolog.NewSessionLogger(logger, "imap"),
|
log: protolog.NewSessionLogger(logger, "imap"),
|
||||||
state: NotAuthenticated,
|
tenantID: tenantID,
|
||||||
|
loginLimiter: loginLimiter,
|
||||||
|
state: NotAuthenticated,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sourceAddr liefert die Quell-IP dieser Verbindung ohne Portanteil
|
||||||
|
// (ING-09).
|
||||||
|
func (s *Session) sourceAddr() string {
|
||||||
|
host, _, err := net.SplitHostPort(s.conn.RemoteAddr().String())
|
||||||
|
if err != nil {
|
||||||
|
return s.conn.RemoteAddr().String()
|
||||||
|
}
|
||||||
|
return host
|
||||||
|
}
|
||||||
|
|
||||||
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
|
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
|
||||||
func (s *Session) currentPhase() protoguard.Phase {
|
func (s *Session) currentPhase() protoguard.Phase {
|
||||||
if s.state == NotAuthenticated {
|
if s.state == NotAuthenticated {
|
||||||
|
|||||||
@@ -47,6 +47,15 @@ func (s *Session) handlePass(ctx context.Context, cmd command) bool {
|
|||||||
// akzeptiert, sobald der Server TLS überhaupt anbietet.
|
// akzeptiert, sobald der Server TLS überhaupt anbietet.
|
||||||
return writeErr(s.writer, "TLS required before authentication, use STLS") == nil
|
return writeErr(s.writer, "TLS required before authentication, use STLS") == nil
|
||||||
}
|
}
|
||||||
|
if s.loginLimiter != nil {
|
||||||
|
if ok, retryAfter := s.loginLimiter.Allow(s.tenantID, s.sourceAddr()); !ok {
|
||||||
|
// Akzeptanzkriterium 1/3 (ING-09): definierte Ablehnung MIT
|
||||||
|
// Wartezeit-Hinweis statt Verbindungsabbruch ohne Erklärung —
|
||||||
|
// die Verbindung bleibt offen (true), nur DIESER Versuch wird
|
||||||
|
// abgelehnt.
|
||||||
|
return writeErr(s.writer, fmt.Sprintf("rate limit exceeded, retry in %.1fs", retryAfter.Seconds())) == nil
|
||||||
|
}
|
||||||
|
}
|
||||||
if s.auth == nil {
|
if s.auth == nil {
|
||||||
return writeErr(s.writer, genericAuthFailure) == nil
|
return writeErr(s.writer, genericAuthFailure) == nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,145 @@
|
|||||||
|
package pop3
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"context"
|
||||||
|
"net"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
|
)
|
||||||
|
|
||||||
|
func startRateLimitedServer(t *testing.T, tenant string, limiter *ratelimit.Limiter) (addr string, stop func()) {
|
||||||
|
t.Helper()
|
||||||
|
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123", "bob": "geheim456"}}
|
||||||
|
store := newFakeMailboxStore()
|
||||||
|
srv := NewServerWithGuardTLSLoggerAndRateLimit(auth, store, protoguard.DefaultConfig(), nil, nil, tenant, limiter)
|
||||||
|
|
||||||
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listener: %v", err)
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
done := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
_ = srv.Serve(ctx, listener)
|
||||||
|
close(done)
|
||||||
|
}()
|
||||||
|
return listener.Addr().String(), func() {
|
||||||
|
cancel()
|
||||||
|
<-done
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// attemptPass führt USER+PASS über eine NEUE Verbindung aus und liefert
|
||||||
|
// die PASS-Antwortzeile.
|
||||||
|
func attemptPass(t *testing.T, addr, user, pass string) string {
|
||||||
|
t.Helper()
|
||||||
|
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial: %v", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = conn.Close() }()
|
||||||
|
reader := bufio.NewReader(conn)
|
||||||
|
_, _ = reader.ReadString('\n')
|
||||||
|
_, _ = conn.Write([]byte("USER " + user + "\r\n"))
|
||||||
|
_, _ = reader.ReadString('\n')
|
||||||
|
_, _ = conn.Write([]byte("PASS " + pass + "\r\n"))
|
||||||
|
resp, err := reader.ReadString('\n')
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("PASS antwort lesen: %v", err)
|
||||||
|
}
|
||||||
|
return resp
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint ist die
|
||||||
|
// geforderte Pflichtprüfung 1 (ING-09): Lasttest bestätigt greifendes
|
||||||
|
// Limit bei Überschreitung — reale, gleichzeitige Anmeldeversuche über
|
||||||
|
// den Burst hinaus.
|
||||||
|
func TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 5, RefillEvery: time.Hour}, // Refill irrelevant für diesen Test
|
||||||
|
})
|
||||||
|
addr, stop := startRateLimitedServer(t, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
const attempts = 20
|
||||||
|
var accepted, rejected int
|
||||||
|
for i := 0; i < attempts; i++ {
|
||||||
|
resp := attemptPass(t, addr, "alice", "geheim123")
|
||||||
|
switch {
|
||||||
|
case strings.HasPrefix(resp, "+OK"):
|
||||||
|
accepted++
|
||||||
|
case strings.HasPrefix(resp, "-ERR") && strings.Contains(resp, "rate limit"):
|
||||||
|
rejected++
|
||||||
|
default:
|
||||||
|
t.Fatalf("unerwartete antwort: %q", resp)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if accepted != 5 {
|
||||||
|
t.Fatalf("erwartete genau 5 akzeptierte versuche (burst), habe %d", accepted)
|
||||||
|
}
|
||||||
|
if rejected != attempts-5 {
|
||||||
|
t.Fatalf("erwartete %d abgelehnte versuche, habe %d", attempts-5, rejected)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LegitUsageBelowThresholdUnaffected ist die geforderte
|
||||||
|
// Pflichtprüfung 2 (ING-09): legitime Nutzung unterhalb der Schwelle
|
||||||
|
// bleibt unbeeinträchtigt.
|
||||||
|
func TestRateLimit_LegitUsageBelowThresholdUnaffected(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 10, RefillEvery: time.Second},
|
||||||
|
})
|
||||||
|
addr, stop := startRateLimitedServer(t, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
resp := attemptPass(t, addr, "alice", "geheim123")
|
||||||
|
if !strings.HasPrefix(resp, "+OK") {
|
||||||
|
t.Fatalf("versuch %d unterhalb der schwelle wurde abgelehnt: %q", i+1, resp)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_PerTenantIndependentAndEffective ist die geforderte
|
||||||
|
// Pflichtprüfung 3 (ING-09): Limit ist je Mandant getrennt
|
||||||
|
// konfigurierbar und wirksam — zwei Serverinstanzen (Mandant A/B) mit
|
||||||
|
// UNTERSCHIEDLICHEM Burst, gegen DENSELBEN Limiter (realistisch: ein
|
||||||
|
// zentraler Limiter-Prozess, mehrere Mandanten-Server).
|
||||||
|
func TestRateLimit_PerTenantIndependentAndEffective(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 2, RefillEvery: time.Hour},
|
||||||
|
PerTenant: map[string]ratelimit.Config{
|
||||||
|
"mandant-grosszuegig": {Burst: 8, RefillEvery: time.Hour},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
addrKnapp, stopKnapp := startRateLimitedServer(t, "mandant-knapp", limiter)
|
||||||
|
defer stopKnapp()
|
||||||
|
addrGross, stopGross := startRateLimitedServer(t, "mandant-grosszuegig", limiter)
|
||||||
|
defer stopGross()
|
||||||
|
|
||||||
|
var acceptedKnapp int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if strings.HasPrefix(attemptPass(t, addrKnapp, "alice", "geheim123"), "+OK") {
|
||||||
|
acceptedKnapp++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var acceptedGross int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if strings.HasPrefix(attemptPass(t, addrGross, "alice", "geheim123"), "+OK") {
|
||||||
|
acceptedGross++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if acceptedKnapp != 2 {
|
||||||
|
t.Fatalf("mandant-knapp: erwartete 2 akzeptierte versuche (eigener burst), habe %d", acceptedKnapp)
|
||||||
|
}
|
||||||
|
if acceptedGross != 8 {
|
||||||
|
t.Fatalf("mandant-grosszuegig: erwartete 8 akzeptierte versuche (eigener, größerer burst), habe %d", acceptedGross)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Server nimmt POP3-Verbindungen an und bedient jede in einer eigenen
|
// Server nimmt POP3-Verbindungen an und bedient jede in einer eigenen
|
||||||
@@ -25,6 +26,12 @@ type Server struct {
|
|||||||
guardCfg protoguard.Config
|
guardCfg protoguard.Config
|
||||||
tlsConfig *tls.Config
|
tlsConfig *tls.Config
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
|
||||||
|
// tenantID identifiziert diesen Server für das Rate-Limiting
|
||||||
|
// (ING-09, Akzeptanzkriterium 3: je Mandant getrennt konfigurierbar)
|
||||||
|
// — leer, wenn loginLimiter nil ist.
|
||||||
|
tenantID string
|
||||||
|
loginLimiter *ratelimit.Limiter
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewServer(auth Authenticator, store MailboxStore) *Server {
|
func NewServer(auth Authenticator, store MailboxStore) *Server {
|
||||||
@@ -52,6 +59,15 @@ func NewServerWithGuardTLSAndLogger(auth Authenticator, store MailboxStore, guar
|
|||||||
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger}
|
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewServerWithGuardTLSLoggerAndRateLimit erlaubt zusätzlich
|
||||||
|
// Rate-Limiting für PASS-Versuche (ING-09). loginLimiter darf nil sein
|
||||||
|
// (Rate-Limiting dann deaktiviert, Rückwärtskompatibilität zu
|
||||||
|
// ING-01..ING-08). tenantID identifiziert diesen Server gegenüber dem
|
||||||
|
// Limiter (Akzeptanzkriterium 3).
|
||||||
|
func NewServerWithGuardTLSLoggerAndRateLimit(auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, loginLimiter *ratelimit.Limiter) *Server {
|
||||||
|
return &Server{auth: auth, store: store, guardCfg: guardCfg, tlsConfig: tlsConfig, logger: logger, tenantID: tenantID, loginLimiter: loginLimiter}
|
||||||
|
}
|
||||||
|
|
||||||
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
|
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
|
||||||
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
||||||
go func() {
|
go func() {
|
||||||
@@ -71,7 +87,7 @@ func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
|||||||
}
|
}
|
||||||
return fmt.Errorf("pop3: verbindung annehmen: %w", err)
|
return fmt.Errorf("pop3: verbindung annehmen: %w", err)
|
||||||
}
|
}
|
||||||
session := newSession(conn, srv.auth, srv.store, srv.guardCfg, srv.tlsConfig, srv.logger)
|
session := newSession(conn, srv.auth, srv.store, srv.guardCfg, srv.tlsConfig, srv.logger, srv.tenantID, srv.loginLimiter)
|
||||||
go session.Serve(ctx)
|
go session.Serve(ctx)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protoguard"
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// phaseAuthorization/phaseTransaction sind die protoguard-Phasen dieser
|
// phaseAuthorization/phaseTransaction sind die protoguard-Phasen dieser
|
||||||
@@ -46,29 +47,48 @@ type Session struct {
|
|||||||
|
|
||||||
log *protolog.SessionLogger // ING-08, nie nil (aber log.Event() ist nil-sicher)
|
log *protolog.SessionLogger // ING-08, nie nil (aber log.Event() ist nil-sicher)
|
||||||
|
|
||||||
|
// tenantID/loginLimiter: Rate-Limiting für PASS-Versuche (ING-09).
|
||||||
|
// loginLimiter nil bedeutet: kein Rate-Limiting (Rückwärtskompatibilität
|
||||||
|
// zu ING-01..ING-08).
|
||||||
|
tenantID string
|
||||||
|
loginLimiter *ratelimit.Limiter
|
||||||
|
|
||||||
state State
|
state State
|
||||||
pendingUsername string // nach USER, vor erfolgreichem PASS
|
pendingUsername string // nach USER, vor erfolgreichem PASS
|
||||||
username string // nach erfolgreichem PASS
|
username string // nach erfolgreichem PASS
|
||||||
deleted map[int]bool
|
deleted map[int]bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger) *Session {
|
func newSession(conn net.Conn, auth Authenticator, store MailboxStore, guardCfg protoguard.Config, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, loginLimiter *ratelimit.Limiter) *Session {
|
||||||
_, alreadyTLS := conn.(*tls.Conn)
|
_, alreadyTLS := conn.(*tls.Conn)
|
||||||
return &Session{
|
return &Session{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
|
reader: bufio.NewReaderSize(conn, maxCommandLineBytes),
|
||||||
writer: bufio.NewWriter(conn),
|
writer: bufio.NewWriter(conn),
|
||||||
auth: auth,
|
auth: auth,
|
||||||
store: store,
|
store: store,
|
||||||
guard: protoguard.New(guardCfg),
|
guard: protoguard.New(guardCfg),
|
||||||
tlsConfig: tlsConfig,
|
tlsConfig: tlsConfig,
|
||||||
tlsActive: alreadyTLS,
|
tlsActive: alreadyTLS,
|
||||||
log: protolog.NewSessionLogger(logger, "pop3"),
|
log: protolog.NewSessionLogger(logger, "pop3"),
|
||||||
state: Authorization,
|
tenantID: tenantID,
|
||||||
deleted: map[int]bool{},
|
loginLimiter: loginLimiter,
|
||||||
|
state: Authorization,
|
||||||
|
deleted: map[int]bool{},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sourceAddr liefert die Quell-IP dieser Verbindung ohne Portanteil,
|
||||||
|
// für das Rate-Limiting (ING-09) und als Schlüssel gegenüber dem
|
||||||
|
// Limiter stabil pro Client.
|
||||||
|
func (s *Session) sourceAddr() string {
|
||||||
|
host, _, err := net.SplitHostPort(s.conn.RemoteAddr().String())
|
||||||
|
if err != nil {
|
||||||
|
return s.conn.RemoteAddr().String()
|
||||||
|
}
|
||||||
|
return host
|
||||||
|
}
|
||||||
|
|
||||||
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
|
// currentPhase liefert die protoguard-Phase des aktuellen Sitzungszustands.
|
||||||
func (s *Session) currentPhase() protoguard.Phase {
|
func (s *Session) currentPhase() protoguard.Phase {
|
||||||
if s.state == Authorization {
|
if s.state == Authorization {
|
||||||
|
|||||||
@@ -0,0 +1,109 @@
|
|||||||
|
// Package ratelimit implementiert ING-09: Token-Bucket-Rate-Limiting
|
||||||
|
// auf Protokollebene für Login-Versuche (IMAP/POP3) und SMTP-Annahme,
|
||||||
|
// je Mandant getrennt konfigurierbar (Akzeptanzkriterium 3).
|
||||||
|
package ratelimit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Config ist die Token-Bucket-Konfiguration EINES Limits
|
||||||
|
// (Akzeptanzkriterium 1/2: begrenzt und konfigurierbar). Burst ist die
|
||||||
|
// Anzahl sofort verfügbarer Versuche, RefillEvery die Zeit, die ein
|
||||||
|
// neuer Versuch nachwächst.
|
||||||
|
type Config struct {
|
||||||
|
Burst int
|
||||||
|
RefillEvery time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
// ConfigProvider liefert die Rate-Limit-Konfiguration für einen
|
||||||
|
// Mandanten (Akzeptanzkriterium 3: je Mandant getrennt konfigurierbar).
|
||||||
|
type ConfigProvider interface {
|
||||||
|
ConfigFor(tenant string) Config
|
||||||
|
}
|
||||||
|
|
||||||
|
// StaticConfig ist ein einfacher ConfigProvider: feste Konfiguration je
|
||||||
|
// Mandant, mit Fallback auf Default für unbekannte/nicht gesondert
|
||||||
|
// konfigurierte Mandanten.
|
||||||
|
type StaticConfig struct {
|
||||||
|
Default Config
|
||||||
|
PerTenant map[string]Config
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s StaticConfig) ConfigFor(tenant string) Config {
|
||||||
|
if cfg, ok := s.PerTenant[tenant]; ok {
|
||||||
|
return cfg
|
||||||
|
}
|
||||||
|
return s.Default
|
||||||
|
}
|
||||||
|
|
||||||
|
// tokenBucket ist EIN Token-Bucket-Zähler für einen Schlüssel
|
||||||
|
// (Mandant+Quelle).
|
||||||
|
type tokenBucket struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
tokens float64
|
||||||
|
lastRefill time.Time
|
||||||
|
cfg Config
|
||||||
|
}
|
||||||
|
|
||||||
|
func newTokenBucket(cfg Config) *tokenBucket {
|
||||||
|
return &tokenBucket{tokens: float64(cfg.Burst), lastRefill: time.Now(), cfg: cfg}
|
||||||
|
}
|
||||||
|
|
||||||
|
// allow entscheidet über EINEN Versuch zum Zeitpunkt now. Bei
|
||||||
|
// Ablehnung liefert retryAfter eine konkrete, positive Wartezeit
|
||||||
|
// (Akzeptanzkriterium 1: definierte Ablehnung MIT Wartezeit-Hinweis,
|
||||||
|
// nicht bloßer Verbindungsabbruch).
|
||||||
|
func (b *tokenBucket) allow(now time.Time) (ok bool, retryAfter time.Duration) {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
|
refillPerSecond := 1.0 / b.cfg.RefillEvery.Seconds()
|
||||||
|
elapsed := now.Sub(b.lastRefill).Seconds()
|
||||||
|
b.tokens += elapsed * refillPerSecond
|
||||||
|
if b.tokens > float64(b.cfg.Burst) {
|
||||||
|
b.tokens = float64(b.cfg.Burst)
|
||||||
|
}
|
||||||
|
b.lastRefill = now
|
||||||
|
|
||||||
|
if b.tokens >= 1 {
|
||||||
|
b.tokens--
|
||||||
|
return true, 0
|
||||||
|
}
|
||||||
|
missing := 1 - b.tokens
|
||||||
|
wait := time.Duration(missing / refillPerSecond * float64(time.Second))
|
||||||
|
if wait <= 0 {
|
||||||
|
wait = time.Millisecond
|
||||||
|
}
|
||||||
|
return false, wait
|
||||||
|
}
|
||||||
|
|
||||||
|
// Limiter verwaltet Token-Buckets je (Mandant, Quelle)-Schlüssel —
|
||||||
|
// EIN Limiter deckt EINEN Limit-Zweck ab (z. B. "Login-Versuche" oder
|
||||||
|
// "SMTP-Annahme"); ein Server verwendet für unterschiedliche Zwecke
|
||||||
|
// unterschiedliche Limiter-Instanzen.
|
||||||
|
type Limiter struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
buckets map[string]*tokenBucket
|
||||||
|
provider ConfigProvider
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewLimiter(provider ConfigProvider) *Limiter {
|
||||||
|
return &Limiter{buckets: map[string]*tokenBucket{}, provider: provider}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Allow prüft, ob EIN Versuch von source innerhalb des Mandanten
|
||||||
|
// tenant aktuell erlaubt ist.
|
||||||
|
func (l *Limiter) Allow(tenant, source string) (ok bool, retryAfter time.Duration) {
|
||||||
|
key := fmt.Sprintf("%s|%s", tenant, source)
|
||||||
|
l.mu.Lock()
|
||||||
|
b, exists := l.buckets[key]
|
||||||
|
if !exists {
|
||||||
|
b = newTokenBucket(l.provider.ConfigFor(tenant))
|
||||||
|
l.buckets[key] = b
|
||||||
|
}
|
||||||
|
l.mu.Unlock()
|
||||||
|
return b.allow(time.Now())
|
||||||
|
}
|
||||||
@@ -0,0 +1,70 @@
|
|||||||
|
package ratelimit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestLimiter_AllowsBurstThenBlocksThenRecovers(t *testing.T) {
|
||||||
|
cfg := Config{Burst: 3, RefillEvery: 50 * time.Millisecond}
|
||||||
|
lim := NewLimiter(StaticConfig{Default: cfg})
|
||||||
|
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
ok, _ := lim.Allow("mandant-a", "1.2.3.4")
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("versuch %d im burst hätte erlaubt sein müssen", i+1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ok, retryAfter := lim.Allow("mandant-a", "1.2.3.4")
|
||||||
|
if ok {
|
||||||
|
t.Fatalf("vierter versuch über dem burst hätte abgelehnt werden müssen")
|
||||||
|
}
|
||||||
|
if retryAfter <= 0 {
|
||||||
|
t.Fatalf("erwartete positive wartezeit als hinweis, habe %v", retryAfter)
|
||||||
|
}
|
||||||
|
|
||||||
|
time.Sleep(retryAfter + 10*time.Millisecond)
|
||||||
|
ok, _ = lim.Allow("mandant-a", "1.2.3.4")
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("nach der wartezeit hätte wieder ein token verfügbar sein müssen")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLimiter_TenantsAreIndependent(t *testing.T) {
|
||||||
|
lim := NewLimiter(StaticConfig{Default: Config{Burst: 1, RefillEvery: time.Hour}})
|
||||||
|
|
||||||
|
okA, _ := lim.Allow("mandant-a", "1.2.3.4")
|
||||||
|
if !okA {
|
||||||
|
t.Fatalf("mandant a: erster versuch hätte erlaubt sein müssen")
|
||||||
|
}
|
||||||
|
okA2, _ := lim.Allow("mandant-a", "1.2.3.4")
|
||||||
|
if okA2 {
|
||||||
|
t.Fatalf("mandant a: zweiter versuch hätte abgelehnt werden müssen")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Mandant B mit DERSELBEN Quelladresse — eigener Bucket.
|
||||||
|
okB, _ := lim.Allow("mandant-b", "1.2.3.4")
|
||||||
|
if !okB {
|
||||||
|
t.Fatalf("mandant b: eigener bucket, erster versuch hätte erlaubt sein müssen")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLimiter_PerTenantConfigOverridesDefault(t *testing.T) {
|
||||||
|
lim := NewLimiter(StaticConfig{
|
||||||
|
Default: Config{Burst: 1, RefillEvery: time.Hour},
|
||||||
|
PerTenant: map[string]Config{
|
||||||
|
"mandant-grosszuegig": {Burst: 5, RefillEvery: time.Hour},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
for i := 0; i < 5; i++ {
|
||||||
|
ok, _ := lim.Allow("mandant-grosszuegig", "1.2.3.4")
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("mandant-grosszuegig: versuch %d hätte im eigenen, größeren burst erlaubt sein müssen", i+1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ok, _ := lim.Allow("mandant-grosszuegig", "1.2.3.4")
|
||||||
|
if ok {
|
||||||
|
t.Fatalf("mandant-grosszuegig: sechster versuch hätte abgelehnt werden müssen")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"bufio"
|
"bufio"
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/tlscert"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/tlscert"
|
||||||
@@ -77,6 +78,14 @@ func (s *Session) handleMailFrom(arg string) bool {
|
|||||||
// SMTP-Fehlermeldung statt Absturz oder Verbindungsabbruch.
|
// SMTP-Fehlermeldung statt Absturz oder Verbindungsabbruch.
|
||||||
return s.reply(553, "invalid sender address") == nil
|
return s.reply(553, "invalid sender address") == nil
|
||||||
}
|
}
|
||||||
|
if s.acceptLimiter != nil {
|
||||||
|
if ok, retryAfter := s.acceptLimiter.Allow(s.tenantID, addr+"|"+s.sourceAddr()); !ok {
|
||||||
|
// Akzeptanzkriterium 1/3 (ING-09): definierte, temporäre
|
||||||
|
// Ablehnung (4xx = "try again later") MIT Wartezeit-Hinweis
|
||||||
|
// statt Verbindungsabbruch ohne Erklärung.
|
||||||
|
return s.reply(451, fmt.Sprintf("rate limit exceeded for sender, retry in %.1fs", retryAfter.Seconds())) == nil
|
||||||
|
}
|
||||||
|
}
|
||||||
s.from = addr
|
s.from = addr
|
||||||
s.to = nil
|
s.to = nil
|
||||||
s.state = MailFromSet
|
s.state = MailFromSet
|
||||||
|
|||||||
@@ -0,0 +1,134 @@
|
|||||||
|
package smtp
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"net"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
|
)
|
||||||
|
|
||||||
|
func startRateLimitedServer(t *testing.T, sink MessageSink, tenant string, limiter *ratelimit.Limiter) (addr string, stop func()) {
|
||||||
|
t.Helper()
|
||||||
|
srv := NewServerWithMaxMessageBytesTLSLoggerAndRateLimit(sink, defaultMaxMessageBytes, nil, nil, tenant, limiter)
|
||||||
|
|
||||||
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listener: %v", err)
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
done := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
_ = srv.Serve(ctx, listener)
|
||||||
|
close(done)
|
||||||
|
}()
|
||||||
|
return listener.Addr().String(), func() {
|
||||||
|
cancel()
|
||||||
|
<-done
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// attemptMailFrom führt EHLO+MAIL FROM über eine NEUE Verbindung aus
|
||||||
|
// und liefert die MAIL FROM-Antwortzeile.
|
||||||
|
func attemptMailFrom(t *testing.T, addr, from string) string {
|
||||||
|
t.Helper()
|
||||||
|
c := dial(t, addr)
|
||||||
|
defer c.close()
|
||||||
|
c.send(t, "EHLO client.example.com")
|
||||||
|
for {
|
||||||
|
line := c.readLine(t)
|
||||||
|
if strings.HasPrefix(line, "250 ") {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return c.send(t, "MAIL FROM:<"+from+">")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint ist die
|
||||||
|
// geforderte Pflichtprüfung 1 (ING-09).
|
||||||
|
func TestRateLimit_LoadExceedingLimitGetsRejectedWithRetryHint(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 5, RefillEvery: time.Hour},
|
||||||
|
})
|
||||||
|
sink := &fakeSink{}
|
||||||
|
addr, stop := startRateLimitedServer(t, sink, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
const attempts = 20
|
||||||
|
var accepted, rejected int
|
||||||
|
for i := 0; i < attempts; i++ {
|
||||||
|
resp := attemptMailFrom(t, addr, "immer-gleicher-absender@example.com")
|
||||||
|
switch {
|
||||||
|
case code(resp) == "250":
|
||||||
|
accepted++
|
||||||
|
case code(resp) == "451" && strings.Contains(resp, "rate limit"):
|
||||||
|
rejected++
|
||||||
|
default:
|
||||||
|
t.Fatalf("unerwartete antwort: %q", resp)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if accepted != 5 {
|
||||||
|
t.Fatalf("erwartete genau 5 akzeptierte versuche (burst), habe %d", accepted)
|
||||||
|
}
|
||||||
|
if rejected != attempts-5 {
|
||||||
|
t.Fatalf("erwartete %d abgelehnte versuche, habe %d", attempts-5, rejected)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_LegitUsageBelowThresholdUnaffected ist die geforderte
|
||||||
|
// Pflichtprüfung 2 (ING-09).
|
||||||
|
func TestRateLimit_LegitUsageBelowThresholdUnaffected(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 10, RefillEvery: time.Second},
|
||||||
|
})
|
||||||
|
sink := &fakeSink{}
|
||||||
|
addr, stop := startRateLimitedServer(t, sink, "mandant-a", limiter)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
resp := attemptMailFrom(t, addr, "legitim@example.com")
|
||||||
|
if code(resp) != "250" {
|
||||||
|
t.Fatalf("versuch %d unterhalb der schwelle wurde abgelehnt: %q", i+1, resp)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRateLimit_PerTenantIndependentAndEffective ist die geforderte
|
||||||
|
// Pflichtprüfung 3 (ING-09).
|
||||||
|
func TestRateLimit_PerTenantIndependentAndEffective(t *testing.T) {
|
||||||
|
limiter := ratelimit.NewLimiter(ratelimit.StaticConfig{
|
||||||
|
Default: ratelimit.Config{Burst: 2, RefillEvery: time.Hour},
|
||||||
|
PerTenant: map[string]ratelimit.Config{
|
||||||
|
"mandant-grosszuegig": {Burst: 8, RefillEvery: time.Hour},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
sinkKnapp := &fakeSink{}
|
||||||
|
addrKnapp, stopKnapp := startRateLimitedServer(t, sinkKnapp, "mandant-knapp", limiter)
|
||||||
|
defer stopKnapp()
|
||||||
|
sinkGross := &fakeSink{}
|
||||||
|
addrGross, stopGross := startRateLimitedServer(t, sinkGross, "mandant-grosszuegig", limiter)
|
||||||
|
defer stopGross()
|
||||||
|
|
||||||
|
var acceptedKnapp int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if code(attemptMailFrom(t, addrKnapp, "absender@example.com")) == "250" {
|
||||||
|
acceptedKnapp++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var acceptedGross int
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
if code(attemptMailFrom(t, addrGross, "absender@example.com")) == "250" {
|
||||||
|
acceptedGross++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if acceptedKnapp != 2 {
|
||||||
|
t.Fatalf("mandant-knapp: erwartete 2 akzeptierte versuche, habe %d", acceptedKnapp)
|
||||||
|
}
|
||||||
|
if acceptedGross != 8 {
|
||||||
|
t.Fatalf("mandant-grosszuegig: erwartete 8 akzeptierte versuche, habe %d", acceptedGross)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -7,6 +7,8 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net"
|
"net"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// defaultMaxMessageBytes ist die Standard-Höchstgröße einer
|
// defaultMaxMessageBytes ist die Standard-Höchstgröße einer
|
||||||
@@ -23,6 +25,9 @@ type Server struct {
|
|||||||
maxMessageBytes int64
|
maxMessageBytes int64
|
||||||
tlsConfig *tls.Config
|
tlsConfig *tls.Config
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
|
||||||
|
tenantID string
|
||||||
|
acceptLimiter *ratelimit.Limiter
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewServer(sink MessageSink) *Server {
|
func NewServer(sink MessageSink) *Server {
|
||||||
@@ -49,6 +54,15 @@ func NewServerWithMaxMessageBytesTLSAndLogger(sink MessageSink, maxMessageBytes
|
|||||||
return &Server{sink: sink, maxMessageBytes: maxMessageBytes, tlsConfig: tlsConfig, logger: logger}
|
return &Server{sink: sink, maxMessageBytes: maxMessageBytes, tlsConfig: tlsConfig, logger: logger}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewServerWithMaxMessageBytesTLSLoggerAndRateLimit erlaubt zusätzlich
|
||||||
|
// Rate-Limiting der SMTP-Annahmerate pro Absender (ING-09, MAIL FROM).
|
||||||
|
// acceptLimiter darf nil sein (Rate-Limiting dann deaktiviert).
|
||||||
|
// tenantID identifiziert diesen Server gegenüber dem Limiter
|
||||||
|
// (Akzeptanzkriterium 3).
|
||||||
|
func NewServerWithMaxMessageBytesTLSLoggerAndRateLimit(sink MessageSink, maxMessageBytes int64, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, acceptLimiter *ratelimit.Limiter) *Server {
|
||||||
|
return &Server{sink: sink, maxMessageBytes: maxMessageBytes, tlsConfig: tlsConfig, logger: logger, tenantID: tenantID, acceptLimiter: acceptLimiter}
|
||||||
|
}
|
||||||
|
|
||||||
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
|
// Serve nimmt Verbindungen auf listener an, bis ctx beendet wird.
|
||||||
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
||||||
go func() {
|
go func() {
|
||||||
@@ -68,7 +82,7 @@ func (srv *Server) Serve(ctx context.Context, listener net.Listener) error {
|
|||||||
}
|
}
|
||||||
return fmt.Errorf("smtp: verbindung annehmen: %w", err)
|
return fmt.Errorf("smtp: verbindung annehmen: %w", err)
|
||||||
}
|
}
|
||||||
session := newSession(conn, srv.sink, srv.maxMessageBytes, srv.tlsConfig, srv.logger)
|
session := newSession(conn, srv.sink, srv.maxMessageBytes, srv.tlsConfig, srv.logger, srv.tenantID, srv.acceptLimiter)
|
||||||
go session.Serve(ctx)
|
go session.Serve(ctx)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/protolog"
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/ratelimit"
|
||||||
)
|
)
|
||||||
|
|
||||||
// maxCommandLineBytes begrenzt eine einzelne Kommando-/DATA-Zeile
|
// maxCommandLineBytes begrenzt eine einzelne Kommando-/DATA-Zeile
|
||||||
@@ -34,12 +35,15 @@ type Session struct {
|
|||||||
|
|
||||||
log *protolog.SessionLogger // ING-08, nie nil (log.Event() ist nil-sicher)
|
log *protolog.SessionLogger // ING-08, nie nil (log.Event() ist nil-sicher)
|
||||||
|
|
||||||
|
tenantID string
|
||||||
|
acceptLimiter *ratelimit.Limiter // ING-09, nil = kein Rate-Limiting
|
||||||
|
|
||||||
state State
|
state State
|
||||||
from string
|
from string
|
||||||
to []string
|
to []string
|
||||||
}
|
}
|
||||||
|
|
||||||
func newSession(conn net.Conn, sink MessageSink, maxMessageBytes int64, tlsConfig *tls.Config, logger *slog.Logger) *Session {
|
func newSession(conn net.Conn, sink MessageSink, maxMessageBytes int64, tlsConfig *tls.Config, logger *slog.Logger, tenantID string, acceptLimiter *ratelimit.Limiter) *Session {
|
||||||
_, alreadyTLS := conn.(*tls.Conn)
|
_, alreadyTLS := conn.(*tls.Conn)
|
||||||
return &Session{
|
return &Session{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
@@ -50,10 +54,22 @@ func newSession(conn net.Conn, sink MessageSink, maxMessageBytes int64, tlsConfi
|
|||||||
tlsConfig: tlsConfig,
|
tlsConfig: tlsConfig,
|
||||||
tlsActive: alreadyTLS,
|
tlsActive: alreadyTLS,
|
||||||
log: protolog.NewSessionLogger(logger, "smtp"),
|
log: protolog.NewSessionLogger(logger, "smtp"),
|
||||||
|
tenantID: tenantID,
|
||||||
|
acceptLimiter: acceptLimiter,
|
||||||
state: Greeting,
|
state: Greeting,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sourceAddr liefert die Quell-IP dieser Verbindung ohne Portanteil
|
||||||
|
// (ING-09).
|
||||||
|
func (s *Session) sourceAddr() string {
|
||||||
|
host, _, err := net.SplitHostPort(s.conn.RemoteAddr().String())
|
||||||
|
if err != nil {
|
||||||
|
return s.conn.RemoteAddr().String()
|
||||||
|
}
|
||||||
|
return host
|
||||||
|
}
|
||||||
|
|
||||||
// State liefert den aktuellen Sitzungszustand (für Tests).
|
// State liefert den aktuellen Sitzungszustand (für Tests).
|
||||||
func (s *Session) State() State { return s.state }
|
func (s *Session) State() State { return s.state }
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user