Compare commits

...
Author SHA1 Message Date
maxlerebourg 9ea2ea2419 🍱 renovate every day 2026-07-26 14:23:05 +02:00
maxlerebourg 1dedaa0571 Renovate update version.go 2026-07-26 14:13:21 +02:00
bb44aef718 feat: Allow cache reading from replicas (#342)
* feat: Allow cache reading from replicas

* 🍱 fix logic

*  add testing for redis with mock

* 🍱 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) <noreply@anthropic.com>

*  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) <noreply@anthropic.com>

* 📝 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) <noreply@anthropic.com>

* 🐛 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) <noreply@anthropic.com>

* 🍱 add test for rotation

* 🐛 readd redis/run.sh

* 🐛 fix redis/run.sh

* 🐛 fix test label; remove log

* 🍱 add tests for roundRobin

* 🐛 fix test

---------

Co-authored-by: maxlerebourg <maxlerebourg@gmail.com>
Co-authored-by: mhx <mathieu@hanotaux.fr>
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-26 13:51:22 +02:00
15 changed files with 232 additions and 82 deletions
-46
View File
@@ -1,46 +0,0 @@
name: Release Version Update
on:
release:
types: [published]
permissions:
contents: write
jobs:
update-version:
name: Update version in source
runs-on: ubuntu-latest
steps:
- name: Checkout code
uses: actions/checkout@v7
with:
ref: main
- name: Extract version from tag
id: get_version
run: |
TAG="${{ github.event.release.tag_name }}"
VERSION="${TAG#v}"
echo "version=$VERSION" >> "$GITHUB_OUTPUT"
echo "tag=$TAG" >> "$GITHUB_OUTPUT"
- name: Update version in version.go
run: |
sed -i 's/pluginVersion = "[^"]*"/pluginVersion = "'"${{ steps.get_version.outputs.version }}"'"/' version.go
cat version.go
- name: Commit, push, and retag
run: |
git config user.name "github-actions[bot]"
git config user.email "github-actions[bot]@users.noreply.github.com"
git add version.go
if git diff --cached --quiet; then
echo "Version already up to date, nothing to commit"
exit 0
fi
git commit -m "⬆️ chore: bump version to ${{ steps.get_version.outputs.version }}"
git push origin main
# Move the release tag to include the version update
git tag -f "${{ steps.get_version.outputs.tag }}"
git push -f origin "${{ steps.get_version.outputs.tag }}"
+2 -2
View File
@@ -1,6 +1,6 @@
name: Renovate
# Self-hosted Renovate: opens dependency-update PRs on a weekly schedule.
# Self-hosted Renovate: opens dependency-update PRs on a daily schedule.
# Config lives in /renovate.json. Requires a repo/org secret RENOVATE_TOKEN
# (a PAT with `repo` + `workflow` scope, or a fine-grained token with
# contents:write + pull-requests:write) so Renovate can push branches and open
@@ -8,7 +8,7 @@ name: Renovate
on:
schedule:
- cron: "0 4 * * 1" # every Monday at 04:00 UTC
- cron: "0 4 * * *" # every day at 04:00 UTC
workflow_dispatch:
inputs:
logLevel:
+1
View File
@@ -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
+2 -3
View File
@@ -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
+10 -2
View File
@@ -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
+1
View File
@@ -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,
)
+43 -24
View File
@@ -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.
+38
View File
@@ -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)
}
}
})
}
}
+2
View File
@@ -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,
+10
View File
@@ -25,6 +25,16 @@
}
],
"customManagers": [
{
"description": "Plugin self-pin in version.go (pluginVersion)",
"customType": "regex",
"managerFilePatterns": ["/^version\\.go$/"],
"matchStrings": [
"pluginVersion\\s*=\\s*\"(?<currentValue>v[0-9]+\\.[0-9]+\\.[0-9]+)\""
],
"depNameTemplate": "maxlerebourg/crowdsec-bouncer-traefik-plugin",
"datasourceTemplate": "github-tags"
},
{
"description": "Plugin self-pin in docker-compose CLI args (--experimental.plugins.bouncer.version=vX)",
"customType": "regex",
+6
View File
@@ -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=$!
+59 -3
View File
@@ -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))
}
@@ -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"
+26
View File
@@ -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
+2 -2
View File
@@ -1,4 +1,4 @@
package crowdsec_bouncer_traefik_plugin //nolint:revive,stylecheck
// pluginVersion is updated automatically by the release workflow.
var pluginVersion = "1.6.X" //nolint:gochecknoglobals
// pluginVersion is updated automatically by the release workflow and Renovate.
var pluginVersion = "v1.6.0" //nolint:gochecknoglobals