Alerts to Slack and ntfy, each destination with its own queue #94

Merged
clawbot merged 1 commits from issue-90-slack-ntfy into next 2026-10-07 05:54:35 +02:00
14 changed files with 1586 additions and 428 deletions
+126 -58
View File
@@ -19,21 +19,21 @@ ledger with the bans you make, keep and lift, the JSON state files with your
edits taken in while it runs and the paths the rate limits do not count, which
come next in the build order, `observe` mode and the rest of the request log's
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 four parts of the
stage after it: the rule files, the first part, with the bans for a clear sign
of attack, the other admin endpoints, the second, alerts to a JSON webhook, the
first of the three destinations alerts go to, and remote log sending.
`smallwebwaf` passes each request to the app and the app's answer back,
unchanged, within its timeouts and size limits, works out each client's address,
bans a client that sends too many requests, not counting those for the paths you
choose, refuses a client that comes from a country you refuse or from a network
you refuse, lets the networks you choose through, checks each request against
the rule files and bans a client whose request is a clear sign of attack, keeps
its bans, each client's counters and history, and GeoJS's answers in JSON files
across restarts, takes in your edits of those files, such as a ban you make,
keep or lift, and of the rule files while it runs, writes a JSON log line for
every request, sends its log lines to a syslog server too if you name one, sends
an alert to a webhook you name for each ban it makes or makes permanent, for
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 other admin endpoints, alerts to all three destinations, a JSON webhook,
Slack and ntfy, and remote log sending. `smallwebwaf` passes each request to the
app and the app's answer back, unchanged, within its timeouts and size limits,
works out each client's address, bans a client that sends too many requests, not
counting those for the paths you choose, refuses a client that comes from a
country you refuse or from a network you refuse, lets the networks you choose
through, checks each request against the rule files and bans a client whose
request is a clear sign of attack, keeps its bans, each client's counters and
history, and GeoJS's answers in JSON files across restarts, takes in your edits
of those files, such as a ban you make, keep or lift, and of the rule files
while it runs, writes a JSON log line for every request, sends its log lines to
a syslog server too if you name one, sends an alert to a webhook, to Slack and
to ntfy, each if you name one, for each ban it makes or makes permanent, for
GeoJS failing and for a rule file or state file with an error, serves Prometheus
metrics to a scraper that holds the metrics token, lets an admin who holds the
admin token list, add and lift bans and ask what it knows of a client, and in
@@ -190,10 +190,12 @@ 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
`SWWAF_LOG_REMOTE_URL` names one (see "Sending the log to a syslog server"
below).
- Sends an alert, as a JSON object, to the webhook `SWWAF_ALERT_WEBHOOK_URL`
names, while it names one, for each ban it makes or makes permanent, for GeoJS
failing, and for a rule file or state file with an error, holding back repeats
and, past an hourly limit, rolling the rest into one summary (see "Alerts"
- Sends an alert for each ban it makes or makes permanent, for GeoJS failing,
and for a rule file or state file with an error, holding back repeats and,
past an hourly limit, rolling the rest into one summary, to each destination
you name: as a JSON object to the webhook `SWWAF_ALERT_WEBHOOK_URL` names, as
a message to the Slack incoming webhook `SWWAF_ALERT_SLACK_WEBHOOK_URL` names,
and as a message to the ntfy topic `SWWAF_ALERT_NTFY_URL` names (see "Alerts"
below).
## Settings
@@ -340,13 +342,29 @@ effective settings are logged at start.
- `SWWAF_ALERT_WEBHOOK_URL` (default unset): the webhook each alert is posted
to, an `http` or `https` URL without a user or a fragment, such as
`https://alerts.example/smallwebwaf` (see "Alerts" below). Unset or empty, no
alert is sent. Since many webhooks carry their secret in the path or the
query, the settings logged at start show `********` in place of them, and a
value that stops the start is not shown.
alert is sent to a webhook. Since many webhooks carry their secret in the path
or the query, the settings logged at start show `********` in place of them,
and a value that stops the start is not shown.
- `SWWAF_ALERT_WEBHOOK_HEADERS` (default empty): headers sent with each alert,
such as one that authenticates it, as a list of a name, `:` and a value, such
as `Authorization:Bearer 0123456789abcdef`. A value cannot hold a comma. The
settings logged at start show `********` in place of each value.
- `SWWAF_ALERT_SLACK_WEBHOOK_URL` (default unset): the Slack incoming webhook
each alert is posted to as a message, such as
`https://hooks.slack.com/services/T0123/B4567/abcdef`. Unset or empty, no
alert is sent to Slack. It is checked and logged as `SWWAF_ALERT_WEBHOOK_URL`
is.
- `SWWAF_ALERT_NTFY_URL` (default unset): the ntfy topic each alert is published
to, as the topic's full URL, such as `https://ntfy.sh/my-alerts`. Unset or
empty, no alert is sent to ntfy. It is checked and logged as
`SWWAF_ALERT_WEBHOOK_URL` is, since anyone who knows a topic on a server open
to all can read it.
- `SWWAF_ALERT_NTFY_TOKEN` (default unset): an ntfy access token, sent to ntfy
with each alert as `Authorization: Bearer <token>`, for a topic that needs
one. The settings logged at start show `********` in its place. A control
character in it, such as the carriage return of a file saved with Windows line
ends, stops the start, and while `SWWAF_ALERT_NTFY_URL` is set, so does one in
`SWWAF_INSTANCE_NAME`, which ntfy is sent in the title.
- `SWWAF_ALERT_EVENTS` (default
`ban,permanent_ban,waf_block,anomaly,reputation_hit,source_failure,file_error`):
the events alerts are sent for. `waf_block`, `anomaly` and `reputation_hit`
@@ -532,11 +550,13 @@ them.
## Alerts
While `SWWAF_ALERT_WEBHOOK_URL` is set, `smallwebwaf` posts each alert to it as
one JSON object, with `Content-Type: application/json` and the headers
`SWWAF_ALERT_WEBHOOK_HEADERS` gives, as "Alert webhook schema" in
[`SPEC.md`](SPEC.md) describes. An alert is for one of these events, and is sent
when `SWWAF_ALERT_EVENTS` names its event:
`smallwebwaf` sends each alert to each destination you name: to
`SWWAF_ALERT_WEBHOOK_URL` it posts the alert as one JSON object, with
`Content-Type: application/json` and the headers `SWWAF_ALERT_WEBHOOK_HEADERS`
gives, as "Alert webhook schema" in [`SPEC.md`](SPEC.md) describes, and to
`SWWAF_ALERT_SLACK_WEBHOOK_URL` and `SWWAF_ALERT_NTFY_URL` a message made from
it, as below. An alert is for one of these events, and is sent when
`SWWAF_ALERT_EVENTS` names its event:
- `ban`: a ban `smallwebwaf` makes, for a broken rate limit or a clear sign of
attack.
@@ -612,6 +632,47 @@ is sent on one line:
- `suppressed_repeats` is how many repeats the cooldown held back before this
alert.
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
for each of the `client`, the `netblock` and the `country`, the `file`, the
`source`, the `error` and the `mode` the `detail` gives, and the
`suppressed_repeats`, leaving out those that are empty or 0. Slack is posted, as
JSON, the title in bold and the text below it, with `&`, `<` and `>` escaped, so
that nothing in them is read as a link or a mention. ntfy is posted the text,
with the title as `Title`, `SWWAF_ALERT_NTFY_TOKEN`, while it is set, as
`Authorization: Bearer <token>`, and the priority and the tag, which ntfy shows
as an emoji, of the alert's event, as `Priority` and `Tags`:
| Event | Priority | Tag |
| ------------------------------ | --------- | -------------------------- |
| `ban` | `default` | `no_entry` |
| `permanent_ban` | `high` | `no_entry` |
| `waf_block` | `default` | `shield` |
| `anomaly` | `high` | `chart_with_upwards_trend` |
| `reputation_hit` | `low` | `label` |
| `source_failure`, `file_error` | `high` | `warning` |
| `summary` | `default` | `bar_chart` |
For the alert above, ntfy is sent these headers and this text:
```text
Title: fsn1app1/gitea: ban
Priority: default
Tags: no_entry
requests per minute over the limit of 1000
client: 203.0.113.9
netblock: 203.0.113.9/32
```
and Slack this JSON object, shown indented; it is sent on one line:
```json
{
"text": "*fsn1app1/gitea: ban*\nrequests per minute over the limit of 1000\nclient: 203.0.113.9\nnetblock: 203.0.113.9/32"
}
```
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
source, less than `SWWAF_ALERT_COOLDOWN` after it, is a repeat: it is held back
@@ -626,18 +687,19 @@ were held back, and 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
alert sent for the same event and netblock, file or source.
The alerts wait in a queue of at most 1000, from which they are sent one at a
time, the oldest first, so a webhook that is slow or down never holds up a
request. The webhook takes an alert by answering with a 2xx status, and refuses
it with a 4xx status other than `408` and `429`: a refused alert is logged,
counted as dropped, and given up, so that the next is sent. Any other answer, a
redirect included, a connection that fails, or no answer within 10 seconds is a
failure: it is logged, without the webhook's URL, and the alert is sent again a
second later, twice as long after each further failure in a row, up to a minute.
With 1000 alerts waiting, the oldest is dropped to make room for a new one. The
cooldowns, the hour under way and the alerts still waiting are kept in
`alerts.json` (see "State files" below), so that after a restart the alerts
waiting are sent, and the cooldowns go on.
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
down holds up neither the others nor any request. A destination takes an alert
by answering with a 2xx status, and refuses it with a 4xx status other than
`408` and `429`: a refused alert is logged, counted as dropped, and given up, so
that the next is sent. Any other answer, a redirect included, a connection that
fails, or no answer within 10 seconds is a failure: it is logged, naming the
destination's setting and not its URL, and the alert is sent again a second
later, twice as long after each further failure in a row, up to a minute. With
1000 alerts waiting for a destination, the oldest is dropped to make room for a
new one. The cooldowns, the hour under way and the alerts still waiting for each
destination are kept in `alerts.json` (see "State files" below), so that after a
restart the alerts waiting are sent, and the cooldowns go on.
## State files
@@ -668,9 +730,13 @@ entries by client address, but for the alerts waiting, with times in UTC.
`source`, or event alone, when the last alert was sent, `sent`, and the
repeats held back since, `suppressed_repeats`; under `hour`, the hour under
way, from its `start`, the alerts `sent` in it and those `held_back` for its
summary, by event; and under `waiting`, the alerts still waiting to be sent,
the oldest first, each as the webhook is sent it. As an hour ends, the
cooldowns that have run out with no repeat held back are dropped.
summary, by event; and under `waiting`, for each destination you name,
`webhook`, `slack` or `ntfy`, the alerts still waiting to be sent to it, the
oldest first, each as the webhook is sent it. As an hour ends, the cooldowns
that have run out with no repeat held back are dropped. 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
through `DELETE /_smallwebwaf/bans/<client>`, or made permanent, with every such
@@ -695,8 +761,9 @@ without a field it needs, named with the entry's place in the file: a ban's
client's `client`, or the `start` of a window in which it has requests; an
answer's `client`, `country`, which is `""` for a client GeoJS cannot place, or
`answered`; a cooldown's `event` or `sent`; an alert waiting's `event` or
`time`. So does a ban whose `cause` is not `limit`, `attack` or `admin`. The AS
number and AS name come with their lookup.
`time`. So does a ban whose `cause` is not `limit`, `attack` or `admin`, and
alerts waiting for a destination that is not `webhook`, `slack` or `ntfy`. The
AS number and AS name come with their lookup.
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
@@ -705,14 +772,14 @@ yours by comparing the file with what it last read or wrote, and before it
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
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 or gives a
ban another `cause`, does not stop the running `smallwebwaf`: it keeps what it
holds, and at the file's next write renames your file to `<name>.bad`, such as
`bans.json.bad`, writes the file again from memory, logs the file and where the
error is, and raises a `file_error` alert for it. It waits for that write
because an editor's file can be read before the editor has finished writing it.
Mend the `.bad` file and move it back. A file you remove is written again at its
next write.
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
`smallwebwaf`: it keeps what it holds, and at the file's next write renames your
file to `<name>.bad`, such as `bans.json.bad`, writes the file again from
memory, logs the file and where the error is, and raises a `file_error` alert
for it. It waits for that write because an editor's file can be read before the
editor has finished writing it. 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`
and its `expires`, `null` for a ban that never ends; its `reason` and its
@@ -876,12 +943,13 @@ other request. No metric carries a client's address.
`smallwebwaf_remote_log_lines_dropped_total`: those dropped, from a full
buffer or because their sending failed; and
`smallwebwaf_remote_log_buffer_depth`: those waiting in the buffer.
- While `SWWAF_ALERT_WEBHOOK_URL` is set, by `destination`, `webhook`:
`smallwebwaf_alerts_sent_total`: the alerts the webhook took;
- For each destination you name, by `destination`, `webhook`, `slack` or `ntfy`:
`smallwebwaf_alerts_sent_total`: the alerts it took;
`smallwebwaf_alerts_failed_total`: the requests to it that failed;
`smallwebwaf_alerts_suppressed_total`: the alerts held back, as repeats or for
an hour's summary; and `smallwebwaf_alerts_dropped_total`: those dropped from
a full queue, or given up as the webhook refused them.
an hour's summary, which are the same for every destination; and
`smallwebwaf_alerts_dropped_total`: those dropped from its full queue, or
given up as it refused them.
- Go's own `go_` metrics and the process's `process_` metrics.
The requests Go's HTTP server ends before `smallwebwaf` sees them (see "Request
@@ -1259,8 +1327,8 @@ addresses are never sent to GeoJS.
standard library alone, whose `log/syslog` writes only the older syslog
format.
- `internal/alerts`: takes the alerts the other parts raise, holds back repeats
and those past the hourly limit, and sends the others to
`SWWAF_ALERT_WEBHOOK_URL` from a queue of its own.
and those past the hourly limit, and sends the others to the webhook, Slack
and ntfy, each from a queue of its own.
- `Dockerfile`: the lint and test phases, then the image, whose last stage
installs Ubuntu's packages, nixpkgs, `runsvinit` and `smallwebwaf`, with
`share/smallwebwaf.run` as runit's `run` script for `smallwebwaf` and
+440 -167
View File
@@ -1,12 +1,16 @@
// Package alerts sends alerts on bans, on a source that fails and on a
// file with an error to the webhook SWWAF_ALERT_WEBHOOK_URL names, each
// as one JSON object, as the "Alert webhook schema" section of SPEC.md
// describes. A repeat within SWWAF_ALERT_COOLDOWN is held back, and so is
// an alert past SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary. The
// others wait in a bounded queue, so that a slow or unreachable webhook
// never holds up a request. The state is written to alerts.json and read
// from it by the state package. Nothing logged names the webhook's URL,
// whose path or query can carry a secret.
// file with an error to each destination set: to the webhook
// 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 SWWAF_ALERT_SLACK_WEBHOOK_URL names, as a message, and to the
// ntfy topic SWWAF_ALERT_NTFY_URL names. A repeat within
// SWWAF_ALERT_COOLDOWN is held back, and so is an alert past
// SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary. The others wait in a
// bounded queue of each destination's own, so that a destination that is
// slow or unreachable holds up neither the others nor any request. The
// state is written to alerts.json and read from it by the state package.
// Nothing logged names a destination's URL, whose path or query can carry
// a secret.
package alerts
import (
@@ -23,6 +27,7 @@ import (
"net/netip"
"net/url"
"slices"
"strings"
"sync"
"sync/atomic"
"time"
@@ -60,36 +65,62 @@ func Events() []string {
}
}
// The destinations alerts are sent to, as the metrics and alerts.json
// name them.
const (
// queueSize is the most alerts that wait to be sent. Past it, the
// oldest is dropped.
// DestinationWebhook is the webhook SWWAF_ALERT_WEBHOOK_URL names.
DestinationWebhook = "webhook"
// DestinationSlack is the Slack incoming webhook
// SWWAF_ALERT_SLACK_WEBHOOK_URL names.
DestinationSlack = "slack"
// DestinationNtfy is the ntfy topic SWWAF_ALERT_NTFY_URL names.
DestinationNtfy = "ntfy"
)
// Destinations returns every destination alerts can be sent to.
func Destinations() []string {
return []string{DestinationWebhook, DestinationSlack, DestinationNtfy}
}
const (
// queueSize is the most alerts that wait to be sent to a destination.
// Past it, the oldest is dropped.
queueSize = 1000
// sendTimeout bounds one request to the webhook.
// sendTimeout bounds one request to a destination.
sendTimeout = 10 * time.Second
// After a request to the webhook fails, the alert is sent again a
// After a request to a destination fails, the alert is sent again a
// second later, and retryDelayFactor times as long after each further
// failure in a row, up to a minute.
firstRetryDelay = time.Second
retryDelayFactor = 2
maxRetryDelay = time.Minute
// maxAnswerBytes is the most of the webhook's answer that is read.
// maxAnswerBytes is the most of a destination's answer that is read.
maxAnswerBytes = 64 << 10
)
var (
errStatus = errors.New("the webhook answered")
// errRefused is a 4xx answer other than 408 and 429: the webhook
errStatus = errors.New("the destination answered")
// errRefused is a 4xx answer other than 408 and 429: the destination
// refuses the alert itself, and would refuse it again.
errRefused = errors.New("the webhook refused the alert, answering")
errRefused = errors.New("the destination refused the alert, answering")
)
// Params are what New needs.
// Params are what New needs. With none of WebhookURL, SlackURL and
// NtfyURL set, no alert is sent.
type Params struct {
// WebhookURL is where each alert is posted (SWWAF_ALERT_WEBHOOK_URL),
// nil while it is unset and no alert is sent. WebhookHeaders are sent
// with each (SWWAF_ALERT_WEBHOOK_HEADERS).
// WebhookURL is where each alert is posted as JSON
// (SWWAF_ALERT_WEBHOOK_URL), nil while it is unset. WebhookHeaders
// are sent with each (SWWAF_ALERT_WEBHOOK_HEADERS).
WebhookURL *url.URL
WebhookHeaders http.Header
// SlackURL is the Slack incoming webhook each alert is posted to as a
// message (SWWAF_ALERT_SLACK_WEBHOOK_URL), nil while it is unset.
SlackURL *url.URL
// NtfyURL is the ntfy topic each alert is published to
// (SWWAF_ALERT_NTFY_URL), nil while it is unset. NtfyToken, unless
// empty, is sent with each as a bearer token (SWWAF_ALERT_NTFY_TOKEN).
NtfyURL *url.URL
NtfyToken string
// Events are the events alerts are sent for (SWWAF_ALERT_EVENTS).
Events []string
// Cooldown is how long a repeat of an alert is held back
@@ -101,7 +132,7 @@ type Params struct {
Instance string
// Now tells the time of an alert, normally time.Now in UTC.
Now func() time.Time
// ProcessLog receives the requests to the webhook that fail.
// ProcessLog receives the requests to a destination that fail.
ProcessLog *slog.Logger
}
@@ -155,32 +186,63 @@ type Hour struct {
}
// State is what alerts.json holds: the cooldowns, the hour under way, and
// the alerts waiting to be sent, oldest first.
// for each destination set, the alerts waiting to be sent to it, oldest
// first.
type State struct {
Cooldowns []Cooldown `json:"cooldowns"`
Hour Hour `json:"hour"`
Waiting []Alert `json:"waiting"`
Cooldowns []Cooldown `json:"cooldowns"`
Hour Hour `json:"hour"`
Waiting map[string][]Alert `json:"waiting"`
}
// Counts are, for a destination, how many alerts it took, how many
// requests to it failed, and how many alerts were dropped from its full
// queue or given up as it refused them.
type Counts struct {
Sent, Failed, Dropped int64
}
// Queue takes the alerts raised, holds back those it must, and sends the
// others to the webhook. It is safe for concurrent use.
// others to each destination set, from a queue of the destination's own.
// It is safe for concurrent use.
type Queue struct {
params Params
// httpClient follows no redirect: a redirect is a failure.
httpClient *http.Client
// queued receives a value when an alert joins the queue, unless one
// waits already, so that Run looks at the queue again.
queued chan struct{}
// destinations are the destinations set, in the order of
// Destinations.
destinations []*destination
mu sync.Mutex
// cooldowns are the alerts last let through, by event and netblock,
// file or source.
cooldowns map[cooldownKey]*Cooldown
hour Hour
suppressed atomic.Int64
}
// destination is a destination set, with the alerts waiting to be sent
// to it. Its mu is taken after the Queue's, never before.
type destination struct {
// name is how the metrics and alerts.json name the destination, and
// setting the setting that is its URL, which the log names in place
// of the URL.
name string
setting string
url *url.URL
// message returns the body an alert is posted with, and the headers
// sent with it.
message func(alert *Alert) ([]byte, http.Header, error)
// httpClient follows no redirect: a redirect is a failure.
httpClient *http.Client
processLog *slog.Logger
// queued receives a value when an alert joins the queue, unless one
// waits already, so that run looks at the queue again.
queued chan struct{}
mu sync.Mutex
// waiting are the alerts waiting to be sent, oldest first.
waiting []*Alert
sent, failed, suppressed, dropped atomic.Int64
sent, failed, dropped atomic.Int64
}
// cooldownKey is what makes an alert a repeat of another: the same event
@@ -203,32 +265,44 @@ func cooldownKeyOf(alert *Alert) cooldownKey {
// New returns a Queue with no alert yet.
func New(params Params) *Queue {
return &Queue{
params: params,
httpClient: &http.Client{
CheckRedirect: func(*http.Request, []*http.Request) error {
return http.ErrUseLastResponse
},
},
queued: make(chan struct{}, 1),
q := &Queue{
params: params,
cooldowns: map[cooldownKey]*Cooldown{},
hour: Hour{HeldBack: map[string]int{}},
}
if params.WebhookURL != nil {
q.addDestination(DestinationWebhook, "SWWAF_ALERT_WEBHOOK_URL",
params.WebhookURL, q.webhookMessage)
}
if params.SlackURL != nil {
q.addDestination(DestinationSlack, "SWWAF_ALERT_SLACK_WEBHOOK_URL",
params.SlackURL, slackMessage)
}
if params.NtfyURL != nil {
q.addDestination(DestinationNtfy, "SWWAF_ALERT_NTFY_URL",
params.NtfyURL, q.ntfyMessage)
}
return q
}
// Raise sends alert, which names its event and what is particular to it,
// unless no webhook is set or SWWAF_ALERT_EVENTS leaves its event 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 counted,
// and the next one let through gives that count. Past MaxPerHour alerts
// let through in the hour under way, by the clock, an alert is held back
// for that hour's summary instead, which is sent once the hour has ended;
// it starts no cooldown, and the repeats held back before it are given by
// the next alert let through. Raise never waits: an alert let through
// joins the queue, from which Run sends it, and with queueSize alerts
// waiting the oldest is dropped.
// 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
// 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
// alerts let through in the hour under way, by the clock, an alert is
// held back for that hour's summary instead, which is sent once the hour
// has ended; it starts no cooldown, and the repeats held back before it
// are given by the next alert let through. Raise never waits: an alert
// let through joins the queue of each destination, from which Run sends
// it, and with queueSize alerts waiting for a destination, the oldest is
// dropped.
func (q *Queue) Raise(alert Alert) {
if q.params.WebhookURL == nil || !slices.Contains(q.params.Events, alert.Event) {
if len(q.destinations) == 0 || !slices.Contains(q.params.Events, alert.Event) {
return
}
@@ -260,12 +334,12 @@ func (q *Queue) Raise(alert Alert) {
}
// WouldSend reports whether Raise would let an alert for event on
// netblock through now: a webhook is set, SWWAF_ALERT_EVENTS chooses
// netblock through now: a destination is set, SWWAF_ALERT_EVENTS chooses
// event, no alert for event on netblock was let through less than
// Cooldown before, and fewer than MaxPerHour alerts have been let through
// in the hour under way. Unlike Raise, it counts nothing.
func (q *Queue) WouldSend(event string, netblock netip.Prefix) bool {
if q.params.WebhookURL == nil || !slices.Contains(q.params.Events, event) {
if len(q.destinations) == 0 || !slices.Contains(q.params.Events, event) {
return false
}
@@ -284,102 +358,70 @@ func (q *Queue) WouldSend(event string, netblock netip.Prefix) bool {
return q.params.MaxPerHour == 0 || q.hour.Sent < q.params.MaxPerHour
}
// Run sends the alerts waiting, oldest first, until ctx is done. An alert
// stays in the queue until the webhook answers it with a 2xx status, or
// refuses it with a 4xx status other than 408 and 429: a refused alert is
// logged, counted as dropped, and given up, so that the next is sent. Any
// other request that fails is logged, and the alert sent again
// firstRetryDelay later, retryDelayFactor times as long after each
// further failure in a row, up to maxRetryDelay. Run also ends each hour
// as Raise does, so that the hour's summary is sent as it ends. With no
// webhook set, it returns at once.
// Run sends the alerts waiting to each destination, from its own queue,
// as destination.run does, until ctx is done. It also ends each hour as
// Raise does, so that the hour's summary is sent as it ends. With no
// destination set, it returns at once.
func (q *Queue) Run(ctx context.Context) {
if q.params.WebhookURL == nil {
if len(q.destinations) == 0 {
return
}
var (
retryDelay time.Duration
retryAt time.Time
)
var sending sync.WaitGroup
for _, d := range q.destinations {
sending.Go(func() { d.run(ctx) })
}
for {
alert, untilHourEnds := q.next()
hourEnds := time.NewTimer(untilHourEnds)
var due <-chan time.Time // nil while no alert waits
if alert != nil {
due = time.After(time.Until(retryAt))
}
q.mu.Lock()
untilHourEnds := q.hour.Start.Add(time.Hour).Sub(q.params.Now())
q.mu.Unlock()
select {
case <-ctx.Done():
hourEnds.Stop()
sending.Wait()
return
case <-q.queued:
case <-hourEnds.C:
case <-time.After(untilHourEnds):
q.mu.Lock()
q.endHour(q.params.Now())
q.mu.Unlock()
case <-due:
err := q.send(ctx, alert)
switch {
case err == nil:
q.remove(alert)
q.sent.Add(1)
retryDelay = 0
retryAt = time.Time{}
case errors.Is(err, errRefused):
q.remove(alert)
q.failed.Add(1)
q.dropped.Add(1)
retryDelay = 0
retryAt = time.Time{}
q.params.ProcessLog.Warn("gave up an alert SWWAF_ALERT_WEBHOOK_URL refused",
"event", alert.Event, "error", err.Error())
case ctx.Err() == nil: // not cut off as smallwebwaf stops
q.failed.Add(1)
retryDelay = min(max(retryDelayFactor*retryDelay, firstRetryDelay),
maxRetryDelay)
retryAt = time.Now().Add(retryDelay)
q.params.ProcessLog.Warn("sending an alert to SWWAF_ALERT_WEBHOOK_URL failed",
"error", err.Error(), "sending_again_in", retryDelay.String())
}
}
hourEnds.Stop()
}
}
// Sent is how many alerts the webhook has taken.
func (q *Queue) Sent() int64 {
return q.sent.Load()
// Counts returns the counts of the destination name, all 0 for one not
// set.
func (q *Queue) Counts(name string) Counts {
for _, d := range q.destinations {
if d.name == name {
return Counts{
Sent: d.sent.Load(), Failed: d.failed.Load(), Dropped: d.dropped.Load(),
}
}
}
return Counts{}
}
// Failed is how many requests to the webhook have failed.
func (q *Queue) Failed() int64 {
return q.failed.Load()
// DestinationsSet returns the destinations set, in the order of
// Destinations.
func (q *Queue) DestinationsSet() []string {
names := make([]string, 0, len(q.destinations))
for _, d := range q.destinations {
names = append(names, d.name)
}
return names
}
// Suppressed is how many alerts were held back: by the cooldown, and past
// MaxPerHour.
// MaxPerHour. No destination is sent such an alert.
func (q *Queue) Suppressed() int64 {
return q.suppressed.Load()
}
// Dropped is how many alerts were dropped from a full queue, or given up
// as the webhook refused them.
func (q *Queue) Dropped() int64 {
return q.dropped.Load()
}
// Snapshot returns the queue's state, as alerts.json holds it, with the
// cooldowns sorted by netblock, then by event, file and source.
func (q *Queue) Snapshot() State {
@@ -389,7 +431,7 @@ func (q *Queue) Snapshot() State {
state := State{
Cooldowns: make([]Cooldown, 0, len(q.cooldowns)),
Hour: q.hour,
Waiting: make([]Alert, 0, len(q.waiting)),
Waiting: map[string][]Alert{},
}
state.Hour.HeldBack = maps.Clone(q.hour.HeldBack)
@@ -402,8 +444,8 @@ func (q *Queue) Snapshot() State {
cmp.Compare(a.File, b.File), cmp.Compare(a.Source, b.Source))
})
for _, alert := range q.waiting {
state.Waiting = append(state.Waiting, *alert)
for _, d := range q.destinations {
state.Waiting[d.name] = d.snapshot()
}
return state
@@ -411,8 +453,9 @@ func (q *Queue) Snapshot() State {
// Load puts state, read from alerts.json, in place of the queue's state.
// Each cooldown's netblock is masked to its length, so that
// 203.0.113.9/24 is 203.0.113.0/24. Past queueSize alerts waiting, the
// oldest are dropped.
// 203.0.113.9/24 is 203.0.113.0/24. The alerts waiting for a destination
// that is not set are dropped, and so are the oldest past queueSize
// alerts waiting for one that is.
func (q *Queue) Load(state State) {
q.mu.Lock()
defer q.mu.Unlock()
@@ -432,13 +475,33 @@ func (q *Queue) Load(state State) {
q.hour.HeldBack = map[string]int{}
}
q.waiting = nil
for _, alert := range state.Waiting {
q.queue(&alert)
for _, d := range q.destinations {
d.load(state.Waiting[d.name])
}
}
// addDestination adds a destination: its name, the setting that gives
// its URL, that URL, target, and message, which makes the messages sent
// to it.
func (q *Queue) addDestination(
name, setting string, target *url.URL,
message func(alert *Alert) ([]byte, http.Header, error),
) {
q.destinations = append(q.destinations, &destination{
name: name,
setting: setting,
url: target,
message: message,
httpClient: &http.Client{
CheckRedirect: func(*http.Request, []*http.Request) error {
return http.ErrUseLastResponse
},
},
processLog: q.params.ProcessLog,
queued: make(chan struct{}, 1),
})
}
// repeat reports whether alert, raised at now, repeats the last one let
// through less than Cooldown before, and counts it if it does.
func (q *Queue) repeat(alert *Alert, now time.Time) bool {
@@ -515,72 +578,163 @@ func (q *Queue) endHour(now time.Time) {
}
}
// queue adds alert to the alerts waiting, first dropping the oldest while
// queueSize wait, and has Run look at the queue again.
// queue adds alert to the alerts waiting for each destination.
func (q *Queue) queue(alert *Alert) {
if len(q.waiting) == queueSize {
q.waiting = slices.Delete(q.waiting, 0, 1)
q.dropped.Add(1)
for _, d := range q.destinations {
d.add(alert)
}
}
// run sends the alerts waiting, oldest first, until ctx is done. An alert
// stays in the queue until the destination answers it with a 2xx status,
// or refuses it with a 4xx status other than 408 and 429: a refused alert
// is logged, counted as dropped, and given up, so that the next is sent.
// Any other request that fails is logged, and the alert sent again
// firstRetryDelay later, retryDelayFactor times as long after each
// further failure in a row, up to maxRetryDelay.
func (d *destination) run(ctx context.Context) {
var (
retryDelay time.Duration
retryAt time.Time
)
for {
alert := d.oldest()
var due <-chan time.Time // nil while no alert waits
if alert != nil {
due = time.After(time.Until(retryAt))
}
select {
case <-ctx.Done():
return
case <-d.queued:
case <-due:
err := d.send(ctx, alert)
switch {
case err == nil:
d.remove(alert)
d.sent.Add(1)
retryDelay = 0
retryAt = time.Time{}
case errors.Is(err, errRefused):
d.remove(alert)
d.failed.Add(1)
d.dropped.Add(1)
retryDelay = 0
retryAt = time.Time{}
d.processLog.Warn("gave up an alert "+d.setting+" refused",
"event", alert.Event, "error", err.Error())
case ctx.Err() == nil: // not cut off as smallwebwaf stops
d.failed.Add(1)
retryDelay = min(max(retryDelayFactor*retryDelay, firstRetryDelay),
maxRetryDelay)
retryAt = time.Now().Add(retryDelay)
d.processLog.Warn("sending an alert to "+d.setting+" failed",
"error", err.Error(), "sending_again_in", retryDelay.String())
}
}
}
}
// add adds alert to the alerts waiting, first dropping the oldest while
// queueSize wait, and has run look at the queue again.
func (d *destination) add(alert *Alert) {
d.mu.Lock()
defer d.mu.Unlock()
if len(d.waiting) == queueSize {
d.waiting = slices.Delete(d.waiting, 0, 1)
d.dropped.Add(1)
}
q.waiting = append(q.waiting, alert)
d.waiting = append(d.waiting, alert)
select {
case q.queued <- struct{}{}:
case d.queued <- struct{}{}:
default: // a value waits already
}
}
// next returns the oldest alert waiting, nil when none waits, and how
// long it is until the hour under way ends.
func (q *Queue) next() (*Alert, time.Duration) {
q.mu.Lock()
defer q.mu.Unlock()
// load puts waiting, read from alerts.json, in place of the alerts
// waiting, as add adds them.
func (d *destination) load(waiting []Alert) {
d.mu.Lock()
d.waiting = nil
d.mu.Unlock()
var oldest *Alert
if len(q.waiting) > 0 {
oldest = q.waiting[0]
for _, alert := range waiting {
d.add(&alert)
}
return oldest, q.hour.Start.Add(time.Hour).Sub(q.params.Now())
}
// remove takes alert, which Run has sent or given up, out of the queue,
// unless it has been dropped from it, or Load has replaced the queue,
// since Run took it. Only the oldest alert is ever dropped, so alert is
// snapshot returns the alerts waiting, oldest first.
func (d *destination) snapshot() []Alert {
d.mu.Lock()
defer d.mu.Unlock()
waiting := make([]Alert, 0, len(d.waiting))
for _, alert := range d.waiting {
waiting = append(waiting, *alert)
}
return waiting
}
// oldest returns the oldest alert waiting, nil when none waits.
func (d *destination) oldest() *Alert {
d.mu.Lock()
defer d.mu.Unlock()
if len(d.waiting) == 0 {
return nil
}
return d.waiting[0]
}
// remove takes alert, which run has sent or given up, out of the queue,
// unless it has been dropped from it, or load has replaced the queue,
// since run took it. Only the oldest alert is ever dropped, so alert is
// the oldest if it is there at all.
func (q *Queue) remove(alert *Alert) {
q.mu.Lock()
defer q.mu.Unlock()
func (d *destination) remove(alert *Alert) {
d.mu.Lock()
defer d.mu.Unlock()
if len(q.waiting) > 0 && q.waiting[0] == alert {
q.waiting = slices.Delete(q.waiting, 0, 1)
if len(d.waiting) > 0 && d.waiting[0] == alert {
d.waiting = slices.Delete(d.waiting, 0, 1)
}
}
// send posts alert to the webhook as JSON, with WebhookHeaders, and
// returns an error unless the webhook answers with a 2xx status: one that
// send posts alert to the destination, as message makes it, and returns
// an error unless the destination answers with a 2xx status: one that
// wraps errRefused for a 4xx status other than 408 and 429. No error
// names the webhook's URL, whose path or query can carry a secret.
func (q *Queue) send(ctx context.Context, alert *Alert) error {
body, err := json.Marshal(alert)
// names the destination's URL, whose path or query can carry a secret.
func (d *destination) send(ctx context.Context, alert *Alert) error {
body, header, err := d.message(alert)
if err != nil {
return fmt.Errorf("encode the alert: %w", err)
return err
}
ctx, cancel := context.WithTimeout(ctx, sendTimeout)
defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
q.params.WebhookURL.String(), bytes.NewReader(body))
req, err := http.NewRequestWithContext(ctx, http.MethodPost, d.url.String(),
bytes.NewReader(body))
if err != nil {
return fmt.Errorf("make the request: %w", err)
}
maps.Copy(req.Header, q.params.WebhookHeaders)
req.Header.Set("Content-Type", "application/json")
maps.Copy(req.Header, header)
res, err := q.httpClient.Do(req)
res, err := d.httpClient.Do(req)
if err != nil {
// The client's error names the URL: only what went wrong is kept.
if urlErr, ok := errors.AsType[*url.Error](err); ok {
@@ -607,3 +761,122 @@ func (q *Queue) send(ctx context.Context, alert *Alert) error {
return fmt.Errorf("%w %s", errStatus, res.Status)
}
}
// webhookMessage returns alert as JSON, for the webhook, and the headers
// sent with it: WebhookHeaders, and its Content-Type.
func (q *Queue) webhookMessage(alert *Alert) ([]byte, http.Header, error) {
body, err := json.Marshal(alert)
if err != nil {
return nil, nil, fmt.Errorf("encode the alert: %w", err)
}
header := http.Header{}
maps.Copy(header, q.params.WebhookHeaders)
header.Set("Content-Type", "application/json")
return body, header, nil
}
// slackMessage returns alert as a message for a Slack incoming webhook,
// in JSON: its title in bold, then its text, with &, < and > escaped, as
// Slack asks, so that nothing in them is read as a link or a mention.
func slackMessage(alert *Alert) ([]byte, http.Header, error) {
escape := strings.NewReplacer("&", "&amp;", "<", "&lt;", ">", "&gt;").Replace
body, err := json.Marshal(map[string]string{
"text": "*" + escape(title(alert)) + "*\n" + escape(text(alert)),
})
if err != nil {
return nil, nil, fmt.Errorf("encode the message: %w", err)
}
return body, http.Header{"Content-Type": {"application/json"}}, nil
}
// ntfyMessage returns alert's text, as the message published to ntfy,
// and the headers sent with it: its title, the priority and the tag of
// its event, and NtfyToken, unless it is empty, as a bearer token.
func (q *Queue) ntfyMessage(alert *Alert) ([]byte, http.Header, error) {
header := http.Header{
"Title": {title(alert)},
"Priority": {ntfyPriority(alert.Event)},
"Tags": {ntfyTag(alert.Event)},
}
if q.params.NtfyToken != "" {
header.Set("Authorization", "Bearer "+q.params.NtfyToken)
}
return []byte(text(alert)), header, nil
}
// ntfyPriority returns the priority an alert for event is published to
// ntfy with: high for an event the admin needs to look at.
func ntfyPriority(event string) string {
switch event {
case EventPermanentBan, EventAnomaly, EventSourceFailure, EventFileError:
return "high"
case EventReputationHit:
return "low"
default: // ban, waf_block and summary
return "default"
}
}
// ntfyTag returns the tag an alert for event is published to ntfy with,
// which ntfy shows as an emoji.
func ntfyTag(event string) string {
switch event {
case EventBan, EventPermanentBan:
return "no_entry"
case EventWAFBlock:
return "shield"
case EventAnomaly:
return "chart_with_upwards_trend"
case EventReputationHit:
return "label"
case EventSourceFailure, EventFileError:
return "warning"
default: // summary
return "bar_chart"
}
}
// title returns the title of alert in Slack and ntfy: the instance and
// the event.
func title(alert *Alert) string {
return alert.Instance + ": " + alert.Event
}
// text returns the text of alert in Slack and ntfy: its reason, then a
// line for each of its client, netblock and country, the file, source,
// error and mode its detail gives, and its suppressed repeats, that it
// has.
func text(alert *Alert) string {
lines := []string{alert.Reason}
if alert.Client.IsValid() {
lines = append(lines, "client: "+alert.Client.String())
}
if alert.Netblock.IsValid() {
lines = append(lines, "netblock: "+alert.Netblock.String())
}
if alert.Country != "" {
lines = append(lines, "country: "+alert.Country)
}
for _, name := range []string{"file", "source", "error", "mode"} {
value, _ := alert.Detail[name].(string)
if value != "" {
lines = append(lines, name+": "+value)
}
}
if alert.SuppressedRepeats > 0 {
lines = append(lines, fmt.Sprintf("suppressed repeats: %d", alert.SuppressedRepeats))
}
return strings.Join(lines, "\n")
}
+495 -33
View File
@@ -26,13 +26,18 @@ import (
// clock of the test's own, which starts at 2000-01-01T00:00:00Z, the start
// of an hour: a wait lasts exactly as long as it should, however slowly
// the test process runs, and synctest.Wait returns once the queue has
// done all it can before time passes. The stand-in for the webhook
// answers without the network, since a request waiting on the network
// would keep that clock from moving on.
// done all it can before time passes. The stand-ins for the webhook,
// Slack and ntfy answer without the network, since a request waiting on
// the network would keep that clock from moving on.
const (
// webhookURL is where the alerts are posted.
// webhookURL is where the alerts are posted, slackURL the Slack
// incoming webhook, and ntfyURL the ntfy topic. ntfyToken is the ntfy
// token of the tests that set one.
webhookURL = "https://alerts.example/smallwebwaf?team=ops"
slackURL = "https://hooks.slack.example/services/T0123/B4567/abcdef"
ntfyURL = "https://ntfy.example/smallwebwaf-alerts"
ntfyToken = "tk_0123456789abcdefghijklmnopq"
// instance is the instance name every alert gives.
instance = "fsn1app1/gitea"
// started is when each test starts, as an alert gives it, and
@@ -135,7 +140,7 @@ func TestWouldSendOnlyForTheChosenEvents(t *testing.T) {
}
}
func TestNothingIsQueuedWithoutAWebhook(t *testing.T) {
func TestRaiseDoesNothingWithoutADestination(t *testing.T) {
t.Parallel()
params := newParams()
@@ -144,12 +149,18 @@ func TestNothingIsQueuedWithoutAWebhook(t *testing.T) {
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
if waiting := q.Snapshot().Waiting; len(waiting) != 0 {
t.Errorf("%d alerts wait, want none", len(waiting))
// No alert waits, no cooldown has started, and the hour counts none.
want := alerts.State{
Cooldowns: []alerts.Cooldown{},
Hour: alerts.Hour{HeldBack: map[string]int{}},
Waiting: map[string][]alerts.Alert{},
}
if got := q.Snapshot(); !reflect.DeepEqual(got, want) {
t.Errorf("state %+v, want %+v", got, want)
}
}
func TestWouldSendNothingWithoutAWebhook(t *testing.T) {
func TestWouldSendNothingWithoutADestination(t *testing.T) {
t.Parallel()
params := newParams()
@@ -161,6 +172,27 @@ func TestWouldSendNothingWithoutAWebhook(t *testing.T) {
}
}
func TestWouldSendWithOnlySlackOrOnlyNtfySet(t *testing.T) {
t.Parallel()
onlySlack := newParams()
onlySlack.WebhookURL = nil
onlySlack.SlackURL = parseURL(slackURL)
onlyNtfy := newParams()
onlyNtfy.WebhookURL = nil
onlyNtfy.NtfyURL = parseURL(ntfyURL)
for setting, params := range map[string]alerts.Params{
"SWWAF_ALERT_SLACK_WEBHOOK_URL": onlySlack,
"SWWAF_ALERT_NTFY_URL": onlyNtfy,
} {
if !alerts.New(params).WouldSend(alerts.EventBan, netblock(1)) {
t.Errorf("a ban alert would not be sent with only %s set", setting)
}
}
}
func TestRepeatWithinTheCooldownIsHeldBackAndCountedInTheNext(t *testing.T) {
t.Parallel()
@@ -265,7 +297,7 @@ func TestFileErrorAndSourceFailureRepeatOnlyForTheSameFileOrSource(t *testing.T)
after.Raise(fileError("/rules.d/50-a.rules"))
after.Raise(fileError("/rules.d/50-c.rules"))
waiting := after.Snapshot().Waiting
waiting := after.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || waiting[0].Detail["file"] != "/rules.d/50-c.rules" {
t.Errorf("after loading, alerts wait %+v, want the one for 50-c.rules", waiting)
}
@@ -444,7 +476,8 @@ func TestFailedRequestIsSentAgainWithBackoff(t *testing.T) {
wantCounts(t, q, 1, int64(len(want)), 0, 0)
if waiting := q.Snapshot().Waiting; len(waiting) != 0 {
waiting := q.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 0 {
t.Errorf("%d alerts still wait, want none", len(waiting))
}
})
@@ -552,7 +585,7 @@ func TestFullQueueDropsTheOldestAndRaiseNeverWaits(t *testing.T) {
wantCounts(t, q, 0, 0, 0, 1)
waiting := q.Snapshot().Waiting
waiting := q.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != alerts.QueueSize || waiting[0].Netblock != netblock(1) {
t.Fatalf("%d alerts wait, the first for %s, want %d, the first for %s",
len(waiting), waiting[0].Netblock, alerts.QueueSize, netblock(1))
@@ -627,7 +660,264 @@ func TestStateLoadedIntoANewQueueCarriesOn(t *testing.T) {
})
}
// How the stand-in for the webhook answers: with a status, or, hanging,
func TestSlackAndNtfyAreSentAMessageForEachEvent(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.WebhookURL = nil
standIns, q := startAll(t, params)
for _, each := range anAlertForEachEvent() {
q.Raise(each.alert)
}
synctest.Wait()
want := anAlertForEachEvent()
slack := standIns[alerts.DestinationSlack].received()
ntfy := standIns[alerts.DestinationNtfy].received()
if len(slack) != len(want) || len(ntfy) != len(want) {
t.Fatalf("Slack was sent %d messages and ntfy %d, want %d each",
len(slack), len(ntfy), len(want))
}
for i, each := range want {
wantSlackMessage(t, slack[i], "*"+each.title+"*\n"+each.text)
wantNtfyMessage(t, ntfy[i], each.title, each.priorityAndTag, each.text)
}
})
}
func TestSlackAndNtfyAreSentTheSummaryAndTheRepeatsHeldBack(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.WebhookURL = nil
params.MaxPerHour = 1
standIns, q := startAll(t, params)
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
// an alert past the limit; once the hour has ended, its summary, and
// the next alert, which gives the repeat.
q.Raise(ban)
q.Raise(ban)
q.Raise(alerts.Alert{Event: alerts.EventFileError, Reason: "a file error"})
time.Sleep(time.Hour)
synctest.Wait()
q.Raise(ban)
synctest.Wait()
slack := standIns[alerts.DestinationSlack].received()
ntfy := standIns[alerts.DestinationNtfy].received()
if len(slack) != 3 || len(ntfy) != 3 {
t.Fatalf("Slack was sent %d messages and ntfy %d, want 3 each",
len(slack), len(ntfy))
}
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"
wantSlackMessage(t, slack[1], "*"+instance+": summary*\n"+summary)
wantNtfyMessage(t, ntfy[1], instance+": summary", "default bar_chart", summary)
wantSlackMessage(t, slack[2],
"*"+instance+": ban*\na ban\nnetblock: 203.0.113.1/32\nsuppressed repeats: 1")
wantNtfyMessage(t, ntfy[2], instance+": ban", "default no_entry",
"a ban\nnetblock: 203.0.113.1/32\nsuppressed repeats: 1")
})
}
func TestSlackMessageEscapesAmpersandsAndAngleBrackets(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.Instance = "app<1>"
standIns, q := startAll(t, params)
q.Raise(alerts.Alert{
Event: alerts.EventFileError, Reason: "<!channel> & <https://x.example|y>",
})
synctest.Wait()
got := standIns[alerts.DestinationSlack].received()
if len(got) != 1 {
t.Fatalf("Slack was sent %d messages, want 1", len(got))
}
wantSlackMessage(t, got[0], "*app&lt;1&gt;: file_error*\n"+
"&lt;!channel&gt; &amp; &lt;https://x.example|y&gt;")
wantBodies(t, standIns[alerts.DestinationNtfy], "<!channel> & <https://x.example|y>")
})
}
func TestNtfyTokenIsSentToNtfyAlone(t *testing.T) {
t.Parallel()
for _, tc := range []struct {
name, token string
want []string
}{
{"with a token", ntfyToken, []string{"Bearer " + ntfyToken}},
{"without one", "", nil},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.NtfyToken = tc.token
standIns, q := startAll(t, params)
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
synctest.Wait()
for destination, want := range map[string][]string{
alerts.DestinationWebhook: nil,
alerts.DestinationSlack: nil,
alerts.DestinationNtfy: tc.want,
} {
got := standIns[destination].received()
if len(got) != 1 {
t.Fatalf("%s was sent %d requests, want 1", destination, len(got))
}
authorization := got[0].header.Values("Authorization")
if !slices.Equal(authorization, want) {
t.Errorf("%s was sent Authorization %q, want %q", destination,
authorization, want)
}
}
})
})
}
}
func TestDestinationThatDoesNotAnswerHoldsUpNeitherOther(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
log := &lockedBuffer{}
params.ProcessLog = slog.New(slog.NewJSONHandler(log, nil))
standIns, q := startAll(t, params)
standIns[alerts.DestinationSlack].set(hanging)
for n := range 3 {
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(n)})
}
synctest.Wait()
// With no time passed, the webhook and ntfy have taken every alert,
// while Slack has not answered the first.
for destination, want := range map[string]alerts.Counts{
alerts.DestinationWebhook: {Sent: 3},
alerts.DestinationSlack: {},
alerts.DestinationNtfy: {Sent: 3},
} {
wantDestinationCounts(t, q, destination, want)
}
if got := standIns[alerts.DestinationSlack].received(); len(got) != 1 {
t.Errorf("Slack was sent %d requests, want 1", len(got))
}
// Slack's request is abandoned after 10 seconds, and logged without
// its URL; Slack, which answers again, is sent every alert a second
// later.
standIns[alerts.DestinationSlack].set(answering)
time.Sleep(11 * time.Second)
synctest.Wait()
wantDestinationCounts(t, q, alerts.DestinationSlack,
alerts.Counts{Sent: 3, Failed: 1})
logged := log.String()
if !strings.Contains(logged,
`"msg":"sending an alert to SWWAF_ALERT_SLACK_WEBHOOK_URL failed"`) ||
strings.Contains(logged, "hooks.slack.example") ||
strings.Contains(logged, "T0123") {
t.Errorf("process log %q names no failure, or names the URL", logged)
}
})
}
func TestDropsAreCountedForEachDestination(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.MaxPerHour = 0
standIns, q := startAll(t, params)
standIns[alerts.DestinationSlack].set(hanging)
standIns[alerts.DestinationNtfy].set(refusing)
// Each alert is raised once each destination has done all it can
// with those before it.
for n := range alerts.QueueSize + 1 {
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(n)})
synctest.Wait()
}
// Slack, which does not answer the first alert, drops it from its
// full queue for the last; ntfy refuses each, which is given up; and
// the webhook takes every one.
const all = alerts.QueueSize + 1
for destination, want := range map[string]alerts.Counts{
alerts.DestinationWebhook: {Sent: all},
alerts.DestinationSlack: {Dropped: 1},
alerts.DestinationNtfy: {Failed: all, Dropped: all},
} {
wantDestinationCounts(t, q, destination, want)
}
})
}
func TestEachDestinationIsSentOnlyTheAlertsWaitingForIt(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
params := withSlackAndNtfy(newParams())
params.WebhookURL = nil
standIns, q := startAll(t, params)
waiting := func(reason string) alerts.Alert {
return alerts.Alert{
Instance: instance, Time: midnight(), Event: alerts.EventFileError,
Reason: reason,
}
}
// As read from alerts.json. The alert waiting for the webhook, which
// is not set, is dropped.
q.Load(roundTrip(t, alerts.State{Waiting: map[string][]alerts.Alert{
alerts.DestinationWebhook: {waiting("first")},
alerts.DestinationSlack: {waiting("second"), waiting("third")},
alerts.DestinationNtfy: {waiting("fourth")},
}}))
synctest.Wait()
wantBodies(t, standIns[alerts.DestinationSlack],
`{"text":"*fsn1app1/gitea: file_error*\nsecond"}`,
`{"text":"*fsn1app1/gitea: file_error*\nthird"}`)
wantBodies(t, standIns[alerts.DestinationNtfy], "fourth")
// alerts.json then lists Slack and ntfy alone, with no alert waiting.
want := map[string][]alerts.Alert{
alerts.DestinationSlack: {}, alerts.DestinationNtfy: {},
}
if got := q.Snapshot().Waiting; !reflect.DeepEqual(got, want) {
t.Errorf("alerts waiting %v, want %v", got, want)
}
})
}
// How a stand-in for a destination answers: with a status, or, hanging,
// not at all, until the request is abandoned.
const (
answering = http.StatusNoContent
@@ -636,7 +926,7 @@ const (
hanging = 0
)
// standIn is a stand-in for the webhook. It notes each request it is
// standIn is a stand-in for a destination. It notes each request it is
// sent.
type standIn struct {
mu sync.Mutex
@@ -644,14 +934,16 @@ type standIn struct {
requests []post
}
// post is a request the webhook was sent: when, its method, URL and
// headers, the alert it carried, and whether the webhook answered it with
// a 2xx status.
// post is a request a destination was sent: when, its method, URL and
// headers, its body, and that body read as a JSON object, which for the
// webhook is the alert it carried, and whether the destination answered
// it with a 2xx status.
type post struct {
at time.Time
method string
url string
header http.Header
body string
alert map[string]any
answered bool
}
@@ -685,7 +977,7 @@ func (s *standIn) ServeHTTP(w http.ResponseWriter, r *http.Request) {
answers := s.answers
s.requests = append(s.requests, post{
at: time.Now(), method: r.Method, url: r.URL.String(), header: r.Header.Clone(),
alert: alert, answered: answers == answering,
body: string(body), alert: alert, answered: answers == answering,
})
s.mu.Unlock()
@@ -739,13 +1031,8 @@ func (b *lockedBuffer) String() string {
// every event, the default cooldown and hourly limit, and the bubble's
// clock in UTC.
func newParams() alerts.Params {
webhook, err := url.Parse(webhookURL)
if err != nil {
panic(err)
}
return alerts.Params{
WebhookURL: webhook,
WebhookURL: parseURL(webhookURL),
Events: alerts.Events(),
Cooldown: cooldown,
MaxPerHour: 60,
@@ -755,14 +1042,48 @@ func newParams() alerts.Params {
}
}
// withSlackAndNtfy returns params with Slack at slackURL and ntfy at
// ntfyURL set as well.
func withSlackAndNtfy(params alerts.Params) alerts.Params {
params.SlackURL = parseURL(slackURL)
params.NtfyURL = parseURL(ntfyURL)
return params
}
// parseURL returns rawURL, parsed.
func parseURL(rawURL string) *url.URL {
parsed, err := url.Parse(rawURL)
if err != nil {
panic(err)
}
return parsed
}
// start returns a stand-in for the webhook that answers, and a Queue that
// sends to it, run until the test ends.
func start(t *testing.T, params alerts.Params) (*standIn, *alerts.Queue) {
t.Helper()
webhook := &standIn{answers: answering}
standIns, q := startAll(t, params)
return standIns[alerts.DestinationWebhook], q
}
// startAll returns, by destination, a stand-in that answers for each
// destination params sets, and a Queue that sends to them, run until the
// test ends.
func startAll(t *testing.T, params alerts.Params) (map[string]*standIn, *alerts.Queue) {
t.Helper()
q := alerts.New(params)
q.SetTransport(webhook)
standIns := map[string]*standIn{}
for _, destination := range q.DestinationsSet() {
standIns[destination] = &standIn{answers: answering}
q.SetTransport(destination, standIns[destination])
}
ctx, stop := context.WithCancel(t.Context())
stopped := make(chan struct{})
@@ -777,7 +1098,7 @@ func start(t *testing.T, params alerts.Params) (*standIn, *alerts.Queue) {
<-stopped
})
return webhook, q
return standIns, q
}
// midnight is when each test starts.
@@ -836,17 +1157,158 @@ func wantEvents(t *testing.T, webhook *standIn, want ...string) {
}
}
// wantCounts checks the alerts q counts as sent, the requests it counts as
// failed, and the alerts it counts as held back and as dropped.
// wantCounts checks the alerts q counts as sent to the webhook, the
// requests to it that it counts as failed, the alerts it counts as held
// back, and those it counts as dropped for the webhook.
func wantCounts(
t *testing.T, q *alerts.Queue, sent, failed, suppressed, dropped int64,
) {
t.Helper()
if q.Sent() != sent || q.Failed() != failed || q.Suppressed() != suppressed ||
q.Dropped() != dropped {
t.Errorf("counts sent %d, failed %d, suppressed %d and dropped %d, "+
"want %d, %d, %d and %d", q.Sent(), q.Failed(), q.Suppressed(), q.Dropped(),
sent, failed, suppressed, dropped)
wantDestinationCounts(t, q, alerts.DestinationWebhook,
alerts.Counts{Sent: sent, Failed: failed, Dropped: dropped})
if q.Suppressed() != suppressed {
t.Errorf("%d alerts held back, want %d", q.Suppressed(), suppressed)
}
}
// wantDestinationCounts checks q's counts for destination.
func wantDestinationCounts(
t *testing.T, q *alerts.Queue, destination string, want alerts.Counts,
) {
t.Helper()
if got := q.Counts(destination); got != want {
t.Errorf("%s counts %+v, want %+v", destination, got, want)
}
}
// eventMessage is an alert, and the message Slack and ntfy are sent for
// it: its title, the priority and the tag ntfy is sent, with a space
// between them, and its text.
type eventMessage struct {
alert alerts.Alert
title, priorityAndTag, text string
}
// anAlertForEachEvent returns an alert for each event SWWAF_ALERT_EVENTS
// names, the first in observe mode, each with its message.
func anAlertForEachEvent() []eventMessage {
return []eventMessage{
{
alerts.Alert{
Event: alerts.EventBan, Client: netip.MustParseAddr("203.0.113.9"),
Netblock: netip.MustParsePrefix("203.0.113.0/24"), Country: "DE",
Reason: "requests per hour over the limit of 10000",
Detail: map[string]any{"mode": "observe"},
},
"fsn1app1/gitea: ban", "default no_entry",
"requests per hour over the limit of 10000\nclient: 203.0.113.9\n" +
"netblock: 203.0.113.0/24\ncountry: DE\nmode: observe",
},
{
alerts.Alert{
Event: alerts.EventPermanentBan, Client: netip.MustParseAddr("198.51.100.7"),
Netblock: netip.MustParsePrefix("198.51.100.7/32"), Country: "FR",
Reason: "matched the rule env-file",
},
"fsn1app1/gitea: permanent_ban", "high no_entry",
"matched the rule env-file\nclient: 198.51.100.7\n" +
"netblock: 198.51.100.7/32\ncountry: FR",
},
{
alerts.Alert{
Event: alerts.EventWAFBlock, Client: netip.MustParseAddr("192.0.2.1"),
Netblock: netip.MustParsePrefix("192.0.2.1/32"),
Reason: "refused by the Core Rule Set",
},
"fsn1app1/gitea: waf_block", "default shield",
"refused by the Core Rule Set\nclient: 192.0.2.1\nnetblock: 192.0.2.1/32",
},
{
alerts.Alert{
Event: alerts.EventAnomaly, Netblock: netip.MustParsePrefix("192.0.2.0/24"),
Reason: "requests per minute over the threshold of 5000",
},
"fsn1app1/gitea: anomaly", "high chart_with_upwards_trend",
"requests per minute over the threshold of 5000\nnetblock: 192.0.2.0/24",
},
{
alerts.Alert{
Event: alerts.EventReputationHit, Client: netip.MustParseAddr("192.0.2.2"),
Netblock: netip.MustParsePrefix("192.0.2.2/32"),
Reason: "listed by a DNS blocklist",
},
"fsn1app1/gitea: reputation_hit", "low label",
"listed by a DNS blocklist\nclient: 192.0.2.2\nnetblock: 192.0.2.2/32",
},
{
alerts.Alert{
Event: alerts.EventSourceFailure, Reason: "asking GeoJS failed",
Detail: map[string]any{"source": "geojs"},
},
"fsn1app1/gitea: source_failure", "high warning",
"asking GeoJS failed\nsource: geojs",
},
{
alerts.Alert{
Event: alerts.EventFileError,
Reason: "a rule file has an error, and the rules stay as they were",
Detail: map[string]any{
"file": "/etc/smallwebwaf/rules.d/50-app.rules", "error": "line 2: no action",
},
},
"fsn1app1/gitea: file_error", "high warning",
"a rule file has an error, and the rules stay as they were\n" +
"file: /etc/smallwebwaf/rules.d/50-app.rules\nerror: line 2: no action",
},
}
}
// wantSlackMessage checks that request, one Slack was sent, posted the
// message text, as JSON.
func wantSlackMessage(t *testing.T, request post, text string) {
t.Helper()
if request.method != http.MethodPost || request.url != slackURL ||
request.header.Get("Content-Type") != "application/json" ||
!reflect.DeepEqual(request.alert, map[string]any{"text": text}) {
t.Errorf("Slack was sent %s %s, Content-Type %q, %s, want POST %s, "+
"application/json, the text %q", request.method, request.url,
request.header.Get("Content-Type"), request.body, slackURL, text)
}
}
// wantNtfyMessage checks that request, one ntfy was sent, posted the
// message text with the title, and with the priority and the tag
// priorityAndTag gives, with a space between them.
func wantNtfyMessage(t *testing.T, request post, title, priorityAndTag, text string) {
t.Helper()
got := []string{
request.method, request.url, request.header.Get("Title"),
request.header.Get("Priority") + " " + request.header.Get("Tags"), request.body,
}
want := []string{http.MethodPost, ntfyURL, title, priorityAndTag, text}
if !slices.Equal(got, want) {
t.Errorf("ntfy was sent the method, URL, title, priority and tag, and text "+
"%q, want %q", got, want)
}
}
// wantBodies checks the bodies of the requests the stand-in was sent, in
// order.
func wantBodies(t *testing.T, s *standIn, want ...string) {
t.Helper()
got := make([]string, 0, len(want))
for _, request := range s.received() {
got = append(got, request.body)
}
if !slices.Equal(got, want) {
t.Errorf("bodies %q, want %q", got, want)
}
}
+9 -5
View File
@@ -2,11 +2,15 @@ package alerts
import "net/http"
// QueueSize is the most alerts that wait to be sent.
// QueueSize is the most alerts that wait to be sent to a destination.
const QueueSize = queueSize
// SetTransport has q's requests to the webhook go through transport
// instead of the network.
func (q *Queue) SetTransport(transport http.RoundTripper) {
q.httpClient.Transport = transport
// SetTransport has q's requests to the destination name go through
// transport instead of the network.
func (q *Queue) SetTransport(name string, transport http.RoundTripper) {
for _, d := range q.destinations {
if d.name == name {
d.httpClient.Transport = transport
}
}
}
+78 -27
View File
@@ -20,6 +20,7 @@ import (
"strconv"
"strings"
"time"
"unicode"
"unicode/utf8"
"sneak.berlin/go/smallwebwaf/internal/alerts"
@@ -160,18 +161,26 @@ type Config struct {
LogRemoteFacility int
LogRemoteAppName string
// AlertWebhookURL is where each alert is posted as JSON
// (SWWAF_ALERT_WEBHOOK_URL), nil while it is unset and no alert is
// sent. AlertWebhookHeaders are sent with each
// (SWWAF_ALERT_WEBHOOK_HEADERS). AlertEvents are the events alerts are
// sent for (SWWAF_ALERT_EVENTS). A repeat of an alert within
// AlertCooldown is held back (SWWAF_ALERT_COOLDOWN), and so is an alert
// past AlertMaxPerHour in an hour, for the hour's summary
// (SWWAF_ALERT_MAX_PER_HOUR); 0 is off for both.
AlertWebhookURL *url.URL
AlertWebhookHeaders http.Header
AlertEvents []string
AlertCooldown time.Duration
AlertMaxPerHour int
// (SWWAF_ALERT_WEBHOOK_URL), nil while it is unset. AlertWebhookHeaders
// are sent with each (SWWAF_ALERT_WEBHOOK_HEADERS).
// AlertSlackWebhookURL is the Slack incoming webhook each alert is
// posted to as a message (SWWAF_ALERT_SLACK_WEBHOOK_URL), and
// AlertNtfyURL the ntfy topic each is published to
// (SWWAF_ALERT_NTFY_URL), each nil while it is unset; AlertNtfyToken,
// unless empty, is sent to ntfy with each (SWWAF_ALERT_NTFY_TOKEN).
// With none of the three URLs set, no alert is sent. AlertEvents are
// the events alerts are sent for (SWWAF_ALERT_EVENTS). A repeat of an
// alert within AlertCooldown is held back (SWWAF_ALERT_COOLDOWN), and
// so is an alert past AlertMaxPerHour in an hour, for the hour's
// summary (SWWAF_ALERT_MAX_PER_HOUR); 0 is off for both.
AlertWebhookURL *url.URL
AlertWebhookHeaders http.Header
AlertSlackWebhookURL *url.URL
AlertNtfyURL *url.URL
AlertNtfyToken string
AlertEvents []string
AlertCooldown time.Duration
AlertMaxPerHour int
// settings are the values read, as given or by default, and the
// files they were read from, for the log line at start.
@@ -252,6 +261,8 @@ var (
errNotWebhookHeader = errors.New(
"is not a header name followed by : and the header's value, " +
"such as Authorization:Bearer <token>")
errControlCharacter = errors.New(
"holds a control character, such as the carriage return of a Windows line end")
errNotAlertEvent = errors.New(
"is not ban, permanent_ban, waf_block, anomaly, reputation_hit, " +
"source_failure or file_error")
@@ -303,17 +314,20 @@ func FromEnvironment(lookupEnv func(string) (string, bool)) (*Config, error) {
StateCounterInterval: env.durationNotOff("SWWAF_STATE_COUNTER_INTERVAL", "15m"),
LogRequestHeaders: env.headerNames("SWWAF_LOG_REQUEST_HEADERS",
"accept,accept-language,accept-encoding,content-type,origin,range"),
AdminToken: env.token("SWWAF_ADMIN_TOKEN"),
MetricsToken: env.token("SWWAF_METRICS_TOKEN"),
MetricsTopN: env.numberNotOff("SWWAF_METRICS_TOP_N", "50"),
RulesDir: env.value("SWWAF_RULES_DIR", "/etc/smallwebwaf/rules.d"),
RulesEnabled: env.boolean("SWWAF_RULES_ENABLED", "true"),
LogRemoteURL: env.logRemoteURL("SWWAF_LOG_REMOTE_URL"),
LogRemoteTLSCAs: env.certificates("SWWAF_LOG_REMOTE_TLS_CA_FILE"),
LogRemoteBuffer: env.numberNotOff("SWWAF_LOG_REMOTE_BUFFER", "10000"),
LogRemoteFacility: env.facility("SWWAF_LOG_REMOTE_FACILITY", "local0"),
AlertWebhookURL: env.webhookURL("SWWAF_ALERT_WEBHOOK_URL"),
AlertWebhookHeaders: env.webhookHeaders("SWWAF_ALERT_WEBHOOK_HEADERS"),
AdminToken: env.token("SWWAF_ADMIN_TOKEN"),
MetricsToken: env.token("SWWAF_METRICS_TOKEN"),
MetricsTopN: env.numberNotOff("SWWAF_METRICS_TOP_N", "50"),
RulesDir: env.value("SWWAF_RULES_DIR", "/etc/smallwebwaf/rules.d"),
RulesEnabled: env.boolean("SWWAF_RULES_ENABLED", "true"),
LogRemoteURL: env.logRemoteURL("SWWAF_LOG_REMOTE_URL"),
LogRemoteTLSCAs: env.certificates("SWWAF_LOG_REMOTE_TLS_CA_FILE"),
LogRemoteBuffer: env.numberNotOff("SWWAF_LOG_REMOTE_BUFFER", "10000"),
LogRemoteFacility: env.facility("SWWAF_LOG_REMOTE_FACILITY", "local0"),
AlertWebhookURL: env.webhookURL("SWWAF_ALERT_WEBHOOK_URL"),
AlertWebhookHeaders: env.webhookHeaders("SWWAF_ALERT_WEBHOOK_HEADERS"),
AlertSlackWebhookURL: env.webhookURL("SWWAF_ALERT_SLACK_WEBHOOK_URL"),
AlertNtfyURL: env.webhookURL("SWWAF_ALERT_NTFY_URL"),
AlertNtfyToken: env.secret("SWWAF_ALERT_NTFY_TOKEN"),
AlertEvents: env.alertEvents("SWWAF_ALERT_EVENTS",
strings.Join(alerts.Events(), ",")),
AlertCooldown: env.duration("SWWAF_ALERT_COOLDOWN", "15m"),
@@ -323,6 +337,8 @@ func FromEnvironment(lookupEnv func(string) (string, bool)) (*Config, error) {
cfg.LogRemoteAppName = env.appName("SWWAF_LOG_REMOTE_APP_NAME",
cfg.InstanceName, cfg.LogRemoteURL != nil)
env.checkInstanceNameForNtfy(cfg.InstanceName, cfg.AlertNtfyURL != nil)
for _, country := range cfg.ExclusivelyAllowedCountries {
if slices.Contains(cfg.DeniedCountries, country) {
env.check("SWWAF_EXCLUSIVELY_ALLOWED_COUNTRIES",
@@ -667,10 +683,23 @@ func (e *environment) appName(name, instanceName string, sending bool) string {
return value
}
// webhookURL reads the setting that is where each alert is posted. Unset
// or empty, it is nil, and no alert is sent. The log shows ******** in
// place of its path and query, and an error shows none of it, since many
// webhooks carry their secret there.
// checkInstanceNameForNtfy refuses an instance name that holds a control
// character while ntfySet, SWWAF_ALERT_NTFY_URL being set: ntfy is sent
// the instance name in a header, which cannot hold one.
func (e *environment) checkInstanceNameForNtfy(instanceName string, ntfySet bool) {
if ntfySet && strings.ContainsFunc(instanceName, unicode.IsControl) {
e.check("SWWAF_INSTANCE_NAME", fmt.Errorf(
"%q %w, and is sent to ntfy in a header while SWWAF_ALERT_NTFY_URL is set",
instanceName, errControlCharacter))
}
}
// webhookURL reads a setting that is a URL each alert is posted to:
// SWWAF_ALERT_WEBHOOK_URL, SWWAF_ALERT_SLACK_WEBHOOK_URL or
// SWWAF_ALERT_NTFY_URL. Unset or empty, it is nil, and no alert is posted
// there. The log shows ******** in place of its path and query, and an
// error shows none of it, since a webhook or an ntfy topic can carry its
// secret there.
func (e *environment) webhookURL(name string) *url.URL {
value, _ := e.lookup(name)
webhook, logged, err := parseWebhookURL(value)
@@ -692,6 +721,28 @@ func (e *environment) webhookHeaders(name string) http.Header {
return headers
}
// secret reads a setting that is a secret another service gave, such as
// an ntfy token, "" while it is unset. It is sent in a header, which
// cannot hold a control character, so one in it is an error. The log
// shows ******** in place of a value that is not empty, and an error
// shows none of it.
func (e *environment) secret(name string) string {
value, _ := e.lookup(name)
logged := ""
if value != "" {
logged = masked
}
e.settings = append(e.settings, slog.String(name, logged))
if strings.ContainsFunc(value, unicode.IsControl) {
e.check(name, errControlCharacter)
}
return value
}
// alertEvents reads the setting that is the events alerts are sent for.
func (e *environment) alertEvents(name, defaultValue string) []string {
events, err := parseAlertEvents(e.value(name, defaultValue))
+122 -18
View File
@@ -66,6 +66,9 @@ const (
logRemoteAppName = "SWWAF_LOG_REMOTE_APP_NAME"
alertWebhookURL = "SWWAF_ALERT_WEBHOOK_URL"
alertWebhookHeaders = "SWWAF_ALERT_WEBHOOK_HEADERS"
alertSlackWebhookURL = "SWWAF_ALERT_SLACK_WEBHOOK_URL"
alertNtfyURL = "SWWAF_ALERT_NTFY_URL"
alertNtfyToken = "SWWAF_ALERT_NTFY_TOKEN" //nolint:gosec // the setting's name
alertEvents = "SWWAF_ALERT_EVENTS"
alertCooldown = "SWWAF_ALERT_COOLDOWN"
alertMaxPerHour = "SWWAF_ALERT_MAX_PER_HOUR"
@@ -498,44 +501,65 @@ func TestAlertSettingsDefaults(t *testing.T) {
cfg := fromEnvironment(t, environment{})
if cfg.AlertWebhookURL != nil || len(cfg.AlertWebhookHeaders) != 0 ||
cfg.AlertSlackWebhookURL != nil || cfg.AlertNtfyURL != nil ||
cfg.AlertNtfyToken != "" ||
strings.Join(cfg.AlertEvents, ",") != defaultAlertEvents ||
cfg.AlertCooldown != 15*time.Minute || cfg.AlertMaxPerHour != 60 {
t.Errorf("alert settings %v, %v, %v, %s and %d, want no URL, no headers, "+
"%s, 15m and 60", cfg.AlertWebhookURL, cfg.AlertWebhookHeaders,
cfg.AlertEvents, cfg.AlertCooldown, cfg.AlertMaxPerHour, defaultAlertEvents)
t.Errorf("alert settings %v, %v, %v, %v, %q, %v, %s and %d, want no URLs, "+
"no headers, no token, %s, 15m and 60", cfg.AlertWebhookURL,
cfg.AlertWebhookHeaders, cfg.AlertSlackWebhookURL, cfg.AlertNtfyURL,
cfg.AlertNtfyToken, cfg.AlertEvents, cfg.AlertCooldown, cfg.AlertMaxPerHour,
defaultAlertEvents)
}
}
func TestAlertSettingsAsSet(t *testing.T) {
t.Parallel()
const webhook = "https://alerts.example:8443/hooks/waf?team=ops"
const (
webhook = "https://alerts.example:8443/hooks/waf?team=ops"
slack = "https://hooks.slack.example/services/T0123/B4567/abcdef"
ntfy = "https://ntfy.example/smallwebwaf-alerts"
token = "tk_0123456789abcdefghijklmnopq"
)
cfg := fromEnvironment(t, environment{
alertWebhookURL: webhook,
alertWebhookHeaders: "Authorization: Bearer abc:def , x-team:ops",
alertEvents: "ban, file_error",
alertCooldown: "1h",
alertMaxPerHour: "10",
alertWebhookURL: webhook,
alertWebhookHeaders: "Authorization: Bearer abc:def , x-team:ops",
alertSlackWebhookURL: slack,
alertNtfyURL: ntfy,
alertNtfyToken: token,
alertEvents: "ban, file_error",
alertCooldown: "1h",
alertMaxPerHour: "10",
})
headers := http.Header{"Authorization": {"Bearer abc:def"}, "X-Team": {"ops"}}
if cfg.AlertWebhookURL.String() != webhook ||
!reflect.DeepEqual(cfg.AlertWebhookHeaders, headers) ||
cfg.AlertSlackWebhookURL.String() != slack || cfg.AlertNtfyURL.String() != ntfy ||
cfg.AlertNtfyToken != token ||
!slices.Equal(cfg.AlertEvents, []string{"ban", "file_error"}) ||
cfg.AlertCooldown != time.Hour || cfg.AlertMaxPerHour != 10 {
t.Errorf("alert settings %v, %v, %v, %s and %d", cfg.AlertWebhookURL,
cfg.AlertWebhookHeaders, cfg.AlertEvents, cfg.AlertCooldown,
cfg.AlertMaxPerHour)
t.Errorf("alert settings %v, %v, %v, %v, %q, %v, %s and %d", cfg.AlertWebhookURL,
cfg.AlertWebhookHeaders, cfg.AlertSlackWebhookURL, cfg.AlertNtfyURL,
cfg.AlertNtfyToken, cfg.AlertEvents, cfg.AlertCooldown, cfg.AlertMaxPerHour)
}
}
cfg = fromEnvironment(t, environment{
alertWebhookURL: "", alertEvents: "", alertCooldown: off, alertMaxPerHour: off,
func TestAlertSettingsSetEmptyOrOff(t *testing.T) {
t.Parallel()
cfg := fromEnvironment(t, environment{
alertWebhookURL: "", alertSlackWebhookURL: "", alertNtfyURL: "",
alertNtfyToken: "", alertEvents: "", alertCooldown: off, alertMaxPerHour: off,
})
if cfg.AlertWebhookURL != nil || len(cfg.AlertEvents) != 0 ||
cfg.AlertCooldown != 0 || cfg.AlertMaxPerHour != 0 {
t.Errorf("set empty or off, alert settings %v, %v, %s and %d",
cfg.AlertWebhookURL, cfg.AlertEvents, cfg.AlertCooldown, cfg.AlertMaxPerHour)
if cfg.AlertWebhookURL != nil || cfg.AlertSlackWebhookURL != nil ||
cfg.AlertNtfyURL != nil || cfg.AlertNtfyToken != "" ||
len(cfg.AlertEvents) != 0 || cfg.AlertCooldown != 0 || cfg.AlertMaxPerHour != 0 {
t.Errorf("set empty or off, alert settings %v, %v, %v, %q, %v, %s and %d",
cfg.AlertWebhookURL, cfg.AlertSlackWebhookURL, cfg.AlertNtfyURL,
cfg.AlertNtfyToken, cfg.AlertEvents, cfg.AlertCooldown, cfg.AlertMaxPerHour)
}
}
@@ -550,6 +574,10 @@ func TestInvalidAlertSettingStopsTheStart(t *testing.T) {
{alertWebhookURL, "https://alerts.example/#top"},
{alertWebhookURL, "https://alerts.example:0/"},
{alertWebhookURL, "https://alerts.example:65536/"},
{alertSlackWebhookURL, "hooks.slack.example/services/T0123"},
{alertSlackWebhookURL, "https://user:password@hooks.slack.example/"},
{alertNtfyURL, "ntfy://ntfy.example/smallwebwaf-alerts"},
{alertNtfyURL, "https://ntfy.example/smallwebwaf-alerts#top"},
{alertWebhookHeaders, "Authorization"},
{alertWebhookHeaders, "X Team:ops"},
{alertWebhookHeaders, ":ops"},
@@ -636,6 +664,79 @@ func TestWebhookURLIsLoggedWithoutItsPathOrQueryAndNeverShown(t *testing.T) {
}
}
func TestSlackAndNtfySettingsAreLoggedWithoutTheirSecrets(t *testing.T) {
t.Parallel()
const token = "tk_0123456789abcdefghijklmnopq"
cfg := fromEnvironment(t, environment{
alertSlackWebhookURL: "https://hooks.slack.example/services/T0123/B4567/abcdef",
alertNtfyURL: "https://ntfy.example/smallwebwaf-alerts",
alertNtfyToken: token,
})
var out bytes.Buffer
slog.New(slog.NewJSONHandler(&out, nil)).Info("starting", "settings", cfg)
logged := out.String()
for _, want := range []string{
`"` + alertSlackWebhookURL + `":"https://hooks.slack.example/********"`,
`"` + alertNtfyURL + `":"https://ntfy.example/********"`,
`"` + alertNtfyToken + `":"********"`,
} {
if !strings.Contains(logged, want) {
t.Errorf("no %s in the settings logged: %s", want, logged)
}
}
for _, secret := range []string{"T0123", "smallwebwaf-alerts", token} {
if strings.Contains(logged, secret) {
t.Errorf("%s in the settings logged: %s", secret, logged)
}
}
}
func TestNtfyTokenWithAControlCharacterStopsTheStartWithoutShowingIt(t *testing.T) {
t.Parallel()
// A file saved with Windows line ends keeps the carriage return.
for name, env := range map[string]environment{
"set": {alertNtfyToken: token + "\r"},
"in a file": {alertNtfyToken + "_FILE": writeFile(t, token+"\r\n")},
} {
_, err := config.FromEnvironment(env.lookupEnv)
want := alertNtfyToken + ": holds a control character, such as the " +
"carriage return of a Windows line end"
if err == nil || err.Error() != want {
t.Errorf("%s: error %v, want %s", name, err, want)
}
}
}
func TestInstanceNameWithAControlCharacterStopsTheStartOnlyWithNtfySet(t *testing.T) {
t.Parallel()
const name = "fsn1app1\r"
_, err := config.FromEnvironment(environment{
instanceName: name, alertNtfyURL: "https://ntfy.example/smallwebwaf-alerts",
}.lookupEnv)
want := instanceName + `: "fsn1app1\r" holds a control character, such as the ` +
`carriage return of a Windows line end, and is sent to ntfy in a header ` +
`while ` + alertNtfyURL + ` is set`
if err == nil || err.Error() != want {
t.Errorf("error %v, want %s", err, want)
}
cfg := fromEnvironment(t, environment{instanceName: name})
if cfg.InstanceName != name {
t.Errorf("not sending to ntfy, %s is %q", instanceName, cfg.InstanceName)
}
}
func TestCodeOnBothCountryListsStopsTheStart(t *testing.T) {
t.Parallel()
@@ -1087,6 +1188,9 @@ func TestLogsEachSettingWithItsValue(t *testing.T) {
logRemoteAppName: hostname,
alertWebhookURL: "",
alertWebhookHeaders: "",
alertSlackWebhookURL: "",
alertNtfyURL: "",
alertNtfyToken: "",
alertEvents: defaultAlertEvents,
alertCooldown: defaultAlertCooldown,
alertMaxPerHour: "60",
+1 -1
View File
@@ -243,7 +243,7 @@ func TestFailureRaisesASourceFailureAlertOncePerCooldown(t *testing.T) {
wantCountry(t, g, clients(), "")
wantRequests(t, geojs, 2)
waiting := queue.Snapshot().Waiting
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || !reflect.DeepEqual(waiting[0], want) {
t.Errorf("alerts waiting %+v, want only %+v", waiting, want)
}
+40 -37
View File
@@ -226,45 +226,48 @@ func (m *Metrics) AddRemoteLog(remote *remotelog.Sender) {
)
}
// AddAlerts adds the metrics of the alerts sent to
// SWWAF_ALERT_WEBHOOK_URL, read from queue as the metrics are asked for,
// with the destination webhook: the alerts sent, the requests to the
// webhook that failed, and the alerts held back and dropped.
// AddAlerts adds the metrics of the alerts sent to each destination set,
// read from queue as the metrics are asked for, by destination: the
// alerts sent, the requests to the destination that failed, the alerts
// held back, which are the same for every destination, and those
// dropped. With no destination set, it adds none.
func (m *Metrics) AddAlerts(queue *alerts.Queue) {
webhook := prometheus.Labels{"destination": "webhook"}
for _, name := range queue.DestinationsSet() {
destination := prometheus.Labels{"destination": name}
m.registry.MustRegister(
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_sent_total",
Help: "Alerts the destination took.",
ConstLabels: webhook,
}, func() float64 {
return float64(queue.Sent())
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_failed_total",
Help: "Requests to the destination that failed.",
ConstLabels: webhook,
}, func() float64 {
return float64(queue.Failed())
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_suppressed_total",
Help: "Alerts held back: repeats within SWWAF_ALERT_COOLDOWN, and " +
"alerts past SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary.",
ConstLabels: webhook,
}, func() float64 {
return float64(queue.Suppressed())
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_dropped_total",
Help: "Alerts dropped, the oldest first, from a full queue, and alerts " +
"given up as the destination refused them.",
ConstLabels: webhook,
}, func() float64 {
return float64(queue.Dropped())
}),
)
m.registry.MustRegister(
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_sent_total",
Help: "Alerts the destination took.",
ConstLabels: destination,
}, func() float64 {
return float64(queue.Counts(name).Sent)
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_failed_total",
Help: "Requests to the destination that failed.",
ConstLabels: destination,
}, func() float64 {
return float64(queue.Counts(name).Failed)
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_suppressed_total",
Help: "Alerts held back: repeats within SWWAF_ALERT_COOLDOWN, and " +
"alerts past SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary.",
ConstLabels: destination,
}, func() float64 {
return float64(queue.Suppressed())
}),
prometheus.NewCounterFunc(prometheus.CounterOpts{
Name: "smallwebwaf_alerts_dropped_total",
Help: "Alerts dropped, the oldest first, from a full queue, and " +
"alerts given up as the destination refused them.",
ConstLabels: destination,
}, func() float64 {
return float64(queue.Counts(name).Dropped)
}),
)
}
}
// ServeHTTP answers with the metrics in the Prometheus text format.
+4 -3
View File
@@ -124,7 +124,7 @@ func TestObserveModeRaisesTheBanAlertsItWouldHave(t *testing.T) {
"for the attack alone, as it was", held, line.BanExpires)
}
waiting := queue.Snapshot().Waiting
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 3 || queue.Suppressed() != 0 {
t.Fatalf("%d alerts wait and %d are held back, want 3 and 0: %+v",
len(waiting), queue.Suppressed(), waiting)
@@ -189,7 +189,8 @@ func TestObserveModeWorksOutABanOnlyWhenItsAlertWouldBeSent(t *testing.T) {
// Had the ban been worked out for any of the requests within the
// cooldown or past the two an hour, its alert would have been raised,
// held back and counted.
if waiting := queue.Snapshot().Waiting; len(waiting) != 2 || queue.Suppressed() != 0 {
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 2 || queue.Suppressed() != 0 {
t.Errorf("%d alerts wait and %d are held back, want 2 and 0: %+v",
len(waiting), queue.Suppressed(), waiting)
}
@@ -250,7 +251,7 @@ func attackAlert(
func wantAlerts(t *testing.T, queue *alerts.Queue, want ...alerts.Alert) {
t.Helper()
got := queue.Snapshot().Waiting
got := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(got) != len(want) {
t.Fatalf("%d alerts wait, want %d: %+v", len(got), len(want), got)
}
+1 -1
View File
@@ -387,7 +387,7 @@ func TestBrokenEditKeepsTheRulesAsTheyWere(t *testing.T) {
wantFileError := func() {
t.Helper()
waiting := queue.Snapshot().Waiting
waiting := queue.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || waiting[0].Event != alerts.EventFileError ||
waiting[0].Reason != hasError || waiting[0].Detail["error"] != want ||
waiting[0].Detail["file"] != filepath.Join(dir, firstFile) {
+6 -5
View File
@@ -122,9 +122,7 @@ func Run(ctx context.Context, params Params) int {
server.Metrics.AddRemoteLog(remote)
}
if cfg.AlertWebhookURL != nil {
server.Metrics.AddAlerts(alertQueue)
}
server.Metrics.AddAlerts(alertQueue)
files, err := loadStateFiles(cfg, server, alertQueue, now, processLog)
if err != nil {
@@ -169,14 +167,17 @@ func loadStateFiles(
})
}
// newAlertQueue returns the queue of the alerts to SWWAF_ALERT_WEBHOOK_URL,
// with the settings for it.
// newAlertQueue returns the queue of the alerts to the webhook, Slack and
// ntfy, with the settings for them.
func newAlertQueue(
cfg *config.Config, now func() time.Time, processLog *slog.Logger,
) *alerts.Queue {
return alerts.New(alerts.Params{
WebhookURL: cfg.AlertWebhookURL,
WebhookHeaders: cfg.AlertWebhookHeaders,
SlackURL: cfg.AlertSlackWebhookURL,
NtfyURL: cfg.AlertNtfyURL,
NtfyToken: cfg.AlertNtfyToken,
Events: cfg.AlertEvents,
Cooldown: cfg.AlertCooldown,
MaxPerHour: cfg.AlertMaxPerHour,
+143 -7
View File
@@ -38,7 +38,8 @@ const (
stateCounterInterval = "SWWAF_STATE_COUNTER_INTERVAL"
rateLimitPerDay = "SWWAF_RATE_LIMIT_PER_DAY"
rulesDir = "SWWAF_RULES_DIR"
adminToken = "SWWAF_ADMIN_TOKEN" //nolint:gosec // the setting's name
adminToken = "SWWAF_ADMIN_TOKEN" //nolint:gosec // the setting's name
metricsToken = "SWWAF_METRICS_TOKEN" //nolint:gosec // the setting's name
// adminSecret is the SWWAF_ADMIN_TOKEN the tests set.
adminSecret = "fedcba9876543210fedcba9876543210"
// greeting is what the tests' app answers.
@@ -136,7 +137,7 @@ func TestShortTokenStopsTheStartUnshown(t *testing.T) {
const token = "a-token-of-31-characters-at-all" //nolint:gosec // too short to use
for _, name := range []string{adminToken, "SWWAF_METRICS_TOKEN"} {
for _, name := range []string{adminToken, metricsToken} {
t.Run(name, func(t *testing.T) {
t.Parallel()
@@ -517,7 +518,7 @@ func TestStalledRemoteLogEndpointHoldsUpNoRequest(t *testing.T) {
rulesDir: t.TempDir(),
"SWWAF_LOG_REMOTE_URL": "syslog+tls://" + endpoint.Addr().String(),
"SWWAF_LOG_REMOTE_BUFFER": "1",
"SWWAF_METRICS_TOKEN": token,
metricsToken: token,
}
out := runUntilStopped(t, env, func(url string) {
@@ -591,7 +592,7 @@ func TestBanIsAlertedAndAnAlertNotSentIsKeptAcrossARestart(t *testing.T) {
// alerts.json keeps it as smallwebwaf stops, and once started again,
// smallwebwaf sends it.
var file struct {
Waiting []struct {
Waiting map[string][]struct {
Event string `json:"event"`
} `json:"waiting"`
}
@@ -603,8 +604,10 @@ func TestBanIsAlertedAndAnAlertNotSentIsKeptAcrossARestart(t *testing.T) {
err = json.Unmarshal(data, &file)
}
if err != nil || len(file.Waiting) != 1 || file.Waiting[0].Event != "permanent_ban" {
t.Fatalf("alerts.json holds %s (%v), want the permanent_ban alert waiting", data, err)
waiting := file.Waiting["webhook"]
if err != nil || len(waiting) != 1 || waiting[0].Event != "permanent_ban" {
t.Fatalf("alerts.json holds %s (%v), want the permanent_ban alert waiting for "+
"the webhook", data, err)
}
// It counts the alert sent in the metrics, read here from a client the
@@ -614,7 +617,7 @@ func TestBanIsAlertedAndAnAlertNotSentIsKeptAcrossARestart(t *testing.T) {
webhook.failing.Store(false)
env["SWWAF_ALLOW_NETS"] = localhost
env["SWWAF_METRICS_TOKEN"] = token
env[metricsToken] = token
runUntilStopped(t, env, func(url string) {
webhook.waitFor(t, "permanent_ban", true)
@@ -639,6 +642,79 @@ func TestBanIsAlertedAndAnAlertNotSentIsKeptAcrossARestart(t *testing.T) {
})
}
func TestBanIsAlertedToSlackAndNtfy(t *testing.T) {
t.Parallel()
const (
ntfyToken = "tk_0123456789abcdefghijklmnopq"
token = "abcdef0123456789abcdef0123456789"
client = "203.0.113.9"
)
slack, ntfy := startDestination(t), startDestination(t)
env := map[string]string{
listenAddr: localhost + ":0",
upstreamURL: startApp(t),
stateDir: t.TempDir(),
rulesDir: t.TempDir(),
trustedProxies: localhost + "/32",
rateLimitPerDay: "1",
// The metrics are read from 127.0.0.1, which no limit counts.
"SWWAF_ALLOW_NETS": localhost + "/32",
metricsToken: token,
"SWWAF_INSTANCE_NAME": "fsn1app1/gitea",
"SWWAF_ALERT_SLACK_WEBHOOK_URL": slack.url,
"SWWAF_ALERT_NTFY_URL": ntfy.url,
"SWWAF_ALERT_NTFY_TOKEN": ntfyToken,
}
runUntilStopped(t, env, func(url string) {
// The client's second request breaks the day limit, and bans it;
// Slack and ntfy are each sent the alert.
wantStatus(t, url, client, http.StatusOK)
wantStatus(t, url, client, http.StatusForbidden)
var message struct {
Text string `json:"text"`
}
slackPost := slack.firstPost(t)
err := json.Unmarshal([]byte(slackPost.body), &message)
if err != nil || !strings.HasPrefix(message.Text, "*fsn1app1/gitea: ban*\n") ||
!strings.Contains(message.Text, "\nclient: "+client+"\n") {
t.Errorf("Slack was sent %s", slackPost.body)
}
ntfyPost := ntfy.firstPost(t)
if ntfyPost.header.Get("Title") != "fsn1app1/gitea: ban" ||
ntfyPost.header.Get("Authorization") != "Bearer "+ntfyToken ||
!strings.Contains(ntfyPost.body, "\nclient: "+client+"\n") {
t.Errorf("ntfy was sent %s, with the headers %v", ntfyPost.body,
ntfyPost.header)
}
// The metrics count it for each, and give no series for the
// webhook, which is not set. As long as that takes, so that a slow
// test process cannot fail the test.
sent := []string{
"\nsmallwebwaf_alerts_sent_total{destination=\"slack\"} 1\n",
"\nsmallwebwaf_alerts_sent_total{destination=\"ntfy\"} 1\n",
}
metrics := metricsText(t, url+"_smallwebwaf/metrics", token)
for !strings.Contains(metrics, sent[0]) || !strings.Contains(metrics, sent[1]) {
time.Sleep(pollInterval)
metrics = metricsText(t, url+"_smallwebwaf/metrics", token)
}
if strings.Contains(metrics, `destination="webhook"`) {
t.Errorf("the metrics give the webhook:\n%s", metrics)
}
})
}
func TestStateFileThatDoesNotParseStopsTheStart(t *testing.T) {
t.Parallel()
@@ -1009,6 +1085,66 @@ func saveUntilAnswered(t *testing.T, path, content, url, from string, status int
}
}
// destination is a stand-in for SWWAF_ALERT_SLACK_WEBHOOK_URL or
// SWWAF_ALERT_NTFY_URL. It notes each request it is sent, and answers
// 200.
type destination struct {
url string
mu sync.Mutex
posts []destinationPost
}
// destinationPost is a request a destination was sent: its headers and
// its body.
type destinationPost struct {
header http.Header
body string
}
// startDestination starts a destination that takes every alert.
func startDestination(t *testing.T) *destination {
t.Helper()
d := &destination{}
server := httptest.NewServer(http.HandlerFunc(
func(_ http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
d.mu.Lock()
d.posts = append(d.posts, destinationPost{
header: r.Header.Clone(), body: string(body),
})
d.mu.Unlock()
}))
t.Cleanup(server.Close)
d.url = server.URL + "/alerts"
return d
}
// firstPost waits until the destination has been sent a request, and
// returns the first. It waits as long as that takes, so that a slow test
// process cannot fail the test.
func (d *destination) firstPost(t *testing.T) destinationPost {
t.Helper()
for {
d.mu.Lock()
if len(d.posts) > 0 {
post := d.posts[0]
d.mu.Unlock()
return post
}
d.mu.Unlock()
time.Sleep(pollInterval)
}
}
// webhook is a stand-in for SWWAF_ALERT_WEBHOOK_URL. It notes each alert
// it is sent, and answers 204, or 503 while failing.
type webhook struct {
+48 -20
View File
@@ -2,11 +2,12 @@
// SWWAF_STATE_DIR, as the "Persistent state" section of SPEC.md describes:
// bans.json holds the bans, clients.json each client's counters and
// history, lookups.json GeoJS's answers, and alerts.json the cooldowns,
// the hour under way and the alerts waiting. Load 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 written outside the
// parts' locks, which are held only to take a snapshot or to put in what
// a file holds, so that no request waits on the disk.
// the hour under way and the alerts waiting for each destination. Load
// 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
// written outside the parts' locks, which are held only to take a
// snapshot or to put in what a file holds, so that no request waits on
// the disk.
package state
import (
@@ -18,9 +19,11 @@ import (
"fmt"
"io/fs"
"log/slog"
"maps"
"net/netip"
"os"
"path/filepath"
"slices"
"sync"
"time"
@@ -50,8 +53,12 @@ const (
var (
errVersion = errors.New("unknown version")
// errMissing is for an entry without a field it needs.
errMissing = errors.New("has no")
errCause = errors.New("is not limit, attack or admin")
errMissing = errors.New("has no")
errCause = errors.New("is not limit, attack or admin")
errDestination = errors.New("is not webhook, slack or ntfy")
errWaitingList = errors.New(`waiting is a list, but now lists the alerts by ` +
`destination: put the list under "webhook", as "waiting": {"webhook": [...]}, ` +
`or remove the file`)
)
// Params are what Load needs.
@@ -129,10 +136,10 @@ type lookupsFile struct {
// alertsFile is alerts.json, indented for an admin to read and edit.
type alertsFile struct {
Version int `json:"version"`
Cooldowns []alerts.Cooldown `json:"cooldowns"`
Hour alerts.Hour `json:"hour"`
Waiting []alerts.Alert `json:"waiting"`
Version int `json:"version"`
Cooldowns []alerts.Cooldown `json:"cooldowns"`
Hour alerts.Hour `json:"hour"`
Waiting map[string][]alerts.Alert `json:"waiting"`
}
// stateFile is the struct of a state file. Once the file is decoded, its
@@ -392,6 +399,17 @@ func (f *Files) takeIn(name string, data []byte, edit bool) (int, error) {
f.params.GeoJS.Load(file.Lookups)
entries = len(file.Lookups)
case alertsJSON:
// waiting was a list, of the alerts waiting for the webhook, before
// alerts went to Slack and ntfy too.
var written struct {
Waiting json.RawMessage `json:"waiting"`
}
if json.Unmarshal(data, &written) == nil &&
bytes.HasPrefix(written.Waiting, []byte("[")) {
return 0, fmt.Errorf("%s: %w", path, errWaitingList)
}
var file alertsFile
err := parse(path, data, &file)
@@ -402,7 +420,10 @@ func (f *Files) takeIn(name string, data []byte, edit bool) (int, error) {
f.params.Alerts.Load(alerts.State{
Cooldowns: file.Cooldowns, Hour: file.Hour, Waiting: file.Waiting,
})
entries = len(file.Waiting)
for _, waiting := range file.Waiting {
entries += len(waiting)
}
}
f.sums[name] = sha256.Sum256(data)
@@ -645,8 +666,9 @@ func (f *lookupsFile) check(data []byte) error {
}
// check refuses a cooldown without its event or when its alert was sent,
// which would hold back no repeat, and an alert waiting without its event
// or its time.
// which would hold back no repeat, alerts waiting for a destination with
// another name than webhook, slack or ntfy, most likely misspelt, and an
// alert waiting without its event or its time.
func (f *alertsFile) check([]byte) error {
for i, cooldown := range f.Cooldowns {
switch {
@@ -657,12 +679,18 @@ func (f *alertsFile) check([]byte) error {
}
}
for i, alert := range f.Waiting {
switch {
case alert.Event == "":
return fmt.Errorf("waiting %w", missing(i, "event"))
case alert.Time.IsZero():
return fmt.Errorf("waiting %w", missing(i, "time"))
for _, destination := range slices.Sorted(maps.Keys(f.Waiting)) {
if !slices.Contains(alerts.Destinations(), destination) {
return fmt.Errorf("waiting %q %w", destination, errDestination)
}
for i, alert := range f.Waiting[destination] {
switch {
case alert.Event == "":
return fmt.Errorf("waiting %s %w", destination, missing(i, "event"))
case alert.Time.IsZero():
return fmt.Errorf("waiting %s %w", destination, missing(i, "time"))
}
}
}
+73 -46
View File
@@ -114,39 +114,41 @@ const filledAlertsJSON = `{
"source_failure": 1
}
},
"waiting": [
{
"instance": "fsn1app1/gitea",
"time": "2026-10-06T00:00:00Z",
"event": "ban",
"client": "203.0.113.9",
"netblock": "203.0.113.9/32",
"asn": "",
"as_name": "",
"country": "DE",
"reason": "requests per minute over the limit of 1",
"detail": {
"cause": "limit"
"waiting": {
"webhook": [
{
"instance": "fsn1app1/gitea",
"time": "2026-10-06T00:00:00Z",
"event": "ban",
"client": "203.0.113.9",
"netblock": "203.0.113.9/32",
"asn": "",
"as_name": "",
"country": "DE",
"reason": "requests per minute over the limit of 1",
"detail": {
"cause": "limit"
},
"suppressed_repeats": 0
},
"suppressed_repeats": 0
},
{
"instance": "fsn1app1/gitea",
"time": "2026-10-06T00:00:00Z",
"event": "file_error",
"client": "",
"netblock": "",
"asn": "",
"as_name": "",
"country": "",
"reason": "writing the state files failed",
"detail": {
"error": "no space left on device",
"file": "/var/lib/smallwebwaf/bans.json"
},
"suppressed_repeats": 0
}
]
{
"instance": "fsn1app1/gitea",
"time": "2026-10-06T00:00:00Z",
"event": "file_error",
"client": "",
"netblock": "",
"asn": "",
"as_name": "",
"country": "",
"reason": "writing the state files failed",
"detail": {
"error": "no space left on device",
"file": "/var/lib/smallwebwaf/bans.json"
},
"suppressed_repeats": 0
}
]
}
}
`
@@ -234,7 +236,7 @@ func TestSourceFailureCooldownKeptInAlertsJSONAcrossARestart(t *testing.T) {
load(t, after)
after.Alerts.Raise(failure)
waiting := after.Alerts.Snapshot().Waiting
waiting := after.Alerts.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || after.Alerts.Suppressed() != 1 {
t.Errorf("%d alerts wait and %d are held back, want the one read back and 1",
len(waiting), after.Alerts.Suppressed())
@@ -270,7 +272,7 @@ func TestMissingFilesAreEmptyState(t *testing.T) {
held := params.Alerts.Snapshot()
if len(params.Ledger.Snapshot()) != 0 || len(params.Limiter.Snapshot()) != 0 ||
len(params.GeoJS.Snapshot()) != 0 || len(held.Cooldowns) != 0 ||
len(held.Waiting) != 0 || held.Hour.Sent != 0 {
len(held.Waiting[alerts.DestinationWebhook]) != 0 || held.Hour.Sent != 0 {
t.Error("state from no files")
}
}
@@ -314,9 +316,14 @@ func TestFileThatDoesNotParseStopsTheStart(t *testing.T) {
},
{
"an unknown field of an alert waiting", alertsJSON,
`{"version": 1, "waiting": [{"event": "ban", "evnet": "ban"}]}`,
`{"version": 1, "waiting": {"webhook": [{"event": "ban", "evnet": "ban"}]}}`,
`: json: unknown field "evnet"`,
},
{
"alerts waiting for an unknown destination", alertsJSON,
`{"version": 1, "waiting": {"webhook": [], "slak": []}}`,
`: waiting "slak" is not webhook, slack or ntfy`,
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
@@ -426,13 +433,13 @@ func TestAlertsJSONEntryWithoutAFieldItNeedsStopsTheStart(t *testing.T) {
},
{
"an alert waiting without its event",
`{"version": 1, "waiting": [{"time": "2026-10-06T00:00:00Z"}]}`,
`: waiting entry 1 has no "event"`,
`{"version": 1, "waiting": {"ntfy": [{"time": "2026-10-06T00:00:00Z"}]}}`,
`: waiting ntfy entry 1 has no "event"`,
},
{
"an alert waiting without its time",
`{"version": 1, "waiting": [{"event": "ban"}]}`,
`: waiting entry 1 has no "time"`,
`{"version": 1, "waiting": {"slack": [{"event": "ban"}]}}`,
`: waiting slack entry 1 has no "time"`,
},
} {
t.Run(tc.name, func(t *testing.T) {
@@ -443,6 +450,23 @@ func TestAlertsJSONEntryWithoutAFieldItNeedsStopsTheStart(t *testing.T) {
}
}
func TestAlertsJSONWithWaitingAsAListStopsTheStartSayingWhatToChange(t *testing.T) {
t.Parallel()
// alerts.json as it was written before alerts went to Slack and ntfy too,
// with no alert waiting, or one.
for _, waiting := range []string{
`[]`,
`[{"event": "ban", "time": "2026-10-06T00:00:00Z"}]`,
} {
wantRefused(t, alertsJSON, `{"version": 1, "cooldowns": [], `+
`"hour": {"start": "2026-10-06T00:00:00Z", "sent": 0, "held_back": {}}, `+
`"waiting": `+waiting+`}`,
`: waiting is a list, but now lists the alerts by destination: put the `+
`list under "webhook", as "waiting": {"webhook": [...]}, or remove the file`)
}
}
func TestBanWithAnotherCauseStopsTheStart(t *testing.T) {
t.Parallel()
@@ -629,7 +653,7 @@ func TestWriteThatFailsWhileRunningRaisesAFileErrorAlertOncePerCooldown(t *testi
time.Sleep(time.Minute)
synctest.Wait()
waiting := params.Alerts.Snapshot().Waiting
waiting := params.Alerts.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 {
t.Fatalf("%d alerts wait, want 1", len(waiting))
}
@@ -648,9 +672,10 @@ func TestWriteThatFailsWhileRunningRaisesAFileErrorAlertOncePerCooldown(t *testi
time.Sleep(time.Minute)
synctest.Wait()
if len(params.Alerts.Snapshot().Waiting) != 1 || params.Alerts.Suppressed() != 1 {
waiting = params.Alerts.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || params.Alerts.Suppressed() != 1 {
t.Errorf("%d alerts wait and %d are held back, want 1 and 1",
len(params.Alerts.Snapshot().Waiting), params.Alerts.Suppressed())
len(waiting), params.Alerts.Suppressed())
}
})
}
@@ -857,7 +882,7 @@ func TestEditOfEachFileTakenIn(t *testing.T) {
// in.
edit(t, dir, alertsJSON, `{"version": 1, "cooldowns": [{"event": "ban", `+
`"netblock": "198.51.100.9/24", "sent": "2026-10-06T00:00:00Z"}], `+
`"waiting": [{"event": "file_error", "time": "2026-10-06T00:00:00Z"}]}`)
`"waiting": {"webhook": [{"event": "file_error", "time": "2026-10-06T00:00:00Z"}]}}`)
wantTakenIn(t, lines, dir, alertsJSON)
want := alerts.State{
@@ -865,8 +890,10 @@ func TestEditOfEachFileTakenIn(t *testing.T) {
Event: alerts.EventBan, Netblock: netip.MustParsePrefix("198.51.100.0/24"),
Sent: midnight(),
}},
Hour: alerts.Hour{HeldBack: map[string]int{}},
Waiting: []alerts.Alert{{Event: alerts.EventFileError, Time: midnight()}},
Hour: alerts.Hour{HeldBack: map[string]int{}},
Waiting: map[string][]alerts.Alert{
alerts.DestinationWebhook: {{Event: alerts.EventFileError, Time: midnight()}},
},
}
if got := params.Alerts.Snapshot(); !reflect.DeepEqual(got, want) {
t.Errorf("%s taken in as\n%+v\nwant\n%+v", alertsJSON, got, want)
@@ -1096,7 +1123,7 @@ func TestBrokenEditSetAsideAtTheNextWrite(t *testing.T) {
}
// It is raised as a file_error alert, with the same file and error.
waiting := params.Alerts.Snapshot().Waiting
waiting := params.Alerts.Snapshot().Waiting[alerts.DestinationWebhook]
if len(waiting) != 1 || waiting[0].Event != alerts.EventFileError ||
waiting[0].Detail["file"] != path+".bad" || waiting[0].Detail["error"] != message {
t.Errorf("alerts waiting %+v, want a file_error alert for %s", waiting, path+".bad")