Alerts to Slack and ntfy, each destination with its own queue (closes #90)
check / check (push) Waiting to run
check / check (push) Waiting to run
Each alert is posted as a message to the Slack incoming webhook SWWAF_ALERT_SLACK_WEBHOOK_URL names, and published to the ntfy topic SWWAF_ALERT_NTFY_URL names, with SWWAF_ALERT_NTFY_TOKEN as a bearer token and a priority and tag by event. The cooldown and the hourly limit stay shared; past them, each destination has its own bounded queue and backoff, and its own sent, failed and dropped counts. alerts.json keeps the alerts waiting by destination; one whose waiting is still a list stops the start, saying what to change. A control character in the ntfy token, or in the instance name ntfy is sent, stops the start. Judgement call: messages also give the detail's file, source, error and mode. Judgement call: alerts_suppressed_total is the same for every destination. Model: opus-5-5
This commit was merged in pull request #94.
This commit is contained in:
+440
-167
@@ -1,12 +1,16 @@
|
||||
// Package alerts sends alerts on bans, on a source that fails and on a
|
||||
// file with an error to the webhook SWWAF_ALERT_WEBHOOK_URL names, each
|
||||
// as one JSON object, as the "Alert webhook schema" section of SPEC.md
|
||||
// describes. A repeat within SWWAF_ALERT_COOLDOWN is held back, and so is
|
||||
// an alert past SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary. The
|
||||
// others wait in a bounded queue, so that a slow or unreachable webhook
|
||||
// never holds up a request. The state is written to alerts.json and read
|
||||
// from it by the state package. Nothing logged names the webhook's URL,
|
||||
// whose path or query can carry a secret.
|
||||
// file with an error to each destination set: to the webhook
|
||||
// SWWAF_ALERT_WEBHOOK_URL names, each as one JSON object, as the "Alert
|
||||
// webhook schema" section of SPEC.md describes, to the Slack incoming
|
||||
// webhook SWWAF_ALERT_SLACK_WEBHOOK_URL names, as a message, and to the
|
||||
// ntfy topic SWWAF_ALERT_NTFY_URL names. A repeat within
|
||||
// SWWAF_ALERT_COOLDOWN is held back, and so is an alert past
|
||||
// SWWAF_ALERT_MAX_PER_HOUR, for the hour's summary. The others wait in a
|
||||
// bounded queue of each destination's own, so that a destination that is
|
||||
// slow or unreachable holds up neither the others nor any request. The
|
||||
// state is written to alerts.json and read from it by the state package.
|
||||
// Nothing logged names a destination's URL, whose path or query can carry
|
||||
// a secret.
|
||||
package alerts
|
||||
|
||||
import (
|
||||
@@ -23,6 +27,7 @@ import (
|
||||
"net/netip"
|
||||
"net/url"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -60,36 +65,62 @@ func Events() []string {
|
||||
}
|
||||
}
|
||||
|
||||
// The destinations alerts are sent to, as the metrics and alerts.json
|
||||
// name them.
|
||||
const (
|
||||
// queueSize is the most alerts that wait to be sent. Past it, the
|
||||
// oldest is dropped.
|
||||
// DestinationWebhook is the webhook SWWAF_ALERT_WEBHOOK_URL names.
|
||||
DestinationWebhook = "webhook"
|
||||
// DestinationSlack is the Slack incoming webhook
|
||||
// SWWAF_ALERT_SLACK_WEBHOOK_URL names.
|
||||
DestinationSlack = "slack"
|
||||
// DestinationNtfy is the ntfy topic SWWAF_ALERT_NTFY_URL names.
|
||||
DestinationNtfy = "ntfy"
|
||||
)
|
||||
|
||||
// Destinations returns every destination alerts can be sent to.
|
||||
func Destinations() []string {
|
||||
return []string{DestinationWebhook, DestinationSlack, DestinationNtfy}
|
||||
}
|
||||
|
||||
const (
|
||||
// queueSize is the most alerts that wait to be sent to a destination.
|
||||
// Past it, the oldest is dropped.
|
||||
queueSize = 1000
|
||||
// sendTimeout bounds one request to the webhook.
|
||||
// sendTimeout bounds one request to a destination.
|
||||
sendTimeout = 10 * time.Second
|
||||
// After a request to the webhook fails, the alert is sent again a
|
||||
// After a request to a destination fails, the alert is sent again a
|
||||
// second later, and retryDelayFactor times as long after each further
|
||||
// failure in a row, up to a minute.
|
||||
firstRetryDelay = time.Second
|
||||
retryDelayFactor = 2
|
||||
maxRetryDelay = time.Minute
|
||||
// maxAnswerBytes is the most of the webhook's answer that is read.
|
||||
// maxAnswerBytes is the most of a destination's answer that is read.
|
||||
maxAnswerBytes = 64 << 10
|
||||
)
|
||||
|
||||
var (
|
||||
errStatus = errors.New("the webhook answered")
|
||||
// errRefused is a 4xx answer other than 408 and 429: the webhook
|
||||
errStatus = errors.New("the destination answered")
|
||||
// errRefused is a 4xx answer other than 408 and 429: the destination
|
||||
// refuses the alert itself, and would refuse it again.
|
||||
errRefused = errors.New("the webhook refused the alert, answering")
|
||||
errRefused = errors.New("the destination refused the alert, answering")
|
||||
)
|
||||
|
||||
// Params are what New needs.
|
||||
// Params are what New needs. With none of WebhookURL, SlackURL and
|
||||
// NtfyURL set, no alert is sent.
|
||||
type Params struct {
|
||||
// WebhookURL is where each alert is posted (SWWAF_ALERT_WEBHOOK_URL),
|
||||
// nil while it is unset and no alert is sent. WebhookHeaders are sent
|
||||
// with each (SWWAF_ALERT_WEBHOOK_HEADERS).
|
||||
// WebhookURL is where each alert is posted as JSON
|
||||
// (SWWAF_ALERT_WEBHOOK_URL), nil while it is unset. WebhookHeaders
|
||||
// are sent with each (SWWAF_ALERT_WEBHOOK_HEADERS).
|
||||
WebhookURL *url.URL
|
||||
WebhookHeaders http.Header
|
||||
// SlackURL is the Slack incoming webhook each alert is posted to as a
|
||||
// message (SWWAF_ALERT_SLACK_WEBHOOK_URL), nil while it is unset.
|
||||
SlackURL *url.URL
|
||||
// NtfyURL is the ntfy topic each alert is published to
|
||||
// (SWWAF_ALERT_NTFY_URL), nil while it is unset. NtfyToken, unless
|
||||
// empty, is sent with each as a bearer token (SWWAF_ALERT_NTFY_TOKEN).
|
||||
NtfyURL *url.URL
|
||||
NtfyToken string
|
||||
// Events are the events alerts are sent for (SWWAF_ALERT_EVENTS).
|
||||
Events []string
|
||||
// Cooldown is how long a repeat of an alert is held back
|
||||
@@ -101,7 +132,7 @@ type Params struct {
|
||||
Instance string
|
||||
// Now tells the time of an alert, normally time.Now in UTC.
|
||||
Now func() time.Time
|
||||
// ProcessLog receives the requests to the webhook that fail.
|
||||
// ProcessLog receives the requests to a destination that fail.
|
||||
ProcessLog *slog.Logger
|
||||
}
|
||||
|
||||
@@ -155,32 +186,63 @@ type Hour struct {
|
||||
}
|
||||
|
||||
// State is what alerts.json holds: the cooldowns, the hour under way, and
|
||||
// the alerts waiting to be sent, oldest first.
|
||||
// for each destination set, the alerts waiting to be sent to it, oldest
|
||||
// first.
|
||||
type State struct {
|
||||
Cooldowns []Cooldown `json:"cooldowns"`
|
||||
Hour Hour `json:"hour"`
|
||||
Waiting []Alert `json:"waiting"`
|
||||
Cooldowns []Cooldown `json:"cooldowns"`
|
||||
Hour Hour `json:"hour"`
|
||||
Waiting map[string][]Alert `json:"waiting"`
|
||||
}
|
||||
|
||||
// Counts are, for a destination, how many alerts it took, how many
|
||||
// requests to it failed, and how many alerts were dropped from its full
|
||||
// queue or given up as it refused them.
|
||||
type Counts struct {
|
||||
Sent, Failed, Dropped int64
|
||||
}
|
||||
|
||||
// Queue takes the alerts raised, holds back those it must, and sends the
|
||||
// others to the webhook. It is safe for concurrent use.
|
||||
// others to each destination set, from a queue of the destination's own.
|
||||
// It is safe for concurrent use.
|
||||
type Queue struct {
|
||||
params Params
|
||||
// httpClient follows no redirect: a redirect is a failure.
|
||||
httpClient *http.Client
|
||||
// queued receives a value when an alert joins the queue, unless one
|
||||
// waits already, so that Run looks at the queue again.
|
||||
queued chan struct{}
|
||||
// destinations are the destinations set, in the order of
|
||||
// Destinations.
|
||||
destinations []*destination
|
||||
|
||||
mu sync.Mutex
|
||||
// cooldowns are the alerts last let through, by event and netblock,
|
||||
// file or source.
|
||||
cooldowns map[cooldownKey]*Cooldown
|
||||
hour Hour
|
||||
|
||||
suppressed atomic.Int64
|
||||
}
|
||||
|
||||
// destination is a destination set, with the alerts waiting to be sent
|
||||
// to it. Its mu is taken after the Queue's, never before.
|
||||
type destination struct {
|
||||
// name is how the metrics and alerts.json name the destination, and
|
||||
// setting the setting that is its URL, which the log names in place
|
||||
// of the URL.
|
||||
name string
|
||||
setting string
|
||||
url *url.URL
|
||||
// message returns the body an alert is posted with, and the headers
|
||||
// sent with it.
|
||||
message func(alert *Alert) ([]byte, http.Header, error)
|
||||
// httpClient follows no redirect: a redirect is a failure.
|
||||
httpClient *http.Client
|
||||
processLog *slog.Logger
|
||||
// queued receives a value when an alert joins the queue, unless one
|
||||
// waits already, so that run looks at the queue again.
|
||||
queued chan struct{}
|
||||
|
||||
mu sync.Mutex
|
||||
// waiting are the alerts waiting to be sent, oldest first.
|
||||
waiting []*Alert
|
||||
|
||||
sent, failed, suppressed, dropped atomic.Int64
|
||||
sent, failed, dropped atomic.Int64
|
||||
}
|
||||
|
||||
// cooldownKey is what makes an alert a repeat of another: the same event
|
||||
@@ -203,32 +265,44 @@ func cooldownKeyOf(alert *Alert) cooldownKey {
|
||||
|
||||
// New returns a Queue with no alert yet.
|
||||
func New(params Params) *Queue {
|
||||
return &Queue{
|
||||
params: params,
|
||||
httpClient: &http.Client{
|
||||
CheckRedirect: func(*http.Request, []*http.Request) error {
|
||||
return http.ErrUseLastResponse
|
||||
},
|
||||
},
|
||||
queued: make(chan struct{}, 1),
|
||||
q := &Queue{
|
||||
params: params,
|
||||
cooldowns: map[cooldownKey]*Cooldown{},
|
||||
hour: Hour{HeldBack: map[string]int{}},
|
||||
}
|
||||
|
||||
if params.WebhookURL != nil {
|
||||
q.addDestination(DestinationWebhook, "SWWAF_ALERT_WEBHOOK_URL",
|
||||
params.WebhookURL, q.webhookMessage)
|
||||
}
|
||||
|
||||
if params.SlackURL != nil {
|
||||
q.addDestination(DestinationSlack, "SWWAF_ALERT_SLACK_WEBHOOK_URL",
|
||||
params.SlackURL, slackMessage)
|
||||
}
|
||||
|
||||
if params.NtfyURL != nil {
|
||||
q.addDestination(DestinationNtfy, "SWWAF_ALERT_NTFY_URL",
|
||||
params.NtfyURL, q.ntfyMessage)
|
||||
}
|
||||
|
||||
return q
|
||||
}
|
||||
|
||||
// Raise sends alert, which names its event and what is particular to it,
|
||||
// unless no webhook is set or SWWAF_ALERT_EVENTS leaves its event out. It
|
||||
// gives alert the instance and the time. An alert that repeats the last
|
||||
// one let through less than Cooldown before is held back and counted,
|
||||
// and the next one let through gives that count. Past MaxPerHour alerts
|
||||
// let through in the hour under way, by the clock, an alert is held back
|
||||
// for that hour's summary instead, which is sent once the hour has ended;
|
||||
// it starts no cooldown, and the repeats held back before it are given by
|
||||
// the next alert let through. Raise never waits: an alert let through
|
||||
// joins the queue, from which Run sends it, and with queueSize alerts
|
||||
// waiting the oldest is dropped.
|
||||
// unless no destination is set or SWWAF_ALERT_EVENTS leaves its event
|
||||
// out. It gives alert the instance and the time. An alert that repeats
|
||||
// the last one let through less than Cooldown before is held back and
|
||||
// counted, and the next one let through gives that count. Past MaxPerHour
|
||||
// alerts let through in the hour under way, by the clock, an alert is
|
||||
// held back for that hour's summary instead, which is sent once the hour
|
||||
// has ended; it starts no cooldown, and the repeats held back before it
|
||||
// are given by the next alert let through. Raise never waits: an alert
|
||||
// let through joins the queue of each destination, from which Run sends
|
||||
// it, and with queueSize alerts waiting for a destination, the oldest is
|
||||
// dropped.
|
||||
func (q *Queue) Raise(alert Alert) {
|
||||
if q.params.WebhookURL == nil || !slices.Contains(q.params.Events, alert.Event) {
|
||||
if len(q.destinations) == 0 || !slices.Contains(q.params.Events, alert.Event) {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -260,12 +334,12 @@ func (q *Queue) Raise(alert Alert) {
|
||||
}
|
||||
|
||||
// WouldSend reports whether Raise would let an alert for event on
|
||||
// netblock through now: a webhook is set, SWWAF_ALERT_EVENTS chooses
|
||||
// netblock through now: a destination is set, SWWAF_ALERT_EVENTS chooses
|
||||
// event, no alert for event on netblock was let through less than
|
||||
// Cooldown before, and fewer than MaxPerHour alerts have been let through
|
||||
// in the hour under way. Unlike Raise, it counts nothing.
|
||||
func (q *Queue) WouldSend(event string, netblock netip.Prefix) bool {
|
||||
if q.params.WebhookURL == nil || !slices.Contains(q.params.Events, event) {
|
||||
if len(q.destinations) == 0 || !slices.Contains(q.params.Events, event) {
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -284,102 +358,70 @@ func (q *Queue) WouldSend(event string, netblock netip.Prefix) bool {
|
||||
return q.params.MaxPerHour == 0 || q.hour.Sent < q.params.MaxPerHour
|
||||
}
|
||||
|
||||
// Run sends the alerts waiting, oldest first, until ctx is done. An alert
|
||||
// stays in the queue until the webhook answers it with a 2xx status, or
|
||||
// refuses it with a 4xx status other than 408 and 429: a refused alert is
|
||||
// logged, counted as dropped, and given up, so that the next is sent. Any
|
||||
// other request that fails is logged, and the alert sent again
|
||||
// firstRetryDelay later, retryDelayFactor times as long after each
|
||||
// further failure in a row, up to maxRetryDelay. Run also ends each hour
|
||||
// as Raise does, so that the hour's summary is sent as it ends. With no
|
||||
// webhook set, it returns at once.
|
||||
// Run sends the alerts waiting to each destination, from its own queue,
|
||||
// as destination.run does, until ctx is done. It also ends each hour as
|
||||
// Raise does, so that the hour's summary is sent as it ends. With no
|
||||
// destination set, it returns at once.
|
||||
func (q *Queue) Run(ctx context.Context) {
|
||||
if q.params.WebhookURL == nil {
|
||||
if len(q.destinations) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
var (
|
||||
retryDelay time.Duration
|
||||
retryAt time.Time
|
||||
)
|
||||
var sending sync.WaitGroup
|
||||
|
||||
for _, d := range q.destinations {
|
||||
sending.Go(func() { d.run(ctx) })
|
||||
}
|
||||
|
||||
for {
|
||||
alert, untilHourEnds := q.next()
|
||||
hourEnds := time.NewTimer(untilHourEnds)
|
||||
|
||||
var due <-chan time.Time // nil while no alert waits
|
||||
if alert != nil {
|
||||
due = time.After(time.Until(retryAt))
|
||||
}
|
||||
q.mu.Lock()
|
||||
untilHourEnds := q.hour.Start.Add(time.Hour).Sub(q.params.Now())
|
||||
q.mu.Unlock()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
hourEnds.Stop()
|
||||
sending.Wait()
|
||||
|
||||
return
|
||||
case <-q.queued:
|
||||
case <-hourEnds.C:
|
||||
case <-time.After(untilHourEnds):
|
||||
q.mu.Lock()
|
||||
q.endHour(q.params.Now())
|
||||
q.mu.Unlock()
|
||||
case <-due:
|
||||
err := q.send(ctx, alert)
|
||||
|
||||
switch {
|
||||
case err == nil:
|
||||
q.remove(alert)
|
||||
q.sent.Add(1)
|
||||
|
||||
retryDelay = 0
|
||||
retryAt = time.Time{}
|
||||
case errors.Is(err, errRefused):
|
||||
q.remove(alert)
|
||||
q.failed.Add(1)
|
||||
q.dropped.Add(1)
|
||||
|
||||
retryDelay = 0
|
||||
retryAt = time.Time{}
|
||||
|
||||
q.params.ProcessLog.Warn("gave up an alert SWWAF_ALERT_WEBHOOK_URL refused",
|
||||
"event", alert.Event, "error", err.Error())
|
||||
case ctx.Err() == nil: // not cut off as smallwebwaf stops
|
||||
q.failed.Add(1)
|
||||
|
||||
retryDelay = min(max(retryDelayFactor*retryDelay, firstRetryDelay),
|
||||
maxRetryDelay)
|
||||
retryAt = time.Now().Add(retryDelay)
|
||||
|
||||
q.params.ProcessLog.Warn("sending an alert to SWWAF_ALERT_WEBHOOK_URL failed",
|
||||
"error", err.Error(), "sending_again_in", retryDelay.String())
|
||||
}
|
||||
}
|
||||
|
||||
hourEnds.Stop()
|
||||
}
|
||||
}
|
||||
|
||||
// Sent is how many alerts the webhook has taken.
|
||||
func (q *Queue) Sent() int64 {
|
||||
return q.sent.Load()
|
||||
// Counts returns the counts of the destination name, all 0 for one not
|
||||
// set.
|
||||
func (q *Queue) Counts(name string) Counts {
|
||||
for _, d := range q.destinations {
|
||||
if d.name == name {
|
||||
return Counts{
|
||||
Sent: d.sent.Load(), Failed: d.failed.Load(), Dropped: d.dropped.Load(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return Counts{}
|
||||
}
|
||||
|
||||
// Failed is how many requests to the webhook have failed.
|
||||
func (q *Queue) Failed() int64 {
|
||||
return q.failed.Load()
|
||||
// DestinationsSet returns the destinations set, in the order of
|
||||
// Destinations.
|
||||
func (q *Queue) DestinationsSet() []string {
|
||||
names := make([]string, 0, len(q.destinations))
|
||||
for _, d := range q.destinations {
|
||||
names = append(names, d.name)
|
||||
}
|
||||
|
||||
return names
|
||||
}
|
||||
|
||||
// Suppressed is how many alerts were held back: by the cooldown, and past
|
||||
// MaxPerHour.
|
||||
// MaxPerHour. No destination is sent such an alert.
|
||||
func (q *Queue) Suppressed() int64 {
|
||||
return q.suppressed.Load()
|
||||
}
|
||||
|
||||
// Dropped is how many alerts were dropped from a full queue, or given up
|
||||
// as the webhook refused them.
|
||||
func (q *Queue) Dropped() int64 {
|
||||
return q.dropped.Load()
|
||||
}
|
||||
|
||||
// Snapshot returns the queue's state, as alerts.json holds it, with the
|
||||
// cooldowns sorted by netblock, then by event, file and source.
|
||||
func (q *Queue) Snapshot() State {
|
||||
@@ -389,7 +431,7 @@ func (q *Queue) Snapshot() State {
|
||||
state := State{
|
||||
Cooldowns: make([]Cooldown, 0, len(q.cooldowns)),
|
||||
Hour: q.hour,
|
||||
Waiting: make([]Alert, 0, len(q.waiting)),
|
||||
Waiting: map[string][]Alert{},
|
||||
}
|
||||
state.Hour.HeldBack = maps.Clone(q.hour.HeldBack)
|
||||
|
||||
@@ -402,8 +444,8 @@ func (q *Queue) Snapshot() State {
|
||||
cmp.Compare(a.File, b.File), cmp.Compare(a.Source, b.Source))
|
||||
})
|
||||
|
||||
for _, alert := range q.waiting {
|
||||
state.Waiting = append(state.Waiting, *alert)
|
||||
for _, d := range q.destinations {
|
||||
state.Waiting[d.name] = d.snapshot()
|
||||
}
|
||||
|
||||
return state
|
||||
@@ -411,8 +453,9 @@ func (q *Queue) Snapshot() State {
|
||||
|
||||
// Load puts state, read from alerts.json, in place of the queue's state.
|
||||
// Each cooldown's netblock is masked to its length, so that
|
||||
// 203.0.113.9/24 is 203.0.113.0/24. Past queueSize alerts waiting, the
|
||||
// oldest are dropped.
|
||||
// 203.0.113.9/24 is 203.0.113.0/24. The alerts waiting for a destination
|
||||
// that is not set are dropped, and so are the oldest past queueSize
|
||||
// alerts waiting for one that is.
|
||||
func (q *Queue) Load(state State) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
@@ -432,13 +475,33 @@ func (q *Queue) Load(state State) {
|
||||
q.hour.HeldBack = map[string]int{}
|
||||
}
|
||||
|
||||
q.waiting = nil
|
||||
|
||||
for _, alert := range state.Waiting {
|
||||
q.queue(&alert)
|
||||
for _, d := range q.destinations {
|
||||
d.load(state.Waiting[d.name])
|
||||
}
|
||||
}
|
||||
|
||||
// addDestination adds a destination: its name, the setting that gives
|
||||
// its URL, that URL, target, and message, which makes the messages sent
|
||||
// to it.
|
||||
func (q *Queue) addDestination(
|
||||
name, setting string, target *url.URL,
|
||||
message func(alert *Alert) ([]byte, http.Header, error),
|
||||
) {
|
||||
q.destinations = append(q.destinations, &destination{
|
||||
name: name,
|
||||
setting: setting,
|
||||
url: target,
|
||||
message: message,
|
||||
httpClient: &http.Client{
|
||||
CheckRedirect: func(*http.Request, []*http.Request) error {
|
||||
return http.ErrUseLastResponse
|
||||
},
|
||||
},
|
||||
processLog: q.params.ProcessLog,
|
||||
queued: make(chan struct{}, 1),
|
||||
})
|
||||
}
|
||||
|
||||
// repeat reports whether alert, raised at now, repeats the last one let
|
||||
// through less than Cooldown before, and counts it if it does.
|
||||
func (q *Queue) repeat(alert *Alert, now time.Time) bool {
|
||||
@@ -515,72 +578,163 @@ func (q *Queue) endHour(now time.Time) {
|
||||
}
|
||||
}
|
||||
|
||||
// queue adds alert to the alerts waiting, first dropping the oldest while
|
||||
// queueSize wait, and has Run look at the queue again.
|
||||
// queue adds alert to the alerts waiting for each destination.
|
||||
func (q *Queue) queue(alert *Alert) {
|
||||
if len(q.waiting) == queueSize {
|
||||
q.waiting = slices.Delete(q.waiting, 0, 1)
|
||||
q.dropped.Add(1)
|
||||
for _, d := range q.destinations {
|
||||
d.add(alert)
|
||||
}
|
||||
}
|
||||
|
||||
// run sends the alerts waiting, oldest first, until ctx is done. An alert
|
||||
// stays in the queue until the destination answers it with a 2xx status,
|
||||
// or refuses it with a 4xx status other than 408 and 429: a refused alert
|
||||
// is logged, counted as dropped, and given up, so that the next is sent.
|
||||
// Any other request that fails is logged, and the alert sent again
|
||||
// firstRetryDelay later, retryDelayFactor times as long after each
|
||||
// further failure in a row, up to maxRetryDelay.
|
||||
func (d *destination) run(ctx context.Context) {
|
||||
var (
|
||||
retryDelay time.Duration
|
||||
retryAt time.Time
|
||||
)
|
||||
|
||||
for {
|
||||
alert := d.oldest()
|
||||
|
||||
var due <-chan time.Time // nil while no alert waits
|
||||
if alert != nil {
|
||||
due = time.After(time.Until(retryAt))
|
||||
}
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-d.queued:
|
||||
case <-due:
|
||||
err := d.send(ctx, alert)
|
||||
|
||||
switch {
|
||||
case err == nil:
|
||||
d.remove(alert)
|
||||
d.sent.Add(1)
|
||||
|
||||
retryDelay = 0
|
||||
retryAt = time.Time{}
|
||||
case errors.Is(err, errRefused):
|
||||
d.remove(alert)
|
||||
d.failed.Add(1)
|
||||
d.dropped.Add(1)
|
||||
|
||||
retryDelay = 0
|
||||
retryAt = time.Time{}
|
||||
|
||||
d.processLog.Warn("gave up an alert "+d.setting+" refused",
|
||||
"event", alert.Event, "error", err.Error())
|
||||
case ctx.Err() == nil: // not cut off as smallwebwaf stops
|
||||
d.failed.Add(1)
|
||||
|
||||
retryDelay = min(max(retryDelayFactor*retryDelay, firstRetryDelay),
|
||||
maxRetryDelay)
|
||||
retryAt = time.Now().Add(retryDelay)
|
||||
|
||||
d.processLog.Warn("sending an alert to "+d.setting+" failed",
|
||||
"error", err.Error(), "sending_again_in", retryDelay.String())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// add adds alert to the alerts waiting, first dropping the oldest while
|
||||
// queueSize wait, and has run look at the queue again.
|
||||
func (d *destination) add(alert *Alert) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
if len(d.waiting) == queueSize {
|
||||
d.waiting = slices.Delete(d.waiting, 0, 1)
|
||||
d.dropped.Add(1)
|
||||
}
|
||||
|
||||
q.waiting = append(q.waiting, alert)
|
||||
d.waiting = append(d.waiting, alert)
|
||||
|
||||
select {
|
||||
case q.queued <- struct{}{}:
|
||||
case d.queued <- struct{}{}:
|
||||
default: // a value waits already
|
||||
}
|
||||
}
|
||||
|
||||
// next returns the oldest alert waiting, nil when none waits, and how
|
||||
// long it is until the hour under way ends.
|
||||
func (q *Queue) next() (*Alert, time.Duration) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
// load puts waiting, read from alerts.json, in place of the alerts
|
||||
// waiting, as add adds them.
|
||||
func (d *destination) load(waiting []Alert) {
|
||||
d.mu.Lock()
|
||||
d.waiting = nil
|
||||
d.mu.Unlock()
|
||||
|
||||
var oldest *Alert
|
||||
if len(q.waiting) > 0 {
|
||||
oldest = q.waiting[0]
|
||||
for _, alert := range waiting {
|
||||
d.add(&alert)
|
||||
}
|
||||
|
||||
return oldest, q.hour.Start.Add(time.Hour).Sub(q.params.Now())
|
||||
}
|
||||
|
||||
// remove takes alert, which Run has sent or given up, out of the queue,
|
||||
// unless it has been dropped from it, or Load has replaced the queue,
|
||||
// since Run took it. Only the oldest alert is ever dropped, so alert is
|
||||
// snapshot returns the alerts waiting, oldest first.
|
||||
func (d *destination) snapshot() []Alert {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
waiting := make([]Alert, 0, len(d.waiting))
|
||||
for _, alert := range d.waiting {
|
||||
waiting = append(waiting, *alert)
|
||||
}
|
||||
|
||||
return waiting
|
||||
}
|
||||
|
||||
// oldest returns the oldest alert waiting, nil when none waits.
|
||||
func (d *destination) oldest() *Alert {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
if len(d.waiting) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
return d.waiting[0]
|
||||
}
|
||||
|
||||
// remove takes alert, which run has sent or given up, out of the queue,
|
||||
// unless it has been dropped from it, or load has replaced the queue,
|
||||
// since run took it. Only the oldest alert is ever dropped, so alert is
|
||||
// the oldest if it is there at all.
|
||||
func (q *Queue) remove(alert *Alert) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
func (d *destination) remove(alert *Alert) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
if len(q.waiting) > 0 && q.waiting[0] == alert {
|
||||
q.waiting = slices.Delete(q.waiting, 0, 1)
|
||||
if len(d.waiting) > 0 && d.waiting[0] == alert {
|
||||
d.waiting = slices.Delete(d.waiting, 0, 1)
|
||||
}
|
||||
}
|
||||
|
||||
// send posts alert to the webhook as JSON, with WebhookHeaders, and
|
||||
// returns an error unless the webhook answers with a 2xx status: one that
|
||||
// send posts alert to the destination, as message makes it, and returns
|
||||
// an error unless the destination answers with a 2xx status: one that
|
||||
// wraps errRefused for a 4xx status other than 408 and 429. No error
|
||||
// names the webhook's URL, whose path or query can carry a secret.
|
||||
func (q *Queue) send(ctx context.Context, alert *Alert) error {
|
||||
body, err := json.Marshal(alert)
|
||||
// names the destination's URL, whose path or query can carry a secret.
|
||||
func (d *destination) send(ctx context.Context, alert *Alert) error {
|
||||
body, header, err := d.message(alert)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode the alert: %w", err)
|
||||
return err
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(ctx, sendTimeout)
|
||||
defer cancel()
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
q.params.WebhookURL.String(), bytes.NewReader(body))
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, d.url.String(),
|
||||
bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return fmt.Errorf("make the request: %w", err)
|
||||
}
|
||||
|
||||
maps.Copy(req.Header, q.params.WebhookHeaders)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
maps.Copy(req.Header, header)
|
||||
|
||||
res, err := q.httpClient.Do(req)
|
||||
res, err := d.httpClient.Do(req)
|
||||
if err != nil {
|
||||
// The client's error names the URL: only what went wrong is kept.
|
||||
if urlErr, ok := errors.AsType[*url.Error](err); ok {
|
||||
@@ -607,3 +761,122 @@ func (q *Queue) send(ctx context.Context, alert *Alert) error {
|
||||
return fmt.Errorf("%w %s", errStatus, res.Status)
|
||||
}
|
||||
}
|
||||
|
||||
// webhookMessage returns alert as JSON, for the webhook, and the headers
|
||||
// sent with it: WebhookHeaders, and its Content-Type.
|
||||
func (q *Queue) webhookMessage(alert *Alert) ([]byte, http.Header, error) {
|
||||
body, err := json.Marshal(alert)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("encode the alert: %w", err)
|
||||
}
|
||||
|
||||
header := http.Header{}
|
||||
maps.Copy(header, q.params.WebhookHeaders)
|
||||
header.Set("Content-Type", "application/json")
|
||||
|
||||
return body, header, nil
|
||||
}
|
||||
|
||||
// slackMessage returns alert as a message for a Slack incoming webhook,
|
||||
// in JSON: its title in bold, then its text, with &, < and > escaped, as
|
||||
// Slack asks, so that nothing in them is read as a link or a mention.
|
||||
func slackMessage(alert *Alert) ([]byte, http.Header, error) {
|
||||
escape := strings.NewReplacer("&", "&", "<", "<", ">", ">").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")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user