Alerts to a JSON webhook, with a cooldown and an hourly summary (closes #26)
check / check (push) Waiting to run
check / check (push) Waiting to run
SWWAF_ALERT_WEBHOOK_URL gets one JSON POST per alert, in SPEC.md's schema, with SWWAF_ALERT_WEBHOOK_HEADERS: ban and permanent_ban, with the ban's notes, in observe mode too, marked mode observe and worked out only when the alert would be sent; source_failure for GeoJS; file_error for a rule or state file with an error. SWWAF_ALERT_EVENTS chooses; SWWAF_ALERT_COOLDOWN holds back repeats by netblock, file or source; past SWWAF_ALERT_MAX_PER_HOUR the hour ends in one summary. A bounded queue, retried with backoff, holds up no request; a 4xx other than 408 and 429 gives the alert up. alerts.json keeps the queue, the cooldowns and the hour. Nothing shows the URL's path or query. Judgement call: the summary's event is summary, which SPEC.md omits. Judgement call: an admin's ban raises no alert. Model: opus-5-5
This commit was merged in pull request #93.
This commit is contained in:
@@ -0,0 +1,609 @@
|
||||
// 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.
|
||||
package alerts
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"cmp"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/netip"
|
||||
"net/url"
|
||||
"slices"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The events an alert is for, as SWWAF_ALERT_EVENTS names them.
|
||||
const (
|
||||
// EventBan is a ban smallwebwaf made.
|
||||
EventBan = "ban"
|
||||
// EventPermanentBan is a permanent ban smallwebwaf made, or a ban it
|
||||
// made permanent.
|
||||
EventPermanentBan = "permanent_ban"
|
||||
// EventWAFBlock, EventAnomaly and EventReputationHit come with the
|
||||
// Core Rule Set, the anomaly thresholds and the reputation sources;
|
||||
// nothing raises them yet.
|
||||
EventWAFBlock = "waf_block"
|
||||
EventAnomaly = "anomaly"
|
||||
EventReputationHit = "reputation_hit"
|
||||
// EventSourceFailure is GeoJS failing or refusing smallwebwaf.
|
||||
EventSourceFailure = "source_failure"
|
||||
// EventFileError is a rule file or state file edited while smallwebwaf
|
||||
// runs that does not parse, or a state file that cannot be written.
|
||||
EventFileError = "file_error"
|
||||
// EventSummary is the summary of the alerts an hour held back past
|
||||
// SWWAF_ALERT_MAX_PER_HOUR. SWWAF_ALERT_EVENTS does not name it.
|
||||
EventSummary = "summary"
|
||||
)
|
||||
|
||||
// Events returns every event SWWAF_ALERT_EVENTS can name, which is its
|
||||
// default.
|
||||
func Events() []string {
|
||||
return []string{
|
||||
EventBan, EventPermanentBan, EventWAFBlock, EventAnomaly,
|
||||
EventReputationHit, EventSourceFailure, EventFileError,
|
||||
}
|
||||
}
|
||||
|
||||
const (
|
||||
// queueSize is the most alerts that wait to be sent. Past it, the
|
||||
// oldest is dropped.
|
||||
queueSize = 1000
|
||||
// sendTimeout bounds one request to the webhook.
|
||||
sendTimeout = 10 * time.Second
|
||||
// After a request to the webhook 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 = 64 << 10
|
||||
)
|
||||
|
||||
var (
|
||||
errStatus = errors.New("the webhook answered")
|
||||
// errRefused is a 4xx answer other than 408 and 429: the webhook
|
||||
// refuses the alert itself, and would refuse it again.
|
||||
errRefused = errors.New("the webhook refused the alert, answering")
|
||||
)
|
||||
|
||||
// Params are what New needs.
|
||||
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 *url.URL
|
||||
WebhookHeaders http.Header
|
||||
// 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
|
||||
// (SWWAF_ALERT_COOLDOWN), 0 for no time. MaxPerHour is the most alerts
|
||||
// sent in an hour (SWWAF_ALERT_MAX_PER_HOUR), 0 for no limit.
|
||||
Cooldown time.Duration
|
||||
MaxPerHour int
|
||||
// Instance is SWWAF_INSTANCE_NAME, which every alert gives.
|
||||
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 *slog.Logger
|
||||
}
|
||||
|
||||
// Alert is one alert, as the webhook is sent it and alerts.json holds it,
|
||||
// with the fields of the "Alert webhook schema" section of SPEC.md. ASN
|
||||
// and ASName are empty until AS numbers are looked up.
|
||||
//
|
||||
//nolint:tagliatelle // SPEC.md's alert webhook schema names its fields in snake_case
|
||||
type Alert struct {
|
||||
Instance string `json:"instance"`
|
||||
Time time.Time `json:"time"`
|
||||
Event string `json:"event"`
|
||||
Client netip.Addr `json:"client"`
|
||||
Netblock netip.Prefix `json:"netblock"`
|
||||
ASN string `json:"asn"`
|
||||
ASName string `json:"as_name"`
|
||||
Country string `json:"country"`
|
||||
// Reason is a short sentence, and Detail what is particular to the
|
||||
// event: for a file_error, its "file", and for a source_failure, its
|
||||
// "source", which the cooldown tells repeats by.
|
||||
Reason string `json:"reason"`
|
||||
Detail map[string]any `json:"detail"`
|
||||
// SuppressedRepeats is how many repeats of the alert the cooldown
|
||||
// held back since the last one let through.
|
||||
SuppressedRepeats int `json:"suppressed_repeats"`
|
||||
}
|
||||
|
||||
// Cooldown is, for an event on a netblock, or about a file or a source,
|
||||
// when the last alert let through was raised, and how many repeats the
|
||||
// cooldown has held back since, as alerts.json holds it.
|
||||
//
|
||||
//nolint:tagliatelle // the state files use snake_case, as the request log does
|
||||
type Cooldown struct {
|
||||
Event string `json:"event"`
|
||||
Netblock netip.Prefix `json:"netblock"`
|
||||
File string `json:"file,omitempty"`
|
||||
Source string `json:"source,omitempty"`
|
||||
Sent time.Time `json:"sent"`
|
||||
SuppressedRepeats int `json:"suppressed_repeats"`
|
||||
}
|
||||
|
||||
// Hour is the hour under way, by the clock, as alerts.json holds it: when
|
||||
// it started, how many alerts were let through in it, and how many were
|
||||
// held back in it past MaxPerHour, by event, for its summary.
|
||||
//
|
||||
//nolint:tagliatelle // the state files use snake_case, as the request log does
|
||||
type Hour struct {
|
||||
Start time.Time `json:"start"`
|
||||
Sent int `json:"sent"`
|
||||
HeldBack map[string]int `json:"held_back"`
|
||||
}
|
||||
|
||||
// State is what alerts.json holds: the cooldowns, the hour under way, and
|
||||
// the alerts waiting to be sent, oldest first.
|
||||
type State struct {
|
||||
Cooldowns []Cooldown `json:"cooldowns"`
|
||||
Hour Hour `json:"hour"`
|
||||
Waiting []Alert `json:"waiting"`
|
||||
}
|
||||
|
||||
// Queue takes the alerts raised, holds back those it must, and sends the
|
||||
// others to the webhook. 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{}
|
||||
|
||||
mu sync.Mutex
|
||||
// cooldowns are the alerts last let through, by event and netblock,
|
||||
// file or source.
|
||||
cooldowns map[cooldownKey]*Cooldown
|
||||
hour Hour
|
||||
// waiting are the alerts waiting to be sent, oldest first.
|
||||
waiting []*Alert
|
||||
|
||||
sent, failed, suppressed, dropped atomic.Int64
|
||||
}
|
||||
|
||||
// cooldownKey is what makes an alert a repeat of another: the same event
|
||||
// on the same netblock, and about the same file or source, as its detail
|
||||
// names them. Each is empty for an alert without one.
|
||||
type cooldownKey struct {
|
||||
event string
|
||||
netblock netip.Prefix
|
||||
file string
|
||||
source string
|
||||
}
|
||||
|
||||
// cooldownKeyOf returns what makes another alert a repeat of alert.
|
||||
func cooldownKeyOf(alert *Alert) cooldownKey {
|
||||
file, _ := alert.Detail["file"].(string)
|
||||
source, _ := alert.Detail["source"].(string)
|
||||
|
||||
return cooldownKey{alert.Event, alert.Netblock, file, source}
|
||||
}
|
||||
|
||||
// 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),
|
||||
cooldowns: map[cooldownKey]*Cooldown{},
|
||||
hour: Hour{HeldBack: map[string]int{}},
|
||||
}
|
||||
}
|
||||
|
||||
// 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.
|
||||
func (q *Queue) Raise(alert Alert) {
|
||||
if q.params.WebhookURL == nil || !slices.Contains(q.params.Events, alert.Event) {
|
||||
return
|
||||
}
|
||||
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
now := q.params.Now()
|
||||
alert.Instance = q.params.Instance
|
||||
alert.Time = now
|
||||
|
||||
if q.repeat(&alert, now) {
|
||||
q.suppressed.Add(1)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
q.endHour(now)
|
||||
|
||||
if q.params.MaxPerHour > 0 && q.hour.Sent >= q.params.MaxPerHour {
|
||||
q.hour.HeldBack[alert.Event]++
|
||||
q.suppressed.Add(1)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
q.startCooldown(&alert, now)
|
||||
q.hour.Sent++
|
||||
q.queue(&alert)
|
||||
}
|
||||
|
||||
// WouldSend reports whether Raise would let an alert for event on
|
||||
// netblock through now: a webhook 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) {
|
||||
return false
|
||||
}
|
||||
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
now := q.params.Now()
|
||||
|
||||
last, found := q.cooldowns[cooldownKey{event: event, netblock: netblock}]
|
||||
if q.params.Cooldown > 0 && found && now.Sub(last.Sent) < q.params.Cooldown {
|
||||
return false
|
||||
}
|
||||
|
||||
q.endHour(now)
|
||||
|
||||
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.
|
||||
func (q *Queue) Run(ctx context.Context) {
|
||||
if q.params.WebhookURL == nil {
|
||||
return
|
||||
}
|
||||
|
||||
var (
|
||||
retryDelay time.Duration
|
||||
retryAt time.Time
|
||||
)
|
||||
|
||||
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))
|
||||
}
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
hourEnds.Stop()
|
||||
|
||||
return
|
||||
case <-q.queued:
|
||||
case <-hourEnds.C:
|
||||
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()
|
||||
}
|
||||
|
||||
// Failed is how many requests to the webhook have failed.
|
||||
func (q *Queue) Failed() int64 {
|
||||
return q.failed.Load()
|
||||
}
|
||||
|
||||
// Suppressed is how many alerts were held back: by the cooldown, and past
|
||||
// MaxPerHour.
|
||||
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 {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
state := State{
|
||||
Cooldowns: make([]Cooldown, 0, len(q.cooldowns)),
|
||||
Hour: q.hour,
|
||||
Waiting: make([]Alert, 0, len(q.waiting)),
|
||||
}
|
||||
state.Hour.HeldBack = maps.Clone(q.hour.HeldBack)
|
||||
|
||||
for _, cooldown := range q.cooldowns {
|
||||
state.Cooldowns = append(state.Cooldowns, *cooldown)
|
||||
}
|
||||
|
||||
slices.SortFunc(state.Cooldowns, func(a, b Cooldown) int {
|
||||
return cmp.Or(a.Netblock.Compare(b.Netblock), cmp.Compare(a.Event, b.Event),
|
||||
cmp.Compare(a.File, b.File), cmp.Compare(a.Source, b.Source))
|
||||
})
|
||||
|
||||
for _, alert := range q.waiting {
|
||||
state.Waiting = append(state.Waiting, *alert)
|
||||
}
|
||||
|
||||
return 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.
|
||||
func (q *Queue) Load(state State) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
q.cooldowns = map[cooldownKey]*Cooldown{}
|
||||
|
||||
for _, cooldown := range state.Cooldowns {
|
||||
cooldown.Netblock = cooldown.Netblock.Masked()
|
||||
key := cooldownKey{cooldown.Event, cooldown.Netblock, cooldown.File, cooldown.Source}
|
||||
q.cooldowns[key] = &cooldown
|
||||
}
|
||||
|
||||
q.hour = state.Hour
|
||||
q.hour.HeldBack = maps.Clone(state.Hour.HeldBack)
|
||||
|
||||
if q.hour.HeldBack == nil {
|
||||
q.hour.HeldBack = map[string]int{}
|
||||
}
|
||||
|
||||
q.waiting = nil
|
||||
|
||||
for _, alert := range state.Waiting {
|
||||
q.queue(&alert)
|
||||
}
|
||||
}
|
||||
|
||||
// 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 {
|
||||
if q.params.Cooldown == 0 {
|
||||
return false
|
||||
}
|
||||
|
||||
last, found := q.cooldowns[cooldownKeyOf(alert)]
|
||||
if !found || now.Sub(last.Sent) >= q.params.Cooldown {
|
||||
return false
|
||||
}
|
||||
|
||||
last.SuppressedRepeats++
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
// startCooldown gives alert, let through at now, the count of the repeats
|
||||
// held back since the last one let through, and notes alert as the last
|
||||
// one let through.
|
||||
func (q *Queue) startCooldown(alert *Alert, now time.Time) {
|
||||
if q.params.Cooldown == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
key := cooldownKeyOf(alert)
|
||||
|
||||
last, found := q.cooldowns[key]
|
||||
if found {
|
||||
alert.SuppressedRepeats = last.SuppressedRepeats
|
||||
}
|
||||
|
||||
q.cooldowns[key] = &Cooldown{
|
||||
Event: alert.Event, Netblock: alert.Netblock, File: key.file, Source: key.source,
|
||||
Sent: now,
|
||||
}
|
||||
}
|
||||
|
||||
// endHour ends the hour under way, if now is past it: it queues that
|
||||
// hour's summary when alerts were held back in it past MaxPerHour, and
|
||||
// forgets the cooldowns that have run out with no repeat held back, which
|
||||
// no alert needs any more.
|
||||
func (q *Queue) endHour(now time.Time) {
|
||||
start := now.Truncate(time.Hour)
|
||||
if !start.After(q.hour.Start) {
|
||||
return
|
||||
}
|
||||
|
||||
heldBack := 0
|
||||
for _, count := range q.hour.HeldBack {
|
||||
heldBack += count
|
||||
}
|
||||
|
||||
if heldBack > 0 {
|
||||
q.queue(&Alert{
|
||||
Instance: q.params.Instance,
|
||||
Time: now,
|
||||
Event: EventSummary,
|
||||
Reason: fmt.Sprintf("%d alerts held back in the hour from %s, past the %d "+
|
||||
"an hour SWWAF_ALERT_MAX_PER_HOUR allows", heldBack,
|
||||
q.hour.Start.Format(time.RFC3339), q.params.MaxPerHour),
|
||||
Detail: map[string]any{
|
||||
"hour": q.hour.Start, "count": heldBack, "events": q.hour.HeldBack,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
q.hour = Hour{Start: start, HeldBack: map[string]int{}}
|
||||
|
||||
for key, cooldown := range q.cooldowns {
|
||||
if now.Sub(cooldown.Sent) >= q.params.Cooldown && cooldown.SuppressedRepeats == 0 {
|
||||
delete(q.cooldowns, key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// queue adds alert to the alerts waiting, first dropping the oldest while
|
||||
// queueSize wait, and has Run look at the queue again.
|
||||
func (q *Queue) queue(alert *Alert) {
|
||||
if len(q.waiting) == queueSize {
|
||||
q.waiting = slices.Delete(q.waiting, 0, 1)
|
||||
q.dropped.Add(1)
|
||||
}
|
||||
|
||||
q.waiting = append(q.waiting, alert)
|
||||
|
||||
select {
|
||||
case q.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()
|
||||
|
||||
var oldest *Alert
|
||||
if len(q.waiting) > 0 {
|
||||
oldest = q.waiting[0]
|
||||
}
|
||||
|
||||
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
|
||||
// the oldest if it is there at all.
|
||||
func (q *Queue) remove(alert *Alert) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
if len(q.waiting) > 0 && q.waiting[0] == alert {
|
||||
q.waiting = slices.Delete(q.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
|
||||
// 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)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode the alert: %w", err)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(ctx, sendTimeout)
|
||||
defer cancel()
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
q.params.WebhookURL.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")
|
||||
|
||||
res, err := q.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 {
|
||||
return urlErr.Err
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
defer func() {
|
||||
_ = res.Body.Close()
|
||||
}()
|
||||
|
||||
// Read, so that the connection can be used again.
|
||||
_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, maxAnswerBytes))
|
||||
|
||||
switch status := res.StatusCode; {
|
||||
case status >= http.StatusOK && status < http.StatusMultipleChoices:
|
||||
return nil
|
||||
case status >= http.StatusBadRequest && status < http.StatusInternalServerError &&
|
||||
status != http.StatusRequestTimeout && status != http.StatusTooManyRequests:
|
||||
return fmt.Errorf("%w %s", errRefused, res.Status)
|
||||
default:
|
||||
return fmt.Errorf("%w %s", errStatus, res.Status)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,852 @@
|
||||
package alerts_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/netip"
|
||||
"net/url"
|
||||
"reflect"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/synctest"
|
||||
"time"
|
||||
|
||||
"sneak.berlin/go/smallwebwaf/internal/alerts"
|
||||
)
|
||||
|
||||
// The tests run in a synctest bubble, where the time package runs on a
|
||||
// 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.
|
||||
|
||||
const (
|
||||
// webhookURL is where the alerts are posted.
|
||||
webhookURL = "https://alerts.example/smallwebwaf?team=ops"
|
||||
// instance is the instance name every alert gives.
|
||||
instance = "fsn1app1/gitea"
|
||||
// started is when each test starts, as an alert gives it, and
|
||||
// anHourOn an hour later.
|
||||
started = "2000-01-01T00:00:00Z"
|
||||
anHourOn = "2000-01-01T01:00:00Z"
|
||||
// cooldown is the cooldown of most tests, the default.
|
||||
cooldown = 15 * time.Minute
|
||||
)
|
||||
|
||||
func TestAlertIsPostedAsJSONWithItsFieldsAndTheHeaders(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.WebhookHeaders = http.Header{
|
||||
"Authorization": {"Bearer 0123456789abcdef"},
|
||||
"X-Team": {"ops"},
|
||||
}
|
||||
webhook, q := start(t, params)
|
||||
|
||||
q.Raise(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 minute over the limit of 1000",
|
||||
Detail: map[string]any{"cause": "limit", "ban_expires": anHourOn},
|
||||
})
|
||||
synctest.Wait()
|
||||
|
||||
got := webhook.received()
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("the webhook had %d requests, want 1", len(got))
|
||||
}
|
||||
|
||||
if got[0].method != http.MethodPost || got[0].url != webhookURL {
|
||||
t.Errorf("request %s %s, want POST %s", got[0].method, got[0].url, webhookURL)
|
||||
}
|
||||
|
||||
for name, want := range map[string]string{
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": "Bearer 0123456789abcdef",
|
||||
"X-Team": "ops",
|
||||
} {
|
||||
if got[0].header.Get(name) != want {
|
||||
t.Errorf("header %s is %q, want %q", name, got[0].header.Get(name), want)
|
||||
}
|
||||
}
|
||||
|
||||
wantAlert(t, got[0].alert, map[string]any{
|
||||
"instance": instance,
|
||||
"time": started,
|
||||
"event": "ban",
|
||||
"client": "203.0.113.9",
|
||||
"netblock": "203.0.113.0/24",
|
||||
"asn": "",
|
||||
"as_name": "",
|
||||
"country": "DE",
|
||||
"reason": "requests per minute over the limit of 1000",
|
||||
"detail": map[string]any{"cause": "limit", "ban_expires": anHourOn},
|
||||
"suppressed_repeats": float64(0),
|
||||
})
|
||||
wantCounts(t, q, 1, 0, 0, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestOnlyTheChosenEventsAreSent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.Events = []string{alerts.EventSourceFailure, alerts.EventFileError}
|
||||
webhook, q := start(t, params)
|
||||
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventFileError, Reason: "a file error"})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventPermanentBan, Netblock: netblock(1)})
|
||||
synctest.Wait()
|
||||
|
||||
wantEvents(t, webhook, alerts.EventFileError)
|
||||
wantCounts(t, q, 1, 0, 0, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestWouldSendOnlyForTheChosenEvents(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
params := newParams()
|
||||
params.Events = []string{alerts.EventSourceFailure, alerts.EventFileError}
|
||||
q := alerts.New(params)
|
||||
|
||||
if q.WouldSend(alerts.EventBan, netblock(1)) ||
|
||||
q.WouldSend(alerts.EventPermanentBan, netblock(1)) {
|
||||
t.Error("a ban alert would be sent, though SWWAF_ALERT_EVENTS leaves it out")
|
||||
}
|
||||
|
||||
if !q.WouldSend(alerts.EventFileError, netip.Prefix{}) {
|
||||
t.Error("a file_error alert would not be sent")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNothingIsQueuedWithoutAWebhook(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
params := newParams()
|
||||
params.WebhookURL = nil
|
||||
q := alerts.New(params)
|
||||
|
||||
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))
|
||||
}
|
||||
}
|
||||
|
||||
func TestWouldSendNothingWithoutAWebhook(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
params := newParams()
|
||||
params.WebhookURL = nil
|
||||
q := alerts.New(params)
|
||||
|
||||
if q.WouldSend(alerts.EventBan, netblock(1)) {
|
||||
t.Error("an alert would be sent with no webhook set")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRepeatWithinTheCooldownIsHeldBackAndCountedInTheNext(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
webhook, q := start(t, newParams())
|
||||
raise := func(event string, n int) {
|
||||
q.Raise(alerts.Alert{Event: event, Netblock: netblock(n)})
|
||||
}
|
||||
|
||||
raise(alerts.EventBan, 1)
|
||||
|
||||
// The same event on the same netblock is a repeat; another netblock
|
||||
// or another event is not.
|
||||
time.Sleep(time.Minute)
|
||||
raise(alerts.EventBan, 1)
|
||||
raise(alerts.EventBan, 2)
|
||||
raise(alerts.EventPermanentBan, 1)
|
||||
|
||||
time.Sleep(cooldown - time.Minute - time.Nanosecond)
|
||||
raise(alerts.EventBan, 1)
|
||||
|
||||
// Once the cooldown has run out, the next one is sent with the
|
||||
// count of those held back.
|
||||
time.Sleep(time.Nanosecond)
|
||||
raise(alerts.EventBan, 1)
|
||||
|
||||
// And starts the cooldown again.
|
||||
time.Sleep(time.Minute)
|
||||
raise(alerts.EventBan, 1)
|
||||
synctest.Wait()
|
||||
|
||||
got := webhook.received()
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventBan, alerts.EventPermanentBan,
|
||||
alerts.EventBan)
|
||||
|
||||
for i, want := range []struct {
|
||||
netblock int
|
||||
repeats float64
|
||||
}{{1, 0}, {2, 0}, {1, 0}, {1, 2}} {
|
||||
alert := got[i].alert
|
||||
if alert["netblock"] != netblock(want.netblock).String() ||
|
||||
alert["suppressed_repeats"] != want.repeats {
|
||||
t.Errorf("alert %d is for %v with %v repeats, want %s with %v", i,
|
||||
alert["netblock"], alert["suppressed_repeats"], netblock(want.netblock),
|
||||
want.repeats)
|
||||
}
|
||||
}
|
||||
|
||||
wantCounts(t, q, 4, 0, 3, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestFileErrorAndSourceFailureRepeatOnlyForTheSameFileOrSource(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
webhook, q := start(t, params)
|
||||
fileError := func(file string) alerts.Alert {
|
||||
return alerts.Alert{
|
||||
Event: alerts.EventFileError,
|
||||
Detail: map[string]any{"file": file, "error": "line 2: an error"},
|
||||
}
|
||||
}
|
||||
sourceFailure := func(source string) alerts.Alert {
|
||||
return alerts.Alert{
|
||||
Event: alerts.EventSourceFailure, Detail: map[string]any{"source": source},
|
||||
}
|
||||
}
|
||||
|
||||
// Another file, or another source, is no repeat.
|
||||
q.Raise(fileError("/rules.d/50-a.rules"))
|
||||
q.Raise(fileError("/rules.d/50-b.rules"))
|
||||
q.Raise(fileError("/rules.d/50-a.rules"))
|
||||
q.Raise(sourceFailure("geojs"))
|
||||
q.Raise(sourceFailure("abuseipdb"))
|
||||
q.Raise(sourceFailure("geojs"))
|
||||
synctest.Wait()
|
||||
|
||||
// Each alert is named by its file, or its source.
|
||||
got := make([]string, 0, len(webhook.received()))
|
||||
for _, request := range webhook.received() {
|
||||
detail, _ := request.alert["detail"].(map[string]any)
|
||||
file, _ := detail["file"].(string)
|
||||
source, _ := detail["source"].(string)
|
||||
got = append(got, file+source)
|
||||
}
|
||||
|
||||
want := []string{
|
||||
"/rules.d/50-a.rules", "/rules.d/50-b.rules", "geojs", "abuseipdb",
|
||||
}
|
||||
if !slices.Equal(got, want) {
|
||||
t.Errorf("the webhook was sent alerts for %v, want %v", got, want)
|
||||
}
|
||||
|
||||
wantCounts(t, q, 4, 0, 2, 0)
|
||||
|
||||
// alerts.json keeps each file's cooldown: a new queue holds back
|
||||
// the next for the first file, and sends the one for a third.
|
||||
after := alerts.New(params)
|
||||
after.Load(roundTrip(t, q.Snapshot()))
|
||||
after.Raise(fileError("/rules.d/50-a.rules"))
|
||||
after.Raise(fileError("/rules.d/50-c.rules"))
|
||||
|
||||
waiting := after.Snapshot().Waiting
|
||||
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)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestNoCooldownSendsEveryRepeat(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.Cooldown = 0
|
||||
webhook, q := start(t, params)
|
||||
|
||||
for range 3 {
|
||||
q.Raise(alerts.Alert{Event: alerts.EventFileError})
|
||||
time.Sleep(time.Minute)
|
||||
}
|
||||
|
||||
synctest.Wait()
|
||||
|
||||
wantEvents(t, webhook, alerts.EventFileError, alerts.EventFileError,
|
||||
alerts.EventFileError)
|
||||
wantCounts(t, q, 3, 0, 0, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestAlertsPastTheHourlyLimitAreRolledIntoOneSummary(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.MaxPerHour = 2
|
||||
webhook, q := start(t, params)
|
||||
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(2)})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(3)})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventPermanentBan, Netblock: netblock(4)})
|
||||
q.Raise(alerts.Alert{Event: alerts.EventFileError})
|
||||
|
||||
// The summary is sent as the hour ends, and not before.
|
||||
time.Sleep(time.Hour - time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventBan)
|
||||
|
||||
time.Sleep(time.Nanosecond)
|
||||
synctest.Wait()
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventBan, alerts.EventSummary)
|
||||
|
||||
summary := webhook.received()[2].alert
|
||||
wantAlert(t, summary, map[string]any{
|
||||
"instance": instance,
|
||||
"time": anHourOn,
|
||||
"event": "summary",
|
||||
"client": "",
|
||||
"netblock": "",
|
||||
"asn": "",
|
||||
"as_name": "",
|
||||
"country": "",
|
||||
"reason": "3 alerts held back in the hour from 2000-01-01T00:00:00Z, " +
|
||||
"past the 2 an hour SWWAF_ALERT_MAX_PER_HOUR allows",
|
||||
"detail": map[string]any{
|
||||
"hour": started,
|
||||
"count": float64(3),
|
||||
"events": map[string]any{
|
||||
"ban": float64(1), "permanent_ban": float64(1), "file_error": float64(1),
|
||||
},
|
||||
},
|
||||
"suppressed_repeats": float64(0),
|
||||
})
|
||||
|
||||
// The next hour sends alerts again, and, with none held back, ends
|
||||
// without a summary.
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(5)})
|
||||
time.Sleep(time.Hour)
|
||||
synctest.Wait()
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventBan, alerts.EventSummary,
|
||||
alerts.EventBan)
|
||||
wantCounts(t, q, 4, 0, 3, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestRepeatsBeforeAnAlertPastTheHourlyLimitAreGivenByTheNextSent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.MaxPerHour = 1
|
||||
webhook, q := start(t, params)
|
||||
raise := func() {
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
}
|
||||
|
||||
// The hour's one alert, and two repeats the cooldown holds back.
|
||||
raise()
|
||||
raise()
|
||||
raise()
|
||||
|
||||
// Once the cooldown has run out, the next is past the hourly limit.
|
||||
time.Sleep(cooldown)
|
||||
raise()
|
||||
|
||||
// The next hour's first alert gives the two repeats, and the summary
|
||||
// the alert past the limit.
|
||||
time.Sleep(time.Hour - cooldown)
|
||||
synctest.Wait()
|
||||
raise()
|
||||
synctest.Wait()
|
||||
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventSummary, alerts.EventBan)
|
||||
|
||||
got := webhook.received()
|
||||
if len(got) == 3 {
|
||||
detail, _ := got[1].alert["detail"].(map[string]any)
|
||||
repeats := got[2].alert["suppressed_repeats"]
|
||||
|
||||
if detail["count"] != float64(1) || repeats != float64(2) {
|
||||
t.Errorf("the summary counts %v alerts, and the last alert gives %v "+
|
||||
"repeats, want 1 and 2", detail["count"], repeats)
|
||||
}
|
||||
}
|
||||
|
||||
wantCounts(t, q, 3, 0, 3, 0)
|
||||
})
|
||||
}
|
||||
|
||||
func TestFailedRequestIsSentAgainWithBackoff(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
log := &lockedBuffer{}
|
||||
params.ProcessLog = slog.New(slog.NewJSONHandler(log, nil))
|
||||
webhook, q := start(t, params)
|
||||
webhook.set(failing)
|
||||
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
|
||||
// A second after the first failure, then twice as long after each
|
||||
// further one, up to a minute.
|
||||
time.Sleep(200 * time.Second)
|
||||
synctest.Wait()
|
||||
|
||||
after := make([]time.Duration, 0, len(webhook.received()))
|
||||
for _, request := range webhook.received() {
|
||||
after = append(after, request.at.Sub(midnight()))
|
||||
}
|
||||
|
||||
want := []time.Duration{
|
||||
0, time.Second, 3 * time.Second, 7 * time.Second, 15 * time.Second,
|
||||
31 * time.Second, 63 * time.Second, 123 * time.Second, 183 * time.Second,
|
||||
}
|
||||
if !slices.Equal(after, want) {
|
||||
t.Errorf("requests at %v, want %v", after, want)
|
||||
}
|
||||
|
||||
wantCounts(t, q, 0, int64(len(want)), 0, 0)
|
||||
|
||||
if !strings.Contains(log.String(),
|
||||
`"msg":"sending an alert to SWWAF_ALERT_WEBHOOK_URL failed"`) {
|
||||
t.Errorf("process log %q names no failure", log.String())
|
||||
}
|
||||
|
||||
// Once the webhook answers, the alert is sent, and leaves the
|
||||
// queue.
|
||||
webhook.set(answering)
|
||||
time.Sleep(time.Minute)
|
||||
synctest.Wait()
|
||||
|
||||
got := webhook.received()
|
||||
if last := got[len(got)-1]; !last.answered ||
|
||||
last.alert["netblock"] != netblock(1).String() {
|
||||
t.Errorf("the last request was not the alert, answered")
|
||||
}
|
||||
|
||||
wantCounts(t, q, 1, int64(len(want)), 0, 0)
|
||||
|
||||
if waiting := q.Snapshot().Waiting; len(waiting) != 0 {
|
||||
t.Errorf("%d alerts still wait, want none", len(waiting))
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestRefusedAlertIsGivenUpAndTheNextSent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
log := &lockedBuffer{}
|
||||
params.ProcessLog = slog.New(slog.NewJSONHandler(log, nil))
|
||||
webhook, q := start(t, params)
|
||||
|
||||
// 429 and 408 are failures, and the alert is sent again; 400 refuses
|
||||
// it, and it is given up.
|
||||
webhook.set(http.StatusTooManyRequests)
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
synctest.Wait()
|
||||
webhook.set(http.StatusRequestTimeout)
|
||||
time.Sleep(time.Second)
|
||||
synctest.Wait()
|
||||
webhook.set(refusing)
|
||||
time.Sleep(2 * time.Second)
|
||||
synctest.Wait()
|
||||
|
||||
// The next alert is sent at once.
|
||||
webhook.set(answering)
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(2)})
|
||||
time.Sleep(time.Minute)
|
||||
synctest.Wait()
|
||||
|
||||
// Each request, by when it was sent, and the netblock of its alert.
|
||||
got := make([]string, 0, len(webhook.received()))
|
||||
for _, request := range webhook.received() {
|
||||
block, _ := request.alert["netblock"].(string)
|
||||
got = append(got, request.at.Sub(midnight()).String()+" "+block)
|
||||
}
|
||||
|
||||
want := []string{
|
||||
"0s " + netblock(1).String(), "1s " + netblock(1).String(),
|
||||
"3s " + netblock(1).String(), "3s " + netblock(2).String(),
|
||||
}
|
||||
if !slices.Equal(got, want) {
|
||||
t.Errorf("requests %v, want %v", got, want)
|
||||
}
|
||||
|
||||
wantCounts(t, q, 1, 3, 0, 1)
|
||||
|
||||
if !strings.Contains(log.String(),
|
||||
`"msg":"gave up an alert SWWAF_ALERT_WEBHOOK_URL refused"`) {
|
||||
t.Errorf("process log %q names no alert given up", log.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestFailedRequestIsLoggedWithoutTheURL(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
log := &lockedBuffer{}
|
||||
params.ProcessLog = slog.New(slog.NewJSONHandler(log, nil))
|
||||
webhook, q := start(t, params)
|
||||
webhook.set(hanging)
|
||||
|
||||
// The request is abandoned after 10 seconds, with an error from the
|
||||
// HTTP client, which names the URL.
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
time.Sleep(11 * time.Second)
|
||||
synctest.Wait()
|
||||
|
||||
logged := log.String()
|
||||
if !strings.Contains(logged,
|
||||
`"msg":"sending an alert to SWWAF_ALERT_WEBHOOK_URL failed"`) ||
|
||||
strings.Contains(logged, "alerts.example") || strings.Contains(logged, "team=ops") {
|
||||
t.Errorf("process log %q names no failure, or names the URL", logged)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestFullQueueDropsTheOldestAndRaiseNeverWaits(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.MaxPerHour = 0
|
||||
webhook, q := start(t, params)
|
||||
webhook.set(hanging)
|
||||
|
||||
// The webhook does not answer the first alert, while one more alert
|
||||
// than the queue holds is raised: none waits, and the oldest, the
|
||||
// one the webhook was sent, is dropped.
|
||||
for n := range alerts.QueueSize + 1 {
|
||||
q.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(n)})
|
||||
|
||||
if n == 0 {
|
||||
synctest.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
if took := time.Since(midnight()); took != 0 {
|
||||
t.Errorf("raising the alerts took %s, want no time", took)
|
||||
}
|
||||
|
||||
wantCounts(t, q, 0, 0, 0, 1)
|
||||
|
||||
waiting := q.Snapshot().Waiting
|
||||
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))
|
||||
}
|
||||
|
||||
// The request is abandoned after 10 seconds, and the webhook, which
|
||||
// answers again, is sent the others, in order, a second later.
|
||||
webhook.set(answering)
|
||||
time.Sleep(11 * time.Second)
|
||||
synctest.Wait()
|
||||
|
||||
got := webhook.received()
|
||||
if len(got) != alerts.QueueSize+1 ||
|
||||
got[0].alert["netblock"] != netblock(0).String() {
|
||||
t.Fatalf("the webhook had %d requests, want %d, the first for %s",
|
||||
len(got), alerts.QueueSize+1, netblock(0))
|
||||
}
|
||||
|
||||
for i, request := range got[1:] {
|
||||
if request.alert["netblock"] != netblock(i+1).String() {
|
||||
t.Fatalf("request %d is for %v, want %s", i+1, request.alert["netblock"],
|
||||
netblock(i+1))
|
||||
}
|
||||
}
|
||||
|
||||
wantCounts(t, q, alerts.QueueSize, 1, 0, 1)
|
||||
})
|
||||
}
|
||||
|
||||
func TestStateLoadedIntoANewQueueCarriesOn(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
params := newParams()
|
||||
params.MaxPerHour = 1
|
||||
before := alerts.New(params)
|
||||
|
||||
// Not sent: Run is not running. The repeat is held back by the
|
||||
// cooldown, and the file error past the hourly limit.
|
||||
before.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
before.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
before.Raise(alerts.Alert{Event: alerts.EventFileError})
|
||||
|
||||
time.Sleep(time.Minute)
|
||||
|
||||
webhook, after := start(t, params)
|
||||
after.Load(roundTrip(t, before.Snapshot()))
|
||||
|
||||
// The new queue sends the alert waiting, holds back the repeat as
|
||||
// the cooldown still runs, and sends the summary of the hour.
|
||||
after.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
synctest.Wait()
|
||||
wantEvents(t, webhook, alerts.EventBan)
|
||||
|
||||
time.Sleep(time.Hour)
|
||||
synctest.Wait()
|
||||
wantEvents(t, webhook, alerts.EventBan, alerts.EventSummary)
|
||||
|
||||
detail, _ := webhook.received()[1].alert["detail"].(map[string]any)
|
||||
if detail["count"] != float64(1) {
|
||||
t.Errorf("the summary counts %v alerts, want 1", detail["count"])
|
||||
}
|
||||
|
||||
// The cooldown has run out, and the next one gives both repeats.
|
||||
after.Raise(alerts.Alert{Event: alerts.EventBan, Netblock: netblock(1)})
|
||||
synctest.Wait()
|
||||
|
||||
got := webhook.received()
|
||||
if repeats := got[len(got)-1].alert["suppressed_repeats"]; repeats != float64(2) {
|
||||
t.Errorf("the last alert gives %v repeats, want 2", repeats)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// How the stand-in for the webhook answers: with a status, or, hanging,
|
||||
// not at all, until the request is abandoned.
|
||||
const (
|
||||
answering = http.StatusNoContent
|
||||
failing = http.StatusServiceUnavailable
|
||||
refusing = http.StatusBadRequest
|
||||
hanging = 0
|
||||
)
|
||||
|
||||
// standIn is a stand-in for the webhook. It notes each request it is
|
||||
// sent.
|
||||
type standIn struct {
|
||||
mu sync.Mutex
|
||||
answers int
|
||||
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.
|
||||
type post struct {
|
||||
at time.Time
|
||||
method string
|
||||
url string
|
||||
header http.Header
|
||||
alert map[string]any
|
||||
answered bool
|
||||
}
|
||||
|
||||
// RoundTrip has the stand-in answer req, in place of the network. A
|
||||
// request abandoned before the stand-in answers fails, as over the
|
||||
// network.
|
||||
func (s *standIn) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||
answer := httptest.NewRecorder()
|
||||
s.ServeHTTP(answer, req)
|
||||
|
||||
_ = req.Body.Close()
|
||||
|
||||
err := req.Context().Err()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return answer.Result(), nil
|
||||
}
|
||||
|
||||
// ServeHTTP notes the request, and answers it as the stand-in is set to.
|
||||
func (s *standIn) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
body, _ := io.ReadAll(r.Body)
|
||||
|
||||
var alert map[string]any
|
||||
|
||||
_ = json.Unmarshal(body, &alert)
|
||||
|
||||
s.mu.Lock()
|
||||
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,
|
||||
})
|
||||
s.mu.Unlock()
|
||||
|
||||
if answers == hanging {
|
||||
<-r.Context().Done()
|
||||
} else {
|
||||
w.WriteHeader(answers)
|
||||
}
|
||||
}
|
||||
|
||||
// set sets how the stand-in answers: with the status answers, or hanging.
|
||||
func (s *standIn) set(answers int) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
s.answers = answers
|
||||
}
|
||||
|
||||
// received returns the requests the stand-in has been sent so far.
|
||||
func (s *standIn) received() []post {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
return slices.Clone(s.requests)
|
||||
}
|
||||
|
||||
// lockedBuffer is a buffer the process log can write to while the test
|
||||
// reads it.
|
||||
type lockedBuffer struct {
|
||||
mu sync.Mutex
|
||||
buf bytes.Buffer
|
||||
}
|
||||
|
||||
// Write adds p to the buffer.
|
||||
func (b *lockedBuffer) Write(p []byte) (int, error) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
|
||||
return b.buf.Write(p)
|
||||
}
|
||||
|
||||
// String returns what was written.
|
||||
func (b *lockedBuffer) String() string {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
|
||||
return b.buf.String()
|
||||
}
|
||||
|
||||
// newParams returns the Params of most tests: the webhook at webhookURL,
|
||||
// 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,
|
||||
Events: alerts.Events(),
|
||||
Cooldown: cooldown,
|
||||
MaxPerHour: 60,
|
||||
Instance: instance,
|
||||
Now: func() time.Time { return time.Now().UTC() },
|
||||
ProcessLog: slog.New(slog.DiscardHandler),
|
||||
}
|
||||
}
|
||||
|
||||
// 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}
|
||||
q := alerts.New(params)
|
||||
q.SetTransport(webhook)
|
||||
|
||||
ctx, stop := context.WithCancel(t.Context())
|
||||
stopped := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
q.Run(ctx)
|
||||
close(stopped)
|
||||
}()
|
||||
|
||||
t.Cleanup(func() {
|
||||
stop()
|
||||
<-stopped
|
||||
})
|
||||
|
||||
return webhook, q
|
||||
}
|
||||
|
||||
// midnight is when each test starts.
|
||||
func midnight() time.Time {
|
||||
return time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC)
|
||||
}
|
||||
|
||||
// netblock returns the n-th netblock of a test, counted from 0.
|
||||
func netblock(n int) netip.Prefix {
|
||||
return netip.MustParsePrefix(fmt.Sprintf("203.0.%d.%d/32", 113+n/256, n%256))
|
||||
}
|
||||
|
||||
// roundTrip returns state once written as JSON and read back, as
|
||||
// alerts.json carries it from one start to the next.
|
||||
func roundTrip(t *testing.T, state alerts.State) alerts.State {
|
||||
t.Helper()
|
||||
|
||||
data, err := json.Marshal(state)
|
||||
if err != nil {
|
||||
t.Fatalf("encode: %v", err)
|
||||
}
|
||||
|
||||
var read alerts.State
|
||||
|
||||
err = json.Unmarshal(data, &read)
|
||||
if err != nil {
|
||||
t.Fatalf("decode: %v", err)
|
||||
}
|
||||
|
||||
return read
|
||||
}
|
||||
|
||||
// wantAlert checks every field of an alert the webhook was sent.
|
||||
func wantAlert(t *testing.T, got, want map[string]any) {
|
||||
t.Helper()
|
||||
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("alert %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// wantEvents checks the events of the alerts the webhook was sent, in
|
||||
// order.
|
||||
func wantEvents(t *testing.T, webhook *standIn, want ...string) {
|
||||
t.Helper()
|
||||
|
||||
got := make([]string, 0, len(webhook.received()))
|
||||
|
||||
for _, request := range webhook.received() {
|
||||
event, _ := request.alert["event"].(string)
|
||||
got = append(got, event)
|
||||
}
|
||||
|
||||
if !slices.Equal(got, want) {
|
||||
t.Errorf("the webhook was sent %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// 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.
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package alerts
|
||||
|
||||
import "net/http"
|
||||
|
||||
// QueueSize is the most alerts that wait to be sent.
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user