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