check / check (push) Waiting to run
Process log lines carry instance, as request lines do; the instance name is read before the other settings, so the line saying a setting is invalid carries it too. Every metric, Go's and the process's included, carries the label instance, set once on the registry. README.md says so, and that Prometheus keeps it as exported_instance unless the scrape sets honor_labels. An instance name that is not valid UTF-8 stops the start, as the metrics library panics on such a label. Tests that read metrics expect the label; one helper replaces the alert tests' loops that wait for them. Judgement call: the label is named instance, as in the log lines and alerts, although Prometheus gives each target a label of that name. Model: opus-5-5
318 lines
8.9 KiB
Go
318 lines
8.9 KiB
Go
// Package smallwebwaf runs the smallwebwaf process: it reads the settings,
|
|
// the rule files and the state files, serves requests until it is told to
|
|
// stop, and then stops in an orderly way, writing the state files.
|
|
package smallwebwaf
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"io"
|
|
"log/slog"
|
|
"net"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
"time"
|
|
|
|
"sneak.berlin/go/smallwebwaf/internal/alerts"
|
|
"sneak.berlin/go/smallwebwaf/internal/config"
|
|
"sneak.berlin/go/smallwebwaf/internal/lookup"
|
|
"sneak.berlin/go/smallwebwaf/internal/proxy"
|
|
"sneak.berlin/go/smallwebwaf/internal/remotelog"
|
|
"sneak.berlin/go/smallwebwaf/internal/requestlog"
|
|
"sneak.berlin/go/smallwebwaf/internal/rules"
|
|
"sneak.berlin/go/smallwebwaf/internal/state"
|
|
)
|
|
|
|
// shutdownTimeout is how long requests in progress may take to finish
|
|
// once smallwebwaf is told to stop, before their connections are closed.
|
|
// runit and docker wait a little longer before they kill the process.
|
|
const shutdownTimeout = 5 * time.Second
|
|
|
|
// remoteLogStopTimeout is how long, as smallwebwaf stops, the log lines
|
|
// still waiting are sent to SWWAF_LOG_REMOTE_URL before they are given
|
|
// up. stdout has carried them.
|
|
const remoteLogStopTimeout = 2 * time.Second
|
|
|
|
// Params are what Run needs from the process.
|
|
type Params struct {
|
|
// Version is the version of the binary, set when it is built.
|
|
Version string
|
|
// LookupEnv reads an environment variable, normally os.LookupEnv.
|
|
LookupEnv func(string) (string, bool)
|
|
// Stdout receives the request log and the process's own messages.
|
|
Stdout io.Writer
|
|
}
|
|
|
|
// Main runs smallwebwaf until SIGTERM or SIGINT, and returns the
|
|
// process's exit status. Run as `smallwebwaf healthcheck`, it is the
|
|
// container's health check instead.
|
|
func Main(version string) int {
|
|
if len(os.Args) > 1 && os.Args[1] == "healthcheck" {
|
|
return HealthCheck(context.Background(), os.Args[2:], os.LookupEnv, os.Stderr)
|
|
}
|
|
|
|
ctx, stop := signal.NotifyContext(context.Background(),
|
|
syscall.SIGTERM, os.Interrupt)
|
|
defer stop()
|
|
|
|
return Run(ctx, Params{
|
|
Version: version,
|
|
LookupEnv: os.LookupEnv,
|
|
Stdout: os.Stdout,
|
|
})
|
|
}
|
|
|
|
// Run reads the settings, the rule files and the state files, then serves
|
|
// requests until ctx is done. It returns the process's exit status, 1
|
|
// when smallwebwaf cannot start.
|
|
func Run(ctx context.Context, params Params) int {
|
|
processLog := requestlog.NewProcessLogger(params.Stdout,
|
|
config.InstanceName(params.LookupEnv))
|
|
|
|
cfg, err := config.FromEnvironment(params.LookupEnv)
|
|
if err != nil {
|
|
processLog.Error("invalid setting", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
// While SWWAF_LOG_REMOTE_URL is set, every line on stdout from here on
|
|
// is sent there too.
|
|
stdout := params.Stdout
|
|
|
|
var remote *remotelog.Sender
|
|
|
|
if cfg.LogRemoteURL != nil {
|
|
remote = newRemoteLogSender(cfg)
|
|
stdout = io.MultiWriter(params.Stdout, remote)
|
|
processLog = requestlog.NewProcessLogger(stdout, cfg.InstanceName)
|
|
|
|
stopSending := startSending(ctx, remote, processLog)
|
|
defer stopSending()
|
|
}
|
|
|
|
// The state files and the alerts give times in UTC.
|
|
now := func() time.Time { return time.Now().UTC() }
|
|
|
|
alertQueue := newAlertQueue(cfg, now, processLog)
|
|
|
|
ruleFiles, err := rules.Load(rules.Params{
|
|
Dir: cfg.RulesDir,
|
|
Enabled: cfg.RulesEnabled,
|
|
ProcessLog: processLog,
|
|
Alerts: alertQueue,
|
|
})
|
|
if err != nil {
|
|
processLog.Error("cannot use the rule files", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
server := proxy.New(proxy.Params{
|
|
Config: cfg,
|
|
RequestLog: stdout,
|
|
ProcessLog: processLog,
|
|
GeoJSURL: lookup.URL,
|
|
Now: now,
|
|
Rules: ruleFiles,
|
|
Alerts: alertQueue,
|
|
})
|
|
if remote != nil {
|
|
server.Metrics.AddRemoteLog(remote)
|
|
}
|
|
|
|
server.Metrics.AddAlerts(alertQueue)
|
|
|
|
files, err := loadStateFiles(cfg, server, alertQueue, now, processLog)
|
|
if err != nil {
|
|
processLog.Error("cannot use the state files", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
listener, err := (&net.ListenConfig{}).Listen(ctx, "tcp", cfg.ListenAddr)
|
|
if err != nil {
|
|
processLog.Error("cannot listen on SWWAF_LISTEN_ADDR",
|
|
"error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
processLog.Info("starting",
|
|
"version", params.Version,
|
|
"address", listener.Addr().String(),
|
|
"settings", cfg)
|
|
|
|
return serve(ctx, server.Server, listener, files, ruleFiles, alertQueue, processLog)
|
|
}
|
|
|
|
// loadStateFiles reads the state files into the parts of server and into
|
|
// alertQueue, as state.Load does.
|
|
func loadStateFiles(
|
|
cfg *config.Config, server *proxy.Server, alertQueue *alerts.Queue,
|
|
now func() time.Time, processLog *slog.Logger,
|
|
) (*state.Files, error) {
|
|
return state.Load(state.Params{
|
|
Dir: cfg.StateDir,
|
|
WriteDelay: cfg.StateWriteDelay,
|
|
CounterInterval: cfg.StateCounterInterval,
|
|
Ledger: server.Ledger,
|
|
Limiter: server.Limiter,
|
|
GeoJS: server.GeoJS,
|
|
Alerts: alertQueue,
|
|
Now: now,
|
|
ProcessLog: processLog,
|
|
Metrics: server.Metrics,
|
|
})
|
|
}
|
|
|
|
// newAlertQueue returns the queue of the alerts to the webhook, Slack and
|
|
// ntfy, with the settings for them.
|
|
func newAlertQueue(
|
|
cfg *config.Config, now func() time.Time, processLog *slog.Logger,
|
|
) *alerts.Queue {
|
|
return alerts.New(alerts.Params{
|
|
WebhookURL: cfg.AlertWebhookURL,
|
|
WebhookHeaders: cfg.AlertWebhookHeaders,
|
|
SlackURL: cfg.AlertSlackWebhookURL,
|
|
NtfyURL: cfg.AlertNtfyURL,
|
|
NtfyToken: cfg.AlertNtfyToken,
|
|
Events: cfg.AlertEvents,
|
|
Cooldown: cfg.AlertCooldown,
|
|
MaxPerHour: cfg.AlertMaxPerHour,
|
|
Instance: cfg.InstanceName,
|
|
Now: now,
|
|
ProcessLog: processLog,
|
|
})
|
|
}
|
|
|
|
// newRemoteLogSender returns a sender of the log lines to
|
|
// SWWAF_LOG_REMOTE_URL, with the settings for it.
|
|
func newRemoteLogSender(cfg *config.Config) *remotelog.Sender {
|
|
return remotelog.New(remotelog.Params{
|
|
URL: cfg.LogRemoteURL,
|
|
RootCAs: cfg.LogRemoteTLSCAs,
|
|
Buffer: cfg.LogRemoteBuffer,
|
|
Facility: cfg.LogRemoteFacility,
|
|
AppName: cfg.LogRemoteAppName,
|
|
})
|
|
}
|
|
|
|
// startSending runs remote until the function it returns is called, which
|
|
// then waits at most remoteLogStopTimeout for the lines still waiting to
|
|
// be sent. Sending goes on after ctx is done, so that the lines written
|
|
// while smallwebwaf stops are sent too.
|
|
func startSending(
|
|
ctx context.Context, remote *remotelog.Sender, processLog *slog.Logger,
|
|
) func() {
|
|
sending, stop := context.WithCancel(context.WithoutCancel(ctx))
|
|
sent := make(chan struct{})
|
|
|
|
go func() {
|
|
remote.Run(sending, processLog)
|
|
close(sent)
|
|
}()
|
|
|
|
return func() {
|
|
stop()
|
|
|
|
select {
|
|
case <-sent:
|
|
case <-time.After(remoteLogStopTimeout):
|
|
}
|
|
}
|
|
}
|
|
|
|
// serve serves requests on listener, writes the state files as they are
|
|
// due, takes in an admin's edits of them, reads the rule files again as
|
|
// they change, and sends the alerts, until ctx is done. Then it gives the
|
|
// requests in progress shutdownTimeout to finish, and writes every state
|
|
// file, alerts.json with the alerts still waiting.
|
|
func serve(
|
|
ctx context.Context, server *http.Server, listener net.Listener,
|
|
files *state.Files, ruleFiles *rules.Files, alertQueue *alerts.Queue,
|
|
processLog *slog.Logger,
|
|
) int {
|
|
served := make(chan error, 1)
|
|
|
|
go func() {
|
|
served <- server.Serve(listener)
|
|
}()
|
|
|
|
writing, stopWriting := context.WithCancel(ctx)
|
|
defer stopWriting()
|
|
|
|
written := inBackground(func() { files.Run(writing) })
|
|
watched := inBackground(func() { files.Watch(writing) })
|
|
rulesWatched := inBackground(func() { ruleFiles.Watch(writing) })
|
|
alertsSent := inBackground(func() { alertQueue.Run(writing) })
|
|
|
|
select {
|
|
case err := <-served:
|
|
processLog.Error("serving failed", "error", err.Error())
|
|
|
|
return 1
|
|
case <-ctx.Done():
|
|
}
|
|
|
|
processLog.Info("stopping")
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx),
|
|
shutdownTimeout)
|
|
defer cancel()
|
|
|
|
err := server.Shutdown(shutdownCtx)
|
|
if err != nil {
|
|
processLog.Warn("requests still in progress were cut off",
|
|
"error", err.Error())
|
|
|
|
_ = server.Close()
|
|
}
|
|
|
|
err = <-served
|
|
if !errors.Is(err, http.ErrServerClosed) {
|
|
processLog.Error("serving failed", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
// Run and Watch have ended, so nothing else reads or writes the
|
|
// files, and no alert is being sent, so that alerts.json keeps every
|
|
// alert not yet sent. Every request has ended too, but for two kinds
|
|
// that Go's server does not wait for: one cut off because Shutdown
|
|
// timed out, and one whose connection switched protocols, such as a
|
|
// WebSocket. Such a request adds to its client's history only as it
|
|
// ends, which can be after this write, and then that request is
|
|
// missing from clients.json.
|
|
<-written
|
|
<-watched
|
|
<-rulesWatched
|
|
<-alertsSent
|
|
|
|
err = files.WriteAll()
|
|
if err != nil {
|
|
processLog.Error("writing the state files failed", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
processLog.Info("stopped")
|
|
|
|
return 0
|
|
}
|
|
|
|
// inBackground runs task on a goroutine of its own, and returns a channel
|
|
// that is closed once task has returned.
|
|
func inBackground(task func()) <-chan struct{} {
|
|
done := make(chan struct{})
|
|
|
|
go func() {
|
|
task()
|
|
close(done)
|
|
}()
|
|
|
|
return done
|
|
}
|