// 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) } }