Anomaly thresholds: alerts for unusual traffic, nothing refused (closes #101)
check / check (push) Waiting to run

SWWAF_ANOMALY_CLIENT_*, _NET_*, _ASN_*, _TOTAL_* and SWWAF_WATCH_* with
SWWAF_WATCH_NETS: requests and bytes per minute and per hour, each off by
default; with all off, nothing is counted. Otherwise every request but the
health check is counted, allow-listed and exempt ones included; a count
over its threshold raises an anomaly alert, with a cooldown per scope. At
most 20,000 counters, kept in alerts.json. A per-AS-number threshold with
lookups off, or a malformed SWWAF_WATCH_NETS, stops the start. A cooldown
that has run out is dropped as the hour ends, whatever it held back; the
hour's summary gives its repeats.

Judgement call: refused requests are counted too.
Judgement call: per-client counters are kept in alerts.json, which SPEC.md does not list.
Judgement call: a request counts for an AS number only if the lookup answered before it ended.

Model: opus-5-5
This commit is contained in:
2026-10-07 14:04:15 +00:00
parent 2421cdc273
commit acde5bde07
17 changed files with 2132 additions and 297 deletions
+204 -115
View File
@@ -22,32 +22,33 @@ fields, which come a little later, and the metrics endpoint and the header size
and the idle time as settings, which come last in it. So are the four parts of and the idle time as settings, which come last in it. So are the four parts of
the stage after it: the rule files, with the bans for a clear sign of attack, the stage after it: the rule files, with the bans for a clear sign of attack,
the other admin endpoints, alerts to all three destinations, a JSON webhook, the other admin endpoints, alerts to all three destinations, a JSON webhook,
Slack and ntfy, and remote log sending. So are three parts of the stage after Slack and ntfy, and remote log sending. So is the stage after that: the AS
that: the AS number and country of every client, looked up through GeoJS or in number and country of every client, looked up through GeoJS or in the IPinfo
the IPinfo Lite database file, the byte limits, and the biased thresholds, lower Lite database file, the byte limits, the biased thresholds, lower limits for the
limits for the AS numbers and countries you list. `smallwebwaf` passes each AS numbers and countries you list, and the anomaly thresholds, alerts for
request to the app and the app's answer back, unchanged, within its timeouts and unusual traffic that refuse nothing. `smallwebwaf` passes each request to the
size limits, works out each client's address, looks up its AS number and country app and the app's answer back, unchanged, within its timeouts and size limits,
unless you switch that off, bans a client that sends too many requests or too works out each client's address, looks up its AS number and country unless you
many bytes, not counting those for the paths you choose, with lower limits for switch that off, bans a client that sends too many requests or too many bytes,
the clients of the AS numbers and countries you list, refuses a client that not counting those for the paths you choose, with lower limits for the clients
comes from a country you refuse or from a network you refuse, lets the networks of the AS numbers and countries you list, refuses a client that comes from a
you choose through, checks each request against the rule files and bans a client country you refuse or from a network you refuse, lets the networks you choose
whose request is a clear sign of attack, keeps its bans, each client's counters through, checks each request against the rule files and bans a client whose
and history, and GeoJS's answers in JSON files across restarts, takes in your request is a clear sign of attack, keeps its bans, each client's counters and
edits of those files, such as a ban you make, keep or lift, and of the rule history, and GeoJS's answers in JSON files across restarts, takes in your edits
files while it runs, writes a JSON log line for every request, sends its log of those files, such as a ban you make, keep or lift, and of the rule files
lines to a syslog server too if you name one, sends an alert to a webhook, to while it runs, writes a JSON log line for every request, sends its log lines to
Slack and to ntfy, each if you name one, for each ban it makes or makes a syslog server too if you name one, sends an alert to a webhook, to Slack and
permanent, for GeoJS failing, for a rule file or state file with an error and to ntfy, each if you name one, for each ban it makes or makes permanent, for
for a replacement of the lookup database it cannot read, serves Prometheus traffic over an anomaly threshold you set, for GeoJS failing, for a rule file or
metrics to a scraper that holds the metrics token, lets an admin who holds the state file with an error and for a replacement of the lookup database it cannot
admin token list, add and lift bans and ask what it knows of a client, and in read, serves Prometheus metrics to a scraper that holds the metrics token, lets
`observe` mode passes on the requests it would refuse, logging what it would an admin who holds the admin token list, add and lift bans and ask what it knows
have done with them. It comes as the image the app's own image is built on. The of a client, and in `observe` mode passes on the requests it would refuse,
rest of the design comes after that, in the order of the build order in logging what it would have done with them. It comes as the image the app's own
[`SPEC.md`](SPEC.md). The survey of existing tools that led to the design is in image is built on. The rest of the design comes after that, in the order of the
[`EVALUATION.md`](EVALUATION.md). build order in [`SPEC.md`](SPEC.md). The survey of existing tools that led to
the design is in [`EVALUATION.md`](EVALUATION.md).
## Getting started ## Getting started
@@ -177,13 +178,13 @@ in `bin/state` unless `SWWAF_STATE_DIR` is set, and the default rule file of
IPinfo Lite database file while `SWWAF_LOOKUP_SOURCE` is `file`, after the IPinfo Lite database file while `SWWAF_LOOKUP_SOURCE` is `file`, after the
static lists and bans, unless `SWWAF_LOOKUP_SOURCE` is `off` (see "Country and static lists and bans, unless `SWWAF_LOOKUP_SOURCE` is `off` (see "Country and
AS number lookup" below), for the request log, the client's history, the notes AS number lookup" below), for the request log, the client's history, the notes
of its bans, their alerts and the metrics. The file answers at once. With of its bans, their alerts, the metrics and the anomaly thresholds per AS
GeoJS, a request waits for its client's first answer only while a setting acts number. The file answers at once. With GeoJS, a request waits for its client's
on it, a country list, `SWWAF_ADD_LOOKUP_HEADERS` or a biased threshold that first answer only while a setting acts on it, a country list,
lowers a limit. Otherwise it goes on at once, and the answer reaches the `SWWAF_ADD_LOOKUP_HEADERS` or a biased threshold that lowers a limit.
client's history and the notes of its bans when it comes, but not the log Otherwise it goes on at once, and the answer reaches the client's history and
lines of the requests that went on without it, nor the alerts already raised the notes of its bans when it comes, but not the log lines of the requests
for those bans. that went on without it, nor the alerts already raised for those bans.
- Refuses a request from a country you refuse with `SWWAF_BAN_RESPONSE`, as soon - Refuses a request from a country you refuse with `SWWAF_BAN_RESPONSE`, as soon
as the client's country is known and before its body is read; such a request as the client's country is known and before its body is read; such a request
is not counted for the rate limits. A client on a private, loopback or is not counted for the rate limits. A client on a private, loopback or
@@ -237,13 +238,30 @@ in `bin/state` unless `SWWAF_STATE_DIR` is set, and the default rule file of
- Sends every line it writes on stdout to a syslog server as well, while - Sends every line it writes on stdout to a syslog server as well, while
`SWWAF_LOG_REMOTE_URL` names one (see "Sending the log to a syslog server" `SWWAF_LOG_REMOTE_URL` names one (see "Sending the log to a syslog server"
below). below).
- Sends an alert for each ban it makes or makes permanent, for GeoJS failing, - Sends an alert for each ban it makes or makes permanent, for a count over an
for a rule file or state file with an error, and for a replacement of the anomaly threshold, for GeoJS failing, for a rule file or state file with an
lookup database it cannot read, holding back repeats and, past an hourly error, and for a replacement of the lookup database it cannot read, holding
limit, rolling the rest into one summary, to each destination you name: as a back repeats and, past an hourly limit, rolling the rest into one summary, to
JSON object to the webhook `SWWAF_ALERT_WEBHOOK_URL` names, as a message to each destination you name: as a JSON object to the webhook
the Slack incoming webhook `SWWAF_ALERT_SLACK_WEBHOOK_URL` names, and as a `SWWAF_ALERT_WEBHOOK_URL` names, as a message to the Slack incoming webhook
message to the ntfy topic `SWWAF_ALERT_NTFY_URL` names (see "Alerts" below). `SWWAF_ALERT_SLACK_WEBHOOK_URL` names, and as a message to the ntfy topic
`SWWAF_ALERT_NTFY_URL` names (see "Alerts" below).
- Counts requests and their bytes over a minute and an hour, per client, per
netblock around a client, per AS number, for the whole service and per named
netblock, and sends an `anomaly` alert for a count over the anomaly threshold
you set for it (see the anomaly thresholds below). These thresholds only
alert: they refuse and ban nothing. A scope whose four thresholds are all off
is not counted, and within a scope only the counts whose threshold is set are
counted. Every request but the health check is counted, whatever is done with
it: one that is refused, one from a client in `SWWAF_ALLOW_NETS` or
`SWWAF_RATE_LIMIT_EXEMPT_NETS`, and one for a path in
`SWWAF_RATE_LIMIT_EXEMPT_PATHS`, with its body bytes once it has ended, those
of the answer, of the request or both, as `SWWAF_BYTES_COUNT` says. A request
is counted for its client's AS number only if the lookup has given that by the
time the request ends: no request waits for it, and a client in
`SWWAF_ALLOW_NETS`, which is not looked up, counts for no AS number. At most
20,000 counters are kept, the one counted least recently dropped first, and
`alerts.json` keeps them across a restart (see "State files" below).
## Settings ## Settings
@@ -330,9 +348,10 @@ effective settings are logged at start.
address of every new visitor, `file`, the IPinfo Lite database file address of every new visitor, `file`, the IPinfo Lite database file
`SWWAF_LOOKUP_DB_PATH` names, or `off`, which looks up no client and sends no `SWWAF_LOOKUP_DB_PATH` names, or `off`, which looks up no client and sends no
address to GeoJS. With `off`, a country list that is not empty, address to GeoJS. With `off`, a country list that is not empty,
`SWWAF_ADD_LOOKUP_HEADERS` set to `true`, or a biased threshold that lowers a `SWWAF_ADD_LOOKUP_HEADERS` set to `true`, a biased threshold that lowers a
limit, a list of them that is not empty or `SWWAF_UNKNOWN_LIMIT_PERCENT` below limit, a list of them that is not empty or `SWWAF_UNKNOWN_LIMIT_PERCENT` below
100, stops the start, with a message naming it and `SWWAF_LOOKUP_SOURCE`. 100, or an anomaly threshold per AS number that is not `off`, stops the start,
with a message naming it and `SWWAF_LOOKUP_SOURCE`.
- `SWWAF_LOOKUP_DB_PATH` (default empty): the IPinfo Lite database file, in its - `SWWAF_LOOKUP_DB_PATH` (default empty): the IPinfo Lite database file, in its
`.mmdb` form, for `SWWAF_LOOKUP_SOURCE=file`. `file` without it, or it with `.mmdb` form, for `SWWAF_LOOKUP_SOURCE=file`. `file` without it, or it with
any other `SWWAF_LOOKUP_SOURCE`, the default included, stops the start, with a any other `SWWAF_LOOKUP_SOURCE`, the default included, stops the start, with a
@@ -467,16 +486,46 @@ effective settings are logged at start.
`SWWAF_INSTANCE_NAME`, which ntfy is sent in the title. `SWWAF_INSTANCE_NAME`, which ntfy is sent in the title.
- `SWWAF_ALERT_EVENTS` (default - `SWWAF_ALERT_EVENTS` (default
`ban,permanent_ban,waf_block,anomaly,reputation_hit,source_failure,file_error`): `ban,permanent_ban,waf_block,anomaly,reputation_hit,source_failure,file_error`):
the events alerts are sent for. `waf_block`, `anomaly` and `reputation_hit` the events alerts are sent for. `waf_block` and `reputation_hit` come with the
come with the features that raise them; nothing raises them yet. features that raise them; nothing raises them yet.
- `SWWAF_ALERT_COOLDOWN` (default `15m`): how long a repeat of an alert is held - `SWWAF_ALERT_COOLDOWN` (default `15m`): how long a repeat of an alert is held
back (see "Alerts" below). back (see "Alerts" below).
- `SWWAF_ALERT_MAX_PER_HOUR` (default `60`): the most alerts sent in an hour; - `SWWAF_ALERT_MAX_PER_HOUR` (default `60`): the most alerts sent in an hour;
the rest of the hour's alerts are rolled into one summary. the rest of the hour's alerts are rolled into one summary.
- `SWWAF_ANOMALY_CLIENT_REQUESTS_PER_MINUTE`,
`SWWAF_ANOMALY_CLIENT_REQUESTS_PER_HOUR`,
`SWWAF_ANOMALY_CLIENT_BYTES_PER_MINUTE` and
`SWWAF_ANOMALY_CLIENT_BYTES_PER_HOUR` (default `off`): the anomaly thresholds
per client, the most requests and the most bytes a client may have counted in
a minute and in an hour before an `anomaly` alert is sent for it. They refuse
and ban nothing. Each scope below has the same four thresholds, their names
ending in `REQUESTS_PER_MINUTE`, `REQUESTS_PER_HOUR`, `BYTES_PER_MINUTE` and
`BYTES_PER_HOUR`, and each is `off` by default, since what is unusual depends
on each service's normal traffic, which the metrics show.
- `SWWAF_ANOMALY_NET_REQUESTS_PER_MINUTE`, `..._PER_HOUR`,
`SWWAF_ANOMALY_NET_BYTES_PER_MINUTE` and `..._PER_HOUR` (default `off`): the
anomaly thresholds per netblock around a client, which is
`SWWAF_ANOMALY_NET_V4_PREFIX` (default `24`) long, from 0 to 32, for an IPv4
client, and `SWWAF_ANOMALY_NET_V6_PREFIX` (default `48`) long, from 0 to 128,
for an IPv6 one.
- `SWWAF_ANOMALY_ASN_REQUESTS_PER_MINUTE`, `..._PER_HOUR`,
`SWWAF_ANOMALY_ASN_BYTES_PER_MINUTE` and `..._PER_HOUR` (default `off`): the
anomaly thresholds per AS number.
- `SWWAF_ANOMALY_TOTAL_REQUESTS_PER_MINUTE`, `..._PER_HOUR`,
`SWWAF_ANOMALY_TOTAL_BYTES_PER_MINUTE` and `..._PER_HOUR` (default `off`): the
anomaly thresholds for the whole service.
- `SWWAF_WATCH_NETS` (default empty): named netblocks, each a name, `=` and a
netblock, such as `office=203.0.113.0/24,scraper-x=198.51.100.0/22`. An item
without `=`, without a name or without a valid netblock, or a name listed
twice, stops the start. `SWWAF_WATCH_REQUESTS_PER_MINUTE`, `..._PER_HOUR`,
`SWWAF_WATCH_BYTES_PER_MINUTE` and `..._PER_HOUR` (default `off`) are the
anomaly thresholds of each named netblock as a whole, which counts every
client inside it; a client inside several is counted in each.
Durations are in Go's syntax, with `d` for days (`90s`, `15m`, `7d`). Sizes are Durations are in Go's syntax, with `d` for days (`90s`, `15m`, `7d`). Sizes are
bytes, with an optional `K`, `M` or `G`, which are powers of 1024 (`1K` is 1024 bytes, with an optional `K`, `M` or `G`, which are powers of 1024 (`1K` is 1024
bytes). Rate limits are whole numbers of requests, and byte limits are sizes. bytes). Rate limits and the anomaly thresholds on requests are whole numbers of
requests, and byte limits and the anomaly thresholds on bytes are sizes.
Netblocks are in CIDR form, and a bare address stands for itself alone. Netblocks are in CIDR form, and a bare address stands for itself alone.
Countries are the two-letter codes ISO 3166-1 assigns today, and `xk` for Countries are the two-letter codes ISO 3166-1 assigns today, and `xk` for
Kosovo, in either case (`de` and `DE` are the same); any other code, such as Kosovo, in either case (`de` and `DE` are the same); any other code, such as
@@ -485,15 +534,16 @@ code on both country lists. AS numbers are `AS` and the number, in either case.
Percentages are whole numbers from 0 to 100, and an entry of a list of them is Percentages are whole numbers from 0 to 100, and an entry of a list of them is
an AS number or a country, `:` and a percentage; an AS number or a country an AS number or a country, `:` and a percentage; an AS number or a country
listed twice in one of them stops the start. `off` switches a timeout, a size listed twice in one of them stops the start. `off` switches a timeout, a size
limit, a rate limit, a byte limit, `SWWAF_ALERT_COOLDOWN` or limit, a rate limit, a byte limit, an anomaly threshold, `SWWAF_ALERT_COOLDOWN`
`SWWAF_ALERT_MAX_PER_HOUR` off; `SWWAF_CLIENT_REQUEST_HEADER_MAX_BYTES`, or `SWWAF_ALERT_MAX_PER_HOUR` off; `SWWAF_CLIENT_REQUEST_HEADER_MAX_BYTES`,
`SWWAF_LOOKUP_TIMEOUT`, `SWWAF_UNKNOWN_LIMIT_PERCENT`, the ban settings, the `SWWAF_LOOKUP_TIMEOUT`, `SWWAF_UNKNOWN_LIMIT_PERCENT`, the ban settings, the
state settings, `SWWAF_METRICS_TOP_N` and `SWWAF_LOG_REMOTE_BUFFER` cannot be state settings, `SWWAF_METRICS_TOP_N`, `SWWAF_LOG_REMOTE_BUFFER`,
off. `SWWAF_ANOMALY_NET_V4_PREFIX` and `SWWAF_ANOMALY_NET_V6_PREFIX` cannot be off.
Several limits are fixed rather than settings. At most 20,000 clients are kept, Several limits are fixed rather than settings. At most 20,000 clients are kept,
with their counters and history, and an IPv6 client is counted by its /64. At with their counters and history, and an IPv6 client is counted by its /64. At
most 100,000 answers from GeoJS are kept, for 7 days each. most 100,000 answers from GeoJS are kept, for 7 days each, and at most 20,000
anomaly counters.
### Settings given as files ### Settings given as files
@@ -684,6 +734,9 @@ it, as below. An alert is for one of these events, and is sent when
clear sign of attack. clear sign of attack.
- `permanent_ban`: a permanent ban it makes, or a ban for a clear sign of attack - `permanent_ban`: a permanent ban it makes, or a ban for a clear sign of attack
that a request made permanent. that a request made permanent.
- `anomaly`: a count of requests or bytes over an anomaly threshold, raised by
each request that ends with the count over it, in `observe` mode as in
`enforce` mode. It refuses and bans nothing.
- `source_failure`: GeoJS failing or refusing `smallwebwaf`. - `source_failure`: GeoJS failing or refusing `smallwebwaf`.
- `file_error`: a rule file edited while it runs that has an error, an edit of a - `file_error`: a rule file edited while it runs that has an error, an edit of a
state file set aside as `<name>.bad`, a state file it could not write while state file set aside as `<name>.bad`, a state file it could not write while
@@ -745,19 +798,29 @@ is sent on one line:
- `instance` is `SWWAF_INSTANCE_NAME`, and `time` when the alert was raised, in - `instance` is `SWWAF_INSTANCE_NAME`, and `time` when the alert was raised, in
UTC. UTC.
- `client` is the address of the client whose request raised the alert, and - `client` is the address of the client whose request raised the alert, and
`netblock` the netblock of the ban; both are empty for `source_failure` and `netblock` the netblock of the ban, or for an `anomaly`, the netblock counted:
the client's own, the netblock around it or a named netblock, and none for an
AS number or the whole service; both are empty for `source_failure` and
`file_error`. `asn`, `as_name` and `country` are, for a ban, the client's as `file_error`. `asn`, `as_name` and `country` are, for a ban, the client's as
the ban's notes give them when the alert is raised: empty, as in this alert, the ban's notes give them when the alert is raised: empty, as in this alert,
when GeoJS had not answered about the client by then. when GeoJS had not answered about the client by then; for an `anomaly`, the
- `reason` is a short sentence; for a ban, the ban's `reason` in `bans.json`. client's as the lookup gave them by the time its request ended.
- `reason` is a short sentence; for a ban, the ban's `reason` in `bans.json`;
for an `anomaly`, what was counted over which threshold, such as
`requests per minute of the netblock 203.0.113.0/24 over the threshold of 1000`.
- `detail` is what is particular to the event: for a ban, its `cause`, when it - `detail` is what is particular to the event: for a ban, its `cause`, when it
ends as `ban_expires`, in the form the request log gives it, and its `notes`, ends as `ban_expires`, in the form the request log gives it, and its `notes`,
as `bans.json` gives them; for `source_failure`, the `source`, `geojs`, the as `bans.json` gives them; for an `anomaly`, the `scope`, `client`, `net`,
`error`, and when GeoJS is asked again, `asking_again_in`; for `file_error`, `asn`, `total` or `watch`, as the settings name them, the `asn` counted for
the `file`, which for an edit set aside is the file it was renamed to, and the `asn` and the `name` of the named netblock for `watch`, the `window`, `minute`
`error`, which for a file that does not parse names where in it the error is. or `hour`, the `kind`, `requests` or `bytes`, the `count`, which is weighted
as the rate limits weigh theirs, and the `threshold`; for `source_failure`,
the `source`, `geojs`, the `error`, and when GeoJS is asked again,
`asking_again_in`; for `file_error`, the `file`, which for an edit set aside
is the file it was renamed to, and the `error`, which for a file that does not
parse names where in it the error is.
- `suppressed_repeats` is how many repeats the cooldown held back before this - `suppressed_repeats` is how many repeats the cooldown held back before this
alert. alert, and for a `summary`, those no other alert gives (see below).
Slack and ntfy are each sent the alert as a message: a title, the instance and Slack and ntfy are each sent the alert as a message: a title, the instance and
the event, such as `fsn1app1/gitea: ban`, and a text, the `reason`, then a line the event, such as `fsn1app1/gitea: ban`, and a text, the `reason`, then a line
@@ -802,17 +865,21 @@ and Slack this JSON object, shown indented; it is sent on one line:
An alert for the same event as the last one sent, on the same netblock, or for a An alert for the same event as the last one sent, on the same netblock, or for a
`file_error` about the same file, or for a `source_failure` about the same `file_error` about the same file, or for a `source_failure` about the same
source, less than `SWWAF_ALERT_COOLDOWN` after it, is a repeat: it is held back source, or for an `anomaly` in the same scope, with the same netblock, AS number
and counted, and the next alert sent for them gives that count as or name, whatever its window and kind, less than `SWWAF_ALERT_COOLDOWN` after
`suppressed_repeats`. it, is a repeat: it is held back and counted, and the next alert sent for them
gives that count as `suppressed_repeats`. As each hour of the clock, in UTC,
ends, the cooldowns that have run out are dropped, and the repeats they held
back, which no alert sent since has given, go in that hour's summary.
Past `SWWAF_ALERT_MAX_PER_HOUR` alerts in an hour of the clock, in UTC, the Past `SWWAF_ALERT_MAX_PER_HOUR` alerts in an hour, the hour's other alerts are
hour's other alerts are held back and counted by event. Once the hour has ended, held back and counted by event. An alert held back this way starts no cooldown.
one alert sums them up: its `event` is `summary`, its `reason` says how many Once the hour has ended, one alert sums up the alerts held back and the repeats
were held back, and its `detail` gives the `hour` as when it started, the of the cooldowns dropped: its `event` is `summary`, its `reason` says how many
`count`, and the count for each event, as `events`. An alert held back this way of each were held back, its `detail` gives the `hour` as when it started, the
starts no cooldown, and the repeats held back before it are given by the next `count` of alerts held back, and the count for each event, as `events`, and its
alert sent for the same event and netblock, file or source. `suppressed_repeats` gives the repeats. An hour with neither ends without a
summary.
Each destination has a queue of its own, of at most 1000 alerts, from which they Each destination has a queue of its own, of at most 1000 alerts, from which they
are sent to it one at a time, the oldest first, so a destination that is slow or are sent to it one at a time, the oldest first, so a destination that is slow or
@@ -833,7 +900,8 @@ restart the alerts waiting are sent, and the cooldowns go on.
`smallwebwaf` keeps its state in memory and a copy of it in four JSON files in `smallwebwaf` keeps its state in memory and a copy of it in four JSON files in
`SWWAF_STATE_DIR`, `/var/lib/smallwebwaf` by default, as "Persistent state" in `SWWAF_STATE_DIR`, `/var/lib/smallwebwaf` by default, as "Persistent state" in
[`SPEC.md`](SPEC.md) describes. Each has a top-level `version`, 1, and lists its [`SPEC.md`](SPEC.md) describes. Each has a top-level `version`, 1, and lists its
entries by client address, but for the alerts waiting, with times in UTC. entries by client address, but for the alerts waiting, and the anomaly counters,
which are listed by scope first, with times in UTC.
- `bans.json`: every ban with its notes, indented to be read. A permanent ban's - `bans.json`: every ban with its notes, indented to be read. A permanent ban's
`expires` is `null`. A ban's `cause` is `limit` for a broken rate limit or `expires` is `null`. A ban's `cause` is `limit` for a broken rate limit or
@@ -861,16 +929,23 @@ entries by client address, but for the alerts waiting, with times in UTC.
number, AS name and country, when GeoJS gave it and when it was last used. number, AS name and country, when GeoJS gave it and when it was last used.
- `alerts.json`: the state of the alerts (see "Alerts" above), indented to be - `alerts.json`: the state of the alerts (see "Alerts" above), indented to be
read: under `cooldowns`, for each event and netblock, or event and `file` or read: under `cooldowns`, for each event and netblock, or event and `file` or
`source`, or event alone, when the last alert was sent, `sent`, and the `source`, or for an `anomaly`, its `scope` with its `netblock`, `asn` or
repeats held back since, `suppressed_repeats`; under `hour`, the hour under `name`, or event alone, when the last alert was sent, `sent`, and the repeats
way, from its `start`, the alerts `sent` in it and those `held_back` for its held back since, `suppressed_repeats`; under `hour`, the hour under way, from
summary, by event; and under `waiting`, for each destination you name, its `start`, the alerts `sent` in it and those `held_back` for its summary, by
`webhook`, `slack` or `ntfy`, the alerts still waiting to be sent to it, the event; under `waiting`, for each destination you name, `webhook`, `slack` or
oldest first, each as the webhook is sent it. As an hour ends, the cooldowns `ntfy`, the alerts still waiting to be sent to it, the oldest first, each as
that have run out with no repeat held back are dropped. As the file is read, the webhook is sent it; and under `anomaly_counters`, each anomaly counter:
the alerts waiting for a destination you no longer name are dropped. A file its `scope`, as an `anomaly` alert names it, with the `netblock` of a client,
whose `waiting` is a list, as it was before alerts went to Slack and ntfy too, of a netblock around a client or of a named netblock, the `asn` of an AS
stops the start: put the list under `"webhook"`, or remove the file. number and the `name` of a named netblock, and its two buckets of requests in
the minute and the hour, `minute` and `hour`, and of bytes, `minute_bytes` and
`hour_bytes`, each left out while it is empty. As an hour ends, the cooldowns
that have run out are dropped, and the hour's summary gives the repeats they
held back. As the file is read, the alerts waiting for a destination you no
longer name are dropped. A file whose `waiting` is a list, as it was before
alerts went to Slack and ntfy too, stops the start: put the list under
`"webhook"`, or remove the file.
`bans.json` is written `SWWAF_STATE_WRITE_DELAY` after a ban is made, lifted `bans.json` is written `SWWAF_STATE_WRITE_DELAY` after a ban is made, lifted
through `DELETE /_smallwebwaf/bans/<client>`, or made permanent, with every such through `DELETE /_smallwebwaf/bans/<client>`, or made permanent, with every such
@@ -882,22 +957,27 @@ whole. A write that fails is logged, raised as a `file_error` alert while
changed since the last write. changed since the last write.
At start the files are read back: each client keeps its counts, so a restart At start the files are read back: each client keeps its counts, so a restart
gives it no fresh allowance, and each ban keeps refusing every client in its gives it no fresh allowance, each anomaly counter keeps its counts, and each ban
netblock until it ends, even after `SWWAF_BAN_SCOPE_V4_PREFIX` has changed. A keeps refusing every client in its netblock until it ends, even after
netblock whose address has bits past its length, such as `203.0.113.9/24`, is `SWWAF_BAN_SCOPE_V4_PREFIX` has changed. A netblock whose address has bits past
read as the netblock it is in, `203.0.113.0/24`. Buckets and answers whose time its length, such as `203.0.113.9/24`, is read as the netblock it is in,
has passed are dropped. A missing file is empty state, as on a first start. A `203.0.113.0/24`. Buckets and answers whose time has passed are dropped, and so
file that does not parse, or has another `version`, stops the start with a is an anomaly counter left with no bucket. A missing file is empty state, as on
message naming the file, and the line and column where Go's JSON decoder gives a first start. A file that does not parse, or has another `version`, stops the
them; so does a state directory `smallwebwaf` cannot write. So does an entry start with a message naming the file, and the line and column where Go's JSON
without a field it needs, named with the entry's place in the file: a ban's decoder gives them; so does a state directory `smallwebwaf` cannot write. So
`netblock`, `start` or `expires`, which is `null` for a permanent ban; a does an entry without a field it needs, named with the entry's place in the
client's `client`, or the `start` of a window in which it has requests or bytes; file: a ban's `netblock`, `start` or `expires`, which is `null` for a permanent
an answer's `client`, `country`, which is `""` for a client GeoJS cannot place, ban; a client's `client`, or the `start` of a window in which it has requests or
or `answered`; a cooldown's `event` or `sent`; an alert waiting's `event` or bytes; an answer's `client`, `country`, which is `""` for a client GeoJS cannot
`time`. So does a ban whose `cause` is not `limit`, `attack` or `admin`, and place, or `answered`; a cooldown's `event` or `sent`; an alert waiting's `event`
alerts waiting for a destination that is not `webhook`, `slack` or `ntfy`. An or `time`; an anomaly counter's `netblock`, unless it counts an AS number or the
answer's `asn` or `as_name` left out reads as empty. whole service, its `asn`, for an AS number, its `name`, for a named netblock, or
the `start` of a window in which it has requests or bytes. So does a ban whose
`cause` is not `limit`, `attack` or `admin`, alerts waiting for a destination
that is not `webhook`, `slack` or `ntfy`, and an anomaly counter whose `scope`
is not `client`, `net`, `asn`, `total` or `watch`. An answer's `asn` or
`as_name` left out reads as empty.
While it runs, `smallwebwaf` watches `SWWAF_STATE_DIR` and takes in your edit of While it runs, `smallwebwaf` watches `SWWAF_STATE_DIR` and takes in your edit of
a state file as soon as you save it: what the file then holds replaces what a state file as soon as you save it: what the file then holds replaces what
@@ -907,13 +987,14 @@ writes a file it takes in any edit made since, so your edit is not overwritten;
a change `smallwebwaf` made after you opened the file, such as a new ban, is a change `smallwebwaf` made after you opened the file, such as a new ban, is
lost when you save over it. An edit that would stop the start, because it does lost when you save over it. An edit that would stop the start, because it does
not parse, has another `version`, leaves out a field an entry needs, gives a ban not parse, has another `version`, leaves out a field an entry needs, gives a ban
another `cause` or names another destination, does not stop the running another `cause`, names another destination or gives an anomaly counter another
`smallwebwaf`: it keeps what it holds, and at the file's next write renames your `scope`, does not stop the running `smallwebwaf`: it keeps what it holds, and at
file to `<name>.bad`, such as `bans.json.bad`, writes the file again from the file's next write renames your file to `<name>.bad`, such as
memory, logs the file and where the error is, and raises a `file_error` alert `bans.json.bad`, writes the file again from memory, logs the file and where the
for it. It waits for that write because an editor's file can be read before the error is, and raises a `file_error` alert for it. It waits for that write
editor has finished writing it. Mend the `.bad` file and move it back. A file because an editor's file can be read before the editor has finished writing it.
you remove is written again at its next write. Mend the `.bad` file and move it back. A file you remove is written again at its
next write.
To ban a netblock, add an entry to `bans.json` with its `netblock`, its `start` To ban a netblock, add an entry to `bans.json` with its `netblock`, its `start`
and its `expires`, `null` for a ban that never ends; its `reason` and its and its `expires`, `null` for a ban that never ends; its `reason` and its
@@ -1304,8 +1385,8 @@ For each request `smallwebwaf`:
- forwards it to the app and streams the response back, within the size and time - forwards it to the app and streams the response back, within the size and time
limits; limits;
- counts the bytes and any refusal by the rule files or the Core Rule Set, bans - counts the bytes and any refusal by the rule files or the Core Rule Set, bans
the client if it broke a limit, updates its history, sends any alerts that are the client if it broke a limit, updates its history and the anomaly counters,
due, and writes the log line. sends any alerts that are due, and writes the log line.
A minimal deployment is the app's own Dockerfile, built on the `smallwebwaf` A minimal deployment is the app's own Dockerfile, built on the `smallwebwaf`
image, with no setting. That image is built on Ubuntu 26.04 LTS, the newest image, with no setting. That image is built on Ubuntu 26.04 LTS, the newest
@@ -1394,14 +1475,15 @@ the metrics, failure behaviour and the build order.
`smallwebwaf` looks up the AS number and country of every client through GeoJS, `smallwebwaf` looks up the AS number and country of every client through GeoJS,
a free web service that needs no account and no file, for the request log, the a free web service that needs no account and no file, for the request log, the
client's history, the notes of its bans, their alerts and the metrics, and for client's history, the notes of its bans, their alerts and the metrics, and for
the country lists when you set them. This means that GeoJS is told the address the country lists and the anomaly thresholds per AS number when you set them.
of every new visitor, whether or not a setting uses the answer, unless you set This means that GeoJS is told the address of every new visitor, whether or not a
`SWWAF_LOOKUP_SOURCE=off`. The only visitors it is not told about are those in setting uses the answer, unless you set `SWWAF_LOOKUP_SOURCE=off`. The only
`SWWAF_ALLOW_NETS` or `SWWAF_DENY_NETS`, those whose netblock a ban covers, and visitors it is not told about are those in `SWWAF_ALLOW_NETS` or
those on a private, loopback or link-local address. An IPv6 visitor is asked `SWWAF_DENY_NETS`, those whose netblock a ban covers, and those on a private,
about by the first address of its /64. Each answer is kept for seven days, in loopback or link-local address. An IPv6 visitor is asked about by the first
memory and in `lookups.json`, so that it survives a restart, and a visitor whose address of its /64. Each answer is kept for seven days, in memory and in
answer is kept is not asked about again. `lookups.json`, so that it survives a restart, and a visitor whose answer is
kept is not asked about again.
A request waits for its client's first answer only while a setting acts on it A request waits for its client's first answer only while a setting acts on it
before the request goes on: a country list, `SWWAF_ADD_LOOKUP_HEADERS`, or a before the request goes on: a country list, `SWWAF_ADD_LOOKUP_HEADERS`, or a
@@ -1470,7 +1552,9 @@ GeoJS.
what it would have refused for noted in the log line. A request under what it would have refused for noted in the log line. A request under
`/_smallwebwaf/` that `check` lets through is answered by `answerAdmin` `/_smallwebwaf/` that `check` lets through is answered by `answerAdmin`
instead of reaching the app. Once the answer to a request passed to the app instead of reaching the app. Once the answer to a request passed to the app
has ended, `countBytes` counts its bytes for the byte limits. has ended, `countBytes` counts its bytes for the byte limits, and once any
request but the health check has ended, `countAnomalies` counts it for the
anomaly thresholds.
- `internal/metrics`: the metrics, counted as the other parts tell it what - `internal/metrics`: the metrics, counted as the other parts tell it what
happened, and served in the Prometheus text format. happened, and served in the Prometheus text format.
- `internal/bans`: the ban ledger: each netblock's bans with their notes, how - `internal/bans`: the ban ledger: each netblock's bans with their notes, how
@@ -1486,6 +1570,10 @@ GeoJS.
- `internal/ratelimit`: the table of clients: counts each client's requests and - `internal/ratelimit`: the table of clients: counts each client's requests and
bytes, tells when they take it over a rate limit or a byte limit, and keeps bytes, tells when they take it over a rate limit or a byte limit, and keeps
each client's history. each client's history.
- `internal/anomaly`: the anomaly counters: counts each request and its bytes
per client, per netblock around a client, per AS number, for the whole service
and per named netblock, in the buckets `internal/ratelimit` counts in, and
raises an `anomaly` alert for a count over its threshold.
- `internal/state`: reads the state files at start, takes in an admin's edit of - `internal/state`: reads the state files at start, takes in an admin's edit of
one while running, and writes them when they are due and at the stop. one while running, and writes them when they are due and at the stop.
- `internal/requestlog`: the lines on stdout: the request log line and the - `internal/requestlog`: the lines on stdout: the request log line and the
@@ -1506,8 +1594,9 @@ GeoJS.
Besides the Go standard library, `github.com/hashicorp/golang-lru/v2` keeps the Besides the Go standard library, `github.com/hashicorp/golang-lru/v2` keeps the
table of clients to 20,000 and the GeoJS answers to 100,000, dropping the least table of clients to 20,000 and the GeoJS answers to 100,000, dropping the least
recently seen, and the banned netblocks in the order they were last seen, from recently seen, the anomaly counters to 20,000, dropping the one counted least
which the ledger picks the ban to drop past `SWWAF_MAX_BANS`, and recently, and the banned netblocks in the order they were last seen, from which
the ledger picks the ban to drop past `SWWAF_MAX_BANS`, and
`github.com/prometheus/client_golang` keeps the metrics and serves them, and `github.com/prometheus/client_golang` keeps the metrics and serves them, and
`github.com/fsnotify/fsnotify` tells `smallwebwaf` when a state file or a rule `github.com/fsnotify/fsnotify` tells `smallwebwaf` when a state file or a rule
file is saved, or the lookup database replaced, and file is saved, or the lookup database replaced, and
+86 -43
View File
@@ -1,5 +1,6 @@
// Package alerts sends alerts on bans, on a source that fails and on a // Package alerts sends alerts on bans, on traffic over an anomaly
// file with an error to each destination set: to the webhook // threshold, on a source that fails and on a file with an error to each
// destination set: to the webhook
// SWWAF_ALERT_WEBHOOK_URL names, each as one JSON object, as the "Alert // SWWAF_ALERT_WEBHOOK_URL names, each as one JSON object, as the "Alert
// webhook schema" section of SPEC.md describes, to the Slack incoming // webhook schema" section of SPEC.md describes, to the Slack incoming
// webhook SWWAF_ALERT_SLACK_WEBHOOK_URL names, as a message, and to the // webhook SWWAF_ALERT_SLACK_WEBHOOK_URL names, as a message, and to the
@@ -40,11 +41,12 @@ const (
// EventPermanentBan is a permanent ban smallwebwaf made, or a ban it // EventPermanentBan is a permanent ban smallwebwaf made, or a ban it
// made permanent. // made permanent.
EventPermanentBan = "permanent_ban" EventPermanentBan = "permanent_ban"
// EventWAFBlock, EventAnomaly and EventReputationHit come with the // EventAnomaly is a count of requests or bytes over an anomaly
// Core Rule Set, the anomaly thresholds and the reputation sources; // threshold.
// nothing raises them yet. EventAnomaly = "anomaly"
// EventWAFBlock and EventReputationHit come with the Core Rule Set and
// the reputation sources; nothing raises them yet.
EventWAFBlock = "waf_block" EventWAFBlock = "waf_block"
EventAnomaly = "anomaly"
EventReputationHit = "reputation_hit" EventReputationHit = "reputation_hit"
// EventSourceFailure is GeoJS failing or refusing smallwebwaf. // EventSourceFailure is GeoJS failing or refusing smallwebwaf.
EventSourceFailure = "source_failure" EventSourceFailure = "source_failure"
@@ -52,8 +54,11 @@ const (
// runs that does not parse, a replacement of the lookup database that // runs that does not parse, a replacement of the lookup database that
// cannot be read, or a state file that cannot be written. // cannot be read, or a state file that cannot be written.
EventFileError = "file_error" EventFileError = "file_error"
// EventSummary is the summary of the alerts an hour held back past // EventSummary is the summary sent as an hour ends: of the alerts held
// SWWAF_ALERT_MAX_PER_HOUR. SWWAF_ALERT_EVENTS does not name it. // back in it past SWWAF_ALERT_MAX_PER_HOUR, and of the repeats held
// back by the cooldowns dropped as it ends, which no alert let through
// has given. It is sent with SWWAF_ALERT_MAX_PER_HOUR off too, for
// those repeats. SWWAF_ALERT_EVENTS does not name it.
EventSummary = "summary" EventSummary = "summary"
) )
@@ -153,18 +158,22 @@ type Alert struct {
ASName string `json:"as_name"` ASName string `json:"as_name"`
Country string `json:"country"` Country string `json:"country"`
// Reason is a short sentence, and Detail what is particular to the // Reason is a short sentence, and Detail what is particular to the
// event: for a file_error, its "file", and for a source_failure, its // event: for a file_error, its "file", for a source_failure, its
// "source", which the cooldown tells repeats by. // "source", and for an anomaly, its "scope", with the "asn" or the
// "name" of some scopes, which the cooldown tells repeats by.
Reason string `json:"reason"` Reason string `json:"reason"`
Detail map[string]any `json:"detail"` Detail map[string]any `json:"detail"`
// SuppressedRepeats is how many repeats of the alert the cooldown // SuppressedRepeats is how many repeats of the alert the cooldown
// held back since the last one let through. // held back since the last one let through. For a summary, it is how
// many the cooldowns dropped as the hour ended had held back that no
// alert let through gave.
SuppressedRepeats int `json:"suppressed_repeats"` SuppressedRepeats int `json:"suppressed_repeats"`
} }
// Cooldown is, for an event on a netblock, or about a file or a source, // Cooldown is, for an event on a netblock, about a file or a source, or
// when the last alert let through was raised, and how many repeats the // for an anomaly in a scope, when the last alert let through was raised,
// cooldown has held back since, as alerts.json holds it. // and how many repeats the cooldown has held back since, as alerts.json
// holds it.
// //
//nolint:tagliatelle // the state files use snake_case, as the request log does //nolint:tagliatelle // the state files use snake_case, as the request log does
type Cooldown struct { type Cooldown struct {
@@ -172,6 +181,9 @@ type Cooldown struct {
Netblock netip.Prefix `json:"netblock"` Netblock netip.Prefix `json:"netblock"`
File string `json:"file,omitempty"` File string `json:"file,omitempty"`
Source string `json:"source,omitempty"` Source string `json:"source,omitempty"`
Scope string `json:"scope,omitempty"`
ASN string `json:"asn,omitempty"`
Name string `json:"name,omitempty"`
Sent time.Time `json:"sent"` Sent time.Time `json:"sent"`
SuppressedRepeats int `json:"suppressed_repeats"` SuppressedRepeats int `json:"suppressed_repeats"`
} }
@@ -214,7 +226,7 @@ type Queue struct {
mu sync.Mutex mu sync.Mutex
// cooldowns are the alerts last let through, by event and netblock, // cooldowns are the alerts last let through, by event and netblock,
// file or source. // file, source or scope.
cooldowns map[cooldownKey]*Cooldown cooldowns map[cooldownKey]*Cooldown
hour Hour hour Hour
@@ -248,21 +260,28 @@ type destination struct {
} }
// cooldownKey is what makes an alert a repeat of another: the same event // cooldownKey is what makes an alert a repeat of another: the same event
// on the same netblock, and about the same file or source, as its detail // on the same netblock, and about the same file or source, or in the same
// names them. Each is empty for an alert without one. // scope with the same AS number or name, as its detail names them. Each
// is empty for an alert without one.
type cooldownKey struct { type cooldownKey struct {
event string event string
netblock netip.Prefix netblock netip.Prefix
file string file string
source string source string
scope string
asn string
name string
} }
// cooldownKeyOf returns what makes another alert a repeat of alert. // cooldownKeyOf returns what makes another alert a repeat of alert.
func cooldownKeyOf(alert *Alert) cooldownKey { func cooldownKeyOf(alert *Alert) cooldownKey {
file, _ := alert.Detail["file"].(string) file, _ := alert.Detail["file"].(string)
source, _ := alert.Detail["source"].(string) source, _ := alert.Detail["source"].(string)
scope, _ := alert.Detail["scope"].(string)
asn, _ := alert.Detail["asn"].(string)
name, _ := alert.Detail["name"].(string)
return cooldownKey{alert.Event, alert.Netblock, file, source} return cooldownKey{alert.Event, alert.Netblock, file, source, scope, asn, name}
} }
// New returns a Queue with no alert yet. // New returns a Queue with no alert yet.
@@ -295,14 +314,14 @@ func New(params Params) *Queue {
// unless no destination is set or SWWAF_ALERT_EVENTS leaves its event // unless no destination is set or SWWAF_ALERT_EVENTS leaves its event
// out. It gives alert the instance and the time. An alert that repeats // out. It gives alert the instance and the time. An alert that repeats
// the last one let through less than Cooldown before is held back and // the last one let through less than Cooldown before is held back and
// counted, and the next one let through gives that count. Past MaxPerHour // counted. The next one let through gives that count, unless an hour of
// alerts let through in the hour under way, by the clock, an alert is // the clock ends first after the cooldown has run out: the cooldown is
// held back for that hour's summary instead, which is sent once the hour // then dropped, and that hour's summary gives the count. Past MaxPerHour
// has ended; it starts no cooldown, and the repeats held back before it // alerts let through in the hour under way, an alert is held back for
// are given by the next alert let through. Raise never waits: an alert // that hour's summary instead, which is sent once the hour has ended; it
// let through joins the queue of each destination, from which Run sends // starts no cooldown. Raise never waits: an alert let through joins the
// it, and with queueSize alerts waiting for a destination, the oldest is // queue of each destination, from which Run sends it, and with queueSize
// dropped. // alerts waiting for a destination, the oldest is dropped.
func (q *Queue) Raise(alert Alert) { func (q *Queue) Raise(alert Alert) {
if len(q.destinations) == 0 || !slices.Contains(q.params.Events, alert.Event) { if len(q.destinations) == 0 || !slices.Contains(q.params.Events, alert.Event) {
return return
@@ -425,7 +444,8 @@ func (q *Queue) Suppressed() int64 {
} }
// Snapshot returns the queue's state, as alerts.json holds it, with the // Snapshot returns the queue's state, as alerts.json holds it, with the
// cooldowns sorted by netblock, then by event, file and source. // cooldowns sorted by netblock, then by event, file, source, scope, AS
// number and name.
func (q *Queue) Snapshot() State { func (q *Queue) Snapshot() State {
q.mu.Lock() q.mu.Lock()
defer q.mu.Unlock() defer q.mu.Unlock()
@@ -443,7 +463,9 @@ func (q *Queue) Snapshot() State {
slices.SortFunc(state.Cooldowns, func(a, b Cooldown) int { slices.SortFunc(state.Cooldowns, func(a, b Cooldown) int {
return cmp.Or(a.Netblock.Compare(b.Netblock), cmp.Compare(a.Event, b.Event), return cmp.Or(a.Netblock.Compare(b.Netblock), cmp.Compare(a.Event, b.Event),
cmp.Compare(a.File, b.File), cmp.Compare(a.Source, b.Source)) cmp.Compare(a.File, b.File), cmp.Compare(a.Source, b.Source),
cmp.Compare(a.Scope, b.Scope), cmp.Compare(a.ASN, b.ASN),
cmp.Compare(a.Name, b.Name))
}) })
for _, d := range q.destinations { for _, d := range q.destinations {
@@ -466,7 +488,10 @@ func (q *Queue) Load(state State) {
for _, cooldown := range state.Cooldowns { for _, cooldown := range state.Cooldowns {
cooldown.Netblock = cooldown.Netblock.Masked() cooldown.Netblock = cooldown.Netblock.Masked()
key := cooldownKey{cooldown.Event, cooldown.Netblock, cooldown.File, cooldown.Source} key := cooldownKey{
cooldown.Event, cooldown.Netblock, cooldown.File, cooldown.Source,
cooldown.Scope, cooldown.ASN, cooldown.Name,
}
q.cooldowns[key] = &cooldown q.cooldowns[key] = &cooldown
} }
@@ -538,46 +563,64 @@ func (q *Queue) startCooldown(alert *Alert, now time.Time) {
q.cooldowns[key] = &Cooldown{ q.cooldowns[key] = &Cooldown{
Event: alert.Event, Netblock: alert.Netblock, File: key.file, Source: key.source, Event: alert.Event, Netblock: alert.Netblock, File: key.file, Source: key.source,
Sent: now, Scope: key.scope, ASN: key.asn, Name: key.name, Sent: now,
} }
} }
// endHour ends the hour under way, if now is past it: it queues that // endHour ends the hour under way, if now is past it. It drops the
// hour's summary when alerts were held back in it past MaxPerHour, and // cooldowns that have run out, whatever repeats they held back, so that
// forgets the cooldowns that have run out with no repeat held back, which // they do not pile up, and queues that hour's summary when alerts were
// no alert needs any more. // held back in it past MaxPerHour, or when a cooldown dropped had held
// back repeats, which no alert let through has given: the summary gives
// them.
func (q *Queue) endHour(now time.Time) { func (q *Queue) endHour(now time.Time) {
start := now.Truncate(time.Hour) start := now.Truncate(time.Hour)
if !start.After(q.hour.Start) { if !start.After(q.hour.Start) {
return return
} }
repeats := 0
for key, cooldown := range q.cooldowns {
if now.Sub(cooldown.Sent) >= q.params.Cooldown {
repeats += cooldown.SuppressedRepeats
delete(q.cooldowns, key)
}
}
heldBack := 0 heldBack := 0
for _, count := range q.hour.HeldBack { for _, count := range q.hour.HeldBack {
heldBack += count heldBack += count
} }
var reasons []string
if heldBack > 0 { if heldBack > 0 {
reasons = append(reasons, fmt.Sprintf("%d alerts held back in the hour from %s, "+
"past the %d an hour SWWAF_ALERT_MAX_PER_HOUR allows", heldBack,
q.hour.Start.Format(time.RFC3339), q.params.MaxPerHour))
}
if repeats > 0 {
reasons = append(reasons, fmt.Sprintf("%d repeats held back by "+
"SWWAF_ALERT_COOLDOWN that no later alert gives", repeats))
}
if len(reasons) > 0 {
q.queue(&Alert{ q.queue(&Alert{
Instance: q.params.Instance, Instance: q.params.Instance,
Time: now, Time: now,
Event: EventSummary, Event: EventSummary,
Reason: fmt.Sprintf("%d alerts held back in the hour from %s, past the %d "+ Reason: strings.Join(reasons, "; "),
"an hour SWWAF_ALERT_MAX_PER_HOUR allows", heldBack,
q.hour.Start.Format(time.RFC3339), q.params.MaxPerHour),
Detail: map[string]any{ Detail: map[string]any{
"hour": q.hour.Start, "count": heldBack, "events": q.hour.HeldBack, "hour": q.hour.Start, "count": heldBack, "events": q.hour.HeldBack,
}, },
SuppressedRepeats: repeats,
}) })
} }
q.hour = Hour{Start: start, HeldBack: map[string]int{}} q.hour = Hour{Start: start, HeldBack: map[string]int{}}
for key, cooldown := range q.cooldowns {
if now.Sub(cooldown.Sent) >= q.params.Cooldown && cooldown.SuppressedRepeats == 0 {
delete(q.cooldowns, key)
}
}
} }
// queue adds alert to the alerts waiting for each destination. // queue adds alert to the alerts waiting for each destination.
+87 -25
View File
@@ -381,7 +381,7 @@ func TestAlertsPastTheHourlyLimitAreRolledIntoOneSummary(t *testing.T) {
}) })
} }
func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheNextSent(t *testing.T) { func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheSummary(t *testing.T) {
t.Parallel() t.Parallel()
synctest.Test(t, func(t *testing.T) { synctest.Test(t, func(t *testing.T) {
@@ -401,8 +401,8 @@ func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheNextSent(t *testing.
time.Sleep(cooldown) time.Sleep(cooldown)
raise() raise()
// The next hour's first alert gives the two repeats, and the summary // The summary gives the alert past the limit and the two repeats,
// the alert past the limit. // and the next hour's first alert none.
time.Sleep(time.Hour - cooldown) time.Sleep(time.Hour - cooldown)
synctest.Wait() synctest.Wait()
raise() raise()
@@ -413,11 +413,14 @@ func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheNextSent(t *testing.
got := webhook.received() got := webhook.received()
if len(got) == 3 { if len(got) == 3 {
detail, _ := got[1].alert["detail"].(map[string]any) detail, _ := got[1].alert["detail"].(map[string]any)
repeats := got[2].alert["suppressed_repeats"] summaryRepeats := got[1].alert["suppressed_repeats"]
lastRepeats := got[2].alert["suppressed_repeats"]
if detail["count"] != float64(1) || repeats != float64(2) { if detail["count"] != float64(1) || summaryRepeats != float64(2) ||
t.Errorf("the summary counts %v alerts, and the last alert gives %v "+ lastRepeats != float64(0) {
"repeats, want 1 and 2", detail["count"], repeats) t.Errorf("the summary counts %v alerts and %v repeats, and the last "+
"alert gives %v repeats, want 1, 2 and 0", detail["count"],
summaryRepeats, lastRepeats)
} }
} }
@@ -425,6 +428,70 @@ func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheNextSent(t *testing.
}) })
} }
func TestCooldownsThatHaveRunOutAreDroppedAndTheirRepeatsSummedUp(t *testing.T) {
t.Parallel()
for name, maxPerHour := range map[string]int{"limit off": 0, "limit set": 60} {
t.Run(name, func(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := newParams()
params.MaxPerHour = maxPerHour
webhook, q := start(t, params)
// For four hours, a netblock of its own each minute is over an
// anomaly threshold twice: an alert, and a repeat the cooldown
// holds back.
netblocks := 0
for range 4 {
for range 60 {
anomaly := alerts.Alert{
Event: alerts.EventAnomaly, Netblock: netblock(netblocks),
Detail: map[string]any{"scope": "net"},
}
q.Raise(anomaly)
q.Raise(anomaly)
netblocks++
time.Sleep(time.Minute)
}
// As the hour ends, only the cooldowns started less than the
// cooldown before are kept, in memory and for alerts.json.
synctest.Wait()
kept := len(q.Snapshot().Cooldowns)
if kept > int(cooldown/time.Minute) {
t.Errorf("after %d netblocks, %d cooldowns are kept, want at most %d",
netblocks, kept, int(cooldown/time.Minute))
}
}
// An hour on, every cooldown has been dropped, and the summaries
// have given every repeat.
time.Sleep(time.Hour)
synctest.Wait()
repeats := 0.0
for _, request := range webhook.received() {
count, _ := request.alert["suppressed_repeats"].(float64)
repeats += count
}
kept := len(q.Snapshot().Cooldowns)
if kept != 0 || repeats != float64(netblocks) {
t.Errorf("%d cooldowns are kept and the webhook was given %v repeats, "+
"want 0 and %d", kept, repeats, netblocks)
}
})
})
}
}
func TestFailedRequestIsSentAgainWithBackoff(t *testing.T) { func TestFailedRequestIsSentAgainWithBackoff(t *testing.T) {
t.Parallel() t.Parallel()
@@ -635,7 +702,8 @@ func TestStateLoadedIntoANewQueueCarriesOn(t *testing.T) {
after.Load(roundTrip(t, before.Snapshot())) after.Load(roundTrip(t, before.Snapshot()))
// The new queue sends the alert waiting, holds back the repeat as // The new queue sends the alert waiting, holds back the repeat as
// the cooldown still runs, and sends the summary of the hour. // the cooldown still runs, and sends the summary of the hour, which
// gives both repeats, as the cooldown has run out.
after.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)}) after.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
synctest.Wait() synctest.Wait()
wantEvents(t, webhook, alerts.EventBan) wantEvents(t, webhook, alerts.EventBan)
@@ -644,18 +712,12 @@ func TestStateLoadedIntoANewQueueCarriesOn(t *testing.T) {
synctest.Wait() synctest.Wait()
wantEvents(t, webhook, alerts.EventBan, alerts.EventSummary) wantEvents(t, webhook, alerts.EventBan, alerts.EventSummary)
detail, _ := webhook.received()[1].alert["detail"].(map[string]any) summary := webhook.received()[1].alert
if detail["count"] != float64(1) { detail, _ := summary["detail"].(map[string]any)
t.Errorf("the summary counts %v alerts, want 1", detail["count"])
}
// The cooldown has run out, and the next one gives both repeats. if detail["count"] != float64(1) || summary["suppressed_repeats"] != float64(2) {
after.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)}) t.Errorf("the summary counts %v alerts and %v repeats, want 1 and 2",
synctest.Wait() detail["count"], summary["suppressed_repeats"])
got := webhook.received()
if repeats := got[len(got)-1].alert["suppressed_repeats"]; repeats != float64(2) {
t.Errorf("the last alert gives %v repeats, want 2", repeats)
} }
}) })
} }
@@ -701,8 +763,8 @@ func TestSlackAndNtfyAreSentTheSummaryAndTheRepeatsHeldBack(t *testing.T) {
ban := alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1), Reason: "a ban"} ban := alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1), Reason: "a ban"}
// The hour's one alert, a repeat of it the cooldown holds back, and // The hour's one alert, a repeat of it the cooldown holds back, and
// an alert past the limit; once the hour has ended, its summary, and // an alert past the limit; once the hour has ended, its summary,
// the next alert, which gives the repeat. // which gives the repeat, and the next alert, which gives none.
q.Raise(ban) q.Raise(ban)
q.Raise(ban) q.Raise(ban)
q.Raise(alerts.Alert{Event: alerts.EventFileError, Reason: "a file error"}) q.Raise(alerts.Alert{Event: alerts.EventFileError, Reason: "a file error"})
@@ -720,14 +782,14 @@ func TestSlackAndNtfyAreSentTheSummaryAndTheRepeatsHeldBack(t *testing.T) {
} }
const summary = "1 alerts held back in the hour from 2000-01-01T00:00:00Z, " + const summary = "1 alerts held back in the hour from 2000-01-01T00:00:00Z, " +
"past the 1 an hour SWWAF_ALERT_MAX_PER_HOUR allows" "past the 1 an hour SWWAF_ALERT_MAX_PER_HOUR allows; 1 repeats held back " +
"by SWWAF_ALERT_COOLDOWN that no later alert gives\nsuppressed repeats: 1"
wantSlackMessage(t, slack[1], "*"+instance+": summary*\n"+summary) wantSlackMessage(t, slack[1], "*"+instance+": summary*\n"+summary)
wantNtfyMessage(t, ntfy[1], instance+": summary", "default bar_chart", summary) wantNtfyMessage(t, ntfy[1], instance+": summary", "default bar_chart", summary)
wantSlackMessage(t, slack[2], wantSlackMessage(t, slack[2], "*"+instance+": ban*\na ban\nnetblock: 203.0.113.1/32")
"*"+instance+": ban*\na ban\nnetblock: 203.0.113.1/32\nsuppressed repeats: 1")
wantNtfyMessage(t, ntfy[2], instance+": ban", "default no_entry", wantNtfyMessage(t, ntfy[2], instance+": ban", "default no_entry",
"a ban\nnetblock: 203.0.113.1/32\nsuppressed repeats: 1") "a ban\nnetblock: 203.0.113.1/32")
}) })
} }
+406
View File
@@ -0,0 +1,406 @@
// Package anomaly counts requests and bytes over a minute and an hour, per
// client, per surrounding netblock, per AS number, for the whole service
// and per named netblock, and raises an anomaly alert for a count over its
// threshold, as "Anomaly thresholds" under "Configuration surface" in
// SPEC.md describes. It refuses and bans nothing. At most 20,000 counters
// are kept, in memory, and written to alerts.json and read from it by the
// state package.
package anomaly
import (
"cmp"
"fmt"
"net/netip"
"slices"
"sync"
"time"
"github.com/hashicorp/golang-lru/v2/simplelru"
"sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/ratelimit"
)
// maxCounters is how many counters are kept. Past it, the counter counted
// least recently is dropped, and starts afresh if it is counted again.
const maxCounters = 20000
// The scopes, what a counter counts, as the settings, alerts.json and the
// alerts name them.
const (
// ScopeClient is one client: an IPv4 address, or an IPv6 /64.
ScopeClient = "client"
// ScopeNet is the netblock around a client, SWWAF_ANOMALY_NET_V4_PREFIX
// or SWWAF_ANOMALY_NET_V6_PREFIX long.
ScopeNet = "net"
// ScopeASN is an AS number.
ScopeASN = "asn"
// ScopeTotal is the whole service.
ScopeTotal = "total"
// ScopeWatch is a named netblock of SWWAF_WATCH_NETS.
ScopeWatch = "watch"
)
// Scopes returns every scope.
func Scopes() []string {
return []string{ScopeClient, ScopeNet, ScopeASN, ScopeTotal, ScopeWatch}
}
// The windows a counter counts in, as the alerts name them.
const (
minute = "minute"
hour = "hour"
)
// Thresholds are the most requests and the most bytes a scope may have
// counted in a minute and in an hour before an alert is raised. Zero is
// off.
type Thresholds struct {
RequestsPerMinute int64
RequestsPerHour int64
BytesPerMinute int64
BytesPerHour int64
}
// NamedNetblock is a netblock SWWAF_WATCH_NETS names.
type NamedNetblock struct {
Name string
Netblock netip.Prefix
}
// Params are what New needs.
type Params struct {
// The thresholds of each scope: SWWAF_ANOMALY_CLIENT_*,
// SWWAF_ANOMALY_NET_*, SWWAF_ANOMALY_ASN_*, SWWAF_ANOMALY_TOTAL_* and
// SWWAF_WATCH_*.
Client, Net, ASN, Total, Watch Thresholds
// NetV4Prefix and NetV6Prefix are the lengths of the netblock around a
// client (SWWAF_ANOMALY_NET_V4_PREFIX and SWWAF_ANOMALY_NET_V6_PREFIX).
NetV4Prefix, NetV6Prefix int
// NamedNetblocks are SWWAF_WATCH_NETS.
NamedNetblocks []NamedNetblock
// Alerts receive the anomaly alerts.
Alerts *alerts.Queue
}
// Counter is one scope's counts, as alerts.json holds them: the scope,
// with the netblock, the AS number or the name that tells it from the
// others in that scope, and its two buckets of requests and of bytes in
// the minute and in the hour. A bucket whose threshold is off counts
// nothing, and is left out.
//
//nolint:tagliatelle // the state files use snake_case, as the request log does
type Counter struct {
Scope string `json:"scope"`
Netblock netip.Prefix `json:"netblock,omitzero"`
ASN string `json:"asn,omitempty"`
Name string `json:"name,omitempty"`
Minute ratelimit.Buckets `json:"minute,omitzero"`
Hour ratelimit.Buckets `json:"hour,omitzero"`
MinuteBytes ratelimit.Buckets `json:"minute_bytes,omitzero"`
HourBytes ratelimit.Buckets `json:"hour_bytes,omitzero"`
}
// Request is a request that has ended, as the counters count it.
type Request struct {
// Client is the client's address, and ClientGroup the client it is
// counted as: its IPv4 address, or its IPv6 /64.
Client netip.Addr
ClientGroup netip.Prefix
// ASN, ASName and Country are the client's as looked up, each "" when
// unknown.
ASN, ASName, Country string
// Bytes are the request's bytes, as SWWAF_BYTES_COUNT counts them.
Bytes int64
}
// Counters counts each request in the scopes it is in. It is safe for
// concurrent use.
type Counters struct {
params Params
mu sync.Mutex
counters *simplelru.LRU[key, *Counter]
}
// key is what tells a counter from the others: its scope, with its
// netblock, AS number or name.
type key struct {
scope string
netblock netip.Prefix
asn string
name string
}
// New returns Counters for params, with nothing counted yet.
func New(params Params) *Counters {
counters, err := simplelru.NewLRU[key, *Counter](maxCounters, nil)
if err != nil {
panic(err) // NewLRU fails only for a size below one
}
return &Counters{params: params, counters: counters}
}
// Count counts r, a request that has ended, at now, in each scope it is
// in whose thresholds are not all off: its client, the netblock around
// it, its AS number once known, the whole service, and each named
// netblock it is in. Only the counts whose threshold is set are counted.
// For each scope whose count is over a threshold, it raises an anomaly
// alert, for the first such count in the order requests and bytes in the
// minute, then in the hour; the alert queue's cooldown holds back the
// repeats. Nothing is refused or banned.
func (c *Counters) Count(now time.Time, r Request) {
var raised []alerts.Alert
c.mu.Lock()
for _, scope := range c.scopesOf(r) {
counter, found := c.counters.Get(scope.key)
if !found {
counter = scope.key.counter()
c.counters.Add(scope.key, counter)
}
over, passed := counter.add(now, r.Bytes, scope.thresholds)
if passed {
raised = append(raised, alertFor(r, scope.key, over))
}
}
c.mu.Unlock()
for _, alert := range raised {
c.params.Alerts.Raise(alert)
}
}
// Snapshot returns every counter, sorted by scope, then by netblock, AS
// number and name, as alerts.json lists them.
func (c *Counters) Snapshot() []Counter {
c.mu.Lock()
counters := make([]Counter, 0, c.counters.Len())
for _, counter := range c.counters.Values() {
counters = append(counters, *counter)
}
c.mu.Unlock()
slices.SortFunc(counters, func(a, b Counter) int {
return cmp.Or(cmp.Compare(a.Scope, b.Scope), a.Netblock.Compare(b.Netblock),
cmp.Compare(a.ASN, b.ASN), cmp.Compare(a.Name, b.Name))
})
return counters
}
// Load puts counters, read from alerts.json, in place of those held, in
// the order they were last counted, as the starts of their buckets tell,
// so that the one counted least recently is dropped first. Each netblock
// is masked to its length, so that 203.0.113.9/24 is 203.0.113.0/24.
// Buckets whose time has passed at now are emptied, and a counter left
// with every bucket empty is dropped.
func (c *Counters) Load(counters []Counter, now time.Time) {
counters = slices.Clone(counters)
slices.SortStableFunc(counters, func(a, b Counter) int {
return a.lastStart().Compare(b.lastStart())
})
c.mu.Lock()
defer c.mu.Unlock()
c.counters.Purge()
for _, counter := range counters {
counter.Netblock = counter.Netblock.Masked()
empty := true
for _, count := range counter.counts() {
if count.buckets.Passed(now, count.length) {
*count.buckets = ratelimit.Buckets{}
}
empty = empty && *count.buckets == ratelimit.Buckets{}
}
if !empty {
c.counters.Add(counter.key(), &counter)
}
}
}
// scope is a scope a request is counted in, and its thresholds.
type scope struct {
key key
thresholds Thresholds
}
// scopesOf returns the scopes r is in whose thresholds are not all off.
func (c *Counters) scopesOf(r Request) []scope {
p := c.params
client := r.Client.Unmap()
all := []scope{
{key{scope: ScopeClient, netblock: r.ClientGroup}, p.Client},
{key{scope: ScopeNet, netblock: c.netAround(client)}, p.Net},
{key{scope: ScopeTotal}, p.Total},
}
if r.ASN != "" {
all = append(all, scope{key{scope: ScopeASN, asn: r.ASN}, p.ASN})
}
for _, named := range p.NamedNetblocks {
if named.Netblock.Contains(client) {
all = append(all, scope{
key{scope: ScopeWatch, netblock: named.Netblock, name: named.Name}, p.Watch,
})
}
}
return slices.DeleteFunc(all, func(s scope) bool {
return s.thresholds == Thresholds{}
})
}
// netAround returns the netblock around client that ScopeNet counts it
// in: NetV4Prefix or NetV6Prefix long.
func (c *Counters) netAround(client netip.Addr) netip.Prefix {
length := c.params.NetV6Prefix
if client.Is4() {
length = c.params.NetV4Prefix
}
return netip.PrefixFrom(client, length).Masked()
}
// overThreshold is a count over its threshold: what it counts, requests or
// bytes, its window, the count and the threshold.
type overThreshold struct {
kind, window string
count float64
threshold int64
}
// add counts a request of bytes at now in each of c's counts whose
// threshold, in thresholds, is set, and returns the first count over its
// threshold, and whether there is one.
func (c *Counter) add(
now time.Time, bytes int64, thresholds Thresholds,
) (overThreshold, bool) {
// In the order of counts.
inOrder := [4]int64{
thresholds.RequestsPerMinute, thresholds.BytesPerMinute,
thresholds.RequestsPerHour, thresholds.BytesPerHour,
}
var (
first overThreshold
passed bool
)
for i, count := range c.counts() {
threshold := inOrder[i]
if threshold == 0 {
continue
}
n := int64(1)
if count.kind == ratelimit.KindBytes {
n = bytes
}
counted := count.buckets.Add(now, count.length, n)
if !passed && counted > float64(threshold) {
first = overThreshold{count.kind, count.window, counted, threshold}
passed = true
}
}
return first, passed
}
// bucketCount is one of a counter's four counts: requests or bytes, in a
// window of length, and the buckets they are counted in.
type bucketCount struct {
kind, window string
length time.Duration
buckets *ratelimit.Buckets
}
// counts returns c's counts: requests and bytes in the minute, then in
// the hour.
func (c *Counter) counts() [4]bucketCount {
return [4]bucketCount{
{ratelimit.KindRequests, minute, time.Minute, &c.Minute},
{ratelimit.KindBytes, minute, time.Minute, &c.MinuteBytes},
{ratelimit.KindRequests, hour, time.Hour, &c.Hour},
{ratelimit.KindBytes, hour, time.Hour, &c.HourBytes},
}
}
// lastStart returns the start of c's latest bucket, which tells, to the
// minute or to the hour, when c was last counted.
func (c *Counter) lastStart() time.Time {
var latest time.Time
for _, count := range c.counts() {
if count.buckets.Start.After(latest) {
latest = count.buckets.Start
}
}
return latest
}
// key returns what tells c from the other counters.
func (c *Counter) key() key {
return key{scope: c.Scope, netblock: c.Netblock, asn: c.ASN, name: c.Name}
}
// counter returns a counter for k, with nothing counted yet.
func (k key) counter() *Counter {
return &Counter{Scope: k.scope, Netblock: k.netblock, ASN: k.asn, Name: k.name}
}
// alertFor returns the anomaly alert for o, a count over its threshold in
// the scope k, which r took over it. It gives r's client, with its AS
// number, AS name and country, and the netblock counted, of a client, the
// netblock around it or a named netblock. Its detail gives the scope, the
// AS number or the name of a scope that has one, the window, what is
// counted, the count and the threshold.
func alertFor(r Request, k key, o overThreshold) alerts.Alert {
detail := map[string]any{
"scope": k.scope, "window": o.window, "kind": o.kind, "count": o.count,
"threshold": o.threshold,
}
var counted string
switch k.scope {
case ScopeClient:
counted = "the client " + k.netblock.String()
case ScopeNet:
counted = "the netblock " + k.netblock.String()
case ScopeASN:
counted = k.asn
detail["asn"] = k.asn
case ScopeTotal:
counted = "the whole service"
default: // watch
counted = "the named netblock " + k.name + ", " + k.netblock.String()
detail["name"] = k.name
}
return alerts.Alert{
Event: alerts.EventAnomaly,
Client: r.Client,
Netblock: k.netblock,
ASN: r.ASN,
ASName: r.ASName,
Country: r.Country,
Reason: fmt.Sprintf("%s per %s of %s over the threshold of %d", o.kind, o.window,
counted, o.threshold),
Detail: detail,
}
}
+238
View File
@@ -0,0 +1,238 @@
package anomaly_test
import (
"encoding/json"
"fmt"
"net/netip"
"net/url"
"reflect"
"slices"
"testing"
"time"
"sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/ratelimit"
)
// maxCounters is how many counters are kept.
const maxCounters = 20000
func TestEachScopeHasACooldownOfItsOwn(t *testing.T) {
t.Parallel()
queue := newQueue()
office := netip.MustParsePrefix("203.0.113.0/24")
overAtTheSecond := anomaly.Thresholds{RequestsPerMinute: 1}
counters := anomaly.New(anomaly.Params{
Client: overAtTheSecond, Net: overAtTheSecond, ASN: overAtTheSecond,
Total: overAtTheSecond, Watch: overAtTheSecond,
// The netblock around a client is the client's own, and two names
// name one netblock.
NetV4Prefix: 32,
NamedNetblocks: []anomaly.NamedNetblock{
{Name: "office", Netblock: office}, {Name: "hq", Netblock: office},
},
Alerts: queue,
})
// The first client's second request is over the threshold in the six
// scopes it is in. The other client's two are both over it in the whole
// service and in each named netblock, three repeats each, and its
// second is over it in the scopes of its own, its client, its netblock
// and its AS number, which are no repeats.
for _, r := range []anomaly.Request{
{Client: netip.MustParseAddr("203.0.113.9"), ASN: "AS64496"},
{Client: netip.MustParseAddr("203.0.113.10"), ASN: "AS64511"},
} {
r.ClientGroup = netip.PrefixFrom(r.Client, 32)
for range 2 {
counters.Count(midnight(), r)
}
}
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 9 || queue.Suppressed() != 6 {
t.Fatalf("%d alerts wait and %d are held back, want 9 and 6: %+v",
len(waiting), queue.Suppressed(), waiting)
}
// alerts.json keeps each scope's cooldown: each alert raised again
// after a restart is a repeat.
data, err := json.Marshal(queue.Snapshot())
if err != nil {
t.Fatalf("encode: %v", err)
}
var read alerts.State
err = json.Unmarshal(data, &read)
if err != nil {
t.Fatalf("decode: %v", err)
}
after := newQueue()
after.Load(read)
for _, alert := range read.Waiting[alerts.DestinationWebhook] {
after.Raise(alert)
}
if after.Suppressed() != 9 {
t.Errorf("after loading, %d alerts are held back, want 9", after.Suppressed())
}
}
func TestKeepsAtMost20000CountersDroppingTheLeastRecentlyCounted(t *testing.T) {
t.Parallel()
counters := newCounters(anomaly.Params{
Client: anomaly.Thresholds{RequestsPerMinute: 1000},
})
for i := range maxCounters {
counters.Count(midnight(), request(i))
}
// Counted again, the first client is the most recently counted, and
// the second is dropped for a new one.
counters.Count(midnight(), request(0))
counters.Count(midnight(), request(maxCounters))
got := counters.Snapshot()
if len(got) != maxCounters || !holds(got, 0) || holds(got, 1) ||
!holds(got, maxCounters) {
t.Errorf("%d counters, holding the first client %v, the second %v and the "+
"new one %v, want %d, the first and the new one", len(got), holds(got, 0),
holds(got, 1), holds(got, maxCounters), maxCounters)
}
}
func TestLoadEmptiesBucketsWhoseTimeHasPassedAndDropsEmptyCounters(t *testing.T) {
t.Parallel()
counters := newCounters(anomaly.Params{
Net: anomaly.Thresholds{RequestsPerMinute: 1000, RequestsPerHour: 1000},
Total: anomaly.Thresholds{RequestsPerMinute: 1000},
NetV4Prefix: 24,
})
halfAnHourOn := midnight().Add(30 * time.Minute)
// Half an hour on, the hour's buckets count still, and the minute's
// do not.
counters.Load([]anomaly.Counter{
{
Scope: anomaly.ScopeNet,
Netblock: netip.MustParsePrefix("203.0.113.9/24"),
Minute: ratelimit.Buckets{Start: midnight(), Current: 5},
Hour: ratelimit.Buckets{Start: midnight(), Current: 7},
},
{
Scope: anomaly.ScopeTotal,
Minute: ratelimit.Buckets{Start: midnight(), Current: 1},
},
}, halfAnHourOn)
// The whole service's counter, left empty, is dropped, and the
// netblock read is masked to its length.
netblock := anomaly.Counter{
Scope: anomaly.ScopeNet,
Netblock: netip.MustParsePrefix("203.0.113.0/24"),
Hour: ratelimit.Buckets{Start: midnight(), Current: 7},
}
if got, want := counters.Snapshot(), []anomaly.Counter{netblock}; !reflect.DeepEqual(
got, want) {
t.Errorf("counters read\n%+v\nwant\n%+v", got, want)
}
// A request from the netblock is counted with the requests read.
counters.Count(halfAnHourOn, anomaly.Request{
Client: netip.MustParseAddr("203.0.113.9"),
ClientGroup: netip.MustParsePrefix("203.0.113.9/32"),
})
netblock.Minute = ratelimit.Buckets{Start: halfAnHourOn, Current: 1}
netblock.Hour.Current = 8
want := []anomaly.Counter{netblock, {
Scope: anomaly.ScopeTotal,
Minute: ratelimit.Buckets{Start: halfAnHourOn, Current: 1},
}}
if got := counters.Snapshot(); !reflect.DeepEqual(got, want) {
t.Errorf("counters after a request\n%+v\nwant\n%+v", got, want)
}
}
func TestLoadDropsTheLeastRecentlyCountedFirst(t *testing.T) {
t.Parallel()
counters := newCounters(anomaly.Params{
Client: anomaly.Thresholds{RequestsPerMinute: 1000},
})
now := midnight().Add(time.Minute)
// The second half of the file was counted in the minute before the
// first half.
read := make([]anomaly.Counter, 0, maxCounters)
for i := range maxCounters {
start := now
if i >= maxCounters/2 {
start = midnight()
}
read = append(read, anomaly.Counter{
Scope: anomaly.ScopeClient, Netblock: request(i).ClientGroup,
Minute: ratelimit.Buckets{Start: start, Current: 1},
})
}
counters.Load(read, now)
counters.Count(now, request(maxCounters))
got := counters.Snapshot()
if !holds(got, 0) || holds(got, maxCounters/2) {
t.Errorf("holding the first client of the file %v, and the first counted in "+
"the minute before %v, want only the first", holds(got, 0),
holds(got, maxCounters/2))
}
}
// midnight is the time of the tests' requests.
func midnight() time.Time {
return time.Date(2026, 10, 6, 0, 0, 0, 0, time.UTC)
}
// newCounters returns Counters for params, whose alerts go nowhere.
func newCounters(params anomaly.Params) *anomaly.Counters {
params.Alerts = alerts.New(alerts.Params{})
return anomaly.New(params)
}
// newQueue returns a queue of alerts to a webhook, with the default
// cooldown, which keeps them waiting, since it is never run.
func newQueue() *alerts.Queue {
return alerts.New(alerts.Params{
WebhookURL: &url.URL{Scheme: "https", Host: "alerts.example"},
Events: alerts.Events(),
Cooldown: 15 * time.Minute,
MaxPerHour: 60,
Now: midnight,
})
}
// request returns a request from client number i, an address in
// 10.0.0.0/8.
func request(i int) anomaly.Request {
client := netip.MustParseAddr(fmt.Sprintf("10.%d.%d.%d", i>>16, i>>8&255, i&255))
return anomaly.Request{Client: client, ClientGroup: netip.PrefixFrom(client, 32)}
}
// holds reports whether counters hold the counter of client number i.
func holds(counters []anomaly.Counter, i int) bool {
return slices.ContainsFunc(counters, func(counter anomaly.Counter) bool {
return counter.Netblock == request(i).ClientGroup
})
}
+115 -5
View File
@@ -24,6 +24,7 @@ import (
"unicode/utf8" "unicode/utf8"
"sneak.berlin/go/smallwebwaf/internal/alerts" "sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/remotelog" "sneak.berlin/go/smallwebwaf/internal/remotelog"
) )
@@ -220,6 +221,23 @@ type Config struct {
AlertEvents []string AlertEvents []string
AlertCooldown time.Duration AlertCooldown time.Duration
AlertMaxPerHour int AlertMaxPerHour int
// The anomaly thresholds, which only raise alerts: the most requests
// and bytes a minute and an hour per client (SWWAF_ANOMALY_CLIENT_*),
// per netblock around a client (SWWAF_ANOMALY_NET_*), per AS number
// (SWWAF_ANOMALY_ASN_*), for the whole service (SWWAF_ANOMALY_TOTAL_*)
// and per named netblock (SWWAF_WATCH_*), each 0 while it is off.
// AnomalyNetV4Prefix and AnomalyNetV6Prefix are the lengths of the
// netblock around a client (SWWAF_ANOMALY_NET_V4_PREFIX and
// SWWAF_ANOMALY_NET_V6_PREFIX), and WatchNets the named netblocks
// (SWWAF_WATCH_NETS).
AnomalyClient anomaly.Thresholds
AnomalyNet anomaly.Thresholds
AnomalyASN anomaly.Thresholds
AnomalyTotal anomaly.Thresholds
AnomalyWatch anomaly.Thresholds
AnomalyNetV4Prefix int
AnomalyNetV6Prefix int
WatchNets []anomaly.NamedNetblock
// settings are the values read, as given or by default, and the // settings are the values read, as given or by default, and the
// files they were read from, for the log line at start. // files they were read from, for the log line at start.
@@ -240,6 +258,7 @@ const (
mebibyte = 1 << 20 mebibyte = 1 << 20
gibibyte = 1 << 30 gibibyte = 1 << 30
ipv4Bits = 32 ipv4Bits = 32
ipv6Bits = 128
// minTokenLength is the fewest characters a token may have. // minTokenLength is the fewest characters a token may have.
minTokenLength = 32 minTokenLength = 32
// masked is what the log shows for a token that is set, and in place of // masked is what the log shows for a token that is set, and in place of
@@ -287,6 +306,10 @@ var (
errNotBanResponse = errors.New("is not 403, 429 or close") errNotBanResponse = errors.New("is not 403, 429 or close")
errNotV4Prefix = errors.New( errNotV4Prefix = errors.New(
"is not the length of an IPv4 netblock, from 0 to 32, such as 24") "is not the length of an IPv4 netblock, from 0 to 32, such as 24")
errNotV6Prefix = errors.New(
"is not the length of an IPv6 netblock, from 0 to 128, such as 48")
errNotNamedNetblock = errors.New(
"is not a name, = and a netblock, such as office=203.0.113.0/24")
errNotAbsolutePath = errors.New( errNotAbsolutePath = errors.New(
"is not an absolute path, such as /var/lib/smallwebwaf") "is not an absolute path, such as /var/lib/smallwebwaf")
errShortToken = errors.New("is shorter than 32 characters") errShortToken = errors.New("is shorter than 32 characters")
@@ -398,8 +421,16 @@ func FromEnvironment(lookupEnv func(string) (string, bool)) (*Config, error) {
AlertNtfyToken: env.secret("SWWAF_ALERT_NTFY_TOKEN"), AlertNtfyToken: env.secret("SWWAF_ALERT_NTFY_TOKEN"),
AlertEvents: env.alertEvents("SWWAF_ALERT_EVENTS", AlertEvents: env.alertEvents("SWWAF_ALERT_EVENTS",
strings.Join(alerts.Events(), ",")), strings.Join(alerts.Events(), ",")),
AlertCooldown: env.duration("SWWAF_ALERT_COOLDOWN", "15m"), AlertCooldown: env.duration("SWWAF_ALERT_COOLDOWN", "15m"),
AlertMaxPerHour: env.numberOrOff("SWWAF_ALERT_MAX_PER_HOUR", "60"), AlertMaxPerHour: env.numberOrOff("SWWAF_ALERT_MAX_PER_HOUR", "60"),
AnomalyClient: env.thresholds("SWWAF_ANOMALY_CLIENT_"),
AnomalyNet: env.thresholds("SWWAF_ANOMALY_NET_"),
AnomalyASN: env.thresholds("SWWAF_ANOMALY_ASN_"),
AnomalyTotal: env.thresholds("SWWAF_ANOMALY_TOTAL_"),
AnomalyWatch: env.thresholds("SWWAF_WATCH_"),
AnomalyNetV4Prefix: env.v4Prefix("SWWAF_ANOMALY_NET_V4_PREFIX", "24"),
AnomalyNetV6Prefix: env.v6Prefix("SWWAF_ANOMALY_NET_V6_PREFIX", "48"),
WatchNets: env.namedNetblocks("SWWAF_WATCH_NETS"),
} }
cfg.LogRemoteAppName = env.appName("SWWAF_LOG_REMOTE_APP_NAME", cfg.LogRemoteAppName = env.appName("SWWAF_LOG_REMOTE_APP_NAME",
@@ -666,9 +697,9 @@ func (e *environment) checkLookupDBPath(cfg *Config) {
// checkCountriesAndLookups refuses a country on both country lists, and, // checkCountriesAndLookups refuses a country on both country lists, and,
// while SWWAF_LOOKUP_SOURCE is off, each setting that needs clients looked // while SWWAF_LOOKUP_SOURCE is off, each setting that needs clients looked
// up: the country lists, SWWAF_ADD_LOOKUP_HEADERS, and the biased // up: the country lists, SWWAF_ADD_LOOKUP_HEADERS, the biased thresholds,
// thresholds, of which SWWAF_UNKNOWN_LIMIT_PERCENT needs them only below // of which SWWAF_UNKNOWN_LIMIT_PERCENT needs them only below 100, where it
// 100, where it lowers a limit. // lowers a limit, and the anomaly thresholds per AS number.
func (e *environment) checkCountriesAndLookups(cfg *Config) { func (e *environment) checkCountriesAndLookups(cfg *Config) {
for _, country := range cfg.ExclusivelyAllowedCountries { for _, country := range cfg.ExclusivelyAllowedCountries {
if slices.Contains(cfg.DeniedCountries, country) { if slices.Contains(cfg.DeniedCountries, country) {
@@ -693,6 +724,10 @@ func (e *environment) checkCountriesAndLookups(cfg *Config) {
{"SWWAF_ASN_BYTES_PERCENT", len(cfg.ASNBytesPercent) > 0}, {"SWWAF_ASN_BYTES_PERCENT", len(cfg.ASNBytesPercent) > 0},
{"SWWAF_COUNTRY_BYTES_PERCENT", len(cfg.CountryBytesPercent) > 0}, {"SWWAF_COUNTRY_BYTES_PERCENT", len(cfg.CountryBytesPercent) > 0},
{"SWWAF_UNKNOWN_LIMIT_PERCENT", cfg.UnknownLimitPercent < 100}, {"SWWAF_UNKNOWN_LIMIT_PERCENT", cfg.UnknownLimitPercent < 100},
{"SWWAF_ANOMALY_ASN_REQUESTS_PER_MINUTE", cfg.AnomalyASN.RequestsPerMinute > 0},
{"SWWAF_ANOMALY_ASN_REQUESTS_PER_HOUR", cfg.AnomalyASN.RequestsPerHour > 0},
{"SWWAF_ANOMALY_ASN_BYTES_PER_MINUTE", cfg.AnomalyASN.BytesPerMinute > 0},
{"SWWAF_ANOMALY_ASN_BYTES_PER_HOUR", cfg.AnomalyASN.BytesPerHour > 0},
} { } {
if setting.set { if setting.set {
e.check(setting.name, fmt.Errorf("is set while SWWAF_LOOKUP_SOURCE is off; %w", e.check(setting.name, fmt.Errorf("is set while SWWAF_LOOKUP_SOURCE is off; %w",
@@ -744,6 +779,35 @@ func (e *environment) v4Prefix(name, defaultValue string) int {
return length return length
} }
// v6Prefix reads a setting that is the length of an IPv6 netblock.
func (e *environment) v6Prefix(name, defaultValue string) int {
length, err := parseV6Prefix(e.value(name, defaultValue))
e.check(name, err)
return length
}
// thresholds reads the four anomaly thresholds whose settings' names
// start with prefix: requests and bytes per minute and per hour. Each is
// off by default.
func (e *environment) thresholds(prefix string) anomaly.Thresholds {
return anomaly.Thresholds{
RequestsPerMinute: e.count(prefix+"REQUESTS_PER_MINUTE", off),
RequestsPerHour: e.count(prefix+"REQUESTS_PER_HOUR", off),
BytesPerMinute: e.size(prefix+"BYTES_PER_MINUTE", off),
BytesPerHour: e.size(prefix+"BYTES_PER_HOUR", off),
}
}
// namedNetblocks reads a setting that is a list of named netblocks. It is
// empty by default.
func (e *environment) namedNetblocks(name string) []anomaly.NamedNetblock {
named, err := parseNamedNetblocks(e.value(name, ""))
e.check(name, err)
return named
}
// absolutePath reads a setting that is an absolute path. // absolutePath reads a setting that is an absolute path.
func (e *environment) absolutePath(name, defaultValue string) string { func (e *environment) absolutePath(name, defaultValue string) string {
path := e.value(name, defaultValue) path := e.value(name, defaultValue)
@@ -1085,6 +1149,16 @@ func parseV4Prefix(value string) (int, error) {
return n, nil return n, nil
} }
// parseV6Prefix reads the length of an IPv6 netblock, from 0 to 128.
func parseV6Prefix(value string) (int, error) {
n, err := strconv.Atoi(value)
if err != nil || n < 0 || n > ipv6Bits {
return 0, fmt.Errorf("%q %w", value, errNotV6Prefix)
}
return n, nil
}
// parseList splits a comma-separated list and trims the spaces around // parseList splits a comma-separated list and trims the spaces around
// each item. An empty value is an empty list. // each item. An empty value is an empty list.
func parseList(value string) ([]string, error) { func parseList(value string) ([]string, error) {
@@ -1144,6 +1218,42 @@ func parseNetblock(value string) (netip.Prefix, error) {
return netip.PrefixFrom(addr, addr.BitLen()), nil return netip.PrefixFrom(addr, addr.BitLen()), nil
} }
// parseNamedNetblocks reads a comma-separated list of named netblocks,
// each a name, = and a netblock, such as office=203.0.113.0/24. An empty
// value is an empty list. A name listed twice is an error.
func parseNamedNetblocks(value string) ([]anomaly.NamedNetblock, error) {
items, err := parseList(value)
if err != nil {
return nil, err
}
named := make([]anomaly.NamedNetblock, 0, len(items))
for _, item := range items {
name, netblockText, found := strings.Cut(item, "=")
name = strings.TrimSpace(name)
if !found || name == "" {
return nil, fmt.Errorf("%q %w", item, errNotNamedNetblock)
}
netblock, err := parseNetblock(strings.TrimSpace(netblockText))
if err != nil {
return nil, err
}
if slices.ContainsFunc(named, func(n anomaly.NamedNetblock) bool {
return n.Name == name
}) {
return nil, fmt.Errorf("%q %w", name, errListedTwice)
}
named = append(named, anomaly.NamedNetblock{Name: name, Netblock: netblock})
}
return named, nil
}
// parsePathPrefixes reads a comma-separated list of path prefixes, each // parsePathPrefixes reads a comma-separated list of path prefixes, each
// starting with /. // starting with /.
func parsePathPrefixes(value string) ([]string, error) { func parsePathPrefixes(value string) ([]string, error) {
+211 -11
View File
@@ -12,10 +12,12 @@ import (
"path/filepath" "path/filepath"
"reflect" "reflect"
"slices" "slices"
"strconv"
"strings" "strings"
"testing" "testing"
"time" "time"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/config" "sneak.berlin/go/smallwebwaf/internal/config"
) )
@@ -85,8 +87,59 @@ const (
alertEvents = "SWWAF_ALERT_EVENTS" alertEvents = "SWWAF_ALERT_EVENTS"
alertCooldown = "SWWAF_ALERT_COOLDOWN" alertCooldown = "SWWAF_ALERT_COOLDOWN"
alertMaxPerHour = "SWWAF_ALERT_MAX_PER_HOUR" alertMaxPerHour = "SWWAF_ALERT_MAX_PER_HOUR"
anomalyNetV4Prefix = "SWWAF_ANOMALY_NET_V4_PREFIX"
anomalyNetV6Prefix = "SWWAF_ANOMALY_NET_V6_PREFIX"
watchNets = "SWWAF_WATCH_NETS"
) )
// The anomaly thresholds: each of the prefixes below, which name a scope,
// followed by each of the four ends.
const (
anomalyClient = "SWWAF_ANOMALY_CLIENT_"
anomalyNet = "SWWAF_ANOMALY_NET_"
anomalyASN = "SWWAF_ANOMALY_ASN_"
anomalyTotal = "SWWAF_ANOMALY_TOTAL_"
watch = "SWWAF_WATCH_"
requestsPerMinute = "REQUESTS_PER_MINUTE"
requestsPerHour = "REQUESTS_PER_HOUR"
bytesPerMinute = "BYTES_PER_MINUTE"
bytesPerHour = "BYTES_PER_HOUR"
)
// anomalyScopes returns the prefixes of the anomaly thresholds, one for
// each scope.
func anomalyScopes() []string {
return []string{anomalyClient, anomalyNet, anomalyASN, anomalyTotal, watch}
}
// anomalyThresholds returns the names of the twenty anomaly thresholds.
func anomalyThresholds() []string {
ends := []string{requestsPerMinute, requestsPerHour, bytesPerMinute, bytesPerHour}
names := make([]string, 0, len(anomalyScopes())*len(ends))
for _, scope := range anomalyScopes() {
for _, end := range ends {
names = append(names, scope+end)
}
}
return names
}
// loggedAnomalyDefaults returns the anomaly settings as the settings
// logged at start give them by default.
func loggedAnomalyDefaults() map[string]string {
logged := map[string]string{
anomalyNetV4Prefix: "24", anomalyNetV6Prefix: "48", watchNets: "",
}
for _, name := range anomalyThresholds() {
logged[name] = off
}
return logged
}
// defaultAlertEvents is the default of SWWAF_ALERT_EVENTS, and // defaultAlertEvents is the default of SWWAF_ALERT_EVENTS, and
// defaultAlertCooldown that of SWWAF_ALERT_COOLDOWN. // defaultAlertCooldown that of SWWAF_ALERT_COOLDOWN.
const ( const (
@@ -897,14 +950,18 @@ func TestSettingNeedingLookupsStopsTheStartWhileTheyAreOff(t *testing.T) {
t.Parallel() t.Parallel()
for name, value := range map[string]string{ for name, value := range map[string]string{
deniedCountries: "kp", deniedCountries: "kp",
allowedCountries: "de", allowedCountries: "de",
addLookupHeaders: enabled, addLookupHeaders: enabled,
asnLimitPercent: "AS64496:50", asnLimitPercent: "AS64496:50",
countryLimitPercent: "cn:25", countryLimitPercent: "cn:25",
asnBytesPercent: "AS64496:50", asnBytesPercent: "AS64496:50",
countryBytesPercent: "cn:25", countryBytesPercent: "cn:25",
unknownLimitPercent: "99", unknownLimitPercent: "99",
anomalyASN + requestsPerMinute: "1000",
anomalyASN + requestsPerHour: "10000",
anomalyASN + bytesPerMinute: "1G",
anomalyASN + bytesPerHour: "10G",
} { } {
t.Run(name, func(t *testing.T) { t.Run(name, func(t *testing.T) {
t.Parallel() t.Parallel()
@@ -920,12 +977,153 @@ func TestSettingNeedingLookupsStopsTheStartWhileTheyAreOff(t *testing.T) {
} }
// Set empty, the lists need nothing looked up, and nor does // Set empty, the lists need nothing looked up, and nor does
// SWWAF_UNKNOWN_LIMIT_PERCENT at 100, which lowers no limit. // SWWAF_UNKNOWN_LIMIT_PERCENT at 100, which lowers no limit, an anomaly
fromEnvironment(t, environment{ // threshold per AS number that is off, or any other anomaly threshold.
env := environment{
lookupSource: off, deniedCountries: "", allowedCountries: "", lookupSource: off, deniedCountries: "", allowedCountries: "",
asnLimitPercent: "", countryLimitPercent: "", asnBytesPercent: "", asnLimitPercent: "", countryLimitPercent: "", asnBytesPercent: "",
countryBytesPercent: "", unknownLimitPercent: "100", countryBytesPercent: "", unknownLimitPercent: "100",
}) }
for _, name := range anomalyThresholds() {
env[name] = "1000"
if strings.HasPrefix(name, anomalyASN) {
env[name] = off
}
}
fromEnvironment(t, env)
}
func TestAnomalySettingsDefaults(t *testing.T) {
t.Parallel()
cfg := fromEnvironment(t, environment{})
wantAllOff(t, cfg)
if cfg.AnomalyNetV4Prefix != 24 || cfg.AnomalyNetV6Prefix != 48 ||
len(cfg.WatchNets) != 0 {
t.Errorf("%s, %s and %s gave %d, %d and %v, want 24, 48 and none",
anomalyNetV4Prefix, anomalyNetV6Prefix, watchNets, cfg.AnomalyNetV4Prefix,
cfg.AnomalyNetV6Prefix, cfg.WatchNets)
}
}
func TestAnomalySettingsAsSet(t *testing.T) {
t.Parallel()
// Each threshold of a scope its own value; bytes are sizes.
env := environment{
anomalyNetV4Prefix: "16",
anomalyNetV6Prefix: "56",
// Spaces around a name or a netblock, and a bare address.
watchNets: "office = 203.0.113.0/24, scraper-x=198.51.100.7,v6=2001:db8::/32",
}
want := map[string]anomaly.Thresholds{}
for i, scope := range anomalyScopes() {
n := int64(i + 1)
env[scope+requestsPerMinute] = strconv.FormatInt(n, 10)
env[scope+requestsPerHour] = strconv.FormatInt(10*n, 10)
env[scope+bytesPerMinute] = strconv.FormatInt(n, 10) + "K"
env[scope+bytesPerHour] = strconv.FormatInt(n, 10) + "G"
want[scope] = anomaly.Thresholds{
RequestsPerMinute: n, RequestsPerHour: 10 * n,
BytesPerMinute: n << 10, BytesPerHour: n << 30,
}
}
cfg := fromEnvironment(t, env)
if got := thresholdsByScope(cfg); !maps.Equal(got, want) {
t.Errorf("thresholds by scope\n%+v\nwant\n%+v", got, want)
}
wantNamed := []anomaly.NamedNetblock{
{Name: "office", Netblock: netip.MustParsePrefix("203.0.113.0/24")},
{Name: "scraper-x", Netblock: netip.MustParsePrefix("198.51.100.7/32")},
{Name: "v6", Netblock: netip.MustParsePrefix("2001:db8::/32")},
}
if cfg.AnomalyNetV4Prefix != 16 || cfg.AnomalyNetV6Prefix != 56 ||
!slices.Equal(cfg.WatchNets, wantNamed) {
t.Errorf("%s, %s and %s gave %d, %d and %v, want 16, 56 and %v",
anomalyNetV4Prefix, anomalyNetV6Prefix, watchNets, cfg.AnomalyNetV4Prefix,
cfg.AnomalyNetV6Prefix, cfg.WatchNets, wantNamed)
}
// off switches each threshold off.
for _, name := range anomalyThresholds() {
env[name] = off
}
wantAllOff(t, fromEnvironment(t, env))
}
// thresholdsByScope returns cfg's anomaly thresholds, each by the prefix
// of its scope's settings.
func thresholdsByScope(cfg *config.Config) map[string]anomaly.Thresholds {
return map[string]anomaly.Thresholds{
anomalyClient: cfg.AnomalyClient, anomalyNet: cfg.AnomalyNet,
anomalyASN: cfg.AnomalyASN, anomalyTotal: cfg.AnomalyTotal,
watch: cfg.AnomalyWatch,
}
}
// wantAllOff checks that every anomaly threshold of cfg is off.
func wantAllOff(t *testing.T, cfg *config.Config) {
t.Helper()
for scope, thresholds := range thresholdsByScope(cfg) {
if thresholds != (anomaly.Thresholds{}) {
t.Errorf("%s* gave %+v, want every one off", scope, thresholds)
}
}
}
func TestInvalidAnomalySettingStopsTheStartSayingWhatIsWrong(t *testing.T) {
t.Parallel()
const (
notCount = " is not a whole number of requests such as 1000, or off"
notSize = " is not a size such as 512K, 100M or 5G, or off"
notPositive = " must be more than zero, or off"
notNamed = " is not a name, = and a netblock, such as office=203.0.113.0/24"
notNetblock = " is not a netblock such as 10.0.0.0/8, or an address"
notV4Prefix = " is not the length of an IPv4 netblock, from 0 to 32, such as 24"
notV6Prefix = " is not the length of an IPv6 netblock, from 0 to 128, such as 48"
officeNetblock = "office=203.0.113.0/24"
scraperNetblock = "scraper=198.51.100.0/24"
)
for _, tc := range []struct{ name, value, want string }{
{anomalyClient + requestsPerMinute, "1K", `"1K"` + notCount},
{anomalyNet + requestsPerHour, "0", `"0"` + notPositive},
{anomalyTotal + bytesPerMinute, "1T", `"1T"` + notSize},
{watch + bytesPerHour, "-1G", `"-1G"` + notPositive},
{anomalyNetV4Prefix, "33", `"33"` + notV4Prefix},
{anomalyNetV4Prefix, off, `"off"` + notV4Prefix},
{anomalyNetV6Prefix, "129", `"129"` + notV6Prefix},
{anomalyNetV6Prefix, "/48", `"/48"` + notV6Prefix},
{watchNets, "office", `"office"` + notNamed},
{watchNets, "=203.0.113.0/24", `"=203.0.113.0/24"` + notNamed},
{watchNets, "office=203.0.113.300/24", `"203.0.113.300/24"` + notNetblock},
{watchNets, officeNetblock + ",", `"` + officeNetblock + `," has an empty item ` +
`in its list`},
{
watchNets, officeNetblock + "," + scraperNetblock + ",office=192.0.2.0/24",
`"office" is listed twice`,
},
} {
t.Run(tc.name+"="+tc.value, func(t *testing.T) {
t.Parallel()
_, err := config.FromEnvironment(environment{tc.name: tc.value}.lookupEnv)
want := tc.name + ": " + tc.want
if err == nil || err.Error() != want {
t.Errorf("error %v, want %s", err, want)
}
})
}
} }
func TestBiasedThresholdsAsSet(t *testing.T) { func TestBiasedThresholdsAsSet(t *testing.T) {
@@ -1461,6 +1659,8 @@ func TestLogsEachSettingWithItsValue(t *testing.T) {
alertCooldown: defaultAlertCooldown, alertCooldown: defaultAlertCooldown,
alertMaxPerHour: "60", alertMaxPerHour: "60",
} }
maps.Copy(want, loggedAnomalyDefaults())
if got := loggedSettings(t, cfg); !maps.Equal(got, want) { if got := loggedSettings(t, cfg); !maps.Equal(got, want) {
t.Errorf("logged settings\n%v\nwant\n%v", got, want) t.Errorf("logged settings\n%v\nwant\n%v", got, want)
} }
+30
View File
@@ -0,0 +1,30 @@
package proxy
import (
"net/netip"
"testing"
"sneak.berlin/go/smallwebwaf/internal/config"
)
func TestWithEveryAnomalyThresholdOffARequestIsNotCounted(t *testing.T) {
t.Parallel()
// A request from a client looked up through GeoJS, with every anomaly
// threshold off. Its handler has neither GeoJS's answers nor the
// anomaly counters, nor a clock, and the request no response: reading
// any of them to count the request panics.
rq := &request{
h: &handler{config: &config.Config{LookupSource: "geojs"}},
client: netip.MustParseAddr("203.0.113.9"),
lookedUp: true,
}
defer func() {
if r := recover(); r != nil {
t.Errorf("counting the request did work, with every threshold off: %v", r)
}
}()
rq.countAnomalies()
}
+377
View File
@@ -0,0 +1,377 @@
package proxy_test
import (
"maps"
"net/http"
"net/netip"
"reflect"
"slices"
"strconv"
"sync/atomic"
"testing"
"time"
"sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/proxy"
"sneak.berlin/go/smallwebwaf/internal/ratelimit"
"sneak.berlin/go/smallwebwaf/internal/requestlog"
)
// The anomaly thresholds: the prefix of a scope followed by the end of a
// count.
const (
anomalyClient = "SWWAF_ANOMALY_CLIENT_"
anomalyNet = "SWWAF_ANOMALY_NET_"
anomalyASN = "SWWAF_ANOMALY_ASN_"
anomalyTotal = "SWWAF_ANOMALY_TOTAL_"
anomalyWatch = "SWWAF_WATCH_"
requestsPerMinute = "REQUESTS_PER_MINUTE"
requestsPerHour = "REQUESTS_PER_HOUR"
bytesPerMinute = "BYTES_PER_MINUTE"
bytesPerHour = "BYTES_PER_HOUR"
)
// The other anomaly settings.
const (
anomalyNetV4Prefix = "SWWAF_ANOMALY_NET_V4_PREFIX"
anomalyNetV6Prefix = "SWWAF_ANOMALY_NET_V6_PREFIX"
watchNets = "SWWAF_WATCH_NETS"
)
const (
// clientsNet is the netblock around client at the default length, and
// office a named netblock of the same.
clientsNet = "203.0.113.0/24"
office = "office=" + clientsNet
// aLot is a threshold no test reaches.
aLot = "1000"
// hour is the window an alert names for a threshold per hour.
hour = "hour"
)
func TestEachScopeAndWindowOverItsThresholdAlertsOncePerCooldown(t *testing.T) {
t.Parallel()
for _, scope := range []struct {
prefix, scope string
// netblock is the alert's, and counted what its reason names. extra
// is what its detail gives besides what every anomaly alert's does.
netblock netip.Prefix
counted string
extra map[string]any
}{
{
anomalyClient, anomaly.ScopeClient, netip.MustParsePrefix(client + "/32"),
"the client " + client + "/32", nil,
},
{
anomalyNet, anomaly.ScopeNet, netip.MustParsePrefix(clientsNet),
"the netblock " + clientsNet, nil,
},
{anomalyASN, anomaly.ScopeASN, netip.Prefix{}, asnDE, map[string]any{"asn": asnDE}},
{anomalyTotal, anomaly.ScopeTotal, netip.Prefix{}, "the whole service", nil},
{
anomalyWatch, anomaly.ScopeWatch, netip.MustParsePrefix(clientsNet),
"the named netblock office, " + clientsNet, map[string]any{"name": "office"},
},
} {
for _, threshold := range []struct {
end, kind, window string
// value is the threshold, which the third upload of 100 bytes
// takes the count over, to count.
value int64
count float64
}{
{requestsPerMinute, ratelimit.KindRequests, minute, 2, 3},
{requestsPerHour, ratelimit.KindRequests, hour, 2, 3},
{bytesPerMinute, ratelimit.KindBytes, minute, 250, 300},
{bytesPerHour, ratelimit.KindBytes, hour, 250, 300},
} {
setting := scope.prefix + threshold.end
value := strconv.FormatInt(threshold.value, 10)
t.Run(setting, func(t *testing.T) {
t.Parallel()
s, clk, server, queue := startWithLookupsAndClock(t, map[string]string{
setting: value, watchNets: office,
})
start := clk.Now()
// The third upload takes the count over the threshold, and the
// fourth, within the cooldown, is held back. Each is passed to
// the app.
for range 4 {
s.uploadFrom(client)
}
detail := map[string]any{
"scope": scope.scope, "window": threshold.window, "kind": threshold.kind,
"count": threshold.count, "threshold": threshold.value,
}
maps.Copy(detail, scope.extra)
wantAlerts(t, queue, alerts.Alert{
Instance: alertInstance,
Time: start,
Event: alerts.EventAnomaly,
Client: netip.MustParseAddr(client),
Netblock: scope.netblock,
ASN: asnDE,
ASName: asNameDE,
Country: "DE",
Reason: threshold.kind + " per " + threshold.window + " of " +
scope.counted + " over the threshold of " + value,
Detail: detail,
})
wantAlertedAgainOnceTheCooldownHasRunOut(t, s, clk, queue)
if held := server.Ledger.Snapshot(); len(held) != 0 {
t.Errorf("the ledger holds %+v, want no ban", held)
}
})
}
}
}
// wantAlertedAgainOnceTheCooldownHasRunOut checks that, once the cooldown
// has run out after a first alert, which held back one repeat, the next
// count over the threshold, at the latest three uploads from client on,
// raises another alert, giving that repeat.
func wantAlertedAgainOnceTheCooldownHasRunOut(
t *testing.T, s *sender, clk *clock, queue *alerts.Queue,
) {
t.Helper()
clk.advance(15 * time.Minute)
for range 3 {
s.uploadFrom(client)
}
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 2 || !waiting[1].Time.Equal(clk.Now()) ||
waiting[1].SuppressedRepeats != 1 {
t.Errorf("alerts wait %+v, want the first and another, with 1 repeat", waiting)
}
}
func TestEveryRequestIsCountedWhateverIsDoneWithIt(t *testing.T) {
t.Parallel()
const (
allowed = "192.0.2.7" // in SWWAF_ALLOW_NETS
exempt = "192.0.2.10" // in SWWAF_RATE_LIMIT_EXEMPT_NETS
denied = "192.0.2.20" // in SWWAF_DENY_NETS
)
s, _, _, queue := startAppWithAlerts(t, readAndAnswer, map[string]string{
anomalyClient + requestsPerMinute: "2",
allowNets: allowed,
rateLimitExemptNets: exempt,
rateLimitExemptPaths: "/static/",
denyNets: denied,
})
// The third request of each takes its client's count over the threshold
// of 2.
for _, sent := range []struct {
from, path string
status int
action string
}{
{allowed, "/", http.StatusOK, requestlog.ActionForward},
{exempt, "/", http.StatusOK, requestlog.ActionForward},
{client, "/static/app.js", http.StatusOK, requestlog.ActionForward},
{denied, "/", http.StatusForbidden, requestlog.ActionDenied},
} {
for range 3 {
s.request(sent.from, sent.path, sent.status, sent.action)
}
}
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
got := make([]string, 0, len(waiting))
for _, alert := range waiting {
got = append(got, alert.Client.String())
}
if want := []string{allowed, exempt, client, denied}; !slices.Equal(got, want) {
t.Errorf("alerts for the clients %v, want %v", got, want)
}
}
func TestThresholdsOffCountNothingAndAlertNothing(t *testing.T) {
t.Parallel()
// With every threshold off, nothing is counted.
s, server, queue := startWithLookups(t, map[string]string{watchNets: office})
for range 5 {
s.uploadFrom(client)
}
if counters := server.Anomalies.Snapshot(); len(counters) != 0 {
t.Errorf("counters %+v, want none", counters)
}
wantAlerts(t, queue)
// With one set, its count alone is counted, in its scope alone.
s, clk, server, queue := startWithLookupsAndClock(t, map[string]string{
anomalyNet + requestsPerMinute: aLot, watchNets: office,
})
for range 5 {
s.uploadFrom(client)
}
want := []anomaly.Counter{{
Scope: anomaly.ScopeNet,
Netblock: netip.MustParsePrefix(clientsNet),
Minute: ratelimit.Buckets{Start: clk.Now(), Current: 5},
}}
if got := server.Anomalies.Snapshot(); !reflect.DeepEqual(got, want) {
t.Errorf("counters\n%+v\nwant\n%+v", got, want)
}
wantAlerts(t, queue)
}
func TestNetblockAroundAClientIsAsLongAsTheSettingsSay(t *testing.T) {
t.Parallel()
for _, tc := range []struct {
name string
env map[string]string
// Each client of sent sends one request, and want gives the
// netblocks they are counted in, each with its requests.
sent []string
want map[string]int64
}{
{
"by default", nil,
[]string{client, "203.0.113.200", "192.0.2.7", ipv6Client, "2001:db8:0:ffff::1"},
map[string]int64{clientsNet: 2, "192.0.2.0/24": 1, "2001:db8::/48": 2},
},
{
"as set", map[string]string{anomalyNetV4Prefix: "16", anomalyNetV6Prefix: "32"},
[]string{client, "203.0.200.1", ipv6Client, "2001:db8:ffff::1"},
map[string]int64{"203.0.0.0/16": 2, "2001:db8::/32": 2},
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
env := map[string]string{anomalyNet + requestsPerMinute: aLot}
maps.Copy(env, tc.env)
s, _, server, _ := startAppWithAlerts(t, readAndAnswer, env)
for _, from := range tc.sent {
s.get(from, http.StatusOK, requestlog.ActionForward)
}
got := map[string]int64{}
for _, counter := range server.Anomalies.Snapshot() {
got[counter.Netblock.String()] = counter.Minute.Current
}
if !maps.Equal(got, tc.want) {
t.Errorf("requests by netblock %v, want %v", got, tc.want)
}
})
}
}
func TestClientIsCountedForItsASNumberOnceTheLookupGivesOne(t *testing.T) {
t.Parallel()
s, server, _ := startWithLookups(t, map[string]string{
anomalyASN + requestsPerMinute: aLot,
})
// The lookup database does not hold unplaced.
for _, from := range []string{fromDE, fromDE, fromKP, noCountry, unplaced} {
s.uploadFrom(from)
}
got := map[string]int64{}
for _, counter := range server.Anomalies.Snapshot() {
got[counter.ASN] = counter.Minute.Current
}
if want := map[string]int64{asnDE: 2, asnKP: 1, "AS64500": 1}; !maps.Equal(got, want) {
t.Errorf("requests by AS number %v, want %v", got, want)
}
}
func TestRequestCountsForTheASNumberGeoJSGivesBeforeItEnds(t *testing.T) {
t.Parallel()
// The stand-in for GeoJS answers only once released, which the app
// does as it answers the request, and then waits until the answer is
// kept.
geojsURL, _, release := startHeldGeoJS(t)
var server atomic.Pointer[proxy.Server]
app := startApp(t, func(http.ResponseWriter, *http.Request) {
release()
waitUntil(func() bool {
_, kept := server.Load().GeoJS.Kept(netip.MustParsePrefix(fromDE + "/32"))
return kept
})
})
clk := &clock{now: time.Date(2026, 10, 6, 0, 0, 0, 0, time.UTC)}
addr, out, started := startProxyWithClock(t, app.URL, geojsURL, clk.Now,
map[string]string{
trustedProxies: trustLocalhost,
lookupTimeout: "1h",
anomalyASN + requestsPerMinute: aLot,
})
server.Store(started)
// The request went on without the answer, and is counted for the AS
// number it gives.
s := &sender{t: t, addr: addr, out: out}
if line := s.get(fromDE, http.StatusOK, requestlog.ActionForward); line.ASN != "" {
t.Errorf("log line has AS number %q, want none: the request waited", line.ASN)
}
want := []anomaly.Counter{{
Scope: anomaly.ScopeASN, ASN: asnDE,
Minute: ratelimit.Buckets{Start: clk.Now(), Current: 1},
}}
if got := started.Anomalies.Snapshot(); !reflect.DeepEqual(got, want) {
t.Errorf("counters\n%+v\nwant\n%+v", got, want)
}
}
func TestEachNamedNetblockCountsTheClientsInIt(t *testing.T) {
t.Parallel()
s, _, server, _ := startAppWithAlerts(t, readAndAnswer, map[string]string{
anomalyWatch + requestsPerMinute: aLot,
watchNets: office + ",wide=203.0.0.0/16,other=198.51.100.0/25",
})
// client is in office and in wide.
for _, from := range []string{client, "203.0.200.1", "192.0.2.7"} {
s.get(from, http.StatusOK, requestlog.ActionForward)
}
got := map[string]int64{}
for _, counter := range server.Anomalies.Snapshot() {
got[counter.Name] = counter.Minute.Current
}
if want := map[string]int64{"office": 1, "wide": 2}; !maps.Equal(got, want) {
t.Errorf("requests by named netblock %v, want %v", got, want)
}
}
+33 -31
View File
@@ -56,44 +56,23 @@ func (rq *request) limitBroken(now time.Time) bool {
return over return over
} }
// countBytes counts the request's bytes for the byte limits, once its // countBytes counts the request's bytes, as countedBytes gives them, for
// response has ended, and notes the client's byte totals for the log line; // the byte limits, once its response has ended, and notes the client's
// its requests stay there as the rate limits counted them. The bytes are // byte totals for the log line; its requests stay there as the rate limits
// the response's body bytes, the request's, or both, as SWWAF_BYTES_COUNT // counted them. Only a request passed to the app has them counted, and
// says; for an upgraded connection, such as a WebSocket, which has closed // only one the rate limits counted; in observe mode, not one that enforce
// by then, what it carried from the app counts with the response's and // mode would have refused. Bytes that take the client over a byte limit,
// what it carried from the client with the request's. Only a request // as its limit percentage for the byte limits lowers it, break it; the
// passed to the app has them counted, and only one the rate limits // response was passed on whole.
// counted; in observe mode, not one that enforce mode would have refused.
// Bytes that take the client over a byte limit, as its limit percentage
// for the byte limits lowers it, break it; the response was passed on
// whole.
func (rq *request) countBytes() { func (rq *request) countBytes() {
if !rq.counted || rq.line.WouldAction != "" { if !rq.counted || rq.line.WouldAction != "" {
return return
} }
response, request := rq.out.bytes, rq.requestBytes()
if rq.upgraded != nil {
response += rq.upgraded.fromApp.Load()
request += rq.upgraded.toApp.Load()
}
var bytes int64
switch rq.h.config.BytesCount {
case "response":
bytes = response
case "request":
bytes = request
default: // both
bytes = response + request
}
now := rq.h.now() now := rq.h.now()
counts, hit, over := rq.h.limiter.CountBytes(clientGroup(rq.client), now, bytes, counts, hit, over := rq.h.limiter.CountBytes(clientGroup(rq.client), now,
rq.bytesPercent.percent) rq.countedBytes(), rq.bytesPercent.percent)
rq.line.Counts.MinuteBytes = counts.MinuteBytes rq.line.Counts.MinuteBytes = counts.MinuteBytes
rq.line.Counts.HourBytes = counts.HourBytes rq.line.Counts.HourBytes = counts.HourBytes
rq.line.Counts.DayBytes = counts.DayBytes rq.line.Counts.DayBytes = counts.DayBytes
@@ -103,6 +82,29 @@ func (rq *request) countBytes() {
} }
} }
// countedBytes returns the request's bytes, once it has ended, as the
// byte limits and the anomaly thresholds count them: the response's body
// bytes, the request's, or both, as SWWAF_BYTES_COUNT says. For an
// upgraded connection, such as a WebSocket, which has closed by then, what
// it carried from the app counts with the response's and what it carried
// from the client with the request's.
func (rq *request) countedBytes() int64 {
response, request := rq.out.bytes, rq.requestBytes()
if rq.upgraded != nil {
response += rq.upgraded.fromApp.Load()
request += rq.upgraded.toApp.Load()
}
switch rq.h.config.BytesCount {
case "response":
return response
case "request":
return request
default: // both
return response + request
}
}
// banForLimit bans the client's netblock at now for a broken limit, the // banForLimit bans the client's netblock at now for a broken limit, the
// one hit names, and notes the offence for the log line. status is what // one hit names, and notes the offence for the log line. status is what
// the client was sent, or is sent: SWWAF_BAN_RESPONSE for a request over // the client was sent, or is sent: SWWAF_BAN_RESPONSE for a request over
+18 -9
View File
@@ -421,17 +421,28 @@ func TestBanForALoweredLimitGivesThePercentageInItsNotesAndItsAlert(t *testing.T
} }
} }
// startWithLookups is startAppWithAlerts in front of readAndAnswer, with // startWithLookups is startWithLookupsAndClock for a test that needs no
// the settings in env on top of clients looked up in a lookup database, // clock.
// which places fromDE and fromKP in the AS numbers and countries the
// stand-in for GeoJS gives them, noCountry in AS64500 and no country, and
// no other address. It returns the sender, the server and the queue of
// the alerts.
func startWithLookups( func startWithLookups(
t *testing.T, env map[string]string, t *testing.T, env map[string]string,
) (*sender, *proxy.Server, *alerts.Queue) { ) (*sender, *proxy.Server, *alerts.Queue) {
t.Helper() t.Helper()
s, _, server, queue := startWithLookupsAndClock(t, env)
return s, server, queue
}
// startWithLookupsAndClock is startAppWithAlerts in front of
// readAndAnswer, with the settings in env on top of clients looked up in a
// lookup database, which places fromDE and fromKP in the AS numbers and
// countries the stand-in for GeoJS gives them, noCountry in AS64500 and no
// country, and no other address.
func startWithLookupsAndClock(
t *testing.T, env map[string]string,
) (*sender, *clock, *proxy.Server, *alerts.Queue) {
t.Helper()
path := filepath.Join(t.TempDir(), "ipinfo_lite.mmdb") path := filepath.Join(t.TempDir(), "ipinfo_lite.mmdb")
lookuptest.Write(t, path, map[string]lookuptest.Network{ lookuptest.Write(t, path, map[string]lookuptest.Network{
fromDE + "/32": {ASN: asnDE, ASName: asNameDE, Country: "DE"}, fromDE + "/32": {ASN: asnDE, ASName: asNameDE, Country: "DE"},
@@ -442,9 +453,7 @@ func startWithLookups(
settings := map[string]string{lookupSource: fileSource, lookupDBPath: path} settings := map[string]string{lookupSource: fileSource, lookupDBPath: path}
maps.Copy(settings, env) maps.Copy(settings, env)
s, _, server, queue := startAppWithAlerts(t, readAndAnswer, settings) return startAppWithAlerts(t, readAndAnswer, settings)
return s, server, queue
} }
// uploadFrom is upload from the client at from. // uploadFrom is upload from the client at from.
+18 -1
View File
@@ -12,6 +12,7 @@ import (
"time" "time"
"sneak.berlin/go/smallwebwaf/internal/alerts" "sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/bans" "sneak.berlin/go/smallwebwaf/internal/bans"
"sneak.berlin/go/smallwebwaf/internal/config" "sneak.berlin/go/smallwebwaf/internal/config"
"sneak.berlin/go/smallwebwaf/internal/lookup" "sneak.berlin/go/smallwebwaf/internal/lookup"
@@ -68,7 +69,8 @@ type Params struct {
// against. // against.
Rules *rules.Files Rules *rules.Files
// Alerts receive the alert for each ban the proxy makes or makes // Alerts receive the alert for each ban the proxy makes or makes
// permanent, and for GeoJS failing. // permanent, for each count over an anomaly threshold, and for GeoJS
// failing.
Alerts *alerts.Queue Alerts *alerts.Queue
} }
@@ -81,6 +83,7 @@ type Server struct {
Ledger *bans.Ledger Ledger *bans.Ledger
Limiter *ratelimit.Limiter Limiter *ratelimit.Limiter
GeoJS *lookup.GeoJS GeoJS *lookup.GeoJS
Anomalies *anomaly.Counters
LookupFile *lookup.File LookupFile *lookup.File
Metrics *metrics.Metrics Metrics *metrics.Metrics
} }
@@ -117,6 +120,17 @@ func New(params Params) *Server {
AttackBanDuration: params.Config.AttackBanDuration, AttackBanDuration: params.Config.AttackBanDuration,
MaxBans: params.Config.MaxBans, MaxBans: params.Config.MaxBans,
}), }),
anomalies: anomaly.New(anomaly.Params{
Client: params.Config.AnomalyClient,
Net: params.Config.AnomalyNet,
ASN: params.Config.AnomalyASN,
Total: params.Config.AnomalyTotal,
Watch: params.Config.AnomalyWatch,
NetV4Prefix: params.Config.AnomalyNetV4Prefix,
NetV6Prefix: params.Config.AnomalyNetV6Prefix,
NamedNetblocks: params.Config.WatchNets,
Alerts: params.Alerts,
}),
lookupFile: params.LookupFile, lookupFile: params.LookupFile,
rules: params.Rules, rules: params.Rules,
alerts: params.Alerts, alerts: params.Alerts,
@@ -154,6 +168,7 @@ func New(params Params) *Server {
Ledger: h.ledger, Ledger: h.ledger,
Limiter: h.limiter, Limiter: h.limiter,
GeoJS: h.geojs, GeoJS: h.geojs,
Anomalies: h.anomalies,
LookupFile: h.lookupFile, LookupFile: h.lookupFile,
Metrics: m, Metrics: m,
} }
@@ -172,6 +187,7 @@ type handler struct {
limiter *ratelimit.Limiter limiter *ratelimit.Limiter
ledger *bans.Ledger ledger *bans.Ledger
geojs *lookup.GeoJS geojs *lookup.GeoJS
anomalies *anomaly.Counters
lookupFile *lookup.File lookupFile *lookup.File
rules *rules.Files rules *rules.Files
alerts *alerts.Queue alerts *alerts.Queue
@@ -211,6 +227,7 @@ func (h *handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// Once the request has ended, before its log line is written. // Once the request has ended, before its log line is written.
defer rq.addToHistory() defer rq.addToHistory()
defer rq.countAnomalies()
refused := rq.check(r.Context()) refused := rq.check(r.Context())
rq.checked = time.Now() rq.checked = time.Now()
+57 -20
View File
@@ -16,6 +16,8 @@ import (
"sync/atomic" "sync/atomic"
"time" "time"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/config"
"sneak.berlin/go/smallwebwaf/internal/lookup" "sneak.berlin/go/smallwebwaf/internal/lookup"
"sneak.berlin/go/smallwebwaf/internal/ratelimit" "sneak.berlin/go/smallwebwaf/internal/ratelimit"
"sneak.berlin/go/smallwebwaf/internal/requestlog" "sneak.berlin/go/smallwebwaf/internal/requestlog"
@@ -513,16 +515,14 @@ func timing(start, end time.Time) *float64 {
} }
// addToHistory adds the request, which has ended, to its client's // addToHistory adds the request, which has ended, to its client's
// history, and then, for a client that was looked up, the lookup // history, and then the lookup's answer about the client, as
// database's answer about it, or the answer from GeoJS kept about it, to // answerAtTheEnd gives it, to that history and to the notes of the bans
// that history and to the notes of the bans on its netblock: an answer // on its netblock: an answer may have come before either was there, and
// may have come before either was there, and one from GeoJS that comes // one from GeoJS that comes later is added when it comes.
// later is added when it comes.
func (rq *request) addToHistory() { func (rq *request) addToHistory() {
forwarded := !rq.upstreamStart.IsZero() forwarded := !rq.upstreamStart.IsZero()
group := clientGroup(rq.client)
rq.h.limiter.AddToHistory(group, rq.h.now(), ratelimit.Request{ rq.h.limiter.AddToHistory(clientGroup(rq.client), rq.h.now(), ratelimit.Request{
Forwarded: forwarded, Forwarded: forwarded,
Refused: !forwarded && rq.refused.Load() != nil, Refused: !forwarded && rq.refused.Load() != nil,
Status: rq.out.status, Status: rq.out.status,
@@ -531,23 +531,60 @@ func (rq *request) addToHistory() {
BrokeLimit: rq.line.Offence == requestlog.OffenceLimit, BrokeLimit: rq.line.Offence == requestlog.OffenceLimit,
}) })
if !rq.lookedUp { answer, found := rq.answerAtTheEnd()
return if found {
}
// The lookup database's answer was there at once.
if rq.h.config.LookupSource == "file" {
rq.h.addLookup(rq.lookupAnswer)
return
}
answer, kept := rq.h.geojs.Kept(group)
if kept {
rq.h.addLookup(answer) rq.h.addLookup(answer)
} }
} }
// countAnomalies counts the request, which has ended, and its bytes, as
// countedBytes gives them, for the anomaly thresholds, whatever was done
// with it: a request refused, one from a client in SWWAF_ALLOW_NETS or
// SWWAF_RATE_LIMIT_EXEMPT_NETS, and one for a path in
// SWWAF_RATE_LIMIT_EXEMPT_PATHS are counted too. It is counted for its
// client's AS number when answerAtTheEnd gives one. With every anomaly
// threshold off, the default, it does nothing.
func (rq *request) countAnomalies() {
if !anomalyThresholdsSet(rq.h.config) {
return
}
answer, _ := rq.answerAtTheEnd()
rq.h.anomalies.Count(rq.h.now(), anomaly.Request{
Client: rq.client,
ClientGroup: clientGroup(rq.client),
ASN: answer.ASN,
ASName: answer.ASName,
Country: answer.Country,
Bytes: rq.countedBytes(),
})
}
// anomalyThresholdsSet reports whether any anomaly threshold is set.
func anomalyThresholdsSet(cfg *config.Config) bool {
off := anomaly.Thresholds{}
return cfg.AnomalyClient != off || cfg.AnomalyNet != off || cfg.AnomalyASN != off ||
cfg.AnomalyTotal != off || cfg.AnomalyWatch != off
}
// answerAtTheEnd returns, for a client that was looked up, the lookup's
// answer about it as the request ends, and whether there is one: the
// lookup database's, which was there at once, or the one GeoJS has given
// by then, which a request does not wait for unless a setting needs it.
func (rq *request) answerAtTheEnd() (lookup.Answer, bool) {
if !rq.lookedUp {
return lookup.Answer{}, false
}
if rq.h.config.LookupSource == "file" {
return rq.lookupAnswer, true
}
return rq.h.geojs.Kept(clientGroup(rq.client))
}
// requestBytes is how many bytes of the request's body have been read. // requestBytes is how many bytes of the request's body have been read.
func (rq *request) requestBytes() int64 { func (rq *request) requestBytes() int64 {
if rq.body == nil { if rq.body == nil {
+17 -11
View File
@@ -361,9 +361,7 @@ func (l *Limiter) Load(clients []Client, now time.Time) {
for _, c := range clients { for _, c := range clients {
for i, w := range l.windows { for i, w := range l.windows {
for _, b := range []*Buckets{c.buckets()[i], c.byteBuckets()[i]} { for _, b := range []*Buckets{c.buckets()[i], c.byteBuckets()[i]} {
// The window that ends at now covers neither bucket once it if b.Passed(now, w.length) {
// begins after the bucket under way has ended.
if !now.Add(-w.length).Before(b.Start.Add(w.length)) {
*b = Buckets{} *b = Buckets{}
} }
} }
@@ -393,8 +391,8 @@ func (l *Limiter) count(
) )
for i, w := range l.windows { for i, w := range l.windows {
requestCounts[i] = requestBuckets[i].add(now, w.length, requests) requestCounts[i] = requestBuckets[i].Add(now, w.length, requests)
byteCounts[i] = byteBuckets[i].add(now, w.length, bytes) byteCounts[i] = byteBuckets[i].Add(now, w.length, bytes)
limit, byteLimit := percentOf(w.limit, percent), percentOf(w.byteLimit, percent) limit, byteLimit := percentOf(w.limit, percent), percentOf(w.byteLimit, percent)
switch { switch {
@@ -460,18 +458,18 @@ func percentOf(limit, percent int64) int64 {
return limit/hundred*percent + limit%hundred*percent/hundred return limit/hundred*percent + limit%hundred*percent/hundred
} }
// add counts n requests, or n bytes, at now in a window of length, and // Add counts n requests, or n bytes, at now in a window of length, and
// returns the client's count in the window that ends at now: what is in // returns the count in the window that ends at now: what is in the bucket
// the bucket under way, and what is in the bucket before it weighted by // under way, and what is in the bucket before it weighted by how much of
// how much of that bucket the window still covers. With n zero it counts // that bucket the window still covers. With n zero it counts nothing, and
// nothing, and returns the count. // returns the count. The anomaly counters count in Buckets too.
// //
// Concurrent requests can be counted out of order, so now can be a moment // Concurrent requests can be counted out of order, so now can be a moment
// before the bucket under way began; such a request is counted in that // before the bucket under way began; such a request is counted in that
// bucket. A request dated more than a second before it means the clock // bucket. A request dated more than a second before it means the clock
// was set back, and the buckets start afresh: otherwise the bucket before // was set back, and the buckets start afresh: otherwise the bucket before
// would keep its full weight until the clock caught up. // would keep its full weight until the clock caught up.
func (b *Buckets) add(now time.Time, length time.Duration, n int64) float64 { func (b *Buckets) Add(now time.Time, length time.Duration, n int64) float64 {
if now.Before(b.Start.Add(-time.Second)) { if now.Before(b.Start.Add(-time.Second)) {
*b = Buckets{} *b = Buckets{}
} }
@@ -496,6 +494,14 @@ func (b *Buckets) add(now time.Time, length time.Duration, n int64) float64 {
return float64(b.Previous)*covered + float64(b.Current) return float64(b.Previous)*covered + float64(b.Current)
} }
// Passed reports whether the window of length that ends at now covers
// neither of b's buckets: it begins after the bucket under way has ended.
// What they hold then counts no more, and a state file read at now drops
// it.
func (b *Buckets) Passed(now time.Time, length time.Duration) bool {
return !now.Add(-length).Before(b.Start.Add(length))
}
// add counts a response with status in its class. A status of 0, for // add counts a response with status in its class. A status of 0, for
// nothing sent, is not a response. // nothing sent, is not a response.
func (r *Responses) add(status int) { func (r *Responses) add(status int) {
+1
View File
@@ -201,6 +201,7 @@ func loadStateFiles(
Limiter: server.Limiter, Limiter: server.Limiter,
GeoJS: server.GeoJS, GeoJS: server.GeoJS,
Alerts: alertQueue, Alerts: alertQueue,
Anomalies: server.Anomalies,
Now: now, Now: now,
ProcessLog: processLog, ProcessLog: processLog,
Metrics: server.Metrics, Metrics: server.Metrics,
+73 -15
View File
@@ -2,7 +2,8 @@
// SWWAF_STATE_DIR, as the "Persistent state" section of SPEC.md describes: // SWWAF_STATE_DIR, as the "Persistent state" section of SPEC.md describes:
// bans.json holds the bans, clients.json each client's counters and // bans.json holds the bans, clients.json each client's counters and
// history, lookups.json GeoJS's answers, and alerts.json the cooldowns, // history, lookups.json GeoJS's answers, and alerts.json the cooldowns,
// the hour under way and the alerts waiting for each destination. Load // the hour under way, the alerts waiting for each destination and the
// anomaly counters. Load
// reads them at start, Watch takes in an admin's edit of one while // reads them at start, Watch takes in an admin's edit of one while
// smallwebwaf runs, and Run and WriteAll write them. The disk is read and // smallwebwaf runs, and Run and WriteAll write them. The disk is read and
// written outside the parts' locks, which are held only to take a // written outside the parts' locks, which are held only to take a
@@ -29,6 +30,7 @@ import (
"github.com/fsnotify/fsnotify" "github.com/fsnotify/fsnotify"
"sneak.berlin/go/smallwebwaf/internal/alerts" "sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/bans" "sneak.berlin/go/smallwebwaf/internal/bans"
"sneak.berlin/go/smallwebwaf/internal/lookup" "sneak.berlin/go/smallwebwaf/internal/lookup"
"sneak.berlin/go/smallwebwaf/internal/metrics" "sneak.berlin/go/smallwebwaf/internal/metrics"
@@ -56,6 +58,7 @@ var (
errMissing = errors.New("has no") errMissing = errors.New("has no")
errCause = errors.New("is not limit, attack or admin") errCause = errors.New("is not limit, attack or admin")
errDestination = errors.New("is not webhook, slack or ntfy") errDestination = errors.New("is not webhook, slack or ntfy")
errScope = errors.New("is not client, net, asn, total or watch")
errWaitingList = errors.New(`waiting is a list, but now lists the alerts by ` + errWaitingList = errors.New(`waiting is a list, but now lists the alerts by ` +
`destination: put the list under "webhook", as "waiting": {"webhook": [...]}, ` + `destination: put the list under "webhook", as "waiting": {"webhook": [...]}, ` +
`or remove the file`) `or remove the file`)
@@ -70,13 +73,14 @@ type Params struct {
// is (SWWAF_STATE_COUNTER_INTERVAL). // is (SWWAF_STATE_COUNTER_INTERVAL).
WriteDelay time.Duration WriteDelay time.Duration
CounterInterval time.Duration CounterInterval time.Duration
// Ledger, Limiter, GeoJS and Alerts hold the state. Alerts also // Ledger, Limiter, GeoJS, Alerts and Anomalies hold the state. Alerts
// receive a file_error alert for an edit set aside, and for a write // also receive a file_error alert for an edit set aside, and for a
// that fails while smallwebwaf runs. // write that fails while smallwebwaf runs.
Ledger *bans.Ledger Ledger *bans.Ledger
Limiter *ratelimit.Limiter Limiter *ratelimit.Limiter
GeoJS *lookup.GeoJS GeoJS *lookup.GeoJS
Alerts *alerts.Queue Alerts *alerts.Queue
Anomalies *anomaly.Counters
// Now tells the time by which the counters' buckets run out, normally // Now tells the time by which the counters' buckets run out, normally
// time.Now in UTC. // time.Now in UTC.
Now func() time.Time Now func() time.Time
@@ -135,11 +139,14 @@ type lookupsFile struct {
} }
// alertsFile is alerts.json, indented for an admin to read and edit. // alertsFile is alerts.json, indented for an admin to read and edit.
//
//nolint:tagliatelle // the state files use snake_case, as the request log does
type alertsFile struct { type alertsFile struct {
Version int `json:"version"` Version int `json:"version"`
Cooldowns []alerts.Cooldown `json:"cooldowns"` Cooldowns []alerts.Cooldown `json:"cooldowns"`
Hour alerts.Hour `json:"hour"` Hour alerts.Hour `json:"hour"`
Waiting map[string][]alerts.Alert `json:"waiting"` Waiting map[string][]alerts.Alert `json:"waiting"`
AnomalyCounters []anomaly.Counter `json:"anomaly_counters"`
} }
// stateFile is the struct of a state file. Once the file is decoded, its // stateFile is the struct of a state file. Once the file is decoded, its
@@ -420,6 +427,7 @@ func (f *Files) takeIn(name string, data []byte, edit bool) (int, error) {
f.params.Alerts.Load(alerts.State{ f.params.Alerts.Load(alerts.State{
Cooldowns: file.Cooldowns, Hour: file.Hour, Waiting: file.Waiting, Cooldowns: file.Cooldowns, Hour: file.Hour, Waiting: file.Waiting,
}) })
f.params.Anomalies.Load(file.AnomalyCounters, f.params.Now())
for _, waiting := range file.Waiting { for _, waiting := range file.Waiting {
entries += len(waiting) entries += len(waiting)
@@ -522,7 +530,7 @@ func (f *Files) encode(name string) ([]byte, error) {
held := f.params.Alerts.Snapshot() held := f.params.Alerts.Snapshot()
file := alertsFile{ file := alertsFile{
Version: version, Cooldowns: held.Cooldowns, Hour: held.Hour, Version: version, Cooldowns: held.Cooldowns, Hour: held.Hour,
Waiting: held.Waiting, Waiting: held.Waiting, AnomalyCounters: f.params.Anomalies.Snapshot(),
} }
data, err := json.MarshalIndent(file, "", " ") data, err := json.MarshalIndent(file, "", " ")
@@ -673,8 +681,9 @@ func (f *lookupsFile) check(data []byte) error {
// check refuses a cooldown without its event or when its alert was sent, // check refuses a cooldown without its event or when its alert was sent,
// which would hold back no repeat, alerts waiting for a destination with // which would hold back no repeat, alerts waiting for a destination with
// another name than webhook, slack or ntfy, most likely misspelt, and an // another name than webhook, slack or ntfy, most likely misspelt, an
// alert waiting without its event or its time. // alert waiting without its event or its time, and an anomaly counter as
// checkAnomalyCounters does.
func (f *alertsFile) check([]byte) error { func (f *alertsFile) check([]byte) error {
for i, cooldown := range f.Cooldowns { for i, cooldown := range f.Cooldowns {
switch { switch {
@@ -700,9 +709,58 @@ func (f *alertsFile) check([]byte) error {
} }
} }
return checkAnomalyCounters(f.AnomalyCounters)
}
// checkAnomalyCounters refuses an anomaly counter whose scope is not
// client, net, asn, total or watch, most likely misspelt, and one without
// a field it needs, as missingFromCounter tells.
func checkAnomalyCounters(counters []anomaly.Counter) error {
for i, counter := range counters {
if !slices.Contains(anomaly.Scopes(), counter.Scope) {
return fmt.Errorf("anomaly_counters entry %d's scope %q %w", i+1,
counter.Scope, errScope)
}
field := missingFromCounter(counter)
if field != "" {
return fmt.Errorf("anomaly_counters %w", missing(i, field))
}
}
return nil return nil
} }
// missingFromCounter returns the first field counter, an anomaly counter,
// needs and has not, or "" when it has them all: what tells it from the
// others in its scope, without which it would never be counted again, the
// netblock of a client, net or watch counter, the AS number of an asn one
// and the name of a watch one; and the start of a window in which it has
// requests or bytes, without which they would be dropped.
func missingFromCounter(counter anomaly.Counter) string {
scope := counter.Scope
switch {
case scope != anomaly.ScopeASN && scope != anomaly.ScopeTotal &&
!counter.Netblock.IsValid():
return "netblock"
case scope == anomaly.ScopeASN && counter.ASN == "":
return "asn"
case scope == anomaly.ScopeWatch && counter.Name == "":
return "name"
case countsWithoutStart(counter.Minute):
return "minute.start"
case countsWithoutStart(counter.Hour):
return "hour.start"
case countsWithoutStart(counter.MinuteBytes):
return "minute_bytes.start"
case countsWithoutStart(counter.HourBytes):
return "hour_bytes.start"
default:
return ""
}
}
// countsWithoutStart reports whether b holds requests, or bytes, but no // countsWithoutStart reports whether b holds requests, or bytes, but no
// start, which places them in time. // start, which places them in time.
func countsWithoutStart(b ratelimit.Buckets) bool { func countsWithoutStart(b ratelimit.Buckets) bool {
+161 -11
View File
@@ -22,6 +22,7 @@ import (
"time" "time"
"sneak.berlin/go/smallwebwaf/internal/alerts" "sneak.berlin/go/smallwebwaf/internal/alerts"
"sneak.berlin/go/smallwebwaf/internal/anomaly"
"sneak.berlin/go/smallwebwaf/internal/bans" "sneak.berlin/go/smallwebwaf/internal/bans"
"sneak.berlin/go/smallwebwaf/internal/lookup" "sneak.berlin/go/smallwebwaf/internal/lookup"
"sneak.berlin/go/smallwebwaf/internal/metrics" "sneak.berlin/go/smallwebwaf/internal/metrics"
@@ -156,7 +157,50 @@ const filledAlertsJSON = `{
"suppressed_repeats": 0 "suppressed_repeats": 0
} }
] ]
} },
"anomaly_counters": [
{
"scope": "asn",
"asn": "AS64496",
"hour_bytes": {
"start": "2026-10-06T00:00:00Z",
"current": 8,
"previous": 0
}
},
{
"scope": "net",
"netblock": "203.0.113.0/24",
"minute": {
"start": "2026-10-06T00:00:00Z",
"current": 1,
"previous": 0
}
},
{
"scope": "total",
"minute": {
"start": "2026-10-06T00:00:00Z",
"current": 1,
"previous": 0
},
"minute_bytes": {
"start": "2026-10-06T00:00:00Z",
"current": 8,
"previous": 0
}
},
{
"scope": "watch",
"netblock": "203.0.113.0/24",
"name": "office",
"hour": {
"start": "2026-10-06T00:00:00Z",
"current": 1,
"previous": 0
}
}
]
} }
` `
@@ -191,6 +235,8 @@ func TestFilesWrittenAndReadBack(t *testing.T) {
t.Errorf("%s read back\n%+v\nwant\n%+v", alertsJSON, got, want) t.Errorf("%s read back\n%+v\nwant\n%+v", alertsJSON, got, want)
} }
wantEqual(t, alertsJSON, after.Anomalies.Snapshot(), before.Anomalies.Snapshot())
// Each one-per-line file lists its entries by client, and nothing // Each one-per-line file lists its entries by client, and nothing
// but the four files is left in the directory. // but the four files is left in the directory.
wantEntries(t, filepath.Join(dir, clientsJSON), "clients", wantEntries(t, filepath.Join(dir, clientsJSON), "clients",
@@ -251,6 +297,49 @@ func TestSourceFailureCooldownKeptInAlertsJSONAcrossARestart(t *testing.T) {
} }
} }
func TestAnomalyCountersKeptInAlertsJSONAcrossARestart(t *testing.T) {
t.Parallel()
dir := t.TempDir()
request := anomaly.Request{
Client: netip.MustParseAddr("203.0.113.9"),
ClientGroup: netip.MustParsePrefix("203.0.113.9/32"),
}
// The whole service may have two requests a minute.
withThreshold := func() state.Params {
params := newParams(dir)
params.Anomalies = anomaly.New(anomaly.Params{
Total: anomaly.Thresholds{RequestsPerMinute: 2}, Alerts: params.Alerts,
})
return params
}
before := withThreshold()
files := load(t, before)
for range 2 {
before.Anomalies.Count(midnight(), request)
}
err := files.WriteAll()
if err != nil {
t.Fatalf("write: %v", err)
}
// After the restart, the third request in the minute is over it.
after := withThreshold()
load(t, after)
after.Anomalies.Count(midnight(), request)
waiting := after.Alerts.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || waiting[0].Event != alerts.EventAnomaly ||
waiting[0].Detail["count"] != float64(3) {
t.Errorf("alerts wait %+v, want one for 3 requests", waiting)
}
}
func TestBansJSONIsIndentedWithNullForAPermanentBan(t *testing.T) { func TestBansJSONIsIndentedWithNullForAPermanentBan(t *testing.T) {
t.Parallel() t.Parallel()
@@ -332,6 +421,13 @@ func TestFileThatDoesNotParseStopsTheStart(t *testing.T) {
`{"version": 1, "waiting": {"webhook": [], "slak": []}}`, `{"version": 1, "waiting": {"webhook": [], "slak": []}}`,
`: waiting "slak" is not webhook, slack or ntfy`, `: waiting "slak" is not webhook, slack or ntfy`,
}, },
{
"an anomaly counter of an unknown scope", alertsJSON,
`{"version": 1, "anomaly_counters": [{"scope": "total"}, ` +
`{"scope": "nett", "netblock": "203.0.113.0/24"}]}`,
`: anomaly_counters entry 2's scope "nett" is not client, net, asn, total ` +
`or watch`,
},
} { } {
t.Run(tc.name, func(t *testing.T) { t.Run(tc.name, func(t *testing.T) {
t.Parallel() t.Parallel()
@@ -455,6 +551,39 @@ func TestAlertsJSONEntryWithoutAFieldItNeedsStopsTheStart(t *testing.T) {
`{"version": 1, "waiting": {"slack": [{"event": "ban"}]}}`, `{"version": 1, "waiting": {"slack": [{"event": "ban"}]}}`,
`: waiting slack entry 1 has no "time"`, `: waiting slack entry 1 has no "time"`,
}, },
{
// The whole service's counter needs nothing to tell it apart.
"an anomaly counter of a netblock without it",
`{"version": 1, "anomaly_counters": [{"scope": "total"}, {"scope": "net"}]}`,
`: anomaly_counters entry 2 has no "netblock"`,
},
{
"an anomaly counter of a client without its netblock",
`{"version": 1, "anomaly_counters": [{"scope": "client"}]}`,
`: anomaly_counters entry 1 has no "netblock"`,
},
{
"an anomaly counter of an AS number without it",
`{"version": 1, "anomaly_counters": [{"scope": "asn"}]}`,
`: anomaly_counters entry 1 has no "asn"`,
},
{
"an anomaly counter of a named netblock without its name",
`{"version": 1, "anomaly_counters": [` +
`{"scope": "watch", "netblock": "203.0.113.0/24"}]}`,
`: anomaly_counters entry 1 has no "name"`,
},
{
"an anomaly counter of a named netblock without its netblock",
`{"version": 1, "anomaly_counters": [{"scope": "watch", "name": "office"}]}`,
`: anomaly_counters entry 1 has no "netblock"`,
},
{
"an anomaly counter with bytes in a window without its start",
`{"version": 1, "anomaly_counters": [` +
`{"scope": "total", "hour_bytes": {"current": 5}}]}`,
`: anomaly_counters entry 1 has no "hour_bytes.start"`,
},
} { } {
t.Run(tc.name, func(t *testing.T) { t.Run(tc.name, func(t *testing.T) {
t.Parallel() t.Parallel()
@@ -1273,10 +1402,19 @@ func midnight() time.Time {
// newParams returns Params for the state files in dir, with parts that // newParams returns Params for the state files in dir, with parts that
// hold nothing yet. GeoJS is never asked, and the alerts, at most two an // hold nothing yet. GeoJS is never asked, and the alerts, at most two an
// hour, are never sent. // hour, are never sent. The anomaly counters count the scopes fill
// counts, with thresholds fill does not reach.
func newParams(dir string) state.Params { func newParams(dir string) state.Params {
discard := slog.New(slog.DiscardHandler) discard := slog.New(slog.DiscardHandler)
m := metrics.New(1, "app") m := metrics.New(1, "app")
queue := alerts.New(alerts.Params{
WebhookURL: &url.URL{Scheme: "https", Host: "alerts.example"},
Events: alerts.Events(),
Cooldown: 15 * time.Minute,
MaxPerHour: 2,
Instance: "fsn1app1/gitea",
Now: midnight,
})
return state.Params{ return state.Params{
Dir: dir, Dir: dir,
@@ -1293,13 +1431,16 @@ func newParams(dir string) state.Params {
GeoJS: lookup.New(lookup.Params{ GeoJS: lookup.New(lookup.Params{
Now: midnight, ProcessLog: discard, Metrics: m, Now: midnight, ProcessLog: discard, Metrics: m,
}), }),
Alerts: alerts.New(alerts.Params{ Alerts: queue,
WebhookURL: &url.URL{Scheme: "https", Host: "alerts.example"}, Anomalies: anomaly.New(anomaly.Params{
Events: alerts.Events(), Net: anomaly.Thresholds{RequestsPerMinute: 1000},
Cooldown: 15 * time.Minute, ASN: anomaly.Thresholds{BytesPerHour: 1 << 30},
MaxPerHour: 2, Total: anomaly.Thresholds{RequestsPerMinute: 1000, BytesPerMinute: 1 << 30},
Instance: "fsn1app1/gitea", Watch: anomaly.Thresholds{RequestsPerHour: 1000},
Now: midnight, NetV4Prefix: 24,
NetV6Prefix: 48,
NamedNetblocks: []anomaly.NamedNetblock{{Name: "office", Netblock: office()}},
Alerts: queue,
}), }),
Now: midnight, Now: midnight,
ProcessLog: discard, ProcessLog: discard,
@@ -1307,10 +1448,15 @@ func newParams(dir string) state.Params {
} }
} }
// office is the named netblock of the anomaly counters of newParams.
func office() netip.Prefix {
return netip.MustParsePrefix("203.0.113.0/24")
}
// fill puts a permanent ban an admin made, a ban for a broken limit and // fill puts a permanent ban an admin made, a ban for a broken limit and
// one for a clear sign of attack, clients with counts and histories, // one for a clear sign of attack, clients with counts and histories,
// GeoJS answers, and alerts, as filledAlertsJSON holds them, into the // GeoJS answers, and alerts and anomaly counters, as filledAlertsJSON
// parts of params. // holds them, into the parts of params.
func fill(params state.Params) { func fill(params state.Params) {
now := midnight() now := midnight()
client := netip.MustParsePrefix("203.0.113.9/32") client := netip.MustParsePrefix("203.0.113.9/32")
@@ -1361,6 +1507,10 @@ func fill(params state.Params) {
params.Alerts.Raise(alerts.Alert{ params.Alerts.Raise(alerts.Alert{
Event: alerts.EventSourceFailure, Reason: "asking GeoJS failed", Event: alerts.EventSourceFailure, Reason: "asking GeoJS failed",
}) })
params.Anomalies.Count(now, anomaly.Request{
Client: client.Addr(), ClientGroup: client, ASN: asn, Bytes: 8,
})
} }
// permanentBan is the ban permanentBansJSON holds. // permanentBan is the ban permanentBansJSON holds.