check / check (push) Waiting to run
SWWAF_BLOCKLIST_URLS names lists of addresses and netblocks, fetched every SWWAF_BLOCKLIST_REFRESH (24h, never under 1h); an IPv4-mapped line stands for its IPv4 address or netblock. reputation.json keeps each list's last try, failed or not, even one cut off by a stop, which a restart waits on as a running instance does, and its last good copy, whole, used while a fetch fails. SWWAF_BLOCKLIST_ACTION denies, limits or only logs a listed client; the log line names the lists, each raises reputation_hit, and a failed fetch raises source_failure. SWWAF_ASN_LIMIT_PERCENT_URL is fetched the same way and counts as SWWAF_ASN_LIMIT_PERCENT does, the lower winning. Judgement call: a failed fetch is retried after the refresh, not sooner. Not done: ban notes do not name the lists yet. Model: opus-5-5
929 lines
28 KiB
Go
929 lines
28 KiB
Go
// Package alerts sends alerts on bans, on traffic over an anomaly
|
|
// threshold, on a source that fails and on a file with an error to each
|
|
// destination set: to the webhook
|
|
// SWWAF_ALERT_WEBHOOK_URL names, each as one JSON object, as the "Alert
|
|
// 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 (
|
|
"bytes"
|
|
"cmp"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"maps"
|
|
"net/http"
|
|
"net/netip"
|
|
"net/url"
|
|
"slices"
|
|
"strings"
|
|
"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"
|
|
// EventAnomaly is a count of requests or bytes over an anomaly
|
|
// threshold.
|
|
EventAnomaly = "anomaly"
|
|
// EventWAFBlock comes with the Core Rule Set; nothing raises it yet.
|
|
EventWAFBlock = "waf_block"
|
|
// EventReputationHit is a request whose client a blocklist lists.
|
|
EventReputationHit = "reputation_hit"
|
|
// EventSourceFailure is GeoJS failing or refusing smallwebwaf, or a
|
|
// fetch of a list failing.
|
|
EventSourceFailure = "source_failure"
|
|
// EventFileError is a rule file or state file edited while smallwebwaf
|
|
// runs that does not parse, a replacement of the lookup database that
|
|
// cannot be read, or a state file that cannot be written.
|
|
EventFileError = "file_error"
|
|
// EventSummary is the summary sent as an hour ends: of the alerts held
|
|
// back in it past SWWAF_ALERT_MAX_PER_HOUR, and of the repeats held
|
|
// back by the cooldowns dropped as it ends, which no alert let through
|
|
// has given. It is sent with SWWAF_ALERT_MAX_PER_HOUR off too, for
|
|
// those repeats. SWWAF_ALERT_EVENTS does not name it.
|
|
EventSummary = "summary"
|
|
)
|
|
|
|
// 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,
|
|
}
|
|
}
|
|
|
|
// The destinations alerts are sent to, as the metrics and alerts.json
|
|
// name them.
|
|
const (
|
|
// 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 a destination.
|
|
sendTimeout = 10 * time.Second
|
|
// 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 a destination's answer that is read.
|
|
maxAnswerBytes = 64 << 10
|
|
)
|
|
|
|
var (
|
|
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 destination refused the alert, answering")
|
|
)
|
|
|
|
// 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 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
|
|
// (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 a destination 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,
|
|
// ASName and Country are, for a ban, the client's as the ban's notes give
|
|
// them.
|
|
//
|
|
//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", for a source_failure, its
|
|
// "source", and for an anomaly, its "scope", with the "asn" or the
|
|
// "name" of some scopes, which the cooldown tells repeats by.
|
|
Reason string `json:"reason"`
|
|
Detail map[string]any `json:"detail"`
|
|
// SuppressedRepeats is how many repeats of the alert the cooldown
|
|
// held back since the last one let through. For a summary, it is how
|
|
// many the cooldowns dropped as the hour ended had held back that no
|
|
// alert let through gave.
|
|
SuppressedRepeats int `json:"suppressed_repeats"`
|
|
}
|
|
|
|
// Cooldown is, for an event on a netblock, about a file or a source, or
|
|
// for an anomaly in a scope, 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"`
|
|
Scope string `json:"scope,omitempty"`
|
|
ASN string `json:"asn,omitempty"`
|
|
Name string `json:"name,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
|
|
// 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 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 each destination set, from a queue of the destination's own.
|
|
// It is safe for concurrent use.
|
|
type Queue struct {
|
|
params Params
|
|
// 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, source or scope.
|
|
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, 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, or in the same
|
|
// scope with the same AS number or name, 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
|
|
scope string
|
|
asn string
|
|
name 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)
|
|
scope, _ := alert.Detail["scope"].(string)
|
|
asn, _ := alert.Detail["asn"].(string)
|
|
name, _ := alert.Detail["name"].(string)
|
|
|
|
return cooldownKey{alert.Event, alert.Netblock, file, source, scope, asn, name}
|
|
}
|
|
|
|
// New returns a Queue with no alert yet.
|
|
func New(params Params) *Queue {
|
|
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 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. The next one let through gives that count, unless an hour of
|
|
// the clock ends first after the cooldown has run out: the cooldown is
|
|
// then dropped, and that hour's summary gives the count. Past MaxPerHour
|
|
// alerts let through in the hour under way, an alert is held back for
|
|
// that hour's summary instead, which is sent once the hour has ended; it
|
|
// starts no cooldown. 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 len(q.destinations) == 0 || !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 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 len(q.destinations) == 0 || !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 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 len(q.destinations) == 0 {
|
|
return
|
|
}
|
|
|
|
var sending sync.WaitGroup
|
|
|
|
for _, d := range q.destinations {
|
|
sending.Go(func() { d.run(ctx) })
|
|
}
|
|
|
|
for {
|
|
q.mu.Lock()
|
|
untilHourEnds := q.hour.Start.Add(time.Hour).Sub(q.params.Now())
|
|
q.mu.Unlock()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
sending.Wait()
|
|
|
|
return
|
|
case <-time.After(untilHourEnds):
|
|
q.mu.Lock()
|
|
q.endHour(q.params.Now())
|
|
q.mu.Unlock()
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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{}
|
|
}
|
|
|
|
// 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. No destination is sent such an alert.
|
|
func (q *Queue) Suppressed() int64 {
|
|
return q.suppressed.Load()
|
|
}
|
|
|
|
// Snapshot returns the queue's state, as alerts.json holds it, with the
|
|
// cooldowns sorted by netblock, then by event, file, source, scope, AS
|
|
// number and name.
|
|
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: map[string][]Alert{},
|
|
}
|
|
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),
|
|
cmp.Compare(a.Scope, b.Scope), cmp.Compare(a.ASN, b.ASN),
|
|
cmp.Compare(a.Name, b.Name))
|
|
})
|
|
|
|
for _, d := range q.destinations {
|
|
state.Waiting[d.name] = d.snapshot()
|
|
}
|
|
|
|
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. 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()
|
|
|
|
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,
|
|
cooldown.Scope, cooldown.ASN, cooldown.Name,
|
|
}
|
|
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{}
|
|
}
|
|
|
|
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 {
|
|
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,
|
|
Scope: key.scope, ASN: key.asn, Name: key.name, Sent: now,
|
|
}
|
|
}
|
|
|
|
// endHour ends the hour under way, if now is past it. It drops the
|
|
// cooldowns that have run out, whatever repeats they held back, so that
|
|
// they do not pile up, and queues that hour's summary when alerts were
|
|
// held back in it past MaxPerHour, or when a cooldown dropped had held
|
|
// back repeats, which no alert let through has given: the summary gives
|
|
// them.
|
|
func (q *Queue) endHour(now time.Time) {
|
|
start := now.Truncate(time.Hour)
|
|
if !start.After(q.hour.Start) {
|
|
return
|
|
}
|
|
|
|
repeats := 0
|
|
|
|
for key, cooldown := range q.cooldowns {
|
|
if now.Sub(cooldown.Sent) >= q.params.Cooldown {
|
|
repeats += cooldown.SuppressedRepeats
|
|
|
|
delete(q.cooldowns, key)
|
|
}
|
|
}
|
|
|
|
heldBack := 0
|
|
for _, count := range q.hour.HeldBack {
|
|
heldBack += count
|
|
}
|
|
|
|
var reasons []string
|
|
|
|
if heldBack > 0 {
|
|
reasons = append(reasons, fmt.Sprintf("%d alerts held back in the hour from %s, "+
|
|
"past the %d an hour SWWAF_ALERT_MAX_PER_HOUR allows", heldBack,
|
|
q.hour.Start.Format(time.RFC3339), q.params.MaxPerHour))
|
|
}
|
|
|
|
if repeats > 0 {
|
|
reasons = append(reasons, fmt.Sprintf("%d repeats held back by "+
|
|
"SWWAF_ALERT_COOLDOWN that no later alert gives", repeats))
|
|
}
|
|
|
|
if len(reasons) > 0 {
|
|
q.queue(&Alert{
|
|
Instance: q.params.Instance,
|
|
Time: now,
|
|
Event: EventSummary,
|
|
Reason: strings.Join(reasons, "; "),
|
|
Detail: map[string]any{
|
|
"hour": q.hour.Start, "count": heldBack, "events": q.hour.HeldBack,
|
|
},
|
|
SuppressedRepeats: repeats,
|
|
})
|
|
}
|
|
|
|
q.hour = Hour{Start: start, HeldBack: map[string]int{}}
|
|
}
|
|
|
|
// queue adds alert to the alerts waiting for each destination.
|
|
func (q *Queue) queue(alert *Alert) {
|
|
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)
|
|
}
|
|
|
|
d.waiting = append(d.waiting, alert)
|
|
|
|
select {
|
|
case d.queued <- struct{}{}:
|
|
default: // a value waits already
|
|
}
|
|
}
|
|
|
|
// 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()
|
|
|
|
for _, alert := range waiting {
|
|
d.add(&alert)
|
|
}
|
|
}
|
|
|
|
// 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 (d *destination) remove(alert *Alert) {
|
|
d.mu.Lock()
|
|
defer d.mu.Unlock()
|
|
|
|
if len(d.waiting) > 0 && d.waiting[0] == alert {
|
|
d.waiting = slices.Delete(d.waiting, 0, 1)
|
|
}
|
|
}
|
|
|
|
// 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 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 err
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(ctx, sendTimeout)
|
|
defer cancel()
|
|
|
|
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, header)
|
|
|
|
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 {
|
|
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)
|
|
}
|
|
}
|
|
|
|
// 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("&", "&", "<", "<", ">", ">").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")
|
|
}
|