Compare commits

...
Author SHA1 Message Date
maxlerebourg 0f8608d770 Merge remote-tracking branch 'origin' into 271-feature-support-decision-scope 2026-08-24 20:51:27 +02:00
maxlerebourgandRenovate Bot ae7481caa5 ⬆️ renovate: Update all (#375)
Co-authored-by: Renovate Bot <22881669+maxlerebourg@users.noreply.github.com>
2026-08-20 12:05:45 +02:00
mathieuHaandClaude Opus 5 955391671c ♻️ review fixes for #368: CIDR key correctness, lookup cost, tests (#372)
* ♻️ cidr: build keys through net.IPNet instead of hand-masking bytes

CIDRKeys masked the address byte by byte and formatted the result with
string concatenation, while SetCIDR/DeleteCIDR format their keys with
net.IPNet.String() via NormalizeCIDR. The two agreed only by coincidence:
any divergence in formatting silently stops every range decision from
matching, with no test covering the invariant.

Mask with net.IP.Mask and format through net.IPNet.String() so both sides
go through the same formatter. Output is byte for byte identical to the
previous implementation (checked against a golden dump of both IPv4 and
IPv6 keys, including ::ffff: forms).

Dropping the inner byte loops also removes the only intrange violation in
the tree, so the linter exclusion added for them is no longer needed, and
the redundant import alias on pkg/ip goes away with it.

*  cidr: only probe the prefix lengths that have a decision

In stream mode nothing caches a negative result per IP, so the exact IP
lookup misses on every legitimate request and each one fell through to
GetCIDR, which probed every possible prefix length: 33 cache reads for an
IPv4 client, 129 for an IPv6 one, even when no range decision existed at
all. On the local cache that is wasted work on the request path; with
redis it is 33 to 129 sequential round trips per request.

Keep the set of prefix lengths that have at least one decision under a
single key, written before the decision itself, and probe only those.
Measured cache reads per request: 1 with no range decision (was 33 / 129),
2 with a single /24 in use, 4 with four prefix lengths in use.

The set only grows, so a deleted or expired decision leaves a length
behind that costs one extra read rather than risking an unmatched
decision, and it is written with an effectively infinite duration since it
has to outlive every decision it describes. If it is ever missing while
decisions live (a redis eviction under maxmemory), range decisions stop
matching until the next one arrives; it is the hottest key of the
namespace, so an LRU policy evicts it last.

* 🔊 cidr: log the decisions dropped for an unparsable CIDR

SetCIDR and DeleteCIDR returned silently when NormalizeCIDR rejected the
value, so a range decision the plugin does not understand is not enforced
and nothing says why. Every other operation of the package logs, and this
one fails open, which is the direction worth shouting about.

Log at Error with the raw value and what the consequence is, so an
unexpected decision format shows up in the logs instead of looking like a
decision that was applied.

*  cidr: cover the invariant the range matching rests on

The helpers were tested in isolation but nothing tied them together, and
what actually has to hold is that the key SetCIDR writes for a decision is
one of the keys GetCIDR looks up for an IP that decision covers. A
formatting change on either side would have silently stopped every range
decision from matching with all tests green.

TestCIDRKeys_MatchNormalizeCIDR pins that both ways, including the cases
worth being explicit about: a decision that is not on a network address,
IPv4 mapped clients against an IPv4 range, and that the two families do not
mix. Dropping the mask in cidrKey fails 9 of its cases.

Also covers CIDRLookupKeys against CIDRKeys length by length, the IPv6 side
of the network address test, and CIDRPrefixLen.

*  cache: cover the CIDR operations and what a lookup costs

pkg/cache had tests for Get, Set and Delete but none for their CIDR
counterparts, so the range keyspace was only exercised end to end by the
e2e scenario.

Test_GetCIDR covers hits, the boundaries of a range, IPv6, IPv4 mapped
clients and invalid input. Test_GetCIDR_MostSpecific pins the precedence
between overlapping decisions, which is deliberate behaviour that nothing
was holding in place: a captcha on a /24 is not overruled by a ban on its
/8. Test_DeleteCIDR checks the wider decision survives a narrower one being
removed, and Test_SetCIDR_InvalidIsNotStored that a rejected value stores
nothing at all.

Test_GetCIDR_Reads counts cache reads through an isolated cacheInterface,
so the cost of a lookup is part of the contract: probing every prefix
length again turns it into 34 reads for IPv4 and 130 for IPv6 and fails.

* 🐛 e2e: stop racing the deadline in the stream failure check

handleStreamTicker compares updateFailure to updateMaxFailure before
incrementing it, so with updateMaxFailure 2 and a 1s interval the bouncer
gives up on the third consecutive failed poll, roughly 3s after the
endpoint starts failing. The check slept 2s and then polled for a 200 for
up to 15s, leaving about a second of margin, and once that window closes it
never reopens: a runner under load turns this into a 15s wait followed by a
failure.

Assert the 200 immediately after the endpoint starts failing, which is
always inside the window, and keep polling for the 403 that follows.

* 🐛 e2e: do not fail a scenario on a single slow response

The new -m 1 is right for the polling helpers, where a timed out request is
just another attempt, but the assertions and the mock control plane have no
retry: one request that takes over a second on a loaded runner fails the
scenario, and since common.sh runs under set -euo pipefail a timed out
lapi_add_decision aborts it before the decision even exists.

Bound the connect at 1s instead and give the whole request 5s in the seven
places that get a single attempt. Nothing waits longer on the happy path.

*  e2e: restore the stream-mode failure check

Revert 99673b1. The check is not ours to change: it must stay as it was.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* 🔥 e2e: drop the comment above wait_for_status

The timeout change in 3f78887 stands; only the comment goes.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-07 21:10:44 +02:00
mathieuHaandmaxlerebourg 9b8d6b937c 🐛 keep the stream lease alive when updateIntervalSeconds is 1 (#371)
* 🐛 keep the stream lease alive when updateIntervalSeconds is 1

handleStreamCache takes a lease so a single node polls LAPI per interval,
and stores it for updateInterval-1 seconds. The e2e stream scenario now
sets updateIntervalSeconds to 1, which makes that a 0 second duration:
golang-ttl-map returns early on a zero ttl (map.go:114) and redis rejects
a non positive EX, so the lease is never stored and the guard silently
does nothing.

Floor the duration at 1 second. At an interval of 1 the lease can survive
a tick that fires slightly early and cost one skipped poll, which is far
better than every node polling every tick against a shared redis.

* Adjust lease duration to prevent cache update conflicts

Updated lease duration logic to ensure a minimum of 1 second.

---------

Co-authored-by: maxlerebourg <maxlerebourg@gmail.com>
2026-08-07 19:11:54 +02:00
maxlerebourgandRenovate Bot d57ead2ec7 ⬆️ renovate: Update actions/setup-go action to v7 (#364)
Co-authored-by: Renovate Bot <22881669+maxlerebourg@users.noreply.github.com>
2026-08-07 12:15:10 +02:00
maxlerebourg 1ba556e918 🍱 fix lint 2026-07-31 14:39:09 +02:00
maxlerebourg 6b0518859d 🐛 Support range decision on stream mode 2026-07-31 14:08:59 +02:00
github-actions[bot]andgithub-actions[bot] bef5dfaadb 🔖 release v1.7.1 (#367)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-07-31 13:05:36 +02:00
99cf9712f4 cicd: bump the version before tagging so releases report their own version (#365)
*  cicd: bump the version before tagging instead of after

The version reported to the Crowdsec LAPI lives in version.go, so it must
be correct in the very commit the tag points at. Every mechanism so far
updated it *after* the tag existed, which cannot work:

- release.yml ran on `release: published` and force-moved the tag. It also
  failed on all four of its runs and was removed in #360.
- The Renovate customManager on version.go uses the github-tags datasource,
  so it can only propose vX once vX is already tagged. The bump always lands
  after the tag.

Result: v1.7.0 is tagged at a commit reading v1.6.0 (#363), same shape as
the earlier #322.

Replace both with a two-step flow that bumps first and tags last, so the
released source always matches its tag.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* 🍱 reduce loc + remove claude code comment

* 🍱 remove useless spellcheck disable

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: maxlerebourg <maxlerebourg@gmail.com>
2026-07-30 20:18:15 +02:00
maxlerebourgandRenovate Bot ed4a9e8262 ⬆️ renovate: Update all (#362)
Co-authored-by: Renovate Bot <22881669+maxlerebourg@users.noreply.github.com>
2026-07-27 21:09:55 +02:00
maxlerebourg f6ef95cf38 🐛 fix default value for CrowdsecAppsecUnreadableBodyBlock (#361) 2026-07-27 08:53:32 +02:00
maxlerebourgandRenovate Bot 9daba9739c ⬆️ renovate: Update all (#350)
Co-authored-by: Renovate Bot <22881669+maxlerebourg@users.noreply.github.com>
2026-07-26 17:06:02 +02:00
maxlerebourg e98b8ed5ba Renovate update version.go (#360)
*  Renovate update version.go

* 🍱 renovate every day
2026-07-26 17:05:43 +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
35 changed files with 1145 additions and 132 deletions
+1 -1
View File
@@ -19,7 +19,7 @@ jobs:
steps: steps:
- uses: actions/checkout@v7 - uses: actions/checkout@v7
- name: Set up Go - name: Set up Go
uses: actions/setup-go@v6 uses: actions/setup-go@v7
with: with:
# Track go.mod (Go 1.22) — the plugin's yaegi-bound floor. Keeps the # Track go.mod (Go 1.22) — the plugin's yaegi-bound floor. Keeps the
# single source of truth and builds the mock on the supported version. # single source of truth and builds the mock on the supported version.
+1 -1
View File
@@ -32,7 +32,7 @@ jobs:
# https://github.com/marketplace/actions/setup-go-environment # https://github.com/marketplace/actions/setup-go-environment
- name: Set up Go ${{ env.GO_VERSION }} - name: Set up Go ${{ env.GO_VERSION }}
uses: actions/setup-go@v6 uses: actions/setup-go@v7
with: with:
go-version: ${{ env.GO_VERSION }} go-version: ${{ env.GO_VERSION }}
+83
View File
@@ -0,0 +1,83 @@
name: Release (1/2) Prepare
# Step 1 of the release process: bump pluginVersion *before* the tag exists.
#
# The version reported to the Crowdsec LAPI lives in version.go, so it has to
# be correct in the very commit the tag points at. Anything that patches
# version.go after the release is published is too late: Traefik's plugin
# service caches the plugin archive per module+version, so users keep the
# source that was there when the tag was first resolved (see #322, #363).
#
# This workflow opens a "release" PR containing only that bump. Merging it
# triggers Release (2/2) Publish, which creates the tag and the GitHub release
# on the merged commit.
on:
workflow_dispatch:
inputs:
version:
description: "Version to release, e.g. v1.7.1 or v1.8.0-alpha"
required: true
type: string
permissions:
contents: write
pull-requests: write
jobs:
prepare:
name: Open release PR for ${{ inputs.version }}
runs-on: ubuntu-latest
steps:
- name: Check out main
uses: actions/checkout@v7
with:
ref: main
fetch-depth: 0
- name: Validate version
env:
VERSION: ${{ inputs.version }}
run: |
if ! [[ "$VERSION" =~ ^v[0-9]+\.[0-9]+\.[0-9]+(-[0-9A-Za-z.]+)?$ ]]; then
echo "::error::'$VERSION' is not a vX.Y.Z / vX.Y.Z-suffix version"
exit 1
fi
if git rev-parse -q --verify "refs/tags/$VERSION" >/dev/null; then
echo "::error::tag $VERSION already exists"
exit 1
fi
- name: Bump version.go
env:
VERSION: ${{ inputs.version }}
run: |
sed -i 's/pluginVersion = "[^"]*"/pluginVersion = "'"$VERSION"'"/' version.go
cat version.go
if git diff --quiet -- version.go; then
echo "::error::version.go already reads $VERSION, nothing to release"
exit 1
fi
- name: Push release branch and open PR
env:
GH_TOKEN: ${{ github.token }}
VERSION: ${{ inputs.version }}
run: |
git config user.name "github-actions[bot]"
git config user.email "github-actions[bot]@users.noreply.github.com"
git switch -c "release/$VERSION"
git commit -am "🔖 release $VERSION"
git push -u origin "release/$VERSION"
cat > /tmp/pr-body.md <<EOF
Bumps \`pluginVersion\` to \`$VERSION\` so the tag carries the version
the plugin reports to the Crowdsec LAPI.
Merging this PR tags \`$VERSION\` on the resulting commit and publishes
the GitHub release automatically.
> Keep the PR title as-is: **Release (2/2) Publish** matches on it.
EOF
gh pr create --base main --head "release/$VERSION" --title "🔖 release $VERSION" --body-file /tmp/pr-body.md
+53
View File
@@ -0,0 +1,53 @@
name: Release (2/2) Publish
# Step 2 of the release process: tag and publish the commit prepared by
# Release (1/2) Prepare.
#
# Triggered by the release PR landing on main. The tag is created on that
# commit, so version.go inside the released source always matches the tag —
# no post-release patching, no force-moved tags.
on:
push:
branches: [main]
paths: ["version.go"]
permissions:
contents: write
jobs:
publish:
name: Tag and publish
runs-on: ubuntu-latest
steps:
- name: Check out the pushed commit
uses: actions/checkout@v7
with:
fetch-depth: 0
- name: Resolve release version
id: resolve
run: |
version="$(git log -1 --format='%B' | grep -oP '🔖 release \Kv[0-9]+\.[0-9]+\.[0-9]+(-[0-9A-Za-z.]+)?' || true)"
[ -z "$version" ] && { echo "version.go changed outside a release commit, nothing to do"; echo "release=false" >> "$GITHUB_OUTPUT"; exit 0; }
in_source="$(sed -n 's/.*pluginVersion = "\([^"]*\)".*/\1/p' version.go)"
[ "$in_source" != "$version" ] && { echo "::error::commit says $version but version.go reads $in_source"; exit 1; }
git rev-parse -q --verify "refs/tags/$version" >/dev/null && { echo "::error::tag $version already exists"; exit 1; }
echo "release=true" >> "$GITHUB_OUTPUT"
echo "version=$version" >> "$GITHUB_OUTPUT"
echo "prerelease=$([[ "$version" == *-* ]] && echo '--prerelease')" >> "$GITHUB_OUTPUT"
- name: Tag and create the GitHub release
if: steps.resolve.outputs.release == 'true'
env:
GH_TOKEN: ${{ github.token }}
VERSION: ${{ steps.resolve.outputs.version }}
PRERELEASE: ${{ steps.resolve.outputs.prerelease }}
run: |
git config user.name "github-actions[bot]"
git config user.email "github-actions[bot]@users.noreply.github.com"
git tag -a "$VERSION" -m "$VERSION"
git push origin "$VERSION"
gh release create "$VERSION" --title "$VERSION" --generate-notes $PRERELEASE
-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 }}"
+3 -3
View File
@@ -1,6 +1,6 @@
name: Renovate 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 # Config lives in /renovate.json. Requires a repo/org secret RENOVATE_TOKEN
# (a PAT with `repo` + `workflow` scope, or a fine-grained token with # (a PAT with `repo` + `workflow` scope, or a fine-grained token with
# contents:write + pull-requests:write) so Renovate can push branches and open # contents:write + pull-requests:write) so Renovate can push branches and open
@@ -8,7 +8,7 @@ name: Renovate
on: on:
schedule: schedule:
- cron: "0 4 * * 1" # every Monday at 04:00 UTC - cron: "0 4 * * *" # every day at 04:00 UTC
workflow_dispatch: workflow_dispatch:
inputs: inputs:
logLevel: logLevel:
@@ -28,7 +28,7 @@ jobs:
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- name: Run Renovate - name: Run Renovate
uses: renovatebot/github-action@v46.1.17 uses: renovatebot/github-action@v46.2.2
with: with:
token: ${{ secrets.RENOVATE_TOKEN }} token: ${{ secrets.RENOVATE_TOKEN }}
env: env:
+1
View File
@@ -41,6 +41,7 @@ linters-settings:
- $test - $test
allow: allow:
- $gostd - $gostd
- github.com/maxlerebourg/simpleredis
- github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger - 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/ip
- github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/configuration - 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. # Binary/mock suite (Traefik binary + mock LAPI). This is what CI runs.
# The local Docker suite (make e2e) lives in a separate PR/branch. # 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 default: lint test
@@ -20,7 +20,7 @@ yaegi_test:
e2e_mock: $(addprefix e2e_mock_,$(E2E_MOCK_SCENARIOS)) e2e_mock: $(addprefix e2e_mock_,$(E2E_MOCK_SCENARIOS))
e2e_mock_%: e2e_mock_%:
./tests/e2e/mock/scenarios/$*/run.sh bash ./tests/e2e/mock/scenarios/$*/run.sh
vendor: vendor:
go mod vendor go mod vendor
@@ -124,4 +124,3 @@ show_metrics:
show_decisions: show_decisions:
docker exec crowdsec cscli decisions list docker exec crowdsec cscli decisions list
+12 -4
View File
@@ -384,8 +384,8 @@ make run
- Transmit only the first number of bytes to Crowdsec Appsec Server. - Transmit only the first number of bytes to Crowdsec Appsec Server.
- CrowdsecAppsecUnreadableBodyBlock - CrowdsecAppsecUnreadableBodyBlock
- bool - bool
- default: false - default: true
- Behaviour when the request body cannot be buffered for inspection (HTTP/2 or HTTP/3 request without a `Content-Length`, typically a bidirectional gRPC stream). When `false` (default) the request is forwarded to the Appsec Server with headers only (the body is left to stream through untouched). When `true` the request is blocked outright. Mirrors the reference bouncers' `APPSEC_DROP_UNREADABLE_BODY` option. - Behaviour when the request body cannot be buffered for inspection (HTTP/2 or HTTP/3 request without a `Content-Length`, typically a bidirectional gRPC stream). When `false` the request is forwarded to the Appsec Server with headers only (the body is left to stream through untouched). When `true` the request is blocked outright. Mirrors the reference bouncers' `APPSEC_DROP_UNREADABLE_BODY` option.
- CrowdsecAppsecKey - CrowdsecAppsecKey
- string - string
- default: value of `CrowdsecLapiKey` - default: value of `CrowdsecLapiKey`
@@ -444,7 +444,12 @@ make run
- RedisCacheHost - RedisCacheHost
- string - string
- default: "redis:6379" - 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 - RedisCachePassword
- string - string
- default: "" - default: ""
@@ -640,7 +645,10 @@ http:
forwardedHeadersCustomName: X-Custom-Header forwardedHeadersCustomName: X-Custom-Header
remediationHeadersCustomName: cs-remediation remediationHeadersCustomName: cs-remediation
redisCacheEnabled: false redisCacheEnabled: false
redisCacheHost: "redis:6379" redisCacheHost: "redis-primary:6379"
redisCacheReadHosts:
- "redis-replica-1:6379"
- "redis-replica-2:6379"
redisCachePassword: password redisCachePassword: password
redisCacheDatabase: "5" redisCacheDatabase: "5"
redisCacheUnreachableBlock: true redisCacheUnreachableBlock: true
+28 -4
View File
@@ -268,6 +268,7 @@ func New(_ context.Context, next http.Handler, config *configuration.Config, nam
log, log,
config.RedisCacheEnabled, config.RedisCacheEnabled,
config.RedisCacheHost, config.RedisCacheHost,
config.RedisCacheReadHosts,
config.RedisCachePassword, config.RedisCachePassword,
config.RedisCacheDatabase, config.RedisCacheDatabase,
) )
@@ -329,7 +330,7 @@ func New(_ context.Context, next http.Handler, config *configuration.Config, nam
// ServeHTTP principal function of plugin. // ServeHTTP principal function of plugin.
// //
//nolint:nestif //nolint:nestif,gocognit
func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) { func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
if !bouncer.enabled { if !bouncer.enabled {
bouncer.next.ServeHTTP(rw, req) bouncer.next.ServeHTTP(rw, req)
@@ -391,6 +392,16 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
// Right here if we cannot join the stream we forbid the request to go on. // Right here if we cannot join the stream we forbid the request to go on.
if bouncer.crowdsecMode == configuration.StreamMode || bouncer.crowdsecMode == configuration.AloneMode { if bouncer.crowdsecMode == configuration.StreamMode || bouncer.crowdsecMode == configuration.AloneMode {
if isCrowdsecStreamHealthy { if isCrowdsecStreamHealthy {
cidrValue, cidrErr := bouncer.cacheClient.GetCIDR(remoteIP)
if cidrErr == nil {
bouncer.log.Debug(fmt.Sprintf("ServeHTTP ip:%s cidr:hit isBanned:%v", remoteIP, cidrValue))
if cidrValue == cache.NoBannedValue {
bouncer.handleNextServeHTTP(rw, req, remoteIP)
} else {
bouncer.handleRemediationServeHTTP(rw, req, remoteIP, cidrValue)
}
return
}
bouncer.handleNextServeHTTP(rw, req, remoteIP) bouncer.handleNextServeHTTP(rw, req, remoteIP)
} else { } else {
bouncer.log.Debug(fmt.Sprintf("ServeHTTP isCrowdsecStreamHealthy:false ip:%s updateFailure:%d", remoteIP, updateFailure)) bouncer.log.Debug(fmt.Sprintf("ServeHTTP isCrowdsecStreamHealthy:false ip:%s updateFailure:%d", remoteIP, updateFailure))
@@ -640,7 +651,12 @@ func handleStreamCache(bouncer *Bouncer) error {
if err.Error() != cache.CacheMiss { if err.Error() != cache.CacheMiss {
return err return err
} }
bouncer.cacheClient.Set(cacheTimeoutKey, cache.NoBannedValue, bouncer.updateInterval-1) // To avoid every instance trying to update the cache, set 1 second at least
leaseDuration := bouncer.updateInterval - 1
if leaseDuration < 1 {
leaseDuration = 1
}
bouncer.cacheClient.Set(cacheTimeoutKey, cache.NoBannedValue, leaseDuration)
streamRouteURL := url.URL{ streamRouteURL := url.URL{
Scheme: bouncer.crowdsecScheme, Scheme: bouncer.crowdsecScheme,
Host: bouncer.crowdsecHost, Host: bouncer.crowdsecHost,
@@ -668,11 +684,19 @@ func handleStreamCache(bouncer *Bouncer) error {
default: default:
bouncer.log.Info("handleStreamCache:unknownType " + decision.Type) bouncer.log.Info("handleStreamCache:unknownType " + decision.Type)
} }
bouncer.cacheClient.Set(decision.Value, value, int64(duration.Seconds())) if strings.Contains(decision.Value, "/") {
bouncer.cacheClient.SetCIDR(decision.Value, value, int64(duration.Seconds()))
} else {
bouncer.cacheClient.Set(decision.Value, value, int64(duration.Seconds()))
}
} }
} }
for _, decision := range stream.Deleted { for _, decision := range stream.Deleted {
bouncer.cacheClient.Delete(decision.Value) if strings.Contains(decision.Value, "/") {
bouncer.cacheClient.DeleteCIDR(decision.Value)
} else {
bouncer.cacheClient.Delete(decision.Value)
}
} }
bouncer.log.Debug("handleStreamCache:updated") bouncer.log.Debug("handleStreamCache:updated")
isCrowdsecStreamStartup = false isCrowdsecStreamStartup = false
+1 -1
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -12,7 +12,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
volumes: volumes:
- "/var/run/docker.sock:/var/run/docker.sock:ro" - "/var/run/docker.sock:/var/run/docker.sock:ro"
# - './ban.html:/ban.html:ro' # - './ban.html:/ban.html:ro'
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+3 -3
View File
@@ -1,6 +1,6 @@
services: services:
cloudflare: cloudflare:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "cloudflare" container_name: "cloudflare"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -19,7 +19,7 @@ services:
- 8080:8080 - 8080:8080
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -33,7 +33,7 @@ services:
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5" - "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
- logs-traefik:/var/log/traefik - logs-traefik:/var/log/traefik
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -14,7 +14,7 @@ services:
- "--entrypoints.web.forwardedheaders.trustedips=172.18.0.0/24" - "--entrypoints.web.forwardedheaders.trustedips=172.18.0.0/24"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+2 -2
View File
@@ -1,5 +1,5 @@
image: image:
tag: v3.7.6 tag: v3.7.11
logs: logs:
general: general:
@@ -15,4 +15,4 @@ experimental:
plugins: plugins:
bouncer: bouncer:
moduleName: "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" moduleName: "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
version: "v1.6.0" version: "v1.7.1"
+3 -3
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
@@ -87,7 +87,7 @@ services:
- "traefik.enable=false" - "traefik.enable=false"
redis-secure: redis-secure:
image: "redis:8.8.0-alpine" image: "redis:8.10.0-alpine"
container_name: "redis-secure" container_name: "redis-secure"
hostname: redis-secure hostname: redis-secure
restart: unless-stopped restart: unless-stopped
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+2 -2
View File
@@ -1,6 +1,6 @@
services: services:
traefik: traefik:
image: "traefik:v3.7.6" image: "traefik:v3.7.11"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped restart: unless-stopped
command: command:
@@ -13,7 +13,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.6.0" - "--experimental.plugins.bouncer.version=v1.7.1"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
+130 -24
View File
@@ -6,9 +6,14 @@ import (
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
"strconv"
"strings"
"sync/atomic"
ttl_map "github.com/leprosus/golang-ttl-map" ttl_map "github.com/leprosus/golang-ttl-map"
simpleredis "github.com/maxlerebourg/simpleredis" simpleredis "github.com/maxlerebourg/simpleredis"
"github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip"
) )
const ( const (
@@ -24,13 +29,21 @@ const (
CacheMiss = "cache:miss" CacheMiss = "cache:miss"
// CacheUnreachable error string when cache is unreachable. // CacheUnreachable error string when cache is unreachable.
CacheUnreachable = "cache:unreachable" CacheUnreachable = "cache:unreachable"
// cidrPrefix namespaces a CIDR decision, one cache key per CIDR.
cidrPrefix = "cidr:"
// cidrPrefixLensKey holds the prefix lengths that have a decision, so a lookup probes
// only those. Its absence means no CIDR decision was ever stored.
cidrPrefixLensKey = "cidrprefixlens"
// cidrPrefixLensSeparator separates the prefix lengths in cidrPrefixLensKey.
cidrPrefixLensSeparator = ","
// cidrPrefixLensDuration has to outlive every decision it describes, so it is
// effectively infinite. It cannot be zero: the local cache ignores a zero duration
// and redis rejects a non positive EX.
cidrPrefixLensDuration = 10 * 365 * 24 * 60 * 60
) )
//nolint:gochecknoglobals //nolint:gochecknoglobals
var ( var cache = ttl_map.New()
redis simpleredis.SimpleRedis
cache = ttl_map.New()
)
type localCache struct{} type localCache struct{}
@@ -52,33 +65,48 @@ func (localCache) delete(key string) {
} }
type redisCache struct { 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) { func (rc *redisCache) nextReader() *simpleredis.SimpleRedis {
value, err := redis.Get(key) 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) valueString := string(value)
if err == nil && len(valueString) > 0 { if len(valueString) > 0 {
return valueString, nil return valueString, nil
} }
errRedisMessage := err.Error() return "", errors.New(CacheMiss)
if errRedisMessage == simpleredis.RedisMiss {
return "", errors.New(CacheMiss)
}
if errRedisMessage == simpleredis.RedisUnreachable {
return "", errors.New(CacheUnreachable)
}
return "", err
} }
func (rc redisCache) set(key, value string, duration int64) { func (rc *redisCache) set(key, value string, duration int64) {
if err := redis.Set(key, []byte(value), duration); err != nil { if err := rc.writer.Set(key, []byte(value), duration); err != nil {
rc.log.Error("cache:setDecisionRedisCache" + err.Error()) rc.log.Error("cache:setDecisionRedisCache" + err.Error())
} }
} }
func (rc redisCache) delete(key string) { func (rc *redisCache) delete(key string) {
if err := redis.Del(key); err != nil { if err := rc.writer.Del(key); err != nil {
rc.log.Error("cache:deleteDecisionRedisCache " + err.Error()) rc.log.Error("cache:deleteDecisionRedisCache " + err.Error())
} }
} }
@@ -96,15 +124,21 @@ type Client struct {
} }
// New Initialize cache client. // 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 c.log = log
if isRedis { if isRedis {
redis.Init(host, pass, database) rc := &redisCache{log: log}
c.cache = &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 { } else {
c.cache = &localCache{} 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. // Delete delete decision in cache.
@@ -125,3 +159,75 @@ func (c *Client) Set(key string, value string, duration int64) {
c.log.Debug(fmt.Sprintf("cache:Set key:%v value:%v duration:%vs", key, value, duration)) c.log.Debug(fmt.Sprintf("cache:Set key:%v value:%v duration:%vs", key, value, duration))
c.cache.set(key, value, duration) c.cache.set(key, value, duration)
} }
// DeleteCIDR removes a CIDR decision from the cache.
func (c *Client) DeleteCIDR(cidr string) {
normalized := ip.NormalizeCIDR(cidr)
if normalized == "" {
c.log.Error(fmt.Sprintf("cache:DeleteCIDR:invalidCIDR cidr:%v decision is left in cache", cidr))
return
}
cidr = normalized
c.cache.delete(cidrPrefix + cidr)
c.log.Debug(fmt.Sprintf("cache:DeleteCIDR cidr:%v", cidr))
}
// GetCIDR checks if an IP matches a CIDR decision in the cache.
// Only probes the prefix lengths that have a decision: it is on the request path.
func (c *Client) GetCIDR(ipStr string) (string, error) {
prefixLens, err := c.cache.get(cidrPrefixLensKey)
if err != nil {
return "", err
}
for _, key := range ip.CIDRLookupKeys(ipStr, parsePrefixLens(prefixLens)) {
value, getErr := c.cache.get(cidrPrefix + key)
if getErr == nil {
return value, nil
}
}
return "", errors.New(CacheMiss)
}
// SetCIDR stores a CIDR decision in the cache.
func (c *Client) SetCIDR(cidr, value string, duration int64) {
normalized := ip.NormalizeCIDR(cidr)
prefixLen := ip.CIDRPrefixLen(cidr)
if normalized == "" || prefixLen < 0 {
c.log.Error(fmt.Sprintf("cache:SetCIDR:invalidCIDR cidr:%v value:%v decision is not enforced", cidr, value))
return
}
// Publish the length first, or a concurrent lookup misses the decision.
c.addCIDRPrefixLen(prefixLen)
c.cache.set(cidrPrefix+normalized, value, duration)
c.log.Debug(fmt.Sprintf("cache:SetCIDR cidr:%v value:%v duration:%vs", normalized, value, duration))
}
// addCIDRPrefixLen records a prefix length in the set probed on lookup. The set only grows:
// a stale length costs one extra read, dropping one too early leaves decisions unmatched.
func (c *Client) addCIDRPrefixLen(prefixLen int) {
prefixLens, err := c.cache.get(cidrPrefixLensKey)
if err == nil {
for _, known := range parsePrefixLens(prefixLens) {
if known == prefixLen {
return
}
}
prefixLens += cidrPrefixLensSeparator + strconv.Itoa(prefixLen)
} else {
prefixLens = strconv.Itoa(prefixLen)
}
c.cache.set(cidrPrefixLensKey, prefixLens, cidrPrefixLensDuration)
c.log.Debug(fmt.Sprintf("cache:addCIDRPrefixLen prefixLens:%v", prefixLens))
}
// parsePrefixLens decodes the set of prefix lengths stored in cidrPrefixLensKey.
func parsePrefixLens(value string) []int {
fields := strings.Split(value, cidrPrefixLensSeparator)
prefixLens := make([]int, 0, len(fields))
for _, field := range fields {
if prefixLen, err := strconv.Atoi(field); err == nil {
prefixLens = append(prefixLens, prefixLen)
}
}
return prefixLens
}
+222
View File
@@ -3,9 +3,11 @@
package cache package cache
import ( import (
"errors"
"testing" "testing"
logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger" logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
simpleredis "github.com/maxlerebourg/simpleredis"
) )
func Test_Get(t *testing.T) { func Test_Get(t *testing.T) {
@@ -122,3 +124,223 @@ 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)
}
}
})
}
}
// countingCache is an isolated cacheInterface recording how many reads a lookup costs,
// so a CIDR lookup can be checked for both its result and its price.
type countingCache struct {
values map[string]string
reads int
}
func newCountingCache() *countingCache {
return &countingCache{values: map[string]string{}}
}
func (c *countingCache) get(key string) (string, error) {
c.reads++
if value, found := c.values[key]; found && value != "" {
return value, nil
}
return "", errors.New(CacheMiss)
}
func (c *countingCache) set(key, value string, _ int64) {
c.values[key] = value
}
func (c *countingCache) delete(key string) {
delete(c.values, key)
}
func newCIDRClient(decisions map[string]string) (*Client, *countingCache) {
counting := newCountingCache()
client := &Client{cache: counting, log: logger.New("INFO", "")}
for cidr, value := range decisions {
client.SetCIDR(cidr, value, 60)
}
return client, counting
}
func Test_GetCIDR(t *testing.T) {
decisions := map[string]string{
"10.0.0.0/24": BannedValue,
"192.168.1.42/24": CaptchaValue, // not a network address, host bits are dropped
"2001:db8::/32": BannedValue,
}
tests := []struct {
name string
clientIP string
want string
wantErr bool
}{
{name: "IP inside a banned range", clientIP: "10.0.0.7", want: BannedValue},
{name: "network address itself", clientIP: "10.0.0.0", want: BannedValue},
{name: "broadcast address of the range", clientIP: "10.0.0.255", want: BannedValue},
{name: "IP just outside the range", clientIP: "10.0.1.0", wantErr: true},
{name: "IP inside a captcha range", clientIP: "192.168.1.7", want: CaptchaValue},
{name: "IP inside an IPv6 range", clientIP: "2001:db8::dead:beef", want: BannedValue},
{name: "IP outside the IPv6 range", clientIP: "2001:db9::1", wantErr: true},
{name: "IPv4 mapped client against an IPv4 range", clientIP: "::ffff:10.0.0.7", want: BannedValue},
{name: "unknown IP", clientIP: "8.8.8.8", wantErr: true},
{name: "invalid IP", clientIP: "not-an-ip", wantErr: true},
{name: "empty IP", clientIP: "", wantErr: true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
client, _ := newCIDRClient(decisions)
got, err := client.GetCIDR(tt.clientIP)
if (err != nil) != tt.wantErr {
t.Fatalf("GetCIDR(%q) error = %v, wantErr %v", tt.clientIP, err, tt.wantErr)
}
if got != tt.want {
t.Errorf("GetCIDR(%q) = %q, want %q", tt.clientIP, got, tt.want)
}
})
}
}
// Test_GetCIDR_MostSpecific pins precedence: the narrowest range wins, so a captcha on
// a /24 is not overruled by a ban on its /8.
func Test_GetCIDR_MostSpecific(t *testing.T) {
client, _ := newCIDRClient(map[string]string{
"10.0.0.0/8": BannedValue,
"10.1.0.0/16": CaptchaValue,
"10.1.2.0/24": BannedValue,
"2001:db8::/32": BannedValue,
"2001:db8::/48": CaptchaValue,
})
tests := []struct {
clientIP string
want string
}{
{clientIP: "10.1.2.3", want: BannedValue},
{clientIP: "10.1.3.3", want: CaptchaValue},
{clientIP: "10.2.3.4", want: BannedValue},
{clientIP: "2001:db8::1", want: CaptchaValue},
{clientIP: "2001:db8:1::1", want: BannedValue},
}
for _, tt := range tests {
t.Run(tt.clientIP, func(t *testing.T) {
got, err := client.GetCIDR(tt.clientIP)
if err != nil {
t.Fatalf("GetCIDR(%q) unexpected error %v", tt.clientIP, err)
}
if got != tt.want {
t.Errorf("GetCIDR(%q) = %q, want %q", tt.clientIP, got, tt.want)
}
})
}
}
func Test_DeleteCIDR(t *testing.T) {
client, _ := newCIDRClient(map[string]string{
"10.0.0.0/8": BannedValue,
"10.1.2.0/24": CaptchaValue,
})
client.DeleteCIDR("10.1.2.0/24")
// The wider decision is untouched and takes over.
if got, err := client.GetCIDR("10.1.2.3"); err != nil || got != BannedValue {
t.Errorf("after deleting the /24, GetCIDR = %q %v, want %q", got, err, BannedValue)
}
client.DeleteCIDR("10.0.0.0/8")
if _, err := client.GetCIDR("10.1.2.3"); err == nil {
t.Error("GetCIDR should miss once every decision is deleted")
}
}
func Test_SetCIDR_InvalidIsNotStored(t *testing.T) {
for _, cidr := range []string{"", "garbage", "10.0.0.1", "10.0.0.0/33", "10.0.0.0/-1"} {
t.Run(cidr, func(t *testing.T) {
_, counting := newCIDRClient(map[string]string{cidr: BannedValue})
if len(counting.values) != 0 {
t.Errorf("SetCIDR(%q) stored %v, want nothing", cidr, counting.values)
}
})
}
}
// Test_GetCIDR_Reads guards the cost of the lookup: it must probe only the prefix
// lengths that have a decision, not every possible one.
func Test_GetCIDR_Reads(t *testing.T) {
tests := []struct {
name string
decisions map[string]string
clientIP string
wantReads int
}{
{name: "no decision at all, IPv4", decisions: nil, clientIP: "10.0.0.1", wantReads: 1},
{name: "no decision at all, IPv6", decisions: nil, clientIP: "2001:db8::1", wantReads: 1},
{
name: "one prefix length, hit",
decisions: map[string]string{"10.0.0.0/24": BannedValue},
clientIP: "10.0.0.1",
wantReads: 2,
},
{
name: "one prefix length, miss",
decisions: map[string]string{"10.0.0.0/24": BannedValue},
clientIP: "11.0.0.1",
wantReads: 2,
},
{
name: "three prefix lengths, miss probes each once",
decisions: map[string]string{"10.0.0.0/8": BannedValue, "10.1.0.0/16": BannedValue, "10.1.2.0/24": BannedValue},
clientIP: "11.0.0.1",
wantReads: 4,
},
{
name: "IPv6 client does not probe every length",
decisions: map[string]string{"2001:db8::/32": BannedValue},
clientIP: "2001:dead::1",
wantReads: 2,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
client, counting := newCIDRClient(tt.decisions)
counting.reads = 0
// Only the number of reads matters here, the result is covered above.
_, _ = client.GetCIDR(tt.clientIP)
if counting.reads != tt.wantReads {
t.Errorf("GetCIDR(%q) did %d cache reads, want %d", tt.clientIP, counting.reads, tt.wantReads)
}
})
}
}
+2
View File
@@ -96,6 +96,7 @@ type Config struct {
ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"` ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"`
RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"` RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"`
RedisCacheHost string `json:"redisCacheHost,omitempty"` RedisCacheHost string `json:"redisCacheHost,omitempty"`
RedisCacheReadHosts []string `json:"redisCacheReadHosts,omitempty"`
RedisCachePassword string `json:"redisCachePassword,omitempty"` RedisCachePassword string `json:"redisCachePassword,omitempty"`
RedisCachePasswordFile string `json:"redisCachePasswordFile,omitempty"` RedisCachePasswordFile string `json:"redisCachePasswordFile,omitempty"`
RedisCacheDatabase string `json:"redisCacheDatabase,omitempty"` RedisCacheDatabase string `json:"redisCacheDatabase,omitempty"`
@@ -172,6 +173,7 @@ func New() *Config {
ClientTrustedIPs: []string{}, ClientTrustedIPs: []string{},
RedisCacheEnabled: false, RedisCacheEnabled: false,
RedisCacheHost: "redis:6379", RedisCacheHost: "redis:6379",
RedisCacheReadHosts: []string{},
RedisCachePassword: "", RedisCachePassword: "",
RedisCacheDatabase: "", RedisCacheDatabase: "",
RedisCacheUnreachableBlock: true, RedisCacheUnreachableBlock: true,
+85
View File
@@ -0,0 +1,85 @@
package ip
import (
"net"
"strings"
)
const (
maxIPv4PrefixLen = 32
maxIPv6PrefixLen = 128
)
// CIDRKeys returns all possible CIDR prefixes of an IP, from the most specific (/32 for IPv4, /128 for IPv6) to the least specific (/0).
func CIDRKeys(ipStr string) []string {
parsed, maxBits := parseForPrefix(ipStr)
if parsed == nil {
return nil
}
keys := make([]string, 0, maxBits+1)
for bits := maxBits; bits >= 0; bits-- {
keys = append(keys, cidrKey(parsed, bits, maxBits))
}
return keys
}
// CIDRLookupKeys returns the keys of the CIDRs containing an IP for the given prefix lengths
// only, most specific first. Duplicates and lengths of the other family are skipped.
func CIDRLookupKeys(ipStr string, prefixLens []int) []string {
parsed, maxBits := parseForPrefix(ipStr)
if parsed == nil {
return nil
}
var wanted [maxIPv6PrefixLen + 1]bool
for _, bits := range prefixLens {
if bits >= 0 && bits <= maxBits {
wanted[bits] = true
}
}
keys := make([]string, 0, len(prefixLens))
for bits := maxBits; bits >= 0; bits-- {
if wanted[bits] {
keys = append(keys, cidrKey(parsed, bits, maxBits))
}
}
return keys
}
// NormalizeCIDR parses a CIDR string and returns its normalized form, or an empty string if invalid.
func NormalizeCIDR(cidrStr string) string {
_, ipNet, err := net.ParseCIDR(strings.TrimSpace(cidrStr))
if err != nil {
return ""
}
return ipNet.String()
}
// CIDRPrefixLen returns the prefix length of a CIDR, or -1 if it is not a valid CIDR.
func CIDRPrefixLen(cidrStr string) int {
_, ipNet, err := net.ParseCIDR(strings.TrimSpace(cidrStr))
if err != nil {
return -1
}
prefixLen, _ := ipNet.Mask.Size()
return prefixLen
}
// parseForPrefix returns the IP in the native form of its family, and that family's bit length.
func parseForPrefix(ipStr string) (net.IP, int) {
parsed := net.ParseIP(ipStr)
if parsed == nil {
return nil, 0
}
if parsed4 := parsed.To4(); parsed4 != nil {
return parsed4, maxIPv4PrefixLen
}
return parsed.To16(), maxIPv6PrefixLen
}
// cidrKey builds the key of the CIDR of bits length containing the IP.
// It formats through net.IPNet like NormalizeCIDR, so writes and lookups agree.
func cidrKey(parsed net.IP, bits, maxBits int) string {
mask := net.CIDRMask(bits, maxBits)
ipNet := net.IPNet{IP: parsed.Mask(mask), Mask: mask}
return ipNet.String()
}
+301
View File
@@ -0,0 +1,301 @@
package ip
import (
"strconv"
"strings"
"testing"
)
func TestCIDRKeys(t *testing.T) {
tests := []struct {
ip string
wantKeys int
checks map[int]string
}{
{
ip: "10.0.0.1",
wantKeys: 33,
checks: map[int]string{0: "10.0.0.1/32", 8: "10.0.0.0/24", 32: "0.0.0.0/0"},
},
{
ip: "2001:db8::1",
wantKeys: 129,
checks: map[int]string{0: "2001:db8::1/128", 32: "2001:db8::/96", 128: "::/0"},
},
}
for _, tt := range tests {
t.Run(tt.ip, func(t *testing.T) {
keys := CIDRKeys(tt.ip)
if keys == nil {
t.Fatal("CIDRKeys returned nil")
}
if len(keys) != tt.wantKeys {
t.Fatalf("expected %d keys, got %d", tt.wantKeys, len(keys))
}
for idx, want := range tt.checks {
if keys[idx] != want {
t.Errorf("keys[%d] should be %s, got %s", idx, want, keys[idx])
}
}
})
}
}
func TestCIDRKeys_MostToLeastSpecific(t *testing.T) {
ips := []string{"10.0.0.1", "2001:db8::1"}
for _, ip := range ips {
t.Run(ip, func(t *testing.T) {
keys := CIDRKeys(ip)
for i := 1; i < len(keys); i++ {
prevBits := strings.Split(keys[i-1], "/")[1]
curBits := strings.Split(keys[i], "/")[1]
prevN, _ := strconv.Atoi(prevBits)
curN, _ := strconv.Atoi(curBits)
if prevN <= curN {
t.Errorf("keys should go from most specific to least specific at index %d: /%d <= /%d", i, prevN, curN)
}
}
})
}
}
func TestCIDRKeys_IPVariants(t *testing.T) {
tests := []struct {
ip string
wantKeys int
}{
{"0.0.0.0", 33},
{"255.255.255.255", 33},
{"1.2.3.4", 33},
{"10.0.0.1", 33},
{"192.168.1.1", 33},
{"invalid", 0},
{"", 0},
{" ", 0},
}
for _, tt := range tests {
t.Run(tt.ip, func(t *testing.T) {
keys := CIDRKeys(tt.ip)
if len(keys) != tt.wantKeys {
t.Errorf("CIDRKeys(%q) returned %d keys, want %d", tt.ip, len(keys), tt.wantKeys)
}
})
}
}
func TestCIDRKeys_VerifyNetworkAddress(t *testing.T) {
keys := CIDRKeys("10.1.2.3")
tests := []struct {
bits int
want string
}{
{24, "10.1.2.0/24"},
{16, "10.1.0.0/16"},
{8, "10.0.0.0/8"},
}
for _, tt := range tests {
if got := keys[32-tt.bits]; got != tt.want {
t.Errorf("/%d network should be %s, got %s", tt.bits, tt.want, got)
}
}
}
func TestCIDRKeys_VerifyNetworkAddressIPv6(t *testing.T) {
keys := CIDRKeys("2001:db8:1:2:3:4:5:6")
tests := []struct {
bits int
want string
}{
{64, "2001:db8:1:2::/64"},
{48, "2001:db8:1::/48"},
{32, "2001:db8::/32"},
}
for _, tt := range tests {
if got := keys[128-tt.bits]; got != tt.want {
t.Errorf("/%d network should be %s, got %s", tt.bits, tt.want, got)
}
}
}
// TestCIDRKeys_MatchNormalizeCIDR covers the invariant range support rests on: the key
// written for a decision is a key looked up for the IPs it covers, and only those.
func TestCIDRKeys_MatchNormalizeCIDR(t *testing.T) {
tests := []struct {
cidr string
ip string
match bool
}{
{cidr: "10.0.0.0/8", ip: "10.1.2.3", match: true},
{cidr: "10.0.0.0/24", ip: "10.0.0.1", match: true},
{cidr: "10.0.0.0/24", ip: "10.0.1.1", match: false},
{cidr: "1.2.3.4/32", ip: "1.2.3.4", match: true},
{cidr: "1.2.3.4/32", ip: "1.2.3.5", match: false},
{cidr: "0.0.0.0/0", ip: "8.8.8.8", match: true},
// LAPI does not have to send a network address, the host bits are dropped.
{cidr: "10.0.0.5/24", ip: "10.0.0.9", match: true},
{cidr: " 192.168.1.0/24 ", ip: "192.168.1.42", match: true},
{cidr: "2001:db8::/32", ip: "2001:db8::1", match: true},
{cidr: "2001:db8::/32", ip: "2001:db9::1", match: false},
{cidr: "::/0", ip: "2001:db8::1", match: true},
// An IPv4 range and an IPv4 mapped client still have to meet.
{cidr: "10.0.0.0/8", ip: "::ffff:10.1.2.3", match: true},
{cidr: "::ffff:10.0.0.0/104", ip: "10.1.2.3", match: true},
// Families do not mix.
{cidr: "::/0", ip: "8.8.8.8", match: false},
{cidr: "2001:db8::/32", ip: "::ffff:10.0.0.1", match: false},
}
for _, tt := range tests {
t.Run(tt.cidr+"_"+tt.ip, func(t *testing.T) {
key := NormalizeCIDR(tt.cidr)
if key == "" {
t.Fatalf("NormalizeCIDR(%q) returned nothing", tt.cidr)
}
found := false
for _, candidate := range CIDRKeys(tt.ip) {
if candidate == key {
found = true
break
}
}
if found != tt.match {
t.Errorf("key %q of %q found in CIDRKeys(%q) = %v, want %v", key, tt.cidr, tt.ip, found, tt.match)
}
})
}
}
func TestCIDRLookupKeys(t *testing.T) {
tests := []struct {
name string
ip string
prefixLens []int
want []string
}{
{
name: "most specific first",
ip: "10.1.2.3",
prefixLens: []int{8, 32, 16},
want: []string{"10.1.2.3/32", "10.1.0.0/16", "10.0.0.0/8"},
},
{
name: "duplicates are dropped",
ip: "10.1.2.3",
prefixLens: []int{24, 24, 24},
want: []string{"10.1.2.0/24"},
},
{
name: "lengths of the other family are skipped",
ip: "10.1.2.3",
prefixLens: []int{48, 64, 24},
want: []string{"10.1.2.0/24"},
},
{
name: "out of range lengths are skipped",
ip: "10.1.2.3",
prefixLens: []int{-1, 33, 129, 8},
want: []string{"10.0.0.0/8"},
},
{
name: "ipv6 keeps its own lengths",
ip: "2001:db8::1",
prefixLens: []int{32, 64},
want: []string{"2001:db8::/64", "2001:db8::/32"},
},
{
name: "no length gives no key",
ip: "10.1.2.3",
prefixLens: []int{},
want: []string{},
},
{
name: "invalid ip gives no key",
ip: "invalid",
prefixLens: []int{24},
want: nil,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := CIDRLookupKeys(tt.ip, tt.prefixLens)
if len(got) != len(tt.want) {
t.Fatalf("CIDRLookupKeys(%q, %v) = %v, want %v", tt.ip, tt.prefixLens, got, tt.want)
}
for i := range got {
if got[i] != tt.want[i] {
t.Errorf("CIDRLookupKeys(%q, %v)[%d] = %q, want %q", tt.ip, tt.prefixLens, i, got[i], tt.want[i])
}
}
})
}
}
// TestCIDRLookupKeys_SubsetOfCIDRKeys: restricting the lengths only removes candidates,
// it never changes the key of a length that is kept.
func TestCIDRLookupKeys_SubsetOfCIDRKeys(t *testing.T) {
for _, ipStr := range []string{"10.1.2.3", "2001:db8::1", "::ffff:10.1.2.3"} {
t.Run(ipStr, func(t *testing.T) {
all := CIDRKeys(ipStr)
maxBits := len(all) - 1
for bits := 0; bits <= maxBits; bits++ {
got := CIDRLookupKeys(ipStr, []int{bits})
if len(got) != 1 {
t.Fatalf("CIDRLookupKeys(%q, [%d]) returned %d keys", ipStr, bits, len(got))
}
if want := all[maxBits-bits]; got[0] != want {
t.Errorf("CIDRLookupKeys(%q, [%d]) = %q, want %q", ipStr, bits, got[0], want)
}
}
})
}
}
func TestCIDRPrefixLen(t *testing.T) {
tests := []struct {
input string
want int
}{
{"10.0.0.0/8", 8},
{"10.0.0.0/32", 32},
{"0.0.0.0/0", 0},
{"10.0.0.5/24", 24},
{" 10.0.0.0/16 ", 16},
{"2001:db8::/32", 32},
{"2001:db8::/128", 128},
{"::/0", 0},
{"10.0.0.1", -1},
{"invalid", -1},
{"", -1},
}
for _, tt := range tests {
t.Run(tt.input, func(t *testing.T) {
if got := CIDRPrefixLen(tt.input); got != tt.want {
t.Errorf("CIDRPrefixLen(%q) = %d, want %d", tt.input, got, tt.want)
}
})
}
}
func TestNormalizeCIDR(t *testing.T) {
tests := []struct {
input string
want string
}{
{"10.0.0.0/8", "10.0.0.0/8"},
{"10.0.0.0/16", "10.0.0.0/16"},
{"192.168.1.0/24", "192.168.1.0/24"},
{"2001:db8::/32", "2001:db8::/32"},
{"0.0.0.0/0", "0.0.0.0/0"},
{"::/0", "::/0"},
{"invalid", ""},
{"", ""},
{" 10.0.0.0/8 ", "10.0.0.0/8"},
}
for _, tt := range tests {
t.Run(tt.input, func(t *testing.T) {
got := NormalizeCIDR(tt.input)
if got != tt.want {
t.Errorf("NormalizeCIDR(%q) = %q, want %q", tt.input, got, tt.want)
}
})
}
}
+24 -10
View File
@@ -14,12 +14,14 @@
set -euo pipefail set -euo pipefail
# Pinned to match the Docker suite (tests/e2e/scenarios/*/docker-compose.yml). # Pinned to match the Docker suite (tests/e2e/scenarios/*/docker-compose.yml).
TRAEFIK_VERSION="${TRAEFIK_VERSION:-v3.7.6}" TRAEFIK_VERSION="${TRAEFIK_VERSION:-v3.7.11}"
WEB_PORT="${WEB_PORT:-8000}" WEB_PORT="${WEB_PORT:-8000}"
LAPI_PORT="${LAPI_PORT:-8090}" LAPI_PORT="${LAPI_PORT:-8090}"
BACKEND_PORT="${BACKEND_PORT:-8091}" BACKEND_PORT="${BACKEND_PORT:-8091}"
APPSEC_PORT="${APPSEC_PORT:-8092}" APPSEC_PORT="${APPSEC_PORT:-8092}"
REDIS_PORT="${REDIS_PORT:-8093}"
REDIS_READ_PORT="${REDIS_READ_PORT:-8094}"
LAPI_KEY="${LAPI_KEY:-e2e-mock-key}" LAPI_KEY="${LAPI_KEY:-e2e-mock-key}"
MOCK_LIB_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" MOCK_LIB_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
@@ -76,11 +78,11 @@ ensure_mock() {
# Poll a URL until it returns the expected status code, or fail. # Poll a URL until it returns the expected status code, or fail.
# Usage: wait_for_status URL CODE [TIMEOUT_SECONDS] [curl args...] # Usage: wait_for_status URL CODE [TIMEOUT_SECONDS] [curl args...]
wait_for_status() { wait_for_status() {
local url="$1" expected="$2" timeout="${3:-30}" local url="$1" expected="$2" timeout="${3:-15}"
shift 3 || true shift 3 || true
local elapsed=0 got="" local elapsed=0 got=""
while (( elapsed < timeout )); do while (( elapsed < timeout )); do
got=$(curl -s -o /dev/null -w '%{http_code}' "$@" "$url" || true) got=$(curl -s -m 1 -o /dev/null -w '%{http_code}' "$@" "$url" || true)
if [[ "$got" == "$expected" ]]; then if [[ "$got" == "$expected" ]]; then
return 0 return 0
fi fi
@@ -96,11 +98,11 @@ wait_for_status() {
# code alone can't tell the states apart (e.g. captcha page vs backend, both 200). # code alone can't tell the states apart (e.g. captcha page vs backend, both 200).
# Usage: wait_for_body_contains URL NEEDLE [TIMEOUT_SECONDS] [curl args...] # Usage: wait_for_body_contains URL NEEDLE [TIMEOUT_SECONDS] [curl args...]
wait_for_body_contains() { wait_for_body_contains() {
local url="$1" needle="$2" timeout="${3:-30}" local url="$1" needle="$2" timeout="${3:-15}"
shift 3 || true shift 3 || true
local elapsed=0 body="" local elapsed=0 body=""
while (( elapsed < timeout )); do while (( elapsed < timeout )); do
body=$(curl -s "$@" "$url" || true) body=$(curl -s -m 1 "$@" "$url" || true)
if grep -q "$needle" <<<"$body"; then if grep -q "$needle" <<<"$body"; then
return 0 return 0
fi fi
@@ -117,7 +119,7 @@ assert_status() {
local url="$1" expected="$2" local url="$1" expected="$2"
shift 2 || true shift 2 || true
local got local got
got=$(curl -s -o /dev/null -w '%{http_code}' "$@" "$url") got=$(curl -s --connect-timeout 1 -m 5 -o /dev/null -w '%{http_code}' "$@" "$url")
if [[ "$got" != "$expected" ]]; then if [[ "$got" != "$expected" ]]; then
echo "assert_status: $url expected $expected, got $got" >&2 echo "assert_status: $url expected $expected, got $got" >&2
return 1 return 1
@@ -130,7 +132,7 @@ assert_header() {
local url="$1" header="$2" expected="$3" local url="$1" header="$2" expected="$3"
shift 3 || true shift 3 || true
local got local got
got=$(curl -s -D - -o /dev/null "$@" "$url" | tr -d '\r' \ got=$(curl -s --connect-timeout 1 -m 5 -D - -o /dev/null "$@" "$url" | tr -d '\r' \
| awk -v h="${header,,}" -F': ' 'tolower($1) == h { print $2; exit }') | awk -v h="${header,,}" -F': ' 'tolower($1) == h { print $2; exit }')
if [[ "$got" != "$expected" ]]; then if [[ "$got" != "$expected" ]]; then
echo "assert_header: $url header $header expected \"$expected\", got \"$got\"" >&2 echo "assert_header: $url header $header expected \"$expected\", got \"$got\"" >&2
@@ -144,7 +146,7 @@ assert_body_contains() {
local url="$1" needle="$2" local url="$1" needle="$2"
shift 2 || true shift 2 || true
local body local body
body=$(curl -s "$@" "$url") body=$(curl -s --connect-timeout 1 -m 5 "$@" "$url")
if ! grep -q "$needle" <<<"$body"; then if ! grep -q "$needle" <<<"$body"; then
echo "assert_body_contains: $url expected to contain \"$needle\", got:" >&2 echo "assert_body_contains: $url expected to contain \"$needle\", got:" >&2
echo "$body" >&2 echo "$body" >&2
@@ -156,12 +158,20 @@ assert_body_contains() {
lapi_add_decision() { lapi_add_decision() {
local ip="$1" type="${2:-ban}" duration="${3:-4h}" local ip="$1" type="${2:-ban}" duration="${3:-4h}"
curl -sS -X POST "http://127.0.0.1:${LAPI_PORT}/admin/decisions?ip=${ip}&type=${type}&duration=${duration}" >/dev/null curl -sS --connect-timeout 1 -m 5 -X POST "http://127.0.0.1:${LAPI_PORT}/admin/decisions?ip=${ip}&type=${type}&duration=${duration}" >/dev/null
} }
lapi_delete_decision() { lapi_delete_decision() {
local ip="$1" local ip="$1"
curl -sS -X DELETE "http://127.0.0.1:${LAPI_PORT}/admin/decisions?ip=${ip}" >/dev/null curl -sS --connect-timeout 1 -m 5 -X DELETE "http://127.0.0.1:${LAPI_PORT}/admin/decisions?ip=${ip}" >/dev/null
}
lapi_set_stream_fail() {
curl -sS --connect-timeout 1 -m 5 -X POST "http://127.0.0.1:${LAPI_PORT}/admin/stream-fail" >/dev/null
}
lapi_clear_stream_fail() {
curl -sS --connect-timeout 1 -m 5 -X DELETE "http://127.0.0.1:${LAPI_PORT}/admin/stream-fail" >/dev/null
} }
# --- stack lifecycle --------------------------------------------------------- # --- stack lifecycle ---------------------------------------------------------
@@ -187,6 +197,8 @@ start_stack() {
-e "s|@@LAPI_HOST@@|127.0.0.1:${LAPI_PORT}|g" \ -e "s|@@LAPI_HOST@@|127.0.0.1:${LAPI_PORT}|g" \
-e "s|@@APPSEC_HOST@@|127.0.0.1:${APPSEC_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|@@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" \ -e "s|@@SCENARIO_DIR@@|${scenario_dir}|g" \
"$scenario_dir/dynamic.yml" > "$WORKDIR/dynamic.yml" "$scenario_dir/dynamic.yml" > "$WORKDIR/dynamic.yml"
@@ -203,6 +215,8 @@ start_stack() {
--lapi-addr "127.0.0.1:${LAPI_PORT}" \ --lapi-addr "127.0.0.1:${LAPI_PORT}" \
--backend-addr "127.0.0.1:${BACKEND_PORT}" \ --backend-addr "127.0.0.1:${BACKEND_PORT}" \
--appsec-addr "127.0.0.1:${APPSEC_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_tls_args[@]}" >"$WORKDIR/mock.log" 2>&1 &
MOCK_PID=$! MOCK_PID=$!
+78 -3
View File
@@ -2,7 +2,8 @@
// suite. It answers only the few LAPI routes the plugin calls — live/none // 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 // 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 // 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 // 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 // CRS, virtual patching) is out of scope. The AppSec endpoint here emulates a
@@ -11,13 +12,16 @@
package main package main
import ( import (
"bufio"
"encoding/json" "encoding/json"
"flag" "flag"
"io" "io"
"log" "log"
"net"
"net/http" "net/http"
"strings" "strings"
"sync" "sync"
"sync/atomic"
) )
// Decision is the subset of a LAPI decision the plugin actually reads. // Decision is the subset of a LAPI decision the plugin actually reads.
@@ -31,6 +35,9 @@ var (
mu sync.Mutex mu sync.Mutex
active = map[string]Decision{} // ip -> decision currently in force active = map[string]Decision{} // ip -> decision currently in force
deleted = map[string]Decision{} // ip -> decision to report in the stream "deleted" list deleted = map[string]Decision{} // ip -> decision to report in the stream "deleted" list
// streamFail makes /v1/decisions/stream return 500 when set, to exercise the
// bouncer's fail-closed behaviour on consecutive stream poll failures.
streamFail atomic.Bool
) )
func writeJSON(w http.ResponseWriter, v any) { func writeJSON(w http.ResponseWriter, v any) {
@@ -46,6 +53,50 @@ func list(m map[string]Decision) []Decision {
return out 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() { func main() {
lapiAddr := flag.String("lapi-addr", "127.0.0.1:8090", "address for the LAPI mock") 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 // The stub upstream Traefik proxies allowed requests to — the binary-suite
@@ -53,6 +104,12 @@ func main() {
backendAddr := flag.String("backend-addr", "127.0.0.1:8091", "address for the stub upstream service") 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. // 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") 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 // 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 // (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. // bouncer's system-trust-store path. Backend and AppSec stay plaintext.
@@ -96,6 +153,9 @@ func main() {
}))) })))
}() }()
go serveRedis(*redisAddr, false)
go serveRedis(*redisReadAddr, true)
mux := http.NewServeMux() mux := http.NewServeMux()
// Readiness probe for the test harness (empty body, 200). // Readiness probe for the test harness (empty body, 200).
@@ -117,6 +177,10 @@ func main() {
// "deleted". Re-sending the same on every poll is harmless — the plugin just // "deleted". Re-sending the same on every poll is harmless — the plugin just
// re-adds to / re-deletes from its cache. // re-adds to / re-deletes from its cache.
mux.HandleFunc("/v1/decisions/stream", func(w http.ResponseWriter, _ *http.Request) { mux.HandleFunc("/v1/decisions/stream", func(w http.ResponseWriter, _ *http.Request) {
if streamFail.Load() {
w.WriteHeader(http.StatusInternalServerError)
return
}
mu.Lock() mu.Lock()
defer mu.Unlock() defer mu.Unlock()
writeJSON(w, map[string][]Decision{"new": list(active), "deleted": list(deleted)}) writeJSON(w, map[string][]Decision{"new": list(active), "deleted": list(deleted)})
@@ -127,6 +191,17 @@ func main() {
w.WriteHeader(http.StatusCreated) w.WriteHeader(http.StatusCreated)
}) })
// Test control plane: make the stream endpoint fail (POST) or recover (DELETE).
mux.HandleFunc("/admin/stream-fail", func(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodPost:
streamFail.Store(true)
case http.MethodDelete:
streamFail.Store(false)
}
w.WriteHeader(http.StatusOK)
})
// Test control plane: add / remove decisions instead of cscli. // Test control plane: add / remove decisions instead of cscli.
mux.HandleFunc("/admin/decisions", func(_ http.ResponseWriter, r *http.Request) { mux.HandleFunc("/admin/decisions", func(_ http.ResponseWriter, r *http.Request) {
q := r.URL.Query() q := r.URL.Query()
@@ -154,9 +229,9 @@ func main() {
}) })
if *lapiTLSCert != "" && *lapiTLSKey != "" { 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.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)) log.Fatal(http.ListenAndServe(*lapiAddr, mux))
} }
+11 -3
View File
@@ -11,17 +11,25 @@ body() {
echo "[$SCENARIO] no decision -> request passes (LAPI queried per request)" echo "[$SCENARIO] no decision -> request passes (LAPI queried per request)"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4" assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4"
echo "[$SCENARIO] adding ban decision for 1.2.3.4" echo "[$SCENARIO] adding ban decision for 1.2.3.4 and 2001:db8::1"
lapi_add_decision 1.2.3.4 ban 5m lapi_add_decision 1.2.3.4 ban 5m
lapi_add_decision "2001:db8::1" ban 5m
echo "[$SCENARIO] none mode has no cache -> next request must be blocked immediately" echo "[$SCENARIO] IP banned must be blocked (HTTP 403)"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.4" assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.4"
echo "[$SCENARIO] IPv6 banned must be blocked (HTTP 403)"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 2001:db8::1"
echo "[$SCENARIO] deleting decision" echo "[$SCENARIO] deleting decision"
lapi_delete_decision 1.2.3.4 lapi_delete_decision 1.2.3.4
lapi_delete_decision "2001:db8::1"
echo "[$SCENARIO] previously banned IP must pass again immediately" echo "[$SCENARIO] previously banned IP must pass again"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4" assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4"
echo "[$SCENARIO] previously banned IPv6 must pass again"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 2001:db8::1"
} }
run_scenario "$SCENARIO" "$HERE" body run_scenario "$SCENARIO" "$HERE" body
@@ -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
@@ -18,7 +18,8 @@ http:
bouncer: bouncer:
enabled: "true" enabled: "true"
crowdsecMode: stream crowdsecMode: stream
updateIntervalSeconds: "2" updateIntervalSeconds: "1"
updateMaxFailure: "2"
crowdsecLapiScheme: http crowdsecLapiScheme: http
crowdsecLapiHost: "@@LAPI_HOST@@" crowdsecLapiHost: "@@LAPI_HOST@@"
crowdsecLapiKey: "@@APIKEY@@" crowdsecLapiKey: "@@APIKEY@@"
+22 -2
View File
@@ -11,8 +11,9 @@ body() {
echo "[$SCENARIO] no decision yet -> request allowed" echo "[$SCENARIO] no decision yet -> request allowed"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4" assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4"
echo "[$SCENARIO] adding ban decision for 1.2.3.4" echo "[$SCENARIO] adding ban decision for 1.2.3.4 and 10.0.0.0/8"
lapi_add_decision 1.2.3.4 ban 5m lapi_add_decision 1.2.3.4 ban 5m
lapi_add_decision 10.0.0.0/24 ban 5m
echo "[$SCENARIO] banned IP must be blocked once the next stream poll lands (HTTP 403)" echo "[$SCENARIO] banned IP must be blocked once the next stream poll lands (HTTP 403)"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 403 15 -H "X-Forwarded-For: 1.2.3.4" wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 403 15 -H "X-Forwarded-For: 1.2.3.4"
@@ -20,11 +21,30 @@ body() {
echo "[$SCENARIO] non-banned IP must still pass (HTTP 200)" echo "[$SCENARIO] non-banned IP must still pass (HTTP 200)"
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 5.6.7.8" assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 5.6.7.8"
echo "[$SCENARIO] deleting ban decision" echo "[$SCENARIO] banned IP in CIDR must be blocked once polled (HTTP 403)"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 403 15 -H "X-Forwarded-For: 10.0.0.1"
echo "[$SCENARIO] deleting ban decision for 1.2.3.4 and 10.0.0.0/8"
lapi_delete_decision 1.2.3.4 lapi_delete_decision 1.2.3.4
lapi_delete_decision 10.0.0.0/24
echo "[$SCENARIO] previously banned IP must pass again once the deletion is polled" echo "[$SCENARIO] previously banned IP must pass again once the deletion is polled"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 1.2.3.4" wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 1.2.3.4"
echo "[$SCENARIO] previously CIDR-banned IP must pass again once deletion is polled"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 10.0.0.1"
echo "[$SCENARIO] making the stream endpoint fail -> bouncer must pass for one more cycle (updateMaxFailure: 2)"
lapi_set_stream_fail
sleep 2 # update cache is every 1 seconds then waiting for minimum 1 cycle
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 8.8.8.8"
echo "[$SCENARIO] bouncer must block everything (isStreamHealthy: false)"
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 403 15 -H "X-Forwarded-For: 8.8.8.8"
echo "[$SCENARIO] restoring the stream endpoint -> bouncer must recover and pass again"
lapi_clear_stream_fail
wait_for_status "http://127.0.0.1:${WEB_PORT}/foo" 200 15 -H "X-Forwarded-For: 8.8.8.8"
} }
run_scenario "$SCENARIO" "$HERE" body run_scenario "$SCENARIO" "$HERE" body
+3 -2
View File
@@ -1,4 +1,5 @@
package crowdsec_bouncer_traefik_plugin //nolint:revive,stylecheck package crowdsec_bouncer_traefik_plugin //nolint:revive,stylecheck
// pluginVersion is updated automatically by the release workflow. // pluginVersion is what the plugin reports to the Crowdsec LAPI.
var pluginVersion = "1.6.X" //nolint:gochecknoglobals // Do not edit by hand: the "Release (1/2) Prepare" workflow bumps it.
var pluginVersion = "v1.7.1" //nolint:gochecknoglobals