mirror of
https://github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin.git
synced 2026-09-02 20:28:50 +02:00
Compare commits
17
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0f8608d770 | ||
|
|
ae7481caa5 | ||
|
|
955391671c | ||
|
|
9b8d6b937c | ||
|
|
d57ead2ec7 | ||
|
|
1ba556e918 | ||
|
|
6b0518859d | ||
|
|
bef5dfaadb | ||
|
|
99cf9712f4 | ||
|
|
ed4a9e8262 | ||
|
|
f6ef95cf38 | ||
|
|
9daba9739c | ||
|
|
e98b8ed5ba | ||
|
|
bb44aef718 | ||
|
|
31874263f2 | ||
|
|
1c98c70f14 | ||
|
|
be13c49144 |
@@ -19,7 +19,7 @@ jobs:
|
||||
steps:
|
||||
- uses: actions/checkout@v7
|
||||
- name: Set up Go
|
||||
uses: actions/setup-go@v6
|
||||
uses: actions/setup-go@v7
|
||||
with:
|
||||
# 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.
|
||||
|
||||
@@ -32,7 +32,7 @@ jobs:
|
||||
|
||||
# https://github.com/marketplace/actions/setup-go-environment
|
||||
- name: Set up Go ${{ env.GO_VERSION }}
|
||||
uses: actions/setup-go@v6
|
||||
uses: actions/setup-go@v7
|
||||
with:
|
||||
go-version: ${{ env.GO_VERSION }}
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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 }}"
|
||||
@@ -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:
|
||||
@@ -28,7 +28,7 @@ jobs:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Run Renovate
|
||||
uses: renovatebot/github-action@v46.1.14
|
||||
uses: renovatebot/github-action@v46.2.2
|
||||
with:
|
||||
token: ${{ secrets.RENOVATE_TOKEN }}
|
||||
env:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -384,8 +384,8 @@ make run
|
||||
- Transmit only the first number of bytes to Crowdsec Appsec Server.
|
||||
- CrowdsecAppsecUnreadableBodyBlock
|
||||
- bool
|
||||
- default: false
|
||||
- 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.
|
||||
- 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` 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
|
||||
- string
|
||||
- default: value of `CrowdsecLapiKey`
|
||||
@@ -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
|
||||
|
||||
+38
-5
@@ -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,
|
||||
)
|
||||
@@ -329,7 +330,7 @@ func New(_ context.Context, next http.Handler, config *configuration.Config, nam
|
||||
|
||||
// ServeHTTP principal function of plugin.
|
||||
//
|
||||
//nolint:nestif
|
||||
//nolint:nestif,gocognit
|
||||
func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
|
||||
if !bouncer.enabled {
|
||||
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.
|
||||
if bouncer.crowdsecMode == configuration.StreamMode || bouncer.crowdsecMode == configuration.AloneMode {
|
||||
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)
|
||||
} else {
|
||||
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 {
|
||||
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{
|
||||
Scheme: bouncer.crowdsecScheme,
|
||||
Host: bouncer.crowdsecHost,
|
||||
@@ -668,12 +684,20 @@ func handleStreamCache(bouncer *Bouncer) error {
|
||||
default:
|
||||
bouncer.log.Info("handleStreamCache:unknownType " + decision.Type)
|
||||
}
|
||||
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 {
|
||||
if strings.Contains(decision.Value, "/") {
|
||||
bouncer.cacheClient.DeleteCIDR(decision.Value)
|
||||
} else {
|
||||
bouncer.cacheClient.Delete(decision.Value)
|
||||
}
|
||||
}
|
||||
bouncer.log.Debug("handleStreamCache:updated")
|
||||
isCrowdsecStreamStartup = false
|
||||
return nil
|
||||
@@ -732,7 +756,17 @@ func crowdsecQuery(bouncer *Bouncer, stringURL string, data []byte) ([]byte, err
|
||||
// a 403. This mirrors the reference lua-cs-bouncer behavior, which refuses to
|
||||
// read the body of an HTTP/2+ request that has no Content-Length.
|
||||
func isBodyUnreadable(httpReq *http.Request) bool {
|
||||
return httpReq.Body != nil && httpReq.ProtoMajor >= 2 && httpReq.ContentLength < 0
|
||||
return httpReq.Body != nil && httpReq.Body != http.NoBody && httpReq.ProtoMajor >= 2 && httpReq.ContentLength < 0
|
||||
}
|
||||
|
||||
// isMethodWithBody used only when isBodyUnreadable returns true but the request method can't have body.
|
||||
func isMethodWithBody(method string) bool {
|
||||
switch method {
|
||||
case http.MethodPost, http.MethodPut, http.MethodPatch, http.MethodDelete:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func appsecQuery(bouncer *Bouncer, ip string, httpReq *http.Request) error {
|
||||
@@ -744,8 +778,7 @@ func appsecQuery(bouncer *Bouncer, ip string, httpReq *http.Request) error {
|
||||
var req *http.Request
|
||||
switch {
|
||||
case isBodyUnreadable(httpReq):
|
||||
if bouncer.appsecUnreadableBodyBlock {
|
||||
// The caller (handleNextServeHTTP) logs this returned error with the IP.
|
||||
if bouncer.appsecUnreadableBodyBlock && isMethodWithBody(httpReq.Method) {
|
||||
return errors.New("appsecQuery:unreadableBody dropped")
|
||||
}
|
||||
req, _ = http.NewRequest(http.MethodGet, routeURL.String(), nil)
|
||||
|
||||
+53
-11
@@ -7,6 +7,7 @@ import (
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
"text/template"
|
||||
"time"
|
||||
@@ -396,29 +397,27 @@ func (b blockingBody) Read(_ []byte) (int, error) {
|
||||
func (blockingBody) Close() error { return nil }
|
||||
|
||||
func Test_isBodyUnreadable(t *testing.T) {
|
||||
realBody := func() io.ReadCloser { return io.NopCloser(strings.NewReader("data")) }
|
||||
tests := []struct {
|
||||
name string
|
||||
protoMajor int
|
||||
contentLength int64
|
||||
hasBody bool
|
||||
body io.ReadCloser
|
||||
want bool
|
||||
}{
|
||||
{name: "http2 grpc stream without content-length", protoMajor: 2, contentLength: -1, hasBody: true, want: true},
|
||||
{name: "http3 stream without content-length", protoMajor: 3, contentLength: -1, hasBody: true, want: true},
|
||||
{name: "http2 with content-length", protoMajor: 2, contentLength: 42, hasBody: true, want: false},
|
||||
{name: "http1.1 chunked without content-length", protoMajor: 1, contentLength: -1, hasBody: true, want: false},
|
||||
{name: "http2 without body", protoMajor: 2, contentLength: -1, hasBody: false, want: false},
|
||||
{name: "http2 grpc stream without content-length", protoMajor: 2, contentLength: -1, body: realBody(), want: true},
|
||||
{name: "http3 stream without content-length", protoMajor: 3, contentLength: -1, body: realBody(), want: true},
|
||||
{name: "http2 with content-length", protoMajor: 2, contentLength: 42, body: realBody(), want: false},
|
||||
{name: "http1.1 chunked without content-length", protoMajor: 1, contentLength: -1, body: realBody(), want: false},
|
||||
{name: "http2 without body", protoMajor: 2, contentLength: -1, body: nil, want: false},
|
||||
{name: "http2 with http.NoBody", protoMajor: 2, contentLength: -1, body: http.NoBody, want: false},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
req, _ := http.NewRequest(http.MethodPost, "http://localhost", nil)
|
||||
req.ProtoMajor = tt.protoMajor
|
||||
req.ContentLength = tt.contentLength
|
||||
if tt.hasBody {
|
||||
req.Body = http.NoBody
|
||||
} else {
|
||||
req.Body = nil
|
||||
}
|
||||
req.Body = tt.body
|
||||
if got := isBodyUnreadable(req); got != tt.want {
|
||||
t.Errorf("isBodyUnreadable() = %v, want %v", got, tt.want)
|
||||
}
|
||||
@@ -513,3 +512,46 @@ func Test_appsecQuery_dropUnreadableBody(t *testing.T) {
|
||||
t.Fatal("appsecQuery() blocked on a streaming request body (issue #323 regression)")
|
||||
}
|
||||
}
|
||||
|
||||
func newUnreadableGetRequest(done <-chan struct{}) *http.Request {
|
||||
req, _ := http.NewRequest(http.MethodGet, "http://localhost/", blockingBody{done: done})
|
||||
req.ProtoMajor = 3
|
||||
req.ContentLength = -1
|
||||
return req
|
||||
}
|
||||
|
||||
// Test_appsecQuery_unreadableBodyGetNotDropped is a regression test for issue #351
|
||||
func Test_appsecQuery_unreadableBodyGetNotDropped(t *testing.T) {
|
||||
appsecServer := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, _ *http.Request) {
|
||||
rw.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer appsecServer.Close()
|
||||
|
||||
appsecURL, _ := url.Parse(appsecServer.URL)
|
||||
bouncer := &Bouncer{
|
||||
appsecScheme: appsecURL.Scheme,
|
||||
appsecHost: appsecURL.Host,
|
||||
appsecPath: "/",
|
||||
appsecBodyLimit: 10485760,
|
||||
appsecUnreadableBodyBlock: true,
|
||||
httpAppsecClient: appsecServer.Client(),
|
||||
log: logger.New("INFO", ""),
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
defer close(done)
|
||||
|
||||
finished := make(chan error, 1)
|
||||
go func() {
|
||||
finished <- appsecQuery(bouncer, "1.2.3.4", newUnreadableGetRequest(done))
|
||||
}()
|
||||
|
||||
select {
|
||||
case err := <-finished:
|
||||
if err != nil {
|
||||
t.Errorf("appsecQuery() on an HTTP/3 GET without content-length returned error: %v", err)
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("appsecQuery() blocked on an HTTP/3 GET request body (issue #351 regression)")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.5.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -80,7 +80,7 @@ services:
|
||||
- "traefik.http.routers.router-bar3.entrypoints=web"
|
||||
- "traefik.http.routers.router-bar3.middlewares=crowdsec2@docker"
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
+3
-3
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -12,7 +12,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
volumes:
|
||||
- "/var/run/docker.sock:/var/run/docker.sock:ro"
|
||||
# - './ban.html:/ban.html:ro'
|
||||
@@ -59,7 +59,7 @@ services:
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.5.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.5.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -47,7 +47,7 @@ services:
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecappsechost=crowdsec:7422"
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
cloudflare:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "cloudflare"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -19,7 +19,7 @@ services:
|
||||
- 8080:8080
|
||||
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -33,7 +33,7 @@ services:
|
||||
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
- logs-traefik:/var/log/traefik
|
||||
@@ -79,7 +79,7 @@ services:
|
||||
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -55,7 +55,7 @@ services:
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.captchaHTMLFilePath=/captcha.html"
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -46,7 +46,7 @@ services:
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.banFilePath=/ban.html"
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.5.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -14,7 +14,7 @@ services:
|
||||
- "--entrypoints.web.forwardedheaders.trustedips=172.18.0.0/24"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.4.5"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -59,7 +59,7 @@ services:
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.captchaHTMLFilePath=/captcha.html"
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
image:
|
||||
tag: v1.6.1-2
|
||||
tag: v1.7.8-2
|
||||
|
||||
agent:
|
||||
acquisition:
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
image:
|
||||
tag: v3.0.0
|
||||
tag: v3.7.11
|
||||
|
||||
logs:
|
||||
general:
|
||||
@@ -15,4 +15,4 @@ experimental:
|
||||
plugins:
|
||||
bouncer:
|
||||
moduleName: "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
version: "v1.3.0"
|
||||
version: "v1.7.1"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -71,7 +71,7 @@ services:
|
||||
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
@@ -87,7 +87,7 @@ services:
|
||||
- "traefik.enable=false"
|
||||
|
||||
redis-secure:
|
||||
image: "redis:8.8.0-alpine"
|
||||
image: "redis:8.10.0-alpine"
|
||||
container_name: "redis-secure"
|
||||
hostname: redis-secure
|
||||
restart: unless-stopped
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.5.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.5.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -71,7 +71,7 @@ services:
|
||||
# Define AppSec host and port informations
|
||||
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecappsechost=crowdsec:7422"
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
services:
|
||||
traefik:
|
||||
image: "traefik:v3.0.0"
|
||||
image: "traefik:v3.7.11"
|
||||
container_name: "traefik"
|
||||
restart: unless-stopped
|
||||
command:
|
||||
@@ -13,7 +13,7 @@ services:
|
||||
- "--entrypoints.web.address=:80"
|
||||
|
||||
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
- "--experimental.plugins.bouncer.version=v1.3.0"
|
||||
- "--experimental.plugins.bouncer.version=v1.7.1"
|
||||
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
@@ -65,7 +65,7 @@ services:
|
||||
|
||||
|
||||
crowdsec:
|
||||
image: crowdsecurity/crowdsec:v1.6.8
|
||||
image: crowdsecurity/crowdsec:v1.7.8
|
||||
container_name: "crowdsec"
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
module github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
|
||||
|
||||
go 1.22
|
||||
go 1.22.12
|
||||
|
||||
require (
|
||||
github.com/leprosus/golang-ttl-map v1.1.7
|
||||
|
||||
Vendored
+128
-22
@@ -6,9 +6,14 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
|
||||
ttl_map "github.com/leprosus/golang-ttl-map"
|
||||
simpleredis "github.com/maxlerebourg/simpleredis"
|
||||
|
||||
"github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -24,13 +29,21 @@ const (
|
||||
CacheMiss = "cache:miss"
|
||||
// CacheUnreachable error string when cache is 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
|
||||
var (
|
||||
redis simpleredis.SimpleRedis
|
||||
cache = ttl_map.New()
|
||||
)
|
||||
var cache = ttl_map.New()
|
||||
|
||||
type localCache struct{}
|
||||
|
||||
@@ -53,32 +66,47 @@ func (localCache) delete(key string) {
|
||||
|
||||
type redisCache struct {
|
||||
log *slog.Logger
|
||||
writer simpleredis.SimpleRedis
|
||||
readers []simpleredis.SimpleRedis
|
||||
counter atomic.Uint64
|
||||
}
|
||||
|
||||
func (redisCache) get(key string) (string, error) {
|
||||
value, err := redis.Get(key)
|
||||
valueString := string(value)
|
||||
if err == nil && len(valueString) > 0 {
|
||||
return valueString, nil
|
||||
func (rc *redisCache) nextReader() *simpleredis.SimpleRedis {
|
||||
n := len(rc.readers)
|
||||
if n == 0 {
|
||||
return &rc.writer
|
||||
}
|
||||
errRedisMessage := err.Error()
|
||||
if errRedisMessage == simpleredis.RedisMiss {
|
||||
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)
|
||||
}
|
||||
if errRedisMessage == simpleredis.RedisUnreachable {
|
||||
case simpleredis.RedisUnreachable:
|
||||
return "", errors.New(CacheUnreachable)
|
||||
}
|
||||
default:
|
||||
return "", err
|
||||
}
|
||||
}
|
||||
valueString := string(value)
|
||||
if len(valueString) > 0 {
|
||||
return valueString, nil
|
||||
}
|
||||
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 +124,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.
|
||||
@@ -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.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
|
||||
}
|
||||
|
||||
Vendored
+222
@@ -3,9 +3,11 @@
|
||||
package cache
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
|
||||
simpleredis "github.com/maxlerebourg/simpleredis"
|
||||
)
|
||||
|
||||
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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -14,12 +14,14 @@
|
||||
set -euo pipefail
|
||||
|
||||
# Pinned to match the Docker suite (tests/e2e/scenarios/*/docker-compose.yml).
|
||||
TRAEFIK_VERSION="${TRAEFIK_VERSION:-v3.7.1}"
|
||||
TRAEFIK_VERSION="${TRAEFIK_VERSION:-v3.7.11}"
|
||||
|
||||
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)"
|
||||
@@ -76,11 +78,11 @@ ensure_mock() {
|
||||
# Poll a URL until it returns the expected status code, or fail.
|
||||
# Usage: wait_for_status URL CODE [TIMEOUT_SECONDS] [curl args...]
|
||||
wait_for_status() {
|
||||
local url="$1" expected="$2" timeout="${3:-30}"
|
||||
local url="$1" expected="$2" timeout="${3:-15}"
|
||||
shift 3 || true
|
||||
local elapsed=0 got=""
|
||||
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
|
||||
return 0
|
||||
fi
|
||||
@@ -96,11 +98,11 @@ wait_for_status() {
|
||||
# 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...]
|
||||
wait_for_body_contains() {
|
||||
local url="$1" needle="$2" timeout="${3:-30}"
|
||||
local url="$1" needle="$2" timeout="${3:-15}"
|
||||
shift 3 || true
|
||||
local elapsed=0 body=""
|
||||
while (( elapsed < timeout )); do
|
||||
body=$(curl -s "$@" "$url" || true)
|
||||
body=$(curl -s -m 1 "$@" "$url" || true)
|
||||
if grep -q "$needle" <<<"$body"; then
|
||||
return 0
|
||||
fi
|
||||
@@ -117,7 +119,7 @@ assert_status() {
|
||||
local url="$1" expected="$2"
|
||||
shift 2 || true
|
||||
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
|
||||
echo "assert_status: $url expected $expected, got $got" >&2
|
||||
return 1
|
||||
@@ -130,7 +132,7 @@ assert_header() {
|
||||
local url="$1" header="$2" expected="$3"
|
||||
shift 3 || true
|
||||
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 }')
|
||||
if [[ "$got" != "$expected" ]]; then
|
||||
echo "assert_header: $url header $header expected \"$expected\", got \"$got\"" >&2
|
||||
@@ -144,7 +146,7 @@ assert_body_contains() {
|
||||
local url="$1" needle="$2"
|
||||
shift 2 || true
|
||||
local body
|
||||
body=$(curl -s "$@" "$url")
|
||||
body=$(curl -s --connect-timeout 1 -m 5 "$@" "$url")
|
||||
if ! grep -q "$needle" <<<"$body"; then
|
||||
echo "assert_body_contains: $url expected to contain \"$needle\", got:" >&2
|
||||
echo "$body" >&2
|
||||
@@ -156,12 +158,20 @@ assert_body_contains() {
|
||||
|
||||
lapi_add_decision() {
|
||||
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() {
|
||||
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 ---------------------------------------------------------
|
||||
@@ -187,6 +197,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 +215,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=$!
|
||||
|
||||
|
||||
@@ -3,4 +3,4 @@
|
||||
// golangci-lint and `go mod vendor`. Stdlib only — no dependencies.
|
||||
module mocklapi
|
||||
|
||||
go 1.22
|
||||
go 1.22.12
|
||||
|
||||
@@ -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,13 +12,16 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"flag"
|
||||
"io"
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
)
|
||||
|
||||
// Decision is the subset of a LAPI decision the plugin actually reads.
|
||||
@@ -31,6 +35,9 @@ var (
|
||||
mu sync.Mutex
|
||||
active = map[string]Decision{} // ip -> decision currently in force
|
||||
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) {
|
||||
@@ -46,6 +53,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 +104,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 +153,9 @@ func main() {
|
||||
})))
|
||||
}()
|
||||
|
||||
go serveRedis(*redisAddr, false)
|
||||
go serveRedis(*redisReadAddr, true)
|
||||
|
||||
mux := http.NewServeMux()
|
||||
|
||||
// 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
|
||||
// re-adds to / re-deletes from its cache.
|
||||
mux.HandleFunc("/v1/decisions/stream", func(w http.ResponseWriter, _ *http.Request) {
|
||||
if streamFail.Load() {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
writeJSON(w, map[string][]Decision{"new": list(active), "deleted": list(deleted)})
|
||||
@@ -127,6 +191,17 @@ func main() {
|
||||
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.
|
||||
mux.HandleFunc("/admin/decisions", func(_ http.ResponseWriter, r *http.Request) {
|
||||
q := r.URL.Query()
|
||||
@@ -154,9 +229,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))
|
||||
}
|
||||
|
||||
@@ -18,7 +18,7 @@ http:
|
||||
bouncer:
|
||||
enabled: "true"
|
||||
# IP bouncing disabled — this scenario exercises AppSec only.
|
||||
crowdsecMode: none
|
||||
crowdsecMode: appsec
|
||||
crowdsecLapiScheme: http
|
||||
crowdsecLapiHost: "@@LAPI_HOST@@"
|
||||
crowdsecLapiKey: "@@APIKEY@@"
|
||||
|
||||
@@ -30,6 +30,15 @@ body() {
|
||||
|
||||
echo "[$SCENARIO] request that send bad body before crowdsecAppsecBodyLimit must pass (AppSec 403)"
|
||||
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.4" -X POST -d "a=0&______"
|
||||
|
||||
echo "[$SCENARIO] request http2 that send no body GET (AppSec 200)"
|
||||
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 200 -H "X-Forwarded-For: 1.2.3.4" --http2-prior-knowledge -H "Content-Length:"
|
||||
|
||||
echo "[$SCENARIO] request http2 that send unreadable body GET (AppSec 403)"
|
||||
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.4" --http2-prior-knowledge -H "Content-Length:" -d "test"
|
||||
|
||||
echo "[$SCENARIO] request http2 that send unreadable body POST (AppSec 403)"
|
||||
assert_status "http://127.0.0.1:${WEB_PORT}/foo" 403 -H "X-Forwarded-For: 1.2.3.4" --http2-prior-knowledge -H "Content-Length:" -X POST -d "test"
|
||||
}
|
||||
|
||||
run_scenario "$SCENARIO" "$HERE" body
|
||||
|
||||
@@ -11,17 +11,25 @@ body() {
|
||||
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"
|
||||
|
||||
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 "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"
|
||||
|
||||
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"
|
||||
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"
|
||||
|
||||
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
|
||||
|
||||
@@ -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"
|
||||
@@ -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:
|
||||
enabled: "true"
|
||||
crowdsecMode: stream
|
||||
updateIntervalSeconds: "2"
|
||||
updateIntervalSeconds: "1"
|
||||
updateMaxFailure: "2"
|
||||
crowdsecLapiScheme: http
|
||||
crowdsecLapiHost: "@@LAPI_HOST@@"
|
||||
crowdsecLapiKey: "@@APIKEY@@"
|
||||
|
||||
@@ -11,8 +11,9 @@ body() {
|
||||
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"
|
||||
|
||||
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 10.0.0.0/24 ban 5m
|
||||
|
||||
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"
|
||||
@@ -20,11 +21,30 @@ body() {
|
||||
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"
|
||||
|
||||
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 10.0.0.0/24
|
||||
|
||||
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"
|
||||
|
||||
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
|
||||
|
||||
+3
-2
@@ -1,4 +1,5 @@
|
||||
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 what the plugin reports to the Crowdsec LAPI.
|
||||
// Do not edit by hand: the "Release (1/2) Prepare" workflow bumps it.
|
||||
var pluginVersion = "v1.7.1" //nolint:gochecknoglobals
|
||||
|
||||
Reference in New Issue
Block a user