From bb44aef7183a167ac8ee1b9c2087e5376a1e4489 Mon Sep 17 00:00:00 2001 From: samtoxie Date: Sun, 26 Jul 2026 13:51:22 +0200 Subject: [PATCH] =?UTF-8?q?=E2=9C=A8=20feat:=20Allow=20cache=20reading=20f?= =?UTF-8?q?rom=20replicas=20(#342)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: Allow cache reading from replicas * :bento: fix logic * :sparkles: add testing for redis with mock * :bento: fix permission * πŸ“ test(e2e/redis): fix swapped IPβ†’verdict comments The mock returns "f" (not banned) for 1.2.3.4 and "t" (banned) for 1.2.3.5, and the run.sh assertions match that. Both doc comments described the opposite mapping; correct them to match the code. Co-Authored-By: Claude Opus 4.8 (1M context) * βœ… test(e2e/redis): exercise read-from-replica path The redis scenario only set redisCacheHost, so it validated the writer but never the round-robin reader path this feature adds. Split the mock into two roles: the primary (--redis-addr) now answers every GET with a miss, while the replica (--redis-read-addr) serves the hardcoded verdicts. The scenario points redisCacheReadHosts at the replica (twice, to drive round-robin), so the banned-IP-blocked assertion only passes if the plugin actually reads decisions from the replica rather than the primary. Co-Authored-By: Claude Opus 4.8 (1M context) * πŸ“ docs: note replicas don't fall back to primary on outage When RedisCacheReadHosts is set, reads are not retried against the primary if the replicas are unreachable. Document that this, combined with the default RedisCacheUnreachableBlock=true, means a replica outage can block traffic while the primary is healthy. Co-Authored-By: Claude Opus 4.8 (1M context) * πŸ› fix(cache): avoid nil-pointer panic on empty redis read redisCache.get fell through to `switch err.Error()` when Get returned a nil error with an empty value, panicking on the nil error. simpleredis never returns that combination today (a miss yields RedisMiss), so it was unreachable in practice β€” but the read path is safer treating an empty, error-free read as a cache miss, which also guarantees err is non-nil before err.Error() is called. Co-Authored-By: Claude Opus 4.8 (1M context) * :bento: add test for rotation * :bug: readd redis/run.sh * :bug: fix redis/run.sh * :bug: fix test label; remove log * :bento: add tests for roundRobin * :bug: fix test --------- Co-authored-by: maxlerebourg Co-authored-by: mhx Co-authored-by: Claude Opus 4.8 (1M context) --- .golangci.yml | 1 + Makefile | 5 +- README.md | 12 +++- bouncer.go | 1 + pkg/cache/cache.go | 67 ++++++++++++++-------- pkg/cache/cache_test.go | 38 ++++++++++++ pkg/configuration/configuration.go | 2 + tests/e2e/mock/lib/common.sh | 6 ++ tests/e2e/mock/mocklapi/main.go | 62 +++++++++++++++++++- tests/e2e/mock/scenarios/redis/dynamic.yml | 30 ++++++++++ tests/e2e/mock/scenarios/redis/run.sh | 26 +++++++++ 11 files changed, 218 insertions(+), 32 deletions(-) create mode 100644 tests/e2e/mock/scenarios/redis/dynamic.yml create mode 100644 tests/e2e/mock/scenarios/redis/run.sh diff --git a/.golangci.yml b/.golangci.yml index 34b62ea..33598d6 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -41,6 +41,7 @@ linters-settings: - $test allow: - $gostd + - github.com/maxlerebourg/simpleredis - github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger - github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip - github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/configuration diff --git a/Makefile b/Makefile index 11a2f61..5dcec89 100644 --- a/Makefile +++ b/Makefile @@ -4,7 +4,7 @@ export GO111MODULE=on # Binary/mock suite (Traefik binary + mock LAPI). This is what CI runs. # The local Docker suite (make e2e) lives in a separate PR/branch. -E2E_MOCK_SCENARIOS := stream-mode live-mode none-mode trusted-ips custom-ban-page captcha appsec tls-system-ca +E2E_MOCK_SCENARIOS := $(notdir $(wildcard tests/e2e/mock/scenarios/*)) default: lint test @@ -20,7 +20,7 @@ yaegi_test: e2e_mock: $(addprefix e2e_mock_,$(E2E_MOCK_SCENARIOS)) e2e_mock_%: - ./tests/e2e/mock/scenarios/$*/run.sh + bash ./tests/e2e/mock/scenarios/$*/run.sh vendor: go mod vendor @@ -124,4 +124,3 @@ show_metrics: show_decisions: docker exec crowdsec cscli decisions list - diff --git a/README.md b/README.md index 4098944..1c069f7 100644 --- a/README.md +++ b/README.md @@ -444,7 +444,12 @@ make run - RedisCacheHost - string - default: "redis:6379" - - hostname and port for the Redis service + - hostname and port for the Redis write host (primary) +- RedisCacheReadHosts + - []string + - default: [] + - List of Redis replica hostnames (host:port) to use for read operations. Reads are distributed round-robin across replicas. Falls back to RedisCacheHost when empty. + - Note: when set, reads are not retried against RedisCacheHost (the primary) if the replicas are unreachable. With RedisCacheUnreachableBlock at its default (true), a replica outage will therefore block/delay requests even though the primary is healthy. - RedisCachePassword - string - default: "" @@ -640,7 +645,10 @@ http: forwardedHeadersCustomName: X-Custom-Header remediationHeadersCustomName: cs-remediation redisCacheEnabled: false - redisCacheHost: "redis:6379" + redisCacheHost: "redis-primary:6379" + redisCacheReadHosts: + - "redis-replica-1:6379" + - "redis-replica-2:6379" redisCachePassword: password redisCacheDatabase: "5" redisCacheUnreachableBlock: true diff --git a/bouncer.go b/bouncer.go index a1746e8..d362996 100644 --- a/bouncer.go +++ b/bouncer.go @@ -268,6 +268,7 @@ func New(_ context.Context, next http.Handler, config *configuration.Config, nam log, config.RedisCacheEnabled, config.RedisCacheHost, + config.RedisCacheReadHosts, config.RedisCachePassword, config.RedisCacheDatabase, ) diff --git a/pkg/cache/cache.go b/pkg/cache/cache.go index e966638..3059f01 100644 --- a/pkg/cache/cache.go +++ b/pkg/cache/cache.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "log/slog" + "sync/atomic" ttl_map "github.com/leprosus/golang-ttl-map" simpleredis "github.com/maxlerebourg/simpleredis" @@ -27,10 +28,7 @@ const ( ) //nolint:gochecknoglobals -var ( - redis simpleredis.SimpleRedis - cache = ttl_map.New() -) +var cache = ttl_map.New() type localCache struct{} @@ -52,33 +50,48 @@ func (localCache) delete(key string) { } type redisCache struct { - log *slog.Logger + log *slog.Logger + writer simpleredis.SimpleRedis + readers []simpleredis.SimpleRedis + counter atomic.Uint64 } -func (redisCache) get(key string) (string, error) { - value, err := redis.Get(key) +func (rc *redisCache) nextReader() *simpleredis.SimpleRedis { + n := len(rc.readers) + if n == 0 { + return &rc.writer + } + idx := rc.counter.Add(1) % uint64(n) + return &rc.readers[idx] +} + +func (rc *redisCache) get(key string) (string, error) { + value, err := rc.nextReader().Get(key) + if err != nil { + switch err.Error() { + case simpleredis.RedisMiss: + return "", errors.New(CacheMiss) + case simpleredis.RedisUnreachable: + return "", errors.New(CacheUnreachable) + default: + return "", err + } + } valueString := string(value) - if err == nil && len(valueString) > 0 { + if len(valueString) > 0 { return valueString, nil } - errRedisMessage := err.Error() - if errRedisMessage == simpleredis.RedisMiss { - return "", errors.New(CacheMiss) - } - if errRedisMessage == simpleredis.RedisUnreachable { - return "", errors.New(CacheUnreachable) - } - return "", err + return "", errors.New(CacheMiss) } -func (rc redisCache) set(key, value string, duration int64) { - if err := redis.Set(key, []byte(value), duration); err != nil { +func (rc *redisCache) set(key, value string, duration int64) { + if err := rc.writer.Set(key, []byte(value), duration); err != nil { rc.log.Error("cache:setDecisionRedisCache" + err.Error()) } } -func (rc redisCache) delete(key string) { - if err := redis.Del(key); err != nil { +func (rc *redisCache) delete(key string) { + if err := rc.writer.Del(key); err != nil { rc.log.Error("cache:deleteDecisionRedisCache " + err.Error()) } } @@ -96,15 +109,21 @@ type Client struct { } // New Initialize cache client. -func (c *Client) New(log *slog.Logger, isRedis bool, host, pass, database string) { +func (c *Client) New(log *slog.Logger, isRedis bool, writeHost string, readHosts []string, pass, database string) { c.log = log if isRedis { - redis.Init(host, pass, database) - c.cache = &redisCache{log: log} + rc := &redisCache{log: log} + rc.writer.Init(writeHost, pass, database) + for _, h := range readHosts { + var r simpleredis.SimpleRedis + r.Init(h, pass, database) + rc.readers = append(rc.readers, r) + } + c.cache = rc } else { c.cache = &localCache{} } - c.log.Debug(fmt.Sprintf("cache:New initialized isRedis:%v", isRedis)) + c.log.Debug(fmt.Sprintf("cache:New initialized isRedis:%v writeHost:%v readHosts:%v", isRedis, writeHost, readHosts)) } // Delete delete decision in cache. diff --git a/pkg/cache/cache_test.go b/pkg/cache/cache_test.go index ce62d6e..3aa3aea 100644 --- a/pkg/cache/cache_test.go +++ b/pkg/cache/cache_test.go @@ -6,6 +6,7 @@ import ( "testing" logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger" + simpleredis "github.com/maxlerebourg/simpleredis" ) func Test_Get(t *testing.T) { @@ -122,3 +123,40 @@ func Test_Delete(t *testing.T) { }) } } + +// indexOfReader returns the position of r inside rc.readers, or -1 when r is the writer (the no-readers fallback). +func indexOfReader(rc *redisCache, r *simpleredis.SimpleRedis) int { + if r == &rc.writer { + return -1 + } + for i := range rc.readers { + if r == &rc.readers[i] { + return i + } + } + return -2 +} + +func Test_nextReader(t *testing.T) { + // The counter starts at 0, so the first Add(1) yields index 1, then 2, 0, 1, ... over n readers. + tests := []struct { + name string + readers int + want []int + }{ + {name: "round-robin over three readers", readers: 3, want: []int{1, 2, 0, 1, 2, 0, 1}}, + {name: "single reader always selected", readers: 1, want: []int{0, 0, 0, 0, 0}}, + {name: "no readers fall back to writer", readers: 0, want: []int{-1, -1, -1}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + rc := &redisCache{log: logger.New("INFO", "")} + rc.readers = make([]simpleredis.SimpleRedis, tt.readers) + for call, want := range tt.want { + if got := indexOfReader(rc, rc.nextReader()); got != want { + t.Errorf("call %d: nextReader() -> reader[%d], want reader[%d]", call, got, want) + } + } + }) + } +} diff --git a/pkg/configuration/configuration.go b/pkg/configuration/configuration.go index fda49d7..6a7a2e5 100644 --- a/pkg/configuration/configuration.go +++ b/pkg/configuration/configuration.go @@ -96,6 +96,7 @@ type Config struct { ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"` RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"` RedisCacheHost string `json:"redisCacheHost,omitempty"` + RedisCacheReadHosts []string `json:"redisCacheReadHosts,omitempty"` RedisCachePassword string `json:"redisCachePassword,omitempty"` RedisCachePasswordFile string `json:"redisCachePasswordFile,omitempty"` RedisCacheDatabase string `json:"redisCacheDatabase,omitempty"` @@ -172,6 +173,7 @@ func New() *Config { ClientTrustedIPs: []string{}, RedisCacheEnabled: false, RedisCacheHost: "redis:6379", + RedisCacheReadHosts: []string{}, RedisCachePassword: "", RedisCacheDatabase: "", RedisCacheUnreachableBlock: true, diff --git a/tests/e2e/mock/lib/common.sh b/tests/e2e/mock/lib/common.sh index 66ad627..d5cdaab 100644 --- a/tests/e2e/mock/lib/common.sh +++ b/tests/e2e/mock/lib/common.sh @@ -20,6 +20,8 @@ WEB_PORT="${WEB_PORT:-8000}" LAPI_PORT="${LAPI_PORT:-8090}" BACKEND_PORT="${BACKEND_PORT:-8091}" APPSEC_PORT="${APPSEC_PORT:-8092}" +REDIS_PORT="${REDIS_PORT:-8093}" +REDIS_READ_PORT="${REDIS_READ_PORT:-8094}" LAPI_KEY="${LAPI_KEY:-e2e-mock-key}" MOCK_LIB_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" @@ -187,6 +189,8 @@ start_stack() { -e "s|@@LAPI_HOST@@|127.0.0.1:${LAPI_PORT}|g" \ -e "s|@@APPSEC_HOST@@|127.0.0.1:${APPSEC_PORT}|g" \ -e "s|@@BACKEND_URL@@|http://127.0.0.1:${BACKEND_PORT}|g" \ + -e "s|@@REDIS_HOST@@|127.0.0.1:${REDIS_PORT}|g" \ + -e "s|@@REDIS_READ_HOST@@|127.0.0.1:${REDIS_READ_PORT}|g" \ -e "s|@@SCENARIO_DIR@@|${scenario_dir}|g" \ "$scenario_dir/dynamic.yml" > "$WORKDIR/dynamic.yml" @@ -203,6 +207,8 @@ start_stack() { --lapi-addr "127.0.0.1:${LAPI_PORT}" \ --backend-addr "127.0.0.1:${BACKEND_PORT}" \ --appsec-addr "127.0.0.1:${APPSEC_PORT}" \ + --redis-addr "127.0.0.1:${REDIS_PORT}" \ + --redis-read-addr "127.0.0.1:${REDIS_READ_PORT}" \ "${mock_tls_args[@]}" >"$WORKDIR/mock.log" 2>&1 & MOCK_PID=$! diff --git a/tests/e2e/mock/mocklapi/main.go b/tests/e2e/mock/mocklapi/main.go index cb0ba11..65f5ad6 100644 --- a/tests/e2e/mock/mocklapi/main.go +++ b/tests/e2e/mock/mocklapi/main.go @@ -2,7 +2,8 @@ // suite. It answers only the few LAPI routes the plugin calls β€” live/none // decision lookups, the stream poll and the usage-metrics push β€” and lets the // test drive decisions through /admin instead of `cscli`. It also serves the -// stub upstream that Traefik proxies allowed requests to. +// stub upstream that Traefik proxies allowed requests to, and a hardcoded Redis +// stand-in for exercising the redis cache path. // // It is NOT a Crowdsec/AppSec conformance harness β€” the real WAF engine (OWASP // CRS, virtual patching) is out of scope. The AppSec endpoint here emulates a @@ -11,10 +12,12 @@ package main import ( + "bufio" "encoding/json" "flag" "io" "log" + "net" "net/http" "strings" "sync" @@ -46,6 +49,50 @@ func list(m map[string]Decision) []Decision { return out } +// --- Redis mock (inline-command wire format, as spoken by simpleredis) --- + +// serveRedis is a hardcoded stand-in. When verdicts is true it plays a replica +// that holds decisions: every line is scanned for known IPs, 1.2.3.4 β†’ "f" +// (clean), 1.2.3.5 β†’ "t" (banned); any other GET is a miss ($-1). When verdicts +// is false it plays the primary and answers every GET with a miss, so a +// scenario can prove reads are served from the replica and not the primary. +// SET, DEL, AUTH, SELECT get +OK (they don't read the response anyway). +func serveRedis(addr string, verdicts bool) { + ln, err := net.Listen("tcp", addr) + if err != nil { + log.Fatal(err) + } + defer ln.Close() + + for { + conn, err := ln.Accept() + if err != nil { + continue + } + go func(conn net.Conn) { + defer conn.Close() + rd := bufio.NewReader(conn) + for { + line, _, err := rd.ReadLine() + if err != nil { + return + } + s := string(line) + switch { + case verdicts && strings.Contains(s, "1.2.3.4"): + conn.Write([]byte("$1\r\nf\r\n")) + case verdicts && strings.Contains(s, "1.2.3.5"): + conn.Write([]byte("$1\r\nt\r\n")) + case strings.HasPrefix(strings.ToUpper(s), "GET "): + conn.Write([]byte("$-1\r\n")) + default: + conn.Write([]byte("+OK\r\n")) + } + } + }(conn) + } +} + func main() { lapiAddr := flag.String("lapi-addr", "127.0.0.1:8090", "address for the LAPI mock") // The stub upstream Traefik proxies allowed requests to β€” the binary-suite @@ -53,6 +100,12 @@ func main() { backendAddr := flag.String("backend-addr", "127.0.0.1:8091", "address for the stub upstream service") // AppSec WAF stand-in (the real engine listens on :7422). Not a CRS engine. appsecAddr := flag.String("appsec-addr", "127.0.0.1:8092", "address for the AppSec mock") + // Redis stand-ins on plain TCP ports, enough to exercise the plugin's redis + // cache path. The primary answers every GET with a miss; the replica serves + // the hardcoded verdicts, so a scenario pointing redisCacheReadHosts at the + // replica proves reads are offloaded to replicas. + redisAddr := flag.String("redis-addr", "127.0.0.1:8093", "address for the Redis primary mock (writes; GET always misses)") + redisReadAddr := flag.String("redis-read-addr", "127.0.0.1:8094", "address for the Redis replica mock (serves cached verdicts)") // Optional TLS for the LAPI: when both are set the LAPI is served over HTTPS // (cert signed by the scenario's throwaway CA) so the suite can exercise the // bouncer's system-trust-store path. Backend and AppSec stay plaintext. @@ -96,6 +149,9 @@ func main() { }))) }() + go serveRedis(*redisAddr, false) + go serveRedis(*redisReadAddr, true) + mux := http.NewServeMux() // Readiness probe for the test harness (empty body, 200). @@ -154,9 +210,9 @@ func main() { }) if *lapiTLSCert != "" && *lapiTLSKey != "" { - log.Printf("mocklapi: LAPI on %s (TLS), backend on %s, appsec on %s", *lapiAddr, *backendAddr, *appsecAddr) + log.Printf("mocklapi: LAPI on %s (TLS), backend on %s, appsec on %s, redis on %s (read %s)", *lapiAddr, *backendAddr, *appsecAddr, *redisAddr, *redisReadAddr) log.Fatal(http.ListenAndServeTLS(*lapiAddr, *lapiTLSCert, *lapiTLSKey, mux)) } - log.Printf("mocklapi: LAPI on %s, backend on %s, appsec on %s", *lapiAddr, *backendAddr, *appsecAddr) + log.Printf("mocklapi: LAPI on %s, backend on %s, appsec on %s, redis on %s (read %s)", *lapiAddr, *backendAddr, *appsecAddr, *redisAddr, *redisReadAddr) log.Fatal(http.ListenAndServe(*lapiAddr, mux)) } diff --git a/tests/e2e/mock/scenarios/redis/dynamic.yml b/tests/e2e/mock/scenarios/redis/dynamic.yml new file mode 100644 index 0000000..317916b --- /dev/null +++ b/tests/e2e/mock/scenarios/redis/dynamic.yml @@ -0,0 +1,30 @@ +http: + routers: + r: + rule: "PathPrefix(`/foo`)" + entryPoints: + - web + service: backend + middlewares: + - bouncer + services: + backend: + loadBalancer: + servers: + - url: "@@BACKEND_URL@@" + middlewares: + bouncer: + plugin: + bouncer: + enabled: "true" + crowdsecMode: live + crowdsecLapiScheme: http + crowdsecLapiHost: "@@LAPI_HOST@@" + crowdsecLapiKey: "@@APIKEY@@" + redisCacheEnabled: "true" + redisCacheHost: "@@REDIS_HOST@@" + redisCacheReadHosts: + - "@@REDIS_READ_HOST@@" + - "@@REDIS_HOST@@" + forwardedHeadersTrustedIps: + - "127.0.0.1/32" diff --git a/tests/e2e/mock/scenarios/redis/run.sh b/tests/e2e/mock/scenarios/redis/run.sh new file mode 100644 index 0000000..65b0c89 --- /dev/null +++ b/tests/e2e/mock/scenarios/redis/run.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash +set -euo pipefail + +HERE="$(cd "$(dirname "$0")" && pwd)" +# shellcheck source=../../lib/common.sh +source "$HERE/../../lib/common.sh" + +SCENARIO=redis + +# The replica mock returns "f" (not banned) for 1.2.3.4 and "t" (banned) for 1.2.3.5. +# The primary mock always misses. +body() { + echo "[$SCENARIO] cached banned IP must not be blocked because call for primary (test rotation)" + assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.5" + + echo "[$SCENARIO] cached banned IP must be blocked" + assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.5" + + echo "[$SCENARIO] cached clean IP must pass" + assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4" + + echo "[$SCENARIO] unknown IP (redis miss) must fall through to LAPI and pass" + assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.6" +} + +run_scenario "$SCENARIO" "$HERE" body