CrowdSec decision list fetched, kept, and its clients banned until the decision ends (closes #106)
check / check (push) Waiting to run
check / check (push) Waiting to run
SWWAF_CROWDSEC_LAPI_URL and SWWAF_CROWDSEC_LAPI_KEY name an engine whose decision list, <url>/v1/decisions, is fetched every minute with the key in X-Api-Key and kept as a blocklist is: used while a fetch fails, and across restarts through reputation.json. Ban decisions on an Ip or a Range end at the fetch time plus their duration. A listed client's request is refused and bans its netblock with the cause crowdsec until the decision ends; bans.json, ban notes and metrics take the cause. Judgement call: fetched every minute, not a setting. Judgement call: a crowdsec ban never lengthens a limit ban. Judgement call: a lifted crowdsec ban is remade while its decision lasts. Model: opus-5-5
This commit is contained in:
@@ -0,0 +1,428 @@
|
||||
package reputation_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/netip"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/synctest"
|
||||
"time"
|
||||
|
||||
"sneak.berlin/go/smallwebwaf/internal/alerts"
|
||||
"sneak.berlin/go/smallwebwaf/internal/reputation"
|
||||
)
|
||||
|
||||
// The tests run in a synctest bubble, as those of the blocklists do, and
|
||||
// fetch the decision list from engine, a stand-in for a CrowdSec engine
|
||||
// that answers without the network.
|
||||
|
||||
const (
|
||||
// decisionsURL is the decision list of the tests' engine, and engineKey
|
||||
// the key it answers.
|
||||
decisionsURL = "http://crowdsec.example:8080/v1/decisions"
|
||||
engineKey = "crowdsec-key-0123456789abcdef"
|
||||
// sshBF and probing are scenarios of the engine's decisions.
|
||||
sshBF = "crowdsecurity/ssh-bf"
|
||||
probing = "crowdsecurity/http-probing"
|
||||
// ban is the type of a decision to ban, as CrowdSec names it.
|
||||
ban = "ban"
|
||||
)
|
||||
|
||||
func TestCrowdSecDecisionBansItsNetblockUntilItEndsEvenWithTheEngineDown(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
began := time.Now()
|
||||
manual := "manual 'ban' from 'localhost'"
|
||||
e := &engine{key: engineKey, decisions: []decision{
|
||||
{"Ip", suspect, ban, manual, began.Add(6 * time.Hour)},
|
||||
// A shorter decision on the same address, which is not the one
|
||||
// used.
|
||||
{"Ip", suspect, ban, sshBF, began.Add(4 * time.Hour)},
|
||||
{"Range", "198.51.100.0/24", ban, probing, began.Add(time.Hour)},
|
||||
{"Ip", "2001:db8::1", ban, sshBF, began.Add(2 * time.Hour)},
|
||||
// Left out: a decision to show a captcha, and one on a country.
|
||||
{"Ip", "192.0.2.50", "captcha", probing, began.Add(time.Hour)},
|
||||
{"Country", "KP", ban, manual, began.Add(time.Hour)},
|
||||
}}
|
||||
lists := start(t, e, crowdSecParams())
|
||||
|
||||
for addr, want := range map[string]reputation.Decision{
|
||||
suspect: {Expires: began.Add(6 * time.Hour), Scenario: manual},
|
||||
"198.51.100.0": {Expires: began.Add(time.Hour), Scenario: probing},
|
||||
"198.51.100.255": {Expires: began.Add(time.Hour), Scenario: probing},
|
||||
"2001:db8::1": {Expires: began.Add(2 * time.Hour), Scenario: sshBF},
|
||||
"203.0.113.10": {},
|
||||
"198.51.101.0": {},
|
||||
"2001:db8::2": {},
|
||||
"192.0.2.50": {},
|
||||
} {
|
||||
wantDecision(t, lists, addr, want)
|
||||
}
|
||||
|
||||
// With the engine down, the copy kept still holds the decision on
|
||||
// 198.51.100.0/24, which no longer bans once it has ended.
|
||||
e.set(func(e *engine) { e.failing = true })
|
||||
time.Sleep(time.Hour - time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantDecision(t, lists, "198.51.100.7",
|
||||
reputation.Decision{Expires: began.Add(time.Hour), Scenario: probing})
|
||||
|
||||
time.Sleep(time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantDecision(t, lists, "198.51.100.7", reputation.Decision{})
|
||||
wantDecision(t, lists, suspect,
|
||||
reputation.Decision{Expires: began.Add(6 * time.Hour), Scenario: manual})
|
||||
})
|
||||
}
|
||||
|
||||
func TestCrowdSecDecisionListFetchedAgainEveryMinute(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
began := time.Now()
|
||||
e := &engine{key: engineKey, decisions: []decision{
|
||||
{"Ip", suspect, ban, sshBF, began.Add(4 * time.Hour)},
|
||||
}}
|
||||
lists := start(t, e, crowdSecParams())
|
||||
wantEngineFetches(t, e, 1)
|
||||
|
||||
added := reputation.Decision{Expires: began.Add(2 * time.Hour), Scenario: probing}
|
||||
|
||||
e.set(func(e *engine) {
|
||||
e.decisions = append(e.decisions,
|
||||
decision{"Ip", "203.0.113.10", ban, probing, added.Expires})
|
||||
})
|
||||
|
||||
time.Sleep(time.Minute - time.Nanosecond)
|
||||
wantEngineFetches(t, e, 1)
|
||||
wantDecision(t, lists, "203.0.113.10", reputation.Decision{})
|
||||
|
||||
time.Sleep(time.Nanosecond)
|
||||
wantEngineFetches(t, e, 2)
|
||||
wantDecision(t, lists, "203.0.113.10", added)
|
||||
})
|
||||
}
|
||||
|
||||
func TestCrowdSecFailureKeepsTheLastGoodCopyAlertsOncePerCooldownAndHidesTheKey(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
for _, tc := range crowdSecFailures() {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
var log bytes.Buffer
|
||||
|
||||
began := time.Now()
|
||||
e := &engine{key: engineKey, decisions: []decision{
|
||||
{"Ip", suspect, ban, sshBF, began.Add(4 * time.Hour)},
|
||||
}}
|
||||
queue := newQueue()
|
||||
p := crowdSecParams()
|
||||
p.ProcessLog = slog.New(slog.NewJSONHandler(&log, nil))
|
||||
p.Alerts = queue
|
||||
lists := start(t, e, p)
|
||||
kept := lists.Snapshot()
|
||||
|
||||
e.set(tc.fail)
|
||||
|
||||
// Each failure is tried again a minute after it.
|
||||
for range 2 {
|
||||
time.Sleep(time.Minute)
|
||||
synctest.Wait()
|
||||
}
|
||||
|
||||
wantEngineFetches(t, e, 3)
|
||||
wantDecision(t, lists, suspect,
|
||||
reputation.Decision{Expires: began.Add(4 * time.Hour), Scenario: sshBF})
|
||||
|
||||
want := kept[0]
|
||||
want.Tried = time.Now()
|
||||
|
||||
if got := lists.Snapshot(); !reflect.DeepEqual(got, []reputation.List{want}) {
|
||||
t.Errorf("lists %+v, want the first copy, last tried now, %+v", got, want)
|
||||
}
|
||||
|
||||
if lists.Failures(decisionsURL) != 2 {
|
||||
t.Errorf("%d failures, want 2", lists.Failures(decisionsURL))
|
||||
}
|
||||
|
||||
// One alert for the first failure; the cooldown holds back the
|
||||
// second.
|
||||
wantAlert(t, queue,
|
||||
fetchFailure(time.Now().Add(-time.Minute), decisionsURL, tc.error))
|
||||
|
||||
if !strings.Contains(log.String(), `"msg":"fetching a list failed",`+
|
||||
`"url":"`+decisionsURL+`","error":"`+tc.error) {
|
||||
t.Errorf("logged\n%s\nwant the failures", log.String())
|
||||
}
|
||||
|
||||
wantKeyNotShown(t, log.String(), lists, queue)
|
||||
})
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// crowdSecFailure is a way for the engine to fail: fail has it answer the
|
||||
// fetches after the first so that they fail with error.
|
||||
type crowdSecFailure struct {
|
||||
name string
|
||||
fail func(e *engine)
|
||||
error string
|
||||
}
|
||||
|
||||
// crowdSecFailures returns the ways the engine can fail.
|
||||
func crowdSecFailures() []crowdSecFailure {
|
||||
const notDecision = " does not give an address or a netblock and a duration, " +
|
||||
"such as 4h0m0s"
|
||||
|
||||
return []crowdSecFailure{
|
||||
{
|
||||
"an answer other than 200",
|
||||
func(e *engine) { e.failing = true },
|
||||
"the server answered 503 Service Unavailable",
|
||||
},
|
||||
{
|
||||
"a key the engine refuses",
|
||||
func(e *engine) { e.key = "another-key-0123456789abcdef" },
|
||||
"the server answered 403 Forbidden",
|
||||
},
|
||||
{
|
||||
"an answer that does not read",
|
||||
func(e *engine) { e.answer = "<html>" },
|
||||
"read the answer: invalid character '<' looking for beginning of value",
|
||||
},
|
||||
{
|
||||
"a decision to ban whose value does not read",
|
||||
func(e *engine) {
|
||||
e.answer = `[{"duration": "4h", "scenario": "` + sshBF + `", ` +
|
||||
`"scope": "Ip", "type": "ban", "value": "203.0.113.300"}]`
|
||||
},
|
||||
"decision 1" + notDecision,
|
||||
},
|
||||
{
|
||||
"a decision to ban whose duration does not read",
|
||||
func(e *engine) {
|
||||
e.answer = `[{"duration": "4h", "scope": "Country", "type": "ban", ` +
|
||||
`"value": "KP"}, {"duration": "four hours", "scope": "Range", ` +
|
||||
`"type": "ban", "value": "198.51.100.0/24"}]`
|
||||
},
|
||||
"decision 2" + notDecision,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// wantKeyNotShown checks that the engine's key is in none of what the
|
||||
// fetches leave behind: log, the process log, the alerts waiting in queue,
|
||||
// and the copies of lists, which reputation.json keeps.
|
||||
func wantKeyNotShown(
|
||||
t *testing.T, log string, lists *reputation.Lists, queue *alerts.Queue,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
shown, err := json.Marshal([]any{lists.Snapshot(), waiting(queue)})
|
||||
if err != nil {
|
||||
t.Fatalf("encode: %v", err)
|
||||
}
|
||||
|
||||
if strings.Contains(log+string(shown), engineKey) {
|
||||
t.Errorf("the key is shown in\n%s\n%s", log, shown)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrowdSecDecisionListKeptAcrossARestartEndsWhenItsDecisionsDo(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
began := time.Now()
|
||||
e := &engine{key: engineKey, decisions: []decision{
|
||||
{"Range", "198.51.100.0/24", ban, probing, began.Add(time.Hour)},
|
||||
}}
|
||||
lists := start(t, e, crowdSecParams())
|
||||
kept := lists.Snapshot()
|
||||
|
||||
// Restarted half an hour later with what reputation.json keeps, and
|
||||
// the engine down, the decision still bans, until the end it had at
|
||||
// the fetch, half an hour on.
|
||||
time.Sleep(30 * time.Minute)
|
||||
|
||||
down := &engine{key: engineKey, failing: true}
|
||||
again := reputation.New(crowdSecParams())
|
||||
again.SetTransport(down)
|
||||
|
||||
err := again.Load(kept)
|
||||
if err != nil {
|
||||
t.Fatalf("load: %v", err)
|
||||
}
|
||||
|
||||
run(t, again)
|
||||
|
||||
want := reputation.Decision{Expires: began.Add(time.Hour), Scenario: probing}
|
||||
wantDecision(t, again, "198.51.100.7", want)
|
||||
|
||||
time.Sleep(30*time.Minute - time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantDecision(t, again, "198.51.100.7", want)
|
||||
|
||||
time.Sleep(time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantDecision(t, again, "198.51.100.7", reputation.Decision{})
|
||||
})
|
||||
}
|
||||
|
||||
func TestLoadTakesACrowdSecListNeverFetchedAndRefusesACopyThatDoesNotRead(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
now := time.Date(2026, 10, 6, 0, 0, 0, 0, time.UTC)
|
||||
lists := reputation.New(crowdSecParams())
|
||||
|
||||
// Tried, and never fetched: there is no copy to read.
|
||||
err := lists.Load([]reputation.List{{URL: decisionsURL, Tried: now}})
|
||||
if err != nil {
|
||||
t.Errorf("load the list never fetched: %v", err)
|
||||
}
|
||||
|
||||
err = lists.Load([]reputation.List{{
|
||||
URL: decisionsURL, Tried: now, Fetched: now, Lines: []string{
|
||||
`[{"duration": "4h", "scope": "Range", "type": "ban", ` +
|
||||
`"value": "198.51.100.0/33"}]`,
|
||||
},
|
||||
}})
|
||||
|
||||
const want = "the copy of " + decisionsURL + ": decision 1 does not give an " +
|
||||
"address or a netblock and a duration, such as 4h0m0s"
|
||||
if err == nil || err.Error() != want {
|
||||
t.Errorf("error %v, want %s", err, want)
|
||||
}
|
||||
}
|
||||
|
||||
// engine is a stand-in for the local API of a CrowdSec engine. It answers
|
||||
// a fetch of the decision list that carries its key in X-Api-Key with its
|
||||
// decisions still in force, each with the time it has left as it answers,
|
||||
// by the bubble's clock, as an engine does, or with answer while that is
|
||||
// not "". It answers 403 to a fetch without its key, as an engine does,
|
||||
// and 503 while failing. It counts the fetches.
|
||||
type engine struct {
|
||||
mu sync.Mutex
|
||||
key string
|
||||
decisions []decision
|
||||
answer string
|
||||
failing bool
|
||||
fetches int
|
||||
}
|
||||
|
||||
// decision is a decision of the engine, which ends at expires.
|
||||
type decision struct {
|
||||
scope, value, kind, scenario string
|
||||
expires time.Time
|
||||
}
|
||||
|
||||
// RoundTrip has the engine answer req, in place of the network.
|
||||
func (e *engine) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||
e.mu.Lock()
|
||||
defer e.mu.Unlock()
|
||||
|
||||
e.fetches++
|
||||
|
||||
status, body := http.StatusOK, e.answer
|
||||
|
||||
switch {
|
||||
case req.URL.String() != decisionsURL || req.Header.Get("X-Api-Key") != e.key:
|
||||
status, body = http.StatusForbidden, `{"message":"access forbidden"}`
|
||||
case e.failing:
|
||||
status, body = http.StatusServiceUnavailable, ""
|
||||
case body == "":
|
||||
body = e.inForce(time.Now())
|
||||
}
|
||||
|
||||
return &http.Response{
|
||||
StatusCode: status,
|
||||
Status: fmt.Sprintf("%d %s", status, http.StatusText(status)),
|
||||
Header: http.Header{},
|
||||
Body: io.NopCloser(strings.NewReader(body)),
|
||||
Request: req,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// inForce returns the decisions in force at now, as the engine answers
|
||||
// them: a JSON list, null for none.
|
||||
func (e *engine) inForce(now time.Time) string {
|
||||
var answer []map[string]string
|
||||
|
||||
for _, d := range e.decisions {
|
||||
if now.Before(d.expires) {
|
||||
answer = append(answer, map[string]string{
|
||||
"duration": d.expires.Sub(now).String(), "origin": "crowdsec",
|
||||
"scenario": d.scenario, "scope": d.scope, "type": d.kind, "value": d.value,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
body, err := json.Marshal(answer)
|
||||
if err != nil {
|
||||
panic(err) // a list of maps of strings always encodes
|
||||
}
|
||||
|
||||
return string(body)
|
||||
}
|
||||
|
||||
// set changes the engine with change.
|
||||
func (e *engine) set(change func(e *engine)) {
|
||||
e.mu.Lock()
|
||||
defer e.mu.Unlock()
|
||||
|
||||
change(e)
|
||||
}
|
||||
|
||||
// crowdSecParams returns the Params of the decision list of the tests'
|
||||
// engine, fetched with its key, by the bubble's clock, with alerts to a
|
||||
// queue that sends none.
|
||||
func crowdSecParams() reputation.Params {
|
||||
p := params()
|
||||
p.CrowdSecDecisionsURL = decisionsURL
|
||||
p.CrowdSecKey = engineKey
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
// wantEngineFetches waits until Run has made the fetches due, and checks
|
||||
// how many the engine has had.
|
||||
func wantEngineFetches(t *testing.T, e *engine, want int) {
|
||||
t.Helper()
|
||||
|
||||
synctest.Wait()
|
||||
|
||||
e.mu.Lock()
|
||||
got := e.fetches
|
||||
e.mu.Unlock()
|
||||
|
||||
if got != want {
|
||||
t.Errorf("%d fetches, want %d", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// wantDecision checks the decision lists says is in force on addr now,
|
||||
// the zero Decision for none.
|
||||
func wantDecision(
|
||||
t *testing.T, lists *reputation.Lists, addr string, want reputation.Decision,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
got, listed := lists.CrowdSecDecision(netip.MustParseAddr(addr), time.Now())
|
||||
if listed != !want.Expires.IsZero() ||
|
||||
listed && (!got.Expires.Equal(want.Expires) || got.Scenario != want.Scenario) {
|
||||
t.Errorf("%s has the decision %+v in force %t, want %+v", addr, got, listed, want)
|
||||
}
|
||||
}
|
||||
@@ -1,8 +1,9 @@
|
||||
// Package reputation fetches the lists the settings name by URL: the
|
||||
// blocklists of SWWAF_BLOCKLIST_URLS, and the file of AS:percent lines
|
||||
// SWWAF_ASN_LIMIT_PERCENT_URL names. It keeps the last good copy of each,
|
||||
// whole, comment lines included, which is used while a fetch fails, and
|
||||
// when each was last tried. It also asks the DNSBL zones of
|
||||
// blocklists of SWWAF_BLOCKLIST_URLS, the file of AS:percent lines
|
||||
// SWWAF_ASN_LIMIT_PERCENT_URL names, and the decision list of the CrowdSec
|
||||
// engine SWWAF_CROWDSEC_LAPI_URL names. It keeps the last good copy of
|
||||
// each, whole, comment lines included, which is used while a fetch fails,
|
||||
// and when each was last tried. It also asks the DNSBL zones of
|
||||
// SWWAF_DNSBL_ZONES about clients, and keeps their verdicts, and checks
|
||||
// clients with AbuseIPDB, and keeps their scores and the checks spent
|
||||
// today. The state package writes all of these to reputation.json and
|
||||
@@ -11,6 +12,7 @@ package reputation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -32,6 +34,10 @@ const (
|
||||
maxListBytes = 16 << 20
|
||||
// fetchTimeout bounds one fetch of a list.
|
||||
fetchTimeout = time.Minute
|
||||
// crowdSecRefresh is how long after the CrowdSec decision list was last
|
||||
// fetched or tried it is fetched again: the engine is the operator's
|
||||
// own, and makes and ends decisions all the time.
|
||||
crowdSecRefresh = time.Minute
|
||||
// mappedBits is the length of ::ffff:0.0.0.0/96, the netblock of every
|
||||
// IPv4-mapped address.
|
||||
mappedBits = 96
|
||||
@@ -43,6 +49,8 @@ var (
|
||||
errNotNetblock = errors.New("is not an address or a netblock, such as 192.0.2.0/24")
|
||||
errNotASNPercent = errors.New(
|
||||
"is not an AS number, : and a percentage, such as AS64496:50")
|
||||
errNotDecision = errors.New(
|
||||
"does not give an address or a netblock and a duration, such as 4h0m0s")
|
||||
)
|
||||
|
||||
// List is a list as reputation.json holds it: the URL it is fetched from,
|
||||
@@ -63,8 +71,14 @@ type Params struct {
|
||||
// (SWWAF_ASN_LIMIT_PERCENT_URL), "" while it is unset.
|
||||
BlocklistURLs []string
|
||||
ASNLimitPercentURL string
|
||||
// CrowdSecDecisionsURL is the CrowdSec decision list, "" while
|
||||
// SWWAF_CROWDSEC_LAPI_URL is unset, fetched with CrowdSecKey
|
||||
// (SWWAF_CROWDSEC_LAPI_KEY).
|
||||
CrowdSecDecisionsURL string
|
||||
CrowdSecKey string
|
||||
// Refresh is how long after a list was last fetched or tried it is
|
||||
// fetched again (SWWAF_BLOCKLIST_REFRESH).
|
||||
// fetched again (SWWAF_BLOCKLIST_REFRESH), but for the CrowdSec decision
|
||||
// list, which is fetched again crowdSecRefresh after.
|
||||
Refresh time.Duration
|
||||
// Now tells the time, normally time.Now in UTC.
|
||||
Now func() time.Time
|
||||
@@ -95,12 +109,22 @@ type list struct {
|
||||
}
|
||||
|
||||
// entries are what the lines of a copy say: for a blocklist, the netblocks
|
||||
// it names, with the lengths among them, and for the file of AS:percent
|
||||
// lines, the percentage it gives each AS number.
|
||||
// it names, with the lengths among them, for the file of AS:percent lines,
|
||||
// the percentage it gives each AS number, and for the CrowdSec decision
|
||||
// list, the decision on each netblock that ends last, with the lengths
|
||||
// among them.
|
||||
type entries struct {
|
||||
netblocks map[netip.Prefix]bool
|
||||
lengths []int
|
||||
percents map[string]int64
|
||||
decisions map[netip.Prefix]Decision
|
||||
}
|
||||
|
||||
// Decision is a decision of the CrowdSec engine to ban a netblock: when
|
||||
// it ends, and the scenario that made it, such as crowdsecurity/ssh-bf.
|
||||
type Decision struct {
|
||||
Expires time.Time
|
||||
Scenario string
|
||||
}
|
||||
|
||||
// New returns the lists, without a copy of any yet.
|
||||
@@ -115,13 +139,18 @@ func New(params Params) *Lists {
|
||||
}
|
||||
|
||||
// URLs returns the URL of every list: the blocklists' in the order
|
||||
// SWWAF_BLOCKLIST_URLS names them, then SWWAF_ASN_LIMIT_PERCENT_URL.
|
||||
// SWWAF_BLOCKLIST_URLS names them, then SWWAF_ASN_LIMIT_PERCENT_URL, then
|
||||
// the CrowdSec decision list's.
|
||||
func (l *Lists) URLs() []string {
|
||||
urls := slices.Clone(l.params.BlocklistURLs)
|
||||
if l.params.ASNLimitPercentURL != "" {
|
||||
urls = append(urls, l.params.ASNLimitPercentURL)
|
||||
}
|
||||
|
||||
if l.params.CrowdSecDecisionsURL != "" {
|
||||
urls = append(urls, l.params.CrowdSecDecisionsURL)
|
||||
}
|
||||
|
||||
return urls
|
||||
}
|
||||
|
||||
@@ -157,6 +186,37 @@ func (l *Lists) ASNLimitPercent(asn string) (int64, bool) {
|
||||
return percent, listed
|
||||
}
|
||||
|
||||
// CrowdSecDecision returns the decision of the copy of the CrowdSec
|
||||
// decision list on a netblock that holds addr and that ends last, and
|
||||
// whether it is still in force at now. A decision that has ended no
|
||||
// longer bans, even before the next fetch drops it.
|
||||
func (l *Lists) CrowdSecDecision(addr netip.Addr, now time.Time) (Decision, bool) {
|
||||
if l.params.CrowdSecDecisionsURL == "" {
|
||||
return Decision{}, false
|
||||
}
|
||||
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
|
||||
kept := l.lists[l.params.CrowdSecDecisionsURL].entries
|
||||
|
||||
var last Decision
|
||||
|
||||
for _, length := range kept.lengths {
|
||||
netblock, err := addr.Prefix(length)
|
||||
if err != nil {
|
||||
continue // an IPv6 netblock's length, past an IPv4 address's 32 bits
|
||||
}
|
||||
|
||||
decision := kept.decisions[netblock]
|
||||
if decision.Expires.After(last.Expires) {
|
||||
last = decision
|
||||
}
|
||||
}
|
||||
|
||||
return last, now.Before(last.Expires)
|
||||
}
|
||||
|
||||
// Fetched returns when the copy in use of the list at listURL was
|
||||
// fetched, or zero while there is none.
|
||||
func (l *Lists) Fetched(listURL string) time.Time {
|
||||
@@ -174,10 +234,9 @@ func (l *Lists) Failures(listURL string) int {
|
||||
return l.lists[listURL].failures
|
||||
}
|
||||
|
||||
// Run fetches each list once Refresh has passed since it was last fetched
|
||||
// or tried, the later of the two, until ctx is done. A list never tried is
|
||||
// fetched at once, and so is one whose last try or copy, read from
|
||||
// reputation.json, is that old.
|
||||
// Run fetches each list once it is due, as due tells, until ctx is done. A
|
||||
// list never tried is fetched at once, and so is one that is due by its
|
||||
// last try or copy read from reputation.json.
|
||||
func (l *Lists) Run(ctx context.Context) {
|
||||
if len(l.lists) == 0 {
|
||||
return
|
||||
@@ -225,11 +284,12 @@ func (l *Lists) Load(lists []List) error {
|
||||
found := make(map[string]entries, len(lists))
|
||||
|
||||
for _, kept := range lists {
|
||||
if _, named := l.lists[kept.URL]; !named {
|
||||
continue
|
||||
_, named := l.lists[kept.URL]
|
||||
if !named || kept.Fetched.IsZero() {
|
||||
continue // dropped, or a list tried but never fetched, without a copy
|
||||
}
|
||||
|
||||
read, err := l.parse(kept.URL, kept.Lines)
|
||||
read, err := l.parse(kept.URL, kept.Lines, kept.Fetched)
|
||||
if err != nil {
|
||||
return fmt.Errorf("the copy of %s: %w", kept.URL, err)
|
||||
}
|
||||
@@ -245,9 +305,8 @@ func (l *Lists) Load(lists []List) error {
|
||||
}
|
||||
|
||||
for _, kept := range lists {
|
||||
read, named := found[kept.URL]
|
||||
if named {
|
||||
l.lists[kept.URL].kept, l.lists[kept.URL].entries = kept, read
|
||||
if _, named := l.lists[kept.URL]; named {
|
||||
l.lists[kept.URL].kept, l.lists[kept.URL].entries = kept, found[kept.URL]
|
||||
}
|
||||
}
|
||||
|
||||
@@ -276,7 +335,8 @@ func (l *Lists) fetchDue(ctx context.Context) time.Time {
|
||||
}
|
||||
|
||||
// due returns when the list at listURL is to be fetched: Refresh after it
|
||||
// was last fetched or tried, the later of the two.
|
||||
// was last fetched or tried, the later of the two, or crowdSecRefresh
|
||||
// after for the CrowdSec decision list.
|
||||
func (l *Lists) due(listURL string) time.Time {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
@@ -288,6 +348,10 @@ func (l *Lists) due(listURL string) time.Time {
|
||||
last = held.kept.Tried
|
||||
}
|
||||
|
||||
if listURL == l.params.CrowdSecDecisionsURL {
|
||||
return last.Add(crowdSecRefresh)
|
||||
}
|
||||
|
||||
return last.Add(l.params.Refresh)
|
||||
}
|
||||
|
||||
@@ -298,14 +362,14 @@ func (l *Lists) due(listURL string) time.Time {
|
||||
// so that a restart waits for it: the server may have had its request.
|
||||
func (l *Lists) fetch(ctx context.Context, listURL string) {
|
||||
lines, err := l.get(ctx, listURL)
|
||||
now := l.params.Now()
|
||||
|
||||
var found entries
|
||||
if err == nil {
|
||||
found, err = l.parse(listURL, lines)
|
||||
found, err = l.parse(listURL, lines, now)
|
||||
}
|
||||
|
||||
cutOff := err != nil && ctx.Err() != nil
|
||||
now := l.params.Now()
|
||||
|
||||
l.mu.Lock()
|
||||
|
||||
@@ -350,8 +414,10 @@ func raiseFailure(queue *alerts.Queue, reason, source string, err error) {
|
||||
})
|
||||
}
|
||||
|
||||
// get fetches the list at listURL, and returns its lines. An answer other
|
||||
// than 200, or a list longer than maxListBytes, is a failure.
|
||||
// get fetches the list at listURL, and returns its lines. The CrowdSec
|
||||
// decision list is fetched with CrowdSecKey in the header X-Api-Key, where
|
||||
// the engine looks for it. An answer other than 200, or a list longer
|
||||
// than maxListBytes, is a failure.
|
||||
func (l *Lists) get(ctx context.Context, listURL string) ([]string, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, fetchTimeout)
|
||||
defer cancel()
|
||||
@@ -361,6 +427,10 @@ func (l *Lists) get(ctx context.Context, listURL string) ([]string, error) {
|
||||
return nil, fmt.Errorf("make the request: %w", err)
|
||||
}
|
||||
|
||||
if listURL == l.params.CrowdSecDecisionsURL {
|
||||
req.Header.Set("X-Api-Key", l.params.CrowdSecKey)
|
||||
}
|
||||
|
||||
res, err := l.httpClient.Do(req)
|
||||
if err != nil {
|
||||
// Do's error names the URL, which the log line and the alert name
|
||||
@@ -393,16 +463,22 @@ func (l *Lists) get(ctx context.Context, listURL string) ([]string, error) {
|
||||
return lines, nil
|
||||
}
|
||||
|
||||
// parse reads the lines of the list at listURL: those of a blocklist, or
|
||||
// of the file of AS:percent lines. Anything after a ; or a # on a line is
|
||||
// parse reads the lines of the list at listURL, fetched at fetched: those
|
||||
// of a blocklist, of the file of AS:percent lines, or of the CrowdSec
|
||||
// decision list. In the first two, anything after a ; or a # on a line is
|
||||
// left out, and so is a line left blank. Any other line that does not read
|
||||
// is an error naming it by its number.
|
||||
func (l *Lists) parse(listURL string, lines []string) (entries, error) {
|
||||
if listURL == l.params.ASNLimitPercentURL {
|
||||
func (l *Lists) parse(
|
||||
listURL string, lines []string, fetched time.Time,
|
||||
) (entries, error) {
|
||||
switch listURL {
|
||||
case l.params.ASNLimitPercentURL:
|
||||
return parsePercents(lines)
|
||||
case l.params.CrowdSecDecisionsURL:
|
||||
return parseDecisions(lines, fetched)
|
||||
default:
|
||||
return parseNetblocks(lines)
|
||||
}
|
||||
|
||||
return parseNetblocks(lines)
|
||||
}
|
||||
|
||||
// parseNetblocks reads a blocklist's lines, each an address or a netblock
|
||||
@@ -486,6 +562,56 @@ func parsePercents(lines []string) (entries, error) {
|
||||
return found, nil
|
||||
}
|
||||
|
||||
// parseDecisions reads the lines of the CrowdSec decision list fetched at
|
||||
// fetched: the engine's answer, a JSON list of its decisions in force,
|
||||
// null while it has none. A decision of the type ban whose scope is Ip or
|
||||
// Range, as CrowdSec names them, bans its value, an address or a netblock
|
||||
// as parseNetblock reads it, until its duration, the time it had left as
|
||||
// the engine answered, has passed since fetched. Any other decision, such
|
||||
// as one to show a captcha or one on a country, is left out. A decision
|
||||
// to ban whose value or duration does not read is an error naming it by
|
||||
// its number.
|
||||
func parseDecisions(lines []string, fetched time.Time) (entries, error) {
|
||||
var answer []struct {
|
||||
Duration string `json:"duration"`
|
||||
Scenario string `json:"scenario"`
|
||||
Scope string `json:"scope"`
|
||||
Type string `json:"type"`
|
||||
Value string `json:"value"`
|
||||
}
|
||||
|
||||
err := json.Unmarshal([]byte(strings.Join(lines, "\n")), &answer)
|
||||
if err != nil {
|
||||
return entries{}, fmt.Errorf("read the answer: %w", err)
|
||||
}
|
||||
|
||||
found := entries{decisions: map[netip.Prefix]Decision{}}
|
||||
|
||||
for i, decision := range answer {
|
||||
if decision.Type != "ban" || (decision.Scope != "Ip" && decision.Scope != "Range") {
|
||||
continue
|
||||
}
|
||||
|
||||
netblock, ok := parseNetblock(decision.Value)
|
||||
|
||||
duration, err := time.ParseDuration(decision.Duration)
|
||||
if !ok || err != nil {
|
||||
return entries{}, fmt.Errorf("decision %d %w", i+1, errNotDecision)
|
||||
}
|
||||
|
||||
expires := fetched.Add(duration)
|
||||
if expires.After(found.decisions[netblock].Expires) {
|
||||
found.decisions[netblock] = Decision{Expires: expires, Scenario: decision.Scenario}
|
||||
}
|
||||
|
||||
if !slices.Contains(found.lengths, netblock.Bits()) {
|
||||
found.lengths = append(found.lengths, netblock.Bits())
|
||||
}
|
||||
}
|
||||
|
||||
return found, nil
|
||||
}
|
||||
|
||||
// withoutComment returns line without anything after a ; or a #, and
|
||||
// without the spaces around what is left.
|
||||
func withoutComment(line string) string {
|
||||
|
||||
@@ -215,12 +215,7 @@ func TestFailedFetchKeepsTheLastGoodCopyAndAlertsOncePerCooldown(t *testing.T) {
|
||||
|
||||
// One alert for the first failure; the cooldown holds back the
|
||||
// second.
|
||||
wantAlert(t, queue, alerts.Alert{
|
||||
Time: time.Now().Add(-refresh),
|
||||
Event: alerts.EventSourceFailure,
|
||||
Reason: "fetching a list failed",
|
||||
Detail: map[string]any{"source": dropURL, "error": tc.error},
|
||||
})
|
||||
wantAlert(t, queue, fetchFailure(time.Now().Add(-refresh), dropURL, tc.error))
|
||||
|
||||
if !strings.Contains(log.String(), `"msg":"fetching a list failed",`+
|
||||
`"url":"`+dropURL+`","error":"`+tc.error) {
|
||||
@@ -535,7 +530,9 @@ func newQueue() *alerts.Queue {
|
||||
|
||||
// start returns the lists of p, fetched through servers by Run, which runs
|
||||
// until the test ends, once Run has fetched those due at start.
|
||||
func start(t *testing.T, servers *standIn, p reputation.Params) *reputation.Lists {
|
||||
func start(
|
||||
t *testing.T, servers http.RoundTripper, p reputation.Params,
|
||||
) *reputation.Lists {
|
||||
t.Helper()
|
||||
|
||||
lists := reputation.New(p)
|
||||
@@ -597,6 +594,17 @@ func waiting(queue *alerts.Queue) []alerts.Alert {
|
||||
return queue.Snapshot().Waiting[alerts.DestinationWebhook]
|
||||
}
|
||||
|
||||
// fetchFailure is the source_failure alert raised at the time raised for
|
||||
// a fetch of the list at listURL that failed with err.
|
||||
func fetchFailure(raised time.Time, listURL, err string) alerts.Alert {
|
||||
return alerts.Alert{
|
||||
Time: raised,
|
||||
Event: alerts.EventSourceFailure,
|
||||
Reason: "fetching a list failed",
|
||||
Detail: map[string]any{"source": listURL, "error": err},
|
||||
}
|
||||
}
|
||||
|
||||
// wantAlert checks that want is the one alert waiting in queue, and that
|
||||
// the cooldown has held back one repeat of it.
|
||||
func wantAlert(t *testing.T, queue *alerts.Queue, want alerts.Alert) {
|
||||
|
||||
Reference in New Issue
Block a user